diff --git a/modules/accord b/modules/accord index d5de79b56f..12da4693b4 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit d5de79b56f792d3b6d868ca46845b830e8906828 +Subproject commit 12da4693b449b30d0420673579a06025b7b3e484 diff --git a/src/java/org/apache/cassandra/service/accord/AccordCache.java b/src/java/org/apache/cassandra/service/accord/AccordCache.java index 839f31f1a6..c96a4e685c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCache.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCache.java @@ -52,6 +52,7 @@ import accord.utils.IntrusiveLinkedList; import accord.utils.Invariants; import accord.utils.QuadFunction; import accord.utils.TriFunction; +import accord.utils.UnhandledEnum; import org.agrona.collections.Object2ObjectHashMap; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.config.DatabaseDescriptor; @@ -255,7 +256,7 @@ public class AccordCache implements CacheSize Status status = node.status(); switch (status) { - default: throw new IllegalStateException("Unhandled status " + status); + default: throw new UnhandledEnum(status); case LOADING: node.loading().loading.cancel(); case WAITING_TO_LOAD: diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 95aba5de55..5806ae391f 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -308,7 +308,13 @@ public class AccordCommandStore extends CommandStore CommandsForKey loadCommandsForKey(RoutableKey key) { - return CommandsForKeyAccessor.load(id, (TokenKey) key); + CommandsForKey cfk = CommandsForKeyAccessor.load(id, (TokenKey) key); + if (cfk == null) + return null; + RedundantBefore.QuickBounds bounds = unsafeGetRedundantBefore().get(key); + if (bounds == null) + return cfk; // TODO (required): I don't think this should be possible? but we hit it on some test + return cfk.withRedundantBeforeAtLeast(bounds.gcBefore, false); } boolean validateCommandsForKey(RoutableKey key, CommandsForKey evicting) diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index ab5100a723..431c1639d3 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -36,8 +36,6 @@ import java.util.stream.Stream; import javax.annotation.Nullable; -import com.google.common.base.Functions; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -463,7 +461,7 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor void submit(BiConsumer sync, Function async, P1 p1) + private void submit(BiConsumer sync, Function async, P1 p1) { submit((e, c, p1a, p2a, p3) -> c.accept(e, p1a), (f, p1a, p2a, p3) -> f.apply(p1a), sync, async, p1, null, null); } - private void submit(TriConsumer sync, BiFunction async, P1 p1, P2 p2) + private void submit(TriConsumer sync, BiFunction async, P1 p1, P2 p2) { submit((e, c, p1a, p2a, p3) -> c.accept(e, p1a, p2a), (f, p1a, p2a, p3) -> f.apply(p1a, p2a), sync, async, p1, p2, null); } - private void submit(QuadConsumer sync, TriFunction async, P1 p1, P2 p2, P3 p3) + private void submit(QuadConsumer sync, TriFunction async, P1 p1, P2 p2, P3 p3) { submit((e, c, p1a, p2a, p3a) -> c.accept(e, p1a, p2a, p3a), TriFunction::apply, sync, async, p1, p2, p3); } - private void submit(QuintConsumer sync, QuadFunction async, P1 p1, P2 p2, P3 p3, P4 p4) + private void submit(QuintConsumer sync, QuadFunction async, P1 p1, P2 p2, P3 p3, P4 p4) { submit(sync, async, p1, p1, p2, p3, p4); } - abstract void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4); + abstract void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4); void submit(AccordTask operation) { - submit(AccordExecutor::submitExclusive, Function.identity(), operation); + submit(AccordExecutor::submitExclusive, i -> i, operation); } void submitExclusive(AccordTask task) @@ -775,7 +773,7 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor i, task); return task; } @@ -898,10 +896,16 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor extends SubmitAsync + private static class CancelAsync extends SubmitAsync { - final AccordTask cancel; + final Task cancel; - private CancelAsync(AccordTask cancel) + private CancelAsync(Task cancel) { this.cancel = cancel; } @@ -1195,7 +1202,7 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor { - SequentialExecutor executor = executor(); - TaskQueue queue = executor == null ? waitingToRun : executor; - if (queue.contains(this)) - { - queue.remove(this); - try { fail(new CancellationException()); } - catch (Throwable t) { agent.onUncaughtException(t); } - completeTaskExclusive(this, true); - } - }); + submit((e, c) -> c.cancelExclusive(e), CancelAsync::new, this); + } + + void cancelExclusive(AccordExecutor owner) + { + SequentialExecutor executor = executor(); + TaskQueue queue = executor == null ? waitingToRun : executor; + if (queue.contains(this)) + { + queue.remove(this); + completeTaskExclusive(this, true); + try { fail(new CancellationException()); } + catch (Throwable t) { agent.onUncaughtException(t); } + } } @Override diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java index 8348e14ee5..31583de8a5 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java @@ -33,7 +33,7 @@ import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_L abstract class AccordExecutorAbstractLockLoop extends AccordExecutor { - final ConcurrentLinkedStack submitted = new ConcurrentLinkedStack<>(); + final ConcurrentLinkedStack submitted = new ConcurrentLinkedStack<>(); boolean isHeldByExecutor; boolean shutdown; @@ -47,17 +47,17 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor abstract void awaitExclusive() throws InterruptedException; abstract AccordExecutorLoops loops(); abstract boolean isInLoop(); - abstract void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4); + abstract void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4); - void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) { // if we're a loop thread, we will poll the waitingToRun queue when we come around // NOTE: this assumes no synchronous blocking tasks are submitted to this executor - if (isInLoop()) submitted.push(async.apply(p1a, p2, p3, p4)); + if (isInLoop() || isOwningThread()) submitted.push(async.apply(p1a, p2, p3, p4)); else submitExternal(sync, async, p1s, p1a, p2, p3, p4); } - void submitExternalExclusive(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + void submitExternalExclusive(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) { try { diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java index e7f01e1cd6..eb0f56fe66 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java @@ -34,7 +34,7 @@ abstract class AccordExecutorAbstractSemiSyncSubmit extends AccordExecutorAbstra abstract void awaitExclusive() throws InterruptedException; - void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) { if (!lock.tryLock()) { diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java index bbeb165d50..202dcb5f86 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java @@ -111,7 +111,7 @@ class AccordExecutorSimple extends AccordExecutor } @Override - void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) { lock.lock(); try diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java index 3e94a21fac..89f836cd2c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java @@ -96,7 +96,7 @@ class AccordExecutorSyncSubmit extends AccordExecutorAbstractLockLoop hasWork.signal(); } - void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + void submitExternal(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) { lock.lock(); try diff --git a/src/java/org/apache/cassandra/service/accord/AccordTask.java b/src/java/org/apache/cassandra/service/accord/AccordTask.java index f04d4ad7c0..1153599ff1 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTask.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTask.java @@ -784,6 +784,11 @@ public abstract class AccordTask extends SubmittableTask implements Function< callback.accept(null, new CancellationException()); } + void cancelExclusive(AccordExecutor owner) + { + owner.cancel(this); + } + public State state() { return state; 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 b43b08fc20..93ef9457bd 100644 --- a/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java +++ b/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java @@ -246,13 +246,13 @@ public class AccordAgent implements Agent Shard shard = node.topology().forEpochIfKnown(homeKey, command.txnId().epoch()); // TODO (expected): make this a configurable calculation on normal request latencies (like ContentionStrategy) + long nowMicros = MILLISECONDS.toMicros(Clock.Global.currentTimeMillis()); long oneSecond = SECONDS.toMicros(1L); long promisedHlc = command.promised().hlc(); - if (promisedHlc == Long.MAX_VALUE) + if (promisedHlc > nowMicros + TimeUnit.MINUTES.toMicros(1)) promisedHlc = 0; long mostRecentStart = Math.max(command.txnId().hlc(), promisedHlc); long waitMicros = recover(txnId).computeWait(retryCount, MICROSECONDS); - long nowMicros = MILLISECONDS.toMicros(Clock.Global.currentTimeMillis()); if (mostRecentStart > nowMicros + SECONDS.toMicros(1L)) logger.warn("max({},{})>{}", command.txnId(), command.promised(), nowMicros); long startTime = mostRecentStart + waitMicros; diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java index 4dadcc6904..cec6ec9cb7 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java @@ -43,6 +43,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.CassandraRelevantProperties; +import org.apache.cassandra.db.commitlog.CommitLog; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.distributed.api.Feature; @@ -258,7 +259,8 @@ public class AccordLoadTest extends AccordTestBase try { i.acceptOnInstance(name -> { - AccordKeyspace.AccordColumnFamilyStores.commandsForKey.forceFlush(UNIT_TESTS); + if (CommitLog.instance.isStarted()) + AccordKeyspace.AccordColumnFamilyStores.commandsForKey.forceFlush(UNIT_TESTS); }, accordTableName); } catch (Throwable t)