From a0dc6f857b27c43d6264ac4ba35c709ba159e2fd Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Tue, 28 Apr 2026 23:21:24 +0100 Subject: [PATCH] Introduce AccordExecutorSignalLoop that aims to reduce lock contention: - consumer and producer threads wait signal without acquiring the lock first - lock owners may prepare more work than they need, moving some dynamically-adjusted portion of the work from the prioritised lock-managed structures onto a non-blocking queue, so that other threads may consume work from there; the portion is continually micro-adjusted to target some available work whenever the lock is acquired. - adopts/supports some features of the SEPExecutor: - threads may auto-adjust the number of running threads based on how much time is collectively spent waiting, to minimise time spent signalling - consumers may (timed-sleep) poll rather than await a signal, reducing the number of kernel interactions needed; relying on the auto-adjustment to bound the time the scheduler spends waking threads with no work to do Also Fix: - Client Result should not be persisted or included in Command object at any time patch by Benedict; reviewed by Alex Petrov and Ariel Weisberg for CASSANDRA-21375 --- .build/checkstyle.xml | 6 + modules/accord | 2 +- .../apache/cassandra/config/AccordConfig.java | 39 + .../org/apache/cassandra/config/Config.java | 1 - .../cassandra/config/DatabaseDescriptor.java | 46 +- .../service/accord/AccordCommandStore.java | 40 +- .../service/accord/AccordCommandStores.java | 61 +- .../service/accord/AccordExecutor.java | 231 ++--- .../AccordExecutorAbstractLockLoop.java | 62 +- .../accord/AccordExecutorAbstractLoop.java | 139 +++ .../AccordExecutorAbstractSemiSyncSubmit.java | 4 +- .../accord/AccordExecutorAsyncSubmit.java | 4 +- .../service/accord/AccordExecutorLoops.java | 6 +- .../accord/AccordExecutorSemiSyncSubmit.java | 4 +- .../accord/AccordExecutorSignalLoop.java | 472 +++++++++++ .../service/accord/AccordExecutorSimple.java | 10 +- .../accord/AccordExecutorSyncSubmit.java | 8 +- .../service/accord/AccordObjectSizes.java | 44 +- .../service/accord/AccordResult.java | 20 +- .../service/accord/AccordService.java | 12 +- .../cassandra/service/accord/AccordTask.java | 24 +- .../service/accord/debug/DebugExecution.java | 99 ++- .../accord/interop/AccordInteropApply.java | 10 +- .../accord/interop/AccordInteropPersist.java | 2 +- .../accord/serializers/ApplySerializers.java | 6 +- .../serializers/CheckStatusSerializers.java | 4 +- .../accord/serializers/ResultSerializers.java | 12 +- .../cassandra/service/accord/txn/TxnData.java | 2 +- .../service/accord/txn/TxnResult.java | 8 + .../utils/concurrent/SignalLock.java | 798 ++++++++++++++++++ .../test/accord/AccordLoadTest.java | 243 +++++- .../fuzz/topology/AccordBootstrapTest.java | 2 +- .../topology/AccordTopologyMixupTest.java | 2 +- .../HarryOnAccordTopologyMixupTest.java | 2 +- .../simulator/test/AccordExecutorTest.java | 164 ++++ .../db/virtual/AccordDebugKeyspaceTest.java | 2 +- .../cassandra/hints/HintsServiceTest.java | 2 +- .../accord/AccordCommandStoreTest.java | 4 +- .../service/accord/AccordTestUtils.java | 9 +- .../service/accord/EpochSyncTest.java | 4 +- .../CommandsForKeySerializerTest.java | 3 +- .../cassandra/utils/AccordGenerators.java | 4 +- 42 files changed, 2244 insertions(+), 373 deletions(-) create mode 100644 src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLoop.java create mode 100644 src/java/org/apache/cassandra/service/accord/AccordExecutorSignalLoop.java create mode 100644 src/java/org/apache/cassandra/utils/concurrent/SignalLock.java create mode 100644 test/simulator/test/org/apache/cassandra/simulator/test/AccordExecutorTest.java diff --git a/.build/checkstyle.xml b/.build/checkstyle.xml index 99faf7335e..8473f4c297 100644 --- a/.build/checkstyle.xml +++ b/.build/checkstyle.xml @@ -63,6 +63,12 @@ + + + + + + diff --git a/modules/accord b/modules/accord index e8c70e097f..a081e19e33 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit e8c70e097f305a5097084079f136d8be7c085b29 +Subproject commit a081e19e33ddeb91040f2ff8f0bef119f57d11fb diff --git a/src/java/org/apache/cassandra/config/AccordConfig.java b/src/java/org/apache/cassandra/config/AccordConfig.java index 7cf899a520..0f502ddc1f 100644 --- a/src/java/org/apache/cassandra/config/AccordConfig.java +++ b/src/java/org/apache/cassandra/config/AccordConfig.java @@ -66,11 +66,15 @@ public class AccordConfig /** * Same number of threads as queue shards, but the shard lock is held only while managing the queue, * so that submitting threads may queue load/save work. + * + * This is incompatible with QueueSubmissionModel.SIGNAL */ THREAD_PER_SHARD, /** * Same number of threads as shards, and the shard lock is held for the duration of serving requests. + * + * This is incompatible with QueueSubmissionModel.SIGNAL */ THREAD_PER_SHARD_SYNC_QUEUE, @@ -103,6 +107,14 @@ public class AccordConfig */ ASYNC, + /** + * Queue workers try to avoid competing for the lock, with the lock owner distributing work to any waiting threads + * and signalling them without them taking the lock + * + * NOTE: EXPERIMENTAL + */ + SIGNAL, + /** * The queue is backed by submission to a single-threaded plain executor. * This implementation does not honor the sharding model option. @@ -151,8 +163,30 @@ public class AccordConfig */ public volatile OptionaldPositiveInt queue_shard_count = OptionaldPositiveInt.UNDEFINED; + /** + * The total number of threads to share between queue shards + */ + public volatile OptionaldPositiveInt queue_thread_count = OptionaldPositiveInt.UNDEFINED; + public QueuePriorityModel queue_priority_model = HLC_FIFO; + /** + * If set, the signal loop does not match park/unpark pairs, but instead consumers perform timed-park spin waits + */ + public DurationSpec.LongMicrosecondsBound queue_spin_interval; + + /** + * If set, the signal loop reduces the number of threads it is using when the time spent parked exceeds real-time + * by this interval. + */ + public DurationSpec.LongMicrosecondsBound queue_stop_check_interval; + + /** + * If set, the signal loop reduces the number of threads it is using when the time spent parked exceeds real-time + * by this interval. + */ + public DurationSpec.LongMicrosecondsBound queue_signal_stop_check_interval_credit; + // yield to other executor threads after executing this many tasks in a row, if there are waiting threads and tasks public int queue_yield_interval = 100; @@ -456,4 +490,9 @@ public class AccordConfig return version.version; } } + + public int commandStoreShardCount() + { + return command_store_shard_count.or(DatabaseDescriptor::getAvailableProcessors); + } } diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index ba50692cfb..3a485fddc5 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -203,7 +203,6 @@ public class Config public int concurrent_reads = 32; public int concurrent_writes = 32; - public int concurrent_accord_operations = 32; public int concurrent_counter_writes = 32; public int concurrent_materialized_view_writes = 32; public OptionaldPositiveInt available_processors = new OptionaldPositiveInt(CASSANDRA_AVAILABLE_PROCESSORS.getInt(OptionaldPositiveInt.UNDEFINED_VALUE)); diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 24a8384b7a..9ee09db707 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -700,9 +700,6 @@ public class DatabaseDescriptor if (conf.concurrent_counter_writes < 2) throw new ConfigurationException("concurrent_counter_writes must be at least 2, but was " + conf.concurrent_counter_writes, false); - if (conf.concurrent_accord_operations < 1) - throw new ConfigurationException("concurrent_accord_operations must be at least 1, but was " + conf.concurrent_accord_operations, false); - if (conf.networking_cache_size == null) conf.networking_cache_size = new DataStorageSpec.IntMebibytesBound(Math.min(128, (int) (Runtime.getRuntime().maxMemory() / (16 * 1048576)))); @@ -2900,7 +2897,7 @@ public class DatabaseDescriptor public static int getAccordConcurrentOps() { - return conf.concurrent_accord_operations; + return conf.accord.queue_thread_count.or(2 * FBUtilities.getAvailableProcessors()); } public static void setConcurrentAccordOps(int concurrent_operations) @@ -2909,7 +2906,7 @@ public class DatabaseDescriptor { throw new IllegalArgumentException("Concurrent accord operations must be non-negative"); } - conf.concurrent_accord_operations = concurrent_operations; + conf.accord.queue_thread_count = new OptionaldPositiveInt(concurrent_operations); } public static int getFlushWriters() @@ -5635,6 +5632,7 @@ public class DatabaseDescriptor return conf.accord; } + // TODO (expected): move all getAccordX into AccordConfig public static AccordConfig.TransactionalRangeMigration getTransactionalRangeMigration() { return conf.accord.range_migration; @@ -5665,44 +5663,6 @@ public class DatabaseDescriptor conf.accord.enabled = b; } - public static AccordConfig.QueueShardModel getAccordQueueShardModel() - { - return conf.accord.queue_shard_model; - } - - public static AccordConfig.QueueSubmissionModel getAccordQueueSubmissionModel() - { - return conf.accord.queue_submission_model; - } - - public static int getAccordQueueShardCount() - { - switch (getAccordQueueShardModel()) - { - default: throw new AssertionError("Unhandled queue_shard_model: " + conf.accord.queue_shard_model); - case THREAD_PER_SHARD: - case THREAD_PER_SHARD_SYNC_QUEUE: - return conf.accord.queue_shard_count.or(DatabaseDescriptor::getAvailableProcessors); - case THREAD_POOL_PER_SHARD: - return conf.accord.queue_shard_count.or(DatabaseDescriptor.getAvailableProcessors()/4); - } - } - - public static int getAccordCommandStoreShardCount() - { - return conf.accord.command_store_shard_count.or(DatabaseDescriptor::getAvailableProcessors); - } - - public static int getAccordMaxQueuedLoadCount() - { - return conf.accord.max_queued_loads.or(getAccordConcurrentOps()); - } - - public static int getAccordMaxQueuedRangeLoadCount() - { - return conf.accord.max_queued_range_loads.or(Math.max(4, getAccordConcurrentOps() / 4)); - } - public static DefaultProgressLog.Config getAccordProgressLogConfig() { return accordProgressLogConfig; diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index d3676ef23b..1b13e1807f 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -245,26 +245,6 @@ public class AccordCommandStore extends CommandStore if (this.progressLog instanceof DefaultProgressLog) ((DefaultProgressLog)this.progressLog).unsafeSetConfig(DatabaseDescriptor.getAccordProgressLogConfig()); - final AccordCache.Type.Instance commands; - final AccordCache.Type.Instance commandsForKey; - try (AccordExecutor.ExclusiveGlobalCaches exclusive = sharedExecutor.lockCaches()) - { - commands = exclusive.commands.newInstance(this); - commandsForKey = exclusive.commandsForKey.newInstance(this); - this.caches = new ExclusiveCaches(sharedExecutor.unsafeLock(), exclusive.global, commands, commandsForKey); - } - - this.exclusiveExecutor = sharedExecutor.executor(id); - { - AccordConfig.RangeIndexMode mode = getAccord().range_index_mode; - switch (mode) - { - default: throw new UnhandledEnum(mode); - case journal_sai: rangeIndex = new JournalRangeIndex(this); break; - case in_memory: rangeIndex = new InMemoryRangeIndex(this); break; - } - } - maybeLoadRedundantBefore(journal.loadRedundantBefore(id())); maybeLoadBootstrapBeganAt(journal.loadBootstrapBeganAt(id())); maybeLoadSafeToRead(journal.loadSafeToRead(id())); @@ -284,6 +264,26 @@ public class AccordCommandStore extends CommandStore return a; }).orElseThrow(() -> Invariants.illegalState("CommandStore %d created with no ranges", id)); + final AccordCache.Type.Instance commands; + final AccordCache.Type.Instance commandsForKey; + try (AccordExecutor.ExclusiveGlobalCaches exclusive = sharedExecutor.lockCaches()) + { + commands = exclusive.commands.newInstance(this); + commandsForKey = exclusive.commandsForKey.newInstance(this); + this.caches = new ExclusiveCaches(sharedExecutor.unsafeLock(), exclusive.global, commands, commandsForKey); + } + this.exclusiveExecutor = sharedExecutor.executor(id); + + { + AccordConfig.RangeIndexMode mode = getAccord().range_index_mode; + switch (mode) + { + default: throw new UnhandledEnum(mode); + case journal_sai: rangeIndex = new JournalRangeIndex(this); break; + case in_memory: rangeIndex = new InMemoryRangeIndex(this); break; + } + } + if (AccordService.isStarted()) progressLog.unsafeStart(); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java index fa007bdab9..f362c960c4 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java @@ -48,16 +48,17 @@ import accord.utils.async.AsyncResults; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.concurrent.Shutdownable; +import org.apache.cassandra.config.AccordConfig; import org.apache.cassandra.config.AccordConfig.QueueShardModel; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.journal.Descriptor; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.accord.AccordCommandStore.DurablyAppliedTo; import org.apache.cassandra.service.accord.AccordExecutor.AccordExecutorFactory; +import org.apache.cassandra.utils.FBUtilities; import static org.apache.cassandra.config.AccordConfig.QueueShardModel.THREAD_PER_SHARD; -import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueShardCount; -import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueSubmissionModel; +import static org.apache.cassandra.config.DatabaseDescriptor.getAccord; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITHOUT_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.constant; @@ -83,10 +84,12 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu AccordCommandStore.factory(id -> executors[id % executors.length])); this.executors = executors; this.mask = Integer.highestOneBit(executors.length) - 1; + cacheSize = DatabaseDescriptor.getAccordCacheSizeInMiB() << 20; workingSetSize = DatabaseDescriptor.getAccordWorkingSetSizeInMiB() << 20; - maxQueuedLoads = DatabaseDescriptor.getAccordMaxQueuedLoadCount(); - maxQueuedRangeLoads = DatabaseDescriptor.getAccordMaxQueuedRangeLoadCount(); + AccordConfig config = DatabaseDescriptor.getAccord(); + maxQueuedLoads = maxQueuedLoads(config); + maxQueuedRangeLoads = maxQueuedRangeLoads(config); shrinkingOn = DatabaseDescriptor.getAccordCacheShrinkingOn(); refreshCapacities(); ScheduledExecutors.scheduledFastTasks.scheduleWithFixedDelay(() -> { @@ -102,15 +105,21 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu static Factory factory() { return (NodeCommandStoreService time, Agent agent, DataStore store, RandomSource random, Journal journal, ShardDistributor shardDistributor, ProgressLog.Factory progressLogFactory, LocalListeners.Factory listenersFactory) -> { - AccordExecutor[] executors = new AccordExecutor[getAccordQueueShardCount()]; + AccordConfig config = getAccord(); + AccordExecutor[] executors = new AccordExecutor[executorShards(config)]; AccordExecutorFactory factory; int maxThreads = Integer.MAX_VALUE; - switch (getAccordQueueSubmissionModel()) + switch (config.queue_submission_model) { - default: throw new AssertionError("Unhandled QueueSubmissionModel: " + getAccordQueueSubmissionModel()); + default: throw new AssertionError("Unhandled QueueSubmissionModel: " + config.queue_submission_model); case SYNC: factory = AccordExecutorSyncSubmit::new; break; case SEMI_SYNC: factory = AccordExecutorSemiSyncSubmit::new; break; case ASYNC: factory = AccordExecutorAsyncSubmit::new; break; + case SIGNAL: + long spinIntervalNanos = config.queue_spin_interval == null ? 0 : config.queue_spin_interval.to(TimeUnit.NANOSECONDS); + long stopCheckIntervalNanos = config.queue_stop_check_interval == null ? 0 : config.queue_stop_check_interval.to(TimeUnit.NANOSECONDS); + factory = (executorId, mode, threads, name, agent0) -> new AccordExecutorSignalLoop(executorId, mode, threads, spinIntervalNanos, stopCheckIntervalNanos, TimeUnit.NANOSECONDS, name, agent0); + break; case EXEC_ST: factory = AccordExecutorSimple::new; maxThreads = 1; @@ -119,9 +128,8 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu for (int id = 0; id < executors.length; id++) { - QueueShardModel shardModel = DatabaseDescriptor.getAccordQueueShardModel(); + QueueShardModel shardModel = config.queue_shard_model; String baseName = AccordExecutor.class.getSimpleName() + '[' + id; - int threads = Math.min(maxThreads, Math.max(DatabaseDescriptor.getAccordConcurrentOps() / getAccordQueueShardCount(), 1)); switch (shardModel) { case THREAD_PER_SHARD: @@ -129,6 +137,7 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu executors[id] = factory.get(id, shardModel == THREAD_PER_SHARD ? RUN_WITHOUT_LOCK : RUN_WITH_LOCK, 1, constant(baseName + ']'), agent); break; case THREAD_POOL_PER_SHARD: + int threads = Math.min(maxThreads, Math.max(DatabaseDescriptor.getAccordConcurrentOps() / executors.length, 1)); executors[id] = factory.get(id, RUN_WITHOUT_LOCK, threads, num -> baseName + ',' + num + ']', agent); break; } @@ -204,8 +213,8 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu { long capacityPerExecutor = cacheSize / executors.length; long workingSetPerExecutor = workingSetSize < 0 ? Long.MAX_VALUE : workingSetSize / executors.length; - int maxLoadsPerExecutor = (maxQueuedLoads + executors.length - 1) / executors.length; - int maxRangeLoadsPerExecutor = (maxQueuedRangeLoads + executors.length - 1) / executors.length; + int maxLoadsPerExecutor = Math.max(1, (maxQueuedLoads + executors.length - 1) / executors.length); + int maxRangeLoadsPerExecutor = Math.max(1, (maxQueuedRangeLoads + executors.length - 1) / executors.length); for (AccordExecutor executor : executors) { executor.executeDirectlyWithLock(() -> { @@ -319,4 +328,34 @@ public class AccordCommandStores extends CommandStores implements CacheSize, Shu } return AsyncChains.allOf(chains); } + + private static int executorShards(AccordConfig config) + { + switch (config.queue_shard_model) + { + default: throw new AssertionError("Unhandled queue_shard_model: " + config.queue_shard_model); + case THREAD_PER_SHARD: + case THREAD_PER_SHARD_SYNC_QUEUE: + return config.queue_shard_count.or(DatabaseDescriptor::getAvailableProcessors); + case THREAD_POOL_PER_SHARD: + return Math.max(1, config.queue_shard_count.or(DatabaseDescriptor.getAvailableProcessors() / 8)); + } + } + + private static int threads(AccordConfig config) + { + return config.queue_thread_count.or(2 * FBUtilities.getAvailableProcessors()); + } + + public static int maxQueuedLoads(AccordConfig config) + { + return config.max_queued_loads.or(FBUtilities.getAvailableProcessors()); + } + + public static int maxQueuedRangeLoads(AccordConfig config) + { + return config.max_queued_range_loads.or(maxQueuedLoads(config) / 4); + } + + } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index 1d70803fe6..b1aff08f13 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -200,7 +200,6 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor running = new TaskQueue<>(RUNNING); private final TaskQueue> waitingToLoadRangeTxns = new TaskQueue<>(WAITING_TO_LOAD); - private final TaskQueue> waitingToLoad = new TaskQueue<>(WAITING_TO_LOAD); private final TaskQueue waitingToRun = new TaskQueue<>(WAITING_TO_RUN); @@ -233,7 +232,6 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor task) { task.unqueueIfQueued(); @@ -606,13 +605,13 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor task) { task.onWaitingToRun(); task.addToQueue(task.commandStore.exclusiveExecutor); } - private void waitingToRun(SubmittableTask task, @Nullable SequentialExecutor queue) + private void waitingToRun(Task task, @Nullable SequentialExecutor queue) { task.onWaitingToRun(); task.addToQueue(queue == null ? waitingToRun : queue); @@ -638,7 +637,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) { @@ -867,7 +866,6 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor= 1, "Must permit at least one load"); + Invariants.requireArgument(range >= 1, "Must permit at least one range load"); maxQueuedLoads = total; maxQueuedRangeLoads = range; } @@ -1112,25 +1107,71 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor ownerUpdater = AtomicReferenceFieldUpdater.newUpdater(SequentialExecutor.class, Thread.class, "owner"); @@ -1232,7 +1271,6 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor IntFunction constant(O out) @@ -1614,12 +1641,12 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor out; } - abstract class Plain extends SubmittableTask implements Cancellable + abstract class Plain extends Task implements Cancellable { abstract SequentialExecutor executor(); @Override - protected void preRunExclusive(Thread assigned) {} + protected void preRunExclusive() {} @Override protected final void addToQueue(TaskQueue queue) @@ -1631,7 +1658,7 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor c.cancelExclusive(e), CancelAsync::new, this); + submit((e, c) -> c.cancelExclusive(e), CancelTask::new, this); } void cancelExclusive(AccordExecutor owner) @@ -1681,14 +1708,14 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor submitted = new ConcurrentLinkedStack<>(); + int runningThreads; boolean shutdown; AccordExecutorAbstractLockLoop(Lock lock, int executorId, Agent agent) @@ -49,14 +46,13 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor abstract void notifyWorkExclusive(); void loopYieldExclusive() throws InterruptedException {} abstract void awaitExclusive() throws InterruptedException; - abstract AccordExecutorLoops loops(); - 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() || isOwningThread()) submitted.push(async.apply(p1a, p2, p3, p4)); + if (isInLoop() || isOwningThread()) push(async.apply(p1a, p2, p3, p4)); else submitExternal(sync, async, p1s, p1a, p2, p3, p4); } @@ -66,7 +62,7 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor { try { - drainSubmittedExclusive(); + drainUnqueuedExclusive(); } catch (Throwable t) { @@ -82,33 +78,6 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor } } - public boolean hasTasks() - { - if (tasks > 0 || !submitted.isEmpty() || runningThreads > 0) - return true; - - lock(); - try - { - return tasks > 0 || !submitted.isEmpty() || runningThreads > 0; - } - finally - { - unlock(); - } - } - - final void updateWaitingToRunExclusive() - { - drainSubmittedExclusive(); - super.updateWaitingToRunExclusive(); - } - - final void drainSubmittedExclusive() - { - submitted.drain(AccordExecutor::consumeExclusive, this, true); - } - final void notifyIfMoreWorkExclusive() { if (hasWaitingToRun()) @@ -149,7 +118,7 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor ++runningThreads; } - LoopTask task(String name, Mode mode) + LoopTask task(int index, String name, Mode mode) { return mode == RUN_WITH_LOCK ? runWithLock(name) : runWithoutLock(name); } @@ -178,7 +147,7 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor setRunning(task); try { - task.preRunExclusive(self); + task.preRunExclusive(); task.runInternal(); } catch (Throwable t) @@ -263,7 +232,7 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor if (task != null) { setRunning(task); - task.preRunExclusive(self); + task.preRunExclusive(); if (DEBUG_EXECUTION) debug.onExitLock(); exitLockLoop(); break; @@ -337,23 +306,10 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor }; } - @Override - public Stream active() - { - return loops().active(); - } - @Override public void shutdown() { shutdown = true; notifyWork(); } - - @Override - public Object shutdownNow() - { - shutdown(); - return null; - } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLoop.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLoop.java new file mode 100644 index 0000000000..a76d1bac18 --- /dev/null +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLoop.java @@ -0,0 +1,139 @@ +/* + * 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; + +import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; +import java.util.concurrent.locks.Lock; +import java.util.stream.Stream; + +import accord.api.Agent; +import accord.utils.Invariants; + +import org.apache.cassandra.concurrent.DebuggableTask.DebuggableTaskRunner; + +abstract class AccordExecutorAbstractLoop extends AccordExecutor +{ + private volatile Task unqueued; + private static final AtomicReferenceFieldUpdater unqueuedUpdater = AtomicReferenceFieldUpdater.newUpdater(AccordExecutorAbstractLoop.class, Task.class, "unqueued"); + + AccordExecutorAbstractLoop(Lock lock, int executorId, Agent agent) + { + super(lock, executorId, agent); + } + + abstract AccordExecutorLoops loops(); + + boolean hasUnqueued() + { + return unqueued != null; + } + + Task unqueued() + { + return unqueued; + } + + final Task push(Task submit) + { + Invariants.require(submit.next == null); + while (true) + { + Task next = unqueued; + submit.next = next; + if (unqueuedUpdater.compareAndSet(this, next, submit)) + return next; + } + } + + @Override + public boolean hasTasks() + { + if (hasUnqueued() || tasks > 0) + return true; + + lock(); + try + { + return hasUnqueued() || tasks > 0; + } + finally + { + unlock(); + } + } + + final void updateWaitingToRunExclusive() + { + drainUnqueuedExclusive(); + super.updateWaitingToRunExclusive(); + } + + final void drainUnqueuedExclusive() + { + Task cur = Task.reverse(acquireUnqueuedExclusive()); + while (cur != null) + cur = enqueueOneExclusive(cur); + } + + final Task acquireUnqueuedExclusive() + { + return unqueuedUpdater.getAndSet(this, null); + } + + final Task enqueueOneExclusive(Task cur) + { + Invariants.require(cur != null); + Task next = cur.next; + cur.next = null; + if (cur.isReadyToCleanup()) completeTaskExclusive(cur); + else cur.submitExclusive(this); + return next; + } + + final Task enqueueOneCleanup(Task cur) + { + Invariants.require(cur != null); + Task next = cur.next; + cur.next = null; + completeTaskExclusive(cur); + return next; + } + + final Task enqueueOneSubmit(Task cur) + { + Invariants.require(cur != null); + Task next = cur.next; + cur.next = null; + cur.submitExclusive(this); + return next; + } + + @Override + public Stream active() + { + return loops().active(); + } + + @Override + public Object shutdownNow() + { + shutdown(); + return null; + } +} diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java index f10caab42f..dce92f103d 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java @@ -33,9 +33,9 @@ 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 (submitted.push(async.apply(p1a, p2, p3, p4)) && !isInLoop()) + if (push(async.apply(p1a, p2, p3, p4)) == null && !isInLoop()) notifyWork(); } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java index b212ccc32c..a2c57a4556 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java @@ -26,7 +26,7 @@ import accord.api.Agent; import org.apache.cassandra.utils.concurrent.LockWithAsyncSignal; // WARNING: experimental - needs more testing -class AccordExecutorAsyncSubmit extends AccordExecutorAbstractSemiSyncSubmit +public class AccordExecutorAsyncSubmit extends AccordExecutorAbstractSemiSyncSubmit { private final AccordExecutorLoops loops; private final LockWithAsyncSignal lock; @@ -47,7 +47,7 @@ class AccordExecutorAsyncSubmit extends AccordExecutorAbstractSemiSyncSubmit void awaitExclusive() throws InterruptedException { lock.clearSignal(); - if (submitted.isEmpty()) + if (!hasUnqueued()) lock.await(); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java index 17baa1b122..3eddaaf9e2 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java @@ -20,13 +20,13 @@ package org.apache.cassandra.service.accord; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; -import java.util.function.BiFunction; import java.util.function.IntFunction; import java.util.stream.Stream; import org.agrona.collections.Long2ObjectHashMap; import accord.utils.Invariants; +import accord.utils.TriFunction; import org.apache.cassandra.concurrent.DebuggableTask.DebuggableTaskRunner; import org.apache.cassandra.service.accord.AccordExecutor.Mode; @@ -54,7 +54,7 @@ class AccordExecutorLoops private final AtomicInteger running = new AtomicInteger(); private final Condition terminated = Condition.newOneTimeCondition(); - public AccordExecutorLoops(Mode mode, int threads, IntFunction loopName, BiFunction loopFactory) + public AccordExecutorLoops(Mode mode, int threads, IntFunction loopName, TriFunction loopFactory) { Invariants.require(mode == RUN_WITH_LOCK ? threads == 1 : threads >= 1); running.addAndGet(threads); @@ -63,7 +63,7 @@ class AccordExecutorLoops for (int i = 0; i < threads; ++i) { String name = loopName.apply(i); - LoopTask task = loopFactory.apply(name, mode); + LoopTask task = loopFactory.apply(i, name, mode); Thread thread = executorFactory().startThread(name, wrap(task), NON_DAEMON, INFINITE_LOOP); Thread conflict = loops.putIfAbsent(thread.getId(), thread); Invariants.require(conflict == null || !conflict.isAlive(), "Allocated two threads with the same threadId!"); diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java index f3e694c1c9..f97e206ef5 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java @@ -26,7 +26,7 @@ import java.util.function.IntFunction; import accord.api.Agent; // WARNING: experimental - needs more testing -class AccordExecutorSemiSyncSubmit extends AccordExecutorAbstractSemiSyncSubmit +public class AccordExecutorSemiSyncSubmit extends AccordExecutorAbstractSemiSyncSubmit { private final AccordExecutorLoops loops; private final ReentrantLock lock; @@ -61,7 +61,7 @@ class AccordExecutorSemiSyncSubmit extends AccordExecutorAbstractSemiSyncSubmit @Override void awaitExclusive() throws InterruptedException { - if (submitted.isEmpty()) + if (!hasUnqueued()) awaitWork(); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSignalLoop.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSignalLoop.java new file mode 100644 index 0000000000..034544de4e --- /dev/null +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSignalLoop.java @@ -0,0 +1,472 @@ +/* + * 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; + +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.TimeUnit; +import java.util.function.IntFunction; + +import javax.annotation.Nullable; + +import accord.api.Agent; +import accord.utils.Invariants; +import accord.utils.QuadFunction; +import accord.utils.QuintConsumer; + +import org.apache.cassandra.service.accord.debug.DebugExecution.DebugTask; +import org.apache.cassandra.utils.concurrent.SignalLock; + +import static accord.utils.Invariants.nonNull; +import static org.apache.cassandra.service.accord.debug.DebugExecution.DEBUG_EXECUTION; + +public class AccordExecutorSignalLoop extends AccordExecutorAbstractLoop +{ + private static class ShutdownException extends RuntimeException {} + + private final SignalLock lock; + private final AccordExecutorLoops loops; + private int readyToRunTarget = 1; + private final int readyToRunLimit; + // TODO (desired): intrusive queue using Task.next, but a little challenging because we reuse SequentialQueueTask so have ABA problem + private final ConcurrentLinkedQueue readyToRun = new ConcurrentLinkedQueue<>(); + + private Task pendingSequentialHead, pendingSequentialTail; + private Task pendingCleanupHead, pendingCleanupTail; + private Task pendingNewHead, pendingNewTail; + private int pendingCount; + + private boolean shutdown; + + public AccordExecutorSignalLoop(int executorId, Mode mode, int threads, long spinInterval, long stopCheckInterval, TimeUnit units, IntFunction name, Agent agent) + { + this(new SignalLock(threads, spinInterval, stopCheckInterval, units), executorId, mode, threads, name, agent); + } + + public AccordExecutorSignalLoop(SignalLock lock, int executorId, Mode mode, int threads, IntFunction name, Agent agent) + { + super(lock, executorId, agent); + Invariants.require(threads < SignalLock.MAX_THREADS); + this.lock = lock; + this.loops = new AccordExecutorLoops(mode, threads, name, this::task); + this.readyToRunLimit = Math.min(threads * 4, SignalLock.MAX_SIGNAL_COUNT); + } + + void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4) + { + Task next = async.apply(p1a, p2, p3, p4); + Task prev = push(next); + if (prev == null) + lock.signalLockWork(); + } + + @Override + final void beforeUnlockExternal() + { + } + + private Task pollReadyToRun() + { + return readyToRun.poll(); + } + + private void addReadyToRun(Task task) + { + readyToRun.add(task); + } + + private boolean hasReadyToRun() + { + return !readyToRun.isEmpty(); + } + + private LoopTask task(int index, String id, Mode mode) + { + return new LoopTask(index, id); + } + + class LoopTask extends AccordExecutorLoops.LoopTask + { + final int index; + + LoopTask(int index, String id) + { + super(id); + this.index = index; + } + + private Task awaitWork() + { + while (true) + { + if (lock.awaitAsyncOrLock(index)) + { + try + { + if (DEBUG_EXECUTION) debug.onEnterLock(); + fetchWorkExclusive(); + } + catch (Throwable t) + { + unlock(); + throw t; + } + + if (!unlockAndAcquire()) + continue; + } + + if (shutdown) + throw new ShutdownException(); + + return nonNull(pollReadyToRun()); + } + } + + private Task cleanupAndMaybeGetWork(@Nullable Task cleanup) + { + if (cleanup == null) + return null; + + if (lock.tryAcquireAsyncWork()) + { + if (shutdown) + throw new ShutdownException(); + + return pushCleanupAndReturn(cleanup, nonNull(pollReadyToRun())); + } + + if (!tryLock()) + return pushCleanupAndReturn(cleanup, null); + + try + { + completeTaskExclusive(cleanup); + clearRunning(); + fetchWorkExclusive(); + } + catch (Throwable t) + { + unlock(); + throw t; + } + + if (unlockAndAcquire()) + { + if (shutdown) + throw new ShutdownException(); + + return nonNull(pollReadyToRun()); + } + return null; + } + + private Task pushCleanupAndReturn(Task cleanup, Task result) + { + clearRunning(); + cleanup.setReadyToCleanup(); + if (push(cleanup) == null) + lock.signalLockWork(); + return result; + } + + final boolean tryLock() + { + return onTryLock(lock.tryLock(index)); + } + + final boolean unlockAndAcquire() + { + if (Invariants.isParanoid()) paranoidUnlockExclusive(); + if (DEBUG_EXECUTION) debug.onExitLock(); + return lock.unlockAndAcquireAsyncWork(); + } + + private boolean enqueueOnePending() + { + if (pendingCount == 0) + return false; + + if (pendingSequentialHead != null) + { + pendingSequentialHead = enqueueOneCleanup(pendingSequentialHead); + if (pendingSequentialHead == null) + pendingSequentialTail = null; + } + else if (pendingCleanupHead != null) + { + pendingCleanupHead = enqueueOneCleanup(pendingCleanupHead); + if (pendingCleanupHead == null) + pendingCleanupTail = null; + } + else + { + pendingNewHead = enqueueOneSubmit(pendingNewHead); + if (pendingNewHead == null) + pendingNewTail = null; + } + --pendingCount; + return true; + } + + private void fetchWorkExclusive() + { + boolean hasReadyToRun = hasReadyToRun(); + boolean hadWaitingToRun = hasAlreadyWaitingToRun(); + { + int prevPendingUnqueued = pendingCount; + updateAndEnqueuePendingUntilHasWaitingToRun(prevPendingUnqueued / 2); + } + + if (hadWaitingToRun && hasReadyToRun) + { + lock.addAndGetEnabledThreadCount(1); + } + else if (hadWaitingToRun) + { + readyToRunTarget = Math.min(readyToRunLimit, readyToRunTarget + (1+readyToRunTarget)/2); + } + else if (hasReadyToRun && readyToRunTarget > 1) + { + --readyToRunTarget; + } + + boolean hasDrainedSignal = false; + while (true) + { + long state = lock.state(); + int signals = SignalLock.asyncSignalCount(state); + int waiters = SignalLock.waitingEnabledThreadCount(state); + if (signals >= readyToRunTarget) + { + if (enqueueOnePending() || (updatePendingUnqueued() && enqueueOnePending())) continue; + else if (hasDrainedSignal) + lock.signalLockWorkExclusive(); + return; + } + else if (waiters > 0 && signals > 1 && SignalLock.activeEnabledThreadCount(state) == 1) + { + // ensure at least one other thread is running if there's enough work for it; it will spin up other threads if necessary + lock.propagateAsyncWorkSignals(1); + } + + Task task = pollAlreadyWaitingToRunExclusive(); + if (task == null) + { + if (updateAndEnqueuePendingUntilHasWaitingToRun(0)) + continue; + + lock.clearLockWork(); + hasDrainedSignal = true; + if (!updateAndEnqueuePendingUntilHasWaitingToRun(0)) + { + if (tasks == 0) + notifyQuiescentExclusive(); + return; + } + } + else + { + try { task.preRunExclusive(); } + catch (Throwable t) + { + try { task.fail(t); } + catch (Throwable t2) { try { t.addSuppressed(t2); } catch (Throwable t3) {} } + try { completeTaskExclusive(task); } + catch (Throwable t2) { try { t.addSuppressed(t2); } catch (Throwable t3) {} } + continue; + } + if (DEBUG_EXECUTION) DebugTask.get(task).onPreRun(); + + addReadyToRun(task); + boolean incremented = lock.incrementAsyncWork(false); + Invariants.require(incremented); + } + } + } + + private boolean updatePendingUnqueued() + { + if (!hasUnqueued()) + return false; + + int count = 0; + Task addSequentialHead = null, addSequentialTail = null; + Task addCleanupHead = null, addCleanupTail = null; + Task addNewHead = null, addNewTail = null; + { + Task cur = Task.reverse(acquireUnqueuedExclusive()); + while (cur != null) + { + Task next = cur.next; + if (!cur.isReadyToCleanup()) + { + if (addNewHead == null) addNewHead = addNewTail = setNextNull(cur); + else addNewHead = reverseOne(addNewHead, cur); + } + else if (cur instanceof SequentialQueueTask) + { + if (addSequentialHead == null) addSequentialHead = addSequentialTail = setNextNull(cur); + else addSequentialHead = reverseOne(addSequentialHead, cur); + } + else + { + if (addCleanupHead == null) addCleanupHead = addCleanupTail = setNextNull(cur); + else addCleanupHead = reverseOne(addCleanupHead, cur); + } + ++count; + cur = next; + } + } + + pendingCount += count; + if (addSequentialHead != null) + { + if (pendingSequentialHead == null) pendingSequentialHead = addSequentialHead; + else pendingSequentialTail.next = addSequentialHead; + pendingSequentialTail = addSequentialTail; + } + if (addCleanupHead != null) + { + if (pendingCleanupHead == null) pendingCleanupHead = addCleanupHead; + else pendingCleanupTail.next = addCleanupHead; + pendingCleanupTail = addCleanupTail; + } + if (addNewHead != null) + { + if (pendingNewHead == null) pendingNewHead = addNewHead; + else pendingNewTail.next = addNewHead; + pendingNewTail = addNewTail; + } + return true; + } + + private Task reverseOne(Task prev, Task cur) + { + cur.next = prev; + return cur; + } + + private Task setNextNull(Task cur) + { + cur.next = null; + return cur; + } + + private boolean updateAndEnqueuePendingUntilHasWaitingToRun(int processAtLeast) + { + updatePendingUnqueued(); + int count = 0; + while (enqueueOnePending()) + { + if (++count >= processAtLeast && hasAlreadyWaitingToRun()) + return true; + } + return hasAlreadyWaitingToRun(); + } + + @Override + public void run() + { + lock.register(index, Thread.currentThread()); + Task task = null; + while (true) + { + try + { + try { task = cleanupAndMaybeGetWork(task); } + catch (Throwable t) { task = null; throw t; } + if (task == null) + task = awaitWork(); + + try + { + setRunning(task); + task.runInternal(); + } + catch (Throwable t) + { + try { task.fail(t); } + catch (Throwable t2) + { + try + { + t2.addSuppressed(t); + agent.onException(t2); + } + catch (Throwable t3) { /* empty to ensure we definitely loop so we cleanup the task */ } + } + } + } + catch (ShutdownException ignore) + { + break; + } + catch (Throwable t) + { + agent.onException(t); + } + } + } + } + + @Override + public void shutdown() + { + lock.lock(); + try + { + shutdown = true; + lock.signalAllRegistered(); + } + finally + { + lock.unlock(); + } + } + + @Override + AccordExecutorLoops loops() + { + return loops; + } + + @Override + boolean isInLoop() + { + return loops.isInLoop(); + } + + @Override + boolean isOwningThread() + { + return lock.isOwner(); + } + + @Override + public boolean isTerminated() + { + return loops.isTerminated(); + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException + { + return loops.awaitTermination(timeout, unit); + } +} diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java index 349c9be620..4e89f068ab 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java @@ -35,7 +35,6 @@ class AccordExecutorSimple extends AccordExecutor { final ExecutorPlus executor; final ReentrantLock lock; - private Task active; public AccordExecutorSimple(int executorId, String name, Agent agent) { @@ -87,30 +86,25 @@ class AccordExecutorSimple extends AccordExecutor lock.lock(); try { - runningThreads = 1; while (true) { Task task = pollWaitingToRunExclusive(); - active = task; if (task == null) { - runningThreads = 0; notifyQuiescentExclusive(); return; } - try { task.preRunExclusive(self); task.runInternal(); } + try { task.preRunExclusive(); task.runInternal(); } catch (Throwable t) { task.fail(t); } finally { completeTaskExclusive(task); - active = null; } } } finally { - runningThreads = 0; if (hasWaitingToRun()) executor.execute(this::run); lock.unlock(); @@ -118,7 +112,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 6547048f67..0d852ccde3 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java @@ -27,7 +27,7 @@ import accord.api.Agent; import accord.utils.QuadFunction; import accord.utils.QuintConsumer; -class AccordExecutorSyncSubmit extends AccordExecutorAbstractLockLoop +public class AccordExecutorSyncSubmit extends AccordExecutorAbstractLockLoop { private final AccordExecutorLoops loops; private final ReentrantLock lock; @@ -89,16 +89,16 @@ 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(); + lock(); try { submitExternalExclusive(sync, p1s, p2, p3, p4); } finally { - lock.unlock(); + unlock(); } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java b/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java index 8f6f374986..ddef3232f3 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java +++ b/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java @@ -22,7 +22,6 @@ import java.util.UUID; import java.util.function.ToLongFunction; import accord.api.Key; -import accord.api.Result; import accord.api.RoutingKey; import accord.local.Command; import accord.local.CommandBuilder; @@ -39,6 +38,7 @@ import accord.primitives.FullKeyRoute; import accord.primitives.FullRangeRoute; import accord.primitives.KeyDeps; import accord.primitives.Keys; +import accord.primitives.Known; import accord.primitives.PartialKeyRoute; import accord.primitives.PartialRangeRoute; import accord.primitives.PartialTxn; @@ -50,6 +50,7 @@ import accord.primitives.RoutingKeys; import accord.primitives.SaveStatus; import accord.primitives.Seekable; import accord.primitives.Seekables; +import accord.primitives.Status; import accord.primitives.Timestamp; import accord.primitives.Txn.Kind; import accord.primitives.TxnId; @@ -67,10 +68,8 @@ import org.apache.cassandra.service.accord.api.TokenKey; import org.apache.cassandra.service.accord.serializers.ResultSerializers; import org.apache.cassandra.service.accord.serializers.TableMetadatasAndKeys; import org.apache.cassandra.service.accord.txn.AccordUpdate; -import org.apache.cassandra.service.accord.txn.TxnData; import org.apache.cassandra.service.accord.txn.TxnQuery; import org.apache.cassandra.service.accord.txn.TxnRead; -import org.apache.cassandra.service.accord.txn.TxnResult; import org.apache.cassandra.service.accord.txn.TxnWrite; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.ObjectSizes; @@ -80,6 +79,7 @@ import static accord.local.cfk.CommandsForKey.InternalStatus.ACCEPTED; import static accord.primitives.SaveStatus.Invalidated; import static accord.primitives.SaveStatus.NotDefined; import static accord.primitives.SaveStatus.PreAccepted; +import static accord.primitives.SaveStatus.Stable; import static accord.primitives.SaveStatus.TruncatedUnapplied; import static accord.primitives.Status.Durability.NotDurable; import static accord.primitives.TxnId.NO_TXNIDS; @@ -293,20 +293,13 @@ public class AccordObjectSizes return size; } - public static long results(Result result) - { - if (result == ResultSerializers.APPLIED) - return 0; - return ((TxnResult) result).estimatedSizeOnHeap(); - } - private static class CommandEmptySizes { private final static PartitionKey EMPTY_KEY = new PartitionKey(EMPTY_ID, new BufferDecoratedKey(new Murmur3Partitioner.LongToken(1), ByteBufferUtil.EMPTY_BYTE_BUFFER)); private final static TokenKey EMPTY_TOKEN_KEY = new TokenKey(EMPTY_ID, new Murmur3Partitioner.LongToken(1)); private final static TxnId EMPTY_TXNID = new TxnId(42, 42, 0, Kind.Read, Domain.Key, new Node.Id(42)); - private static Command build(SaveStatus saveStatus, boolean hasDeps, boolean hasTxn, boolean executes) + private static Command build(SaveStatus saveStatus, boolean executes) { Keys keys = Keys.of(EMPTY_KEY); FullKeyRoute route = new FullKeyRoute(EMPTY_TOKEN_KEY, new RoutingKey[]{ EMPTY_TOKEN_KEY }); @@ -315,30 +308,30 @@ public class AccordObjectSizes .durability(NotDurable) .executeAt(EMPTY_TXNID) .promised(Ballot.ZERO); - if (hasDeps) + if (saveStatus.known.hasAnyDeps()) builder.partialDeps(new Deps(KeyDeps.none(route.toParticipants()), RangeDeps.NONE).intersecting(route)); - if (hasTxn) + if (saveStatus.known.isDefinitionKnown()) builder.partialTxn(new PartialTxn.InMemory(Kind.Read, keys, TxnRead.empty(Domain.Key), null, null, TableMetadatasAndKeys.none(Domain.Key))); - if (executes) - { + if (saveStatus.compareTo(Stable) >= 0 && !saveStatus.hasBeen(Status.Truncated)) builder.waitingOn(WaitingOn.empty(Domain.Key)); - builder.result(new TxnData()); - } + + if (saveStatus.known.is(Known.Outcome.Apply)) + builder.result(ResultSerializers.APPLIED); return builder.build(saveStatus); } - final static long NOT_DEFINED = measure(build(NotDefined, false, false, false)); - final static long PREACCEPTED = measure(build(PreAccepted, false, true, false)); - final static long NOTACCEPTED = measure(build(SaveStatus.AcceptedInvalidate, false, false, false)); - final static long ACCEPTED = measure(build(SaveStatus.AcceptedMedium, true, false, false)); - final static long COMMITTED = measure(build(SaveStatus.Committed, true, true, false)); - final static long EXECUTED = measure(build(SaveStatus.Applied, true, true, true)); + final static long NOT_DEFINED = measure(build(NotDefined, false)); + final static long PREACCEPTED = measure(build(PreAccepted, false)); + final static long NOTACCEPTED = measure(build(SaveStatus.AcceptedInvalidate, false)); + final static long ACCEPTED = measure(build(SaveStatus.AcceptedMedium, false)); + final static long COMMITTED = measure(build(SaveStatus.Committed, false)); + final static long EXECUTED = measure(build(SaveStatus.Applied, true)); // TODO (expected): TruncatedAwaitsOnlyDeps - final static long TRUNCATED = measure(build(TruncatedUnapplied, false, false, false).participants()); - final static long INVALIDATED = measure(build(Invalidated, false, false, false).participants()); + final static long TRUNCATED = measure(build(TruncatedUnapplied, false).participants()); + final static long INVALIDATED = measure(build(Invalidated, false).participants()); private static void touch() {} @@ -407,7 +400,6 @@ public class AccordObjectSizes size += sizeNullable(command.partialDeps(), AccordObjectSizes::dependencies); size += sizeNullable(command.acceptedOrCommitted(), AccordObjectSizes::timestamp); size += sizeNullable(command.writes(), AccordObjectSizes::writes); - size += sizeNullable(command.result(), AccordObjectSizes::results); size += sizeNullable(command.waitingOn(), AccordObjectSizes::waitingOn); return size; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordResult.java b/src/java/org/apache/cassandra/service/accord/AccordResult.java index 013af99b87..033ccdd88a 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordResult.java +++ b/src/java/org/apache/cassandra/service/accord/AccordResult.java @@ -66,8 +66,9 @@ public class AccordResult extends AsyncFuture implements BiConsumer keysOrRanges, RequestBookkeeping bookkeeping, long startedAtNanos, long deadlineAtNanos, boolean isTxnRequest) + public AccordResult(@Nullable TxnId txnId, Seekables keysOrRanges, RequestBookkeeping bookkeeping, long startedAtNanos, long deadlineAtNanos, boolean isTxnRequest, @Nullable accord.api.Tracing tracing) { this.txnId = txnId; this.keysOrRanges = keysOrRanges; @@ -75,6 +76,7 @@ public class AccordResult extends AsyncFuture implements BiConsumer extends AsyncFuture implements BiConsumer extends AsyncFuture implements BiConsumer dataStore, - new KeyspaceSplitter(new EvenSplit<>(getAccordCommandStoreShardCount(), getPartitioner().accordSplitter())), + new KeyspaceSplitter(new EvenSplit<>(getAccord().commandStoreShardCount(), getPartitioner().accordSplitter())), agent, new DefaultRandom(), scheduler, @@ -1025,7 +1026,7 @@ public class AccordService implements IAccordService, Shutdownable public static V getBlocking(AsyncChain async, @Nullable TxnId txnId, Seekables keysOrRanges, RequestBookkeeping bookkeeping, long startedAt, long deadline, boolean isTxnRequest) { - AccordResult result = new AccordResult<>(txnId, keysOrRanges, bookkeeping, startedAt, deadline, isTxnRequest); + AccordResult result = new AccordResult<>(txnId, keysOrRanges, bookkeeping, startedAt, deadline, isTxnRequest, null); async.begin(result); return result.awaitAndGet(); } @@ -1037,7 +1038,7 @@ public class AccordService implements IAccordService, Shutdownable public static V getBlocking(AsyncResult async, @Nullable TxnId txnId, Seekables keysOrRanges, RequestBookkeeping bookkeeping, long startedAt, long deadline, boolean isTxnRequest) { - AccordResult result = new AccordResult<>(txnId, keysOrRanges, bookkeeping, startedAt, deadline, isTxnRequest); + AccordResult result = new AccordResult<>(txnId, keysOrRanges, bookkeeping, startedAt, deadline, isTxnRequest, null); async.invoke(result); return result.awaitAndGet(); } @@ -1142,7 +1143,8 @@ public class AccordService implements IAccordService, Shutdownable ClientRequestBookkeeping bookkeeping = txn.isWrite() ? accordWriteBookkeeping : accordReadBookkeeping; bookkeeping.metrics.keySize.update(txn.keys().size()); long deadlineNanos = requestTime.computeDeadline(timeout); - AccordResult result = new AccordResult<>(txnId, txn.keys(), bookkeeping, requestTime.startedAtNanos(), deadlineNanos, true); + Tracing tracing = agent().tracing().trace(txnId, txn.keys(), Client); + AccordResult result = new AccordResult<>(txnId, txn.keys(), bookkeeping, requestTime.startedAtNanos(), deadlineNanos, true, tracing); node.coordinate(txnId, txn, minEpoch, deadlineNanos).begin((BiConsumer) result); return result; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordTask.java b/src/java/org/apache/cassandra/service/accord/AccordTask.java index 5ca534d4d6..e3fe9656e0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTask.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTask.java @@ -70,10 +70,11 @@ import org.apache.cassandra.concurrent.DebuggableTask; import org.apache.cassandra.metrics.LogLinearDecayingHistograms; import org.apache.cassandra.service.accord.AccordCacheEntry.Status; import org.apache.cassandra.service.accord.AccordCommandStore.Caches; -import org.apache.cassandra.service.accord.AccordExecutor.SubmittableTask; +import org.apache.cassandra.service.accord.AccordExecutor.Task; import org.apache.cassandra.service.accord.AccordExecutor.TaskQueue; import org.apache.cassandra.service.accord.AccordKeyspace.CommandsForKeyAccessor; import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.service.accord.debug.DebugExecution.DebugTask; import org.apache.cassandra.service.accord.serializers.CommandSerializers; import org.apache.cassandra.tracing.Tracing; import org.apache.cassandra.utils.Clock; @@ -100,9 +101,8 @@ import static org.apache.cassandra.service.accord.AccordTask.State.WAITING_TO_RU import static org.apache.cassandra.service.accord.AccordTask.State.WAITING_TO_SCAN_RANGES; import static org.apache.cassandra.service.accord.debug.DebugExecution.DEBUG_EXECUTION; import static org.apache.cassandra.service.accord.debug.DebugExecution.DebugTask.SANITY_CHECK; -import static org.apache.cassandra.utils.Clock.Global.nanoTime; -public abstract class AccordTask extends SubmittableTask implements Function, Cancellable, DebuggableTask +public abstract class AccordTask extends Task implements Function, Cancellable, DebuggableTask { private static final Logger logger = LoggerFactory.getLogger(AccordTask.class); private static final NoSpamLogger noSpamLogger = NoSpamLogger.getLogger(logger, 1, TimeUnit.MINUTES); @@ -629,6 +629,7 @@ public abstract class AccordTask extends SubmittableTask implements Function< { if (SANITY_CHECK) { + DebugTask debug = DebugTask.get(this); if (debug.sanityCheck == null) debug.sanityCheck = new ArrayList<>(commands.size()); debug.sanityCheck.add(safeCommand.current()); @@ -637,13 +638,13 @@ public abstract class AccordTask extends SubmittableTask implements Function< private void save(List diffs, Runnable onFlush) { - if (SANITY_CHECK && debug.sanityCheck != null) + if (SANITY_CHECK && DebugTask.get(this).sanityCheck != null) { Condition condition = Condition.newOneTimeCondition(); this.commandStore.appendCommands(diffs, condition::signal); condition.awaitUninterruptibly(); - for (Command check : debug.sanityCheck) + for (Command check : DebugTask.get(this).sanityCheck) this.commandStore.sanityCheckCommand(commandStore.unsafeGetRedundantBefore(), check); if (onFlush != null) onFlush.run(); @@ -655,7 +656,7 @@ public abstract class AccordTask extends SubmittableTask implements Function< } @Override - protected void preRunExclusive(Thread assigned) + protected void preRunExclusive() { state(ASSIGNED); queued = null; @@ -673,10 +674,10 @@ public abstract class AccordTask extends SubmittableTask implements Function< @Override public void runInternal() { - runningAt = nanoTime(); + onRunning(); logger.trace("Running {} with state {}", this, state); AccordSafeCommandStore safeStore = null; - try (Closeable close = locals.get()) + try (Closeable close = resources.get()) { if (Tracing.isTracing()) Tracing.trace(preLoadContext.describe()); @@ -717,7 +718,7 @@ public abstract class AccordTask extends SubmittableTask implements Function< commandStore.complete(safeStore); safeStore = null; - if (DEBUG_EXECUTION) debug.onRunComplete(); + onRunComplete(); if (!flush) finish(result, null); } @@ -757,9 +758,9 @@ public abstract class AccordTask extends SubmittableTask implements Function< @Override protected void cleanupExclusive(AccordExecutor executor) { - super.cleanupExclusive(executor); Invariants.expect(state.isExecuted()); releaseResources(commandStore.cachesExclusive()); + super.cleanupExclusive(executor); executor.keys.increment(commandsForKey == null ? 0 : commandsForKey.size(), runningAt); if (histogramBuffer != null) { @@ -825,18 +826,21 @@ public abstract class AccordTask extends SubmittableTask implements Function< { rangeScanner.cleanup(caches); rangeScanner = null; + if (DEBUG_EXECUTION) DebugTask.get(this).onReleasedRangeScanner(); } if (commands != null) { commands.forEach((k, v) -> caches.commands().release(v, this)); commands.clear(); commands = null; + if (DEBUG_EXECUTION) DebugTask.get(this).onReleasedCommands(); } if (commandsForKey != null) { commandsForKey.forEach((k, v) -> caches.commandsForKeys().release(v, this)); commandsForKey.clear(); commandsForKey = null; + if (DEBUG_EXECUTION) DebugTask.get(this).onReleasedCommandsForKeys(); } if (waitingToLoad != null) { diff --git a/src/java/org/apache/cassandra/service/accord/debug/DebugExecution.java b/src/java/org/apache/cassandra/service/accord/debug/DebugExecution.java index a6aa88cd0e..d42f9bb338 100644 --- a/src/java/org/apache/cassandra/service/accord/debug/DebugExecution.java +++ b/src/java/org/apache/cassandra/service/accord/debug/DebugExecution.java @@ -29,7 +29,9 @@ import accord.local.Command; import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.metrics.LogLinearHistogram; -import org.apache.cassandra.service.accord.AccordExecutor; +import org.apache.cassandra.service.accord.AccordExecutor.Task; +import org.apache.cassandra.utils.Closeable; +import org.apache.cassandra.utils.WithResources; import static org.apache.cassandra.config.CassandraRelevantProperties.DTEST_ACCORD_JOURNAL_SANITY_CHECK_ENABLED; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -38,9 +40,10 @@ public class DebugExecution { private static final Logger logger = LoggerFactory.getLogger(DebugExecution.class); public static final boolean DEBUG_EXECUTION = CassandraRelevantProperties.ACCORD_DEBUG_EXECUTION.getBoolean(false); - private static final long REPORT_MIN_LATENCY_MICROS = 10_000; + private static final long REPORT_MIN_LATENCY_MICROS = 20_000; private static final long REPORT_CPU_RATIO = 2; private static final long REPORT_MAX_LATENCY_MICROS = 50_000; + private static final long REPORT_CPU_MICROS = 10_000; // TODO (expected): use sharded histogram so we can report global stats public static class DebugExecutor @@ -54,8 +57,8 @@ public class DebugExecution final LogLinearHistogram sequentialExecutorWaitingToRunLatency = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); final LogLinearHistogram sequentialExecutorSetHeadToRunLatency = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); final LogLinearHistogram pollToRun = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); - final LogLinearHistogram applying = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); final LogLinearHistogram running = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); + final LogLinearHistogram runToCleanup = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); final LogLinearHistogram cleanup = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); final LogLinearHistogram taskTotal = new LogLinearHistogram(REPORT_MAX_LATENCY_MICROS); @@ -84,6 +87,7 @@ public class DebugExecution public void onExitLock() { + // TODO (expected): specialise this for despatch loop, which can reasonably handle longer lock hold periods unlockedAt = nanoTime(); unlockedAtCpu = nowCpu(); long lockedForMicros = (unlockedAt - lockedAt)/1000; @@ -138,7 +142,7 @@ public class DebugExecution final int commandStoreId; long setTaskAt, waitingAt; - AccordExecutor.Task prev; + Task prev; public DebugSequentialExecutor(DebugExecutor owner, int commandStoreId) { @@ -146,13 +150,13 @@ public class DebugExecution this.commandStoreId = commandStoreId; } - public void onSetTask(AccordExecutor.Task next) + public void onSetTask(Task next) { if (next == null) setTaskAt = 0; else setTaskAt = nanoTime(); } - public void onComplete(AccordExecutor.Task completed) + public void onComplete(Task completed) { long readyAt = setTaskAt; if (waitingAt > setTaskAt) @@ -178,55 +182,102 @@ public class DebugExecution } } - public static class DebugTask + public static final class DebugTask implements WithResources { public static final boolean SANITY_CHECK = DTEST_ACCORD_JOURNAL_SANITY_CHECK_ENABLED.getBoolean(); private static final boolean DEBUG = DEBUG_EXECUTION || SANITY_CHECK; - public static DebugTask maybeDebug() { return DEBUG ? new DebugTask() : null; } + + public static WithResources maybeDebug(WithResources resources, Task task) { return DEBUG ? new DebugTask(resources, task) : resources; } + public static DebugTask get(Task task) { return (DebugTask)task.unwrap().resources; } + + final WithResources resources; + final Task task; + public DebugTask(WithResources resources, Task task) + { + this.resources = resources; + this.task = task; + } public List sanityCheck; // for AccordTask only - long polledAt, appliedAt, completedAt; - long polledAtCpu, completedAtCpu; + long polledAt, preRunAt, runCompleteAt, completedAt; + long releasedRangeScannerAt, releasedCommandsAt, releasedCommandsForKeyAt; + long runningAtCpu, runCompleteAtCpu; + Thread thread; public void onPolled() { polledAt = nanoTime(); - polledAtCpu = ManagementFactory.getThreadMXBean().getCurrentThreadCpuTime(); + } + + public void onPreRun() + { + preRunAt = nanoTime(); + } + + public void onRunning() + { + thread = Thread.currentThread(); + runningAtCpu = nowCpu(); } public void onRunComplete() { - appliedAt = nanoTime(); + runCompleteAtCpu = nowCpu(); + runCompleteAt = nanoTime(); } - public void onCompleted(AccordExecutor.Task task, DebugExecutor owner) + public void onReleasedRangeScanner() + { + releasedRangeScannerAt = nanoTime(); + } + + public void onReleasedCommands() + { + releasedCommandsAt = nanoTime(); + } + + public void onReleasedCommandsForKeys() + { + releasedCommandsForKeyAt = nanoTime(); + } + + public void onCompleted(DebugExecutor owner) { completedAt = nanoTime(); - completedAtCpu = ManagementFactory.getThreadMXBean().getCurrentThreadCpuTime(); if (task.runningAt > 0 && polledAt > 0) { long pollToRunMicros = (task.runningAt - polledAt)/1000; owner.pollToRun.increment(pollToRunMicros); - long applyingMicros = -1; - if (appliedAt > 0) + long runningMicros = -1; + if (runCompleteAt > 0) { - applyingMicros = (appliedAt - task.runningAt)/1000; - owner.applying.increment(applyingMicros); + runningMicros = (runCompleteAt - task.runningAt) / 1000; + owner.running.increment(runningMicros); } - long runningMicros = (task.cleanupAt - task.runningAt)/1000; - owner.running.increment(runningMicros); + long runToCleanMicros = (task.cleanupAt - runCompleteAt)/1000; + owner.runToCleanup.increment(runToCleanMicros); long cleanupMicros = (completedAt - task.cleanupAt)/1000; owner.cleanup.increment(cleanupMicros); long totalMicros = (completedAt - polledAt)/1000; owner.taskTotal.increment(totalMicros); - long totalCpu = (completedAtCpu - polledAtCpu)/1000; - if (totalMicros > REPORT_MAX_LATENCY_MICROS || (totalMicros > REPORT_MIN_LATENCY_MICROS && (totalMicros/totalCpu) >= REPORT_CPU_RATIO)) + long totalCpu = (runCompleteAtCpu - runningAtCpu)/1000; + if (totalMicros > REPORT_MAX_LATENCY_MICROS || totalCpu > REPORT_CPU_MICROS || (totalMicros > REPORT_MIN_LATENCY_MICROS && (totalCpu == 0 || totalMicros/totalCpu >= REPORT_CPU_RATIO))) { - report("{}: total {}us {}cpu, running {}us{}, cleanup {}us, pollToRun {}us", task, totalMicros, totalCpu, - runningMicros, (applyingMicros >= 0 ? ", applying " + applyingMicros + "us" : ""), cleanupMicros, pollToRunMicros); + String reason = ""; + if (totalMicros > REPORT_MAX_LATENCY_MICROS) reason += "LONG TIME "; + if (totalCpu > REPORT_CPU_MICROS) reason += "HIGH CPU "; + if ((totalMicros > REPORT_MIN_LATENCY_MICROS && (totalMicros/totalCpu) >= REPORT_CPU_RATIO)) reason += "LOW RATIO "; + report("{}{}: total {}us cpu:{}us ({}), pollToRun {}us, running {}us, runToClean {}us, cleanup {}us", + reason, task, totalMicros, totalCpu, thread, pollToRunMicros, runningMicros, runToCleanMicros, cleanupMicros); } } } + + @Override + public Closeable get() + { + return resources.get(); + } } private static void report(String message, Object ... params) diff --git a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java index 5f09d68b51..3aaaf7dd27 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java @@ -23,7 +23,7 @@ import javax.annotation.Nullable; import org.agrona.collections.Int2ObjectHashMap; import accord.api.LocalListeners; -import accord.api.Result; +import accord.api.Result.PersistableResult; import accord.coordinate.ExecuteFlag.ExecuteFlags; import accord.local.Command; import accord.local.Node.Id; @@ -70,7 +70,7 @@ public class AccordInteropApply extends Apply implements LocalListeners.ComplexL public static final Apply.Factory FACTORY = new Apply.Factory() { @Override - public Apply create(Kind kind, Id to, Topologies participates, TxnId txnId, Ballot ballot, Route route, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, ExecuteFlags flags) + public Apply create(Kind kind, Id to, Topologies participates, TxnId txnId, Ballot ballot, Route route, Txn txn, Timestamp executeAt, Deps deps, Writes writes, PersistableResult result, FullRoute fullRoute, ExecuteFlags flags) { checkArgument(kind != Kind.Maximal || to.equals(AccordService.nodeId()), "Shouldn't need to send a maximal commit with interop support"); ConsistencyLevel commitCL = txn.update() instanceof AccordUpdate ? ((AccordUpdate) txn.update()).cassandraCommitCL() : null; @@ -84,7 +84,7 @@ public class AccordInteropApply extends Apply implements LocalListeners.ComplexL public static final IVersionedSerializer serializer = new ApplySerializer() { @Override - protected AccordInteropApply deserializeApply(TxnId txnId, Ballot ballot, Route scope, long minEpoch, long waitForEpoch, long maxEpoch, Apply.Kind kind, Timestamp executeAt, PartialDeps deps, PartialTxn txn, @Nullable FullRoute fullRoute, Writes writes, Result result, ExecuteFlags flags) + protected AccordInteropApply deserializeApply(TxnId txnId, Ballot ballot, Route scope, long minEpoch, long waitForEpoch, long maxEpoch, Apply.Kind kind, Timestamp executeAt, PartialDeps deps, PartialTxn txn, @Nullable FullRoute fullRoute, Writes writes, PersistableResult result, ExecuteFlags flags) { return new AccordInteropApply(kind, txnId, ballot, scope, minEpoch, waitForEpoch, maxEpoch, executeAt, deps, txn, fullRoute, writes, result, flags); } @@ -94,12 +94,12 @@ public class AccordInteropApply extends Apply implements LocalListeners.ComplexL transient Int2ObjectHashMap listeners; boolean failed; - private AccordInteropApply(Kind kind, TxnId txnId, Ballot ballot, Route route, long minEpoch, long waitForEpoch, long maxEpoch, Timestamp executeAt, PartialDeps deps, @Nullable PartialTxn txn, @Nullable FullRoute fullRoute, Writes writes, Result result, ExecuteFlags flags) + private AccordInteropApply(Kind kind, TxnId txnId, Ballot ballot, Route route, long minEpoch, long waitForEpoch, long maxEpoch, Timestamp executeAt, PartialDeps deps, @Nullable PartialTxn txn, @Nullable FullRoute fullRoute, Writes writes, PersistableResult result, ExecuteFlags flags) { super(kind, txnId, ballot, route, minEpoch, waitForEpoch, maxEpoch, executeAt, deps, txn, fullRoute, writes, result, flags); } - private AccordInteropApply(Kind kind, Id to, Topologies participates, TxnId txnId, Ballot ballot, Route route, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, ExecuteFlags flags) + private AccordInteropApply(Kind kind, Id to, Topologies participates, TxnId txnId, Ballot ballot, Route route, Txn txn, Timestamp executeAt, Deps deps, Writes writes, PersistableResult result, FullRoute fullRoute, ExecuteFlags flags) { super(kind, to, participates, txnId, ballot, route, txn, executeAt, deps, writes, result, fullRoute, flags); } diff --git a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java index 9be2857898..7cdb7032a8 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropPersist.java @@ -55,7 +55,7 @@ public class AccordInteropPersist extends Persist { public AccordInteropPersist(Node node, SequentialAsyncExecutor executor, Topologies topologies, TxnId txnId, Route sendTo, Ballot ballot, Txn txn, Timestamp executeAt, Deps deps, Writes writes, Result result, FullRoute fullRoute, ConsistencyLevel consistencyLevel, ExecuteFlag.CoordinationFlags flags, boolean informDurableOnDone, Apply.Kind applyKind, BiConsumer callback) { - super(node, executor, topologies, txnId, ballot, sendTo, txn, executeAt, deps, writes, result, fullRoute, flags, informDurableOnDone, AccordInteropApply.FACTORY, applyKind, consistencyLevel == ALL ? AllTracker::new : QuorumTracker::new, (ignore, fail) -> { + super(node, executor, topologies, txnId, ballot, sendTo, txn, executeAt, deps, writes, result.toPersistable(), fullRoute, flags, informDurableOnDone, AccordInteropApply.FACTORY, applyKind, consistencyLevel == ALL ? AllTracker::new : QuorumTracker::new, (ignore, fail) -> { if (fail != null) callback.accept(null, fail); else callback.accept(result, null); }); 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 29704e48cb..291acace92 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java @@ -20,7 +20,7 @@ package org.apache.cassandra.service.accord.serializers; import java.io.IOException; -import accord.api.Result; +import accord.api.Result.PersistableResult; import accord.coordinate.ExecuteFlag.ExecuteFlags; import accord.messages.Apply; import accord.messages.Apply.ApplyReply; @@ -64,7 +64,7 @@ public class ApplySerializers } protected abstract A deserializeApply(TxnId txnId, Ballot ballot, Route scope, long minEpoch, long waitForEpoch, long maxEpoch, Apply.Kind kind, - Timestamp executeAt, PartialDeps deps, PartialTxn txn, FullRoute fullRoute, Writes writes, Result result, ExecuteFlags flags); + Timestamp executeAt, PartialDeps deps, PartialTxn txn, FullRoute fullRoute, Writes writes, PersistableResult result, ExecuteFlags flags); @Override public A deserializeBody(DataInputPlus in, Version version, TxnId txnId, Route scope, long waitForEpoch) throws IOException @@ -104,7 +104,7 @@ public class ApplySerializers { @Override protected Apply deserializeApply(TxnId txnId, Ballot ballot, Route scope, long minEpoch, long waitForEpoch, long maxEpoch, Apply.Kind kind, - Timestamp executeAt, PartialDeps deps, PartialTxn txn, FullRoute fullRoute, Writes writes, Result result, ExecuteFlags flags) + Timestamp executeAt, PartialDeps deps, PartialTxn txn, FullRoute fullRoute, Writes writes, PersistableResult result, ExecuteFlags flags) { return Apply.SerializationSupport.create(txnId, ballot, scope, minEpoch, waitForEpoch, maxEpoch, kind, executeAt, deps, txn, fullRoute, writes, result, flags); } 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 5fe95408fe..c3131bd4ba 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/CheckStatusSerializers.java @@ -21,7 +21,7 @@ package org.apache.cassandra.service.accord.serializers; import java.io.IOException; import java.util.Objects; -import accord.api.Result; +import accord.api.Result.PersistableResult; import accord.api.RoutingKey; import accord.coordinate.Infer; import accord.messages.CheckStatus; @@ -228,7 +228,7 @@ public class CheckStatusSerializers PartialDeps committedDeps = DepsSerializers.nullablePartialDeps.deserialize(in); Writes writes = CommandSerializers.nullableWrites.deserialize(in, version); - Result result = null; + PersistableResult result = null; if (maxKnowledgeStatus.known.outcome().isOrWasApply()) result = ResultSerializers.APPLIED; diff --git a/src/java/org/apache/cassandra/service/accord/serializers/ResultSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/ResultSerializers.java index be9e68cb6e..2a9e36b255 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/ResultSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/ResultSerializers.java @@ -18,7 +18,7 @@ package org.apache.cassandra.service.accord.serializers; -import accord.api.Result; +import accord.api.Result.PersistableResult; import org.apache.cassandra.io.UnversionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; @@ -27,17 +27,17 @@ import org.apache.cassandra.io.util.DataOutputPlus; public class ResultSerializers { // TODO (desired): this is meant to encode e.g. whether the transaction's condition met or not for clients to later query - public static final Result APPLIED = new Result(){}; + public static final PersistableResult APPLIED = new PersistableResult(){}; - public static final UnversionedSerializer result = new UnversionedSerializer<>() + public static final UnversionedSerializer result = new UnversionedSerializer<>() { - public void serialize(Result t, DataOutputPlus out) { } - public Result deserialize(DataInputPlus in) + public void serialize(PersistableResult t, DataOutputPlus out) { } + public PersistableResult deserialize(DataInputPlus in) { return APPLIED; } - public long serializedSize(Result t) + public long serializedSize(PersistableResult t) { return 0; } diff --git a/src/java/org/apache/cassandra/service/accord/txn/TxnData.java b/src/java/org/apache/cassandra/service/accord/txn/TxnData.java index 868143393a..2791ecce30 100644 --- a/src/java/org/apache/cassandra/service/accord/txn/TxnData.java +++ b/src/java/org/apache/cassandra/service/accord/txn/TxnData.java @@ -116,7 +116,7 @@ public class TxnData extends Int2ObjectHashMap implements TxnResul private TxnData(int size) { - super(size, 0.65f); + super(size, 0.65f, false); } public static TxnData of(int key, TxnDataValue value) diff --git a/src/java/org/apache/cassandra/service/accord/txn/TxnResult.java b/src/java/org/apache/cassandra/service/accord/txn/TxnResult.java index b3bf2b3fc8..10d97d6f22 100644 --- a/src/java/org/apache/cassandra/service/accord/txn/TxnResult.java +++ b/src/java/org/apache/cassandra/service/accord/txn/TxnResult.java @@ -20,6 +20,8 @@ package org.apache.cassandra.service.accord.txn; import accord.api.Result; +import static org.apache.cassandra.service.accord.serializers.ResultSerializers.APPLIED; + public interface TxnResult extends Result { enum Kind @@ -40,4 +42,10 @@ public interface TxnResult extends Result Kind kind(); long estimatedSizeOnHeap(); + + default PersistableResult toPersistable() + { + // TODO (required): should we persist the Kind? + return APPLIED; + } } diff --git a/src/java/org/apache/cassandra/utils/concurrent/SignalLock.java b/src/java/org/apache/cassandra/utils/concurrent/SignalLock.java new file mode 100644 index 0000000000..bced18c198 --- /dev/null +++ b/src/java/org/apache/cassandra/utils/concurrent/SignalLock.java @@ -0,0 +1,798 @@ +/* + * 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.utils.concurrent; + +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicLongArray; +import java.util.concurrent.atomic.AtomicLongFieldUpdater; +import java.util.concurrent.locks.Condition; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.LockSupport; + +import accord.utils.Invariants; + +import org.apache.cassandra.utils.Nemesis; + +import static org.apache.cassandra.utils.Clock.Global.nanoTime; + +/** + * A lock that supports semi-asynchronous operation; a thread may take the lock synchronously, + * or it may await a signal that it should acquire the lock (to perform synchronous work), + * or a signal that work has been made available for it outside the lock. + * + * In some systems it may be that LockSupport.unpark() often results in the calling thread becoming the callee, + * which can be harmful if the caller owns the lock. To avoid this problem, we support making available permits + * that can be acquired asynchronously, and only on releasing the lock do we attempt to actually directly wake + * any threads. In this mode of operation it is up to the user to unlock when there is sufficient work to signal. + * + * TODO (expected): once we support a high enough JDK, port to compareAndExchange/weakCompareAndSet as appropriate + */ +public final class SignalLock implements Lock +{ + private static final AtomicLongFieldUpdater stateUpdater = AtomicLongFieldUpdater.newUpdater(SignalLock.class, "state"); + + private static final long HAS_LOCK_WORK = 0x8000000000000000L; + // if there is a LOCK_SIGNAL, to avoid unnecessary wakeups we pick a waiting thread to unpark and store it in LOCK_SIGNALLED_MASK bits + private static final int LOCK_THREAD_SHIFT = 64 - 7; + private static final int LOCK_OWNED_SHIFT = LOCK_THREAD_SHIFT - 1; + private static final int ASYNC_SIGNAL_COUNT_SHIFT = LOCK_OWNED_SHIFT - 6; + private static final int ENABLED_THREAD_COUNT_SHIFT = ASYNC_SIGNAL_COUNT_SHIFT - 6; + public static final int MAX_THREADS = ENABLED_THREAD_COUNT_SHIFT; + public static final int THREAD_ID_MASK = 0x3f; + public static final int MAX_SIGNAL_COUNT = THREAD_ID_MASK; + private static final long LOCK_THREAD_MASK = (long) THREAD_ID_MASK << LOCK_THREAD_SHIFT; + private static final long LOCK_OWNED = 1L << LOCK_OWNED_SHIFT; + private static final long LOCK_SIGNALLED = 0L; + private static final long ASYNC_SIGNAL_INCREMENT = 1L << ASYNC_SIGNAL_COUNT_SHIFT; + private static final long ENABLED_THREAD_COUNT_INCREMENT = 1L << ENABLED_THREAD_COUNT_SHIFT; + private static final long ENABLED_THREAD_COUNT_MASK = (long)THREAD_ID_MASK << ENABLED_THREAD_COUNT_SHIFT; + private static final long WAITING_THREADS_MASK = -1L >>> (64 - MAX_THREADS); + + private final boolean allowStealing; + private final Thread[] registered; + @Nemesis private volatile long state; + private Thread owner; + private final ConcurrentLinkedQueue locking = new ConcurrentLinkedQueue<>(); + private final AtomicLongArray stopLog; + private final AtomicLong stopCheck; + private final long stopCheckIntervalNanos; + private final long spinIntervalNanos; + private boolean preferRegistered; + private boolean autoPreferRegistered = true; + + private int depth; + + public SignalLock(int threadCount) + { + this(true, threadCount, 0, 0, null); + } + + public SignalLock(int threadCount, long spinInterval, long stopCheckInterval, TimeUnit units) + { + this(true, threadCount, spinInterval, stopCheckInterval, units); + } + + public SignalLock(boolean allowStealing, int threadCount, long spinInterval, long stopCheckInterval, TimeUnit units) + { + Invariants.requireArgument(threadCount <= MAX_THREADS, "Supports at most " + MAX_THREADS + " registered threads"); + this.allowStealing = allowStealing; + this.registered = new Thread[threadCount]; + this.state = (long)threadCount << ENABLED_THREAD_COUNT_SHIFT; + this.spinIntervalNanos = spinInterval <= 0 ? 0 : units.toNanos(spinInterval); + if (stopCheckInterval > 0) + { + this.stopLog = new AtomicLongArray(threadCount); + this.stopCheck = new AtomicLong(); + this.stopCheckIntervalNanos = units.toNanos(stopCheckInterval); + } + else + { + this.stopLog = null; + this.stopCheck = null; + this.stopCheckIntervalNanos = -1; + } + } + + public void register(int index, Thread thread) + { + registered[index] = thread; + } + + public Thread registeredThread(int index) + { + return registered[index]; + } + + public void lock() + { + Thread self = Thread.currentThread(); + if (tryLock(self)) + return; + + locking.add(self); + while (true) + { + long cur = state; + if (isLockAvailable(cur, ANONYMOUS_OWNER, self)) + { + long upd = setLockThread(clearLockThread(cur), ANONYMOUS_OWNER, LOCK_OWNED); + if (!stateUpdater.compareAndSet(this, cur, upd)) + continue; + + Invariants.require(depth == 0); + Invariants.require(owner == null); + depth = 1; + owner = self; + + Invariants.require(locking.peek() == self); + locking.poll(); + return; + } + + LockSupport.park(); + } + } + + @Override + public boolean tryLock() + { + return tryLock(Thread.currentThread()); + } + + private boolean tryLock(Thread self) + { + return tryLock(ANONYMOUS_OWNER, self); + } + + public boolean tryLock(int thread) + { + Invariants.require(thread >= 0 && thread < registered.length); + Thread self = registeredThread(thread); + return tryLock(thread, self); + } + + private boolean tryLock(int threadOrAnonymous, Thread self) + { + if (owner != self) + { + while (true) + { + long cur = state; + if (!isLockAvailable(cur, threadOrAnonymous, null)) + return false; + + long upd = setLockThread(clearLockThread(cur), threadOrAnonymous, LOCK_OWNED); + if (!stateUpdater.compareAndSet(this, cur, upd)) + continue; + + Invariants.require(depth == 0); + Invariants.require(owner == null); + owner = self; + break; + } + } + + ++depth; + return true; + } + + /** + * Wait for the lock to be signalled and acquire it, or wait for additional work to be signalled. + * Return true if we acquired the lock. + */ + public boolean awaitAsyncOrLock(int thread) + { + return awaitAsyncOrLock(thread, 1); + } + + public boolean awaitAsyncOrLock(int thread, int burstSignalOnAcquire) + { + long threadBit = threadBit(thread); + Thread self = registered[thread]; + Invariants.require(self == Thread.currentThread()); + Invariants.require(owner != self); + Invariants.require((state & threadBit) == 0); + + boolean hasPaused = false; + if (spinIntervalNanos > 0) + { + while (true) + { + long cur = state; + if (hasLockWork(cur) && isLockAvailable(cur, thread, null)) + { + if (tryTakeLockInAwaitLoop(cur, thread, self, hasPaused)) return true; + else continue; + } + else if (tryAcquireAsyncInAwaitLoop(cur, thread, 0L, hasPaused, burstSignalOnAcquire)) + { + return false; + } + + int enabledCount = enabledThreadCount(cur); + if (thread >= enabledCount || enabledCount <= 1) + break; // if we're a disabled thread, or the last thread, move to blocking loop + + hasPaused = pause(thread, hasPaused); + int waitingEnabledThreadCount = waitingEnabledThreadCount(cur); + int multiplier = 1 + waitingEnabledThreadCount == 0 ? 0 : ThreadLocalRandom.current().nextInt(2 * waitingEnabledThreadCount(cur)); + LockSupport.parkNanos(spinIntervalNanos * multiplier); + } + } + + long cur = stateUpdater.addAndGet(this, threadBit); + while (true) + { + do // dummy do-while-false loop, to ensure we refresh cur + { + boolean isSignalled = (cur & threadBit) == 0; + if (isSignalled) + { + Invariants.require(lockThread(cur) != thread); + unpause(thread, hasPaused); + propagateAsyncWorkSignals(cur, burstSignalOnAcquire); + return false; + } + else if (hasLockWork(cur) && isLockAvailable(cur, thread, null)) + { + if (tryTakeLockInAwaitLoop(cur, thread, self, hasPaused)) return true; + else continue; + } + else if (tryAcquireAsyncInAwaitLoop(cur, thread, threadBit, hasPaused, burstSignalOnAcquire)) + { + return false; + } + + hasPaused = pause(thread, hasPaused); + LockSupport.park(); + } + while (false); + + cur = state; + } + } + + private boolean tryTakeLockInAwaitLoop(long cur, int thread, Thread self, boolean hasPaused) + { + long upd = setLockThread(clearLockThread(cur), thread, LOCK_OWNED) ^ threadBit(thread); + if (!stateUpdater.compareAndSet(this, cur, upd)) + return false; + + Invariants.require(depth == 0); + Invariants.require(owner == null); + owner = self; + depth = 1; + unpause(thread, hasPaused); + return true; + } + + private boolean pause(int thread, boolean hasPaused) + { + if (!hasPaused && stopLog != null) + stopLog.set(thread, nanoTime()); + return true; + } + + private void unpause(int thread, boolean hasPaused) + { + if (hasPaused && stopLog != null) + { + long start = stopLog.getAndSet(thread, 0); + long stop = nanoTime(); + updateStopCheck(stop - start, stop); + } + } + + private void updateStopCheck(long stoppedFor, long now) + { + long current = stopCheck.addAndGet(stoppedFor); + if (current > now && stopCheck.compareAndSet(current, now - stopCheckIntervalNanos)) + decrementEnabledThreadCount(); + else if (current < now - 2 * stopCheckIntervalNanos) + stopCheck.compareAndSet(current, now - stopCheckIntervalNanos); + } + + @Override + public void unlock() + { + unlock(1, false); + } + + public void unlock(int burstSignal) + { + unlock(burstSignal, false); + } + + public boolean unlockAndAcquireAsyncWork() + { + return unlockAndAcquireAsyncWork(1); + } + + public boolean unlockAndAcquireAsyncWork(int burstSignal) + { + return unlock(burstSignal, true); + } + + private boolean unlock(int burstSignal, boolean acquire) + { + Thread self = owner; + Invariants.require(self == Thread.currentThread()); + Invariants.require(!acquire || depth == 1, "Cannot acquire async work with reentrancy (depth " + depth + ')'); + return --depth == 0 && releaseAndSignalExclusive(burstSignal, acquire); + } + + private boolean releaseAndSignalExclusive(int burstSignal, boolean acquire) + { + boolean acquired = acquire && tryAcquireAsyncWork(); + owner = null; + Thread next = locking.peek(); + boolean preferRegistered = this.preferRegistered; + while (true) + { + long cur = state; + boolean wakeupLockWork = hasLockWork(cur) && hasMoreWaitersThanSignals(cur); + if (next == null && !wakeupLockWork) + { + if (stateUpdater.compareAndSet(this, cur, clearLockThread(cur))) + { + next = locking.peek(); + if (next != null) + LockSupport.unpark(next); + break; + } + continue; + } + + int thread = ANONYMOUS_OWNER; + if (wakeupLockWork) + { + if (preferRegistered || next == null) + thread = pickThreadToSignalForLock(cur); + if (next != null && autoPreferRegistered) + this.preferRegistered = !preferRegistered; + } + + long upd = setLockThread(clearLockThread(cur), thread, LOCK_SIGNALLED); + if (stateUpdater.compareAndSet(this, cur, upd)) + { + if (thread >= 0) unpark(thread); + else LockSupport.unpark(next); + break; + } + + // cancel the flip of prefer registered (if any) + this.preferRegistered = preferRegistered; + } + + propagateAsyncWorkSignals(burstSignal); + return acquired; + } + + private boolean isLockAvailable(long cur, int thread, Thread waiting) + { + if (isLockOwned(cur)) + return thread >= 0 && lockThread(cur) == thread; + + if (thread >= 0 || waiting == null) + return allowStealing || !isLockOwnedOrSignalled(cur); + + return locking.peek() == waiting && lockThread(cur) < 0; + } + + private boolean tryAcquireAsyncInAwaitLoop(long cur, int thread, long threadBitIfWaiting, boolean hasPaused, int burstSignalOnAcquire) + { + while (true) + { + int signals = asyncSignalCount(cur); + if (signals == 0) + return false; + + if ((cur & threadBitIfWaiting) != 0) + return false; + + long upd = (cur - ASYNC_SIGNAL_INCREMENT) ^ threadBitIfWaiting; + if (!stateUpdater.compareAndSet(this, cur, upd)) + { + cur = state; + continue; + } + + propagateAsyncWorkSignals(upd, burstSignalOnAcquire); + unpause(thread, hasPaused); + return true; + } + } + + public boolean incrementAsyncWork(boolean signalIfWaiting) + { + while (true) + { + long cur = state; + int count = asyncSignalCount(cur); + if (count == MAX_THREADS) + return false; + + if (signalIfWaiting) + { + int thread = pickThreadToSignalAsync(cur); + if (thread != NO_THREAD) + { + long upd = cur ^ threadBit(thread); + if (stateUpdater.compareAndSet(this, cur, upd)) + return unpark(thread); + continue; + } + } + + long upd = cur + ASYNC_SIGNAL_INCREMENT; + if (stateUpdater.compareAndSet(this, cur, upd)) + return true; + } + } + + public boolean tryAcquireAsyncWork() + { + while (true) + { + long cur = state; + int count = asyncSignalCount(cur); + if (count == 0) + return false; + + long upd = cur - ASYNC_SIGNAL_INCREMENT; + if (stateUpdater.compareAndSet(this, cur, upd)) + return true; + } + } + + public void propagateAsyncWorkSignals() + { + propagateAsyncWorkSignals(Integer.MAX_VALUE); + } + + public void propagateAsyncWorkSignals(int burstSignal) + { + propagateAsyncWorkSignals(state, burstSignal); + } + + private void propagateAsyncWorkSignals(long cur, int burstSignal) + { + while (burstSignal > 0) + { + int signals = asyncSignalCount(cur); + int thread = pickThreadToSignalAsync(cur); + if (signals == 0 || thread == NO_THREAD) + return; + + long threadBit = threadBit(thread); + Invariants.require((threadBit & cur) != 0); + long upd = (cur - ASYNC_SIGNAL_INCREMENT) ^ threadBit; + if (!stateUpdater.compareAndSet(this, cur, upd)) cur = state; + else + { + unpark(thread); + --burstSignal; + cur = upd; + } + } + } + + /** + * Signal a waiting registered thread to try and acquire the lock. + * Return true if the lock is held, or the signal has been propagated + */ + private void wakeupForLockWork() + { + while (true) + { + long cur = state; + + if (!hasLockWork(cur) || isLockOwnedOrSignalled(cur)) + return; + + // signal lock owners from end backwards, and async work from front forwards + int thread = pickThreadToSignalForLock(cur); + if (thread < 0) + return; + + Invariants.require(!isLockOwned(cur)); + long upd = setLockThread(cur, thread, LOCK_SIGNALLED); + if (stateUpdater.compareAndSet(this, cur, upd)) + { + unpark(thread); + return; + } + } + } + + private boolean unpark(int thread) + { + LockSupport.unpark(registered[thread]); + return true; + } + + /** + * Signal that there is work to do on the lock + * If the lock is free and there are any waiting threads, wake one + */ + public void signalLockWork() + { + if (setHasLockWork()) + wakeupForLockWork(); + } + + /** + * Signal that there is work to do on the lock + * If the lock is free and there are any waiting threads, wake one + */ + public void signalLockWorkExclusive() + { + setHasLockWork(); + } + + private boolean setHasLockWork() + { + return 0 == (stateUpdater.getAndUpdate(this, v -> v | HAS_LOCK_WORK) & HAS_LOCK_WORK); + } + + public void signalAllRegistered() + { + stateUpdater.updateAndGet(this, v -> v & ~WAITING_THREADS_MASK); + for (Thread thread : registered) + { + if (thread != null) + LockSupport.unpark(thread); + } + } + + /** + * Work that requires the lock has been finished. + * This must be called prior to draining all work that has been asynchronously submitted for processing with the lock. + * Otherwise, the caller must check whether new work has been submitted prior to calling unlock, + * and if so first call signalLockExclusive(). + */ + public void clearLockWork() + { + stateUpdater.updateAndGet(this, v -> v & ~HAS_LOCK_WORK); + } + + public long state() + { + return state; + } + + public int asyncSignalCount() + { + return asyncSignalCount(state); + } + + public int waitingThreadCount() + { + return waitingThreadCount(state); + } + + public int enabledThreadCount() + { + return enabledThreadCount(state); + } + + // if enabledThreadCount() == 1 return false; otherwise reduce by one + private boolean decrementEnabledThreadCount() + { + while (true) + { + long cur = state; + if (enabledThreadCount(cur) == 1) + return false; + + if (stateUpdater.compareAndSet(this, cur, cur - ENABLED_THREAD_COUNT_INCREMENT)) + return true; + } + } + + public int addAndGetEnabledThreadCount(int delta) + { + Invariants.require(delta > 0 && delta < threadCount()); + while (true) + { + long cur = state; + int count = enabledThreadCount(cur); + int newCount = Math.min(threadCount(), count + delta); + if (count == newCount) + return 0; + + if (stateUpdater.compareAndSet(this, cur, setEnabledThreadCount(cur, newCount))) + { + propagateAsyncWorkSignals(1); + return newCount - count; + } + } + } + + public void setEnabledThreadCount(int count) + { + Invariants.require(count > 0 && count <= threadCount()); + stateUpdater.accumulateAndGet(this, count, SignalLock::setEnabledThreadCount); + propagateAsyncWorkSignals(1); + } + + public int threadCount() + { + return registered.length; + } + + public int waitingEnabledThreadCount() + { + return waitingEnabledThreadCount(state); + } + + public int activeThreadCount() + { + long cur = state; + int active = registered.length - waitingThreadCount(cur); + int lockThread = lockThread(cur); + if (lockThread >= 0 && (cur & threadBit(lockThread)) != 0) + ++active; + return active; + } + + public boolean hasLockWork() + { + return hasLockWork(state); + } + + public boolean hasOwner() + { + return isLockOwned(state); + } + + public boolean isOwner() + { + return Thread.currentThread() == owner; + } + + public void setAutoPreferRegistered(boolean autoPreferRegistered) + { + this.autoPreferRegistered = autoPreferRegistered; + } + + public void setPreferRegistered(boolean preferRegistered) + { + this.preferRegistered = preferRegistered; + } + + private static final int NO_THREAD = 64; + private static final int ANONYMOUS_OWNER = -1; + private static int pickThreadToSignalAsync(long state) + { + int lockThread = lockThread(state); + if (lockThread >= 0) + state &= ~threadBit(lockThread); + return Long.numberOfTrailingZeros(state & threadMask(enabledThreadCount(state))); + } + + private static int pickThreadToSignalForLock(long state) + { + return 63 - Long.numberOfLeadingZeros(state & threadMask(enabledThreadCount(state))); + } + + private static int lockThread(long state) + { + return (int) ((state & LOCK_THREAD_MASK) >>> LOCK_THREAD_SHIFT) - 2; + } + + private static boolean isLockOwnedOrSignalled(long state) + { + return 0 != (state & LOCK_THREAD_MASK); + } + + private static boolean isLockOwned(long state) + { + return 0 != (state & LOCK_OWNED); + } + + private static long setLockThread(long state, int thread, long owned) + { + Invariants.require(!isLockOwnedOrSignalled(state)); + state |= (long) (thread + 2) << LOCK_THREAD_SHIFT; + return state | owned; + } + + private static long clearLockThread(long state) + { + state &= ~(LOCK_THREAD_MASK | LOCK_OWNED); + return state; + } + + public static int asyncSignalCount(long state) + { + return (int) ((state >>> ASYNC_SIGNAL_COUNT_SHIFT) & THREAD_ID_MASK); + } + + public static int waitingThreadCount(long state) + { + return Long.bitCount(state & WAITING_THREADS_MASK); + } + + public static int enabledThreadCount(long state) + { + return (int) ((state >>> ENABLED_THREAD_COUNT_SHIFT) & THREAD_ID_MASK); + } + + private static long setEnabledThreadCount(long state, long count) + { + Invariants.require(count <= MAX_THREADS); + return (state & ~ENABLED_THREAD_COUNT_MASK) | ((long)count << ENABLED_THREAD_COUNT_SHIFT); + } + + public static int activeEnabledThreadCount(long state) + { + return activeThreadCount(enabledThreadCount(state), state); + } + + public static int activeThreadCount(int threadCount, long state) + { + return threadCount - Long.bitCount(state & threadMask(threadCount)); + } + + private static long threadMask(int threadCount) + { + return (1L << threadCount) - 1; + } + + public static int waitingEnabledThreadCount(long state) + { + return Long.bitCount(state & threadMask(enabledThreadCount(state))); + } + + public static boolean hasMoreWaitersThanSignals(long state) + { + return waitingEnabledThreadCount(state) > asyncSignalCount(state); + } + + private static long threadBit(int threadIndex) + { + return 1L << threadIndex; + } + + private static boolean hasLockWork(long state) + { + return (state & HAS_LOCK_WORK) != 0; + } + + public void lockInterruptibly() + { + throw new UnsupportedOperationException(); + } + + @Override + public boolean tryLock(long time, TimeUnit unit) throws InterruptedException + { + throw new UnsupportedOperationException(); + } + + @Override + public Condition newCondition() + { + throw new UnsupportedOperationException(); + } +} 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 5687c25276..b0f339523c 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java @@ -23,8 +23,10 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.BitSet; import java.util.Comparator; +import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.PriorityQueue; import java.util.Random; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; @@ -53,6 +55,11 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import accord.local.CommandStore; +import accord.local.PreLoadContext; +import accord.local.SafeCommand; +import accord.primitives.PartialDeps; +import accord.primitives.TxnId; import accord.utils.Functions; import org.apache.cassandra.config.CassandraRelevantProperties; @@ -81,13 +88,17 @@ import org.apache.cassandra.service.accord.debug.CoordinationKinds; import org.apache.cassandra.service.accord.debug.TxnKindsAndDomains; import org.apache.cassandra.utils.EstimatedHistogram; import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; +import org.apache.cassandra.utils.concurrent.WaitQueue; +import static accord.coordinate.Coordination.CoordinationKind.Client; +import static accord.coordinate.Coordination.CoordinationKind.Execute; import static java.lang.System.currentTimeMillis; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.cassandra.db.ColumnFamilyStore.FlushReason.UNIT_TESTS; import static org.apache.cassandra.service.accord.debug.AccordTracing.BucketMode.LEAKY; +import static org.apache.cassandra.service.accord.debug.AccordTracing.BucketMode.RING; import static org.apache.cassandra.service.accord.debug.AccordTracing.BucketMode.SLOWEST; public class AccordLoadTest extends AccordTestBase @@ -104,10 +115,11 @@ public class AccordLoadTest extends AccordTestBase .set("accord.shard_durability_target_splits", "8") .set("accord.shard_durability_max_splits", "16") .set("accord.shard_durability_cycle", "1m") - .set("accord.queue_submission_model", "SEMI_SYNC") + .set("accord.queue_submission_model", "SIGNAL") +// .set("accord.queue_submission_model", "SEMI_SYNC") .set("accord.command_store_shard_count", "8") - .set("concurrent_accord_operations", "8") - .set("accord.queue_shard_count", "2") + .set("accord.queue_thread_count", "4") + .set("accord.queue_shard_count", "1") .set("accord.replica_execution", "ALL") .set("accord.send_stable", "TO_ALL_REPLICA_EXECUTABLE_ELSE_FOR_READS") .set("accord.send_minimal", "false") @@ -138,6 +150,7 @@ public class AccordLoadTest extends AccordTestBase final IntSupplier keySelector; final boolean readBeforeWrite; final float traceSlowest; + final int traceLast; final int[][] artificialLatencies; Settings(SettingsBuilder builder) @@ -163,6 +176,7 @@ public class AccordLoadTest extends AccordTestBase this.readBeforeWrite = builder.readBeforeWrite; this.artificialLatencies = builder.artificialLatencies; this.traceSlowest = builder.traceSlowest; + this.traceLast = builder.traceLast; } } @@ -189,6 +203,7 @@ public class AccordLoadTest extends AccordTestBase IntSupplier keySelector; boolean readBeforeWrite; float traceSlowest; + int traceLast; int[][] artificialLatencies; public SettingsBuilder setRepairInterval(int repairInterval) @@ -305,6 +320,12 @@ public class AccordLoadTest extends AccordTestBase return this; } + public SettingsBuilder setTraceLast(int traceLast) + { + this.traceLast = traceLast; + return this; + } + public SettingsBuilder setKeySelector(IntSupplier keySelector) { this.keySelector = keySelector; @@ -475,12 +496,29 @@ public class AccordLoadTest extends AccordTestBase } } + if (settings.traceLast > 0) + { + int traceLast = settings.traceLast; + for (int i = 0 ; i < cluster.size() ; ++i) + { + cluster.get(i + 1).runOnInstance(() -> { + AccordTracing tracing = ((AccordAgent) AccordService.unsafeInstance().agent()).tracing(); + tracing.setPattern(2, pattern -> pattern.withKinds(TxnKindsAndDomains.parse("{KW}")) + .withTraceNew(CoordinationKinds.ALL), + RING, -1, traceLast, LEAKY, 10, 1, CoordinationKinds.ALL); + }); + } + } + final AtomicBoolean stop = new AtomicBoolean(); + final AtomicBoolean pauseOrStop = new AtomicBoolean(); + final WaitQueue waitQueue = WaitQueue.newWaitQueue(); Random random = new Random(); Semaphore completed = new Semaphore(0); AtomicIntegerArray coordinatorIndexes = new AtomicIntegerArray(clientCount); final List> clients = new ArrayList<>(); final AtomicReferenceArray rateLimiters = new AtomicReferenceArray<>(clientCount); + final ConcurrentHashMap debugLatency = new ConcurrentHashMap<>(); final AtomicReference readHistogram = new AtomicReference<>(new EstimatedHistogram(200)); final AtomicReference writeHistogram = new AtomicReference<>(new EstimatedHistogram(200)); if (settings.clients >= cluster.size()) @@ -496,8 +534,18 @@ public class AccordLoadTest extends AccordTestBase coordinatorIndexes.set(client, client + 1); clients.add(clientExecutor.submit(() -> { final Semaphore inFlight = new Semaphore(settings.clientConcurrency); - while (!stop.get()) + while (true) { + while (pauseOrStop.get()) + { + if (stop.get()) + break; + + WaitQueue.Signal signal = waitQueue.register(); + if (pauseOrStop.get()) signal.awaitThrowUncheckedOnInterrupt(); + else signal.cancel(); + } + int coordinatorIdx = coordinatorIndexes.get(clientIndex); ICoordinator coordinator = cluster.coordinator(coordinatorIdx); try @@ -519,7 +567,8 @@ public class AccordLoadTest extends AccordTestBase completed.release(); if (fail == null) { - writeHistogram.get().add(NANOSECONDS.toMicros(System.nanoTime() - commandStart)); + long elapsed = System.nanoTime() - commandStart; + writeHistogram.get().add(NANOSECONDS.toMicros(elapsed)); synchronized (initialised) { keys.forEachInt(initialised::set); @@ -549,7 +598,10 @@ public class AccordLoadTest extends AccordTestBase inFlight.release(); completed.release(); if (fail == null) - writeHistogram.get().add(NANOSECONDS.toMicros(System.nanoTime() - commandStart)); + { + long elapsed = System.nanoTime() - commandStart; + writeHistogram.get().add(NANOSECONDS.toMicros(elapsed)); + } else logger.error("{}", fail.toString()); }, "BEGIN TRANSACTION\n" + @@ -713,10 +765,9 @@ public class AccordLoadTest extends AccordTestBase Long nowMillis = System.currentTimeMillis(); EstimatedHistogram reads = readHistogram.getAndSet(new EstimatedHistogram(200)); EstimatedHistogram writes = writeHistogram.getAndSet(new EstimatedHistogram(200)); - float traceSlowest = settings.traceSlowest; - if (traceSlowest > 0f) + if (settings.traceSlowest > 0f) { - cluster.forEach(() -> { + cluster.forEach(() -> { AccordTracing tracing = ((AccordAgent)AccordService.instance().agent()).tracing(); tracing.forEach(Functions.alwaysTrue(), (txnId, state) -> { @@ -734,6 +785,135 @@ public class AccordLoadTest extends AccordTestBase tracing.eraseAll(); }); } + if (settings.traceLast > 0) + { + pauseOrStop.set(true); + Map>> print = new HashMap<>(); + for (int i = 1 ; i <= cluster.size() ; ++i) + { + cluster.get(i).acceptOnInstance(out -> { + AccordService service = (AccordService)AccordService.instance(); + AccordTracing tracing = ((AccordAgent)AccordService.instance().agent()).tracing(); + PriorityQueue candidates = new PriorityQueue<>(Comparator.comparingLong(c -> -c.elapsedMicros)); + tracing.forEach(Functions.alwaysTrue(), (txnId, events) -> { + events.forEach(event -> { + if (event.kind == Client) + { + long doneAtMicros = event.doneAtMicros(); + long elapsedMicros = doneAtMicros - event.txnId().hlc(); + if (elapsedMicros > 350000 && elapsedMicros < 390000) + candidates.add(new SortedByElapsed(txnId, elapsedMicros)); + } + }); + }); + + AtomicInteger storeId = new AtomicInteger(); + while (!candidates.isEmpty()) + { + SortedByElapsed sortedCandidate = candidates.poll(); + if (sortedCandidate.elapsedMicros < 300000) + return; + + TxnId candidate = sortedCandidate.txnId; + storeId.lazySet(-1); + tracing.forEach(candidate, events -> { + events.forEach(event -> { + if (storeId.get() >= 0) + return; + for (Message message : event.messages()) + { + if (message.nodeId < 0 && message.commandStoreId >= 0) + { + storeId.set(message.commandStoreId); + break; + } + } + }); + }); + + if (storeId.get() >= 0) + { + CommandStore commandStore = service.node().commandStores().forId(storeId.get()); + List> result = AccordService.getBlocking(commandStore.submit(PreLoadContext.contextFor(candidate, "LoadTest"), safeStore -> { + SafeCommand safeCommand = safeStore.unsafeGet(candidate); + PartialDeps deps = safeCommand.current().partialDeps(); + if (deps == null) + return null; + List> infos = new ArrayList<>(); + for (TxnId txnId : deps.txnIds()) + { + List info = new ArrayList<>(); + info.add(txnId.toString()); + infos.add(info); + } + List info = new ArrayList<>(); + info.add(candidate.toString()); + infos.add(info); + return infos; + })); + + if (result != null) + { + for (List info : result) + { + TxnId txnId = TxnId.parse(info.get(0)); + AccordService.getBlocking(commandStore.execute(PreLoadContext.contextFor(txnId, "LoadTest"), safeStore -> { + SafeCommand safeCommand = safeStore.unsafeGet(txnId); + if (safeCommand.current().executeAt != null) + info.add(safeCommand.current().executeAt.toString()); + })); + } + + out.put(candidate.toString(), result); + return; + } + } + } + }, print); + } + + for (int i = 1 ; i <= cluster.size() ; ++i) + { + cluster.get(i).acceptOnInstance(out -> { + AccordTracing tracing = ((AccordAgent)AccordService.instance().agent()).tracing(); + for (Map.Entry>> e : out.entrySet()) + { + TxnId parentId = TxnId.parse(e.getKey()); + for (List infos : e.getValue()) + { + TxnId depId = TxnId.parse(infos.get(0)); + tracing.forEach(depId, events -> { + events.forEach(event -> { + infos.add(event.kind + ": [" + (event.idMicros - parentId.hlc()) + "..." + (event.doneAtMicros() - parentId.hlc()) + "][" + (event.idMicros - depId.hlc()) + "..." + (event.doneAtMicros() - depId.hlc()) + "]"); + if (event.kind == Execute) + { + for (Message message : event.messages()) + { + if (message.nodeId == parentId.node.id) + { + long atMicros = (event.idMicros + (message.atNanos - event.atNanos)/1000) - parentId.hlc(); + infos.add(atMicros + ": " + message.message); + } + } + } + }); + }); + } + } + }, print); + } + + for (Map.Entry>> e : print.entrySet()) + { + System.out.println("======" + e.getKey() + "======"); + for (List infos : e.getValue()) + System.out.println(infos); + } + if (!print.isEmpty()) + System.out.println(); + pauseOrStop.set(false); + waitQueue.signalAll(); + } cluster.forEach(() -> { refresh(AccordExecutorMetrics.INSTANCE.elapsedRunning); refresh(AccordExecutorMetrics.INSTANCE.elapsed); @@ -879,16 +1059,17 @@ public class AccordLoadTest extends AccordTestBase ws[i] = iw; } System.out.println(Arrays.toString(ws)); + Arrays.fill(ws, 0); for (int i = 0 ; i < qs.length ; ++i) { - int wj = i == 0 ? 1 : 0; - for (int j = 1 ; j < qs.length ; ++j) + for (int j = 0 ; j < qs.length ; ++j) { if (j == i) continue; - if (qs[j] > qs[wj]) - wj = j; + if (qs[j] > 2*qs[i]) continue; + int w = qs[i] + 4*qs[j] + LATENCIES[i][j]; + if (w > ws[i]) + ws[i] = w; } - ws[i] = qs[i] + 4*qs[wj] + LATENCIES[i][wj]; } System.out.println(Arrays.toString(ws)); } @@ -910,11 +1091,33 @@ public class AccordLoadTest extends AccordTestBase DistributedTestBase.beforeClass(); AccordLoadTest.setUp(); AccordLoadTest test = new AccordLoadTest(); - test.setup(); - test.testLoad(withArtificialLatencies(ycsbA(new SettingsBuilder(), 100_000) - .setRatePerSecond(1600).setMinRatePerSecond(200) - .setIncreaseRatePerSecondInterval(5000) -// .setTraceSlowest(0.5f) - ).build()); + try + { + test.setup(); + test.testLoad(withArtificialLatencies(ycsbA(new SettingsBuilder(), 100_000) +// .setRatePerSecond(400).setMinRatePerSecond(200) +// .setRatePerSecond(800).setMinRatePerSecond(200) + .setRatePerSecond(1600).setMinRatePerSecond(200) + .setIncreaseRatePerSecondInterval(5000) +// .setTraceLast(5000) + ).build()); + } + finally + { + test.tearDown(); + } } + + static class SortedByElapsed + { + final TxnId txnId; + final long elapsedMicros; + + SortedByElapsed(TxnId txnId, long elapsedMicros) + { + this.txnId = txnId; + this.elapsedMicros = elapsedMicros; + } + } + } diff --git a/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java b/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java index 2d3474f36d..6749d22555 100644 --- a/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java +++ b/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java @@ -65,7 +65,7 @@ public class AccordBootstrapTest extends FuzzTestBase .withConfig((config) -> config.with(Feature.NETWORK, Feature.GOSSIP) .set("write_request_timeout", "2s") .set("request_timeout", "5s") - .set("concurrent_accord_operations", 2) + .set("accord.queue_thread_count", 2) .set("accord.shard_durability_target_splits", "1") .set("accord.shard_durability_max_splits", "4") .set("accord.catchup_on_start_fail_latency", "60s") diff --git a/test/distributed/org/apache/cassandra/fuzz/topology/AccordTopologyMixupTest.java b/test/distributed/org/apache/cassandra/fuzz/topology/AccordTopologyMixupTest.java index 3f5a8ef12c..834d95179c 100644 --- a/test/distributed/org/apache/cassandra/fuzz/topology/AccordTopologyMixupTest.java +++ b/test/distributed/org/apache/cassandra/fuzz/topology/AccordTopologyMixupTest.java @@ -276,7 +276,7 @@ public class AccordTopologyMixupTest extends TopologyMixupTestBase new AccordExecutorSignalLoop(1, RUN_WITHOUT_LOCK, THREAD_COUNT, -1, -1, TimeUnit.MICROSECONDS, i ->"Loop" + i, new AccordAgent())); + } + + @Test + public void signalSpinLoopTest() + { + executorTest(() -> new AccordExecutorSignalLoop(1, RUN_WITHOUT_LOCK, THREAD_COUNT, 10, 100, TimeUnit.MICROSECONDS, i ->"Loop" + i, new AccordAgent())); + } + + @Test + public void ayncSubmitTest() + { + executorTest(() -> new AccordExecutorAsyncSubmit(1, RUN_WITHOUT_LOCK, THREAD_COUNT, i -> "Loop" + i, new AccordAgent())); + } + + public void executorTest(SerializableSupplier supplier) + { + simulate(arr(() -> { + try + { + DatabaseDescriptor.daemonInitialization(); + AccordExecutor executor = supplier.get(); + Lock lock = executor.unsafeLock(); + SequentialAsyncExecutor sequentialExecutor = executor.newSequentialExecutor(); + Executor lockExecutor = ExecutorFactory.Global.executorFactory().sequential("lock"); + ConcurrentLinkedQueue> await = new ConcurrentLinkedQueue<>(); + + for (float sleepChance : new float[] { 0f, 0.01f, 0.1f }) + for (float lockChance : new float[] { 0f, 0.01f, 0.1f }) + submitLoop(lock, executor, sequentialExecutor, lockExecutor, 200, 10, await, sleepChance, lockChance); + } + catch (Throwable t) + { + throw new RuntimeException(t); + } + }), + () -> {}, 1L); + } + + private static void submitLoop(Lock lock, AccordExecutor executor, SequentialAsyncExecutor sequentialExecutor, Executor lockExecutor, int outerLoop, int innerLoop, ConcurrentLinkedQueue> await, float sleepChance, float lockChance) throws ExecutionException, InterruptedException + { + while (outerLoop-- > 0) + { + for (int i = 0; i < innerLoop; ++i) + submitRecursive(lock, executor, sequentialExecutor, 1 + i, await, sleepChance, lockChance); + + AtomicBoolean done = new AtomicBoolean(); + submitUntil(lock, lockExecutor, sleepChance, done::get); + while (!await.isEmpty()) + await.poll().get(); + done.set(true); + System.out.println("Loop " + (1 + outerLoop)); + } + } + + private static void submitRecursive(Lock lock, AccordExecutor executor, SequentialAsyncExecutor sequentialExecutor, int count, Collection> await, float sleepChance, float lockChance) + { + AsyncExecutor submitTo = ThreadLocalRandom.current().nextBoolean() ? executor : sequentialExecutor; + await.add(toFuture(submitTo.chain(() -> { + ThreadLocalRandom rnd = ThreadLocalRandom.current(); + boolean locked = false; + if (rnd.nextFloat() < lockChance) + { + if (rnd.nextBoolean()) locked = lock.tryLock(); + else { locked = true; lock.lock(); } + } + try + { + if (count > 1) + submitRecursive(lock, executor, sequentialExecutor, count -1, await, sleepChance, lockChance); + if (rnd.nextFloat() < sleepChance) + LockSupport.parkNanos(rnd.nextInt(10000, 100000)); + } + finally + { + if (locked) + lock.unlock(); + } + }).beginAsResult())); + } + + private static void submitUntil(Lock lock, Executor executor, float sleepChance, BooleanSupplier done) + { + if (done.getAsBoolean()) + return; + + executor.execute(() -> { + + ThreadLocalRandom rnd = ThreadLocalRandom.current(); + boolean tryLock = rnd.nextBoolean(); + boolean locked = !tryLock; + if (tryLock) locked = lock.tryLock(); + else lock.lock(); + try + { + if (rnd.nextFloat() < sleepChance) + LockSupport.parkNanos(rnd.nextInt(10000, 100000)); + + submitUntil(lock, executor, sleepChance, done); + } + finally + { + if (locked) + lock.unlock(); + } + }); + } +} diff --git a/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java b/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java index a316b883ae..aa726aaee3 100644 --- a/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java +++ b/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java @@ -256,7 +256,7 @@ public class AccordDebugKeyspaceTest extends CQLTester Config.setOverrideLoadConfig(() -> { Config config = new YamlConfigurationLoader().loadConfig(); config.accord.queue_shard_count = new OptionaldPositiveInt(1); - config.concurrent_accord_operations = 1; + config.accord.queue_thread_count = new OptionaldPositiveInt(1); config.accord.command_store_shard_count = new OptionaldPositiveInt(1); config.accord.enable_virtual_debug_only_keyspace = true; config.accord.permit_fast_quorum_medium_path = true; diff --git a/test/unit/org/apache/cassandra/hints/HintsServiceTest.java b/test/unit/org/apache/cassandra/hints/HintsServiceTest.java index ddfe4a2c8b..610b534082 100644 --- a/test/unit/org/apache/cassandra/hints/HintsServiceTest.java +++ b/test/unit/org/apache/cassandra/hints/HintsServiceTest.java @@ -269,7 +269,7 @@ public class HintsServiceTest { accordTxnCount.incrementAndGet(); TxnId txnId = AccordTestUtils.txnId(42, 43, 44); - AccordResult result = new AccordResult<>(txnId, Keys.EMPTY, new LatencyRequestBookkeeping(null), requestTime.startedAtNanos(), requestTime.startedAtNanos(), true); + AccordResult result = new AccordResult<>(txnId, Keys.EMPTY, new LatencyRequestBookkeeping(null), requestTime.startedAtNanos(), requestTime.startedAtNanos(), true, null); result.accept(new TxnData(), null); return result; } diff --git a/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java b/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java index a6a8ea13ad..102485011d 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordCommandStoreTest.java @@ -30,7 +30,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import accord.api.Key; -import accord.api.Result; +import accord.api.Result.PersistableResult; import accord.local.Command; import accord.local.PreLoadContext; import accord.local.StoreParticipants; @@ -134,7 +134,7 @@ public class AccordCommandStoreTest LargeBitSet waitingOnApply = new LargeBitSet(3); waitingOnApply.set(1); Command.WaitingOn waitingOn = new Command.WaitingOn(dependencies.keyDeps.keys(), dependencies.rangeDeps, new ImmutableBitSet(waitingOnApply), new ImmutableBitSet(2)); - Pair result = getBlocking(AccordTestUtils.processTxnResult(commandStore, txnId, txn, executeAt)); + Pair result = getBlocking(AccordTestUtils.processTxnResult(commandStore, txnId, txn, executeAt)); Command expected = Command.Executed.executed(txnId, SaveStatus.Applied, AllQuorums, StoreParticipants.all(route), promised, executeAt, txn, dependencies, accepted, diff --git a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java index 907df14190..dcc631d776 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java @@ -41,6 +41,7 @@ import accord.api.Journal; import accord.api.ProgressLog.NoOpProgressLog; import accord.api.RemoteListeners.NoOpRemoteListeners; import accord.api.Result; +import accord.api.Result.PersistableResult; import accord.api.RoutingKey; import accord.api.Timeouts; import accord.coordinate.Coordinations; @@ -240,15 +241,15 @@ public class AccordTestUtils return Ballot.fromValues(epoch, hlc, new Node.Id(node)); } - public static AsyncChain> processTxnResult(AccordCommandStore commandStore, TxnId txnId, PartialTxn txn, Timestamp executeAt) throws Throwable + public static AsyncChain> processTxnResult(AccordCommandStore commandStore, TxnId txnId, PartialTxn txn, Timestamp executeAt) throws Throwable { - AtomicReference>> result = new AtomicReference<>(); + AtomicReference>> result = new AtomicReference<>(); getBlocking(commandStore.execute((PreLoadContext.Empty)() -> "Test", safeStore -> result.set(processTxnResultDirect(safeStore, txnId, txn, executeAt)))); return result.get(); } - public static AsyncChain> processTxnResultDirect(SafeCommandStore safeStore, TxnId txnId, PartialTxn txn, Timestamp executeAt) + public static AsyncChain> processTxnResultDirect(SafeCommandStore safeStore, TxnId txnId, PartialTxn txn, Timestamp executeAt) { TxnRead read = (TxnRead) txn.read(); return AsyncChains.allOf(read.keys().stream().map(key -> read.read(safeStore, key, executeAt)) @@ -256,7 +257,7 @@ public class AccordTestUtils .map(list -> { Data data = list.stream().reduce(Data::merge).orElse(new TxnData()); return Pair.create(txnId.is(Write) ? txn.execute(txnId, executeAt, data) : null, - txn.query().compute(txnId, executeAt, txn.keys(), data, txn.read(), txn.update())); + txn.query().compute(txnId, executeAt, txn.keys(), data, txn.read(), txn.update()).toPersistable()); }); } diff --git a/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java b/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java index b5d56053e0..d1982cfa56 100644 --- a/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java +++ b/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java @@ -121,7 +121,7 @@ import org.apache.cassandra.utils.Pair; import static accord.utils.Property.commands; import static accord.utils.Property.stateful; -import static org.apache.cassandra.config.DatabaseDescriptor.getAccordCommandStoreShardCount; +import static org.apache.cassandra.config.DatabaseDescriptor.getAccord; import static org.apache.cassandra.config.DatabaseDescriptor.getPartitioner; public class EpochSyncTest @@ -671,7 +671,7 @@ public class EpochSyncTest this.token = token; this.epoch = epoch; MockCluster.Clock clock = new MockCluster.Clock(0); - Node node = Utils.createNode(id, ignore -> AsyncResults.settable(), new MessageSink.NoOpSink(), clock, new TestAgent(clock), new TokenKey.KeyspaceSplitter(new ShardDistributor.EvenSplit<>(getAccordCommandStoreShardCount(), getPartitioner().accordSplitter()))); + Node node = Utils.createNode(id, ignore -> AsyncResults.settable(), new MessageSink.NoOpSink(), clock, new TestAgent(clock), new TokenKey.KeyspaceSplitter(new ShardDistributor.EvenSplit<>(getAccord().commandStoreShardCount(), getPartitioner().accordSplitter()))); // TODO (review): Should there be a real scheduler here? Is it possible to adapt the Scheduler interface to scheduler used in this test? TimeService time = TimeService.ofNonMonotonic(globalExecutor::currentTimeMillis, TimeUnit.MILLISECONDS); this.topologyService = new AccordTopologyService(id, mapper, messagingService, scheduler); 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 92cc500f60..ca1c1a4e30 100644 --- a/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java +++ b/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java @@ -120,7 +120,6 @@ import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.service.accord.api.AccordAgent; import org.apache.cassandra.service.accord.api.TokenKey; -import org.apache.cassandra.service.accord.txn.TxnData; import org.apache.cassandra.service.accord.txn.TxnWrite; import org.apache.cassandra.simulator.RandomSource.Choices; import org.apache.cassandra.utils.AccordGenerators; @@ -222,7 +221,7 @@ public class CommandsForKeySerializerTest { if (txnId.is(Kind.Write)) builder.writes(new Writes(txnId, executeAt, txn.keys(), new TxnWrite(TableMetadatas.none(), Collections.emptyList(), SimpleBitSets.allSet(1)))); - builder.result(new TxnData()); + builder.result(ResultSerializers.APPLIED); } return builder; } diff --git a/test/unit/org/apache/cassandra/utils/AccordGenerators.java b/test/unit/org/apache/cassandra/utils/AccordGenerators.java index 67d23a4126..7019e05c68 100644 --- a/test/unit/org/apache/cassandra/utils/AccordGenerators.java +++ b/test/unit/org/apache/cassandra/utils/AccordGenerators.java @@ -84,9 +84,9 @@ import org.apache.cassandra.service.accord.AccordTestUtils; import org.apache.cassandra.service.accord.TokenRange; import org.apache.cassandra.service.accord.api.PartitionKey; import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.service.accord.serializers.ResultSerializers; import org.apache.cassandra.service.accord.serializers.TableMetadatas; import org.apache.cassandra.service.accord.topology.FetchTopologies; -import org.apache.cassandra.service.accord.txn.TxnData; import org.apache.cassandra.service.accord.txn.TxnWrite; import static accord.local.CommandStores.RangesForEpoch; @@ -292,7 +292,7 @@ public class AccordGenerators { if (txnId.is(Write)) builder.writes(new Writes(txnId, executeAt, keysOrRanges, new TxnWrite(TableMetadatas.none(), Collections.emptyList(), SimpleBitSets.allSet(1)))); - builder.result(new TxnData()); + builder.result(ResultSerializers.APPLIED); } return builder.build(saveStatus); }