diff --git a/modules/accord b/modules/accord index 20df7e5e5a..d5de79b56f 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 20df7e5e5ad7bae442ff6fe4ee716eb7104336c1 +Subproject commit d5de79b56f792d3b6d868ca46845b830e8906828 diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 972a7488e0..95aba5de55 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -55,11 +55,9 @@ import accord.local.RedundantBefore; import accord.local.SafeCommandStore; import accord.local.cfk.CommandsForKey; import accord.primitives.PartialTxn; -import accord.primitives.RangeDeps; import accord.primitives.Ranges; import accord.primitives.RoutableKey; import accord.primitives.Route; -import accord.primitives.Status; import accord.primitives.Timestamp; import accord.primitives.TxnId; import accord.utils.Invariants; @@ -226,14 +224,6 @@ public class AccordCommandStore extends CommandStore return commandsForRanges; } - @Override - public void markShardDurable(SafeCommandStore safeStore, TxnId globalSyncId, Ranges ranges, Status.Durability durability) - { - super.markShardDurable(safeStore, globalSyncId, ranges, durability); - if (durability.compareTo(Status.Durability.UniversalOrInvalidated) >= 0) - commandsForRanges.gcBefore(globalSyncId, ranges); - } - @Override public boolean inStore() { @@ -400,43 +390,6 @@ public class AccordCommandStore extends CommandStore { } - public void registerTransitive(SafeCommandStore safeStore, RangeDeps rangeDeps) - { - if (rangeDeps.isEmpty()) - return; - - RedundantBefore redundantBefore = unsafeGetRedundantBefore(); - CommandStores.RangesForEpoch ranges = safeStore.ranges(); - // used in places such as accord.local.CommandStore.fetchMajorityDeps - // We find a set of dependencies for a range then update CommandsFor to know about them - Ranges allRanges = safeStore.ranges().all(); - Ranges coordinateRanges = Ranges.EMPTY; - long coordinateEpoch = -1; - try (ExclusiveCaches caches = lockCaches()) - { - for (int i = 0; i < rangeDeps.txnIdCount(); i++) - { - TxnId txnId = rangeDeps.txnId(i); - AccordCacheEntry state = caches.commands().getUnsafe(txnId); - if (state != null && state.isLoaded() && state.getExclusive() != null && state.getExclusive().known().isDefinitionKnown()) - continue; - - Ranges addRanges = rangeDeps.ranges(i).slice(allRanges); - if (addRanges.isEmpty()) continue; - - if (coordinateEpoch != txnId.epoch()) - { - coordinateEpoch = txnId.epoch(); - coordinateRanges = ranges.allAt(txnId.epoch()); - } - if (addRanges.intersects(coordinateRanges)) continue; - addRanges = redundantBefore.removeGcBefore(txnId, txnId, addRanges); - if (addRanges.isEmpty()) continue; - commandsForRanges().mergeTransitive(txnId, addRanges, Ranges::with); - } - } - } - public void appendCommands(List diffs, Runnable onFlush) { for (int i = 0; i < diffs.size(); i++) @@ -668,8 +621,8 @@ public class AccordCommandStore extends CommandStore StringBuilder sb = new StringBuilder("["); if (metadata != null) sb.append(metadata).append('|'); - sb.append(tableId.lsb()); - sb.append(',').append(id).append(',').append(node.id().id).append(']'); + sb.append(tableId); + sb.append('|').append(id).append(',').append(node.id().id).append(']'); return sb.toString(); } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java index ff9169d08b..803839de94 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java @@ -85,13 +85,6 @@ class AccordExecutorLoops }; } - AccordExecutor.Task activeTaskForCurrentThread() - { - Thread thread = Thread.currentThread(); - LoopTask task = tasks.get(thread.getId()); - return task == null ? null : task.runningTask(); - } - public Stream active() { return tasks.values().stream(); diff --git a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java index 356200da23..f0dfcf4c66 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java +++ b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java @@ -302,7 +302,8 @@ public class AccordKeyspace if (current == null) return null; - CommandsForKey updated = current.withRedundantBeforeAtLeast(redundantBefore.gcBefore()); + // TODO (desired): consider whether better to not compact any validation failures, since we expect is already overwritten + CommandsForKey updated = current.withRedundantBeforeAtLeast(redundantBefore.gcBefore(), false); if (current == updated) return row; diff --git a/src/java/org/apache/cassandra/service/accord/AccordTask.java b/src/java/org/apache/cassandra/service/accord/AccordTask.java index 7a4fd0e8cf..f04d4ad7c0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTask.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTask.java @@ -1038,7 +1038,7 @@ public abstract class AccordTask extends SubmittableTask implements Function< { CommandsForRanges.Summary summary = summaryLoader.ifRelevant(state); if (summary != null) - summaries.put(summary.txnId, summary); + summaries.put(summary.plainTxnId(), summary); } } @@ -1059,7 +1059,10 @@ public abstract class AccordTask extends SubmittableTask implements Function< CommandsForRanges.Summary summary = summaryLoader.load(txnId); if (summary != null) + { summaries.putIfAbsent(txnId, summary); + summaryLoader.maybeRecordFutureRx(summary); + } }); } @@ -1086,7 +1089,7 @@ public abstract class AccordTask extends SubmittableTask implements Function< void startInternal(Caches caches) { summaryLoader = commandStore.commandsForRanges().loader(preLoadContext.primaryTxnId(), preLoadContext.loadKeysFor(), keysOrRanges); - summaryLoader.forEachInCache(keysOrRanges, summary -> summaries.put(summary.txnId, summary), caches); + summaryLoader.forEachInCache(keysOrRanges, summary -> summaries.put(summary.plainTxnId(), summary), caches); caches.commands().register(commandWatcher); } diff --git a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java index eddafc8398..b905ec519a 100644 --- a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java +++ b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java @@ -21,18 +21,14 @@ package org.apache.cassandra.service.accord; import java.util.ArrayList; import java.util.Comparator; import java.util.HashMap; -import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.NavigableMap; import java.util.TreeMap; -import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; -import java.util.function.BiFunction; import java.util.function.Consumer; -import java.util.function.UnaryOperator; import javax.annotation.Nullable; import accord.api.RoutingKey; @@ -49,7 +45,6 @@ import accord.primitives.RangeRoute; import accord.primitives.Ranges; import accord.primitives.Routable; import accord.primitives.Timestamp; -import accord.primitives.Txn; import accord.primitives.Txn.Kind.Kinds; import accord.primitives.TxnId; import accord.primitives.Unseekable; @@ -65,7 +60,6 @@ import org.apache.cassandra.utils.btree.BTreeSet; import org.apache.cassandra.utils.btree.IntervalBTree; import org.apache.cassandra.utils.concurrent.IntrusiveStack; -import static accord.local.CommandSummaries.SummaryStatus.NOT_DIRECTLY_WITNESSED; import static org.apache.cassandra.utils.btree.IntervalBTree.InclusiveEndHelper.endWithStart; import static org.apache.cassandra.utils.btree.IntervalBTree.InclusiveEndHelper.keyEndWithStart; import static org.apache.cassandra.utils.btree.IntervalBTree.InclusiveEndHelper.keyStartWithEnd; @@ -176,7 +170,6 @@ public class CommandsForRanges extends TreeMap implements Co private final AccordCommandStore commandStore; private final RangeSearcher searcher; - private final AtomicReference> transitive = new AtomicReference<>(new TreeMap<>()); // TODO (desired): manage memory consumed by this auxillary information private final Object2ObjectHashMap cachedRangeTxnsById = new Object2ObjectHashMap<>(); private Object[] cachedRangeTxnsByRange = IntervalBTree.empty(); @@ -356,80 +349,28 @@ public class CommandsForRanges extends TreeMap implements Co public CommandsForRanges.Loader loader(@Nullable TxnId primaryTxnId, LoadKeysFor loadKeysFor, Unseekables keysOrRanges) { RedundantBefore redundantBefore = commandStore.unsafeGetRedundantBefore(); - return Loader.loader(redundantBefore, primaryTxnId, loadKeysFor, keysOrRanges, this::newLoader); + MaxDecidedRX maxDecidedRX = commandStore.unsafeGetMaxDecidedRX(); + return SummaryLoader.loader(redundantBefore, maxDecidedRX, primaryTxnId, loadKeysFor, keysOrRanges, this::newLoader); } - private Loader newLoader(@Nullable TxnId primaryTxnId, Unseekables searchKeysOrRanges, RedundantBefore redundantBefore, Kinds testKind, TxnId minTxnId, Timestamp maxTxnId, @Nullable TxnId findAsDep) + private Loader newLoader(RedundantBefore redundantBefore, MaxDecidedRX maxDecidedRX, @Nullable TxnId primaryTxnId, Unseekables searchKeysOrRanges, Kinds testKind, TxnId minTxnId, Timestamp maxTxnId, @Nullable TxnId findAsDep) { - MaxDecidedRX maxDecidedRX = null; - if (primaryTxnId != null && primaryTxnId.is(Txn.Kind.ExclusiveSyncPoint) && findAsDep == null) - maxDecidedRX = commandStore.unsafeGetMaxDecidedRX(); - return new Loader(this, primaryTxnId, searchKeysOrRanges, redundantBefore, testKind, minTxnId, maxTxnId, findAsDep, maxDecidedRX); - } - - private void updateTransitive(UnaryOperator> update) - { - while (true) - { - NavigableMap prev = transitive.get(); - NavigableMap next = update.apply(prev); - if (next == null || prev == next) - return; - if (transitive.compareAndSet(prev, next)) - return; - } - } - - public void mergeTransitive(TxnId txnId, Ranges ranges, BiFunction remappingFunction) - { - updateTransitive(transitive -> { - NavigableMap next = new TreeMap<>(transitive); - next.merge(txnId, ranges, remappingFunction); - return next; - }); - } - - public void gcBefore(TxnId gcBefore, Ranges ranges) - { - updateTransitive(transitive -> { - NavigableMap next = null; - Iterator> iterator = transitive.headMap(gcBefore).entrySet().iterator(); - while (iterator.hasNext()) - { - Map.Entry e = iterator.next(); - Ranges newRanges = e.getValue().without(ranges); - if (!newRanges.isEmpty()) - { - if (next == null) - next = new TreeMap<>(); - next.put(e.getKey(), newRanges); - } - } - return next; - }); + return new Loader(this, redundantBefore, maxDecidedRX, primaryTxnId, searchKeysOrRanges, testKind, minTxnId, maxTxnId, findAsDep); } } - public static class Loader extends Summary.Loader + public static class Loader extends SummaryLoader { private final Manager manager; - private final MaxDecidedRX maxDecidedRX; - private final TxnId primaryTxnId; - @Nullable - private final TxnId minDecidedId; - public Loader(Manager manager, TxnId primaryTxnId, Unseekables searchKeysOrRanges, RedundantBefore redundantBefore, Kinds testKinds, TxnId minTxnId, Timestamp maxTxnId, @Nullable TxnId findAsDep, @Nullable MaxDecidedRX maxDecidedRX) + public Loader(Manager manager, RedundantBefore redundantBefore, MaxDecidedRX maxDecidedRX, TxnId primaryTxnId, Unseekables searchKeysOrRanges, Kinds testKinds, TxnId minTxnId, Timestamp maxTxnId, @Nullable TxnId findAsDep) { - super(primaryTxnId, searchKeysOrRanges, redundantBefore, testKinds, minTxnId, maxTxnId, findAsDep); + super(redundantBefore, maxDecidedRX, primaryTxnId, searchKeysOrRanges, testKinds, minTxnId, maxTxnId, findAsDep); this.manager = manager; - this.maxDecidedRX = maxDecidedRX; - this.primaryTxnId = primaryTxnId; - this.minDecidedId = MaxDecidedRX.minDecidedDependencyId(maxDecidedRX, searchKeysOrRanges, primaryTxnId); } public void intersects(Consumer forEach) { - // TODO (expected): use the ranges we find to filter results by MaxDecidedRX (don't just consume the TxnId) switch (searchKeysOrRanges.domain()) { case Range: @@ -440,38 +381,11 @@ public class CommandsForRanges extends TreeMap implements Co for (Unseekable key : searchKeysOrRanges) manager.searcher.search(manager.commandStore.id(), (TokenKey) key, minTxnId, maxTxnId, minDecidedId).consume(forEach); } - - NavigableMap transitive = manager.transitive.get(); - if (!transitive.isEmpty()) - { - for (Map.Entry e : transitive.tailMap(minTxnId, true).entrySet()) - { - if (e.getValue().intersects(searchKeysOrRanges)) - forEach.accept(e.getKey()); - } - } } - boolean isRelevant(TxnIdInterval txnIdInterval) + boolean isMaybeRelevant(TxnIdInterval txnIdInterval) { - if (maxDecidedRX == null) - return true; - - if (!isMaybeRelevant(txnIdInterval.txnId)) - return false; - - TxnId minRelevantId = MaxDecidedRX.minDecidedDependencyId(maxDecidedRX, Ranges.of(txnIdInterval), primaryTxnId); - return isRelevant(minRelevantId, primaryTxnId); - } - - private boolean isRelevant(@Nullable TxnId minRelevantId, TxnId txnId) - { - return minRelevantId == null || minRelevantId.compareTo(txnId) <= 0; - } - - boolean isMaybeRelevant(TxnId txnId) - { - return isRelevant(minDecidedId, txnId); + return isMaybeRelevant(txnIdInterval.txnId, null, Ranges.of(txnIdInterval)); } public void forEachInCache(Unseekables keysOrRanges, Consumer forEach, AccordCommandStore.Caches caches) @@ -484,10 +398,9 @@ public class CommandsForRanges extends TreeMap implements Co for (RoutingKey key : (AbstractUnseekableKeys)keysOrRanges) { IntervalBTree.accumulate(manager.cachedRangeTxnsByRange(), KEY_COMPARATORS, key, (f, s, i, c) -> { - TxnIdInterval interval = (TxnIdInterval)i; - if (isRelevant(interval)) + if (isMaybeRelevant(i)) { - TxnId txnId = ((TxnIdInterval)i).txnId; + TxnId txnId = i.txnId; Summary summary = ifRelevant(c.getUnsafe(txnId)); if (summary != null) f.accept(summary); @@ -502,7 +415,7 @@ public class CommandsForRanges extends TreeMap implements Co for (Range range : (AbstractRanges)keysOrRanges) { IntervalBTree.accumulate(manager.cachedRangeTxnsByRange(), COMPARATORS, new TxnIdInterval(range.start(), range.end(), TxnId.NONE), (f, s, i, c) -> { - if (isRelevant(i)) + if (isMaybeRelevant(i)) { TxnId txnId = i.txnId; AccordCacheEntry entry = c.getUnsafe(txnId); @@ -540,15 +453,7 @@ public class CommandsForRanges extends TreeMap implements Co return ifRelevant(cmd); } - Ranges ranges = manager.transitive.get().get(txnId); - if (ranges == null) - return null; - - ranges = ranges.intersecting(searchKeysOrRanges); - if (ranges.isEmpty()) - return null; - - return new Summary(txnId, txnId, NOT_DIRECTLY_WITNESSED, ranges, null, null); + return null; } public Summary ifRelevant(AccordCacheEntry state) @@ -581,15 +486,10 @@ public class CommandsForRanges extends TreeMap implements Co return ifRelevant(command); } - public Summary ifRelevant(Command cmd) - { - return ifRelevant(cmd.txnId(), cmd.executeAt(), cmd.saveStatus(), cmd.participants(), cmd.partialDeps()); - } - public Summary ifRelevant(Command.Minimal cmd) { Invariants.require(findAsDep == null); - return ifRelevant(cmd.txnId, cmd.executeAt, cmd.saveStatus, cmd.participants, null); + return ifRelevant(cmd.txnId, cmd.executeAt, cmd.saveStatus, cmd.durability, cmd.participants, null); } } } diff --git a/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java b/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java index 9ebeba295f..b43b08fc20 100644 --- a/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java +++ b/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java @@ -253,7 +253,8 @@ public class AccordAgent implements Agent long mostRecentStart = Math.max(command.txnId().hlc(), promisedHlc); long waitMicros = recover(txnId).computeWait(retryCount, MICROSECONDS); long nowMicros = MILLISECONDS.toMicros(Clock.Global.currentTimeMillis()); - Invariants.expect(mostRecentStart <= nowMicros + SECONDS.toMicros(1L), "max(%s,%s)>%d", command.txnId(), command.promised(), nowMicros); + if (mostRecentStart > nowMicros + SECONDS.toMicros(1L)) + logger.warn("max({},{})>{}", command.txnId(), command.promised(), nowMicros); long startTime = mostRecentStart + waitMicros; if (startTime < nowMicros) startTime = nowMicros + waitMicros/2; diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java index 9c9835e4a8..945b8d4537 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java @@ -23,9 +23,12 @@ import com.google.common.base.Throwables; import accord.api.RoutingKey; import accord.local.CommandStores; import accord.local.LoadKeys; +import accord.local.Node; import accord.local.PreLoadContext; import accord.local.cfk.CommandsForKey; import accord.primitives.Ranges; +import accord.primitives.Routable; +import accord.primitives.Txn; import accord.primitives.TxnId; import accord.utils.async.AsyncChains; import org.apache.cassandra.distributed.Cluster; @@ -129,7 +132,8 @@ public class AccordDropTableBase extends TestBaseImpl inst.runOnInstance(() -> { TableId tableId = TableId.fromString(s); AccordService accord = (AccordService) AccordService.instance(); - PreLoadContext ctx = PreLoadContext.contextFor(Ranges.single(TokenRange.fullRange(tableId, getPartitioner())), LoadKeys.SYNC, READ_WRITE, "Test"); + TxnId syntheticTxnId = new TxnId(TxnId.MAX_EPOCH, 0, Txn.Kind.ExclusiveSyncPoint, Routable.Domain.Range, new Node.Id(1)); + PreLoadContext ctx = PreLoadContext.contextFor(syntheticTxnId, Ranges.single(TokenRange.fullRange(tableId, getPartitioner())), LoadKeys.SYNC, READ_WRITE, "Test"); CommandStores stores = accord.node().commandStores(); for (int storeId : stores.ids()) { diff --git a/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java b/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java index c567dd69c8..5129bff029 100644 --- a/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java +++ b/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java @@ -656,7 +656,6 @@ public class CommandsForKeySerializerTest @Override public AsyncChain build(PreLoadContext context, Consumer consumer) { return null; } @Override public AsyncChain build(PreLoadContext context, Function apply) { throw new UnsupportedOperationException(); } @Override public void shutdown() { } - @Override protected void registerTransitive(SafeCommandStore safeStore, RangeDeps deps){ } @Override public AsyncChain build(Callable task) { throw new UnsupportedOperationException(); } @Override public void onInconsistentTimestamp(Command command, Timestamp prev, Timestamp next) { throw new UnsupportedOperationException(); } @Override public void onFailedBootstrap(int attempts, String phase, Ranges ranges, Runnable retry, Throwable failure) { throw new UnsupportedOperationException(); }