From 914e81bab52621c2828c11e9b90efcc39bd7fb65 Mon Sep 17 00:00:00 2001 From: Ethan Dickson Date: Fri, 28 Aug 2026 09:05:04 +0000 Subject: [PATCH] fix(coderd): make chatd stream sync poller a callback fanout to fix send-on-closed-channel panic The stream sync poller delivered hints over per-subscriber channels that unregister closed under the poller mutex, while pollOnce sent on them from a lock-free snapshot. A subscriber unregistering mid-poll made pollOnce panic with "send on closed channel", crashing coderd. Any authenticated user could trigger this by opening and closing chat streams (GHSA-7x3x-59xg-4hrc). Convert the poller to a callback fanout, mirroring pubsub.Subscribe: Register now takes a deliver callback and returns only an unregister func, unregister is delete-only (nothing is ever closed), and the consumer funnels poller hints into its existing consumer-owned updateCh with the same streamCtx guard already used for pubsub hints. This makes the panic unrepresentable and also removes the nil-poller closed-channel return that instantly terminated streams via the ok-check. Adds a regression race test churning register/unregister against pollOnce; the previous implementation panics under it. --- coderd/x/chatd/stream_subscribe.go | 14 ++-- coderd/x/chatd/stream_sync_poller.go | 26 +++---- .../chatd/stream_sync_poller_internal_test.go | 74 +++++++++++++++++++ 3 files changed, 92 insertions(+), 22 deletions(-) create mode 100644 coderd/x/chatd/stream_sync_poller_internal_test.go diff --git a/coderd/x/chatd/stream_subscribe.go b/coderd/x/chatd/stream_subscribe.go index 9545b707c11..248122e2fd3 100644 --- a/coderd/x/chatd/stream_subscribe.go +++ b/coderd/x/chatd/stream_subscribe.go @@ -58,7 +58,12 @@ func (p *Server) subscribeStreamLoop( return subscribeWithInitialError(chatID, "failed to subscribe to chat updates") } - pollerCh, unregisterPoller := p.streamSyncPoller.Register(chatID) + unregisterPoller := p.streamSyncPoller.Register(chatID, func(hint streamSyncHint) { + select { + case updateCh <- hint: + case <-streamCtx.Done(): + } + }) loop := newStreamLoop(chat, p.db, logger, afterMessageID) // The immediate sync builds the initial snapshot returned to the caller // and the relay target for the forwarder. Hints only fire on state @@ -97,13 +102,6 @@ func (p *Server) subscribeStreamLoop( if !p.runStreamSync(streamCtx, loop, relay, events, hint) { return } - case hint, ok := <-pollerCh: - if !ok { - return - } - if !p.runStreamSync(streamCtx, loop, relay, events, hint) { - return - } case part, ok := <-relay.Parts(): if !ok { return diff --git a/coderd/x/chatd/stream_sync_poller.go b/coderd/x/chatd/stream_sync_poller.go index 11e9171687e..b3b1601fdc6 100644 --- a/coderd/x/chatd/stream_sync_poller.go +++ b/coderd/x/chatd/stream_sync_poller.go @@ -27,8 +27,8 @@ type streamSyncPoller struct { } type streamSyncPollerSubscriber struct { - chatID uuid.UUID - hints chan streamSyncHint + chatID uuid.UUID + deliver func(streamSyncHint) } func newStreamSyncPoller( @@ -68,15 +68,17 @@ func (p *streamSyncPoller) Close() { p.cancel() } -func (p *streamSyncPoller) Register(chatID uuid.UUID) (<-chan streamSyncHint, func()) { +// Register subscribes deliver to poll hints for chatID until the returned +// unregister func is called. deliver is invoked outside the poller's mutex and +// may race with unregister, so it must remain safe to call after unregister +// returns (e.g. by guarding on the subscriber's own context). +func (p *streamSyncPoller) Register(chatID uuid.UUID, deliver func(streamSyncHint)) (unregister func()) { if p == nil { - ch := make(chan streamSyncHint) - close(ch) - return ch, func() {} + return func() {} } subscriber := &streamSyncPollerSubscriber{ - chatID: chatID, - hints: make(chan streamSyncHint, 1), + chatID: chatID, + deliver: deliver, } p.mu.Lock() if p.subscribers[chatID] == nil { @@ -85,7 +87,7 @@ func (p *streamSyncPoller) Register(chatID uuid.UUID) (<-chan streamSyncHint, fu p.subscribers[chatID][subscriber] = struct{}{} p.mu.Unlock() - return subscriber.hints, func() { + return func() { p.unregister(subscriber) } } @@ -101,7 +103,6 @@ func (p *streamSyncPoller) unregister(subscriber *streamSyncPollerSubscriber) { if len(chatSubscribers) == 0 { delete(p.subscribers, subscriber.chatID) } - close(subscriber.hints) } func (p *streamSyncPoller) loop() { @@ -132,10 +133,7 @@ func (p *streamSyncPoller) pollOnce() { for _, row := range rows { hint := streamSyncHintFromPollRow(row) for _, subscriber := range subscribers[row.ID] { - select { - case subscriber.hints <- hint: - default: - } + subscriber.deliver(hint) } } } diff --git a/coderd/x/chatd/stream_sync_poller_internal_test.go b/coderd/x/chatd/stream_sync_poller_internal_test.go new file mode 100644 index 00000000000..63b2c01c192 --- /dev/null +++ b/coderd/x/chatd/stream_sync_poller_internal_test.go @@ -0,0 +1,74 @@ +package chatd + +import ( + "context" + "sync" + "testing" + + "github.com/google/uuid" + "go.uber.org/mock/gomock" + + "cdr.dev/slog/v3/sloggers/slogtest" + "github.com/coder/coder/v2/coderd/database" + "github.com/coder/coder/v2/coderd/database/dbmock" +) + +// TestStreamSyncPollerConcurrentRegisterUnregister churns subscriber +// registration while pollOnce delivers hints concurrently. Before the poller +// became a callback fanout, unregister closed the subscriber's hint channel +// while pollOnce could still be mid-send on its lock-free snapshot, panicking +// with "send on closed channel" (GHSA-7x3x-59xg-4hrc). Run with -race. +func TestStreamSyncPollerConcurrentRegisterUnregister(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + db := dbmock.NewMockStore(ctrl) + db.EXPECT().GetChatStreamSyncRows(gomock.Any(), gomock.Any()).AnyTimes().DoAndReturn( + func(_ context.Context, ids []uuid.UUID) ([]database.GetChatStreamSyncRowsRow, error) { + rows := make([]database.GetChatStreamSyncRowsRow, 0, len(ids)) + for _, id := range ids { + rows = append(rows, database.GetChatStreamSyncRowsRow{ID: id}) + } + return rows, nil + }, + ) + + poller := newStreamSyncPoller(context.Background(), db, nil, slogtest.Make(t, nil)) + defer poller.Close() + + chatID := uuid.New() + done := make(chan struct{}) + var wg sync.WaitGroup + for range 8 { + wg.Add(1) + go func() { + defer wg.Done() + for { + select { + case <-done: + return + default: + } + unregister := poller.Register(chatID, func(streamSyncHint) {}) + unregister() + } + }() + } + for range 1000 { + poller.pollOnce() + } + close(done) + wg.Wait() +} + +// TestStreamSyncPollerNilRegister verifies a nil poller degrades to a no-op +// registration instead of terminating subscribers. +func TestStreamSyncPollerNilRegister(t *testing.T) { + t.Parallel() + + var poller *streamSyncPoller + unregister := poller.Register(uuid.New(), func(streamSyncHint) { + t.Fatal("nil poller must never deliver hints") + }) + unregister() +}