Repository navigation
Conversation
… deadline The Go client reused the write stream of a shard as long as it had not failed and no leader hint arrived, whatever leader the shard assignments named. A leader that was replaced, but is still reachable (asymmetric partition, saturated node), keeps the stream open without applying the writes: every write of the client to the shard went to it and timed out, indefinitely, while the new leader was serving. The stream was also set up with the long-lived context of the provider: connecting to an unreachable leader could take up to the gRPC connect timeout (20 s), past the request timeout, and the request was then sent after its deadline. Modifications: - A write stream knows its target, and a write goes to a new stream when the shard assignments name another leader. A leader hint only comes with the failure of the stream, so it needs no check of its own anymore. The write in flight to the old leader is not sent again: it waits for its response until its deadline, as the old leader may still apply it. - The stream setup is cancelled when the request context ends, and a request is not handed to gRPC once its context is done. - The context of a stream is released once the stream ends, and a stream is marked as failed before its requests fail, so that their retries use a new stream. Signed-off-by: Matteo Merli <mmerli@apache.org>
With a single write batch in flight per shard, a SyncClient write can wait behind a batch that the leader does not answer, or retry while the shard has no leader. When the context of the caller ended, the SyncClient returned the context error, but the write stayed queued, and it was sent and applied later: an application issuing the write again after the error created a duplicate. Now the SyncClient gives up on a write when its context ends: unless the write was handed to gRPC already, it is dropped, and it fails with an error telling that it was not sent. A write that was sent may still be applied, and it fails with the plain context error, as before. Modifications: - The SyncClient ties its put, delete and delete range calls to the context of the caller (model.CallContext): a call is either sent or dropped, atomically. - A write batch sends or drops its calls when the request is first handed to gRPC, not when the batch starts, so that a call whose batch retries without sending, e.g. during a leader election, is dropped too. For that, Executor.ExecuteWrite gets the request from a function it invokes right before sending. - Plumbing to tell the definite failures from the ambiguous ones: the failure of a call that was never handed to gRPC is marked (model.NotSent), e.g. a batch that times out while the shard has no leader. It is not exposed in the public API yet, and errors.Is still finds the context error. - A request that was sent is still never abandoned to be sent again: for conditional and sequence puts that would be at-least-once, which is for the sequence dedup work (oxia-db#1223) to address. Signed-off-by: Matteo Merli <mmerli@apache.org>
merlimat
requested review from
RobertIndie,
coderzc and
mattisonchao
as code owners
September 30, 2026 00:12
This branch has not been deployed
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.
Problem
When the leader of a shard changes or stops answering, the Go client can keep sending the writes to the wrong node, send them after their deadline, and report a failure for a write that it applies later:
Put,DeleteandDeleteRangereturn the context error, but the write stays queued, and it is sent and applied later. An application that issues the write again applies it twice, e.g. a sequence put creates two keys.In all these cases the error is
context.DeadlineExceeded(orcontext.Canceled), whether the write was never sent or may still be applied, so the application cannot tell if issuing it again is safe.Example
A client with a 2 s request timeout writes to shard 0, led by node A. Right after A receives the put
w1, B is elected: A stays reachable, but it does not answer anymore.Before:
w1is sent to A, which does not answer.Put(ctx, "w2"), with a 1 s context, waits behindw1.context deadline exceeded. ThenPut(ctx, "w3"), with a 5 s context.w1fails withcontext deadline exceeded. Thenw2, which already failed, is sent to A, on the same stream.w3is sent to A too, and it fails at 6.00 s. B never receives a write.After:
w1is sent to A, which does not answer.Put(ctx, "w2"), with a 1 s context, waits behindw1.context deadline exceeded, andw2is dropped: it is never sent. ThenPut(ctx, "w3"), with a 5 s context.w1fails as before. A may still have applied it, so it is not sent again to B.w3goes to B, on a new stream, and succeeds.If A is only slow instead (it is still the leader, and it answers
w1after 2 s), and the application issuesw2again when it fails at 1.10 s, A appliesw2twice at 2.00 s before the fix: the put that failed, and the one issued again. After the fix, it applies it once.With a leader whose address accepts the connection only after 3 s, a write with a 200 ms deadline failed after 3.0 s, and the leader then received it. It now fails when its deadline expires, and it is never sent.
Modification
Put,DeleteandDeleteRangetie each write to the context of the caller (model.CallContext). When the request of a write batch is first handed to gRPC, each write is atomically either marked as sent or, if its caller gave up on it, dropped: a write is either sent, or it fails and is never sent. This happens right before sending, not when the batch starts, so that a batch retrying while the shard has no leader can still drop its writes. For that,Executor.ExecuteWritetakes a function that returns the request, invoked right before sending.model.NotSent). The public API does not change: the error message is the same, anderrors.Is(err, context.DeadlineExceeded)still holds. Exposing the difference between "not applied" and "outcome unknown" is left to the sequence dedup design.The tests reproduce the example, and each fails on main:
TestRpcProvider_WritesGoToTheNewLeader: the leader changes while a write is in flight. On main, the old leader receives both writes and the new leader none.TestRpcProvider_WriteStreamSetupRespectsTheDeadlineandTestStreamWrapper_DoesNotSendAfterTheDeadline: the connection that takes 3 s.TestSyncClient_QueuedWriteIsNotSentAfterItsContextEndsandTestSyncClient_WriteIsNotSentAfterItsContextEndsWithoutLeader: a put, delete or delete range that failed must never be sent, when it waits behind a write that the leader does not answer, or while the shard has no leader.CallContextand of the batch dropping the writes.Notes
w3waits forw1to time out. Failing it early when the leader changes, with an ambiguous error and without sending it again, would free the queue.