From 10ed20bdc3feba3ebe63f7b1c6f742c70d1828e5 Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Mon, 9 Jun 2025 11:43:22 +0100 Subject: [PATCH] Fixes - Don't assume stillExecutes applies to remote request - must mark unavailable any keys we don't have the txn definition for - Don't exit notifyManagedPreBootstrap notify loop early, as could have later transactions with lower txnId - CoordinateEphemeralRead must retry if insufficient responses from replicas that still own the range - Burn test terminates while non-recurring tasks pending on command store queues - Infinite recovery due to not sending InformDurable when partially truncated - Incorrect participants when invoking removeRedundantDependencies - Infinite bootstrap loop on retired topology Improve - Do not perform linear filters of CommandStores or Topologies - Some default implementations of forEach/iterator - TopologyRetiredException messages patch by Benedict; reviewed by Alex Petrov for CASSANDRA-20707 --- modules/accord | 2 +- .../db/virtual/AccordDebugKeyspace.java | 2 +- .../index/accord/RouteSSTableManager.java | 17 ++++++++++++++++- .../service/accord/AccordService.java | 14 ++++++++------ .../service/accord/AccordVerbHandler.java | 4 +++- .../accord/interop/AccordInteropAdapter.java | 10 +++++----- .../accord/interop/AccordInteropPersist.java | 4 ++-- .../GetEphmrlReadDepsSerializers.java | 19 ++++++++++++++----- .../paxos/AccordSimulationRunner.java | 2 ++ 9 files changed, 52 insertions(+), 22 deletions(-) diff --git a/modules/accord b/modules/accord index c5a984cfe4..bf85660dcc 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit c5a984cfe41bb8d8f1d7a4cb446c194f829a5dd1 +Subproject commit bf85660dccaece9fcb3dd319b6f82265d780378f diff --git a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java index 6e359a74e4..4bab8d2c64 100644 --- a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java @@ -431,7 +431,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace SimpleDataSet ds = new SimpleDataSet(metadata()); for (CommandStore commandStore : commandStores.all()) { - DefaultProgressLog.ImmutableView view = (DefaultProgressLog.ImmutableView) commandStore.unsafeProgressLog(); + DefaultProgressLog.ImmutableView view = ((DefaultProgressLog) commandStore.unsafeProgressLog()).immutableView(); TableId tableId = ((AccordCommandStore)commandStore).tableId(); TableMetadata tableMetadata = tableMetadata(tableId); while (view.advance()) diff --git a/src/java/org/apache/cassandra/index/accord/RouteSSTableManager.java b/src/java/org/apache/cassandra/index/accord/RouteSSTableManager.java index 3fdb884494..1321d9dc8b 100644 --- a/src/java/org/apache/cassandra/index/accord/RouteSSTableManager.java +++ b/src/java/org/apache/cassandra/index/accord/RouteSSTableManager.java @@ -28,8 +28,13 @@ import java.util.Map; import java.util.NavigableSet; import java.util.TreeSet; +import org.apache.cassandra.io.FSReadError; import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.util.File; import org.apache.cassandra.schema.TableId; +import org.apache.cassandra.utils.ByteArrayUtil; + +import static org.apache.cassandra.index.accord.IndexDescriptor.IndexComponent.CINTIA_SORTED_LIST; public class RouteSSTableManager implements SSTableManager { @@ -85,7 +90,17 @@ public class RouteSSTableManager implements SSTableManager Group group = new Group(storeId, tableId); TreeSet matches = new TreeSet<>(); for (SSTableIndex index : sstables.values()) - matches.addAll(index.search(group, start, startInclusive, end, endInclusive)); + { + try + { + matches.addAll(index.search(group, start, startInclusive, end, endInclusive)); + } + catch (Throwable t) + { + File file = index.id.fileFor(CINTIA_SORTED_LIST); + throw new FSReadError("Failed to search range index " + file + " for " + (startInclusive ? "[" : "(") + ByteArrayUtil.bytesToHex(start) + "..." + ByteArrayUtil.bytesToHex(end) + (endInclusive ? "]" : ")"), t, file); + } + } return matches; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index f85aee62af..4bf89c9b1e 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -77,6 +77,7 @@ import accord.primitives.Status; import accord.primitives.Timestamp; import accord.primitives.Txn; import accord.primitives.TxnId; +import accord.topology.Shard; import accord.topology.Topology; import accord.topology.TopologyManager; import accord.utils.DefaultRandom; @@ -962,12 +963,13 @@ public class AccordService implements IAccordService, Shutdownable { // Need to make sure no existing txn are still being processed for this table... this is only used by DROP TABLE so NEW txn are expected to be blocked, so just need to "wait" for existing ones to complete Topology topology = node.topology().current(); - List rangeList = topology.reduce(new ArrayList<>(), - s -> ((TokenRange) s.range).table().equals(id), - (accum, s) -> { - accum.add((TokenRange) s.range); - return accum; - }); + List rangeList = new ArrayList<>(); + for (Shard shard : topology.shards()) + { + TokenRange range = (TokenRange) shard.range; + if (id.equals(range.table())) + rangeList.add(range); + } if (rangeList.isEmpty()) return; // nothing to see here Ranges ranges = Ranges.of(rangeList.toArray(accord.primitives.Range[]::new)); diff --git a/src/java/org/apache/cassandra/service/accord/AccordVerbHandler.java b/src/java/org/apache/cassandra/service/accord/AccordVerbHandler.java index 9562aea86b..2283a366d0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordVerbHandler.java +++ b/src/java/org/apache/cassandra/service/accord/AccordVerbHandler.java @@ -63,10 +63,12 @@ public class AccordVerbHandler implements IVerbHandler Node.Id fromNodeId = endpointMapper.mappedId(message.from()); long waitForEpoch = request.waitForEpoch(); if (node.topology().hasAtLeastEpoch(waitForEpoch)) + { request.process(node, fromNodeId, message.header); + } else { - node.withEpoch(waitForEpoch, (ignored, withEpochFailure) -> { + node.withEpochAtLeast(waitForEpoch, (ignored, withEpochFailure) -> { if (withEpochFailure != null) throw new RuntimeException("Timed out waiting for epoch when processing message from " + fromNodeId + " to " + node + " message " + message, withEpochFailure); request.process(node, fromNodeId, message.header); diff --git a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropAdapter.java b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropAdapter.java index ca466105eb..009c592fcd 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropAdapter.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropAdapter.java @@ -94,12 +94,12 @@ public class AccordInteropAdapter extends TxnAdapter } @Override - public void persist(Node node, Topologies any, Route require, Route sendTo, SelectNodeOwnership selectSendTo, FullRoute route, Ballot ballot, TxnId txnId, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, BiConsumer callback) + public void persist(Node node, Topologies any, Route require, Route sendTo, SelectNodeOwnership selectSendTo, FullRoute route, Ballot ballot, TxnId txnId, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, boolean informDurableOnDone, BiConsumer callback) { - if (applyKind == Minimal && doInteropPersist(node, any, require, sendTo, selectSendTo, ballot, txnId, txn, executeAt, deps, writes, result, route, callback)) + if (applyKind == Minimal && doInteropPersist(node, any, require, sendTo, selectSendTo, ballot, txnId, txn, executeAt, deps, writes, result, route, informDurableOnDone, callback)) return; - super.persist(node, any, require, sendTo, selectSendTo, route, ballot, txnId, txn, executeAt, deps, writes, result, callback); + super.persist(node, any, require, sendTo, selectSendTo, route, ballot, txnId, txn, executeAt, deps, writes, result, informDurableOnDone, callback); } @@ -116,7 +116,7 @@ public class AccordInteropAdapter extends TxnAdapter return true; } - private boolean doInteropPersist(Node node, Topologies any, Route require, Route sendTo, SelectNodeOwnership selectSendTo, Ballot ballot, TxnId txnId, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, BiConsumer callback) + private boolean doInteropPersist(Node node, Topologies any, Route require, Route sendTo, SelectNodeOwnership selectSendTo, Ballot ballot, TxnId txnId, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, boolean informDurableOnDone, BiConsumer callback) { Update update = txn.update(); ConsistencyLevel consistencyLevel = update instanceof AccordUpdate ? ((AccordUpdate) update).cassandraCommitCL() : null; @@ -124,7 +124,7 @@ public class AccordInteropAdapter extends TxnAdapter return false; Topologies all = execution(node, any, sendTo, selectSendTo, fullRoute, txnId, executeAt); - new AccordInteropPersist(node, all, txnId, require, ballot, txn, executeAt, deps, writes, result, fullRoute, consistencyLevel, callback) + new AccordInteropPersist(node, all, txnId, require, ballot, txn, executeAt, deps, writes, result, fullRoute, consistencyLevel, informDurableOnDone, callback) .start(Minimal, any, writes, result); return true; } diff --git a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java index 967438ea56..bad44822ac 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java @@ -109,9 +109,9 @@ public class AccordInteropPersist extends Persist private final ConsistencyLevel consistencyLevel; private CallbackHolder callback; - public AccordInteropPersist(Node node, Topologies topologies, TxnId txnId, Route sendTo, Ballot ballot, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, ConsistencyLevel consistencyLevel, BiConsumer clientCallback) + public AccordInteropPersist(Node node, Topologies topologies, TxnId txnId, Route sendTo, Ballot ballot, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, ConsistencyLevel consistencyLevel, boolean informDurableOnDone, BiConsumer clientCallback) { - super(node, topologies, txnId, ballot, sendTo, txn, executeAt, deps, writes, result, fullRoute, AccordInteropApply.FACTORY); + super(node, topologies, txnId, ballot, sendTo, txn, executeAt, deps, writes, result, fullRoute, informDurableOnDone, AccordInteropApply.FACTORY); Invariants.requireArgument(consistencyLevel == ConsistencyLevel.QUORUM || consistencyLevel == ConsistencyLevel.ALL || consistencyLevel == ConsistencyLevel.SERIAL || consistencyLevel == ConsistencyLevel.ONE); this.consistencyLevel = consistencyLevel; registerClientCallback(result, clientCallback); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/GetEphmrlReadDepsSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/GetEphmrlReadDepsSerializers.java index 9f90274a94..670a8cdee3 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/GetEphmrlReadDepsSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/GetEphmrlReadDepsSerializers.java @@ -60,16 +60,21 @@ public class GetEphmrlReadDepsSerializers @Override public void serialize(GetEphemeralReadDepsOk reply, DataOutputPlus out) throws IOException { - DepsSerializers.deps.serialize(reply.deps, out); out.writeUnsignedVInt(reply.latestEpoch); + out.writeBoolean(reply.deps != null); + if (reply.deps == null) + return; + DepsSerializers.deps.serialize(reply.deps, out); out.writeUnsignedVInt32(reply.flags.bits()); } @Override public GetEphemeralReadDepsOk deserialize(DataInputPlus in) throws IOException { - Deps deps = DepsSerializers.deps.deserialize(in); long latestEpoch = in.readUnsignedVInt(); + if (!in.readBoolean()) + return new GetEphemeralReadDepsOk(latestEpoch); + Deps deps = DepsSerializers.deps.deserialize(in); ExecuteFlags flags = ExecuteFlags.get(in.readUnsignedVInt32()); return new GetEphemeralReadDepsOk(deps, latestEpoch, flags); } @@ -77,9 +82,13 @@ public class GetEphmrlReadDepsSerializers @Override public long serializedSize(GetEphemeralReadDepsOk reply) { - return DepsSerializers.deps.serializedSize(reply.deps) - + TypeSizes.sizeofUnsignedVInt(reply.latestEpoch) - + TypeSizes.sizeofUnsignedVInt(reply.flags.bits()); + long size = 1 + TypeSizes.sizeofUnsignedVInt(reply.latestEpoch); + if (reply.deps != null) + { + size += DepsSerializers.deps.serializedSize(reply.deps) + + TypeSizes.sizeofUnsignedVInt(reply.flags.bits()); + } + return size; } }; } diff --git a/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java b/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java index e782848ea4..147cd3ce59 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java +++ b/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java @@ -29,6 +29,7 @@ import org.slf4j.LoggerFactory; import io.airlift.airline.Cli; import io.airlift.airline.Command; import org.apache.cassandra.config.CassandraRelevantProperties; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.simulator.SimulationRunner; import org.apache.cassandra.simulator.SimulatorUtils; import org.apache.cassandra.utils.StorageCompatibilityMode; @@ -41,6 +42,7 @@ public class AccordSimulationRunner extends SimulationRunner public static void beforeAll() { CassandraRelevantProperties.JUNIT_STORAGE_COMPATIBILITY_MODE.setString(StorageCompatibilityMode.NONE.toString()); + DatabaseDescriptor.clientInitialization(); } @Command(name = "run")