From a2e4c059b0e5739a639680053dd5813237a4954a Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Sat, 26 Sep 2026 02:01:08 +0800 Subject: [PATCH] fix(kv-server): complete the restore session future on synchronous start failure restoreFrom() swallows any Throwable from startRestore()/messenger.send() with just a log line, leaving the session's doneFuture pending forever. There is no timeout on this path, and because the session remains cached and un-done, a retry with the same snapshot (same leader) hits the reuse branch and returns the same dead future - the range is permanently stuck after a failed snapshot install. Complete the future exceptionally in the catch block, so callers observe the failure and a retry starts a fresh session. Control experiment: restoreFromStartFailureCompletesFutureExceptionally fails on the old code and passes with the fix; existing restore expectations unchanged (9/9). Fixes #303 --- .../basekv/store/range/KVRangeRestorer.java | 5 +++++ .../store/range/KVRangeRestorerTest.java | 20 +++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/base-kv/base-kv-store-server/src/main/java/org/apache/bifromq/basekv/store/range/KVRangeRestorer.java b/base-kv/base-kv-store-server/src/main/java/org/apache/bifromq/basekv/store/range/KVRangeRestorer.java index a122e2d26..a3ee4710e 100644 --- a/base-kv/base-kv-store-server/src/main/java/org/apache/bifromq/basekv/store/range/KVRangeRestorer.java +++ b/base-kv/base-kv-store-server/src/main/java/org/apache/bifromq/basekv/store/range/KVRangeRestorer.java @@ -138,6 +138,11 @@ public CompletableFuture restoreFrom(String leader, KVRangeSnapshot rangeS } } catch (Throwable t) { log.error("Unexpected error", t); + // A synchronous failure here (startRestore/send throwing) must complete the session + // future exceptionally. Leaving it pending hangs the caller forever (there is no timeout on + // this path) and, because the session stays cached and un-done, a retry with the same snapshot + // would reuse the same dead session - permanently stuck range. + onDone.completeExceptionally(new KVRangeStoreException("Snapshot restore failed to start", t)); } return onDone; } diff --git a/base-kv/base-kv-store-server/src/test/java/org/apache/bifromq/basekv/store/range/KVRangeRestorerTest.java b/base-kv/base-kv-store-server/src/test/java/org/apache/bifromq/basekv/store/range/KVRangeRestorerTest.java index f6dbf4815..a2db9098f 100644 --- a/base-kv/base-kv-store-server/src/test/java/org/apache/bifromq/basekv/store/range/KVRangeRestorerTest.java +++ b/base-kv/base-kv-store-server/src/test/java/org/apache/bifromq/basekv/store/range/KVRangeRestorerTest.java @@ -250,6 +250,26 @@ public void startNewSessionWhenLeaderChanges() { secondRestore.cancel(true); // ensure no lingering scheduling } + @Test + public void restoreFromStartFailureCompletesFutureExceptionally() { + // A synchronous failure from startRestore must not leave the future (and the cached + // session) pending forever - range would be stuck, and a same-snapshot retry would reuse the dead + // session. + when(range.startRestore(eq(snapshot), any())).thenThrow(new IllegalStateException("boom")); + + KVRangeRestorer restorer = new KVRangeRestorer(snapshot, range, messenger, metricManager, executor, 10); + CompletableFuture restoreFuture = restorer.restoreFrom("leader", snapshot); + + assertTrue(restoreFuture.isCompletedExceptionally()); + + // a retry with the same snapshot must start a fresh session instead of reusing the dead one + when(range.startRestore(eq(snapshot), any())).thenReturn(mock(IKVRangeRestoreSession.class)); + CompletableFuture retry = restorer.restoreFrom("leader", snapshot); + assertFalse(retry.isDone()); + verify(range, times(2)).startRestore(eq(snapshot), any()); + retry.cancel(true); + } + @Test public void reuseAfterDoneStartsNew() { IKVRangeRestoreSession firstRS = mock(IKVRangeRestoreSession.class);