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);