Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 44 additions & 4 deletions manager/logbroker/broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,28 @@ var (
errNotRunning = errors.New("broker is not running")
)

const (
// subscriptionQueueLimit bounds how many subscription events may be
// buffered for a single ListenSubscriptions watcher.
//
// A watcher stops draining its channel whenever stream.Send blocks, which
// happens when the agent's gRPC stream stalls -- an unresponsive or
// partitioned node, for example. Without a limit, every subscription
// published from that point on is retained by the watcher's queue, along
// with its SubscriptionMessage, LogSelector, LogSubscriptionOptions and
// cancel context, even after the subscription has been unregistered.
//
// Tearing the watcher down instead is recoverable: ListenSubscriptions
// returns an error, the agent reconnects, and watchSubscriptions replays
// the currently registered subscriptions for that node.
subscriptionQueueLimit = 1000

// logQueueLimit bounds how many log messages may be buffered for a single
// SubscribeLogs client. A client that cannot keep up has its log stream
// terminated rather than growing the manager's heap without bound.
logQueueLimit = 10000
)

type logMessage struct {
*api.PublishLogsMessage
completed bool
Expand Down Expand Up @@ -65,8 +87,12 @@ func (lb *LogBroker) Start(ctx context.Context) error {
}

lb.pctx, lb.cancelAll = context.WithCancel(ctx)
lb.logQueue = watch.NewQueue()
lb.subscriptionQueue = watch.NewQueue()
// Both queues are bounded and close their output channel on teardown, so a
// consumer that stops draining is disconnected instead of being buffered
// without bound. Callers must handle a closed channel; see SubscribeLogs
// and ListenSubscriptions.
lb.logQueue = watch.NewQueue(watch.WithLimit(logQueueLimit), watch.WithCloseOutChan())
lb.subscriptionQueue = watch.NewQueue(watch.WithLimit(subscriptionQueueLimit), watch.WithCloseOutChan())
lb.registeredSubscriptions = make(map[string]*subscription)
lb.subscriptionsByNode = make(map[string]map[*subscription]struct{})
return nil
Expand Down Expand Up @@ -257,7 +283,13 @@ func (lb *LogBroker) SubscribeLogs(request *api.SubscribeLogsRequest, stream api
return ctx.Err()
case <-pctx.Done():
return pctx.Err()
case event := <-publishCh:
case event, ok := <-publishCh:
if !ok {
// The queue tore the watcher down because this client fell
// further behind than logQueueLimit.
logger.Error("log stream terminated: client is too far behind")
return status.Errorf(codes.ResourceExhausted, "log stream terminated: client is too far behind")
}
publish := event.(*logMessage)
if publish.completed {
return publish.err
Expand Down Expand Up @@ -349,7 +381,15 @@ func (lb *LogBroker) ListenSubscriptions(_ *api.ListenSubscriptionsRequest, stre
// Send down new subscriptions.
for {
select {
case v := <-subscriptionCh:
case v, ok := <-subscriptionCh:
if !ok {
// The queue tore the watcher down because this node fell
// further behind than subscriptionQueueLimit. Returning an
// error lets the agent reconnect, at which point
// watchSubscriptions replays the current subscriptions.
logger.Error("subscription stream terminated: node is too far behind")
return status.Errorf(codes.ResourceExhausted, "subscription stream terminated: node is too far behind")
}
sub := v.(*subscription)

if sub.Closed() {
Expand Down
123 changes: 123 additions & 0 deletions manager/logbroker/broker_queue_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
package logbroker

import (
"context"
"runtime"
"sync/atomic"
"testing"
"time"

"github.com/moby/swarmkit/v2/api"
"github.com/moby/swarmkit/v2/manager/state/store"
"github.com/stretchr/testify/require"
)

// subscriptionRetention drives `count` subscriptions through their full,
// correct lifecycle (Run -> register -> unregister -> Stop) against a broker
// whose only ListenSubscriptions watcher behaves as described by `drain`, and
// reports how many of those subscriptions the runtime was able to reclaim
// afterwards.
//
// A subscription that has been unregistered is no longer referenced by any of
// the broker's bookkeeping maps, so a correctly behaving broker must allow it
// to be collected regardless of what the watcher is doing.
func subscriptionRetention(t *testing.T, count int, drain bool) (reclaimed int64) {
t.Helper()

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

s := store.NewMemoryStore(nil)
require.NotNil(t, s)
defer s.Close()

broker := New(s)
require.NoError(t, broker.Start(ctx))
defer broker.Stop()

const nodeID = "node-stalled"
broker.nodeConnected(nodeID)

// Stand in for ListenSubscriptions: it registers a watch on the
// subscriptionQueue and, in the failure case, stops reading from the
// channel. That is what happens in ListenSubscriptions when stream.Send
// blocks on a worker whose gRPC stream has stalled -- the loop never gets
// back around to `case v := <-subscriptionCh`.
_, subscriptionCh, cancelWatch := broker.watchSubscriptions(nodeID)
defer cancelWatch()

stopDraining := make(chan struct{})
defer close(stopDraining)
if drain {
go func() {
for {
select {
case <-subscriptionCh:
case <-stopDraining:
return
}
}
}()
}

var collected atomic.Int64
for range count {
sub := broker.newSubscription(
&api.LogSelector{NodeIDs: []string{nodeID}},
&api.LogSubscriptionOptions{},
)
runtime.SetFinalizer(sub, func(*subscription) { collected.Add(1) })

// Mirror SubscribeLogs exactly, including every cleanup step it
// performs on return.
sub.Run(ctx)
broker.registerSubscription(sub)
broker.unregisterSubscription(sub)
sub.Stop()
}

// Nothing in this function still references the subscriptions. Give the
// collector several opportunities to reclaim them and to run finalizers.
for range 5 {
runtime.GC()
time.Sleep(50 * time.Millisecond)
}
runtime.GC()
time.Sleep(100 * time.Millisecond)

return collected.Load()
}

// TestLogBrokerSubscriptionQueueBounded asserts that a ListenSubscriptions
// watcher which stops draining its channel cannot pin an unbounded number of
// subscriptions.
//
// A watcher stops draining whenever stream.Send blocks, which happens when the
// agent's gRPC stream stalls. registerSubscription and unregisterSubscription
// both publish the *subscription to subscriptionQueue, so an unbounded queue
// lets a single stalled watcher retain every subscription that has passed
// through the broker -- along with its SubscriptionMessage, LogSelector,
// LogSubscriptionOptions and cancel context -- even though each subscription
// was unregistered correctly.
//
// Because the backlog accumulates in a container/list rather than in blocked
// goroutines, such a leak is invisible to goroutine-count based monitoring and
// shows up only as unexplained heap growth in the manager.
func TestLogBrokerSubscriptionQueueBounded(t *testing.T) {
// Push well past the limit so a bounded queue has to shed load.
const count = 3 * subscriptionQueueLimit

t.Run("draining watcher", func(t *testing.T) {
reclaimed := subscriptionRetention(t, count, true)
t.Logf("reclaimed %d/%d subscriptions", reclaimed, count)
require.EqualValues(t, count, reclaimed,
"a watcher that drains its channel must not pin unregistered subscriptions")
})

t.Run("stalled watcher", func(t *testing.T) {
reclaimed := subscriptionRetention(t, count, false)
t.Logf("reclaimed %d/%d subscriptions", reclaimed, count)
require.GreaterOrEqual(t, reclaimed, int64(count-subscriptionQueueLimit),
"a stalled watcher must not retain more than subscriptionQueueLimit subscriptions")
})
}