From 172d0a6daff853d21cb09a7e41771b86252e0e18 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Sat, 26 Sep 2026 02:00:18 +0800 Subject: [PATCH] fix(kv-client): close mutation pipelines dropped by a route refresh refreshMutPipelines() rebuilds the per-store/per-range mutation pipeline map and assigns it, but only closed pipelines of stores that vanished entirely. A pipeline whose range is no longer led by that store (the regular outcome of a RangeLeaderBalancer leadership transfer or a merge) was simply dropped: its RX subscription, its dedicated gRPC bidi stream and the server-side MutatePipeline leaked on every topology event - unreachable from close() - giving monotonically growing resource usage on long-running brokers. Replace the store-level-only clear loop with a general pass that closes every pipeline instance absent from (or replaced in) the new map, extracted as a package-private static helper for direct unit testing. Tests: range-leadership-move / store-disappear / instance-replaced cases (new BaseKVStoreClientPipelineRefreshTest). Fixes #302 --- .../basekv/client/BaseKVStoreClient.java | 23 +++++- .../BaseKVStoreClientPipelineRefreshTest.java | 76 +++++++++++++++++++ 2 files changed, 96 insertions(+), 3 deletions(-) create mode 100644 base-kv/base-kv-store-client/src/test/java/org/apache/bifromq/basekv/client/BaseKVStoreClientPipelineRefreshTest.java diff --git a/base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java b/base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java index 14a954bd7..ed68a192a 100644 --- a/base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java +++ b/base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java @@ -564,9 +564,26 @@ private void refreshMutPipelines(Map storeDescri } } mutPplns = nextMutPplns; - // clear mut pipelines targeting non-exist storeId; - for (String storeId : Sets.difference(currentMutPplns.keySet(), nextMutPplns.keySet())) { - currentMutPplns.get(storeId).values().forEach(IMutationPipeline::close); + closeDroppedMutationPipelines(currentMutPplns, nextMutPplns); + } + + /** + * Close every mutation pipeline that this refresh drops - either because its range is no longer + * led by that store (leadership transfer/merge) or because the whole store disappeared. Assigning + * mutPplns without closing the dropped pipelines leaked each of them (RX subscription + a dedicated + * gRPC bidi stream + a server-side MutatePipeline) on every leadership transfer, and they are + * unreachable from {@link #close()} once dropped. + */ + static void closeDroppedMutationPipelines(Map> current, + Map> next) { + Map> safeNext = next == null ? emptyMap() : next; + for (Map.Entry> byStore : current.entrySet()) { + Map nextRanges = safeNext.getOrDefault(byStore.getKey(), emptyMap()); + for (Map.Entry byRange : byStore.getValue().entrySet()) { + if (nextRanges.get(byRange.getKey()) != byRange.getValue()) { + byRange.getValue().close(); + } + } } } diff --git a/base-kv/base-kv-store-client/src/test/java/org/apache/bifromq/basekv/client/BaseKVStoreClientPipelineRefreshTest.java b/base-kv/base-kv-store-client/src/test/java/org/apache/bifromq/basekv/client/BaseKVStoreClientPipelineRefreshTest.java new file mode 100644 index 000000000..8b5f59915 --- /dev/null +++ b/base-kv/base-kv-store-client/src/test/java/org/apache/bifromq/basekv/client/BaseKVStoreClientPipelineRefreshTest.java @@ -0,0 +1,76 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.bifromq.basekv.client; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import java.util.Map; +import org.apache.bifromq.basekv.proto.KVRangeId; +import org.testng.annotations.Test; + +/** + * Regression: mutation pipelines dropped by a route refresh must be closed, not silently leaked. + */ +public class BaseKVStoreClientPipelineRefreshTest { + + private static KVRangeId rangeId(long id) { + return KVRangeId.newBuilder().setEpoch(id).setId(id).build(); + } + + @Test + public void closesPipelineWhenRangeLeadershipMovesAway() { + IMutationPipeline kept = mock(IMutationPipeline.class); + IMutationPipeline dropped = mock(IMutationPipeline.class); + Map> current = Map.of( + "storeA", Map.of(rangeId(1), kept, rangeId(2), dropped)); + Map> next = Map.of( + "storeA", Map.of(rangeId(1), kept)); + + BaseKVStoreClient.closeDroppedMutationPipelines(current, next); + + verify(dropped).close(); + verify(kept, never()).close(); + } + + @Test + public void closesPipelinesWhenStoreDisappears() { + IMutationPipeline ppln = mock(IMutationPipeline.class); + Map> current = Map.of("storeB", Map.of(rangeId(3), ppln)); + + BaseKVStoreClient.closeDroppedMutationPipelines(current, Map.of()); + + verify(ppln).close(); + } + + @Test + public void closesReplacedPipelineInstance() { + IMutationPipeline oldPpln = mock(IMutationPipeline.class); + IMutationPipeline newPpln = mock(IMutationPipeline.class); + Map> current = Map.of("storeA", Map.of(rangeId(1), oldPpln)); + Map> next = Map.of("storeA", Map.of(rangeId(1), newPpln)); + + BaseKVStoreClient.closeDroppedMutationPipelines(current, next); + + verify(oldPpln).close(); + verify(newPpln, never()).close(); + } +}