Environment
- BiFroMQ
4.0.0 / main @ eef5e3af
- Module:
base-rpc/base-rpc-client
- File:
base-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java (close())
Summary
close() cancels only preflightTaskQueue (requests not yet sent). It never touches inflightTaskQueue (sent, awaiting a response). After super.close() disposes the stream, 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(maxBurstLatency) fallback.
Pipelines are replaced via old.close() on store→server route changes, so one range migration with a query in flight triggers this.
Suggested 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"));
}
Verification note
A regression test requires a full RPC stack (client + server + stream disposal); we verified the path statically and in a rolling multi-node deployment (see below). A reproduction using the in-process transport can be provided if useful.
Production validation
A 6-node BifroMQ 4.0.0 cluster (standalone.sh, systemd, mTLS clientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5, $oshare subscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.
The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale lib/ replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.
Post-deploy checks (all measured after the rollout):
| Check |
Result |
| Cluster mesh (inter-node 8898/8899 connections, per node) |
408–417 established on every node |
| Live client connections (3 bridge replicas + 10 backend services over mTLS) |
13/13 reconnected and stable |
| Cross-node delivery end-to-end |
publish on one node → bridge (connected elsewhere) consumed it → forwarded (counter 0 → 1) |
| Client-visible protocol smoke |
connect / subscribe / publish with QoS 1 all CONNACK 0 / PUBACK 0 |
No non-batched caller hangs were observed in the rollout window; the fix removes the latent path.
Environment
4.0.0/main @ eef5e3afbase-rpc/base-rpc-clientbase-rpc/base-rpc-client/src/main/java/org/apache/bifromq/baserpc/client/ManagedRequestPipeline.java(close())Summary
close()cancels onlypreflightTaskQueue(requests not yet sent). It never touchesinflightTaskQueue(sent, awaiting a response). Aftersuper.close()disposes the stream, responses can never arrive andonStreamError()— the only path that callscancelInflightTasks()— will not fire again, so the futures of already-sent requests never complete. Callers on the direct (non-batched)query/executepath hang forever; the batched path is only saved byBatcher'sorTimeout(maxBurstLatency)fallback.Pipelines are replaced via
old.close()on store→server route changes, so one range migration with a query in flight triggers this.Suggested fix
Abort in-flight tasks on close, mirroring
onStreamError:Verification note
A regression test requires a full RPC stack (client + server + stream disposal); we verified the path statically and in a rolling multi-node deployment (see below). A reproduction using the in-process transport can be provided if useful.
Production validation
A 6-node BifroMQ 4.0.0 cluster (
standalone.sh, systemd, mTLSclientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5,$osharesubscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale
lib/replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.Post-deploy checks (all measured after the rollout):
No non-batched caller hangs were observed in the rollout window; the fix removes the latent path.