vendor: update all dependencies
This commit is contained in:
+162
-241
@@ -15,8 +15,11 @@
|
||||
package pubsub
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"fmt"
|
||||
"log"
|
||||
"math/rand"
|
||||
"os"
|
||||
"reflect"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -25,151 +28,199 @@ import (
|
||||
"golang.org/x/net/context"
|
||||
|
||||
"cloud.google.com/go/internal/testutil"
|
||||
"google.golang.org/api/iterator"
|
||||
"google.golang.org/api/option"
|
||||
)
|
||||
|
||||
const timeout = time.Minute * 10
|
||||
const ackDeadline = time.Second * 10
|
||||
|
||||
const batchSize = 100
|
||||
const batches = 100
|
||||
const nMessages = 1e4
|
||||
|
||||
// messageCounter keeps track of how many times a given message has been received.
|
||||
type messageCounter struct {
|
||||
mu sync.Mutex
|
||||
counts map[string]int
|
||||
// A value is sent to recv each time Inc is called.
|
||||
recv chan struct{}
|
||||
}
|
||||
// Buffer log messages to debug failures.
|
||||
var logBuf bytes.Buffer
|
||||
|
||||
func (mc *messageCounter) Inc(msgID string) {
|
||||
mc.mu.Lock()
|
||||
mc.counts[msgID] += 1
|
||||
mc.mu.Unlock()
|
||||
mc.recv <- struct{}{}
|
||||
}
|
||||
|
||||
// process pulls messages from an iterator and records them in mc.
|
||||
func process(t *testing.T, it *MessageIterator, mc *messageCounter) {
|
||||
for {
|
||||
m, err := it.Next()
|
||||
if err == iterator.Done {
|
||||
return
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
t.Errorf("unexpected err from iterator: %v", err)
|
||||
return
|
||||
}
|
||||
mc.Inc(m.ID)
|
||||
// Simulate time taken to process m, while continuing to process more messages.
|
||||
go func() {
|
||||
// Some messages will need to have their ack deadline extended due to this delay.
|
||||
delay := rand.Intn(int(ackDeadline * 3))
|
||||
time.After(time.Duration(delay))
|
||||
m.Done(true)
|
||||
}()
|
||||
// TestEndToEnd pumps many messages into a topic and tests that they are all
|
||||
// delivered to each subscription for the topic. It also tests that messages
|
||||
// are not unexpectedly redelivered.
|
||||
func TestEndToEnd(t *testing.T) {
|
||||
t.Parallel()
|
||||
if testing.Short() {
|
||||
t.Skip("Integration tests skipped in short mode")
|
||||
}
|
||||
log.SetOutput(&logBuf)
|
||||
ctx := context.Background()
|
||||
ts := testutil.TokenSource(ctx, ScopePubSub, ScopeCloudPlatform)
|
||||
if ts == nil {
|
||||
t.Skip("Integration tests skipped. See CONTRIBUTING.md for details")
|
||||
}
|
||||
}
|
||||
|
||||
// newIter constructs a new MessageIterator.
|
||||
func newIter(t *testing.T, ctx context.Context, sub *Subscription) *MessageIterator {
|
||||
it, err := sub.Pull(ctx)
|
||||
now := time.Now()
|
||||
topicName := fmt.Sprintf("endtoend-%d", now.UnixNano())
|
||||
subPrefix := fmt.Sprintf("endtoend-%d", now.UnixNano())
|
||||
|
||||
client, err := NewClient(ctx, testutil.ProjID(), option.WithTokenSource(ts))
|
||||
if err != nil {
|
||||
t.Fatalf("error constructing iterator: %v", err)
|
||||
t.Fatalf("Creating client error: %v", err)
|
||||
}
|
||||
return it
|
||||
}
|
||||
|
||||
// launchIter launches a number of goroutines to pull from the supplied MessageIterator.
|
||||
func launchIter(t *testing.T, ctx context.Context, it *MessageIterator, mc *messageCounter, n int, wg *sync.WaitGroup) {
|
||||
for j := 0; j < n; j++ {
|
||||
var topic *Topic
|
||||
if topic, err = client.CreateTopic(ctx, topicName); err != nil {
|
||||
t.Fatalf("CreateTopic error: %v", err)
|
||||
}
|
||||
defer topic.Delete(ctx)
|
||||
|
||||
// Two subscriptions to the same topic.
|
||||
var subs [2]*Subscription
|
||||
for i := 0; i < len(subs); i++ {
|
||||
subs[i], err = client.CreateSubscription(ctx, fmt.Sprintf("%s-%d", subPrefix, i), SubscriptionConfig{
|
||||
Topic: topic,
|
||||
AckDeadline: ackDeadline,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("CreateSub error: %v", err)
|
||||
}
|
||||
defer subs[i].Delete(ctx)
|
||||
}
|
||||
|
||||
ids, err := publish(ctx, topic, nMessages)
|
||||
topic.Stop()
|
||||
if err != nil {
|
||||
t.Fatalf("publish: %v", err)
|
||||
}
|
||||
wantCounts := make(map[string]int)
|
||||
for _, id := range ids {
|
||||
wantCounts[id] = 1
|
||||
}
|
||||
|
||||
// recv provides an indication that messages are still arriving.
|
||||
recv := make(chan struct{})
|
||||
// We have two subscriptions to our topic.
|
||||
// Each subscription will get a copy of each published message.
|
||||
var wg sync.WaitGroup
|
||||
cctx, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
|
||||
consumers := []*consumer{
|
||||
{counts: make(map[string]int), recv: recv, durations: []time.Duration{time.Hour}},
|
||||
{counts: make(map[string]int), recv: recv,
|
||||
durations: []time.Duration{ackDeadline, ackDeadline, ackDeadline / 2, ackDeadline / 2, time.Hour}},
|
||||
}
|
||||
for i, con := range consumers {
|
||||
con := con
|
||||
sub := subs[i]
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
process(t, it, mc)
|
||||
con.consume(t, cctx, sub)
|
||||
}()
|
||||
}
|
||||
}
|
||||
// Wait for a while after the last message before declaring quiescence.
|
||||
// We wait a multiple of the ack deadline, for two reasons:
|
||||
// 1. To detect if messages are redelivered after having their ack
|
||||
// deadline extended.
|
||||
// 2. To wait for redelivery of messages that were en route when a Receive
|
||||
// is canceled. This can take considerably longer than the ack deadline.
|
||||
quiescenceDur := ackDeadline * 6
|
||||
quiescenceTimer := time.NewTimer(quiescenceDur)
|
||||
|
||||
// iteratorLifetime controls how long iterators live for before they are stopped.
|
||||
type iteratorLifetimes interface {
|
||||
// lifetimeChan should be called when an iterator is started. The
|
||||
// returned channel will send when the iterator should be stopped.
|
||||
lifetimeChan() <-chan time.Time
|
||||
}
|
||||
loop:
|
||||
for {
|
||||
select {
|
||||
case <-recv:
|
||||
// Reset timer so we wait quiescenceDur after the last message.
|
||||
// See https://godoc.org/time#Timer.Reset for why the Stop
|
||||
// and channel drain are necessary.
|
||||
if !quiescenceTimer.Stop() {
|
||||
<-quiescenceTimer.C
|
||||
}
|
||||
quiescenceTimer.Reset(quiescenceDur)
|
||||
|
||||
var immortal = &explicitLifetimes{}
|
||||
case <-quiescenceTimer.C:
|
||||
cancel()
|
||||
log.Println("quiesced")
|
||||
break loop
|
||||
|
||||
// explicitLifetimes implements iteratorLifetime with hard-coded lifetimes, falling back
|
||||
// to indefinite lifetimes when no explicit lifetimes remain.
|
||||
type explicitLifetimes struct {
|
||||
mu sync.Mutex
|
||||
lifetimes []time.Duration
|
||||
}
|
||||
|
||||
func (el *explicitLifetimes) lifetimeChan() <-chan time.Time {
|
||||
el.mu.Lock()
|
||||
defer el.mu.Unlock()
|
||||
if len(el.lifetimes) == 0 {
|
||||
return nil
|
||||
case <-cctx.Done():
|
||||
t.Fatal("timed out")
|
||||
}
|
||||
}
|
||||
lifetime := el.lifetimes[0]
|
||||
el.lifetimes = el.lifetimes[1:]
|
||||
return time.After(lifetime)
|
||||
wg.Wait()
|
||||
ok := true
|
||||
for i, con := range consumers {
|
||||
if got, want := con.counts, wantCounts; !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("%d: message counts: %v\n", i, diff(got, want))
|
||||
ok = false
|
||||
}
|
||||
}
|
||||
if !ok {
|
||||
logBuf.WriteTo(os.Stdout)
|
||||
}
|
||||
}
|
||||
|
||||
// publish publishes n messages to topic, and returns the published message IDs.
|
||||
func publish(ctx context.Context, topic *Topic, n int) ([]string, error) {
|
||||
var rs []*PublishResult
|
||||
for i := 0; i < n; i++ {
|
||||
m := &Message{Data: []byte(fmt.Sprintf("msg %d", i))}
|
||||
rs = append(rs, topic.Publish(ctx, m))
|
||||
}
|
||||
var ids []string
|
||||
for _, r := range rs {
|
||||
id, err := r.Get(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ids = append(ids, id)
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
// consumer consumes messages according to its configuration.
|
||||
type consumer struct {
|
||||
// How many goroutines should pull from the subscription.
|
||||
iteratorsInFlight int
|
||||
// How many goroutines should pull from each iterator.
|
||||
concurrencyPerIterator int
|
||||
durations []time.Duration
|
||||
|
||||
lifetimes iteratorLifetimes
|
||||
// A value is sent to recv each time Inc is called.
|
||||
recv chan struct{}
|
||||
|
||||
mu sync.Mutex
|
||||
counts map[string]int
|
||||
total int
|
||||
}
|
||||
|
||||
// consume reads messages from a subscription, and keeps track of what it receives in mc.
|
||||
// After consume returns, the caller should wait on wg to ensure that no more updates to mc will be made.
|
||||
func (c *consumer) consume(t *testing.T, ctx context.Context, sub *Subscription, mc *messageCounter, wg *sync.WaitGroup, stop <-chan struct{}) {
|
||||
for i := 0; i < c.iteratorsInFlight; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for {
|
||||
it := newIter(t, ctx, sub)
|
||||
launchIter(t, ctx, it, mc, c.concurrencyPerIterator, wg)
|
||||
|
||||
select {
|
||||
case <-c.lifetimes.lifetimeChan():
|
||||
it.Stop()
|
||||
case <-stop:
|
||||
it.Stop()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
}()
|
||||
func (c *consumer) consume(t *testing.T, ctx context.Context, sub *Subscription) {
|
||||
for _, dur := range c.durations {
|
||||
ctx2, cancel := context.WithTimeout(ctx, dur)
|
||||
defer cancel()
|
||||
id := sub.name[len(sub.name)-2:]
|
||||
log.Printf("%s: start receive", id)
|
||||
prev := c.total
|
||||
err := sub.Receive(ctx2, c.process)
|
||||
log.Printf("%s: end receive; read %d", id, c.total-prev)
|
||||
if err != nil {
|
||||
t.Errorf("error from Receive: %v", err)
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// publish publishes many messages to topic, and returns the published message ids.
|
||||
func publish(t *testing.T, ctx context.Context, topic *Topic) []string {
|
||||
var published []string
|
||||
msgs := make([]*Message, batchSize)
|
||||
for i := 0; i < batches; i++ {
|
||||
for j := 0; j < batchSize; j++ {
|
||||
text := fmt.Sprintf("msg %02d-%02d", i, j)
|
||||
msgs[j] = &Message{Data: []byte(text)}
|
||||
}
|
||||
ids, err := topic.Publish(ctx, msgs...)
|
||||
if err != nil {
|
||||
t.Errorf("Publish error: %v", err)
|
||||
}
|
||||
published = append(published, ids...)
|
||||
}
|
||||
return published
|
||||
// process handles a message and records it in mc.
|
||||
func (c *consumer) process(_ context.Context, m *Message) {
|
||||
c.mu.Lock()
|
||||
c.counts[m.ID] += 1
|
||||
c.total++
|
||||
c.mu.Unlock()
|
||||
c.recv <- struct{}{}
|
||||
// Simulate time taken to process m, while continuing to process more messages.
|
||||
// Some messages will need to have their ack deadline extended due to this delay.
|
||||
delay := rand.Intn(int(ackDeadline * 3))
|
||||
time.AfterFunc(time.Duration(delay), m.Ack)
|
||||
}
|
||||
|
||||
// diff returns counts of the differences between got and want.
|
||||
@@ -192,133 +243,3 @@ func diff(got, want map[string]int) map[string]int {
|
||||
}
|
||||
return gotWantCount
|
||||
}
|
||||
|
||||
// TestEndToEnd pumps many messages into a topic and tests that they are all delivered to each subscription for the topic.
|
||||
// It also tests that messages are not unexpectedly redelivered.
|
||||
func TestEndToEnd(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Integration tests skipped in short mode")
|
||||
}
|
||||
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("endtoend-%d", now.Unix())
|
||||
subPrefix := fmt.Sprintf("endtoend-%d", now.Unix())
|
||||
|
||||
client, err := NewClient(ctx, testutil.ProjID(), option.WithTokenSource(ts))
|
||||
if err != nil {
|
||||
t.Fatalf("Creating client error: %v", err)
|
||||
}
|
||||
|
||||
var topic *Topic
|
||||
if topic, err = client.CreateTopic(ctx, topicName); err != nil {
|
||||
t.Fatalf("CreateTopic error: %v", err)
|
||||
}
|
||||
defer topic.Delete(ctx)
|
||||
|
||||
// Three subscriptions to the same topic.
|
||||
var subA, subB, subC *Subscription
|
||||
if subA, err = client.CreateSubscription(ctx, subPrefix+"-a", topic, ackDeadline, nil); err != nil {
|
||||
t.Fatalf("CreateSub error: %v", err)
|
||||
}
|
||||
defer subA.Delete(ctx)
|
||||
|
||||
if subB, err = client.CreateSubscription(ctx, subPrefix+"-b", topic, ackDeadline, nil); err != nil {
|
||||
t.Fatalf("CreateSub error: %v", err)
|
||||
}
|
||||
defer subB.Delete(ctx)
|
||||
|
||||
if subC, err = client.CreateSubscription(ctx, subPrefix+"-c", topic, ackDeadline, nil); err != nil {
|
||||
t.Fatalf("CreateSub error: %v", err)
|
||||
}
|
||||
defer subC.Delete(ctx)
|
||||
|
||||
expectedCounts := make(map[string]int)
|
||||
for _, id := range publish(t, ctx, topic) {
|
||||
expectedCounts[id] = 1
|
||||
}
|
||||
|
||||
// recv provides an indication that messages are still arriving.
|
||||
recv := make(chan struct{})
|
||||
|
||||
// Keep track of the number of times each message (by message id) was
|
||||
// seen from each subscription.
|
||||
mcA := &messageCounter{counts: make(map[string]int), recv: recv}
|
||||
mcB := &messageCounter{counts: make(map[string]int), recv: recv}
|
||||
mcC := &messageCounter{counts: make(map[string]int), recv: recv}
|
||||
|
||||
stopC := make(chan struct{})
|
||||
|
||||
// We have three subscriptions to our topic.
|
||||
// Each subscription will get a copy of each pulished message.
|
||||
//
|
||||
// subA has just one iterator, while subB has two. The subB iterators
|
||||
// will each process roughly half of the messages for subB. All of
|
||||
// these iterators live until all messages have been consumed. subC is
|
||||
// processed by a series of short-lived iterators.
|
||||
|
||||
var wg sync.WaitGroup
|
||||
|
||||
con := &consumer{
|
||||
concurrencyPerIterator: 1,
|
||||
iteratorsInFlight: 2,
|
||||
lifetimes: immortal,
|
||||
}
|
||||
con.consume(t, ctx, subA, mcA, &wg, stopC)
|
||||
|
||||
con = &consumer{
|
||||
concurrencyPerIterator: 1,
|
||||
iteratorsInFlight: 2,
|
||||
lifetimes: immortal,
|
||||
}
|
||||
con.consume(t, ctx, subB, mcB, &wg, stopC)
|
||||
|
||||
con = &consumer{
|
||||
concurrencyPerIterator: 1,
|
||||
iteratorsInFlight: 2,
|
||||
lifetimes: &explicitLifetimes{
|
||||
lifetimes: []time.Duration{ackDeadline, ackDeadline, ackDeadline / 2, ackDeadline / 2},
|
||||
},
|
||||
}
|
||||
con.consume(t, ctx, subC, mcC, &wg, stopC)
|
||||
|
||||
go func() {
|
||||
timeoutC := time.After(timeout)
|
||||
// Every time this ticker ticks, we will check if we have received any
|
||||
// messages since the last time it ticked. We check less frequently
|
||||
// than the ack deadline, so that we can detect if messages are
|
||||
// redelivered after having their ack deadline extended.
|
||||
checkQuiescence := time.NewTicker(ackDeadline * 3)
|
||||
defer checkQuiescence.Stop()
|
||||
|
||||
var received bool
|
||||
for {
|
||||
select {
|
||||
case <-recv:
|
||||
received = true
|
||||
case <-checkQuiescence.C:
|
||||
if received {
|
||||
received = false
|
||||
} else {
|
||||
close(stopC)
|
||||
return
|
||||
}
|
||||
case <-timeoutC:
|
||||
t.Errorf("timed out")
|
||||
close(stopC)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
wg.Wait()
|
||||
|
||||
for _, mc := range []*messageCounter{mcA, mcB, mcC} {
|
||||
if got, want := mc.counts, expectedCounts; !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("message counts: %v\n", diff(got, want))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user