fix: keep notification subscriptions live on idle shards and cancel them on close - #385
Merged
Merged
Conversation
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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Motivation
Two problems with the notification subscriptions, fixed by one commit each.
Resumed subscriptions time out on idle shards. Since #370, the
GetNotificationsestablishment barrier fails withTIMEOUTwhen no message arrives withinrequestTimeout. 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 aclosedflag: the live call was not cancelled, andonNextdid not check the flag.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.close()shuts down the channel, withshutdownNowwhile streams are still open. That fails an idle stream (one whose current attempt has received no message yet) withUNAVAILABLE, which the retry policy treats as retryable, so closing the client loggedWARN 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.Modifications
GrpcRpcProvider.getNotificationsappliesrequestTimeoutonly 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 gRPCContext, after itsTIMEOUTis reported, so that it is not left open on the server.RpcProvider.getNotificationstakes aCancelableStreamObserver, likegetSequenceUpdates,listandrangeScan, andGrpcRpcProviderinjects the call of each attempt into it withtoBarrierClientResponseObserver.ShardNotificationReceivercreates one observer per subscription, asSequenceUpdatesdoes, andclose()cancels it. The call then ends withCANCELLED, 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 afterclose()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 realGrpcRpcProvider: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
subscriptionMaxAgeof 2s, so that every subscription had resumed and was idle:Retrying get notifications ... Channel shutdownNow invokedwarnings onclose()Closing a client from inside its own notification callback also completed, in 3 runs out of 3.