Skip to content

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

Description

@ImDanXie

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions