vendor: update all dependencies to latest versions
This commit is contained in:
-159
@@ -1,159 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
// ackBuffer stores the pending ack IDs and notifies the Dirty channel when it becomes non-empty.
|
||||
type ackBuffer struct {
|
||||
Dirty chan struct{}
|
||||
// Close done when ackBuffer is no longer needed.
|
||||
Done chan struct{}
|
||||
|
||||
mu sync.Mutex
|
||||
pending []string
|
||||
send bool
|
||||
}
|
||||
|
||||
// Add adds ackID to the buffer.
|
||||
func (buf *ackBuffer) Add(ackID string) {
|
||||
buf.mu.Lock()
|
||||
defer buf.mu.Unlock()
|
||||
buf.pending = append(buf.pending, ackID)
|
||||
|
||||
// If we are transitioning into a non-empty notification state.
|
||||
if buf.send && len(buf.pending) == 1 {
|
||||
buf.notify()
|
||||
}
|
||||
}
|
||||
|
||||
// RemoveAll removes all ackIDs from the buffer and returns them.
|
||||
func (buf *ackBuffer) RemoveAll() []string {
|
||||
buf.mu.Lock()
|
||||
defer buf.mu.Unlock()
|
||||
|
||||
ret := buf.pending
|
||||
buf.pending = nil
|
||||
return ret
|
||||
}
|
||||
|
||||
// SendNotifications enables sending dirty notification on empty -> non-empty transitions.
|
||||
// If the buffer is already non-empty, a notification will be sent immediately.
|
||||
func (buf *ackBuffer) SendNotifications() {
|
||||
buf.mu.Lock()
|
||||
defer buf.mu.Unlock()
|
||||
|
||||
buf.send = true
|
||||
// If we are transitioning into a non-empty notification state.
|
||||
if len(buf.pending) > 0 {
|
||||
buf.notify()
|
||||
}
|
||||
}
|
||||
|
||||
func (buf *ackBuffer) notify() {
|
||||
go func() {
|
||||
select {
|
||||
case buf.Dirty <- struct{}{}:
|
||||
case <-buf.Done:
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// acker acks messages in batches.
|
||||
type acker struct {
|
||||
s service
|
||||
Ctx context.Context // The context to use when acknowledging messages.
|
||||
Sub string // The full name of the subscription.
|
||||
AckTick <-chan time.Time // AckTick supplies the frequency with which to make ack requests.
|
||||
|
||||
// Notify is called with an ack ID after the message with that ack ID
|
||||
// has been processed. An ackID is considered to have been processed
|
||||
// if at least one attempt has been made to acknowledge it.
|
||||
Notify func(string)
|
||||
|
||||
ackBuffer
|
||||
|
||||
wg sync.WaitGroup
|
||||
done chan struct{}
|
||||
}
|
||||
|
||||
// Start intiates processing of ackIDs which are added via Add.
|
||||
// Notify is called with each ackID once it has been processed.
|
||||
func (a *acker) Start() {
|
||||
a.done = make(chan struct{})
|
||||
a.ackBuffer.Dirty = make(chan struct{})
|
||||
a.ackBuffer.Done = a.done
|
||||
|
||||
a.wg.Add(1)
|
||||
go func() {
|
||||
defer a.wg.Done()
|
||||
for {
|
||||
select {
|
||||
case <-a.ackBuffer.Dirty:
|
||||
a.ack(a.ackBuffer.RemoveAll())
|
||||
case <-a.AckTick:
|
||||
a.ack(a.ackBuffer.RemoveAll())
|
||||
case <-a.done:
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
}()
|
||||
}
|
||||
|
||||
// Ack adds an ack id to be acked in the next batch.
|
||||
func (a *acker) Ack(ackID string) {
|
||||
a.ackBuffer.Add(ackID)
|
||||
}
|
||||
|
||||
// FastMode switches acker into a mode which acks messages as they arrive, rather than waiting
|
||||
// for a.AckTick.
|
||||
func (a *acker) FastMode() {
|
||||
a.ackBuffer.SendNotifications()
|
||||
}
|
||||
|
||||
// Stop drops all pending messages, and releases resources before returning.
|
||||
func (a *acker) Stop() {
|
||||
close(a.done)
|
||||
a.wg.Wait()
|
||||
}
|
||||
|
||||
const maxAckAttempts = 2
|
||||
|
||||
// ack acknowledges the supplied ackIDs.
|
||||
// After the acknowledgement request has completed (regardless of its success
|
||||
// or failure), ids will be passed to a.Notify.
|
||||
func (a *acker) ack(ids []string) {
|
||||
head, tail := a.s.splitAckIDs(ids)
|
||||
for len(head) > 0 {
|
||||
for i := 0; i < maxAckAttempts; i++ {
|
||||
if a.s.acknowledge(a.Ctx, a.Sub, head) == nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
// NOTE: if retry gives up and returns an error, we simply drop
|
||||
// those ack IDs. The messages will be redelivered and this is
|
||||
// a documented behaviour of the API.
|
||||
head, tail = a.s.splitAckIDs(tail)
|
||||
}
|
||||
for _, id := range ids {
|
||||
a.Notify(id)
|
||||
}
|
||||
}
|
||||
-262
@@ -1,262 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"reflect"
|
||||
"sort"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
func TestAcker(t *testing.T) {
|
||||
tick := make(chan time.Time)
|
||||
s := &testService{acknowledgeCalled: make(chan acknowledgeCall)}
|
||||
|
||||
processed := make(chan string, 10)
|
||||
acker := &acker{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
AckTick: tick,
|
||||
Notify: func(ackID string) { processed <- ackID },
|
||||
}
|
||||
acker.Start()
|
||||
|
||||
checkAckProcessed := func(ackIDs []string) {
|
||||
got := <-s.acknowledgeCalled
|
||||
sort.Strings(got.ackIDs)
|
||||
|
||||
want := acknowledgeCall{
|
||||
subName: "subname",
|
||||
ackIDs: ackIDs,
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("acknowledge: got:\n%v\nwant:\n%v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
acker.Ack("a")
|
||||
acker.Ack("b")
|
||||
tick <- time.Time{}
|
||||
checkAckProcessed([]string{"a", "b"})
|
||||
acker.Ack("c")
|
||||
tick <- time.Time{}
|
||||
checkAckProcessed([]string{"c"})
|
||||
acker.Stop()
|
||||
|
||||
// all IDS should have been sent to processed.
|
||||
close(processed)
|
||||
processedIDs := []string{}
|
||||
for id := range processed {
|
||||
processedIDs = append(processedIDs, id)
|
||||
}
|
||||
sort.Strings(processedIDs)
|
||||
want := []string{"a", "b", "c"}
|
||||
if !reflect.DeepEqual(processedIDs, want) {
|
||||
t.Errorf("acker processed: got:\n%v\nwant:\n%v", processedIDs, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAckerFastMode(t *testing.T) {
|
||||
tick := make(chan time.Time)
|
||||
s := &testService{acknowledgeCalled: make(chan acknowledgeCall)}
|
||||
|
||||
processed := make(chan string, 10)
|
||||
acker := &acker{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
AckTick: tick,
|
||||
Notify: func(ackID string) { processed <- ackID },
|
||||
}
|
||||
acker.Start()
|
||||
|
||||
checkAckProcessed := func(ackIDs []string) {
|
||||
got := <-s.acknowledgeCalled
|
||||
sort.Strings(got.ackIDs)
|
||||
|
||||
want := acknowledgeCall{
|
||||
subName: "subname",
|
||||
ackIDs: ackIDs,
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("acknowledge: got:\n%v\nwant:\n%v", got, want)
|
||||
}
|
||||
}
|
||||
// No ticks are sent; fast mode doesn't need them.
|
||||
acker.Ack("a")
|
||||
acker.Ack("b")
|
||||
acker.FastMode()
|
||||
checkAckProcessed([]string{"a", "b"})
|
||||
acker.Ack("c")
|
||||
checkAckProcessed([]string{"c"})
|
||||
acker.Stop()
|
||||
|
||||
// all IDS should have been sent to processed.
|
||||
close(processed)
|
||||
processedIDs := []string{}
|
||||
for id := range processed {
|
||||
processedIDs = append(processedIDs, id)
|
||||
}
|
||||
sort.Strings(processedIDs)
|
||||
want := []string{"a", "b", "c"}
|
||||
if !reflect.DeepEqual(processedIDs, want) {
|
||||
t.Errorf("acker processed: got:\n%v\nwant:\n%v", processedIDs, want)
|
||||
}
|
||||
}
|
||||
|
||||
// TestAckerStop checks that Stop returns immediately.
|
||||
func TestAckerStop(t *testing.T) {
|
||||
tick := make(chan time.Time)
|
||||
s := &testService{acknowledgeCalled: make(chan acknowledgeCall, 10)}
|
||||
|
||||
processed := make(chan string)
|
||||
acker := &acker{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
AckTick: tick,
|
||||
Notify: func(ackID string) { processed <- ackID },
|
||||
}
|
||||
|
||||
acker.Start()
|
||||
|
||||
stopped := make(chan struct{})
|
||||
|
||||
acker.Ack("a")
|
||||
|
||||
go func() {
|
||||
acker.Stop()
|
||||
stopped <- struct{}{}
|
||||
}()
|
||||
|
||||
// Stopped should have been written to by the time this sleep completes.
|
||||
time.Sleep(time.Millisecond)
|
||||
|
||||
// Receiving from processed should cause Stop to subsequently return,
|
||||
// so it should never be possible to read from stopped before
|
||||
// processed.
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-processed:
|
||||
t.Errorf("acker.Stop processed an ack id before returning")
|
||||
case <-time.After(time.Millisecond):
|
||||
t.Errorf("acker.Stop never returned")
|
||||
}
|
||||
}
|
||||
|
||||
type ackCallResult struct {
|
||||
ackIDs []string
|
||||
err error
|
||||
}
|
||||
|
||||
type ackService struct {
|
||||
service
|
||||
|
||||
calls []ackCallResult
|
||||
|
||||
t *testing.T // used for error logging.
|
||||
}
|
||||
|
||||
func (as *ackService) acknowledge(ctx context.Context, subName string, ackIDs []string) error {
|
||||
if len(as.calls) == 0 {
|
||||
as.t.Fatalf("unexpected call to acknowledge: ackIDs: %v", ackIDs)
|
||||
}
|
||||
call := as.calls[0]
|
||||
as.calls = as.calls[1:]
|
||||
|
||||
if got, want := ackIDs, call.ackIDs; !reflect.DeepEqual(got, want) {
|
||||
as.t.Errorf("unexpected arguments to acknowledge: got: %v ; want: %v", got, want)
|
||||
}
|
||||
return call.err
|
||||
}
|
||||
|
||||
// Test implementation returns the first 2 elements as head, and the rest as tail.
|
||||
func (as *ackService) splitAckIDs(ids []string) ([]string, []string) {
|
||||
if len(ids) < 2 {
|
||||
return ids, nil
|
||||
}
|
||||
return ids[:2], ids[2:]
|
||||
}
|
||||
|
||||
func TestAckerSplitsBatches(t *testing.T) {
|
||||
type testCase struct {
|
||||
calls []ackCallResult
|
||||
}
|
||||
for _, tc := range []testCase{
|
||||
{
|
||||
calls: []ackCallResult{
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"e", "f"},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
calls: []ackCallResult{
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// On error we retry once.
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// We give up after failing twice, so we move on to the next set, "c" and "d"
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// Again, we retry once.
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"e", "f"},
|
||||
},
|
||||
},
|
||||
},
|
||||
} {
|
||||
s := &ackService{
|
||||
t: t,
|
||||
calls: tc.calls,
|
||||
}
|
||||
|
||||
acker := &acker{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
Notify: func(string) {},
|
||||
}
|
||||
|
||||
acker.ack([]string{"a", "b", "c", "d", "e", "f"})
|
||||
|
||||
if len(s.calls) != 0 {
|
||||
t.Errorf("expected ack calls did not occur: %v", s.calls)
|
||||
}
|
||||
}
|
||||
}
|
||||
+1
-2
@@ -35,8 +35,7 @@ func insertXGoog(ctx context.Context, val []string) context.Context {
|
||||
return metadata.NewOutgoingContext(ctx, md)
|
||||
}
|
||||
|
||||
// DefaultAuthScopes reports the authentication scopes required
|
||||
// by this package.
|
||||
// DefaultAuthScopes reports the default set of authentication scopes to use with this package.
|
||||
func DefaultAuthScopes() []string {
|
||||
return []string{
|
||||
"https://www.googleapis.com/auth/cloud-platform",
|
||||
|
||||
+152
@@ -75,6 +75,18 @@ func (s *mockPublisherServer) CreateTopic(ctx context.Context, req *pubsubpb.Top
|
||||
return s.resps[0].(*pubsubpb.Topic), nil
|
||||
}
|
||||
|
||||
func (s *mockPublisherServer) UpdateTopic(ctx context.Context, req *pubsubpb.UpdateTopicRequest) (*pubsubpb.Topic, error) {
|
||||
md, _ := metadata.FromIncomingContext(ctx)
|
||||
if xg := md["x-goog-api-client"]; len(xg) == 0 || !strings.Contains(xg[0], "gl-go/") {
|
||||
return nil, fmt.Errorf("x-goog-api-client = %v, expected gl-go key", xg)
|
||||
}
|
||||
s.reqs = append(s.reqs, req)
|
||||
if s.err != nil {
|
||||
return nil, s.err
|
||||
}
|
||||
return s.resps[0].(*pubsubpb.Topic), nil
|
||||
}
|
||||
|
||||
func (s *mockPublisherServer) Publish(ctx context.Context, req *pubsubpb.PublishRequest) (*pubsubpb.PublishResponse, error) {
|
||||
md, _ := metadata.FromIncomingContext(ctx)
|
||||
if xg := md["x-goog-api-client"]; len(xg) == 0 || !strings.Contains(xg[0], "gl-go/") {
|
||||
@@ -358,6 +370,18 @@ func (s *mockSubscriberServer) CreateSnapshot(ctx context.Context, req *pubsubpb
|
||||
return s.resps[0].(*pubsubpb.Snapshot), nil
|
||||
}
|
||||
|
||||
func (s *mockSubscriberServer) UpdateSnapshot(ctx context.Context, req *pubsubpb.UpdateSnapshotRequest) (*pubsubpb.Snapshot, error) {
|
||||
md, _ := metadata.FromIncomingContext(ctx)
|
||||
if xg := md["x-goog-api-client"]; len(xg) == 0 || !strings.Contains(xg[0], "gl-go/") {
|
||||
return nil, fmt.Errorf("x-goog-api-client = %v, expected gl-go key", xg)
|
||||
}
|
||||
s.reqs = append(s.reqs, req)
|
||||
if s.err != nil {
|
||||
return nil, s.err
|
||||
}
|
||||
return s.resps[0].(*pubsubpb.Snapshot), nil
|
||||
}
|
||||
|
||||
func (s *mockSubscriberServer) DeleteSnapshot(ctx context.Context, req *pubsubpb.DeleteSnapshotRequest) (*emptypb.Empty, error) {
|
||||
md, _ := metadata.FromIncomingContext(ctx)
|
||||
if xg := md["x-goog-api-client"]; len(xg) == 0 || !strings.Contains(xg[0], "gl-go/") {
|
||||
@@ -474,6 +498,69 @@ func TestPublisherCreateTopicError(t *testing.T) {
|
||||
}
|
||||
_ = resp
|
||||
}
|
||||
func TestPublisherUpdateTopic(t *testing.T) {
|
||||
var name string = "name3373707"
|
||||
var expectedResponse = &pubsubpb.Topic{
|
||||
Name: name,
|
||||
}
|
||||
|
||||
mockPublisher.err = nil
|
||||
mockPublisher.reqs = nil
|
||||
|
||||
mockPublisher.resps = append(mockPublisher.resps[:0], expectedResponse)
|
||||
|
||||
var topic *pubsubpb.Topic = &pubsubpb.Topic{}
|
||||
var updateMask *field_maskpb.FieldMask = &field_maskpb.FieldMask{}
|
||||
var request = &pubsubpb.UpdateTopicRequest{
|
||||
Topic: topic,
|
||||
UpdateMask: updateMask,
|
||||
}
|
||||
|
||||
c, err := NewPublisherClient(context.Background(), clientOpt)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
resp, err := c.UpdateTopic(context.Background(), request)
|
||||
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if want, got := request, mockPublisher.reqs[0]; !proto.Equal(want, got) {
|
||||
t.Errorf("wrong request %q, want %q", got, want)
|
||||
}
|
||||
|
||||
if want, got := expectedResponse, resp; !proto.Equal(want, got) {
|
||||
t.Errorf("wrong response %q, want %q)", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPublisherUpdateTopicError(t *testing.T) {
|
||||
errCode := codes.PermissionDenied
|
||||
mockPublisher.err = gstatus.Error(errCode, "test error")
|
||||
|
||||
var topic *pubsubpb.Topic = &pubsubpb.Topic{}
|
||||
var updateMask *field_maskpb.FieldMask = &field_maskpb.FieldMask{}
|
||||
var request = &pubsubpb.UpdateTopicRequest{
|
||||
Topic: topic,
|
||||
UpdateMask: updateMask,
|
||||
}
|
||||
|
||||
c, err := NewPublisherClient(context.Background(), clientOpt)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
resp, err := c.UpdateTopic(context.Background(), request)
|
||||
|
||||
if st, ok := gstatus.FromError(err); !ok {
|
||||
t.Errorf("got error %v, expected grpc error", err)
|
||||
} else if c := st.Code(); c != errCode {
|
||||
t.Errorf("got error code %q, want %q", c, errCode)
|
||||
}
|
||||
_ = resp
|
||||
}
|
||||
func TestPublisherPublish(t *testing.T) {
|
||||
var messageIdsElement string = "messageIdsElement-744837059"
|
||||
var messageIds = []string{messageIdsElement}
|
||||
@@ -1581,6 +1668,71 @@ func TestSubscriberCreateSnapshotError(t *testing.T) {
|
||||
}
|
||||
_ = resp
|
||||
}
|
||||
func TestSubscriberUpdateSnapshot(t *testing.T) {
|
||||
var name string = "name3373707"
|
||||
var topic string = "topic110546223"
|
||||
var expectedResponse = &pubsubpb.Snapshot{
|
||||
Name: name,
|
||||
Topic: topic,
|
||||
}
|
||||
|
||||
mockSubscriber.err = nil
|
||||
mockSubscriber.reqs = nil
|
||||
|
||||
mockSubscriber.resps = append(mockSubscriber.resps[:0], expectedResponse)
|
||||
|
||||
var snapshot *pubsubpb.Snapshot = &pubsubpb.Snapshot{}
|
||||
var updateMask *field_maskpb.FieldMask = &field_maskpb.FieldMask{}
|
||||
var request = &pubsubpb.UpdateSnapshotRequest{
|
||||
Snapshot: snapshot,
|
||||
UpdateMask: updateMask,
|
||||
}
|
||||
|
||||
c, err := NewSubscriberClient(context.Background(), clientOpt)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
resp, err := c.UpdateSnapshot(context.Background(), request)
|
||||
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if want, got := request, mockSubscriber.reqs[0]; !proto.Equal(want, got) {
|
||||
t.Errorf("wrong request %q, want %q", got, want)
|
||||
}
|
||||
|
||||
if want, got := expectedResponse, resp; !proto.Equal(want, got) {
|
||||
t.Errorf("wrong response %q, want %q)", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriberUpdateSnapshotError(t *testing.T) {
|
||||
errCode := codes.PermissionDenied
|
||||
mockSubscriber.err = gstatus.Error(errCode, "test error")
|
||||
|
||||
var snapshot *pubsubpb.Snapshot = &pubsubpb.Snapshot{}
|
||||
var updateMask *field_maskpb.FieldMask = &field_maskpb.FieldMask{}
|
||||
var request = &pubsubpb.UpdateSnapshotRequest{
|
||||
Snapshot: snapshot,
|
||||
UpdateMask: updateMask,
|
||||
}
|
||||
|
||||
c, err := NewSubscriberClient(context.Background(), clientOpt)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
resp, err := c.UpdateSnapshot(context.Background(), request)
|
||||
|
||||
if st, ok := gstatus.FromError(err); !ok {
|
||||
t.Errorf("got error %v, expected grpc error", err)
|
||||
} else if c := st.Code(); c != errCode {
|
||||
t.Errorf("got error code %q, want %q", c, errCode)
|
||||
}
|
||||
_ = resp
|
||||
}
|
||||
func TestSubscriberDeleteSnapshot(t *testing.T) {
|
||||
var expectedResponse *emptypb.Empty = &emptypb.Empty{}
|
||||
|
||||
|
||||
+37
-24
@@ -32,14 +32,10 @@ import (
|
||||
"google.golang.org/grpc/codes"
|
||||
)
|
||||
|
||||
var (
|
||||
publisherProjectPathTemplate = gax.MustCompilePathTemplate("projects/{project}")
|
||||
publisherTopicPathTemplate = gax.MustCompilePathTemplate("projects/{project}/topics/{topic}")
|
||||
)
|
||||
|
||||
// PublisherCallOptions contains the retry settings for each method of PublisherClient.
|
||||
type PublisherCallOptions struct {
|
||||
CreateTopic []gax.CallOption
|
||||
UpdateTopic []gax.CallOption
|
||||
Publish []gax.CallOption
|
||||
GetTopic []gax.CallOption
|
||||
ListTopics []gax.CallOption
|
||||
@@ -88,6 +84,7 @@ func defaultPublisherCallOptions() *PublisherCallOptions {
|
||||
}
|
||||
return &PublisherCallOptions{
|
||||
CreateTopic: retry[[2]string{"default", "idempotent"}],
|
||||
UpdateTopic: retry[[2]string{"default", "idempotent"}],
|
||||
Publish: retry[[2]string{"messaging", "one_plus_delivery"}],
|
||||
GetTopic: retry[[2]string{"default", "idempotent"}],
|
||||
ListTopics: retry[[2]string{"default", "idempotent"}],
|
||||
@@ -152,25 +149,20 @@ func (c *PublisherClient) SetGoogleClientInfo(keyval ...string) {
|
||||
|
||||
// PublisherProjectPath returns the path for the project resource.
|
||||
func PublisherProjectPath(project string) string {
|
||||
path, err := publisherProjectPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
""
|
||||
}
|
||||
|
||||
// PublisherTopicPath returns the path for the topic resource.
|
||||
func PublisherTopicPath(project, topic string) string {
|
||||
path, err := publisherTopicPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
"topic": topic,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
"/topics/" +
|
||||
topic +
|
||||
""
|
||||
}
|
||||
|
||||
func (c *PublisherClient) SubscriptionIAM(subscription *pubsubpb.Subscription) *iam.Handle {
|
||||
@@ -197,9 +189,30 @@ func (c *PublisherClient) CreateTopic(ctx context.Context, req *pubsubpb.Topic,
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// Publish adds one or more messages to the topic. Returns `NOT_FOUND` if the topic
|
||||
// UpdateTopic updates an existing topic. Note that certain properties of a topic are not
|
||||
// modifiable. Options settings follow the style guide:
|
||||
// NOTE: The style guide requires body: "topic" instead of body: "*".
|
||||
// Keeping the latter for internal consistency in V1, however it should be
|
||||
// corrected in V2. See
|
||||
// https://cloud.google.com/apis/design/standard_methods#update for details.
|
||||
func (c *PublisherClient) UpdateTopic(ctx context.Context, req *pubsubpb.UpdateTopicRequest, opts ...gax.CallOption) (*pubsubpb.Topic, error) {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.UpdateTopic[0:len(c.CallOptions.UpdateTopic):len(c.CallOptions.UpdateTopic)], opts...)
|
||||
var resp *pubsubpb.Topic
|
||||
err := gax.Invoke(ctx, func(ctx context.Context, settings gax.CallSettings) error {
|
||||
var err error
|
||||
resp, err = c.publisherClient.UpdateTopic(ctx, req, settings.GRPC...)
|
||||
return err
|
||||
}, opts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// Publish adds one or more messages to the topic. Returns NOT_FOUND if the topic
|
||||
// does not exist. The message payload must not be empty; it must contain
|
||||
// either a non-empty data field, or at least one attribute.
|
||||
// either a non-empty data field, or at least one attribute.
|
||||
func (c *PublisherClient) Publish(ctx context.Context, req *pubsubpb.PublishRequest, opts ...gax.CallOption) (*pubsubpb.PublishResponse, error) {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.Publish[0:len(c.CallOptions.Publish):len(c.CallOptions.Publish)], opts...)
|
||||
@@ -301,11 +314,11 @@ func (c *PublisherClient) ListTopicSubscriptions(ctx context.Context, req *pubsu
|
||||
return it
|
||||
}
|
||||
|
||||
// DeleteTopic deletes the topic with the given name. Returns `NOT_FOUND` if the topic
|
||||
// DeleteTopic deletes the topic with the given name. Returns NOT_FOUND if the topic
|
||||
// does not exist. After a topic is deleted, a new topic may be created with
|
||||
// the same name; this is an entirely new topic with none of the old
|
||||
// configuration or subscriptions. Existing subscriptions to this topic are
|
||||
// not deleted, but their `topic` field is set to `_deleted-topic_`.
|
||||
// not deleted, but their topic field is set to _deleted-topic_.
|
||||
func (c *PublisherClient) DeleteTopic(ctx context.Context, req *pubsubpb.DeleteTopicRequest, opts ...gax.CallOption) error {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.DeleteTopic[0:len(c.CallOptions.DeleteTopic):len(c.CallOptions.DeleteTopic)], opts...)
|
||||
|
||||
+18
@@ -85,6 +85,24 @@ func ExamplePublisherClient_CreateTopic() {
|
||||
_ = resp
|
||||
}
|
||||
|
||||
func ExamplePublisherClient_UpdateTopic() {
|
||||
ctx := context.Background()
|
||||
c, err := pubsub.NewPublisherClient(ctx)
|
||||
if err != nil {
|
||||
// TODO: Handle error.
|
||||
}
|
||||
|
||||
req := &pubsubpb.UpdateTopicRequest{
|
||||
// TODO: Fill request struct fields.
|
||||
}
|
||||
resp, err := c.UpdateTopic(ctx, req)
|
||||
if err != nil {
|
||||
// TODO: Handle error.
|
||||
}
|
||||
// TODO: Use resp.
|
||||
_ = resp
|
||||
}
|
||||
|
||||
func ExamplePublisherClient_Publish() {
|
||||
ctx := context.Background()
|
||||
c, err := pubsub.NewPublisherClient(ctx)
|
||||
|
||||
+83
-57
@@ -32,13 +32,6 @@ import (
|
||||
"google.golang.org/grpc/codes"
|
||||
)
|
||||
|
||||
var (
|
||||
subscriberProjectPathTemplate = gax.MustCompilePathTemplate("projects/{project}")
|
||||
subscriberSnapshotPathTemplate = gax.MustCompilePathTemplate("projects/{project}/snapshots/{snapshot}")
|
||||
subscriberSubscriptionPathTemplate = gax.MustCompilePathTemplate("projects/{project}/subscriptions/{subscription}")
|
||||
subscriberTopicPathTemplate = gax.MustCompilePathTemplate("projects/{project}/topics/{topic}")
|
||||
)
|
||||
|
||||
// SubscriberCallOptions contains the retry settings for each method of SubscriberClient.
|
||||
type SubscriberCallOptions struct {
|
||||
CreateSubscription []gax.CallOption
|
||||
@@ -53,6 +46,7 @@ type SubscriberCallOptions struct {
|
||||
ModifyPushConfig []gax.CallOption
|
||||
ListSnapshots []gax.CallOption
|
||||
CreateSnapshot []gax.CallOption
|
||||
UpdateSnapshot []gax.CallOption
|
||||
DeleteSnapshot []gax.CallOption
|
||||
Seek []gax.CallOption
|
||||
}
|
||||
@@ -93,6 +87,21 @@ func defaultSubscriberCallOptions() *SubscriberCallOptions {
|
||||
})
|
||||
}),
|
||||
},
|
||||
{"streaming_messaging", "pull"}: {
|
||||
gax.WithRetry(func() gax.Retryer {
|
||||
return gax.OnCodes([]codes.Code{
|
||||
codes.Canceled,
|
||||
codes.DeadlineExceeded,
|
||||
codes.ResourceExhausted,
|
||||
codes.Internal,
|
||||
codes.Unavailable,
|
||||
}, gax.Backoff{
|
||||
Initial: 100 * time.Millisecond,
|
||||
Max: 60000 * time.Millisecond,
|
||||
Multiplier: 1.3,
|
||||
})
|
||||
}),
|
||||
},
|
||||
}
|
||||
return &SubscriberCallOptions{
|
||||
CreateSubscription: retry[[2]string{"default", "idempotent"}],
|
||||
@@ -103,10 +112,11 @@ func defaultSubscriberCallOptions() *SubscriberCallOptions {
|
||||
ModifyAckDeadline: retry[[2]string{"default", "non_idempotent"}],
|
||||
Acknowledge: retry[[2]string{"messaging", "non_idempotent"}],
|
||||
Pull: retry[[2]string{"messaging", "pull"}],
|
||||
StreamingPull: retry[[2]string{"messaging", "pull"}],
|
||||
StreamingPull: retry[[2]string{"streaming_messaging", "pull"}],
|
||||
ModifyPushConfig: retry[[2]string{"default", "non_idempotent"}],
|
||||
ListSnapshots: retry[[2]string{"default", "idempotent"}],
|
||||
CreateSnapshot: retry[[2]string{"default", "idempotent"}],
|
||||
UpdateSnapshot: retry[[2]string{"default", "idempotent"}],
|
||||
DeleteSnapshot: retry[[2]string{"default", "idempotent"}],
|
||||
Seek: retry[[2]string{"default", "non_idempotent"}],
|
||||
}
|
||||
@@ -130,7 +140,7 @@ type SubscriberClient struct {
|
||||
// NewSubscriberClient creates a new subscriber client.
|
||||
//
|
||||
// The service that an application uses to manipulate subscriptions and to
|
||||
// consume messages from a subscription via the `Pull` method.
|
||||
// consume messages from a subscription via the Pull method.
|
||||
func NewSubscriberClient(ctx context.Context, opts ...option.ClientOption) (*SubscriberClient, error) {
|
||||
conn, err := transport.DialGRPC(ctx, append(defaultSubscriberClientOptions(), opts...)...)
|
||||
if err != nil {
|
||||
@@ -168,49 +178,40 @@ func (c *SubscriberClient) SetGoogleClientInfo(keyval ...string) {
|
||||
|
||||
// SubscriberProjectPath returns the path for the project resource.
|
||||
func SubscriberProjectPath(project string) string {
|
||||
path, err := subscriberProjectPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
""
|
||||
}
|
||||
|
||||
// SubscriberSnapshotPath returns the path for the snapshot resource.
|
||||
func SubscriberSnapshotPath(project, snapshot string) string {
|
||||
path, err := subscriberSnapshotPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
"snapshot": snapshot,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
"/snapshots/" +
|
||||
snapshot +
|
||||
""
|
||||
}
|
||||
|
||||
// SubscriberSubscriptionPath returns the path for the subscription resource.
|
||||
func SubscriberSubscriptionPath(project, subscription string) string {
|
||||
path, err := subscriberSubscriptionPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
"subscription": subscription,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
"/subscriptions/" +
|
||||
subscription +
|
||||
""
|
||||
}
|
||||
|
||||
// SubscriberTopicPath returns the path for the topic resource.
|
||||
func SubscriberTopicPath(project, topic string) string {
|
||||
path, err := subscriberTopicPathTemplate.Render(map[string]string{
|
||||
"project": project,
|
||||
"topic": topic,
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return path
|
||||
return "" +
|
||||
"projects/" +
|
||||
project +
|
||||
"/topics/" +
|
||||
topic +
|
||||
""
|
||||
}
|
||||
|
||||
func (c *SubscriberClient) SubscriptionIAM(subscription *pubsubpb.Subscription) *iam.Handle {
|
||||
@@ -222,13 +223,13 @@ func (c *SubscriberClient) TopicIAM(topic *pubsubpb.Topic) *iam.Handle {
|
||||
}
|
||||
|
||||
// CreateSubscription creates a subscription to a given topic.
|
||||
// If the subscription already exists, returns `ALREADY_EXISTS`.
|
||||
// If the corresponding topic doesn't exist, returns `NOT_FOUND`.
|
||||
// If the subscription already exists, returns ALREADY_EXISTS.
|
||||
// If the corresponding topic doesn't exist, returns NOT_FOUND.
|
||||
//
|
||||
// If the name is not provided in the request, the server will assign a random
|
||||
// name for this subscription on the same project as the topic, conforming
|
||||
// to the
|
||||
// [resource name format](https://cloud.google.com/pubsub/docs/overview#names).
|
||||
// resource name format (at https://cloud.google.com/pubsub/docs/overview#names).
|
||||
// The generated name is populated in the returned Subscription object.
|
||||
// Note that for REST API requests, you must specify a name in the request.
|
||||
func (c *SubscriberClient) CreateSubscription(ctx context.Context, req *pubsubpb.Subscription, opts ...gax.CallOption) (*pubsubpb.Subscription, error) {
|
||||
@@ -264,6 +265,10 @@ func (c *SubscriberClient) GetSubscription(ctx context.Context, req *pubsubpb.Ge
|
||||
|
||||
// UpdateSubscription updates an existing subscription. Note that certain properties of a
|
||||
// subscription, such as its topic, are not modifiable.
|
||||
// NOTE: The style guide requires body: "subscription" instead of body: "*".
|
||||
// Keeping the latter for internal consistency in V1, however it should be
|
||||
// corrected in V2. See
|
||||
// https://cloud.google.com/apis/design/standard_methods#update for details.
|
||||
func (c *SubscriberClient) UpdateSubscription(ctx context.Context, req *pubsubpb.UpdateSubscriptionRequest, opts ...gax.CallOption) (*pubsubpb.Subscription, error) {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.UpdateSubscription[0:len(c.CallOptions.UpdateSubscription):len(c.CallOptions.UpdateSubscription)], opts...)
|
||||
@@ -315,8 +320,8 @@ func (c *SubscriberClient) ListSubscriptions(ctx context.Context, req *pubsubpb.
|
||||
}
|
||||
|
||||
// DeleteSubscription deletes an existing subscription. All messages retained in the subscription
|
||||
// are immediately dropped. Calls to `Pull` after deletion will return
|
||||
// `NOT_FOUND`. After a subscription is deleted, a new one may be created with
|
||||
// are immediately dropped. Calls to Pull after deletion will return
|
||||
// NOT_FOUND. After a subscription is deleted, a new one may be created with
|
||||
// the same name, but the new one has no association with the old
|
||||
// subscription or its topic unless the same topic is specified.
|
||||
func (c *SubscriberClient) DeleteSubscription(ctx context.Context, req *pubsubpb.DeleteSubscriptionRequest, opts ...gax.CallOption) error {
|
||||
@@ -334,7 +339,7 @@ func (c *SubscriberClient) DeleteSubscription(ctx context.Context, req *pubsubpb
|
||||
// to indicate that more time is needed to process a message by the
|
||||
// subscriber, or to make the message available for redelivery if the
|
||||
// processing was interrupted. Note that this does not modify the
|
||||
// subscription-level `ackDeadlineSeconds` used for subsequent messages.
|
||||
// subscription-level ackDeadlineSeconds used for subsequent messages.
|
||||
func (c *SubscriberClient) ModifyAckDeadline(ctx context.Context, req *pubsubpb.ModifyAckDeadlineRequest, opts ...gax.CallOption) error {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.ModifyAckDeadline[0:len(c.CallOptions.ModifyAckDeadline):len(c.CallOptions.ModifyAckDeadline)], opts...)
|
||||
@@ -346,8 +351,8 @@ func (c *SubscriberClient) ModifyAckDeadline(ctx context.Context, req *pubsubpb.
|
||||
return err
|
||||
}
|
||||
|
||||
// Acknowledge acknowledges the messages associated with the `ack_ids` in the
|
||||
// `AcknowledgeRequest`. The Pub/Sub system can remove the relevant messages
|
||||
// Acknowledge acknowledges the messages associated with the ack_ids in the
|
||||
// AcknowledgeRequest. The Pub/Sub system can remove the relevant messages
|
||||
// from the subscription.
|
||||
//
|
||||
// Acknowledging a message whose ack deadline has expired may succeed,
|
||||
@@ -365,7 +370,7 @@ func (c *SubscriberClient) Acknowledge(ctx context.Context, req *pubsubpb.Acknow
|
||||
}
|
||||
|
||||
// Pull pulls messages from the server. Returns an empty list if there are no
|
||||
// messages available in the backlog. The server may return `UNAVAILABLE` if
|
||||
// messages available in the backlog. The server may return UNAVAILABLE if
|
||||
// there are too many concurrent pull requests pending for the given
|
||||
// subscription.
|
||||
func (c *SubscriberClient) Pull(ctx context.Context, req *pubsubpb.PullRequest, opts ...gax.CallOption) (*pubsubpb.PullResponse, error) {
|
||||
@@ -390,9 +395,9 @@ func (c *SubscriberClient) Pull(ctx context.Context, req *pubsubpb.PullRequest,
|
||||
// Establishes a stream with the server, which sends messages down to the
|
||||
// client. The client streams acknowledgements and ack deadline modifications
|
||||
// back to the server. The server will close the stream and return the status
|
||||
// on any error. The server may close the stream with status `OK` to reassign
|
||||
// on any error. The server may close the stream with status OK to reassign
|
||||
// server-side resources, in which case, the client should re-establish the
|
||||
// stream. `UNAVAILABLE` may also be returned in the case of a transient error
|
||||
// stream. UNAVAILABLE may also be returned in the case of a transient error
|
||||
// (e.g., a server restart). These should also be retried by the client. Flow
|
||||
// control can be achieved by configuring the underlying RPC channel.
|
||||
func (c *SubscriberClient) StreamingPull(ctx context.Context, opts ...gax.CallOption) (pubsubpb.Subscriber_StreamingPullClient, error) {
|
||||
@@ -410,12 +415,12 @@ func (c *SubscriberClient) StreamingPull(ctx context.Context, opts ...gax.CallOp
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// ModifyPushConfig modifies the `PushConfig` for a specified subscription.
|
||||
// ModifyPushConfig modifies the PushConfig for a specified subscription.
|
||||
//
|
||||
// This may be used to change a push subscription to a pull one (signified by
|
||||
// an empty `PushConfig`) or vice versa, or change the endpoint URL and other
|
||||
// an empty PushConfig) or vice versa, or change the endpoint URL and other
|
||||
// attributes of a push subscription. Messages will accumulate for delivery
|
||||
// continuously through the call regardless of changes to the `PushConfig`.
|
||||
// continuously through the call regardless of changes to the PushConfig.
|
||||
func (c *SubscriberClient) ModifyPushConfig(ctx context.Context, req *pubsubpb.ModifyPushConfigRequest, opts ...gax.CallOption) error {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.ModifyPushConfig[0:len(c.CallOptions.ModifyPushConfig):len(c.CallOptions.ModifyPushConfig)], opts...)
|
||||
@@ -463,13 +468,13 @@ func (c *SubscriberClient) ListSnapshots(ctx context.Context, req *pubsubpb.List
|
||||
}
|
||||
|
||||
// CreateSnapshot creates a snapshot from the requested subscription.
|
||||
// If the snapshot already exists, returns `ALREADY_EXISTS`.
|
||||
// If the requested subscription doesn't exist, returns `NOT_FOUND`.
|
||||
// If the snapshot already exists, returns ALREADY_EXISTS.
|
||||
// If the requested subscription doesn't exist, returns NOT_FOUND.
|
||||
//
|
||||
// If the name is not provided in the request, the server will assign a random
|
||||
// name for this snapshot on the same project as the subscription, conforming
|
||||
// to the
|
||||
// [resource name format](https://cloud.google.com/pubsub/docs/overview#names).
|
||||
// resource name format (at https://cloud.google.com/pubsub/docs/overview#names).
|
||||
// The generated name is populated in the returned Snapshot object.
|
||||
// Note that for REST API requests, you must specify a name in the request.
|
||||
func (c *SubscriberClient) CreateSnapshot(ctx context.Context, req *pubsubpb.CreateSnapshotRequest, opts ...gax.CallOption) (*pubsubpb.Snapshot, error) {
|
||||
@@ -487,6 +492,27 @@ func (c *SubscriberClient) CreateSnapshot(ctx context.Context, req *pubsubpb.Cre
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// UpdateSnapshot updates an existing snapshot. Note that certain properties of a snapshot
|
||||
// are not modifiable.
|
||||
// NOTE: The style guide requires body: "snapshot" instead of body: "*".
|
||||
// Keeping the latter for internal consistency in V1, however it should be
|
||||
// corrected in V2. See
|
||||
// https://cloud.google.com/apis/design/standard_methods#update for details.
|
||||
func (c *SubscriberClient) UpdateSnapshot(ctx context.Context, req *pubsubpb.UpdateSnapshotRequest, opts ...gax.CallOption) (*pubsubpb.Snapshot, error) {
|
||||
ctx = insertXGoog(ctx, c.xGoogHeader)
|
||||
opts = append(c.CallOptions.UpdateSnapshot[0:len(c.CallOptions.UpdateSnapshot):len(c.CallOptions.UpdateSnapshot)], opts...)
|
||||
var resp *pubsubpb.Snapshot
|
||||
err := gax.Invoke(ctx, func(ctx context.Context, settings gax.CallSettings) error {
|
||||
var err error
|
||||
resp, err = c.subscriberClient.UpdateSnapshot(ctx, req, settings.GRPC...)
|
||||
return err
|
||||
}, opts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// DeleteSnapshot removes an existing snapshot. All messages retained in the snapshot
|
||||
// are immediately dropped. After a snapshot is deleted, a new one may be
|
||||
// created with the same name, but the new one has no association with the old
|
||||
|
||||
+18
@@ -305,6 +305,24 @@ func ExampleSubscriberClient_CreateSnapshot() {
|
||||
_ = resp
|
||||
}
|
||||
|
||||
func ExampleSubscriberClient_UpdateSnapshot() {
|
||||
ctx := context.Background()
|
||||
c, err := pubsub.NewSubscriberClient(ctx)
|
||||
if err != nil {
|
||||
// TODO: Handle error.
|
||||
}
|
||||
|
||||
req := &pubsubpb.UpdateSnapshotRequest{
|
||||
// TODO: Fill request struct fields.
|
||||
}
|
||||
resp, err := c.UpdateSnapshot(ctx, req)
|
||||
if err != nil {
|
||||
// TODO: Handle error.
|
||||
}
|
||||
// TODO: Use resp.
|
||||
_ = resp
|
||||
}
|
||||
|
||||
func ExampleSubscriberClient_DeleteSnapshot() {
|
||||
ctx := context.Background()
|
||||
c, err := pubsub.NewSubscriberClient(ctx)
|
||||
|
||||
+3
-2
@@ -18,7 +18,7 @@ messages, hiding the the details of the underlying server RPCs. Google Cloud
|
||||
Pub/Sub is a many-to-many, asynchronous messaging system that decouples senders
|
||||
and receivers.
|
||||
|
||||
Note: This package is experimental and may make backwards-incompatible changes.
|
||||
Note: This package is in beta. Some backwards-incompatible changes may occur.
|
||||
|
||||
More information about Google Cloud Pub/Sub is available at
|
||||
https://cloud.google.com/pubsub/docs
|
||||
@@ -54,7 +54,8 @@ that is published to the topic will be delivered to all of its subscriptions.
|
||||
|
||||
Subsciptions may be created like so:
|
||||
|
||||
sub, err := pubsubClient.CreateSubscription(context.Background(), "sub-name", topic, 0, nil)
|
||||
sub, err := pubsubClient.CreateSubscription(context.Background(), "sub-name",
|
||||
pubsub.SubscriptionConfig{Topic: topic})
|
||||
|
||||
Messages are then consumed from a subscription via callback.
|
||||
|
||||
|
||||
+12
@@ -49,6 +49,18 @@ func ExampleClient_CreateTopic() {
|
||||
_ = topic // TODO: use the topic.
|
||||
}
|
||||
|
||||
// Use TopicInProject to refer to a topic that is not in the client's project, such
|
||||
// as a public topic.
|
||||
func ExampleClient_TopicInProject() {
|
||||
ctx := context.Background()
|
||||
client, err := pubsub.NewClient(ctx, "project-id")
|
||||
if err != nil {
|
||||
// TODO: Handle error.
|
||||
}
|
||||
topic := client.TopicInProject("topicName", "another-project-id")
|
||||
_ = topic // TODO: use the topic.
|
||||
}
|
||||
|
||||
func ExampleClient_CreateSubscription() {
|
||||
ctx := context.Background()
|
||||
client, err := pubsub.NewClient(ctx, "project-id")
|
||||
|
||||
+3
-10
@@ -76,15 +76,8 @@ func (s *fakeServer) wait() {
|
||||
}
|
||||
|
||||
func (s *fakeServer) StreamingPull(stream pb.Subscriber_StreamingPullServer) error {
|
||||
// Receive initial request.
|
||||
_, err := stream.Recv()
|
||||
if err == io.EOF {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Consume and ignore subsequent requests.
|
||||
s.wg.Add(1)
|
||||
defer s.wg.Done()
|
||||
errc := make(chan error, 1)
|
||||
s.wg.Add(1)
|
||||
go func() {
|
||||
@@ -124,7 +117,7 @@ func (s *fakeServer) StreamingPull(stream pb.Subscriber_StreamingPullServer) err
|
||||
// Add a slight delay to ensure the server receives any
|
||||
// messages en route from the client before shutting down the stream.
|
||||
// This reduces flakiness of tests involving retry.
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
if pr.err == io.EOF {
|
||||
return nil
|
||||
|
||||
+49
-37
@@ -31,6 +31,11 @@ import (
|
||||
"google.golang.org/api/option"
|
||||
)
|
||||
|
||||
var (
|
||||
topicIDs = testutil.NewUIDSpace("topic")
|
||||
subIDs = testutil.NewUIDSpace("sub")
|
||||
)
|
||||
|
||||
// messageData is used to hold the contents of a message so that it can be compared against the contents
|
||||
// of another message without regard to irrelevant fields.
|
||||
type messageData struct {
|
||||
@@ -47,8 +52,7 @@ func extractMessageData(m *Message) *messageData {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAll(t *testing.T) {
|
||||
t.Parallel()
|
||||
func integrationTestClient(t *testing.T, ctx context.Context) *Client {
|
||||
if testing.Short() {
|
||||
t.Skip("Integration tests skipped in short mode")
|
||||
}
|
||||
@@ -56,30 +60,31 @@ func TestAll(t *testing.T) {
|
||||
if projID == "" {
|
||||
t.Skip("Integration tests skipped. See CONTRIBUTING.md for details")
|
||||
}
|
||||
ctx := context.Background()
|
||||
ts := testutil.TokenSource(ctx, ScopePubSub, ScopeCloudPlatform)
|
||||
if ts == nil {
|
||||
t.Skip("Integration tests skipped. See CONTRIBUTING.md for details")
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
topicName := fmt.Sprintf("topic-%d", now.Unix())
|
||||
subName := fmt.Sprintf("subscription-%d", now.Unix())
|
||||
|
||||
client, err := NewClient(ctx, projID, option.WithTokenSource(ts))
|
||||
if err != nil {
|
||||
t.Fatalf("Creating client error: %v", err)
|
||||
}
|
||||
return client
|
||||
}
|
||||
|
||||
func TestAll(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
client := integrationTestClient(t, ctx)
|
||||
defer client.Close()
|
||||
|
||||
var topic *Topic
|
||||
if topic, err = client.CreateTopic(ctx, topicName); err != nil {
|
||||
topic, err := client.CreateTopic(ctx, topicIDs.New())
|
||||
if err != nil {
|
||||
t.Errorf("CreateTopic error: %v", err)
|
||||
}
|
||||
defer topic.Stop()
|
||||
|
||||
var sub *Subscription
|
||||
if sub, err = client.CreateSubscription(ctx, subName, SubscriptionConfig{Topic: topic}); err != nil {
|
||||
if sub, err = client.CreateSubscription(ctx, subIDs.New(), SubscriptionConfig{Topic: topic}); err != nil {
|
||||
t.Errorf("CreateSub error: %v", err)
|
||||
}
|
||||
|
||||
@@ -88,7 +93,7 @@ func TestAll(t *testing.T) {
|
||||
t.Fatalf("TopicExists error: %v", err)
|
||||
}
|
||||
if !exists {
|
||||
t.Errorf("topic %s should exist, but it doesn't", topic)
|
||||
t.Errorf("topic %v should exist, but it doesn't", topic)
|
||||
}
|
||||
|
||||
exists, err = sub.Exists(ctx)
|
||||
@@ -96,10 +101,10 @@ func TestAll(t *testing.T) {
|
||||
t.Fatalf("SubExists error: %v", err)
|
||||
}
|
||||
if !exists {
|
||||
t.Errorf("subscription %s should exist, but it doesn't", subName)
|
||||
t.Errorf("subscription %s should exist, but it doesn't", sub.ID())
|
||||
}
|
||||
|
||||
msgs := []*Message{}
|
||||
var msgs []*Message
|
||||
for i := 0; i < 10; i++ {
|
||||
text := fmt.Sprintf("a message with an index %d", i)
|
||||
attrs := make(map[string]string)
|
||||
@@ -275,37 +280,18 @@ func testIAM(ctx context.Context, h *iam.Handle, permission string) (msg string,
|
||||
func TestSubscriptionUpdate(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
if testing.Short() {
|
||||
t.Skip("Integration tests skipped in short mode")
|
||||
}
|
||||
projID := testutil.ProjID()
|
||||
if projID == "" {
|
||||
t.Skip("Integration tests skipped. See CONTRIBUTING.md for details.")
|
||||
}
|
||||
ts := testutil.TokenSource(ctx, ScopePubSub, ScopeCloudPlatform)
|
||||
if ts == nil {
|
||||
t.Skip("Integration tests skipped. See CONTRIBUTING.md for details")
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
topicName := fmt.Sprintf("topic-modify-%d", now.Unix())
|
||||
subName := fmt.Sprintf("subscription-modify-%d", now.Unix())
|
||||
|
||||
client, err := NewClient(ctx, projID, option.WithTokenSource(ts))
|
||||
if err != nil {
|
||||
t.Fatalf("Creating client error: %v", err)
|
||||
}
|
||||
client := integrationTestClient(t, ctx)
|
||||
defer client.Close()
|
||||
|
||||
var topic *Topic
|
||||
if topic, err = client.CreateTopic(ctx, topicName); err != nil {
|
||||
topic, err := client.CreateTopic(ctx, topicIDs.New())
|
||||
if err != nil {
|
||||
t.Fatalf("CreateTopic error: %v", err)
|
||||
}
|
||||
defer topic.Stop()
|
||||
defer topic.Delete(ctx)
|
||||
|
||||
var sub *Subscription
|
||||
if sub, err = client.CreateSubscription(ctx, subName, SubscriptionConfig{Topic: topic}); err != nil {
|
||||
if sub, err = client.CreateSubscription(ctx, subIDs.New(), SubscriptionConfig{Topic: topic}); err != nil {
|
||||
t.Fatalf("CreateSub error: %v", err)
|
||||
}
|
||||
defer sub.Delete(ctx)
|
||||
@@ -318,6 +304,7 @@ func TestSubscriptionUpdate(t *testing.T) {
|
||||
t.Fatalf("got %+v, want empty PushConfig")
|
||||
}
|
||||
// Add a PushConfig.
|
||||
projID := testutil.ProjID()
|
||||
pc := PushConfig{
|
||||
Endpoint: "https://" + projID + ".appspot.com/_ah/push-handlers/push",
|
||||
Attributes: map[string]string{"x-goog-version": "v1"},
|
||||
@@ -349,3 +336,28 @@ func TestSubscriptionUpdate(t *testing.T) {
|
||||
t.Fatal("got nil, wanted error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPublicTopic(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
client := integrationTestClient(t, ctx)
|
||||
defer client.Close()
|
||||
|
||||
sub, err := client.CreateSubscription(ctx, subIDs.New(), SubscriptionConfig{
|
||||
Topic: client.TopicInProject("taxirides-realtime", "pubsub-public-data"),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer sub.Delete(ctx)
|
||||
// Confirm that Receive works. It doesn't matter if we actually get any
|
||||
// messages.
|
||||
ctxt, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
err = sub.Receive(ctxt, func(_ context.Context, msg *Message) {
|
||||
msg.Ack()
|
||||
cancel()
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
+66
-319
@@ -15,203 +15,21 @@
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
"google.golang.org/api/iterator"
|
||||
"google.golang.org/api/support/bundler"
|
||||
pb "google.golang.org/genproto/googleapis/pubsub/v1"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
)
|
||||
|
||||
type messageIterator struct {
|
||||
impl interface {
|
||||
next() (*Message, error)
|
||||
stop()
|
||||
}
|
||||
}
|
||||
|
||||
type pollingMessageIterator struct {
|
||||
// kaTicker controls how often we send an ack deadline extension request.
|
||||
kaTicker *time.Ticker
|
||||
// ackTicker controls how often we acknowledge a batch of messages.
|
||||
ackTicker *time.Ticker
|
||||
|
||||
ka *keepAlive
|
||||
acker *acker
|
||||
nacker *bundler.Bundler
|
||||
puller *puller
|
||||
|
||||
// mu ensures that cleanup only happens once, and concurrent Stop
|
||||
// invocations block until cleanup completes.
|
||||
mu sync.Mutex
|
||||
|
||||
// closed is used to signal that Stop has been called.
|
||||
closed chan struct{}
|
||||
}
|
||||
|
||||
var useStreamingPull = false
|
||||
|
||||
// newMessageIterator starts a new messageIterator. Stop must be called on the messageIterator
|
||||
// newMessageIterator starts a new streamingMessageIterator. Stop must be called on the messageIterator
|
||||
// when it is no longer needed.
|
||||
// subName is the full name of the subscription to pull messages from.
|
||||
// ctx is the context to use for acking messages and extending message deadlines.
|
||||
func newMessageIterator(ctx context.Context, s service, subName string, po *pullOptions) *messageIterator {
|
||||
if !useStreamingPull {
|
||||
return &messageIterator{
|
||||
impl: newPollingMessageIterator(ctx, s, subName, po),
|
||||
}
|
||||
}
|
||||
func newMessageIterator(ctx context.Context, s service, subName string, po *pullOptions) *streamingMessageIterator {
|
||||
sp := s.newStreamingPuller(ctx, subName, int32(po.ackDeadline.Seconds()))
|
||||
err := sp.open()
|
||||
if grpc.Code(err) == codes.Unimplemented {
|
||||
log.Println("pubsub: streaming pull unimplemented; falling back to legacy pull")
|
||||
return &messageIterator{
|
||||
impl: newPollingMessageIterator(ctx, s, subName, po),
|
||||
}
|
||||
}
|
||||
// TODO(jba): handle other non-nil error?
|
||||
log.Println("using streaming pull")
|
||||
return &messageIterator{
|
||||
impl: newStreamingMessageIterator(ctx, sp, po),
|
||||
}
|
||||
}
|
||||
|
||||
func newPollingMessageIterator(ctx context.Context, s service, subName string, po *pullOptions) *pollingMessageIterator {
|
||||
// TODO: make kaTicker frequency more configurable.
|
||||
// (ackDeadline - 5s) is a reasonable default for now, because the minimum ack period is 10s. This gives us 5s grace.
|
||||
keepAlivePeriod := po.ackDeadline - 5*time.Second
|
||||
kaTicker := time.NewTicker(keepAlivePeriod) // Stopped in it.Stop
|
||||
|
||||
// TODO: make ackTicker more configurable. Something less than
|
||||
// kaTicker is a reasonable default (there's no point extending
|
||||
// messages when they could be acked instead).
|
||||
ackTicker := time.NewTicker(keepAlivePeriod / 2) // Stopped in it.Stop
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: ctx,
|
||||
Sub: subName,
|
||||
ExtensionTick: kaTicker.C,
|
||||
Deadline: po.ackDeadline,
|
||||
MaxExtension: po.maxExtension,
|
||||
}
|
||||
|
||||
ack := &acker{
|
||||
s: s,
|
||||
Ctx: ctx,
|
||||
Sub: subName,
|
||||
AckTick: ackTicker.C,
|
||||
Notify: ka.Remove,
|
||||
}
|
||||
|
||||
nacker := bundler.NewBundler("", func(ackIDs interface{}) {
|
||||
// NACK by setting the ack deadline to zero, to make the message
|
||||
// immediately available for redelivery.
|
||||
//
|
||||
// If the RPC fails, nothing we can do about it. In the worst case, the
|
||||
// deadline for these messages will expire and they will still get
|
||||
// redelivered.
|
||||
_ = s.modifyAckDeadline(ctx, subName, 0, ackIDs.([]string))
|
||||
})
|
||||
nacker.DelayThreshold = keepAlivePeriod / 10 // nack promptly
|
||||
nacker.BundleCountThreshold = 10
|
||||
|
||||
pull := newPuller(s, subName, ctx, po.maxPrefetch, ka.Add, ka.Remove)
|
||||
|
||||
ka.Start()
|
||||
ack.Start()
|
||||
return &pollingMessageIterator{
|
||||
kaTicker: kaTicker,
|
||||
ackTicker: ackTicker,
|
||||
ka: ka,
|
||||
acker: ack,
|
||||
nacker: nacker,
|
||||
puller: pull,
|
||||
closed: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
// Next returns the next Message to be processed. The caller must call
|
||||
// Message.Done when finished with it.
|
||||
// Once Stop has been called, calls to Next will return iterator.Done.
|
||||
func (it *messageIterator) Next() (*Message, error) {
|
||||
return it.impl.next()
|
||||
}
|
||||
|
||||
func (it *pollingMessageIterator) next() (*Message, error) {
|
||||
m, err := it.puller.Next()
|
||||
if err == nil {
|
||||
m.doneFunc = it.done
|
||||
return m, nil
|
||||
}
|
||||
|
||||
select {
|
||||
// If Stop has been called, we return Done regardless the value of err.
|
||||
case <-it.closed:
|
||||
return nil, iterator.Done
|
||||
default:
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// Client code must call Stop on a messageIterator when finished with it.
|
||||
// Stop will block until Done has been called on all Messages that have been
|
||||
// returned by Next, or until the context with which the messageIterator was created
|
||||
// is cancelled or exceeds its deadline.
|
||||
// Stop need only be called once, but may be called multiple times from
|
||||
// multiple goroutines.
|
||||
func (it *messageIterator) Stop() {
|
||||
it.impl.stop()
|
||||
}
|
||||
|
||||
func (it *pollingMessageIterator) stop() {
|
||||
it.mu.Lock()
|
||||
defer it.mu.Unlock()
|
||||
|
||||
select {
|
||||
case <-it.closed:
|
||||
// Cleanup has already been performed.
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
// We close this channel before calling it.puller.Stop to ensure that we
|
||||
// reliably return iterator.Done from Next.
|
||||
close(it.closed)
|
||||
|
||||
// Stop the puller. Once this completes, no more messages will be added
|
||||
// to it.ka.
|
||||
it.puller.Stop()
|
||||
|
||||
// Start acking messages as they arrive, ignoring ackTicker. This will
|
||||
// result in it.ka.Stop, below, returning as soon as possible.
|
||||
it.acker.FastMode()
|
||||
|
||||
// This will block until
|
||||
// (a) it.ka.Ctx is done, or
|
||||
// (b) all messages have been removed from keepAlive.
|
||||
// (b) will happen once all outstanding messages have been either ACKed or NACKed.
|
||||
it.ka.Stop()
|
||||
|
||||
// There are no more live messages, so kill off the acker.
|
||||
it.acker.Stop()
|
||||
it.nacker.Flush()
|
||||
it.kaTicker.Stop()
|
||||
it.ackTicker.Stop()
|
||||
}
|
||||
|
||||
func (it *pollingMessageIterator) done(ackID string, ack bool) {
|
||||
if ack {
|
||||
it.acker.Ack(ackID)
|
||||
// There's no need to call it.ka.Remove here, as acker will
|
||||
// call it via its Notify function.
|
||||
} else {
|
||||
it.ka.Remove(ackID)
|
||||
_ = it.nacker.Add(ackID, len(ackID)) // ignore error; this is just an optimization
|
||||
}
|
||||
_ = sp.open() // error stored in sp
|
||||
return newStreamingMessageIterator(ctx, sp, po)
|
||||
}
|
||||
|
||||
type streamingMessageIterator struct {
|
||||
@@ -224,7 +42,6 @@ type streamingMessageIterator struct {
|
||||
failed chan struct{} // closed on stream error
|
||||
stopped chan struct{} // closed when Stop is called
|
||||
drained chan struct{} // closed when stopped && no more pending messages
|
||||
msgc chan *Message
|
||||
wg sync.WaitGroup
|
||||
|
||||
mu sync.Mutex
|
||||
@@ -240,110 +57,40 @@ func newStreamingMessageIterator(ctx context.Context, sp *streamingPuller, po *p
|
||||
keepAlivePeriod := po.ackDeadline - 5*time.Second
|
||||
kaTicker := time.NewTicker(keepAlivePeriod)
|
||||
|
||||
// TODO: make ackTicker more configurable. Something less than
|
||||
// kaTicker is a reasonable default (there's no point extending
|
||||
// messages when they could be acked instead).
|
||||
ackTicker := time.NewTicker(keepAlivePeriod / 2)
|
||||
nackTicker := time.NewTicker(keepAlivePeriod / 10)
|
||||
// Ack promptly so users don't lose work if client crashes.
|
||||
ackTicker := time.NewTicker(100 * time.Millisecond)
|
||||
nackTicker := time.NewTicker(100 * time.Millisecond)
|
||||
it := &streamingMessageIterator{
|
||||
ctx: ctx,
|
||||
sp: sp,
|
||||
po: po,
|
||||
kaTicker: kaTicker,
|
||||
ackTicker: ackTicker,
|
||||
nackTicker: nackTicker,
|
||||
failed: make(chan struct{}),
|
||||
stopped: make(chan struct{}),
|
||||
drained: make(chan struct{}),
|
||||
// use maxPrefetch as the channel's buffer size.
|
||||
msgc: make(chan *Message, po.maxPrefetch),
|
||||
ctx: ctx,
|
||||
sp: sp,
|
||||
po: po,
|
||||
kaTicker: kaTicker,
|
||||
ackTicker: ackTicker,
|
||||
nackTicker: nackTicker,
|
||||
failed: make(chan struct{}),
|
||||
stopped: make(chan struct{}),
|
||||
drained: make(chan struct{}),
|
||||
keepAliveDeadlines: map[string]time.Time{},
|
||||
pendingReq: &pb.StreamingPullRequest{},
|
||||
}
|
||||
it.wg.Add(2)
|
||||
go it.receiver()
|
||||
it.wg.Add(1)
|
||||
go it.sender()
|
||||
return it
|
||||
}
|
||||
|
||||
func (it *streamingMessageIterator) next() (*Message, error) {
|
||||
// If ctx has been cancelled or the iterator is done, return straight
|
||||
// away (even if there are buffered messages available).
|
||||
select {
|
||||
case <-it.ctx.Done():
|
||||
return nil, it.ctx.Err()
|
||||
|
||||
case <-it.failed:
|
||||
break
|
||||
|
||||
case <-it.stopped:
|
||||
break
|
||||
|
||||
default:
|
||||
// Wait for a message, but also for one of the above conditions.
|
||||
select {
|
||||
case msg := <-it.msgc:
|
||||
// Since active select cases are chosen at random, this can return
|
||||
// nil (from the channel close) even if it.failed or it.stopped is
|
||||
// closed.
|
||||
if msg == nil {
|
||||
break
|
||||
}
|
||||
msg.doneFunc = it.done
|
||||
return msg, nil
|
||||
|
||||
case <-it.ctx.Done():
|
||||
return nil, it.ctx.Err()
|
||||
|
||||
case <-it.failed:
|
||||
break
|
||||
|
||||
case <-it.stopped:
|
||||
break
|
||||
}
|
||||
}
|
||||
// Here if the iterator is done.
|
||||
it.mu.Lock()
|
||||
defer it.mu.Unlock()
|
||||
return nil, it.err
|
||||
}
|
||||
|
||||
// Subscription.receive will call stop on its messageIterator when finished with it.
|
||||
// Stop will block until Done has been called on all Messages that have been
|
||||
// returned by Next, or until the context with which the messageIterator was created
|
||||
// is cancelled or exceeds its deadline.
|
||||
func (it *streamingMessageIterator) stop() {
|
||||
it.mu.Lock()
|
||||
select {
|
||||
case <-it.stopped:
|
||||
it.mu.Unlock()
|
||||
it.wg.Wait()
|
||||
return
|
||||
default:
|
||||
close(it.stopped)
|
||||
}
|
||||
if it.err == nil {
|
||||
it.err = iterator.Done
|
||||
}
|
||||
// Before reading from the channel, see if we're already drained.
|
||||
it.checkDrained()
|
||||
it.mu.Unlock()
|
||||
// Nack all the pending messages.
|
||||
// Grab the lock separately for each message to allow the receiver
|
||||
// and sender goroutines to make progress.
|
||||
// Why this will eventually terminate:
|
||||
// - If the receiver is not blocked on a stream Recv, then
|
||||
// it will write all the messages it has received to the channel,
|
||||
// then exit, closing the channel.
|
||||
// - If the receiver is blocked, then this loop will eventually
|
||||
// nack all the messages in the channel. Once done is called
|
||||
// on the remaining messages, the iterator will be marked as drained,
|
||||
// which will trigger the sender to terminate. When it does, it
|
||||
// performs a CloseSend on the stream, which will result in the blocked
|
||||
// stream Recv returning.
|
||||
for m := range it.msgc {
|
||||
it.mu.Lock()
|
||||
delete(it.keepAliveDeadlines, m.ackID)
|
||||
it.addDeadlineMod(m.ackID, 0)
|
||||
it.checkDrained()
|
||||
it.mu.Unlock()
|
||||
}
|
||||
it.wg.Wait()
|
||||
}
|
||||
|
||||
@@ -398,52 +145,40 @@ func (it *streamingMessageIterator) fail(err error) {
|
||||
it.mu.Unlock()
|
||||
}
|
||||
|
||||
// receiver runs in a goroutine and handles all receives from the stream.
|
||||
func (it *streamingMessageIterator) receiver() {
|
||||
defer it.wg.Done()
|
||||
defer close(it.msgc)
|
||||
for {
|
||||
// Stop retrieving messages if the context is done, the stream
|
||||
// failed, or the iterator's Stop method was called.
|
||||
select {
|
||||
case <-it.ctx.Done():
|
||||
return
|
||||
case <-it.failed:
|
||||
return
|
||||
case <-it.stopped:
|
||||
return
|
||||
default:
|
||||
}
|
||||
// Receive messages from stream. This may block indefinitely.
|
||||
msgs, err := it.sp.fetchMessages()
|
||||
|
||||
// The streamingPuller handles retries, so any error here
|
||||
// is fatal to the iterator.
|
||||
if err != nil {
|
||||
it.fail(err)
|
||||
return
|
||||
}
|
||||
// We received some messages. Remember them so we can
|
||||
// keep them alive.
|
||||
deadline := time.Now().Add(it.po.maxExtension)
|
||||
it.mu.Lock()
|
||||
for _, m := range msgs {
|
||||
it.keepAliveDeadlines[m.ackID] = deadline
|
||||
}
|
||||
it.mu.Unlock()
|
||||
// Deliver the messages to the channel.
|
||||
for _, m := range msgs {
|
||||
select {
|
||||
case <-it.ctx.Done():
|
||||
return
|
||||
case <-it.failed:
|
||||
return
|
||||
// Don't return if stopped. We want to send the remaining
|
||||
// messages on the channel, where they will be nacked.
|
||||
case it.msgc <- m:
|
||||
}
|
||||
}
|
||||
// receive makes a call to the stream's Recv method and returns
|
||||
// its messages.
|
||||
func (it *streamingMessageIterator) receive() ([]*Message, error) {
|
||||
// Stop retrieving messages if the context is done, the stream
|
||||
// failed, or the iterator's Stop method was called.
|
||||
select {
|
||||
case <-it.ctx.Done():
|
||||
return nil, it.ctx.Err()
|
||||
default:
|
||||
}
|
||||
it.mu.Lock()
|
||||
err := it.err
|
||||
it.mu.Unlock()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// Receive messages from stream. This may block indefinitely.
|
||||
msgs, err := it.sp.fetchMessages()
|
||||
// The streamingPuller handles retries, so any error here
|
||||
// is fatal.
|
||||
if err != nil {
|
||||
it.fail(err)
|
||||
return nil, err
|
||||
}
|
||||
// We received some messages. Remember them so we can
|
||||
// keep them alive.
|
||||
deadline := time.Now().Add(it.po.maxExtension)
|
||||
it.mu.Lock()
|
||||
for _, m := range msgs {
|
||||
m.doneFunc = it.done
|
||||
it.keepAliveDeadlines[m.ackID] = deadline
|
||||
}
|
||||
it.mu.Unlock()
|
||||
return msgs, nil
|
||||
}
|
||||
|
||||
// sender runs in a goroutine and handles all sends to the stream.
|
||||
@@ -522,3 +257,15 @@ func (it *streamingMessageIterator) handleKeepAlives() bool {
|
||||
it.checkDrained()
|
||||
return len(live) > 0
|
||||
}
|
||||
|
||||
func getKeepAliveAckIDs(items map[string]time.Time) (live, expired []string) {
|
||||
now := time.Now()
|
||||
for id, expiry := range items {
|
||||
if expiry.Before(now) {
|
||||
expired = append(expired, id)
|
||||
} else {
|
||||
live = append(live, id)
|
||||
}
|
||||
}
|
||||
return live, expired
|
||||
}
|
||||
|
||||
-338
@@ -1,338 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
|
||||
"google.golang.org/api/iterator"
|
||||
)
|
||||
|
||||
func TestReturnsDoneOnStop(t *testing.T) {
|
||||
if useStreamingPull {
|
||||
t.Skip("iterator tests are for polling pull only")
|
||||
}
|
||||
type testCase struct {
|
||||
abort func(*messageIterator, context.CancelFunc)
|
||||
want error
|
||||
}
|
||||
|
||||
for _, tc := range []testCase{
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
it.Stop()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
cancel()
|
||||
},
|
||||
want: context.Canceled,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
it.Stop()
|
||||
cancel()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
cancel()
|
||||
it.Stop()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
} {
|
||||
s := &blockingFetch{}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
it := newMessageIterator(ctx, s, "subname", &pullOptions{ackDeadline: time.Second * 10, maxExtension: time.Hour})
|
||||
defer it.Stop()
|
||||
tc.abort(it, cancel)
|
||||
|
||||
_, err := it.Next()
|
||||
if err != tc.want {
|
||||
t.Errorf("iterator Next error after abort: got:\n%v\nwant:\n%v", err, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// blockingFetch implements message fetching by not returning until its context is cancelled.
|
||||
type blockingFetch struct {
|
||||
service
|
||||
}
|
||||
|
||||
func (s *blockingFetch) fetchMessages(ctx context.Context, subName string, maxMessages int32) ([]*Message, error) {
|
||||
<-ctx.Done()
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
|
||||
func (s *blockingFetch) newStreamingPuller(ctx context.Context, subName string, ackDeadline int32) *streamingPuller {
|
||||
return nil
|
||||
}
|
||||
|
||||
// justInTimeFetch simulates the situation where the iterator is aborted just after the fetch RPC
|
||||
// succeeds, so the rest of puller.Next will continue to execute and return sucessfully.
|
||||
type justInTimeFetch struct {
|
||||
service
|
||||
}
|
||||
|
||||
func (s *justInTimeFetch) fetchMessages(ctx context.Context, subName string, maxMessages int32) ([]*Message, error) {
|
||||
<-ctx.Done()
|
||||
// The context was cancelled, but let's pretend that this happend just after our RPC returned.
|
||||
|
||||
var result []*Message
|
||||
for i := 0; i < int(maxMessages); i++ {
|
||||
val := fmt.Sprintf("msg%v", i)
|
||||
result = append(result, &Message{Data: []byte(val), ackID: val})
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s *justInTimeFetch) splitAckIDs(ids []string) ([]string, []string) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *justInTimeFetch) modifyAckDeadline(ctx context.Context, subName string, deadline time.Duration, ackIDs []string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *justInTimeFetch) newStreamingPuller(ctx context.Context, subName string, ackDeadline int32) *streamingPuller {
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestAfterAbortReturnsNoMoreThanOneMessage(t *testing.T) {
|
||||
// Each test case is excercised by making two concurrent blocking calls on a
|
||||
// messageIterator, and then aborting the iterator.
|
||||
// The result should be one call to Next returning a message, and the other returning an error.
|
||||
t.Skip(`This test has subtle timing dependencies, making it flaky.
|
||||
It is not worth fixing because iterators will be removed shortly.`)
|
||||
type testCase struct {
|
||||
abort func(*messageIterator, context.CancelFunc)
|
||||
// want is the error that should be returned from one Next invocation.
|
||||
want error
|
||||
}
|
||||
for n := 1; n < 3; n++ {
|
||||
for _, tc := range []testCase{
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
it.Stop()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
cancel()
|
||||
},
|
||||
want: context.Canceled,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
it.Stop()
|
||||
cancel()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
{
|
||||
abort: func(it *messageIterator, cancel context.CancelFunc) {
|
||||
cancel()
|
||||
it.Stop()
|
||||
},
|
||||
want: iterator.Done,
|
||||
},
|
||||
} {
|
||||
s := &justInTimeFetch{}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
// if maxPrefetch == 1, there will be no messages in the puller buffer when Next is invoked the second time.
|
||||
// if maxPrefetch == 2, there will be 1 message in the puller buffer when Next is invoked the second time.
|
||||
po := &pullOptions{
|
||||
ackDeadline: time.Second * 10,
|
||||
maxExtension: time.Hour,
|
||||
maxPrefetch: int32(n),
|
||||
}
|
||||
it := newMessageIterator(ctx, s, "subname", po)
|
||||
defer it.Stop()
|
||||
|
||||
type result struct {
|
||||
m *Message
|
||||
err error
|
||||
}
|
||||
results := make(chan *result, 2)
|
||||
|
||||
for i := 0; i < 2; i++ {
|
||||
go func() {
|
||||
m, err := it.Next()
|
||||
results <- &result{m, err}
|
||||
if err == nil {
|
||||
m.Nack()
|
||||
}
|
||||
}()
|
||||
}
|
||||
// Wait for goroutines to block on it.Next().
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
tc.abort(it, cancel)
|
||||
|
||||
result1 := <-results
|
||||
result2 := <-results
|
||||
|
||||
// There should be one error result, and one non-error result.
|
||||
// Make result1 be the non-error result.
|
||||
if result1.err != nil {
|
||||
result1, result2 = result2, result1
|
||||
}
|
||||
|
||||
if string(result1.m.Data) != "msg0" {
|
||||
t.Errorf("After abort, got message: %v, want %v", result1.m.Data, "msg0")
|
||||
}
|
||||
if result1.err != nil {
|
||||
t.Errorf("After abort, got : %v, want nil", result1.err)
|
||||
}
|
||||
if result2.m != nil {
|
||||
t.Errorf("After abort, got message: %v, want nil", result2.m)
|
||||
}
|
||||
if result2.err != tc.want {
|
||||
t.Errorf("After abort, got err: %v, want %v", result2.err, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type fetcherServiceWithModifyAckDeadline struct {
|
||||
fetcherService
|
||||
events chan string
|
||||
}
|
||||
|
||||
func (f *fetcherServiceWithModifyAckDeadline) modifyAckDeadline(_ context.Context, _ string, d time.Duration, ids []string) error {
|
||||
// Different versions of Go use different representations for time.Duration(0).
|
||||
var ds string
|
||||
if d == 0 {
|
||||
ds = "0s"
|
||||
} else {
|
||||
ds = d.String()
|
||||
}
|
||||
f.events <- fmt.Sprintf("modAck(%v, %s)", ids, ds)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fetcherServiceWithModifyAckDeadline) splitAckIDs(ackIDs []string) ([]string, []string) {
|
||||
return ackIDs, nil
|
||||
}
|
||||
|
||||
func (f *fetcherServiceWithModifyAckDeadline) newStreamingPuller(ctx context.Context, subName string, ackDeadline int32) *streamingPuller {
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestMultipleStopCallsBlockUntilMessageDone(t *testing.T) {
|
||||
t.Skip(`This test has subtle timing dependencies, making it flaky.
|
||||
It is not worth fixing because iterators will be removed shortly.`)
|
||||
events := make(chan string, 3)
|
||||
s := &fetcherServiceWithModifyAckDeadline{
|
||||
fetcherService{
|
||||
results: []fetchResult{
|
||||
{
|
||||
msgs: []*Message{{ackID: "a"}, {ackID: "b"}},
|
||||
},
|
||||
},
|
||||
},
|
||||
events,
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
it := newMessageIterator(ctx, s, "subname", &pullOptions{ackDeadline: time.Second * 10, maxExtension: 0})
|
||||
|
||||
m, err := it.Next()
|
||||
if err != nil {
|
||||
t.Errorf("error calling Next: %v", err)
|
||||
}
|
||||
|
||||
go func() {
|
||||
it.Stop()
|
||||
events <- "stopped"
|
||||
}()
|
||||
go func() {
|
||||
it.Stop()
|
||||
events <- "stopped"
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-events:
|
||||
t.Fatal("Stop is not blocked")
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
m.Nack()
|
||||
|
||||
got := []string{<-events, <-events, <-events}
|
||||
want := []string{"modAck([a], 0s)", "stopped", "stopped"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("stopping iterator, got: %v ; want: %v", got, want)
|
||||
}
|
||||
|
||||
// The iterator is stopped, so should not return another message.
|
||||
m, err = it.Next()
|
||||
if m != nil {
|
||||
t.Errorf("message got: %v ; want: nil", m)
|
||||
}
|
||||
if err != iterator.Done {
|
||||
t.Errorf("err got: %v ; want: %v", err, iterator.Done)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFastNack(t *testing.T) {
|
||||
if useStreamingPull {
|
||||
t.Skip("iterator tests are for polling pull only")
|
||||
}
|
||||
events := make(chan string, 3)
|
||||
s := &fetcherServiceWithModifyAckDeadline{
|
||||
fetcherService{
|
||||
results: []fetchResult{
|
||||
{
|
||||
msgs: []*Message{{ackID: "a"}, {ackID: "b"}},
|
||||
},
|
||||
},
|
||||
},
|
||||
events,
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
it := newMessageIterator(ctx, s, "subname", &pullOptions{
|
||||
ackDeadline: time.Second * 6,
|
||||
maxExtension: time.Second * 10,
|
||||
})
|
||||
// Get both messages.
|
||||
_, err := it.Next()
|
||||
if err != nil {
|
||||
t.Errorf("error calling Next: %v", err)
|
||||
}
|
||||
m2, err := it.Next()
|
||||
if err != nil {
|
||||
t.Errorf("error calling Next: %v", err)
|
||||
}
|
||||
// Ignore the first, nack the second.
|
||||
m2.Nack()
|
||||
|
||||
got := []string{<-events, <-events}
|
||||
// The nack should happen before the deadline extension.
|
||||
want := []string{"modAck([b], 0s)", "modAck([a], 6s)"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("got: %v ; want: %v", got, want)
|
||||
}
|
||||
}
|
||||
-182
@@ -1,182 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
// keepAlive keeps track of which Messages need to have their deadline extended, and
|
||||
// periodically extends them.
|
||||
// Messages are tracked by Ack ID.
|
||||
type keepAlive struct {
|
||||
s service
|
||||
Ctx context.Context // The context to use when extending deadlines.
|
||||
Sub string // The full name of the subscription.
|
||||
ExtensionTick <-chan time.Time // ExtensionTick supplies the frequency with which to make extension requests.
|
||||
Deadline time.Duration // How long to extend messages for each time they are extended. Should be greater than ExtensionTick frequency.
|
||||
MaxExtension time.Duration // How long to keep extending each message's ack deadline before automatically removing it.
|
||||
|
||||
mu sync.Mutex
|
||||
// key: ackID; value: time at which ack deadline extension should cease.
|
||||
items map[string]time.Time
|
||||
dr drain
|
||||
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// Start initiates the deadline extension loop. Stop must be called once keepAlive is no longer needed.
|
||||
func (ka *keepAlive) Start() {
|
||||
ka.items = make(map[string]time.Time)
|
||||
ka.dr = drain{Drained: make(chan struct{})}
|
||||
ka.wg.Add(1)
|
||||
go func() {
|
||||
defer ka.wg.Done()
|
||||
for {
|
||||
select {
|
||||
case <-ka.Ctx.Done():
|
||||
// Don't bother waiting for items to be removed: we can't extend them any more.
|
||||
return
|
||||
case <-ka.dr.Drained:
|
||||
return
|
||||
case <-ka.ExtensionTick:
|
||||
live, expired := ka.getAckIDs()
|
||||
ka.wg.Add(1)
|
||||
go func() {
|
||||
defer ka.wg.Done()
|
||||
ka.extendDeadlines(live)
|
||||
}()
|
||||
|
||||
for _, id := range expired {
|
||||
ka.Remove(id)
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// Add adds an ack id to be kept alive.
|
||||
// It should not be called after Stop.
|
||||
func (ka *keepAlive) Add(ackID string) {
|
||||
ka.mu.Lock()
|
||||
defer ka.mu.Unlock()
|
||||
|
||||
ka.items[ackID] = time.Now().Add(ka.MaxExtension)
|
||||
ka.dr.SetPending(true)
|
||||
}
|
||||
|
||||
// Remove removes ackID from the list to be kept alive.
|
||||
func (ka *keepAlive) Remove(ackID string) {
|
||||
ka.mu.Lock()
|
||||
defer ka.mu.Unlock()
|
||||
|
||||
// Note: If users NACKs a message after it has been removed due to
|
||||
// expiring, Remove will be called twice with same ack id. This is OK.
|
||||
delete(ka.items, ackID)
|
||||
ka.dr.SetPending(len(ka.items) != 0)
|
||||
}
|
||||
|
||||
// Stop waits until all added ackIDs have been removed, and cleans up resources.
|
||||
// Stop may only be called once.
|
||||
func (ka *keepAlive) Stop() {
|
||||
ka.mu.Lock()
|
||||
ka.dr.Drain()
|
||||
ka.mu.Unlock()
|
||||
|
||||
ka.wg.Wait()
|
||||
}
|
||||
|
||||
// getAckIDs returns the set of ackIDs that are being kept alive.
|
||||
// The set is divided into two lists: one with IDs that should continue to be kept alive,
|
||||
// and the other with IDs that should be dropped.
|
||||
func (ka *keepAlive) getAckIDs() (live, expired []string) {
|
||||
ka.mu.Lock()
|
||||
defer ka.mu.Unlock()
|
||||
return getKeepAliveAckIDs(ka.items)
|
||||
}
|
||||
|
||||
func getKeepAliveAckIDs(items map[string]time.Time) (live, expired []string) {
|
||||
now := time.Now()
|
||||
for id, expiry := range items {
|
||||
if expiry.Before(now) {
|
||||
expired = append(expired, id)
|
||||
} else {
|
||||
live = append(live, id)
|
||||
}
|
||||
}
|
||||
return live, expired
|
||||
}
|
||||
|
||||
const maxExtensionAttempts = 2
|
||||
|
||||
func (ka *keepAlive) extendDeadlines(ackIDs []string) {
|
||||
head, tail := ka.s.splitAckIDs(ackIDs)
|
||||
for len(head) > 0 {
|
||||
for i := 0; i < maxExtensionAttempts; i++ {
|
||||
if ka.s.modifyAckDeadline(ka.Ctx, ka.Sub, ka.Deadline, head) == nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
// NOTE: Messages whose deadlines we fail to extend will
|
||||
// eventually be redelivered and this is a documented behaviour
|
||||
// of the API.
|
||||
//
|
||||
// NOTE: If we fail to extend deadlines here, this
|
||||
// implementation will continue to attempt extending the
|
||||
// deadlines for those ack IDs the next time the extension
|
||||
// ticker ticks. By then the deadline will have expired.
|
||||
// Re-extending them is harmless, however.
|
||||
//
|
||||
// TODO: call Remove for ids which fail to be extended.
|
||||
|
||||
head, tail = ka.s.splitAckIDs(tail)
|
||||
}
|
||||
}
|
||||
|
||||
// A drain (once started) indicates via a channel when there is no work pending.
|
||||
type drain struct {
|
||||
started bool
|
||||
pending bool
|
||||
|
||||
// Drained is closed once there are no items outstanding if Drain has been called.
|
||||
Drained chan struct{}
|
||||
}
|
||||
|
||||
// Drain starts the drain process. This cannot be undone.
|
||||
func (d *drain) Drain() {
|
||||
d.started = true
|
||||
d.closeIfDrained()
|
||||
}
|
||||
|
||||
// SetPending sets whether there is work pending or not. It may be called multiple times before or after Drain.
|
||||
func (d *drain) SetPending(pending bool) {
|
||||
d.pending = pending
|
||||
d.closeIfDrained()
|
||||
}
|
||||
|
||||
func (d *drain) closeIfDrained() {
|
||||
if !d.pending && d.started {
|
||||
// Check to see if d.Drained is closed before closing it.
|
||||
// This allows SetPending(false) to be safely called multiple times.
|
||||
select {
|
||||
case <-d.Drained:
|
||||
default:
|
||||
close(d.Drained)
|
||||
}
|
||||
}
|
||||
}
|
||||
-319
@@ -1,319 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"reflect"
|
||||
"sort"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
func TestKeepAliveExtendsDeadline(t *testing.T) {
|
||||
ticker := make(chan time.Time)
|
||||
deadline := time.Nanosecond * 15
|
||||
s := &testService{modDeadlineCalled: make(chan modDeadlineCall)}
|
||||
|
||||
checkModDeadlineCall := func(ackIDs []string) {
|
||||
got := <-s.modDeadlineCalled
|
||||
sort.Strings(got.ackIDs)
|
||||
|
||||
want := modDeadlineCall{
|
||||
subName: "subname",
|
||||
deadline: deadline,
|
||||
ackIDs: ackIDs,
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("keepalive: got:\n%v\nwant:\n%v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
ExtensionTick: ticker,
|
||||
Deadline: deadline,
|
||||
MaxExtension: time.Hour,
|
||||
}
|
||||
ka.Start()
|
||||
|
||||
ka.Add("a")
|
||||
ka.Add("b")
|
||||
ticker <- time.Time{}
|
||||
checkModDeadlineCall([]string{"a", "b"})
|
||||
ka.Add("c")
|
||||
ka.Remove("b")
|
||||
ticker <- time.Time{}
|
||||
checkModDeadlineCall([]string{"a", "c"})
|
||||
ka.Remove("a")
|
||||
ka.Remove("c")
|
||||
ka.Add("d")
|
||||
ticker <- time.Time{}
|
||||
checkModDeadlineCall([]string{"d"})
|
||||
|
||||
ka.Remove("d")
|
||||
ka.Stop()
|
||||
}
|
||||
|
||||
func TestKeepAliveStopsWhenNoItem(t *testing.T) {
|
||||
ticker := make(chan time.Time)
|
||||
stopped := make(chan bool)
|
||||
s := &testService{modDeadlineCalled: make(chan modDeadlineCall, 3)}
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
ExtensionTick: ticker,
|
||||
}
|
||||
|
||||
ka.Start()
|
||||
|
||||
// There should be no call to modifyAckDeadline since there is no item.
|
||||
ticker <- time.Time{}
|
||||
|
||||
go func() {
|
||||
ka.Stop() // No items; should not block
|
||||
if len(s.modDeadlineCalled) > 0 {
|
||||
t.Errorf("unexpected extension to non-existent items: %v", <-s.modDeadlineCalled)
|
||||
}
|
||||
close(stopped)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-time.After(time.Second):
|
||||
t.Errorf("keepAlive timed out waiting for stop")
|
||||
}
|
||||
}
|
||||
|
||||
func TestKeepAliveStopsWhenItemsExpired(t *testing.T) {
|
||||
ticker := make(chan time.Time)
|
||||
stopped := make(chan bool)
|
||||
s := &testService{modDeadlineCalled: make(chan modDeadlineCall, 2)}
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
ExtensionTick: ticker,
|
||||
MaxExtension: time.Duration(0), // Should expire items at the first tick.
|
||||
}
|
||||
|
||||
ka.Start()
|
||||
ka.Add("a")
|
||||
ka.Add("b")
|
||||
|
||||
// Wait until the clock advances. Without this loop, this test fails on
|
||||
// Windows because the clock doesn't advance at all between ka.Add and the
|
||||
// expiration check after the tick is received.
|
||||
begin := time.Now()
|
||||
for time.Now().Equal(begin) {
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
|
||||
// There should be no call to modifyAckDeadline since both items are expired.
|
||||
ticker <- time.Time{}
|
||||
|
||||
go func() {
|
||||
ka.Stop() // No live items; should not block.
|
||||
if len(s.modDeadlineCalled) > 0 {
|
||||
t.Errorf("unexpected extension to expired items")
|
||||
}
|
||||
close(stopped)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-time.After(time.Second):
|
||||
t.Errorf("timed out waiting for stop")
|
||||
}
|
||||
}
|
||||
|
||||
func TestKeepAliveBlocksUntilAllItemsRemoved(t *testing.T) {
|
||||
ticker := make(chan time.Time)
|
||||
eventc := make(chan string, 3)
|
||||
s := &testService{modDeadlineCalled: make(chan modDeadlineCall)}
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
ExtensionTick: ticker,
|
||||
MaxExtension: time.Hour, // Should not expire.
|
||||
}
|
||||
|
||||
ka.Start()
|
||||
ka.Add("a")
|
||||
ka.Add("b")
|
||||
|
||||
go func() {
|
||||
ticker <- time.Time{}
|
||||
|
||||
// We expect a call since both items should be extended.
|
||||
select {
|
||||
case args := <-s.modDeadlineCalled:
|
||||
sort.Strings(args.ackIDs)
|
||||
got := args.ackIDs
|
||||
want := []string{"a", "b"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("mismatching IDs:\ngot %v\nwant %v", got, want)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Errorf("timed out waiting for deadline extend call")
|
||||
}
|
||||
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
eventc <- "pre-remove-b"
|
||||
// Remove one item, Stop should still be waiting.
|
||||
ka.Remove("b")
|
||||
|
||||
ticker <- time.Time{}
|
||||
|
||||
// We expect a call since the item is still alive.
|
||||
select {
|
||||
case args := <-s.modDeadlineCalled:
|
||||
got := args.ackIDs
|
||||
want := []string{"a"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("mismatching IDs:\ngot %v\nwant %v", got, want)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Errorf("timed out waiting for deadline extend call")
|
||||
}
|
||||
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
eventc <- "pre-remove-a"
|
||||
// Remove the last item so that Stop can proceed.
|
||||
ka.Remove("a")
|
||||
}()
|
||||
|
||||
go func() {
|
||||
ka.Stop() // Should block all item are removed.
|
||||
eventc <- "post-stop"
|
||||
}()
|
||||
|
||||
for i, want := range []string{"pre-remove-b", "pre-remove-a", "post-stop"} {
|
||||
select {
|
||||
case got := <-eventc:
|
||||
if got != want {
|
||||
t.Errorf("event #%d:\ngot %v\nwant %v", i, got, want)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Errorf("time out waiting for #%d event: want %v", i, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// extendCallResult contains a list of ackIDs which are expected in an ackID
|
||||
// extension request, along with the result that should be returned.
|
||||
type extendCallResult struct {
|
||||
ackIDs []string
|
||||
err error
|
||||
}
|
||||
|
||||
// extendService implements modifyAckDeadline using a hard-coded list of extendCallResults.
|
||||
type extendService struct {
|
||||
service
|
||||
|
||||
calls []extendCallResult
|
||||
|
||||
t *testing.T // used for error logging.
|
||||
}
|
||||
|
||||
func (es *extendService) modifyAckDeadline(ctx context.Context, subName string, deadline time.Duration, ackIDs []string) error {
|
||||
if len(es.calls) == 0 {
|
||||
es.t.Fatalf("unexpected call to modifyAckDeadline: ackIDs: %v", ackIDs)
|
||||
}
|
||||
call := es.calls[0]
|
||||
es.calls = es.calls[1:]
|
||||
|
||||
if got, want := ackIDs, call.ackIDs; !reflect.DeepEqual(got, want) {
|
||||
es.t.Errorf("unexpected arguments to modifyAckDeadline: got: %v ; want: %v", got, want)
|
||||
}
|
||||
return call.err
|
||||
}
|
||||
|
||||
// Test implementation returns the first 2 elements as head, and the rest as tail.
|
||||
func (es *extendService) splitAckIDs(ids []string) ([]string, []string) {
|
||||
if len(ids) < 2 {
|
||||
return ids, nil
|
||||
}
|
||||
return ids[:2], ids[2:]
|
||||
}
|
||||
|
||||
func TestKeepAliveSplitsBatches(t *testing.T) {
|
||||
type testCase struct {
|
||||
calls []extendCallResult
|
||||
}
|
||||
for _, tc := range []testCase{
|
||||
{
|
||||
calls: []extendCallResult{
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"e", "f"},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
calls: []extendCallResult{
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// On error we retry once.
|
||||
{
|
||||
ackIDs: []string{"a", "b"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// We give up after failing twice, so we move on to the next set, "c" and "d".
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
err: errors.New("bang"),
|
||||
},
|
||||
// Again, we retry once.
|
||||
{
|
||||
ackIDs: []string{"c", "d"},
|
||||
},
|
||||
{
|
||||
ackIDs: []string{"e", "f"},
|
||||
},
|
||||
},
|
||||
},
|
||||
} {
|
||||
s := &extendService{
|
||||
t: t,
|
||||
calls: tc.calls,
|
||||
}
|
||||
|
||||
ka := &keepAlive{
|
||||
s: s,
|
||||
Ctx: context.Background(),
|
||||
Sub: "subname",
|
||||
}
|
||||
|
||||
ka.extendDeadlines([]string{"a", "b", "c", "d", "e", "f"})
|
||||
|
||||
if len(s.calls) != 0 {
|
||||
t.Errorf("expected extend calls did not occur: %v", s.calls)
|
||||
}
|
||||
}
|
||||
}
|
||||
+7
-3
@@ -31,7 +31,7 @@ import (
|
||||
|
||||
"cloud.google.com/go/internal/testutil"
|
||||
"cloud.google.com/go/pubsub"
|
||||
"google.golang.org/api/transport"
|
||||
gtransport "google.golang.org/api/transport/grpc"
|
||||
pb "google.golang.org/genproto/googleapis/pubsub/v1"
|
||||
)
|
||||
|
||||
@@ -100,9 +100,13 @@ func perfClient(pubDelay time.Duration, nConns int, f interface {
|
||||
if err != nil {
|
||||
f.Fatal(err)
|
||||
}
|
||||
conn, err := transport.DialGRPCInsecure(ctx,
|
||||
conn, err := gtransport.DialInsecure(ctx,
|
||||
option.WithEndpoint(srv.Addr),
|
||||
option.WithGRPCConnectionPool(nConns))
|
||||
option.WithGRPCConnectionPool(nConns),
|
||||
|
||||
// TODO(grpc/grpc-go#1388) using connection pool without WithBlock
|
||||
// can cause RPCs to fail randomly. We can delete this after the issue is fixed.
|
||||
option.WithGRPCDialOption(grpc.WithBlock()))
|
||||
if err != nil {
|
||||
f.Fatal(err)
|
||||
}
|
||||
|
||||
+2
@@ -22,6 +22,7 @@ import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"log"
|
||||
"runtime"
|
||||
"strconv"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -150,6 +151,7 @@ func (s *SubServer) Start(ctx context.Context, req *pb.StartRequest) (*pb.StartR
|
||||
// Load test API doesn't define any way to stop right now.
|
||||
go func() {
|
||||
sub := c.Subscription(req.GetPubsubOptions().Subscription)
|
||||
sub.ReceiveSettings.NumGoroutines = 10 * runtime.GOMAXPROCS(0)
|
||||
err := sub.Receive(context.Background(), s.callback)
|
||||
log.Fatal(err)
|
||||
}()
|
||||
|
||||
+10
@@ -18,10 +18,12 @@ import (
|
||||
"fmt"
|
||||
"os"
|
||||
"runtime"
|
||||
"time"
|
||||
|
||||
"google.golang.org/api/iterator"
|
||||
"google.golang.org/api/option"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/keepalive"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
@@ -62,6 +64,14 @@ func NewClient(ctx context.Context, projectID string, opts ...option.ClientOptio
|
||||
o = []option.ClientOption{
|
||||
// Create multiple connections to increase throughput.
|
||||
option.WithGRPCConnectionPool(runtime.GOMAXPROCS(0)),
|
||||
|
||||
// TODO(grpc/grpc-go#1388) using connection pool without WithBlock
|
||||
// can cause RPCs to fail randomly. We can delete this after the issue is fixed.
|
||||
option.WithGRPCDialOption(grpc.WithBlock()),
|
||||
|
||||
option.WithGRPCDialOption(grpc.WithKeepaliveParams(keepalive.ClientParameters{
|
||||
Time: 5 * time.Minute,
|
||||
})),
|
||||
}
|
||||
}
|
||||
o = append(o, opts...)
|
||||
|
||||
-115
@@ -1,115 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
// puller fetches messages from the server in a batch.
|
||||
type puller struct {
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
|
||||
// keepAlive takes ownership of the lifetime of the message identified
|
||||
// by ackID, ensuring that its ack deadline does not expire. It should
|
||||
// be called each time a new message is fetched from the server, even
|
||||
// if it is not yet returned from Next.
|
||||
keepAlive func(ackID string)
|
||||
|
||||
// abandon should be called for each message which has previously been
|
||||
// passed to keepAlive, but will never be returned by Next.
|
||||
abandon func(ackID string)
|
||||
|
||||
// fetch fetches a batch of messages from the server.
|
||||
fetch func() ([]*Message, error)
|
||||
|
||||
mu sync.Mutex
|
||||
buf []*Message
|
||||
}
|
||||
|
||||
// newPuller constructs a new puller.
|
||||
// batchSize is the maximum number of messages to fetch at once.
|
||||
// No more than batchSize messages will be outstanding at any time.
|
||||
func newPuller(s service, subName string, ctx context.Context, batchSize int32, keepAlive, abandon func(ackID string)) *puller {
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
return &puller{
|
||||
cancel: cancel,
|
||||
keepAlive: keepAlive,
|
||||
abandon: abandon,
|
||||
ctx: ctx,
|
||||
fetch: func() ([]*Message, error) { return s.fetchMessages(ctx, subName, batchSize) },
|
||||
}
|
||||
}
|
||||
|
||||
const maxPullAttempts = 2
|
||||
|
||||
// Next returns the next message from the server, fetching a new batch if necessary.
|
||||
// keepAlive is called with the ackIDs of newly fetched messages.
|
||||
// If p.Ctx has already been cancelled before Next is called, no new messages
|
||||
// will be fetched.
|
||||
func (p *puller) Next() (*Message, error) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
// If ctx has been cancelled, return straight away (even if there are buffered messages available).
|
||||
select {
|
||||
case <-p.ctx.Done():
|
||||
return nil, p.ctx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
for len(p.buf) == 0 {
|
||||
var buf []*Message
|
||||
var err error
|
||||
|
||||
for i := 0; i < maxPullAttempts; i++ {
|
||||
// Once Stop has completed, all future calls to Next will immediately fail at this point.
|
||||
buf, err = p.fetch()
|
||||
if err == nil || err == context.Canceled || err == context.DeadlineExceeded {
|
||||
break
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
for _, m := range buf {
|
||||
p.keepAlive(m.ackID)
|
||||
}
|
||||
p.buf = buf
|
||||
}
|
||||
|
||||
m := p.buf[0]
|
||||
p.buf = p.buf[1:]
|
||||
return m, nil
|
||||
}
|
||||
|
||||
// Stop aborts any pending calls to Next, and prevents any future ones from succeeding.
|
||||
// Stop also abandons any messages that have been pre-fetched.
|
||||
// Once Stop completes, no calls to Next will succeed.
|
||||
func (p *puller) Stop() {
|
||||
// Next may be executing in another goroutine. Cancel it, and then wait until it terminates.
|
||||
p.cancel()
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
for _, m := range p.buf {
|
||||
p.abandon(m.ackID)
|
||||
}
|
||||
p.buf = nil
|
||||
}
|
||||
-154
@@ -1,154 +0,0 @@
|
||||
// Copyright 2016 Google Inc. All Rights Reserved.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
type fetchResult struct {
|
||||
msgs []*Message
|
||||
err error
|
||||
}
|
||||
|
||||
type fetcherService struct {
|
||||
service
|
||||
results []fetchResult
|
||||
unexpectedCall bool
|
||||
}
|
||||
|
||||
func (s *fetcherService) fetchMessages(ctx context.Context, subName string, maxMessages int32) ([]*Message, error) {
|
||||
if len(s.results) == 0 {
|
||||
s.unexpectedCall = true
|
||||
return nil, errors.New("bang")
|
||||
}
|
||||
ret := s.results[0]
|
||||
s.results = s.results[1:]
|
||||
return ret.msgs, ret.err
|
||||
}
|
||||
|
||||
func TestPuller(t *testing.T) {
|
||||
s := &fetcherService{
|
||||
results: []fetchResult{
|
||||
{
|
||||
msgs: []*Message{{ackID: "a"}, {ackID: "b"}},
|
||||
},
|
||||
{},
|
||||
{
|
||||
msgs: []*Message{{ackID: "c"}, {ackID: "d"}},
|
||||
},
|
||||
{
|
||||
msgs: []*Message{{ackID: "e"}},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
pulled := make(chan string, 10)
|
||||
|
||||
pull := newPuller(s, "subname", context.Background(), 2, func(ackID string) { pulled <- ackID }, func(string) {})
|
||||
|
||||
got := []string{}
|
||||
for i := 0; i < 5; i++ {
|
||||
m, err := pull.Next()
|
||||
got = append(got, m.ackID)
|
||||
if err != nil {
|
||||
t.Errorf("unexpected err from pull.Next: %v", err)
|
||||
}
|
||||
}
|
||||
_, err := pull.Next()
|
||||
if err == nil {
|
||||
t.Errorf("unexpected err from pull.Next: %v", err)
|
||||
}
|
||||
|
||||
want := []string{"a", "b", "c", "d", "e"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("pulled ack ids: got: %v ; want: %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPullerAddsToKeepAlive(t *testing.T) {
|
||||
s := &fetcherService{
|
||||
results: []fetchResult{
|
||||
{
|
||||
msgs: []*Message{{ackID: "a"}, {ackID: "b"}},
|
||||
},
|
||||
{
|
||||
msgs: []*Message{{ackID: "c"}, {ackID: "d"}},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
pulled := make(chan string, 10)
|
||||
|
||||
pull := newPuller(s, "subname", context.Background(), 2, func(ackID string) { pulled <- ackID }, func(string) {})
|
||||
|
||||
got := []string{}
|
||||
for i := 0; i < 3; i++ {
|
||||
m, err := pull.Next()
|
||||
got = append(got, m.ackID)
|
||||
if err != nil {
|
||||
t.Errorf("unexpected err from pull.Next: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
want := []string{"a", "b", "c"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("pulled ack ids: got: %v ; want: %v", got, want)
|
||||
}
|
||||
|
||||
close(pulled)
|
||||
// We should have seen "d" written to the channel too, even though it hasn't been returned yet.
|
||||
pulledIDs := []string{}
|
||||
for id := range pulled {
|
||||
pulledIDs = append(pulledIDs, id)
|
||||
}
|
||||
|
||||
want = append(want, "d")
|
||||
if !reflect.DeepEqual(pulledIDs, want) {
|
||||
t.Errorf("pulled ack ids: got: %v ; want: %v", pulledIDs, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPullerRetriesOnce(t *testing.T) {
|
||||
bang := errors.New("bang")
|
||||
s := &fetcherService{
|
||||
results: []fetchResult{
|
||||
{
|
||||
err: bang,
|
||||
},
|
||||
{
|
||||
err: bang,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
pull := newPuller(s, "subname", context.Background(), 2, func(string) {}, func(string) {})
|
||||
|
||||
_, err := pull.Next()
|
||||
if err != bang {
|
||||
t.Errorf("pull.Next err got: %v, want: %v", err, bang)
|
||||
}
|
||||
|
||||
if s.unexpectedCall {
|
||||
t.Errorf("unexpected retry")
|
||||
}
|
||||
if len(s.results) != 0 {
|
||||
t.Errorf("outstanding calls: got: %v, want: 0", len(s.results))
|
||||
}
|
||||
}
|
||||
+2
-2
@@ -384,8 +384,9 @@ func (p *streamingPuller) openLocked() {
|
||||
return
|
||||
}
|
||||
// No opens in flight; start one.
|
||||
// Keep the lock held, to avoid a race where we
|
||||
// close the old stream while opening a new one.
|
||||
p.inFlight = true
|
||||
p.c.L.Unlock()
|
||||
spc, err := p.subc.StreamingPull(p.ctx, gax.WithGRPCOptions(grpc.MaxCallRecvMsgSize(maxSendRecvBytes)))
|
||||
if err == nil {
|
||||
err = spc.Send(&pb.StreamingPullRequest{
|
||||
@@ -393,7 +394,6 @@ func (p *streamingPuller) openLocked() {
|
||||
StreamAckDeadlineSeconds: p.ackDeadlineSecs,
|
||||
})
|
||||
}
|
||||
p.c.L.Lock()
|
||||
p.spc = spc
|
||||
p.err = err
|
||||
p.inFlight = false
|
||||
|
||||
+6
-23
@@ -58,9 +58,6 @@ func TestStreamingPullMultipleFetches(t *testing.T) {
|
||||
}
|
||||
|
||||
func testStreamingPullIteration(t *testing.T, client *Client, server *fakeServer, msgs []*pb.ReceivedMessage) {
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
sub := client.Subscription("s")
|
||||
gotMsgs, err := pullN(context.Background(), sub, len(msgs), func(_ context.Context, m *Message) {
|
||||
id, err := strconv.Atoi(m.ackID)
|
||||
@@ -116,13 +113,13 @@ func TestStreamingPullError(t *testing.T) {
|
||||
// If an RPC to the service returns a non-retryable error, Pull should
|
||||
// return after all callbacks return, without waiting for messages to be
|
||||
// acked.
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
client, server := newFake(t)
|
||||
server.addStreamingPullMessages(testMessages[:1])
|
||||
server.addStreamingPullError(grpc.Errorf(codes.Internal, ""))
|
||||
server.addStreamingPullError(grpc.Errorf(codes.Unknown, ""))
|
||||
sub := client.Subscription("s")
|
||||
// Use only one goroutine, since the fake server is configured to
|
||||
// return only one error.
|
||||
sub.ReceiveSettings.NumGoroutines = 1
|
||||
callbackDone := make(chan struct{})
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
err := sub.Receive(ctx, func(ctx context.Context, m *Message) {
|
||||
@@ -137,7 +134,7 @@ func TestStreamingPullError(t *testing.T) {
|
||||
default:
|
||||
t.Fatal("Receive returned but callback was not done")
|
||||
}
|
||||
if want := codes.Internal; grpc.Code(err) != want {
|
||||
if want := codes.Unknown; grpc.Code(err) != want {
|
||||
t.Fatalf("got <%v>, want code %v", err, want)
|
||||
}
|
||||
}
|
||||
@@ -145,9 +142,6 @@ func TestStreamingPullError(t *testing.T) {
|
||||
func TestStreamingPullCancel(t *testing.T) {
|
||||
// If Receive's context is canceled, it should return after all callbacks
|
||||
// return and all messages have been acked.
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
client, server := newFake(t)
|
||||
server.addStreamingPullMessages(testMessages)
|
||||
sub := client.Subscription("s")
|
||||
@@ -157,6 +151,7 @@ func TestStreamingPullCancel(t *testing.T) {
|
||||
atomic.AddInt32(&n, 1)
|
||||
defer atomic.AddInt32(&n, -1)
|
||||
cancel()
|
||||
m.Ack()
|
||||
})
|
||||
if got := atomic.LoadInt32(&n); got != 0 {
|
||||
t.Errorf("Receive returned with %d callbacks still running", got)
|
||||
@@ -167,9 +162,6 @@ func TestStreamingPullCancel(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestStreamingPullRetry(t *testing.T) {
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
// Check that we retry on io.EOF or Unavailable.
|
||||
client, server := newFake(t)
|
||||
server.addStreamingPullMessages(testMessages[:1])
|
||||
@@ -185,9 +177,6 @@ func TestStreamingPullRetry(t *testing.T) {
|
||||
|
||||
func TestStreamingPullOneActive(t *testing.T) {
|
||||
// Only one call to Pull can be active at a time.
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
client, srv := newFake(t)
|
||||
srv.addStreamingPullMessages(testMessages[:1])
|
||||
sub := client.Subscription("s")
|
||||
@@ -206,9 +195,6 @@ func TestStreamingPullOneActive(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestStreamingPullConcurrent(t *testing.T) {
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
newMsg := func(i int) *pb.ReceivedMessage {
|
||||
return &pb.ReceivedMessage{
|
||||
AckId: strconv.Itoa(i),
|
||||
@@ -245,9 +231,6 @@ func TestStreamingPullConcurrent(t *testing.T) {
|
||||
|
||||
func TestStreamingPullFlowControl(t *testing.T) {
|
||||
// Callback invocations should not occur if flow control limits are exceeded.
|
||||
if !useStreamingPull {
|
||||
t.SkipNow()
|
||||
}
|
||||
client, server := newFake(t)
|
||||
server.addStreamingPullMessages(testMessages)
|
||||
sub := client.Subscription("s")
|
||||
|
||||
+46
-25
@@ -17,7 +17,7 @@ package pubsub
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"runtime"
|
||||
"io"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -25,7 +25,6 @@ import (
|
||||
"cloud.google.com/go/iam"
|
||||
"golang.org/x/net/context"
|
||||
"golang.org/x/sync/errgroup"
|
||||
"google.golang.org/api/iterator"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
)
|
||||
@@ -153,6 +152,12 @@ type ReceiveSettings struct {
|
||||
// NumGoroutines is the number of goroutines Receive will spawn to pull
|
||||
// messages concurrently. If NumGoroutines is less than 1, it will be treated
|
||||
// as if it were DefaultReceiveSettings.NumGoroutines.
|
||||
//
|
||||
// NumGoroutines does not limit the number of messages that can be processed
|
||||
// concurrently. Even with one goroutine, many messages might be processed at
|
||||
// once, because that goroutine may continually receive messages and invoke the
|
||||
// function passed to Receive on them. To limit the number of messages being
|
||||
// processed concurrently, set MaxOutstandingMessages.
|
||||
NumGoroutines int
|
||||
}
|
||||
|
||||
@@ -161,7 +166,7 @@ var DefaultReceiveSettings = ReceiveSettings{
|
||||
MaxExtension: 10 * time.Minute,
|
||||
MaxOutstandingMessages: 1000,
|
||||
MaxOutstandingBytes: 1e9, // 1G
|
||||
NumGoroutines: 10 * runtime.GOMAXPROCS(0),
|
||||
NumGoroutines: 1,
|
||||
}
|
||||
|
||||
// Delete deletes the subscription.
|
||||
@@ -262,7 +267,7 @@ var errReceiveInProgress = errors.New("pubsub: Receive already in progress for t
|
||||
//
|
||||
// If the service returns a non-retryable error, Receive returns that error after
|
||||
// all of the outstanding calls to f have returned. If ctx is done, Receive
|
||||
// returns either nil after all of the outstanding calls to f have returned and
|
||||
// returns nil after all of the outstanding calls to f have returned and
|
||||
// all messages have been acknowledged or have expired.
|
||||
//
|
||||
// Receive calls f concurrently from multiple goroutines. It is encouraged to
|
||||
@@ -326,13 +331,13 @@ func (s *Subscription) Receive(ctx context.Context, f func(context.Context, *Mes
|
||||
group, gctx := errgroup.WithContext(ctx)
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
group.Go(func() error {
|
||||
return s.receive(gctx, group, po, fc, f)
|
||||
return s.receive(gctx, po, fc, f)
|
||||
})
|
||||
}
|
||||
return group.Wait()
|
||||
}
|
||||
|
||||
func (s *Subscription) receive(ctx context.Context, group *errgroup.Group, po *pullOptions, fc *flowController, f func(context.Context, *Message)) error {
|
||||
func (s *Subscription) receive(ctx context.Context, po *pullOptions, fc *flowController, f func(context.Context, *Message)) error {
|
||||
// Cancel a sub-context when we return, to kick the context-aware callbacks
|
||||
// and the goroutine below.
|
||||
ctx2, cancel := context.WithCancel(ctx)
|
||||
@@ -343,34 +348,50 @@ func (s *Subscription) receive(ctx context.Context, group *errgroup.Group, po *p
|
||||
// that context would immediately stop the iterator without waiting for unacked
|
||||
// messages.
|
||||
iter := newMessageIterator(context.Background(), s.s, s.name, po)
|
||||
group.Go(func() error {
|
||||
|
||||
// We cannot use errgroup from Receive here. Receive might already be calling group.Wait,
|
||||
// and group.Wait cannot be called concurrently with group.Go. We give each receive() its
|
||||
// own WaitGroup instead.
|
||||
// Since wg.Add is only called from the main goroutine, wg.Wait is guaranteed
|
||||
// to be called after all Adds.
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
<-ctx2.Done()
|
||||
iter.Stop()
|
||||
return nil
|
||||
})
|
||||
iter.stop()
|
||||
wg.Done()
|
||||
}()
|
||||
defer wg.Wait()
|
||||
|
||||
defer cancel()
|
||||
for {
|
||||
msg, err := iter.Next()
|
||||
if err == iterator.Done {
|
||||
msgs, err := iter.receive()
|
||||
if err == io.EOF {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// TODO(jba): call acquire closer to when the message is allocated.
|
||||
if err := fc.acquire(ctx, len(msg.Data)); err != nil {
|
||||
// TODO(jba): test that this "orphaned" message is nacked immediately when ctx is done.
|
||||
msg.Nack()
|
||||
return nil
|
||||
for i, msg := range msgs {
|
||||
msg := msg
|
||||
// TODO(jba): call acquire closer to when the message is allocated.
|
||||
if err := fc.acquire(ctx, len(msg.Data)); err != nil {
|
||||
// TODO(jba): test that these "orphaned" messages are nacked immediately when ctx is done.
|
||||
for _, m := range msgs[i:] {
|
||||
m.Nack()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
// TODO(jba): call release when the message is available for GC.
|
||||
// This considers the message to be released when
|
||||
// f is finished, but f may ack early or not at all.
|
||||
defer wg.Done()
|
||||
defer fc.release(len(msg.Data))
|
||||
f(ctx2, msg)
|
||||
}()
|
||||
}
|
||||
group.Go(func() error {
|
||||
// TODO(jba): call release when the message is available for GC.
|
||||
// This considers the message to be released when
|
||||
// f is finished, but f may ack early or not at all.
|
||||
defer fc.release(len(msg.Data))
|
||||
f(ctx2, msg)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+12
-2
@@ -106,14 +106,24 @@ func (c *Client) CreateTopic(ctx context.Context, id string) (*Topic, error) {
|
||||
return t, err
|
||||
}
|
||||
|
||||
// Topic creates a reference to a topic.
|
||||
// Topic creates a reference to a topic in the client's project.
|
||||
//
|
||||
// If a Topic's Publish method is called, it has background goroutines
|
||||
// associated with it. Clean them up by calling Topic.Stop.
|
||||
//
|
||||
// Avoid creating many Topic instances if you use them to publish.
|
||||
func (c *Client) Topic(id string) *Topic {
|
||||
return newTopic(c.s, fmt.Sprintf("projects/%s/topics/%s", c.projectID, id))
|
||||
return c.TopicInProject(id, c.projectID)
|
||||
}
|
||||
|
||||
// TopicInProject creates a reference to a topic in the given project.
|
||||
//
|
||||
// If a Topic's Publish method is called, it has background goroutines
|
||||
// associated with it. Clean them up by calling Topic.Stop.
|
||||
//
|
||||
// Avoid creating many Topic instances if you use them to publish.
|
||||
func (c *Client) TopicInProject(id, projectID string) *Topic {
|
||||
return newTopic(c.s, fmt.Sprintf("projects/%s/topics/%s", projectID, id))
|
||||
}
|
||||
|
||||
func newTopic(s service, name string) *Topic {
|
||||
|
||||
Reference in New Issue
Block a user