Skip to content

fix(rpc): abort in-flight requests when a request pipeline closes — fixes #305 - #306

Open
ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/rpc-close-abort-inflight
Open

ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/rpc-close-abort-inflight

Conversation

@ImDanXie

Copy link
Copy Markdown
Contributor

Summary

ManagedRequestPipeline.close() only cancels preflightTaskQueue (requests not yet sent). It never touches inflightTaskQueue (sent, awaiting a response).

Root cause

After super.close() disposes the stream, responses can never arrive, and onStreamError() — the only path that calls cancelInflightTasks() — will not fire again. The futures of already-sent requests therefore never complete:

  • callers on the direct (non-batched) query/execute path hang forever;
  • the batched path is only saved by Batcher's orTimeout(maxBurstLatency) fallback;
  • the inflight tasks themselves leak.

Pipelines are replaced via old.close() on store→server route changes, so a single range migration with a query in flight triggers this.

Fix

Abort in-flight tasks on close, mirroring onStreamError:

if (isClosed.compareAndSet(false, true)) {
    super.close();
    meter.recordCount(RPCMetric.ReqPipelineCompleteCount);
    cancelPreflightTasks(new RequestAbortException("Pipeline has closed"));
    cancelInflightTasks(new RequestAbortException("Pipeline has closed"));   // added
}

Verification

  • Static analysis of the close path (the audit trail, including the reasoning for why onStreamError cannot fire after super.close(), is in [BUG] ManagedRequestPipeline.close() leaves in-flight request futures pending forever #305).
  • Deployed in a rolling 6-node cluster (wholesale lib/ replacement, one node at a time, zero rollbacks); cluster mesh, all 13 live client connections and cross-node delivery verified post-deploy. No non-batched caller hangs observed.
  • A unit-level regression test requires a full RPC stack (client + server + stream disposal); a reproduction using the in-process transport can be added if the maintainers would like one in-tree.

Fixes #305

close() only cancelled preflight (unsent) tasks. Once the stream is disposed
by super.close(), responses can never arrive and onStreamError() - the only
path that calls cancelInflightTasks() - will not fire again, so the futures
of already-sent requests never complete. Callers on the direct (non-batched)
query/execute path hang forever; the batched path is only saved by Batcher's
orTimeout fallback. Pipelines are replaced via old.close() on store->server
route changes, so one range migration with a query in flight triggers it.

Abort inflight tasks with RequestAbortException, mirroring onStreamError.

Fixes apache#305
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.

[BUG] ManagedRequestPipeline.close() leaves in-flight request futures pending forever

1 participant