From 3ae74e250a245485c5895d493b4ecc3c6cf61073 Mon Sep 17 00:00:00 2001 From: Aleksey Yeshchenko Date: Thu, 27 Aug 2026 14:56:38 +0100 Subject: [PATCH 1/4] Eliminate redundant advances for previously repaired ranges; clean up logging --- .../MutationTrackingMigrationState.java | 4 ++- .../MutationTrackingRepairHandler.java | 27 ++++++++++++++----- .../AdvanceMutationTrackingMigration.java | 10 +++---- .../AdvanceMutationTrackingMigrationTest.java | 15 +++-------- 4 files changed, 31 insertions(+), 25 deletions(-) diff --git a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationState.java b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationState.java index 1c500c6b1009..68ac0ea2a184 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationState.java +++ b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationState.java @@ -177,8 +177,10 @@ public MutationTrackingMigrationState withRangesRepairedForTable(@Nonnull String if (info == null) return this; - // Subtract repaired ranges from table's pending set + // subtract repaired ranges from table's pending set; noop is nothing's changed KeyspaceMigrationInfo updated = info.withRangesRepairedForTable(epoch, tableId, repairedRanges); + if (info == updated) + return this; // if all tables fully repaired, remove keyspace (migration complete) if (updated.isComplete()) diff --git a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java index b55c362a3021..ee8e43879cbc 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java +++ b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java @@ -63,8 +63,8 @@ public void onSuccess(RepairResult repairResult) if (migrationInfo == null) { - logger.info("Repair session {} (parent session {}) completed for {}.{} but the keyspace is not migrating, not advancing mutation tracking migration", - desc.sessionId, desc.parentSessionId, keyspace, tableName); + logger.debug("Repair session {} (parent session {}) completed for {}.{} but the keyspace is not migrating, not advancing mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName); return; } @@ -78,10 +78,13 @@ public void onSuccess(RepairResult repairResult) return; } - if (migrationInfo.getPendingRangesForTable(tableMetadata.id).isEmpty()) + NormalizedRanges pendingRanges = migrationInfo.getPendingRangesForTable(tableMetadata.id); + NormalizedRanges repairedPendingRanges = pendingRanges.intersection(NormalizedRanges.normalizedRanges(repairedRanges)); + + if (repairedPendingRanges.isEmpty()) { - logger.info("Repair session {} (parent session {}) completed for {}.{} but the table has no ranges left to migrate, not advancing mutation tracking migration", - desc.sessionId, desc.parentSessionId, keyspace, tableName); + logger.info("Repair session {} (parent session {}) completed for {}.{} but none of the repaired ranges {} are still pending migration (pending: {}), not advancing mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName, repairedRanges, pendingRanges); return; } @@ -104,7 +107,17 @@ public void onSuccess(RepairResult repairResult) } ClusterMetadata committed = ClusterMetadataService.instance().commit( - new AdvanceMutationTrackingMigration(keyspace, tableMetadata.id, repairedRanges)); + new AdvanceMutationTrackingMigration(keyspace, tableMetadata.id, repairedPendingRanges), + ignore -> ignore, + (code, message) -> + { + logger.info("Repair session {} (parent session {}) did not advance mutation tracking migration of {}.{}: {} ({})", + desc.sessionId, desc.parentSessionId, keyspace, tableName, message, code); + return null; + }); + + if (committed == null) + return; // Report from the metadata commit returned, not current(), which races with other epochs KeyspaceMigrationInfo advanced = committed.mutationTrackingMigrationState.getKeyspaceInfo(keyspace); @@ -119,7 +132,7 @@ public void onSuccess(RepairResult repairResult) "contributed {} range(s) {}; {} range(s) remain to be repaired {}; {} range(s) already repaired {}; " + "{} table(s) in the keyspace still migrating", desc.sessionId, desc.parentSessionId, keyspace, tableName, committed.epoch, - repairedRanges.size(), repairedRanges, + repairedPendingRanges.size(), repairedPendingRanges, pending.size(), pending, repaired.size(), repaired, keyspaceComplete ? 0 : advanced.pendingRangesPerTable.size()); diff --git a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java index 9245a113494d..600651ec9e52 100644 --- a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java +++ b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java @@ -94,10 +94,7 @@ public Result execute(ClusterMetadata prev) KeyspaceMigrationInfo ksInfo = prev.mutationTrackingMigrationState.getKeyspaceInfo(keyspace); if (ksInfo == null) - { - logger.warn("Attempted to advance mutation tracking migration for keyspace {} table {} which is not migrating", keyspace, tableId); return new Rejected(INVALID, String.format("Keyspace %s is not migrating", keyspace)); - } Transformer transformer = prev.transformer(); @@ -105,8 +102,11 @@ public Result execute(ClusterMetadata prev) MutationTrackingMigrationState newState = prev.mutationTrackingMigrationState .withRangesRepairedForTable(keyspace, tableId, repairedRanges, transformer.epoch()); - logger.info("Advanced mutation tracking migration for keyspace {}, table {}: {} ranges repaired", - keyspace, tableId, repairedRanges.size()); + if (newState == prev.mutationTrackingMigrationState) + { + return new Rejected(INVALID, String.format("Keyspace %s table %s has no pending ranges intersecting %s", + keyspace, tableId, repairedRanges)); + } return Transformation.success( transformer.with(newState), diff --git a/test/unit/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigrationTest.java b/test/unit/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigrationTest.java index da33dc9958c9..948395f40a83 100644 --- a/test/unit/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigrationTest.java +++ b/test/unit/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigrationTest.java @@ -193,18 +193,9 @@ public void testAdvanceRangesForWrongTable() Transformation.Result result = transformation.execute(prev); - // confirm noop - assertTrue(result.isSuccess()); - ClusterMetadata updated = result.success().metadata; - - KeyspaceMigrationInfo expected = createExpectedInfo( - "test_ks", - testTableId, - Collections.singleton(fullRing()), - epoch1 - ); - - assertEquals(expected, updated.mutationTrackingMigrationState.getKeyspaceInfo("test_ks")); + // confirm rejection + assertTrue(result.isRejected()); + assertTrue(result.rejected().reason.contains("no pending ranges intersecting")); } @Test From ceb7f6412cb6f24175276b5d1a860d8c99e6c1b8 Mon Sep 17 00:00:00 2001 From: Aleksey Yeshchenko Date: Thu, 27 Aug 2026 17:32:13 +0100 Subject: [PATCH 2/4] Add a virtual table to expose pending and migrated ranges --- .../db/virtual/MutationTrackingTables.java | 95 ++++++++++++++++++- .../migration/KeyspaceMigrationInfo.java | 11 ++- .../AdvanceMutationTrackingMigration.java | 4 - 3 files changed, 101 insertions(+), 9 deletions(-) diff --git a/src/java/org/apache/cassandra/db/virtual/MutationTrackingTables.java b/src/java/org/apache/cassandra/db/virtual/MutationTrackingTables.java index 406693d23e81..e0498028ee7c 100644 --- a/src/java/org/apache/cassandra/db/virtual/MutationTrackingTables.java +++ b/src/java/org/apache/cassandra/db/virtual/MutationTrackingTables.java @@ -22,15 +22,21 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.marshal.BooleanType; import org.apache.cassandra.db.marshal.Int32Type; +import org.apache.cassandra.db.marshal.ListType; import org.apache.cassandra.db.marshal.LongType; import org.apache.cassandra.db.marshal.UTF8Type; +import org.apache.cassandra.db.marshal.UUIDType; import org.apache.cassandra.dht.LocalPartitioner; +import org.apache.cassandra.dht.NormalizedRanges; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; import org.apache.cassandra.journal.ActiveSegment; import org.apache.cassandra.journal.Segment; import org.apache.cassandra.replication.CoordinatorLog; @@ -39,12 +45,16 @@ import org.apache.cassandra.replication.MutationTrackingService; import org.apache.cassandra.replication.Shard; import org.apache.cassandra.replication.ShortMutationId; +import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.service.replication.migration.KeyspaceMigrationInfo; +import org.apache.cassandra.tcm.ClusterMetadata; public class MutationTrackingTables { public static final String MUTATION_JOURNAL = "mutation_journal"; public static final String MUTATION_TRACKING_SHARDS = "mutation_tracking_shards"; + public static final String MUTATION_TRACKING_MIGRATION_STATE = "mutation_tracking_migration_state"; private MutationTrackingTables() {} @@ -53,7 +63,9 @@ public static Collection getAll(String keyspace) if (!DatabaseDescriptor.getMutationTrackingEnabled()) return Collections.emptyList(); - return List.of(new MutationJournalTable(keyspace), new MutationTrackingShardsTable(keyspace)); + return List.of(new MutationJournalTable(keyspace), + new MutationTrackingShardsTable(keyspace), + new MutationTrackingMigrationStateTable(keyspace)); } public static final class MutationJournalTable extends AbstractVirtualTable @@ -183,4 +195,85 @@ public DataSet data(DecoratedKey key) return result; } } + + /** + * Mutation tracking migration progress (held in {@link ClusterMetadata}). + */ + public static class MutationTrackingMigrationStateTable extends AbstractVirtualTable + { + private static final String KEYSPACE_NAME = "keyspace_name"; + private static final String TABLE_NAME = "table_name"; + private static final String TABLE_ID = "table_id"; + private static final String STARTED_AT_EPOCH = "started_at_epoch"; + private static final String PENDING_RANGES = "pending_ranges"; + private static final String MIGRATED_RANGES = "migrated_ranges"; + + private static final ListType STRING_LIST_TYPE = ListType.getInstance(UTF8Type.instance, false); + + MutationTrackingMigrationStateTable(String keyspace) + { + super(TableMetadata.builder(keyspace, MUTATION_TRACKING_MIGRATION_STATE) + .comment("ranges still to be repaired for in-progress mutation tracking migrations") + .kind(TableMetadata.Kind.VIRTUAL) + .partitioner(new LocalPartitioner(UTF8Type.instance)) + .addPartitionKeyColumn(KEYSPACE_NAME, UTF8Type.instance) + .addClusteringColumn(TABLE_NAME, UTF8Type.instance) + .addRegularColumn(TABLE_ID, UUIDType.instance) + .addRegularColumn(STARTED_AT_EPOCH, LongType.instance) + .addRegularColumn(PENDING_RANGES, STRING_LIST_TYPE) + .addRegularColumn(MIGRATED_RANGES, STRING_LIST_TYPE) + .build()); + } + + @Override + public DataSet data() + { + SimpleDataSet result = new SimpleDataSet(metadata()); + ClusterMetadata metadata = ClusterMetadata.current(); + + for (KeyspaceMigrationInfo info : metadata.mutationTrackingMigrationState.keyspaceInfo.values()) + addTableRows(metadata, info, result); + + return result; + } + + @Override + public DataSet data(DecoratedKey key) + { + String keyspaceName = UTF8Type.instance.compose(key.getKey()); + SimpleDataSet result = new SimpleDataSet(metadata()); + ClusterMetadata metadata = ClusterMetadata.current(); + + KeyspaceMigrationInfo info = metadata.mutationTrackingMigrationState.getKeyspaceInfo(keyspaceName); + if (info != null) + addTableRows(metadata, info, result); + + return result; + } + + private static void addTableRows(ClusterMetadata metadata, KeyspaceMigrationInfo info, SimpleDataSet result) + { + NormalizedRanges fullRing = KeyspaceMigrationInfo.fullRing(); + for (Map.Entry> entry : info.pendingRangesPerTable.entrySet()) + { + TableId tid = entry.getKey(); + NormalizedRanges pendingRanges = entry.getValue(); + + TableMetadata tm = metadata.schema.getTableMetadata(tid); + if (tm == null) + continue; + + result.row(info.keyspace, tm.name) + .column(TABLE_ID, tid.asUUID()) + .column(STARTED_AT_EPOCH, info.startedAtEpoch.getEpoch()) + .column(PENDING_RANGES, rangesToStrings(pendingRanges)) + .column(MIGRATED_RANGES, rangesToStrings(fullRing.subtract(pendingRanges))); + } + } + + private static List rangesToStrings(NormalizedRanges ranges) + { + return ranges.stream().map(Range::toString).collect(Collectors.toList()); + } + } } diff --git a/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java b/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java index 65b08d0d78ac..ad5fe763833b 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java +++ b/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java @@ -174,14 +174,17 @@ public KeyspaceMigrationInfo withRangesRepairedForTable(@Nonnull Epoch repairSta @Nonnull TableId tableId, @Nonnull Collection> repairedRanges) { - if (repairStartedEpoch.isBefore(startedAtEpoch)) - return this; + // TODO (expected): do something about this? nuke or serialize the correct epoch alongised the transformation? + // this was dead code; repairStartedEpoch as passed was always next transformation's epoch, + // and it was always > startedAtEpoch, guarding against nothing; + // there is an epoch eligibility check in MutationTrackingRepairHandler in onSuccess(), but it is + // insufficient in face of potential race conditions (AY) + // if (repairStartedEpoch.isBefore(startedAtEpoch)) + // return this; NormalizedRanges currentPendingForTable = pendingRangesPerTable.get(tableId); if (currentPendingForTable == null) - { return this; - } NormalizedRanges normalizedRepaired = NormalizedRanges.normalizedRanges(repairedRanges); NormalizedRanges remainingForTable = currentPendingForTable.subtract(normalizedRepaired); diff --git a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java index 600651ec9e52..987ff528f232 100644 --- a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java +++ b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java @@ -23,9 +23,6 @@ import javax.annotation.Nonnull; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; @@ -57,7 +54,6 @@ */ public class AdvanceMutationTrackingMigration implements Transformation { - private static final Logger logger = LoggerFactory.getLogger(AdvanceMutationTrackingMigration.class); public static final Serializer serializer = new Serializer(); @Nonnull From 5dc67b545ea1e3e3bc4e00d3cbd7fdd79a04b62c Mon Sep 17 00:00:00 2001 From: Aleksey Yeshchenko Date: Wed, 2 Sep 2026 16:18:42 +0100 Subject: [PATCH 3/4] Fix certain repair sessions incorrectly advancing migration state --- .../org/apache/cassandra/repair/RepairJob.java | 2 +- .../MutationTrackingMigrationRepairResult.java | 16 +++++++++++++++- 2 files changed, 16 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/repair/RepairJob.java b/src/java/org/apache/cassandra/repair/RepairJob.java index db7079efc1be..35f764c48992 100644 --- a/src/java/org/apache/cassandra/repair/RepairJob.java +++ b/src/java/org/apache/cassandra/repair/RepairJob.java @@ -295,7 +295,7 @@ public void onSuccess(List stats) cfs.metric.repairsCompleted.inc(); logger.info("Completing repair with excludedDeadNodes {}", session.excludedDeadNodes); ConsensusMigrationRepairResult cmrs = ConsensusMigrationRepairResult.fromRepair(repairStartingEpoch, getUnchecked(accordRepair), session.repairData, doPaxosRepair, doAccordRepair, session.excludedDeadNodes, session.isIncremental); - MutationTrackingMigrationRepairResult mtmrs = MutationTrackingMigrationRepairResult.fromRepair(repairStartingEpoch, session.excludedDeadNodes, session.previewKind.isPreview()); + MutationTrackingMigrationRepairResult mtmrs = MutationTrackingMigrationRepairResult.fromRepair(repairStartingEpoch, session.repairData, session.allReplicas, session.pullRepair, session.excludedDeadNodes, session.previewKind.isPreview()); trySuccess(new RepairResult(desc, stats, cmrs, mtmrs)); } diff --git a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java index e0ec93438249..fb5009f4fb22 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java +++ b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java @@ -33,6 +33,12 @@ public class MutationTrackingMigrationRepairResult new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "dead nodes were excluded from the repair"); private static final MutationTrackingMigrationRepairResult PREVIEW = new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "the repair was a preview"); + private static final MutationTrackingMigrationRepairResult NO_DATA_REPAIR = + new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "the repair did not repair data (paxos-only or accord-only repair)"); + private static final MutationTrackingMigrationRepairResult NOT_ALL_REPLICAS = + new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "the repair did not include all replicas (-local, -dc, or -hosts repair)"); + private static final MutationTrackingMigrationRepairResult PULL_REPAIR = + new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "the repair only streamed data one way (-pull repair)"); public final Epoch minEpoch; public final boolean eligible; @@ -48,10 +54,18 @@ private MutationTrackingMigrationRepairResult(Epoch minEpoch, boolean eligible, this.ineligibleReason = ineligibleReason; } - public static MutationTrackingMigrationRepairResult fromRepair(Epoch minEpoch, boolean deadNodesExcluded, boolean isPreview) + public static MutationTrackingMigrationRepairResult fromRepair(Epoch minEpoch, + boolean dataRepaired, + boolean allReplicas, + boolean pullRepair, + boolean deadNodesExcluded, + boolean isPreview) { if (deadNodesExcluded) return DEAD_NODES_EXCLUDED; if (isPreview) return PREVIEW; + if (!dataRepaired) return NO_DATA_REPAIR; + if (!allReplicas) return NOT_ALL_REPLICAS; + if (pullRepair) return PULL_REPAIR; return new MutationTrackingMigrationRepairResult(minEpoch, true, null); } } From 1d00312cecb99a046cf204d180cc06a7fc961e77 Mon Sep 17 00:00:00 2001 From: Aleksey Yeshchenko Date: Tue, 1 Sep 2026 11:18:09 +0100 Subject: [PATCH 4/4] Add missing AdvanceMutationTrackingMigration#toString() --- .../AdvanceMutationTrackingMigration.java | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java index 987ff528f232..89e635669d54 100644 --- a/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java +++ b/src/java/org/apache/cassandra/tcm/transformations/AdvanceMutationTrackingMigration.java @@ -109,6 +109,16 @@ public Result execute(ClusterMetadata prev) LockedRanges.AffectedRanges.EMPTY); } + @Override + public String toString() + { + return "AdvanceMutationTrackingMigration{" + + "keyspace='" + keyspace + '\'' + + ", tableId=" + tableId + + ", repairedRanges=" + repairedRanges + + '}'; + } + public static class Serializer implements AsymmetricMetadataSerializer { @Override