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")); } }