Skip to content

fix: keep notification subscriptions live on idle shards and cancel them on close - #385

Merged
merlimat merged 4 commits into
oxia-db:mainfrom
merlimat:fix/notification-subscriptions
Sep 25, 2026
Merged

merlimat merged 4 commits into
oxia-db:mainfrom
merlimat:fix/notification-subscriptions

Conversation

@merlimat

Copy link
Copy Markdown
Collaborator

Motivation

Two problems with the notification subscriptions, fixed by one commit each.

Resumed subscriptions time out on idle shards. Since #370, the GetNotifications establishment barrier fails with TIMEOUT when no message arrives within requestTimeout. The server confirms only a new subscription, with a first "dummy" batch: one that resumes from an offset gets nothing until a new notification is written. Since #372 renews every subscription at its max age, from the first renewal on, each resumed subscription on an idle shard timed out and was re-established after the receiver's backoff, which only resets on a message and grew to 60s. The timed-out calls were not cancelled and stayed open on the server until their max age, and while the receiver waited out its backoff no subscription was live, so a notification could be delayed by up to about a minute.

Closing a receiver leaves its call open. ShardNotificationReceiver.close() only set a closed flag: the live call was not cancelled, and onNext did not check the flag.

  • A client built on shared resources does not own the connections, so after client.close() its notification streams stayed open on the server until the subscription max age (5-10 min by default), and every write on those shards still invoked the closed client's callback.
  • A standalone client's close() shuts down the channel, with shutdownNow while streams are still open. That fails an idle stream (one whose current attempt has received no message yet) with UNAVAILABLE, which the retry policy treats as retryable, so closing the client logged WARN Retrying get notifications after PT0.1S. attempt=2 exception=Channel shutdownNow invoked. With resumed subscriptions no longer timed out, it did so for every idle shard.
  • When a shard is removed or reassigned, the old receiver's stream was left open.

Modifications

  • GrpcRpcProvider.getNotifications applies requestTimeout only to a new subscription, which the server answers right away. Like the sequence updates, a resumed subscription is bounded by the subscription max age and the connection health checks. An attempt abandoned by the timeout is cancelled through a per-attempt gRPC Context, after its TIMEOUT is reported, so that it is not left open on the server.
  • RpcProvider.getNotifications takes a CancelableStreamObserver, like getSequenceUpdates, list and rangeScan, and GrpcRpcProvider injects the call of each attempt into it with toBarrierClientResponseObserver.
  • ShardNotificationReceiver creates one observer per subscription, as SequenceUpdates does, and close() cancels it. The call then ends with CANCELLED, which is not retried, before any channel shutdown, and later batches are ignored. Cancelling waits for a batch that is being delivered, so the offset that a reassignment reads after close() is final.
  • NotificationManager.close() closes the receivers sequentially instead of with a parallel stream. A client closed from a notification callback has to close that callback's receiver on the callback's thread: a pool thread would wait for the batch being delivered, while the callback waits for the pool.

Verification

  • New ShardNotificationReceiverGrpcTest, with a fake server and a real GrpcRpcProvider:

    • resumedSubscriptionOnIdleShardStaysLive: after its first renewal, the resumed subscription on an idle shard is not timed out and re-established, no abandoned call is left open on the server, and a notification written later is delivered right away.
    • closeCancelsLiveSubscription, for a new subscription and for a resumed one on an idle shard: closing the receiver while the provider stays open, as in shared mode, cancels the call on the server, a later write reaches no callback, and no new attempt is made. Both cases fail with only the first commit, where the call is never cancelled.
  • GrpcRpcProviderTest: a silent new subscription still times out, and its call is now cancelled on the server; a silent resumed one no longer times out.

  • All unit tests pass (./gradlew :client:test --tests '*Test'). The Docker-based ITs don't run locally.

  • End to end against an Oxia v0.17.1 standalone server with 3 shards and a subscriptionMaxAge of 2s, so that every subscription had resumed and was idle:

    First commit only Both commits
    Shared resources: client A's notification streams on the server, 1s after closing A 2 of 3 (the third had reached its max age) 0
    Shared resources: notifications delivered to the closed client A after client B wrote 30 keys 19 0
    Standalone: Retrying get notifications ... Channel shutdownNow invoked warnings on close() 3 0

    Closing a client from inside its own notification callback also completed, in 3 runs out of 3.

Since oxia-db#370, the GetNotifications establishment barrier fails with TIMEOUT
when no message arrives within requestTimeout. The server confirms only a
new subscription, with a first "dummy" batch: one that resumes from an
offset gets nothing until a new notification is written. Since oxia-db#372
renews every subscription at its max age, from the first renewal on each
resumed subscription on an idle shard timed out and was re-established
after the receiver's backoff, which only resets on a message and grew to
60s. The timed-out calls were not cancelled and stayed open on the server
until their max age, and while the receiver waited out its backoff no
subscription was live, so a notification could be delayed by up to about
a minute.

Apply requestTimeout only to a new subscription, which the server answers
right away. Like the sequence updates, a resumed subscription is bounded
by the subscription max age and the connection health checks. Also cancel
an attempt abandoned by the timeout, so that it is not left open on the
server.

Signed-off-by: Matteo Merli <mmerli@apache.org>
ShardNotificationReceiver.close() only set a flag: it did not cancel the
live GetNotifications call, and a batch received afterwards was still
delivered. A client built on shared resources does not own the
connections, so after close() its notification streams stayed open on the
server until their max age, and every write on those shards still invoked
its callback. A standalone client's close() shuts down the channel with
shutdownNow, which fails an idle stream, one whose current attempt has
received no message yet, with UNAVAILABLE. That is retryable, so closing
the client logged "Retrying get notifications ... Channel shutdownNow
invoked", and now that resumed subscriptions are not timed out, it did so
for every idle shard. The stream of a receiver closed on a shard removal
or reassignment was also left open.

Like getSequenceUpdates, list and rangeScan, RpcProvider.getNotifications
now takes a CancelableStreamObserver, into which GrpcRpcProvider injects
the call of each attempt. The receiver creates one observer per
subscription, as SequenceUpdates does, and close() cancels it: the call
ends with CANCELLED, which is not retried, before any channel shutdown,
and later batches are ignored. Cancelling waits for a batch being
delivered, so the offset that a reassignment reads after close() is final.

Because closing a receiver can wait for its callback, NotificationManager
closes the receivers sequentially: a client closed from a notification
callback must close that callback's receiver on the same thread, and a
parallel stream could hand it to a pool thread and deadlock.

Signed-off-by: Matteo Merli <mmerli@apache.org>
Adapt the getNotifications call that oxia-db#383 added to GrpcRpcProviderTest to
the CancelableStreamObserver parameter.

Signed-off-by: Matteo Merli <mmerli@apache.org>
main now has the idle-shard fix as oxia-db#386, identical to the first commit of
this branch. Keep this branch's side of the conflicting hunks, where the
close-cancel change edits the same lines.

Signed-off-by: Matteo Merli <mmerli@apache.org>
@merlimat
merlimat merged commit b6c3bf3 into oxia-db:main Sep 25, 2026
2 checks passed
@merlimat
merlimat deleted the fix/notification-subscriptions branch September 25, 2026 16:07
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