diff --git a/.gitmodules b/.gitmodules index 6e00943162..616dacf610 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,4 +1,4 @@ [submodule "modules/accord"] path = modules/accord - url = https://github.com/apache/cassandra-accord + url = https://github.com/apache/cassandra-accord.git branch = trunk diff --git a/modules/accord b/modules/accord index 6c6872270e..746dabe0b4 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 6c6872270e16d2e777f1fa2c510b8f15396be3f3 +Subproject commit 746dabe0b43bf719badbd605e68a76037d01256d diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index ccefefbb8b..dbf0d0e242 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -880,7 +880,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte if (redundantBeforeEntry == null) return row; - TxnId redundantBeforeTxnId = redundantBeforeEntry.redundantBefore; + TxnId redundantBeforeTxnId = redundantBeforeEntry.shardRedundantBefore(); Cell lastExecuteMicrosCell = row.getCell(last_executed_micros); Long last_execute_micros = null; @@ -937,7 +937,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte if (redundantBeforeEntry == null) return row; - TxnId redundantBeforeTxnId = redundantBeforeEntry.redundantBefore; + TxnId redundantBeforeTxnId = redundantBeforeEntry.shardRedundantBefore(); Timestamp timestamp = CommandsForKeyRows.getTimestamp(row); if (timestamp != null && timestamp.compareTo(redundantBeforeTxnId) < 0) return null; diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 696cfe227c..fe9c397358 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -452,7 +452,7 @@ public class AccordCommandStore extends CommandStore implements CacheSize current = null; } - O mapReduceForRange(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) + O mapReduceForRange(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) { keysOrRanges = keysOrRanges.slice(slice, Routables.Slice.Minimal); switch (keysOrRanges.domain()) @@ -552,6 +552,9 @@ public class AccordCommandStore extends CommandStore implements CacheSize commandsForRanges.prune(globalSyncId, ranges); } + public NavigableMap bootstrapBeganAt() { return super.bootstrapBeganAt(); } + public NavigableMap safeToRead() { return super.safeToRead(); } + MessageProvider makeMessageProvider(TxnId txnId) { return journal.makeMessageProvider(txnId); diff --git a/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java b/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java index 6c9622cffd..fd9880d3e2 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java +++ b/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java @@ -139,7 +139,7 @@ public class AccordMessageSink implements MessageSink builder.put(MessageType.WAIT_ON_COMMIT_REQ, Verb.ACCORD_WAIT_ON_COMMIT_REQ); builder.put(MessageType.WAIT_ON_COMMIT_RSP, Verb.ACCORD_WAIT_ON_COMMIT_RSP); builder.put(MessageType.WAIT_UNTIL_APPLIED_REQ, Verb.ACCORD_WAIT_UNTIL_APPLIED_REQ); - builder.put(MessageType.APPLY_AND_WAIT_UNTIL_APPLIED_REQ, Verb.ACCORD_APPLY_AND_WAIT_UNTIL_APPLIED_REQ); + builder.put(MessageType.APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, Verb.ACCORD_APPLY_AND_WAIT_UNTIL_APPLIED_REQ); builder.put(MessageType.INFORM_OF_TXN_REQ, Verb.ACCORD_INFORM_OF_TXN_REQ); builder.put(MessageType.INFORM_DURABLE_REQ, Verb.ACCORD_INFORM_DURABLE_REQ); builder.put(MessageType.INFORM_HOME_DURABLE_REQ, Verb.ACCORD_INFORM_HOME_DURABLE_REQ); diff --git a/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java index 4fb9bc203c..122ae52a13 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java @@ -51,6 +51,7 @@ import accord.primitives.Routables; import accord.primitives.Seekable; import accord.primitives.Seekables; import accord.primitives.Timestamp; +import accord.primitives.Txn; import accord.primitives.TxnId; import org.apache.cassandra.service.accord.serializers.CommandsForKeySerializer; @@ -201,7 +202,7 @@ public class AccordSafeCommandStore extends AbstractSafeCommandStore O mapReduce(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) + private O mapReduce(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) { accumulate = commandStore.mapReduceForRange(keysOrRanges, slice, map, accumulate, terminate); if (terminate.test(accumulate)) @@ -209,7 +210,7 @@ public class AccordSafeCommandStore extends AbstractSafeCommandStore O mapReduceForKey(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) + private O mapReduceForKey(Routables keysOrRanges, Ranges slice, BiFunction map, O accumulate, Predicate terminate) { switch (keysOrRanges.domain()) { @@ -252,14 +253,7 @@ public class AccordSafeCommandStore extends AbstractSafeCommandStore T mapReduce(Seekables keysOrRanges, Ranges slice, TestKind testKind, TestTimestamp testTimestamp, Timestamp timestamp, TestDep testDep, @Nullable TxnId depId, @Nullable Status minStatus, @Nullable Status maxStatus, CommandFunction map, T accumulate, T terminalValue) - { - Predicate terminate = Predicates.equalTo(terminalValue); - return mapReduceWithTerminate(keysOrRanges, slice, testKind, testTimestamp, timestamp, testDep, depId, minStatus, maxStatus, map, accumulate, terminate); - } - - @Override - public T mapReduceWithTerminate(Seekables keysOrRanges, Ranges slice, TestKind testKind, TestTimestamp testTimestamp, Timestamp timestamp, TestDep testDep, @Nullable TxnId depId, @Nullable Status minStatus, @Nullable Status maxStatus, CommandFunction map, T accumulate, Predicate terminate) { + public T mapReduce(Seekables keysOrRanges, Ranges slice, Txn.Kind.Kinds testKind, TestTimestamp testTimestamp, Timestamp timestamp, TestDep testDep, @Nullable TxnId depId, @Nullable Status minStatus, @Nullable Status maxStatus, CommandFunction map, P1 p1, T accumulate, Predicate terminate) { accumulate = mapReduce(keysOrRanges, slice, (forKey, prev) -> { CommandTimeseries timeseries; switch (testTimestamp) @@ -285,7 +279,7 @@ public class AccordSafeCommandStore extends AbstractSafeCommandStore execute(SafeCommandStore safeStore, Timestamp executeAt, PartialTxn txn) + protected AsyncChain execute(SafeCommandStore safeStore, Timestamp executeAt, PartialTxn txn, Ranges unavailable) { + // TODO (required): subtract unavailable ranges, either from read or from response (or on coordinator) return AsyncChains.ofCallable(Stage.READ.executor(), () -> new LocalReadData(ReadCommandVerbHandler.instance.doRead(command, false))); } diff --git a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropReadRepair.java b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropReadRepair.java index c16c99e33e..00aeb0f244 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropReadRepair.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropReadRepair.java @@ -134,8 +134,9 @@ public class AccordInteropReadRepair extends AbstractExecute } @Override - protected AsyncChain execute(SafeCommandStore safeStore, Timestamp executeAt, PartialTxn txn) + protected AsyncChain execute(SafeCommandStore safeStore, Timestamp executeAt, PartialTxn txn, Ranges unavailable) { + // TODO (required): subtract unavailable ranges, either from read or from response (or on coordinator) return AsyncChains.ofCallable(Verb.READ_REPAIR_REQ.stage.executor(), () -> { ReadRepairVerbHandler.instance.applyMutation(mutation); return Data.NOOP_DATA; diff --git a/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java index 4decbdf1f8..c273e76be5 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java @@ -63,7 +63,7 @@ public class ApplySerializers DepsSerializer.partialDeps.deserialize(in, version), CommandSerializers.nullablePartialTxn.deserialize(in, version), CommandSerializers.writes.deserialize(in, version), - Result.APPLIED); + CommandSerializers.APPLIED); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java index 734186ea0b..0d15b74f52 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java @@ -22,18 +22,20 @@ import java.io.IOException; import accord.api.Result; import accord.api.RoutingKey; +import accord.coordinate.Infer; import accord.local.SaveStatus; -import accord.local.Status; import accord.local.Status.Durability; +import accord.local.Status.Known; import accord.messages.CheckStatus; import accord.messages.CheckStatus.CheckStatusNack; import accord.messages.CheckStatus.CheckStatusOk; import accord.messages.CheckStatus.CheckStatusOkFull; import accord.messages.CheckStatus.CheckStatusReply; +import accord.messages.CheckStatus.FoundKnown; +import accord.messages.CheckStatus.FoundKnownMap; import accord.primitives.Ballot; import accord.primitives.PartialDeps; import accord.primitives.PartialTxn; -import accord.primitives.Ranges; import accord.primitives.Route; import accord.primitives.Timestamp; import accord.primitives.TxnId; @@ -48,7 +50,74 @@ import static accord.messages.CheckStatus.SerializationSupport.createOk; public class CheckStatusSerializers { - public static final IVersionedSerializer request = new IVersionedSerializer() + public static final IVersionedSerializer foundKnown = new IVersionedSerializer<>() + { + @Override + public void serialize(FoundKnown known, DataOutputPlus out, int version) throws IOException + { + CommandSerializers.known.serialize(known, out, version); + CommandSerializers.invalidIfNot.serialize(known.invalidIfNot, out, version); + CommandSerializers.isPreempted.serialize(known.isPreempted, out, version); + } + + @Override + public FoundKnown deserialize(DataInputPlus in, int version) throws IOException + { + Known known = CommandSerializers.known.deserialize(in, version); + Infer.InvalidIfNot invalidIfNot = CommandSerializers.invalidIfNot.deserialize(in, version); + Infer.IsPreempted isPreempted = CommandSerializers.isPreempted.deserialize(in, version); + return new FoundKnown(known, invalidIfNot, isPreempted); + } + + @Override + public long serializedSize(FoundKnown known, int version) + { + return CommandSerializers.known.serializedSize(known, version) + + CommandSerializers.invalidIfNot.serializedSize(known.invalidIfNot, version) + + CommandSerializers.isPreempted.serializedSize(known.isPreempted, version); + } + }; + + public static final IVersionedSerializer foundKnownMap = new IVersionedSerializer<>() + { + @Override + public void serialize(FoundKnownMap knownMap, DataOutputPlus out, int version) throws IOException + { + int size = knownMap.size(); + out.writeUnsignedVInt32(size); + for (int i = 0 ; i <= size ; ++i) + KeySerializers.routingKey.serialize(knownMap.startAt(i), out, version); + for (int i = 0 ; i < size ; ++i) + foundKnown.serialize(knownMap.valueAt(i), out, version); + } + + @Override + public FoundKnownMap deserialize(DataInputPlus in, int version) throws IOException + { + int size = in.readUnsignedVInt32(); + RoutingKey[] starts = new RoutingKey[size + 1]; + for (int i = 0 ; i <= size ; ++i) + starts[i] = KeySerializers.routingKey.deserialize(in, version); + FoundKnown[] values = new FoundKnown[size]; + for (int i = 0 ; i < size ; ++i) + values[i] = foundKnown.deserialize(in, version); + return FoundKnownMap.SerializerSupport.create(true, starts, values); + } + + @Override + public long serializedSize(FoundKnownMap knownMap, int version) + { + int size = knownMap.size(); + long result = TypeSizes.sizeofUnsignedVInt(size); + for (int i = 0 ; i <= size ; ++i) + result += KeySerializers.routingKey.serializedSize(knownMap.startAt(i), version); + for (int i = 0 ; i < size ; ++i) + result += foundKnown.serializedSize(knownMap.valueAt(i), version); + return result; + } + }; + + public static final IVersionedSerializer request = new IVersionedSerializer<>() { final CheckStatus.IncludeInfo[] infos = CheckStatus.IncludeInfo.values(); @@ -81,7 +150,7 @@ public class CheckStatusSerializers } }; - public static final IVersionedSerializer reply = new IVersionedSerializer() + public static final IVersionedSerializer reply = new IVersionedSerializer<>() { private static final byte OK = 0x00; private static final byte FULL = 0x01; @@ -98,9 +167,8 @@ public class CheckStatusSerializers CheckStatusOk ok = (CheckStatusOk) reply; out.write(reply instanceof CheckStatusOkFull ? FULL : OK); - KeySerializers.ranges.serialize(ok.truncated, out, version); - CommandSerializers.status.serialize(ok.invalidIfNotAtLeast, out, version); - CommandSerializers.saveStatus.serialize(ok.saveStatus, out, version); + foundKnownMap.serialize(ok.map, out, version); + CommandSerializers.saveStatus.serialize(ok.maxKnowledgeSaveStatus, out, version); CommandSerializers.saveStatus.serialize(ok.maxSaveStatus, out, version); CommandSerializers.ballot.serialize(ok.promised, out, version); CommandSerializers.ballot.serialize(ok.accepted, out, version); @@ -130,9 +198,8 @@ public class CheckStatusSerializers return CheckStatusNack.NotOwned; case OK: case FULL: - Ranges truncated = KeySerializers.ranges.deserialize(in, version); - Status invalidIfNotAtLeast = CommandSerializers.status.deserialize(in, version); - SaveStatus status = CommandSerializers.saveStatus.deserialize(in, version); + FoundKnownMap map = foundKnownMap.deserialize(in, version); + SaveStatus maxKnowledgeStatus = CommandSerializers.saveStatus.deserialize(in, version); SaveStatus maxStatus = CommandSerializers.saveStatus.deserialize(in, version); Ballot promised = CommandSerializers.ballot.deserialize(in, version); Ballot accepted = CommandSerializers.ballot.deserialize(in, version); @@ -143,7 +210,7 @@ public class CheckStatusSerializers RoutingKey homeKey = KeySerializers.nullableRoutingKey.deserialize(in, version); if (kind == OK) - return createOk(truncated, invalidIfNotAtLeast, status, maxStatus, promised, accepted, executeAt, + return createOk(map, maxKnowledgeStatus, maxStatus, promised, accepted, executeAt, isCoordinating, durability, route, homeKey); PartialTxn partialTxn = CommandSerializers.nullablePartialTxn.deserialize(in, version); @@ -151,13 +218,11 @@ public class CheckStatusSerializers Writes writes = CommandSerializers.nullableWrites.deserialize(in, version); Result result = null; - if (status == SaveStatus.PreApplied || status == SaveStatus.Applied - || status == SaveStatus.TruncatedApply || status == SaveStatus.TruncatedApplyWithOutcome || status == SaveStatus.TruncatedApplyWithDeps) - result = Result.APPLIED; - else if (status == SaveStatus.Invalidated) - result = Result.INVALIDATED; + if (maxKnowledgeStatus == SaveStatus.PreApplied || maxKnowledgeStatus == SaveStatus.Applied + || maxKnowledgeStatus == SaveStatus.TruncatedApply || maxKnowledgeStatus == SaveStatus.TruncatedApplyWithOutcome || maxKnowledgeStatus == SaveStatus.TruncatedApplyWithDeps) + result = CommandSerializers.APPLIED; - return createOk(truncated, invalidIfNotAtLeast, status, maxStatus, promised, accepted, executeAt, + return createOk(map, maxKnowledgeStatus, maxStatus, promised, accepted, executeAt, isCoordinating, durability, route, homeKey, partialTxn, committedDeps, writes, result); } } @@ -170,9 +235,8 @@ public class CheckStatusSerializers return size; CheckStatusOk ok = (CheckStatusOk) reply; - size += KeySerializers.ranges.serializedSize(ok.truncated, version); - size += CommandSerializers.status.serializedSize(ok.invalidIfNotAtLeast, version); - size += CommandSerializers.saveStatus.serializedSize(ok.saveStatus, version); + size += foundKnownMap.serializedSize(ok.map, version); + size += CommandSerializers.saveStatus.serializedSize(ok.maxKnowledgeSaveStatus, version); size += CommandSerializers.saveStatus.serializedSize(ok.maxSaveStatus, version); size += CommandSerializers.ballot.serializedSize(ok.promised, version); size += CommandSerializers.ballot.serializedSize(ok.accepted, version); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java index 0c15515206..46232a811e 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java @@ -24,7 +24,9 @@ import com.google.common.base.Preconditions; import accord.api.Query; import accord.api.Read; +import accord.api.Result; import accord.api.Update; +import accord.coordinate.Infer; import accord.local.Node; import accord.local.SaveStatus; import accord.local.Status; @@ -32,6 +34,7 @@ import accord.local.Status.Durability; import accord.local.Status.Known; import accord.primitives.Ballot; import accord.primitives.PartialTxn; +import accord.primitives.ProgressToken; import accord.primitives.Ranges; import accord.primitives.Seekables; import accord.primitives.Timestamp; @@ -54,6 +57,16 @@ public class CommandSerializers { private CommandSerializers() {} + // TODO (expected): this is meant to encode e.g. whether the transaction's condition met or not + public static final Result APPLIED = new Result() + { + @Override + public ProgressToken asProgressToken() + { + return ProgressToken.APPLIED; + } + }; + public static final TimestampSerializer txnId = new TimestampSerializer<>(TxnId::fromBits); public static final TimestampSerializer timestamp = new TimestampSerializer<>(Timestamp::fromBits); public static final IVersionedSerializer nullableTimestamp = NullableSerializer.wrap(timestamp); @@ -228,16 +241,20 @@ public class CommandSerializers public static final IVersionedSerializer nullableWrites = NullableSerializer.wrap(writes); + public static final EnumSerializer route = new EnumSerializer<>(Status.KnownRoute.class); public static final EnumSerializer definition = new EnumSerializer<>(Status.Definition.class); public static final EnumSerializer knownExecuteAt = new EnumSerializer<>(Status.KnownExecuteAt.class); public static final EnumSerializer knownDeps = new EnumSerializer<>(Status.KnownDeps.class); public static final EnumSerializer outcome = new EnumSerializer<>(Status.Outcome.class); + public static final EnumSerializer invalidIfNot = new EnumSerializer<>(Infer.InvalidIfNot.class); + public static final EnumSerializer isPreempted = new EnumSerializer<>(Infer.IsPreempted.class); - public static final IVersionedSerializer known = new IVersionedSerializer() + public static final IVersionedSerializer known = new IVersionedSerializer<>() { @Override public void serialize(Known known, DataOutputPlus out, int version) throws IOException { + route.serialize(known.route, out, version); definition.serialize(known.definition, out, version); knownExecuteAt.serialize(known.executeAt, out, version); knownDeps.serialize(known.deps, out, version); @@ -247,7 +264,8 @@ public class CommandSerializers @Override public Known deserialize(DataInputPlus in, int version) throws IOException { - return new Known(definition.deserialize(in, version), + return new Known(route.deserialize(in, version), + definition.deserialize(in, version), knownExecuteAt.deserialize(in, version), knownDeps.deserialize(in, version), outcome.deserialize(in, version)); @@ -256,7 +274,8 @@ public class CommandSerializers @Override public long serializedSize(Known known, int version) { - return definition.serializedSize(known.definition, version) + return route.serializedSize(known.route, version) + + definition.serializedSize(known.definition, version) + knownExecuteAt.serializedSize(known.executeAt, version) + knownDeps.serializedSize(known.deps, version) + outcome.serializedSize(known.outcome, version); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/CommandStoreSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/CommandStoreSerializers.java index 1dd60db796..34dfced93d 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/CommandStoreSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/CommandStoreSerializers.java @@ -139,8 +139,10 @@ public class CommandStoreSerializers out.writeUnsignedVInt(t.startEpoch); if (t.endEpoch == Long.MAX_VALUE) out.writeUnsignedVInt(0L); else out.writeUnsignedVInt(1 + t.endEpoch - t.startEpoch); - CommandSerializers.txnId.serialize(t.redundantBefore, out, version); + CommandSerializers.txnId.serialize(t.locallyAppliedOrInvalidatedBefore, out, version); + CommandSerializers.txnId.serialize(t.shardAppliedOrInvalidatedBefore, out, version); CommandSerializers.txnId.serialize(t.bootstrappedAt, out, version); + CommandSerializers.nullableTimestamp.serialize(t.staleUntilAtLeast, out, version); } @Override @@ -151,9 +153,11 @@ public class CommandStoreSerializers long endEpoch = in.readUnsignedVInt(); if (endEpoch == 0) endEpoch = Long.MAX_VALUE; else endEpoch = startEpoch + 1 + endEpoch; - TxnId redundantBefore = CommandSerializers.txnId.deserialize(in, version); TxnId bootstrappedAt = CommandSerializers.txnId.deserialize(in, version); - return new RedundantBefore.Entry(range, startEpoch, endEpoch, redundantBefore, bootstrappedAt); + TxnId locallyAppliedOrInvalidatedBefore = CommandSerializers.txnId.deserialize(in, version); + TxnId shardAppliedOrInvalidatedBefore = CommandSerializers.txnId.deserialize(in, version); + Timestamp staleUntilAtLeast = CommandSerializers.nullableTimestamp.deserialize(in, version); + return new RedundantBefore.Entry(range, startEpoch, endEpoch, locallyAppliedOrInvalidatedBefore, shardAppliedOrInvalidatedBefore, bootstrappedAt, staleUntilAtLeast); } @Override @@ -162,8 +166,10 @@ public class CommandStoreSerializers long size = TokenRange.serializer.serializedSize((TokenRange) t.range, version); size += TypeSizes.sizeofUnsignedVInt(t.startEpoch); size += TypeSizes.sizeofUnsignedVInt(t.endEpoch == Long.MAX_VALUE ? 0 : 1 + t.endEpoch - t.startEpoch); - size += CommandSerializers.txnId.serializedSize(t.redundantBefore, version); + size += CommandSerializers.txnId.serializedSize(t.locallyAppliedOrInvalidatedBefore, version); + size += CommandSerializers.txnId.serializedSize(t.shardAppliedOrInvalidatedBefore, version); size += CommandSerializers.txnId.serializedSize(t.bootstrappedAt, version); + size += CommandSerializers.nullableTimestamp.serializedSize(t.staleUntilAtLeast, version); return size; } }), RedundantBefore.Entry[]::new, RedundantBefore.SerializerSupport::create); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/FetchSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/FetchSerializers.java index 4b184b49fc..0bdae00bb5 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/FetchSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/FetchSerializers.java @@ -28,6 +28,7 @@ import accord.impl.AbstractFetchCoordinator.FetchResponse; import accord.local.SaveStatus; import accord.local.Status.Durability; import accord.local.Status.Known; +import accord.messages.CheckStatus; import accord.messages.Propagate; import accord.messages.ReadData; import accord.messages.ReadData.ReadReply; @@ -145,16 +146,18 @@ public class FetchSerializers { CommandSerializers.txnId.serialize(p.txnId, out, version); KeySerializers.route.serialize(p.route, out, version); - CommandSerializers.saveStatus.serialize(p.saveStatus, out, version); + CommandSerializers.saveStatus.serialize(p.maxKnowledgeSaveStatus, out, version); CommandSerializers.saveStatus.serialize(p.maxSaveStatus, out, version); CommandSerializers.durability.serialize(p.durability, out, version); KeySerializers.nullableRoutingKey.serialize(p.homeKey, out, version); KeySerializers.nullableRoutingKey.serialize(p.progressKey, out, version); CommandSerializers.known.serialize(p.achieved, out, version); + CheckStatusSerializers.foundKnownMap.serialize(p.known, out, version); + out.writeBoolean(p.isTruncated); CommandSerializers.nullablePartialTxn.serialize(p.partialTxn, out, version); - DepsSerializer.nullablePartialDeps.serialize(p.partialDeps, out, version); + DepsSerializer.nullablePartialDeps.serialize(p.committedDeps, out, version); out.writeLong(p.toEpoch); - CommandSerializers.nullableTimestamp.serialize(p.executeAt, out, version); + CommandSerializers.nullableTimestamp.serialize(p.committedExecuteAt, out, version); CommandSerializers.nullableWrites.serialize(p.writes, out, version); } @@ -169,6 +172,8 @@ public class FetchSerializers RoutingKey homeKey = KeySerializers.nullableRoutingKey.deserialize(in, version); RoutingKey progressKey = KeySerializers.nullableRoutingKey.deserialize(in, version); Known achieved = CommandSerializers.known.deserialize(in, version); + CheckStatus.FoundKnownMap known = CheckStatusSerializers.foundKnownMap.deserialize(in, version); + boolean isTruncated = in.readBoolean(); PartialTxn partialTxn = CommandSerializers.nullablePartialTxn.deserialize(in, version); PartialDeps partialDeps = DepsSerializer.nullablePartialDeps.deserialize(in, version); long toEpoch = in.readLong(); @@ -184,14 +189,11 @@ public class FetchSerializers case TruncatedApply: case TruncatedApplyWithOutcome: case TruncatedApplyWithDeps: - result = Result.APPLIED; - break; - case Invalidated: - result = Result.INVALIDATED; + result = CommandSerializers.APPLIED; break; } - return Propagate.SerializerSupport.create(txnId, route, saveStatus, maxSaveStatus, durability, homeKey, progressKey, achieved, partialTxn, partialDeps, toEpoch, executeAt, writes, result); + return Propagate.SerializerSupport.create(txnId, route, saveStatus, maxSaveStatus, durability, homeKey, progressKey, achieved, known, isTruncated, partialTxn, partialDeps, toEpoch, executeAt, writes, result); } @Override @@ -199,16 +201,18 @@ public class FetchSerializers { return CommandSerializers.txnId.serializedSize(p.txnId, version) + KeySerializers.route.serializedSize(p.route, version) - + CommandSerializers.saveStatus.serializedSize(p.saveStatus, version) + + CommandSerializers.saveStatus.serializedSize(p.maxKnowledgeSaveStatus, version) + CommandSerializers.saveStatus.serializedSize(p.maxSaveStatus, version) + CommandSerializers.durability.serializedSize(p.durability, version) + KeySerializers.nullableRoutingKey.serializedSize(p.homeKey, version) + KeySerializers.nullableRoutingKey.serializedSize(p.progressKey, version) + CommandSerializers.known.serializedSize(p.achieved, version) + + CheckStatusSerializers.foundKnownMap.serializedSize(p.known, version) + + TypeSizes.BOOL_SIZE + CommandSerializers.nullablePartialTxn.serializedSize(p.partialTxn, version) - + DepsSerializer.nullablePartialDeps.serializedSize(p.partialDeps, version) + + DepsSerializer.nullablePartialDeps.serializedSize(p.committedDeps, version) + TypeSizes.sizeof(p.toEpoch) - + CommandSerializers.nullableTimestamp.serializedSize(p.executeAt, version) + + CommandSerializers.nullableTimestamp.serializedSize(p.committedExecuteAt, version) + CommandSerializers.nullableWrites.serializedSize(p.writes, version) ; } diff --git a/src/java/org/apache/cassandra/service/accord/serializers/RecoverySerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/RecoverySerializers.java index 99e54b5b45..adf60212cc 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/RecoverySerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/RecoverySerializers.java @@ -129,9 +129,7 @@ public class RecoverySerializers Result result = null; if (status == Status.PreApplied || status == Status.Applied || status == Status.Truncated) - result = Result.APPLIED; - else if (status == Status.Invalidated) - result = Result.INVALIDATED; + result = CommandSerializers.APPLIED; return deserializeOk(id, status, diff --git a/src/java/org/apache/cassandra/service/accord/txn/TxnWrite.java b/src/java/org/apache/cassandra/service/accord/txn/TxnWrite.java index 27406036d9..3848908216 100644 --- a/src/java/org/apache/cassandra/service/accord/txn/TxnWrite.java +++ b/src/java/org/apache/cassandra/service/accord/txn/TxnWrite.java @@ -374,7 +374,7 @@ public class TxnWrite extends AbstractKeySorted implements Writ // cfk into memory by retaining at all times in memory key ranges that are dirty and must use this logic; // any that aren't can just use executeAt.hlc AccordSafeCommandsForKey cfk = ((AccordSafeCommandStore) safeStore).commandsForKey((RoutableKey) key); - cfk.updateLastExecutionTimestamps(executeAt, true); + cfk.updateLastExecutionTimestamps(safeStore, executeAt, true); long timestamp = cfk.timestampMicrosFor(executeAt, true); // TODO (low priority - do we need to compute nowInSeconds, or can we just use executeAt?) int nowInSeconds = cfk.nowInSecondsFor(executeAt, true); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java index 377a919940..8dff79cd86 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java @@ -28,7 +28,6 @@ import java.util.function.Consumer; import org.junit.Assert; import org.junit.Test; -import accord.local.CommandStore; import accord.local.PreLoadContext; import accord.primitives.Timestamp; import accord.topology.TopologyManager; @@ -48,6 +47,7 @@ import org.apache.cassandra.distributed.test.TestBaseImpl; import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.service.accord.AccordCommandStore; import org.apache.cassandra.service.accord.AccordConfigurationService; import org.apache.cassandra.service.accord.AccordConfigurationService.EpochSnapshot; import org.apache.cassandra.service.accord.AccordService; @@ -271,7 +271,7 @@ public class AccordBootstrapTest extends TestBaseImpl }); awaitUninterruptiblyAndRethrow(service().node().commandStores().forEach(safeStore -> { - CommandStore commandStore = safeStore.commandStore(); + AccordCommandStore commandStore = (AccordCommandStore) safeStore.commandStore(); Assert.assertEquals(Timestamp.NONE, getOnlyElement(commandStore.bootstrapBeganAt().keySet())); Assert.assertEquals(Timestamp.NONE, getOnlyElement(commandStore.safeToRead().keySet())); // @@ -316,7 +316,7 @@ public class AccordBootstrapTest extends TestBaseImpl awaitUninterruptiblyAndRethrow(service().node().commandStores().forEach(safeStore -> { if (safeStore.ranges().currentRanges().contains(partitionKey)) { - CommandStore commandStore = safeStore.commandStore(); + AccordCommandStore commandStore = (AccordCommandStore) safeStore.commandStore(); Assert.assertFalse(commandStore.bootstrapBeganAt().isEmpty()); Assert.assertFalse(commandStore.safeToRead().isEmpty()); @@ -458,7 +458,7 @@ public class AccordBootstrapTest extends TestBaseImpl safeStore -> { if (!safeStore.ranges().allAt(preMove).contains(partitionKey)) { - CommandStore commandStore = safeStore.commandStore(); + AccordCommandStore commandStore = (AccordCommandStore) safeStore.commandStore(); Assert.assertFalse(commandStore.bootstrapBeganAt().isEmpty()); Assert.assertFalse(commandStore.safeToRead().isEmpty()); diff --git a/test/simulator/main/org/apache/cassandra/simulator/RandomSource.java b/test/simulator/main/org/apache/cassandra/simulator/RandomSource.java index 14d7ad9b1d..4e429e418f 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/RandomSource.java +++ b/test/simulator/main/org/apache/cassandra/simulator/RandomSource.java @@ -20,13 +20,17 @@ package org.apache.cassandra.simulator; import java.lang.reflect.Array; import java.util.Arrays; +import java.util.List; import java.util.Map; import java.util.Random; +import java.util.Set; import java.util.function.IntSupplier; import java.util.function.LongSupplier; import java.util.stream.IntStream; import java.util.stream.LongStream; +import com.google.common.collect.Iterators; + import org.apache.cassandra.utils.Shared; import static org.apache.cassandra.utils.Shared.Scope.SIMULATION; @@ -46,11 +50,20 @@ public interface RandomSource } public T choose(RandomSource random) + { + return choose(random.uniformFloat()); + } + + public T choose(accord.utils.RandomSource random) + { + return choose(random.nextFloat()); + } + + private T choose(float choose) { if (options.length == 0) return null; - float choose = random.uniformFloat(); int i = Arrays.binarySearch(cumulativeProbabilities, choose); if (i < 0) i = -1 - i; @@ -131,6 +144,41 @@ public interface RandomSource Arrays.fill(nonCumulativeProbabilities, 1f / options.length); return new Choices<>(cumulativeProbabilities(nonCumulativeProbabilities), options); } + + public static T choose(RandomSource rs, Set set) + { + return choose(rs.uniform(0, set.size()), set); + } + + public static T choose(accord.utils.RandomSource rs, Set set) + { + return choose(rs.nextInt(set.size()), set); + } + + private static T choose(int i, Set set) + { + return Iterators.get(set.iterator(), i); + } + + public static T choose(RandomSource rs, List list) + { + return list.get(rs.uniform(0, list.size())); + } + + public static T choose(accord.utils.RandomSource rs, List list) + { + return list.get(rs.nextInt(list.size())); + } + + public static T choose(RandomSource rs, T ... array) + { + return array[rs.uniform(0, array.length)]; + } + + public static T choose(accord.utils.RandomSource rs, T ... array) + { + return array[rs.nextInt(array.length)]; + } } public static abstract class Abstract implements RandomSource diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java index 33472cf09f..6d180f82ec 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java @@ -350,7 +350,7 @@ public class CompactionAccordIteratorsTest private static RedundantBefore redundantBefore(TxnId txnId) { Ranges ranges = AccordTestUtils.fullRange(AccordTestUtils.keys(table, 42)); - return RedundantBefore.create(ranges, Long.MIN_VALUE, Long.MAX_VALUE, txnId, LT_TXN_ID); + return RedundantBefore.create(ranges, Long.MIN_VALUE, Long.MAX_VALUE, txnId, txnId, LT_TXN_ID); } enum DurableBeforeType diff --git a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java index 50716c80b9..9eafb67d9b 100644 --- a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java +++ b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java @@ -464,13 +464,12 @@ public abstract class FuzzTestBase extends CQLTester.InMemory { if (repair.state.isComplete()) throw new IllegalStateException("Repair is completed! " + repair.state.getResult()); - List participaents = new ArrayList<>(repair.state.getNeighborsAndRanges().participants.size() + 1); - if (rs.nextBoolean()) participaents.add(coordinator.broadcastAddressAndPort()); - participaents.addAll(repair.state.getNeighborsAndRanges().participants); - participaents.sort(Comparator.naturalOrder()); + List participants = new ArrayList<>(repair.state.getNeighborsAndRanges().participants.size() + 1); + if (rs.nextBoolean()) participants.add(coordinator.broadcastAddressAndPort()); + participants.addAll(repair.state.getNeighborsAndRanges().participants); + participants.sort(Comparator.naturalOrder()); - InetAddressAndPort selected = rs.pick(participaents); - return selected; + return participants.get(rs.nextInt(participants.size())); } static void addMismatch(RandomSource rs, ColumnFamilyStore cfs, Validator validator) diff --git a/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java b/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java index e3cb042cb1..928daf6696 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java @@ -54,6 +54,7 @@ import org.apache.cassandra.schema.KeyspaceParams; import org.apache.cassandra.schema.SchemaConstants; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.service.accord.api.PartitionKey; +import org.apache.cassandra.service.accord.serializers.CommandSerializers; import org.apache.cassandra.service.accord.serializers.CommandsForKeySerializer; import org.apache.cassandra.utils.Pair; @@ -126,9 +127,8 @@ public class AccordCommandStoreTest attrs.addListener(new Command.ProxyListener(oldTxnId1)); Pair result = AccordTestUtils.processTxnResult(commandStore, txnId, txn, executeAt); - Command command = Command.SerializerSupport.executed(attrs, SaveStatus.Applied, executeAt, promised, accepted, - waitingOn, result.left, Result.APPLIED); + waitingOn, result.left, CommandSerializers.APPLIED); AccordSafeCommand safeCommand = new AccordSafeCommand(loaded(txnId, null)); safeCommand.set(command); @@ -142,7 +142,7 @@ public class AccordCommandStoreTest dependencies, txn, result.left, - Result.APPLIED); + CommandSerializers.APPLIED); commandStore.appendToJournal(apply); AccordKeyspace.getCommandMutation(commandStore, safeCommand, commandStore.nextSystemTimestampMicros()).apply(); @@ -172,10 +172,10 @@ public class AccordCommandStoreTest cfk.initialize(CommandsForKeySerializer.loader); cfk.updateMax(maxTimestamp); - cfk.updateLastExecutionTimestamps(txnId1, true); + cfk.updateLastExecutionTimestamps(null, txnId1, true); Assert.assertEquals(txnId1.hlc(), cfk.timestampMicrosFor(txnId1, true)); - cfk.updateLastExecutionTimestamps(txnId2, true); + cfk.updateLastExecutionTimestamps(null, txnId2, true); Assert.assertEquals(txnId2.hlc(), cfk.timestampMicrosFor(txnId2, true)); Assert.assertEquals(txnId2, cfk.current().lastExecutedTimestamp()); diff --git a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java index f3bb71707d..bc694eee64 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java @@ -75,6 +75,7 @@ import org.apache.cassandra.utils.concurrent.Future; import org.assertj.core.api.Assertions; import static accord.utils.Property.qt; +import static org.apache.cassandra.simulator.RandomSource.Choices.choose; public class AccordSyncPropagatorTest { @@ -121,7 +122,7 @@ public class AccordSyncPropagatorTest { for (Range range : ranges) { - Cluster.Instace inst = cluster.node(rs.pick(nodes)); + Cluster.Instace inst = cluster.node(choose(rs, nodes)); scheduler.schedule(() -> { Ranges subrange = Ranges.of(range); inst.propagator.reportClosed(epoch, nodes, subrange); diff --git a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java index 8a64409cca..64907c0209 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java @@ -53,6 +53,7 @@ import accord.local.SaveStatus; import accord.local.SaveStatus.LocalExecution; import accord.primitives.Ballot; import accord.primitives.FullKeyRoute; +import accord.primitives.FullRoute; import accord.primitives.Keys; import accord.primitives.PartialDeps; import accord.primitives.PartialTxn; @@ -108,6 +109,7 @@ public class AccordTestUtils { CommonAttributes.Mutable attrs = new CommonAttributes.Mutable(txnId); attrs.partialTxn(txn); + attrs.route(route(txn)); return Command.SerializerSupport.preaccepted(attrs, executeAt, Ballot.ZERO); } @@ -115,9 +117,7 @@ public class AccordTestUtils { CommonAttributes.Mutable attrs = new CommonAttributes.Mutable(txnId).partialDeps(PartialDeps.NONE); attrs.partialTxn(txn); - Seekable key = txn.keys().get(0); - RoutingKey routingKey = key.asKey().toUnseekable(); - attrs.route(new FullKeyRoute(routingKey, true, new RoutingKey[]{ routingKey})); + attrs.route(route(txn)); return Command.SerializerSupport.committed(attrs, SaveStatus.Committed, executeAt, @@ -125,6 +125,13 @@ public class AccordTestUtils Ballot.ZERO, Command.WaitingOn.EMPTY); } + + private static FullRoute route(PartialTxn txn) + { + Seekable key = txn.keys().get(0); + RoutingKey routingKey = key.asKey().toUnseekable(); + return new FullKeyRoute(routingKey, true, new RoutingKey[]{ routingKey }); + } } public static CommandsForKey commandsForKey(Key key) @@ -171,6 +178,7 @@ public class AccordTestUtils @Override public void unwitnessed(TxnId txnId, ProgressShard progressShard) {} @Override public void preaccepted(Command command, ProgressShard progressShard) {} @Override public void accepted(Command command, ProgressShard progressShard) {} + @Override public void precommitted(Command command) {} @Override public void committed(Command command, ProgressShard progressShard) {} @Override public void readyToExecute(Command command) {} @Override public void executed(Command command, ProgressShard progressShard) {} diff --git a/test/unit/org/apache/cassandra/service/accord/CommandsForRangesTest.java b/test/unit/org/apache/cassandra/service/accord/CommandsForRangesTest.java index fcdd5ad1ba..c86e0a769f 100644 --- a/test/unit/org/apache/cassandra/service/accord/CommandsForRangesTest.java +++ b/test/unit/org/apache/cassandra/service/accord/CommandsForRangesTest.java @@ -45,6 +45,7 @@ import org.apache.cassandra.utils.Interval; import org.apache.cassandra.utils.IntervalTree; import static accord.utils.Property.qt; +import static org.apache.cassandra.simulator.RandomSource.Choices.choose; import static org.assertj.core.api.Assertions.assertThat; public class CommandsForRangesTest @@ -91,7 +92,7 @@ public class CommandsForRangesTest private static Gen cfr() { // TODO (coverage): once all partitioners work with regard to splitting, then should test all - Gen partitionerGen = rs -> rs.pick(Murmur3Partitioner.instance, RandomPartitioner.instance); + Gen partitionerGen = rs -> choose(rs, Murmur3Partitioner.instance, RandomPartitioner.instance); Gen statusGen = Gens.enums().all(SaveStatus.class); return rs -> { IPartitioner partitioner = partitionerGen.next(rs); diff --git a/test/unit/org/apache/cassandra/utils/AccordGenerators.java b/test/unit/org/apache/cassandra/utils/AccordGenerators.java index f42dcd2300..7017e98664 100644 --- a/test/unit/org/apache/cassandra/utils/AccordGenerators.java +++ b/test/unit/org/apache/cassandra/utils/AccordGenerators.java @@ -26,6 +26,7 @@ import java.util.Set; import accord.local.Command; import accord.primitives.Deps; +import accord.primitives.FullRoute; import accord.primitives.KeyDeps; import accord.primitives.PartialTxn; import accord.primitives.Range; @@ -74,6 +75,8 @@ public class AccordGenerators //TODO goes against fuzz testing, and also limits to a very specific table existing... // There is a branch that can generate random transactions, so maybe look into that? PartialTxn txn = createPartialTxn(0); + FullRoute route = txn.keys().toRoute(txn.keys().get(0).someIntersectingRoutingKey(null)); + return rs -> { TxnId id = ids.next(rs); Timestamp executeAt = id;