From e827e932d91b0c2c0dbc4fe00a2cf70ded955cb3 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Sat, 26 Sep 2026 01:54:38 +0800 Subject: [PATCH] fix(rpc): abort in-flight requests when a request pipeline closes 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 #305 --- .../apache/bifromq/baserpc/client/ManagedRequestPipeline.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/base-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java b/base-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java index a351ae5b1..284f868c7 100644 --- a/base-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java +++ b/base-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java @@ -142,6 +142,10 @@ public void close() { super.close(); meter.recordCount(RPCMetric.ReqPipelineCompleteCount); cancelPreflightTasks(new RequestAbortException("Pipeline has closed")); + // Also abort in-flight requests: super.close() disposes the stream, so their responses + // can never arrive and onStreamError() will not fire afterwards - leaving the futures pending + // hangs direct (non-batched) callers forever and leaks the inflight tasks. + cancelInflightTasks(new RequestAbortException("Pipeline has closed")); } }