Skip to content

fix: send the client writes to the current shard leader, and never after their deadline - #1399

Draft
merlimat wants to merge 2 commits into
oxia-db:mainfrom
merlimat:fix-client-write-leader-routing
Draft

merlimat wants to merge 2 commits into
oxia-db:mainfrom
merlimat:fix-client-write-leader-routing

Conversation

@merlimat

Copy link
Copy Markdown
Collaborator

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:

  1. The writes stick to the old leader. The client keeps one write stream per shard, and reuses it as long as it has not failed and no leader hint arrived, whatever leader the shard assignments name. A leader that was replaced but is still reachable (asymmetric partition, saturated node) keeps the stream open without answering: every write of the client to the shard goes to it and times out, while the new leader serves the other clients.
  2. The stream setup ignores the request timeout. The stream is set up with the long-lived context of the client, so connecting to an unreachable leader can take up to the gRPC connect timeout (20 s). Once connected, the request is sent anyway, after its deadline.
  3. The SyncClient reports a failure for a write that it sends later. With one write batch in flight per shard, a 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 ends, Put, Delete and DeleteRange return 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 (or context.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:

  1. 0.00 s: w1 is sent to A, which does not answer.
  2. 0.10 s: Put(ctx, "w2"), with a 1 s context, waits behind w1.
  3. 1.10 s: it returns context deadline exceeded. Then Put(ctx, "w3"), with a 5 s context.
  4. 2.00 s: w1 fails with context deadline exceeded. Then w2, which already failed, is sent to A, on the same stream.
  5. 4.00 s: w3 is sent to A too, and it fails at 6.00 s. B never receives a write.

After:

  1. 0.00 s: w1 is sent to A, which does not answer.
  2. 0.10 s: Put(ctx, "w2"), with a 1 s context, waits behind w1.
  3. 1.10 s: it returns context deadline exceeded, and w2 is dropped: it is never sent. Then Put(ctx, "w3"), with a 5 s context.
  4. 2.00 s: w1 fails as before. A may still have applied it, so it is not sent again to B.
  5. 2.00 s: w3 goes to B, on a new stream, and succeeds.

If A is only slow instead (it is still the leader, and it answers w1 after 2 s), and the application issues w2 again when it fails at 1.10 s, A applies w2 twice 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

  1. The writes follow the shard assignments. A write stream knows the node it goes to. An attempt reuses the stream of the shard only if it goes to the leader of that attempt: the one the shard assignments name, or the one hinted by the previous attempt. Otherwise, a new stream replaces it. The write in flight to the old leader is not abandoned to be sent again: it waits for its response until its deadline, because the old leader may still apply it. Sending it again would make the conditional and sequence puts at-least-once, which is for the sequence dedup work (Server-side deduplication of retried sequence requests (exactly-once assignment) #1223) to address.
  2. Nothing is sent after its deadline. The stream setup is cancelled when the context of the request ends (the stream itself still outlives the request), and a request is not handed to gRPC once its context is done.
  3. The SyncClient drops the writes it gave up on. Put, Delete and DeleteRange tie 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.ExecuteWrite takes a function that returns the request, invoked right before sending.
  4. Plumbing for the definite errors. The failure of a write that was never handed to gRPC is marked internally (model.NotSent). The public API does not change: the error message is the same, and errors.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_WriteStreamSetupRespectsTheDeadline and TestStreamWrapper_DoesNotSendAfterTheDeadline: the connection that takes 3 s.
  • TestSyncClient_QueuedWriteIsNotSentAfterItsContextEnds and TestSyncClient_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.
  • Unit tests of CallContext and of the batch dropping the writes.

Notes

… 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>
Copilot AI balanced review requested due to automatic review settings September 30, 2026 00:12

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

This branch has not been deployed

No deployments
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.

2 participants