Skip to content

bus/natsbus: avoid closing the direct NATS callback channel - #141

Open
herefindalex wants to merge 2 commits into
livekit:mainfrom
herefindalex:fix/nats-subscription-shutdown-race
Open

herefindalex wants to merge 2 commits into
livekit:mainfrom
herefindalex:fix/nats-subscription-shutdown-race

Conversation

@herefindalex

@herefindalex herefindalex commented Sep 11, 2026 •

Copy link
Copy Markdown
Contributor

Summary

NATS Unsubscribe removes a subscription but does not wait for an asynchronous callback that has already been selected. Direct PSRPC subscriptions currently close their handoff channel immediately afterward, so a late callback can race with the close and panic with send on closed channel.

Keep the direct callback channel open and make the raw reader observe context cancellation, matching the existing writer behavior.

Routed subscriptions are unchanged because router.mu already orders child writes, route removal, and channel closure.

Root cause

The following sequence is possible:

  1. NATS selects a pending message for an asynchronous subscription and observes that the subscription is still open.
  2. subscription.Close cancels the context and calls Unsubscribe.
  3. Unsubscribe removes the subscription and returns without joining the selected callback.
  4. Close closes msgChan.
  5. The callback resumes and attempts to send to msgChan, causing panic: send on closed channel.

Cancellation alone does not make closing the channel safe. If the canceled-context receive and a send to a closed channel are both selectable, Go may choose the send and panic.

Fix

  • Do not close the direct NATS callback channel.
  • Let Read() and write() terminate through context cancellation.
  • Preserve the existing routed-subscription shutdown behavior.

The resulting ownership invariant is:

While a NATS callback can still reference a direct subscription, its send destination remains open.

Typed Subscription.Close still waits for its forwarding goroutine to exit and closes the typed output channel before returning.

Tests

Three focused tests cover the shutdown contracts:

  • TestNATSSubscriptionReadCancellation: synctest establishes that Read() is blocked before cancellation, then checks that it returns false.
  • TestNATSSubscriptionWriteCancellation: synctest establishes that an unbuffered send is blocked, then checks that cancellation releases it.
  • TestNATSSubscriptionCloseWithInflightCallback: a real NATS callback pauses before write; Close() returns and removes the subscription; the callback's channel is checked to be open before resuming the callback and waiting for the NATS callback loop to exit.

The channel-open assertion makes the ownership regression deterministic. Relying only on the resumed callback would allow Go to select cancellation instead of the invalid send to a closed channel. A single buffered callback channel exercises this boundary; buffer-state matrices and subscription churn add no distinct shutdown invariant.

Verification

  • Focused shutdown tests: 3 passed.
  • go test -race ./pkg/bus/natsbus ./pkg/bus -count=1 -timeout=120s: 58 passed.
  • Running the same focused tests against the unpatched main implementation deterministically fails the Read cancellation and channel ownership assertions; the existing write cancellation behavior passes.

Behavior notes

This is cancellation-based shutdown, not a drain or callback join. A callback selected before Unsubscribe may still enqueue a bounded buffered message if the send is ready. Owners must also continue calling Close; canceling the parent context alone does not unregister the NATS subscription.

@herefindalex herefindalex changed the title bus: avoid closing the direct NATS callback channel ‎internal/bus: avoid closing the direct NATS callback channel Sep 11, 2026
Unsubscribe removes interest but does not join an asynchronous callback
that NATS has already selected. Closing the direct subscription's msgChan
can therefore race with a callback send and panic.

Keep the callback channel open and let the raw reader observe cancellation,
as the writer already does. Typed Close can still wait for its forwarder
to exit. Routed children retain their existing mutex-protected shutdown.

Add a deterministic regression that pauses a real NATS callback across
Close and checks that its destination remains open before resuming it.
Wait for NATS callback completion, cover canceled readers and blocked
writers, and exercise typed shutdown and direct/queue subscription churn.
@herefindalex
herefindalex force-pushed the fix/nats-subscription-shutdown-race branch from 55600a6 to 0d0828e Compare September 21, 2026 03:26
@herefindalex herefindalex changed the title ‎internal/bus: avoid closing the direct NATS callback channel bus/natsbus: avoid closing the direct NATS callback channel Sep 21, 2026
@herefindalex

Copy link
Copy Markdown
Contributor Author

@paulwe I rebased this onto v0.8.0 after #148 and adapted the fix and regression tests to pkg/bus/natsbus.

Verification on the rebased commit:

  • go test ./pkg/bus/natsbus -count=1
  • go test -race ./pkg/bus/natsbus -count=1
  • go test -race ./pkg/bus/natsbus -run 'TestNATSSubscription' -count=10
  • go test ./pkg/bus -count=1
  • go tool mage testall
  • go test -race ./... -count=1

Could you take a look when you have a chance?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant