From 5896cd90739666fd7683dd41f83b8b118590fb65 Mon Sep 17 00:00:00 2001 From: Alan Wang Date: Fri, 17 Jul 2026 16:04:00 -0700 Subject: [PATCH] fix --- .../service/accord/CoordinatedTransfer.java | 26 +++++++++++-------- 1 file changed, 15 insertions(+), 11 deletions(-) diff --git a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java index b52369330f..2702e70b4f 100644 --- a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java @@ -177,8 +177,8 @@ public class CoordinatedTransfer private void performImportTxn() { + Invariants.require(!streamResults.containsValue(SingleTransferResult.Init())); TimeUUID[] planIds = streamResults.values().stream() - .filter(result -> result.planId != null) .map(result -> result.planId) .toArray(TimeUUID[]::new); @@ -284,9 +284,9 @@ public class CoordinatedTransfer if (state.hasFailedSession() || state.hasAbortedSession()) throw new StreamException(state, "Stream failed due to failed or aborted sessions"); - // If the SSTable doesn't contain any rows in the provided range, no streams delivered, nothing to activate - if (state.sessions().isEmpty()) - return SingleTransferResult.Noop(); + // Since we have filtered the SSTables using getPositionsForRanges prior, all sessions created + // should contain at least 1 SSTable + Invariants.require(!state.sessions().isEmpty()); return SingleTransferResult.StreamComplete(plan.planId()); } @@ -392,7 +392,6 @@ public class CoordinatedTransfer enum State { INIT, - STREAM_NOOP, STREAM_FAILED, STREAM_COMPLETE; } @@ -418,12 +417,6 @@ public class CoordinatedTransfer return new SingleTransferResult(State.STREAM_COMPLETE, planId); } - @VisibleForTesting - static SingleTransferResult Noop() - { - return new SingleTransferResult(State.STREAM_NOOP, null); - } - public SingleTransferResult streamFailed(TimeUUID planId) { return new SingleTransferResult(State.STREAM_FAILED, planId); @@ -434,6 +427,17 @@ public class CoordinatedTransfer return planId; } + @Override + public boolean equals(Object that) + { + return that instanceof SingleTransferResult && equals((SingleTransferResult) that); + } + + public boolean equals(SingleTransferResult that) + { + return state == that.state && planId.equals(that.planId); + } + @Override public String toString() {