Skip to content

[BUG] BaseKVStoreClient leaks a mutation pipeline and its gRPC stream on every range-leader change #302

Description

@ImDanXie

Environment

  • BiFroMQ 4.0.0 / main @ eef5e3af
  • Module: base-kv/base-kv-store-client
  • File: base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java (refreshMutPipelines)

Summary

refreshMutPipelines() rebuilds the per-store/per-range mutation pipeline map on every route refresh, reusing pipelines whose range is still led by that store and dropping the rest. The only cleanup loop closes pipelines of stores that disappeared entirely; a pipeline whose range is no longer led by that store (the routine outcome of a RangeLeaderBalancer leadership transfer or a merge) is silently dropped from the map without close(). The dropped ManagedMutationPipeline keeps its RX subscription and its dedicated gRPC bidi stream (plus a server-side MutatePipeline), and it is unreachable from BaseKVStoreClient.close() once dropped.

Impact

RangeLeaderBalancer performs leadership transfers routinely (it is registered by default in dist-worker / inbox-store deployments). Each transfer leaks one pipeline + one gRPC bidi stream + one server-side pipeline: monotonically growing resource usage on long-running brokers.

Suggested fix

Close every pipeline instance absent from (or replaced in) the new map — one pass covers all three forms (range moved away, store disappeared, instance replaced):

mutPplns = nextMutPplns;
closeDroppedMutationPipelines(currentMutPplns, nextMutPplns);
static void closeDroppedMutationPipelines(Map<String, Map<KVRangeId, IMutationPipeline>> current,
                                          Map<String, Map<KVRangeId, IMutationPipeline>> next) {
    Map<String, Map<KVRangeId, IMutationPipeline>> safeNext = next == null ? emptyMap() : next;
    for (Map.Entry<String, Map<KVRangeId, IMutationPipeline>> byStore : current.entrySet()) {
        Map<KVRangeId, IMutationPipeline> nextRanges = safeNext.getOrDefault(byStore.getKey(), emptyMap());
        for (Map.Entry<KVRangeId, IMutationPipeline> byRange : byStore.getValue().entrySet()) {
            if (nextRanges.get(byRange.getKey()) != byRange.getValue()) {
                byRange.getValue().close();
            }
        }
    }
}

Test

New test class with three cases: range-leadership-move / store-disappear / instance-replaced (mocks verify the dropped pipelines are closed and the kept ones are not).

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

Baselines captured per node for leak tracking (JVM threads 216–229, RSS 2.8–5.1 GB, gRPC inter-node connections 408–417); the leak manifests as monotonic growth across further leadership transfers, which is being watched.

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