From 1b3c3f32a4a3130a98022f14ad796190ea42f499 Mon Sep 17 00:00:00 2001 From: Caleb Rackliffe Date: Wed, 8 May 2024 12:24:33 -0500 Subject: [PATCH] post-rebase fixes, mostly around CASSANDRA-19341 and CASSANDRA-19567 --- .gitmodules | 2 +- modules/accord | 2 +- .../apache/cassandra/config/AccordSpec.java | 60 +++++ .../cassandra/config/DatabaseDescriptor.java | 5 + .../dht/IPartitionerDependentSerializer.java | 25 +-- src/java/org/apache/cassandra/dht/Token.java | 16 ++ .../cassandra/index/accord/IndexMetrics.java | 4 +- .../io/IVersionedAsymmetricSerializer.java | 6 +- .../org/apache/cassandra/journal/Metrics.java | 6 +- .../metrics/CassandraMetricsRegistry.java | 4 + src/java/org/apache/cassandra/net/Verb.java | 8 +- .../repair/messages/SyncResponse.java | 9 +- .../service/accord/AccordJournal.java | 124 +++++++---- .../service/accord/AccordKeyspace.java | 12 +- .../service/accord/AccordService.java | 7 +- .../service/accord/async/AsyncOperation.java | 4 +- .../service/accord/async/ExecutionOrder.java | 131 ++++++----- .../migration/ConsensusRequestRouter.java | 2 +- .../cassandra/streaming/SessionSummary.java | 4 +- .../streaming/StreamDeserializingTask.java | 3 +- .../cassandra/streaming/StreamSummary.java | 20 +- .../streaming/messages/CompleteMessage.java | 3 +- .../messages/IncomingStreamMessage.java | 3 +- .../streaming/messages/KeepAliveMessage.java | 3 +- .../messages/OutgoingStreamMessage.java | 3 +- .../streaming/messages/PrepareAckMessage.java | 3 +- .../messages/PrepareSynAckMessage.java | 5 +- .../streaming/messages/PrepareSynMessage.java | 5 +- .../streaming/messages/ReceivedMessage.java | 3 +- .../messages/SessionFailedMessage.java | 3 +- .../streaming/messages/StreamInitMessage.java | 3 +- .../streaming/messages/StreamMessage.java | 7 +- .../test/tcm/RepairMetadataKeyspaceTest.java | 2 + .../test/AccordJournalSimulationTest.java | 3 +- .../config/DatabaseDescriptorRefTest.java | 4 + .../CompactionAccordIteratorsTest.java | 8 +- .../RepairMessageSerializationsTest.java | 19 +- .../cassandra/service/SerializationsTest.java | 8 +- .../service/accord/AccordTestUtils.java | 3 +- .../cassandra/service/accord/MockJournal.java | 28 +++ .../accord/SimulatedAccordCommandStore.java | 13 +- .../async/SimulatedAsyncOperationTest.java | 207 ++++++++++++++++++ .../async/StreamingInboundHandlerTest.java | 4 +- .../MetadataSnapshotListenerTest.java | 4 +- 44 files changed, 598 insertions(+), 200 deletions(-) create mode 100644 test/unit/org/apache/cassandra/service/accord/async/SimulatedAsyncOperationTest.java diff --git a/.gitmodules b/.gitmodules index 60a9510e7a..616dacf610 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,4 +1,4 @@ [submodule "modules/accord"] path = modules/accord - url = ../cassandra-accord.git + url = https://github.com/apache/cassandra-accord.git branch = trunk diff --git a/modules/accord b/modules/accord index 202e673583..256b35e27d 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 202e67358396a1e413e29498bea71047bd586d06 +Subproject commit 256b35e27d170db9fcd8024d5678b4f6e9d3a956 diff --git a/src/java/org/apache/cassandra/config/AccordSpec.java b/src/java/org/apache/cassandra/config/AccordSpec.java index d6fb1a5011..b035b0b9b5 100644 --- a/src/java/org/apache/cassandra/config/AccordSpec.java +++ b/src/java/org/apache/cassandra/config/AccordSpec.java @@ -18,6 +18,8 @@ package org.apache.cassandra.config; +import com.fasterxml.jackson.annotation.JsonIgnore; +import org.apache.cassandra.journal.Params; import org.apache.cassandra.service.consensus.TransactionalMode; public class AccordSpec @@ -71,4 +73,62 @@ public class AccordSpec public TransactionalMode default_transactional_mode = TransactionalMode.off; public boolean ephemeralReadEnabled = false; public boolean state_cache_listener_jfr_enabled = true; + public final JournalSpec journal = new JournalSpec(); + + public static class JournalSpec implements Params + { + public int segmentSize = 32 << 20; + public FailurePolicy failurePolicy = FailurePolicy.STOP; + public FlushMode flushMode = FlushMode.BATCH; + public DurationSpec.IntMillisecondsBound flushPeriod; // pulls default from 'commitlog_sync_period' + public DurationSpec.IntMillisecondsBound periodicFlushLagBlock = new DurationSpec.IntMillisecondsBound("1500ms"); + + @Override + public int segmentSize() + { + return segmentSize; + } + + @Override + public FailurePolicy failurePolicy() + { + return failurePolicy; + } + + @Override + public FlushMode flushMode() + { + return flushMode; + } + + @JsonIgnore + @Override + public int flushPeriodMillis() + { + return flushPeriod == null ? DatabaseDescriptor.getCommitLogSyncPeriod() + : flushPeriod.toMilliseconds(); + } + + @JsonIgnore + @Override + public int periodicFlushLagBlock() + { + return periodicFlushLagBlock.toMilliseconds(); + } + + /** + * This is required by the journal, but we don't have multiple versions, so block it from showing up, so we don't need to worry about maintaining it + */ + @JsonIgnore + @Override + public int userVersion() + { + /* + * NOTE: when accord journal version gets bumped, expose it via yaml. + * This way operators can force previous version on upgrade, temporarily, + * to allow easier downgrades if something goes wrong. + */ + return 1; + } + } } diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index b026ad5286..202469fbff 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -3661,6 +3661,11 @@ public class DatabaseDescriptor return conf.paxos_topology_repair_strict_each_quorum; } + public static AccordSpec getAccord() + { + return conf.accord; + } + public static AccordSpec.TransactionalRangeMigration getTransactionalRangeMigration() { return conf.accord.range_migration; diff --git a/src/java/org/apache/cassandra/dht/IPartitionerDependentSerializer.java b/src/java/org/apache/cassandra/dht/IPartitionerDependentSerializer.java index 5c75788c0b..a70eb83771 100644 --- a/src/java/org/apache/cassandra/dht/IPartitionerDependentSerializer.java +++ b/src/java/org/apache/cassandra/dht/IPartitionerDependentSerializer.java @@ -19,8 +19,8 @@ package org.apache.cassandra.dht; import java.io.IOException; +import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; -import org.apache.cassandra.io.util.DataOutputPlus; /** * Versioned serializer where the serialization depends on partitioner. @@ -28,18 +28,8 @@ import org.apache.cassandra.io.util.DataOutputPlus; * On serialization the partitioner is given by the entity being serialized. To deserialize the partitioner used must * be known to the calling method. */ -public interface IPartitionerDependentSerializer +public interface IPartitionerDependentSerializer extends IVersionedSerializer { - /** - * Serialize the specified type into the specified DataOutputStream instance. - * - * @param t type that needs to be serialized - * @param out DataOutput into which serialization needs to happen. - * @param version protocol version - * @throws java.io.IOException if serialization fails - */ - public void serialize(T t, DataOutputPlus out, int version) throws IOException; - /** * Deserialize into the specified DataInputStream instance. * @param in DataInput from which deserialization needs to happen. @@ -51,11 +41,8 @@ public interface IPartitionerDependentSerializer */ public T deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException; - /** - * Calculate serialized size of object without actually serializing. - * @param t object to calculate serialized size - * @param version protocol version - * @return serialized size of object t - */ - public long serializedSize(T t, int version); + default T deserialize(DataInputPlus in, int version) throws IOException + { + return deserialize(in, null, version); + } } diff --git a/src/java/org/apache/cassandra/dht/Token.java b/src/java/org/apache/cassandra/dht/Token.java index 1df7171f9c..b3fee2a080 100644 --- a/src/java/org/apache/cassandra/dht/Token.java +++ b/src/java/org/apache/cassandra/dht/Token.java @@ -20,6 +20,12 @@ package org.apache.cassandra.dht; import java.io.IOException; import java.io.Serializable; import java.nio.ByteBuffer; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +import com.google.common.collect.Sets; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.apache.cassandra.db.PartitionPosition; import org.apache.cassandra.db.TypeSizes; @@ -35,6 +41,8 @@ import org.apache.cassandra.utils.bytecomparable.ByteSourceInverse; public abstract class Token implements RingPosition, Serializable { + private static final Logger logger = LoggerFactory.getLogger(Token.class); + private static final long serialVersionUID = 1L; public static final TokenSerializer serializer = new TokenSerializer(); @@ -178,11 +186,17 @@ public abstract class Token implements RingPosition, Serializable } } + public static volatile boolean logPartitioner = false; + public static final Set> serializePartitioners = Sets.newSetFromMap(new ConcurrentHashMap<>()); + public static final Set> deserializePartitioners = Sets.newSetFromMap(new ConcurrentHashMap<>()); + public static class CompactTokenSerializer implements IPartitionerDependentSerializer { public void serialize(Token token, DataOutputPlus out, int version) throws IOException { IPartitioner p = token.getPartitioner(); + if (logPartitioner && serializePartitioners.add(p.getClass())) + logger.debug("Serializing token with partitioner " + p); if (!p.isFixedLength()) out.writeUnsignedVInt32(p.getTokenFactory().byteSize(token)); p.getTokenFactory().serialize(token, out); @@ -191,6 +205,8 @@ public abstract class Token implements RingPosition, Serializable public Token deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException { int size = p.isFixedLength() ? p.getMaxTokenSize() : in.readUnsignedVInt32(); + if (logPartitioner && deserializePartitioners.add(p.getClass())) + logger.debug("Deserializing token with partitioner " + p); byte[] bytes = new byte[size]; in.readFully(bytes); return p.getTokenFactory().fromByteArray(ByteBuffer.wrap(bytes)); diff --git a/src/java/org/apache/cassandra/index/accord/IndexMetrics.java b/src/java/org/apache/cassandra/index/accord/IndexMetrics.java index 2956082235..d992c8a16b 100644 --- a/src/java/org/apache/cassandra/index/accord/IndexMetrics.java +++ b/src/java/org/apache/cassandra/index/accord/IndexMetrics.java @@ -30,7 +30,7 @@ import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics; // Stolen from org.apache.cassandra.index.sai.metrics.AbstractMetrics public class IndexMetrics { - private static final String TYPE = "RouteIndex"; + public static final String TYPE = "RouteIndex"; private static final String SCOPE = "IndexMetrics"; private final List tracked = new ArrayList<>(); @@ -61,7 +61,7 @@ public class IndexMetrics { metricScope += '.' + indexName; } - metricScope += '.' + SCOPE + '.' + name; + metricScope += '.' + SCOPE; CassandraMetricsRegistry.MetricName metricName = new CassandraMetricsRegistry.MetricName(DefaultNameFactory.GROUP_NAME, TYPE, name, metricScope, createMBeanName(name, SCOPE)); diff --git a/src/java/org/apache/cassandra/io/IVersionedAsymmetricSerializer.java b/src/java/org/apache/cassandra/io/IVersionedAsymmetricSerializer.java index 8ad2c285c3..ff89110e33 100644 --- a/src/java/org/apache/cassandra/io/IVersionedAsymmetricSerializer.java +++ b/src/java/org/apache/cassandra/io/IVersionedAsymmetricSerializer.java @@ -32,7 +32,7 @@ public interface IVersionedAsymmetricSerializer * @param version protocol version * @throws IOException if serialization fails */ - public void serialize(In t, DataOutputPlus out, int version) throws IOException; + void serialize(In t, DataOutputPlus out, int version) throws IOException; /** * Deserialize into the specified DataInputStream instance. @@ -41,7 +41,7 @@ public interface IVersionedAsymmetricSerializer * @return the type that was deserialized * @throws IOException if deserialization fails */ - public Out deserialize(DataInputPlus in, int version) throws IOException; + Out deserialize(DataInputPlus in, int version) throws IOException; /** * Calculate serialized size of object without actually serializing. @@ -49,5 +49,5 @@ public interface IVersionedAsymmetricSerializer * @param version protocol version * @return serialized size of object t */ - public long serializedSize(In t, int version); + long serializedSize(In t, int version); } diff --git a/src/java/org/apache/cassandra/journal/Metrics.java b/src/java/org/apache/cassandra/journal/Metrics.java index befc3c2ddb..4bca57c77c 100644 --- a/src/java/org/apache/cassandra/journal/Metrics.java +++ b/src/java/org/apache/cassandra/journal/Metrics.java @@ -23,8 +23,10 @@ import org.apache.cassandra.metrics.CassandraMetricsRegistry; import org.apache.cassandra.metrics.DefaultNameFactory; import org.apache.cassandra.metrics.MetricNameFactory; -final class Metrics +public final class Metrics { + public static final String TYPE_NAME = "Journal"; + private static final String WAITING_ON_FLUSH = "WaitingOnFlush"; private static final String WAITING_ON_ALLOCATION = "WaitingOnSegmentAllocation"; private static final String WRITTEN_ENTRIES = "WrittenEntries"; @@ -49,7 +51,7 @@ final class Metrics Metrics(String name) { - this.factory = new DefaultNameFactory("Journal", name); + this.factory = new DefaultNameFactory(TYPE_NAME, name); } void register(Flusher flusher) diff --git a/src/java/org/apache/cassandra/metrics/CassandraMetricsRegistry.java b/src/java/org/apache/cassandra/metrics/CassandraMetricsRegistry.java index 8cf83f5208..b114748759 100644 --- a/src/java/org/apache/cassandra/metrics/CassandraMetricsRegistry.java +++ b/src/java/org/apache/cassandra/metrics/CassandraMetricsRegistry.java @@ -114,6 +114,8 @@ public class CassandraMetricsRegistry extends MetricRegistry // for virtual tables. metricGroups = ImmutableSet.builder() .add(AbstractMetrics.TYPE) + .add(AccordMetrics.ACCORD_COORDINATOR) + .add(AccordMetrics.ACCORD_REPLICA) .add(BatchMetrics.TYPE_NAME) .add(BufferPoolMetrics.TYPE_NAME) .add(CIDRAuthorizerMetrics.TYPE_NAME) @@ -130,8 +132,10 @@ public class CassandraMetricsRegistry extends MetricRegistry .add(DroppedMessageMetrics.TYPE) .add(HintedHandoffMetrics.TYPE_NAME) .add(HintsServiceMetrics.TYPE_NAME) + .add(org.apache.cassandra.index.accord.IndexMetrics.TYPE) .add(InternodeInboundMetrics.TYPE_NAME) .add(InternodeOutboundMetrics.TYPE_NAME) + .add(org.apache.cassandra.journal.Metrics.TYPE_NAME) .add(KeyspaceMetrics.TYPE_NAME) .add(MemtablePool.TYPE_NAME) .add(MessagingMetrics.TYPE_NAME) diff --git a/src/java/org/apache/cassandra/net/Verb.java b/src/java/org/apache/cassandra/net/Verb.java index 89b30f8818..adaa602a9a 100644 --- a/src/java/org/apache/cassandra/net/Verb.java +++ b/src/java/org/apache/cassandra/net/Verb.java @@ -346,15 +346,15 @@ public enum Verb ACCORD_SYNC_NOTIFY_REQ (151, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SIMPLE_RSP ), - ACCORD_APPLY_AND_WAIT_REQ (152, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.readData, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP), + ACCORD_APPLY_AND_WAIT_REQ (152, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.readData, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP), CONSENSUS_KEY_MIGRATION (153, P1, writeTimeout, MUTATION, () -> ConsensusKeyMigrationFinished.serializer,() -> ConsensusKeyMigrationState.consensusKeyMigrationFinishedHandler), ACCORD_INTEROP_READ_RSP (154, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.replySerializer, RESPONSE_HANDLER), - ACCORD_INTEROP_READ_REQ (155, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.requestSerializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_RSP), - ACCORD_INTEROP_COMMIT_REQ (156, P2, writeTimeout, IMMEDIATE, () -> AccordInteropCommit.serializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_RSP), + ACCORD_INTEROP_READ_REQ (155, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.requestSerializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_RSP), + ACCORD_INTEROP_COMMIT_REQ (156, P2, writeTimeout, IMMEDIATE, () -> AccordInteropCommit.serializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_RSP), ACCORD_INTEROP_READ_REPAIR_RSP (157, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.replySerializer, RESPONSE_HANDLER), - ACCORD_INTEROP_READ_REPAIR_REQ (158, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.requestSerializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_REPAIR_RSP), + ACCORD_INTEROP_READ_REPAIR_REQ (158, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.requestSerializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_REPAIR_RSP), ACCORD_INTEROP_APPLY_REQ (160, P2, writeTimeout, IMMEDIATE, () -> AccordInteropApply.serializer, AccordService::verbHandlerOrNoop, ACCORD_APPLY_RSP), // generic failure response diff --git a/src/java/org/apache/cassandra/repair/messages/SyncResponse.java b/src/java/org/apache/cassandra/repair/messages/SyncResponse.java index 0c528a3796..e7b5446bad 100644 --- a/src/java/org/apache/cassandra/repair/messages/SyncResponse.java +++ b/src/java/org/apache/cassandra/repair/messages/SyncResponse.java @@ -24,7 +24,7 @@ import java.util.Objects; import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.dht.IPartitioner; -import org.apache.cassandra.io.IVersionedSerializer; +import org.apache.cassandra.dht.IPartitionerDependentSerializer; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.locator.InetAddressAndPort; @@ -79,7 +79,7 @@ public class SyncResponse extends RepairMessage return Objects.hash(desc, success, nodes, summaries); } - public static final IVersionedSerializer serializer = new IVersionedSerializer() + public static final IPartitionerDependentSerializer serializer = new IPartitionerDependentSerializer() { public void serialize(SyncResponse message, DataOutputPlus out, int version) throws IOException { @@ -94,7 +94,8 @@ public class SyncResponse extends RepairMessage } } - public SyncResponse deserialize(DataInputPlus in, int version) throws IOException + @Override + public SyncResponse deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException { RepairJobDesc desc = RepairJobDesc.serializer.deserialize(in, version); SyncNodePair nodes = SyncNodePair.serializer.deserialize(in, version); @@ -104,7 +105,7 @@ public class SyncResponse extends RepairMessage List summaries = new ArrayList<>(numSummaries); for (int i=0; i keyCRCBytes = ThreadLocal.withInitial(() -> new byte[21]); - static final Params PARAMS = new Params() - { - @Override - public int segmentSize() - { - return 32 << 20; - } - - @Override - public FailurePolicy failurePolicy() - { - return FailurePolicy.STOP; - } - - @Override - public FlushMode flushMode() - { - return FlushMode.BATCH; - } - - @Override - public int flushPeriodMillis() - { - return DatabaseDescriptor.getCommitLogSyncPeriod(); - } - - @Override - public int periodicFlushLagBlock() - { - return 1500; - } - - @Override - public int userVersion() - { - /* - * NOTE: when accord journal version gets bumped, expose it via yaml. - * This way operators can force previous version on upgrade, temporarily, - * to allow easier downgrades if something goes wrong. - */ - return 1; - } - }; - private final File directory; private final Journal journal; private final AccordEndpointMapper endpointMapper; @@ -219,10 +176,10 @@ public class AccordJournal implements IJournal, Shutdownable private final FrameApplicator frameApplicator = new FrameApplicator(); @VisibleForTesting - public AccordJournal(AccordEndpointMapper endpointMapper) + public AccordJournal(AccordEndpointMapper endpointMapper, Params params) { this.directory = new File(DatabaseDescriptor.getAccordJournalDirectory()); - this.journal = new Journal<>("AccordJournal", directory, PARAMS, new JournalCallbacks(), Key.SUPPORT, RECORD_SERIALIZER); + this.journal = new Journal<>("AccordJournal", directory, params, new JournalCallbacks(), Key.SUPPORT, RECORD_SERIALIZER); this.endpointMapper = endpointMapper; } @@ -969,6 +926,22 @@ public class AccordJournal implements IJournal, Shutdownable } } msgTypeToSynonymousTypesMap = ImmutableListMultimap.copyOf(msgTypeToSynonymousTypes); + + //TODO (now): enable as this shows we are currently missing a message +// IllegalStateException e = null; +// for (MessageType t : MessageType.values) +// { +// if (!t.hasSideEffects()) continue; +// Type matches = msgTypeToTypeMap.get(t); +// if (matches == null) +// { +// IllegalStateException ise = new IllegalStateException("Missing MessageType " + t); +// if (e == null) e = ise; +// else e.addSuppressed(ise); +// } +// } +// if (e != null) +// throw e; } static Type fromId(int id) @@ -1164,7 +1137,7 @@ public class AccordJournal implements IJournal, Shutdownable while (null != (request = unframedRequests.poll())) { long waitForEpoch = request.waitForEpoch; - if (!node.topology().hasEpoch(waitForEpoch)) + if (waitForEpoch != 0 && !node.topology().hasEpoch(waitForEpoch)) { delayedRequests.computeIfAbsent(waitForEpoch, ignore -> new ArrayList<>()).add(request); if (!waitForEpochs.containsLong(waitForEpoch)) @@ -1394,6 +1367,12 @@ public class AccordJournal implements IJournal, Shutdownable this.txnId = txnId; } + @Override + public TxnId txnId() + { + return txnId; + } + @Override public Set test(Set messages) { @@ -1464,6 +1443,12 @@ public class AccordJournal implements IJournal, Shutdownable return readMessage(txnId, STABLE_FAST_PATH_REQ, Commit.class); } + @Override + public Commit stableSlowPath() + { + return readMessage(txnId, STABLE_SLOW_PATH_REQ, Commit.class); + } + @Override public Commit stableMaximal() { @@ -1493,6 +1478,18 @@ public class AccordJournal implements IJournal, Shutdownable { return readMessage(txnId, PROPAGATE_APPLY_MSG, Propagate.class); } + + @Override + public Propagate propagateOther() + { + return readMessage(txnId, PROPAGATE_OTHER_MSG, Propagate.class); + } + + @Override + public ApplyThenWaitUntilApplied applyThenWaitUntilApplied() + { + return readMessage(txnId, APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, ApplyThenWaitUntilApplied.class); + } } private final class LoggingMessageProvider implements SerializerSupport.MessageProvider @@ -1506,6 +1503,12 @@ public class AccordJournal implements IJournal, Shutdownable this.provider = provider; } + @Override + public TxnId txnId() + { + return txnId; + } + @Override public Set test(Set messages) { @@ -1587,6 +1590,15 @@ public class AccordJournal implements IJournal, Shutdownable return commit; } + @Override + public Commit stableSlowPath() + { + logger.debug("Fetching {} message for {}", STABLE_SLOW_PATH_REQ, txnId); + Commit commit = provider.stableSlowPath(); + logger.debug("Fetched {} message for {}: {}", STABLE_SLOW_PATH_REQ, txnId, commit); + return commit; + } + @Override public Commit stableMaximal() { @@ -1631,5 +1643,23 @@ public class AccordJournal implements IJournal, Shutdownable logger.debug("Fetched {} message for {}: {}", PROPAGATE_APPLY_MSG, txnId, propagate); return propagate; } + + @Override + public Propagate propagateOther() + { + logger.debug("Fetching {} message for {}", PROPAGATE_OTHER_MSG, txnId); + Propagate propagate = provider.propagateOther(); + logger.debug("Fetched {} message for {}: {}", PROPAGATE_OTHER_MSG, txnId, propagate); + return propagate; + } + + @Override + public ApplyThenWaitUntilApplied applyThenWaitUntilApplied() + { + logger.debug("Fetching {} message for {}", APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, txnId); + ApplyThenWaitUntilApplied apply = provider.applyThenWaitUntilApplied(); + logger.debug("Fetched {} message for {}: {}", APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, txnId, apply); + return apply; + } } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java index 7138eb9ab2..5298670e13 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java +++ b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java @@ -250,6 +250,7 @@ public class AccordKeyspace + format("execute_at %s,", TIMESTAMP_TUPLE) + format("promised_ballot %s,", TIMESTAMP_TUPLE) + format("accepted_ballot %s,", TIMESTAMP_TUPLE) + + format("execute_atleast %s,", TIMESTAMP_TUPLE) + "waiting_on blob," + "listeners set, " + "PRIMARY KEY((store_id, domain, txn_id))" @@ -298,6 +299,7 @@ public class AccordKeyspace public static final ColumnMetadata execute_at = getColumn(Commands, "execute_at"); static final ColumnMetadata promised_ballot = getColumn(Commands, "promised_ballot"); static final ColumnMetadata accepted_ballot = getColumn(Commands, "accepted_ballot"); + static final ColumnMetadata execute_atleast = getColumn(Commands, "execute_atleast"); static final ColumnMetadata waiting_on = getColumn(Commands, "waiting_on"); static final ColumnMetadata listeners = getColumn(Commands, "listeners"); @@ -857,6 +859,8 @@ public class AccordKeyspace addCellIfModified(CommandsColumns.execute_at, Command::executeAt, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command); addCellIfModified(CommandsColumns.promised_ballot, Command::promised, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command); addCellIfModified(CommandsColumns.accepted_ballot, Command::acceptedOrCommitted, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command); + if (command.txnId().kind().awaitsOnlyDeps()) + addCellIfModified(CommandsColumns.execute_atleast, Command::executesAtLeast, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command); if (command.isStable() && !command.isTruncated()) { @@ -1230,11 +1234,12 @@ public class AccordKeyspace Timestamp executeAt = deserializeExecuteAtOrNull(row); Ballot promised = deserializePromisedOrNull(row); Ballot accepted = deserializeAcceptedOrNull(row); + Timestamp executeAtLeast = status.is(Status.Truncated) && txnId.kind().awaitsOnlyDeps() ? deserializeExecuteAtLeastOrNull(row) : null; WaitingOnProvider waitingOn = deserializeWaitingOn(txnId, row); MessageProvider messages = commandStore.makeMessageProvider(txnId); - return SerializerSupport.reconstruct(commandStore.unsafeRangesForEpoch(), attrs, status, executeAt, promised, accepted, waitingOn, messages); + return SerializerSupport.reconstruct(commandStore.unsafeRangesForEpoch(), attrs, status, executeAt, executeAtLeast, promised, accepted, waitingOn, messages); } catch (Throwable t) { @@ -1307,6 +1312,11 @@ public class AccordKeyspace return deserializeTimestampOrNull(row, "execute_at", Timestamp::fromBits); } + public static Timestamp deserializeExecuteAtLeastOrNull(UntypedResultSet.Row row) + { + return deserializeTimestampOrNull(row, "execute_atleast", Timestamp::fromBits); + } + public static Ballot deserializePromisedOrNull(UntypedResultSet.Row row) { return deserializeTimestampOrNull(row.getBlob("promised_ballot"), Ballot::fromBits); diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index a248344e8d..3ee157e436 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -43,6 +43,7 @@ import accord.impl.CoordinateDurabilityScheduling; import accord.primitives.SyncPoint; import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.cql3.statements.RequestValidations; +import org.apache.cassandra.exceptions.RequestExecutionException; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.service.accord.interop.AccordInteropAdapter.AccordInteropFactory; @@ -315,7 +316,7 @@ public class AccordService implements IAccordService, Shutdownable this.scheduler = new AccordScheduler(); this.dataStore = new AccordDataStore(); this.configuration = new AccordConfiguration(DatabaseDescriptor.getRawConfig()); - this.journal = new AccordJournal(configService); + this.journal = new AccordJournal(configService, DatabaseDescriptor.getAccord().journal); this.node = new Node(localId, messageSink, this::handleLocalRequest, @@ -447,7 +448,7 @@ public class AccordService implements IAccordService, Shutdownable private long doWithRetries(LongSupplier action, int retryAttempts, long initialBackoffMillis, long maxBackoffMillis) throws InterruptedException { // Since we could end up having the barrier transaction or the transaction it listens to invalidated - CoordinationFailed existingFailures = null; + RuntimeException existingFailures = null; Long success = null; long backoffMillis = 0; for (int attempt = 0; attempt < retryAttempts; attempt++) @@ -468,7 +469,7 @@ public class AccordService implements IAccordService, Shutdownable success = action.getAsLong(); break; } - catch (CoordinationFailed newFailures) + catch (RequestExecutionException | CoordinationFailed newFailures) { existingFailures = Throwables.merge(existingFailures, newFailures); } diff --git a/src/java/org/apache/cassandra/service/accord/async/AsyncOperation.java b/src/java/org/apache/cassandra/service/accord/async/AsyncOperation.java index 51041aab58..72e673f48b 100644 --- a/src/java/org/apache/cassandra/service/accord/async/AsyncOperation.java +++ b/src/java/org/apache/cassandra/service/accord/async/AsyncOperation.java @@ -215,7 +215,7 @@ public abstract class AsyncOperation extends AsyncChains.Head implements R commandStore.abortCurrentOperation(); case LOADING: context.releaseResources(commandStore); - commandStore.executionOrder().unregister(this); + commandStore.executionOrder().unregisterOutOfOrder(this); case INITIALIZED: break; // nothing to clean up, call callback } @@ -239,6 +239,8 @@ public abstract class AsyncOperation extends AsyncChains.Head implements R default: throw new IllegalStateException("Unexpected state " + state); case INITIALIZED: canRun = commandStore.executionOrder().register(this); + if (Invariants.isParanoid()) + Invariants.checkState(canRun.booleanValue() == commandStore.executionOrder().canRun(this), "Register of %s returned canRun=%s but canRun returned %s!", this, canRun, !canRun); state(LOADING); case LOADING: if (null == canRun) diff --git a/src/java/org/apache/cassandra/service/accord/async/ExecutionOrder.java b/src/java/org/apache/cassandra/service/accord/async/ExecutionOrder.java index 03527f6553..715703acde 100644 --- a/src/java/org/apache/cassandra/service/accord/async/ExecutionOrder.java +++ b/src/java/org/apache/cassandra/service/accord/async/ExecutionOrder.java @@ -21,6 +21,7 @@ import java.util.ArrayDeque; import java.util.ArrayList; import java.util.IdentityHashMap; import java.util.List; +import java.util.function.Consumer; import accord.api.Key; import accord.api.RoutingKey; @@ -82,31 +83,9 @@ public class ExecutionOrder } } - Conflicts remove(AsyncOperation operation) + Conflicts remove(AsyncOperation operation, boolean allowOutOfOrder) { - if (operationOrQueue instanceof AsyncOperation) - { - Invariants.checkState(operationOrQueue == operation); - rangeQueues.remove(range); - } - else - { - @SuppressWarnings("unchecked") - ArrayDeque> queue = (ArrayDeque>) operationOrQueue; - AsyncOperation head = queue.poll(); - Invariants.checkState(head == operation); - - if (queue.isEmpty()) - { - rangeQueues.remove(range); - } - else - { - head = queue.peek(); - if (canRun(head)) - head.onUnblocked(); - } - } + unregister("range", range, operationOrQueue, operation, allowOutOfOrder, () -> rangeQueues.remove(range)); return operationToConflicts.remove(operation); } @@ -182,20 +161,12 @@ public class ExecutionOrder result.rangeConflicts.add(e.getKey()); } RangeState state = e.getValue(); - Object operationOrQueue = state.operationOrQueue; - if (operationOrQueue instanceof AsyncOperation) - { - ArrayDeque> queue = new ArrayDeque<>(4); - queue.add((AsyncOperation) operationOrQueue); - queue.add(operation); - state.operationOrQueue = queue; - } - else - { - @SuppressWarnings("unchecked") - ArrayDeque> queue = (ArrayDeque>) operationOrQueue; - queue.add(operation); - } + // a single range could conflict with multiple other ranges, so it is possible that the operation + // exists in the queue already due to another range in the txn... simple example is + // keys = (0, 10], (12, 15] + // e.getKey() == (-100, 100] + // in this case the operation would attempt to double add since it has 2 keys that conflict with this single range + register(state.operationOrQueue, operation, q -> state.operationOrQueue = q); }); if (result.sameRange != null) { @@ -205,7 +176,7 @@ public class ExecutionOrder { rangeQueues.add(range, new RangeState(range, keyConflicts, result.rangeConflicts, operation)); } - return keyConflicts == null && result.rangeConflicts == null; + return keyConflicts == null && result.rangeConflicts == null && result.sameRange == null; } /** @@ -221,12 +192,19 @@ public class ExecutionOrder return true; } + register(operationOrQueue, operation, q -> queues.put(keyOrTxnId, q)); + return false; + } + + private void register(Object operationOrQueue, AsyncOperation operation, Consumer>> onCreateQueue) + { if (operationOrQueue instanceof AsyncOperation) { + Invariants.checkState(operationOrQueue != operation, "Attempted to double register operation %s", operation); ArrayDeque> queue = new ArrayDeque<>(4); queue.add((AsyncOperation) operationOrQueue); queue.add(operation); - queues.put(keyOrTxnId, queue); + onCreateQueue.accept(queue); } else { @@ -234,23 +212,35 @@ public class ExecutionOrder ArrayDeque> queue = (ArrayDeque>) operationOrQueue; queue.add(operation); } - return false; + } + + /** + * Unregister the operation as being a dependency for its keys and TxnIds, but do so even if it is unable to run now. + */ + void unregisterOutOfOrder(AsyncOperation operation) + { + unregister(operation, true); } /** * Unregister the operation as being a dependency for its keys and TxnIds */ void unregister(AsyncOperation operation) + { + unregister(operation, false); + } + + private void unregister(AsyncOperation operation, boolean allowOutOfOrder) { for (Seekable seekable : operation.keys()) { switch (seekable.domain()) { case Key: - unregister(seekable.asKey(), operation); + unregister(seekable.asKey(), operation, allowOutOfOrder); break; case Range: - unregister(seekable.asRange(), operation); + unregister(seekable.asRange(), operation, allowOutOfOrder); break; default: throw new AssertionError("Unexpected domain: " + seekable.domain()); @@ -259,48 +249,69 @@ public class ExecutionOrder } TxnId primaryTxnId = operation.primaryTxnId(); if (null != primaryTxnId) - unregister(primaryTxnId, operation); + unregister(primaryTxnId, operation, allowOutOfOrder); } - private void unregister(Range range, AsyncOperation operation) + private void unregister(Range range, AsyncOperation operation, boolean allowOutOfOrder) { RangeState state = state(range); - Conflicts conflicts = state.remove(operation); + Conflicts conflicts = state.remove(operation, allowOutOfOrder); if (conflicts.rangeConflicts != null) - conflicts.rangeConflicts.forEach(r -> state(r).remove(operation)); + conflicts.rangeConflicts.forEach(r -> state(r).remove(operation, allowOutOfOrder)); if (conflicts.keyConflicts != null) - conflicts.keyConflicts.forEach(k -> unregister(k, operation)); + conflicts.keyConflicts.forEach(k -> unregister(k, operation, allowOutOfOrder)); } /** * Unregister the operation as being a dependency for key or TxnId */ - private void unregister(Object keyOrTxnId, AsyncOperation operation) + private void unregister(Object keyOrTxnId, AsyncOperation operation, boolean allowOutOfOrder) { Object operationOrQueue = queues.get(keyOrTxnId); Invariants.nonNull(operationOrQueue); + unregister("Key or TxnId", keyOrTxnId, operationOrQueue, operation, allowOutOfOrder, () -> queues.remove(keyOrTxnId)); + } + + private void unregister(String name, Object key, Object operationOrQueue, AsyncOperation operation, boolean allowOutOfOrder, Runnable onEmpty) + { if (operationOrQueue instanceof AsyncOperation) { - Invariants.checkState(operationOrQueue == operation); - queues.remove(keyOrTxnId); + Invariants.checkState(operationOrQueue == operation, "Only single operation present and was not %s; %s %s", name, key); + onEmpty.run(); } else { @SuppressWarnings("unchecked") ArrayDeque> queue = (ArrayDeque>) operationOrQueue; - AsyncOperation head = queue.poll(); - Invariants.checkState(head == operation); - - if (queue.isEmpty()) + if (allowOutOfOrder) { - queues.remove(keyOrTxnId); + Invariants.checkState(queue.remove(operation), "Operation %s was not found in queue: %s; %s %s", operation, queue, name, key); } else { - head = queue.peek(); - if (canRun(head)) - head.onUnblocked(); + Invariants.checkState(queue.peek() == operation, "Operation %s is not at the top of the queue; %s; %s %s", operation, queue, name, key); + queue.poll(); + } + + if (queue.isEmpty()) + { + onEmpty.run(); + } + else + { + AsyncOperation next = queue.peek(); + if (next == operation) + { + // a single range could conflict with multiple other ranges, so it is possible that the operation + // exists in the queue already due to another range in the txn... simple example is + // keys = (0, 10], (12, 15] + // e.getKey() == (-100, 100] + // in this case the operation would attempt to double add since it has 2 keys that conflict with this single range + return; + } + if (canRun(next)) + next.onUnblocked(); } } } @@ -357,7 +368,7 @@ public class ExecutionOrder private RangeState state(Range range) { List list = rangeQueues.get(range); - assert list.size() == 1 : String.format("Expected 1 element but saw list %s", list); + assert list.size() == 1 : String.format("Expected 1 element for range %s but saw list %s", range, list); return list.get(0); } diff --git a/src/java/org/apache/cassandra/service/consensus/migration/ConsensusRequestRouter.java b/src/java/org/apache/cassandra/service/consensus/migration/ConsensusRequestRouter.java index 1188076c90..70c2c39660 100644 --- a/src/java/org/apache/cassandra/service/consensus/migration/ConsensusRequestRouter.java +++ b/src/java/org/apache/cassandra/service/consensus/migration/ConsensusRequestRouter.java @@ -121,7 +121,7 @@ public class ConsensusRequestRouter ClusterMetadata cm = ClusterMetadata.current(); TableMetadata metadata = cm.schema.getTableMetadata(tableId); if (metadata == null) - throw new IllegalStateException("Can't route consensus request for nonexistent table %s".format(tableId.toString())); + throw new IllegalStateException(String.format("Can't route consensus request for nonexistent table %s", tableId)); if (!mayWriteThroughAccord(metadata)) return false; diff --git a/src/java/org/apache/cassandra/streaming/SessionSummary.java b/src/java/org/apache/cassandra/streaming/SessionSummary.java index f5bcfa31be..8bb1a1eb81 100644 --- a/src/java/org/apache/cassandra/streaming/SessionSummary.java +++ b/src/java/org/apache/cassandra/streaming/SessionSummary.java @@ -108,14 +108,14 @@ public class SessionSummary List receivingSummaries = new ArrayList<>(numRcvd); for (int i=0; i sendingSummaries = new ArrayList<>(numRcvd); for (int i=0; i serializer = new StreamSummarySerializer(); + public static final IVersionedSerializer serializer = new StreamSummarySerializer(); public final TableId tableId; public final List> ranges; @@ -86,25 +87,34 @@ public class StreamSummary implements Serializable return sb.toString(); } - public static class StreamSummarySerializer implements IPartitionerDependentSerializer + public static class StreamSummarySerializer implements IVersionedSerializer { public void serialize(StreamSummary summary, DataOutputPlus out, int version) throws IOException { summary.tableId.serialize(out); out.writeInt(summary.files); out.writeLong(summary.totalSize); + Token.logPartitioner = true; if (version >= MessagingService.VERSION_51) CollectionSerializers.serializeCollection(summary.ranges, out, version, Range.rangeSerializer); + Token.logPartitioner = false; } - public StreamSummary deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException + public StreamSummary deserialize(DataInputPlus in, int version) throws IOException { TableId tableId = TableId.deserialize(in); + int files = in.readInt(); long totalSize = in.readLong(); List> ranges = ImmutableList.of(); if (version >= MessagingService.VERSION_51) + { + TableMetadata tableMetadata = Schema.instance.getTableMetadata(tableId); + IPartitioner p = tableMetadata != null ? tableMetadata.partitioner : IPartitioner.global(); + Token.logPartitioner = true; ranges = CollectionSerializers.deserializeList(in, p, version, Range.rangeSerializer); + Token.logPartitioner = false; + } return new StreamSummary(tableId, ranges, files, totalSize); } diff --git a/src/java/org/apache/cassandra/streaming/messages/CompleteMessage.java b/src/java/org/apache/cassandra/streaming/messages/CompleteMessage.java index 86620c3859..afb1c6c7b4 100644 --- a/src/java/org/apache/cassandra/streaming/messages/CompleteMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/CompleteMessage.java @@ -17,7 +17,6 @@ */ package org.apache.cassandra.streaming.messages; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamingDataOutputPlus; @@ -26,7 +25,7 @@ public class CompleteMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public CompleteMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) + public CompleteMessage deserialize(DataInputPlus in, int version) { return new CompleteMessage(); } diff --git a/src/java/org/apache/cassandra/streaming/messages/IncomingStreamMessage.java b/src/java/org/apache/cassandra/streaming/messages/IncomingStreamMessage.java index 4ee726ee83..e48d115e35 100644 --- a/src/java/org/apache/cassandra/streaming/messages/IncomingStreamMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/IncomingStreamMessage.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.Objects; import org.apache.cassandra.db.ColumnFamilyStore; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.IncomingStream; import org.apache.cassandra.streaming.StreamManager; @@ -34,7 +33,7 @@ public class IncomingStreamMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public IncomingStreamMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException + public IncomingStreamMessage deserialize(DataInputPlus input, int version) throws IOException { StreamMessageHeader header = StreamMessageHeader.serializer.deserialize(input, version); StreamSession session = StreamManager.instance.findSession(header.sender, header.planId, header.sessionIndex, header.sendByFollower); diff --git a/src/java/org/apache/cassandra/streaming/messages/KeepAliveMessage.java b/src/java/org/apache/cassandra/streaming/messages/KeepAliveMessage.java index 928783f401..42be1e99a1 100644 --- a/src/java/org/apache/cassandra/streaming/messages/KeepAliveMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/KeepAliveMessage.java @@ -17,7 +17,6 @@ */ package org.apache.cassandra.streaming.messages; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamingDataOutputPlus; @@ -38,7 +37,7 @@ public class KeepAliveMessage extends StreamMessage public static Serializer serializer = new Serializer() { - public KeepAliveMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) + public KeepAliveMessage deserialize(DataInputPlus in, int version) { return new KeepAliveMessage(); } diff --git a/src/java/org/apache/cassandra/streaming/messages/OutgoingStreamMessage.java b/src/java/org/apache/cassandra/streaming/messages/OutgoingStreamMessage.java index dcd3b755e8..b83d7863fc 100644 --- a/src/java/org/apache/cassandra/streaming/messages/OutgoingStreamMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/OutgoingStreamMessage.java @@ -21,7 +21,6 @@ import java.io.IOException; import com.google.common.annotations.VisibleForTesting; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.streaming.OutgoingStream; @@ -33,7 +32,7 @@ public class OutgoingStreamMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public OutgoingStreamMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) + public OutgoingStreamMessage deserialize(DataInputPlus in, int version) { throw new UnsupportedOperationException("Not allowed to call deserialize on an outgoing stream"); } diff --git a/src/java/org/apache/cassandra/streaming/messages/PrepareAckMessage.java b/src/java/org/apache/cassandra/streaming/messages/PrepareAckMessage.java index 72d61d29cb..f93b5afe30 100644 --- a/src/java/org/apache/cassandra/streaming/messages/PrepareAckMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/PrepareAckMessage.java @@ -20,7 +20,6 @@ package org.apache.cassandra.streaming.messages; import java.io.IOException; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamingDataOutputPlus; @@ -34,7 +33,7 @@ public class PrepareAckMessage extends StreamMessage //nop } - public PrepareAckMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException + public PrepareAckMessage deserialize(DataInputPlus in, int version) throws IOException { return new PrepareAckMessage(); } diff --git a/src/java/org/apache/cassandra/streaming/messages/PrepareSynAckMessage.java b/src/java/org/apache/cassandra/streaming/messages/PrepareSynAckMessage.java index e29e651824..e052f4c301 100644 --- a/src/java/org/apache/cassandra/streaming/messages/PrepareSynAckMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/PrepareSynAckMessage.java @@ -22,7 +22,6 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collection; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamSummary; @@ -39,12 +38,12 @@ public class PrepareSynAckMessage extends StreamMessage StreamSummary.serializer.serialize(summary, out, version); } - public PrepareSynAckMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException + public PrepareSynAckMessage deserialize(DataInputPlus input, int version) throws IOException { PrepareSynAckMessage message = new PrepareSynAckMessage(); int numSummaries = input.readInt(); for (int i = 0; i < numSummaries; i++) - message.summaries.add(StreamSummary.serializer.deserialize(input, partitioner, version)); + message.summaries.add(StreamSummary.serializer.deserialize(input, version)); return message; } diff --git a/src/java/org/apache/cassandra/streaming/messages/PrepareSynMessage.java b/src/java/org/apache/cassandra/streaming/messages/PrepareSynMessage.java index e901365e5e..c856f46983 100644 --- a/src/java/org/apache/cassandra/streaming/messages/PrepareSynMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/PrepareSynMessage.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collection; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamRequest; import org.apache.cassandra.streaming.StreamSession; @@ -32,7 +31,7 @@ public class PrepareSynMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public PrepareSynMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException + public PrepareSynMessage deserialize(DataInputPlus input, int version) throws IOException { PrepareSynMessage message = new PrepareSynMessage(); // requests @@ -42,7 +41,7 @@ public class PrepareSynMessage extends StreamMessage // summaries int numSummaries = input.readInt(); for (int i = 0; i < numSummaries; i++) - message.summaries.add(StreamSummary.serializer.deserialize(input, partitioner, version)); + message.summaries.add(StreamSummary.serializer.deserialize(input, version)); return message; } diff --git a/src/java/org/apache/cassandra/streaming/messages/ReceivedMessage.java b/src/java/org/apache/cassandra/streaming/messages/ReceivedMessage.java index c6b7a0f638..6755959614 100644 --- a/src/java/org/apache/cassandra/streaming/messages/ReceivedMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/ReceivedMessage.java @@ -19,7 +19,6 @@ package org.apache.cassandra.streaming.messages; import java.io.IOException; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.streaming.StreamSession; @@ -29,7 +28,7 @@ public class ReceivedMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public ReceivedMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException + public ReceivedMessage deserialize(DataInputPlus input, int version) throws IOException { return new ReceivedMessage(TableId.deserialize(input), input.readInt()); } diff --git a/src/java/org/apache/cassandra/streaming/messages/SessionFailedMessage.java b/src/java/org/apache/cassandra/streaming/messages/SessionFailedMessage.java index f05be58aa6..7fa82d8f67 100644 --- a/src/java/org/apache/cassandra/streaming/messages/SessionFailedMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/SessionFailedMessage.java @@ -17,7 +17,6 @@ */ package org.apache.cassandra.streaming.messages; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamingDataOutputPlus; @@ -26,7 +25,7 @@ public class SessionFailedMessage extends StreamMessage { public static Serializer serializer = new Serializer() { - public SessionFailedMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) + public SessionFailedMessage deserialize(DataInputPlus in, int version) { return new SessionFailedMessage(); } diff --git a/src/java/org/apache/cassandra/streaming/messages/StreamInitMessage.java b/src/java/org/apache/cassandra/streaming/messages/StreamInitMessage.java index 2fd65d7dff..e78442334b 100644 --- a/src/java/org/apache/cassandra/streaming/messages/StreamInitMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/StreamInitMessage.java @@ -20,7 +20,6 @@ package org.apache.cassandra.streaming.messages; import java.io.IOException; import org.apache.cassandra.db.TypeSizes; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.streaming.PreviewKind; @@ -94,7 +93,7 @@ public class StreamInitMessage extends StreamMessage out.writeInt(message.previewKind.getSerializationVal()); } - public StreamInitMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException + public StreamInitMessage deserialize(DataInputPlus in, int version) throws IOException { InetAddressAndPort from = inetAddressAndPortSerializer.deserialize(in, version); int sessionIndex = in.readInt(); diff --git a/src/java/org/apache/cassandra/streaming/messages/StreamMessage.java b/src/java/org/apache/cassandra/streaming/messages/StreamMessage.java index 6e5dc08f88..186ac3274a 100644 --- a/src/java/org/apache/cassandra/streaming/messages/StreamMessage.java +++ b/src/java/org/apache/cassandra/streaming/messages/StreamMessage.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.HashMap; import java.util.Map; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.streaming.StreamSession; import org.apache.cassandra.streaming.StreamingChannel; @@ -45,16 +44,16 @@ public abstract class StreamMessage return 1 + message.type.outSerializer.serializedSize(message, version); } - public static StreamMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException + public static StreamMessage deserialize(DataInputPlus in, int version) throws IOException { Type type = Type.lookupById(in.readByte()); - return type.inSerializer.deserialize(in, partitioner, version); + return type.inSerializer.deserialize(in, version); } /** StreamMessage serializer */ public static interface Serializer { - V deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException; + V deserialize(DataInputPlus in, int version) throws IOException; void serialize(V message, StreamingDataOutputPlus out, int version, StreamSession session) throws IOException; long serializedSize(V message, int version) throws IOException; } diff --git a/test/distributed/org/apache/cassandra/distributed/test/tcm/RepairMetadataKeyspaceTest.java b/test/distributed/org/apache/cassandra/distributed/test/tcm/RepairMetadataKeyspaceTest.java index 074a64913f..70d1acfddd 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/tcm/RepairMetadataKeyspaceTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/tcm/RepairMetadataKeyspaceTest.java @@ -28,6 +28,7 @@ import java.util.stream.Collectors; import org.junit.Test; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.IInvokableInstance; import org.apache.cassandra.distributed.shared.ClusterUtils; @@ -63,6 +64,7 @@ public class RepairMetadataKeyspaceTest extends TestBaseImpl IInvokableInstance toRepair = cluster.get(3); stopUnchecked(toRepair); + DatabaseDescriptor.clientInitialization(); String targetDir = DistributedMetadataLogKeyspace.TABLE_NAME + '-' + DistributedMetadataLogKeyspace.LOG_TABLE_ID.toHexString(); for (File datadir : getDataDirectories(toRepair)) { diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java index cd98c71f18..22dd55bb31 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java @@ -26,6 +26,7 @@ import javax.annotation.Nullable; import com.google.common.collect.ImmutableMap; import accord.topology.TopologyUtils; +import org.apache.cassandra.config.AccordSpec; import org.apache.cassandra.schema.*; import org.junit.Ignore; import org.junit.Test; @@ -166,7 +167,7 @@ public class AccordJournalSimulationTest extends SimulationTestBase } } private static final ExecutorPlus executor = ExecutorFactory.Global.executorFactory().pooled("name", 10); - private static final AccordJournal journal = new AccordJournal(null); + private static final AccordJournal journal = new AccordJournal(null, new AccordSpec.JournalSpec()); private static final int events = 100; private static final CountDownLatch eventsWritten = CountDownLatch.newCountDownLatch(events); private static final CountDownLatch eventsDurable = CountDownLatch.newCountDownLatch(events); diff --git a/test/unit/org/apache/cassandra/config/DatabaseDescriptorRefTest.java b/test/unit/org/apache/cassandra/config/DatabaseDescriptorRefTest.java index c580207634..697214f76a 100644 --- a/test/unit/org/apache/cassandra/config/DatabaseDescriptorRefTest.java +++ b/test/unit/org/apache/cassandra/config/DatabaseDescriptorRefTest.java @@ -78,6 +78,7 @@ public class DatabaseDescriptorRefTest "org.apache.cassandra.auth.INetworkAuthorizer", "org.apache.cassandra.auth.IRoleManager", "org.apache.cassandra.config.AccordSpec", + "org.apache.cassandra.config.AccordSpec$JournalSpec", "org.apache.cassandra.config.AccordSpec$TransactionalRangeMigration", "org.apache.cassandra.config.CassandraRelevantProperties", "org.apache.cassandra.config.CassandraRelevantProperties$PropertyConverter", @@ -275,6 +276,9 @@ public class DatabaseDescriptorRefTest "org.apache.cassandra.io.util.PathUtils$IOToLongFunction", "org.apache.cassandra.io.util.RebufferingInputStream", "org.apache.cassandra.io.util.SpinningDiskOptimizationStrategy", + "org.apache.cassandra.journal.Params", + "org.apache.cassandra.journal.Params$FailurePolicy", + "org.apache.cassandra.journal.Params$FlushMode", "org.apache.cassandra.locator.Endpoint", "org.apache.cassandra.locator.IEndpointSnitch", "org.apache.cassandra.locator.InetAddressAndPort", diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java index ee38ffceb9..5cbf0915b5 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java @@ -91,6 +91,7 @@ import org.apache.cassandra.service.accord.api.PartitionKey; import org.apache.cassandra.service.accord.serializers.CommandsForKeySerializer; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; +import org.assertj.core.api.Assertions; import static accord.impl.TimestampsForKey.NO_LAST_EXECUTED_HLC; import static accord.local.KeyHistory.COMMANDS; @@ -369,9 +370,14 @@ public class CompactionAccordIteratorsTest assertEquals(1, Iterators.size(partition.unfilteredIterator())); ByteBuffer[] partitionKeyComponents = CommandRows.splitPartitionKey(partition.partitionKey()); Row row = (Row)partition.unfilteredIterator().next(); - assertEquals(commands.metadata().regularColumns().size(), row.columnCount()); + + // execute_atleast is null, so when we read from the scanner the column won't be present in the partition + Assertions.assertThat(new ArrayList<>(row.columns())).isEqualTo(commands.metadata().regularColumns().stream().filter(c -> !c.name.toString().equals("execute_atleast")).collect(Collectors.toList())); for (ColumnMetadata cm : commands.metadata().regularColumns()) + { + if (cm.name.toString().equals("execute_atleast")) continue; assertNotNull(row.getColumnData(cm)); + } assertEquals(TXN_ID, CommandRows.getTxnId(partitionKeyComponents)); assertEquals(SaveStatus.Applied, AccordKeyspace.CommandRows.getStatus(row)); }; diff --git a/test/unit/org/apache/cassandra/repair/messages/RepairMessageSerializationsTest.java b/test/unit/org/apache/cassandra/repair/messages/RepairMessageSerializationsTest.java index 9e8080f903..750c614455 100644 --- a/test/unit/org/apache/cassandra/repair/messages/RepairMessageSerializationsTest.java +++ b/test/unit/org/apache/cassandra/repair/messages/RepairMessageSerializationsTest.java @@ -25,6 +25,8 @@ import java.util.List; import java.util.UUID; import com.google.common.collect.Lists; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.dht.*; import org.junit.Assert; import org.junit.BeforeClass; import org.junit.Test; @@ -32,10 +34,7 @@ import org.junit.Test; import org.apache.cassandra.CassandraTestBase; import org.apache.cassandra.CassandraTestBase.DDDaemonInitialization; import org.apache.cassandra.CassandraTestBase.UseMurmur3Partitioner; -import org.apache.cassandra.dht.Murmur3Partitioner; import org.apache.cassandra.dht.Murmur3Partitioner.LongToken; -import org.apache.cassandra.dht.Range; -import org.apache.cassandra.dht.Token; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputBuffer; @@ -116,7 +115,19 @@ public class RepairMessageSerializationsTest extends CassandraTestBase buf.flip(); DataInputPlus in = new DataInputBuffer(buf, false); - T deserialized = serializer.deserialize(in, PROTOCOL_VERSION); + + T deserialized = null; + + if (serializer instanceof IPartitionerDependentSerializer) + { + IPartitionerDependentSerializer pds = (IPartitionerDependentSerializer) serializer; + deserialized = pds.deserialize(in, DatabaseDescriptor.getPartitioner(), PROTOCOL_VERSION); + } + else + { + deserialized = serializer.deserialize(in, PROTOCOL_VERSION); + } + Assert.assertEquals(msg, deserialized); Assert.assertEquals(msg.hashCode(), deserialized.hashCode()); return deserialized; diff --git a/test/unit/org/apache/cassandra/service/SerializationsTest.java b/test/unit/org/apache/cassandra/service/SerializationsTest.java index 9251b590bb..00864f5a34 100644 --- a/test/unit/org/apache/cassandra/service/SerializationsTest.java +++ b/test/unit/org/apache/cassandra/service/SerializationsTest.java @@ -23,7 +23,6 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; -import java.util.UUID; import com.google.common.collect.Lists; import org.junit.AfterClass; @@ -60,6 +59,7 @@ import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.streaming.PreviewKind; import org.apache.cassandra.streaming.SessionSummary; import org.apache.cassandra.streaming.StreamSummary; +import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.utils.Clock; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.MerkleTrees; @@ -70,6 +70,7 @@ import static java.util.Collections.emptyList; public class SerializationsTest extends AbstractSerializationsTester { private static PartitionerSwitcher partitionerSwitcher; + private static TableId TABLE_ID; private static TimeUUID RANDOM_UUID; private static Range FULL_RANGE; private static RepairJobDesc DESC; @@ -84,6 +85,7 @@ public class SerializationsTest extends AbstractSerializationsTester ClusterMetadataTestHelper.setInstanceForTest(); SchemaTestUtil.addOrUpdateKeyspace(KeyspaceMetadata.create("Keyspace1", KeyspaceParams.simple(3))); SchemaTestUtil.announceNewTable(TableMetadata.minimal("Keyspace1", "Standard1")); + TABLE_ID = ClusterMetadata.current().schema.getKeyspaceMetadata("Keyspace1").getTableOrViewNullable("Standard1").id(); RANDOM_UUID = TimeUUID.fromString("743325d0-4c4b-11ec-8a88-2d67081686db"); FULL_RANGE = new Range<>(Util.testPartitioner().getMinimumToken(), Util.testPartitioner().getMinimumToken()); DESC = new RepairJobDesc(RANDOM_UUID, RANDOM_UUID, "Keyspace1", "Standard1", Arrays.asList(FULL_RANGE)); @@ -223,8 +225,8 @@ public class SerializationsTest extends AbstractSerializationsTester // sync success List summaries = new ArrayList<>(); summaries.add(new SessionSummary(src, dest, - Lists.newArrayList(new StreamSummary(TableId.fromUUID(UUID.randomUUID()), emptyList(), 5, 100)), - Lists.newArrayList(new StreamSummary(TableId.fromUUID(UUID.randomUUID()), emptyList(), 500, 10)) + Lists.newArrayList(new StreamSummary(TABLE_ID, emptyList(), 5, 100)), + Lists.newArrayList(new StreamSummary(TABLE_ID, emptyList(), 500, 10)) )); SyncResponse success = new SyncResponse(DESC, src, dest, true, summaries); // sync fail diff --git a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java index a608f41f8d..0d3b3d005c 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java @@ -76,6 +76,7 @@ import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.ImmediateExecutor; import org.apache.cassandra.concurrent.ManualExecutor; import org.apache.cassandra.concurrent.Stage; +import org.apache.cassandra.config.AccordSpec; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.cql3.statements.TransactionStatement; @@ -399,7 +400,7 @@ public class AccordTestUtils public long unix(TimeUnit timeUnit) { return NodeTimeService.unixWrapper(TimeUnit.MICROSECONDS, this::now).applyAsLong(timeUnit); } }; - AccordJournal journal = new AccordJournal(null); + AccordJournal journal = new AccordJournal(null, new AccordSpec.JournalSpec()); journal.start(null); SingleEpochRanges holder = new SingleEpochRanges(topology.rangesForNode(node)); diff --git a/test/unit/org/apache/cassandra/service/accord/MockJournal.java b/test/unit/org/apache/cassandra/service/accord/MockJournal.java index 8a68163ede..dc22f540ed 100644 --- a/test/unit/org/apache/cassandra/service/accord/MockJournal.java +++ b/test/unit/org/apache/cassandra/service/accord/MockJournal.java @@ -28,6 +28,7 @@ import com.google.common.collect.Sets; import accord.local.SerializerSupport; import accord.messages.Accept; import accord.messages.Apply; +import accord.messages.ApplyThenWaitUntilApplied; import accord.messages.BeginRecovery; import accord.messages.Commit; import accord.messages.Message; @@ -43,15 +44,18 @@ import org.apache.cassandra.service.accord.AccordJournal.Type; import static accord.messages.MessageType.ACCEPT_REQ; import static accord.messages.MessageType.APPLY_MAXIMAL_REQ; import static accord.messages.MessageType.APPLY_MINIMAL_REQ; +import static accord.messages.MessageType.APPLY_THEN_WAIT_UNTIL_APPLIED_REQ; import static accord.messages.MessageType.BEGIN_RECOVER_REQ; import static accord.messages.MessageType.COMMIT_MAXIMAL_REQ; import static accord.messages.MessageType.COMMIT_SLOW_PATH_REQ; import static accord.messages.MessageType.PRE_ACCEPT_REQ; import static accord.messages.MessageType.PROPAGATE_APPLY_MSG; +import static accord.messages.MessageType.PROPAGATE_OTHER_MSG; import static accord.messages.MessageType.PROPAGATE_PRE_ACCEPT_MSG; import static accord.messages.MessageType.PROPAGATE_STABLE_MSG; import static accord.messages.MessageType.STABLE_FAST_PATH_REQ; import static accord.messages.MessageType.STABLE_MAXIMAL_REQ; +import static accord.messages.MessageType.STABLE_SLOW_PATH_REQ; public class MockJournal implements IJournal { @@ -61,6 +65,12 @@ public class MockJournal implements IJournal { return new SerializerSupport.MessageProvider() { + @Override + public TxnId txnId() + { + return txnId; + } + @Override public Set test(Set messages) { @@ -146,6 +156,12 @@ public class MockJournal implements IJournal return get(STABLE_FAST_PATH_REQ); } + @Override + public Commit stableSlowPath() + { + return get(STABLE_SLOW_PATH_REQ); + } + @Override public Commit stableMaximal() { @@ -175,6 +191,18 @@ public class MockJournal implements IJournal { return get(PROPAGATE_APPLY_MSG); } + + @Override + public Propagate propagateOther() + { + return get(PROPAGATE_OTHER_MSG); + } + + @Override + public ApplyThenWaitUntilApplied applyThenWaitUntilApplied() + { + return get(APPLY_THEN_WAIT_UNTIL_APPLIED_REQ); + } }; } diff --git a/test/unit/org/apache/cassandra/service/accord/SimulatedAccordCommandStore.java b/test/unit/org/apache/cassandra/service/accord/SimulatedAccordCommandStore.java index 6cb71b68db..128843258e 100644 --- a/test/unit/org/apache/cassandra/service/accord/SimulatedAccordCommandStore.java +++ b/test/unit/org/apache/cassandra/service/accord/SimulatedAccordCommandStore.java @@ -25,6 +25,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.function.BooleanSupplier; import java.util.function.Function; +import java.util.function.Predicate; import java.util.function.ToLongFunction; import accord.impl.SizeOfIntersectionSorter; @@ -76,7 +77,7 @@ import static org.apache.cassandra.db.ColumnFamilyStore.FlushReason.UNIT_TESTS; import static org.apache.cassandra.schema.SchemaConstants.ACCORD_KEYSPACE_NAME; import static org.apache.cassandra.utils.AccordGenerators.fromQT; -class SimulatedAccordCommandStore implements AutoCloseable +public class SimulatedAccordCommandStore implements AutoCloseable { private final List failures = new ArrayList<>(); private final SimulatedExecutorFactory globalExecutor; @@ -90,8 +91,9 @@ class SimulatedAccordCommandStore implements AutoCloseable public final MockJournal journal; public final ScheduledExecutorPlus unorderedScheduled; public final List evictions = new ArrayList<>(); + public Predicate ignoreExceptions = ignore -> false; - SimulatedAccordCommandStore(RandomSource rs) + public SimulatedAccordCommandStore(RandomSource rs) { globalExecutor = new SimulatedExecutorFactory(accord.utilsfork.RandomSource.wrap(rs).fork(), fromQT(Generators.TIMESTAMP_GEN.map(java.sql.Timestamp::getTime)).mapToLong(TimeUnit.MILLISECONDS::toNanos).next(rs), failures::add); this.unorderedScheduled = globalExecutor.scheduled("ignored"); @@ -151,6 +153,13 @@ class SimulatedAccordCommandStore implements AutoCloseable { return false; } + + @Override + public void onUncaughtException(Throwable t) + { + if (ignoreExceptions.test(t)) return; + super.onUncaughtException(t); + } }, null, ignore -> AccordTestUtils.NOOP_PROGRESS_LOG, diff --git a/test/unit/org/apache/cassandra/service/accord/async/SimulatedAsyncOperationTest.java b/test/unit/org/apache/cassandra/service/accord/async/SimulatedAsyncOperationTest.java new file mode 100644 index 0000000000..6e216ff56d --- /dev/null +++ b/test/unit/org/apache/cassandra/service/accord/async/SimulatedAsyncOperationTest.java @@ -0,0 +1,207 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.service.accord.async; + +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.BiConsumer; +import java.util.function.BooleanSupplier; + +import org.junit.Before; +import org.junit.Test; + +import accord.api.Key; +import accord.impl.basic.SimulatedFault; +import accord.local.PreLoadContext; +import accord.local.SafeCommandStore; +import accord.primitives.Keys; +import accord.primitives.Range; +import accord.primitives.Ranges; +import accord.primitives.Seekables; +import accord.utils.Gen; +import accord.utils.Gens; +import accord.utils.RandomSource; +import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.dht.Murmur3Partitioner.LongToken; +import org.apache.cassandra.schema.TableId; +import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.service.accord.AccordCommandStore; +import org.apache.cassandra.service.accord.AccordKeyspace; +import org.apache.cassandra.service.accord.SimulatedAccordCommandStore; +import org.apache.cassandra.service.accord.SimulatedAccordCommandStoreTestBase; +import org.apache.cassandra.service.accord.TokenRange; +import org.apache.cassandra.service.accord.api.AccordRoutingKey.TokenKey; +import org.apache.cassandra.service.accord.api.PartitionKey; +import org.assertj.core.api.Assertions; + +import static accord.utils.Property.qt; + +public class SimulatedAsyncOperationTest extends SimulatedAccordCommandStoreTestBase +{ + @Before + public void precondition() + { + Assertions.assertThat(intTbl.partitioner).isEqualTo(Murmur3Partitioner.instance); + } + + @Test + public void happyPath() + { + qt().withExamples(100).check(rs -> test(rs, 100, intTbl, ignore -> Action.SUCCESS)); + } + + @Test + public void fuzz() + { + Gen actionGen = Gens.enums().allWithWeights(Action.class, 10, 1, 1); + qt().withExamples(100).check(rs -> test(rs, 100, intTbl, actionGen)); + } + + private static void test(RandomSource rs, int numSamples, TableMetadata tbl, Gen actionGen) throws Exception + { + AccordKeyspace.unsafeClear(); + + int numKeys = rs.nextInt(20, 1000); + long minToken = 0; + long maxToken = numKeys; + + Gen keyGen = Gens.longs().between(minToken + 1, maxToken).map(t -> new PartitionKey(tbl.id, tbl.partitioner.decorateKey(LongToken.keyForToken(t)))); + + + Gen keysGen = Gens.lists(keyGen).unique().ofSizeBetween(1, 10).map(l -> Keys.of(l)); + Gen rangesGen = Gens.lists(rangeInsideRange(tbl.id, minToken, maxToken)).uniqueBestEffort().ofSizeBetween(1, 10).map(l -> Ranges.of(l.toArray(Range[]::new))); + Gen> seekablesGen = Gens.oneOf(keysGen, rangesGen); + + try (var instance = new SimulatedAccordCommandStore(rs)) + { + instance.ignoreExceptions = t -> t instanceof SimulatedFault; + Counter counter = new Counter(); + for (int i = 0; i < numSamples; i++) + { + PreLoadContext ctx = PreLoadContext.contextFor(seekablesGen.next(rs)); + operation(instance, ctx, actionGen.next(rs), rs::nextBoolean).begin((ignore, failure) -> { + counter.counter++; + if (failure != null && !(failure instanceof SimulatedFault)) throw new AssertionError("Unexpected error", failure); + }); + } + instance.processAll(); + Assertions.assertThat(counter.counter).isEqualTo(numSamples); + } + } + + private static Gen rangeInsideRange(TableId tableId, long minToken, long maxToken) + { + if (minToken + 1 == maxToken) + { + // only one range is possible... + return Gens.constant(range(tableId, minToken, maxToken)); + } + return rs -> { + long a = rs.nextLong(minToken, maxToken + 1); + long b = rs.nextLong(minToken, maxToken + 1); + while (a == b) + b = rs.nextLong(minToken, maxToken + 1); + if (a > b) + { + long tmp = a; + a = b; + b = tmp; + } + return range(tableId, a, b); + }; + } + + private static TokenRange range(TableId tableId, long start, long end) + { + return new TokenRange(new TokenKey(tableId, new LongToken(start)), new TokenKey(tableId, new LongToken(end))); + } + + private enum Action {SUCCESS, FAILURE, LOAD_FAILURE} + + private static AsyncOperation operation(SimulatedAccordCommandStore instance, PreLoadContext ctx, Action action, BooleanSupplier delay) + { + return new SimulatedOperation(instance.store, ctx, action == Action.FAILURE ? SimulatedOperation.Action.FAILURE : SimulatedOperation.Action.SUCCESS) + { + @Override + AsyncLoader createAsyncLoader(AccordCommandStore commandStore, PreLoadContext preLoadContext) + { + return new SimulatedLoader(action == SimulatedAsyncOperationTest.Action.LOAD_FAILURE ? SimulatedLoader.Action.FAILURE : SimulatedLoader.Action.SUCCESS, delay.getAsBoolean(), instance.unorderedScheduled); + } + }; + } + + private static class Counter + { + int counter = 0; + } + + private static class SimulatedOperation extends AsyncOperation + { + enum Action { SUCCESS, FAILURE} + private final Action action; + + public SimulatedOperation(AccordCommandStore commandStore, PreLoadContext preLoadContext, Action action) + { + super(commandStore, preLoadContext); + this.action = action; + } + + @Override + public Void apply(SafeCommandStore safe) + { + if (action == Action.FAILURE) + throw new SimulatedFault("Operation failed for keys " + keys()); + return null; + } + } + + private static class SimulatedLoader extends AsyncLoader + { + + enum Action { SUCCESS, FAILURE} + + private final Action action; + private boolean delay; + private final ScheduledExecutorService executor; + SimulatedLoader(Action action, boolean delay, ScheduledExecutorService executor) + { + super(null, null, null, null); + this.action = action; + this.delay = delay; + this.executor = executor; + } + + @Override + public boolean load(AsyncOperation.Context context, BiConsumer callback) + { + if (delay) + { + executor.schedule(() -> { + callback.accept(null, action == Action.FAILURE ? new SimulatedFault("Failure loading " + context) : null); + }, 1, TimeUnit.SECONDS); + delay = false; + return false; + } + if (action == Action.FAILURE) + throw new SimulatedFault("Failure loading " + context); + + return true; + } + } +} diff --git a/test/unit/org/apache/cassandra/streaming/async/StreamingInboundHandlerTest.java b/test/unit/org/apache/cassandra/streaming/async/StreamingInboundHandlerTest.java index 069d0fb58d..904272f7b8 100644 --- a/test/unit/org/apache/cassandra/streaming/async/StreamingInboundHandlerTest.java +++ b/test/unit/org/apache/cassandra/streaming/async/StreamingInboundHandlerTest.java @@ -30,7 +30,6 @@ import org.junit.Test; import io.netty.buffer.ByteBuf; import io.netty.channel.embedded.EmbeddedChannel; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.util.DataInputBuffer; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputBuffer; @@ -56,7 +55,6 @@ import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; public class StreamingInboundHandlerTest { - private NettyStreamingChannel streamingChannel; private EmbeddedChannel channel; private ByteBuf buf; @@ -125,7 +123,7 @@ public class StreamingInboundHandlerTest temp.flip(); DataInputPlus in = new DataInputBuffer(temp, false); // session not found - IncomingStreamMessage.serializer.deserialize(in, IPartitioner.global(), MessagingService.current_version); + IncomingStreamMessage.serializer.deserialize(in, MessagingService.current_version); } @Test diff --git a/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java b/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java index 2f962f1ab2..4c57a76c8b 100644 --- a/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java +++ b/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java @@ -20,6 +20,7 @@ package org.apache.cassandra.tcm.listeners; import java.util.Random; +import org.apache.cassandra.config.DatabaseDescriptor; import org.junit.Before; import org.junit.BeforeClass; import org.junit.Test; @@ -50,7 +51,7 @@ import static org.junit.Assert.assertNull; public class MetadataSnapshotListenerTest { private static final Logger logger = LoggerFactory.getLogger(MetadataSnapshotListenerTest.class); - private IPartitioner partitioner = Murmur3Partitioner.instance; + private final IPartitioner partitioner = Murmur3Partitioner.instance; private Random r; @BeforeClass @@ -59,6 +60,7 @@ public class MetadataSnapshotListenerTest // Set this so that we don't attempt to sort the random placements as this depends on a populated // TokenMap. This is a temporary element of ClusterMetadata, at least in the current form CassandraRelevantProperties.TCM_SORT_REPLICA_GROUPS.setBoolean(false); + DatabaseDescriptor.daemonInitialization(); } @Before