Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,11 @@ public CompletableFuture<Void> 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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Void> 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<Void> 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);
Expand Down