diff --git a/modules/accord b/modules/accord index a1c1ed91cf..51f36f788b 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit a1c1ed91cfeb904be9b490d498c1c52aba2253c6 +Subproject commit 51f36f788b8095c0da9efaf1ea91bb9a7dd31181 diff --git a/src/java/org/apache/cassandra/config/AccordSpec.java b/src/java/org/apache/cassandra/config/AccordSpec.java index 817532e23d..ec70398f64 100644 --- a/src/java/org/apache/cassandra/config/AccordSpec.java +++ b/src/java/org/apache/cassandra/config/AccordSpec.java @@ -57,14 +57,11 @@ public class AccordSpec /** * 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. - * - * The global READ and WRITE stages are used for IO. */ THREAD_PER_SHARD, /** * Same number of threads as shards, and the shard lock is held for the duration of serving requests. - * The global READ and WRITE stages are used for IO. */ THREAD_PER_SHARD_SYNC_QUEUE, @@ -73,12 +70,6 @@ public class AccordSpec * Fewer shards is generally better, until queue-contention is encountered. */ THREAD_POOL_PER_SHARD, - - /** - * More threads than shards. Threads update transaction state only, relying on READ and WRITE stages for IO. - * Fewer shards is generally better, until queue-contention is encountered. - */ - THREAD_POOL_PER_SHARD_EXCLUDES_IO, } public enum QueueSubmissionModel @@ -130,6 +121,7 @@ public class AccordSpec public volatile OptionaldPositiveInt max_queued_loads = OptionaldPositiveInt.UNDEFINED; public volatile OptionaldPositiveInt max_queued_range_loads = OptionaldPositiveInt.UNDEFINED; + public volatile OptionaldPositiveInt max_progress_log_concurrency = OptionaldPositiveInt.UNDEFINED; public DataStorageSpec.LongMebibytesBound cache_size = null; public DataStorageSpec.LongMebibytesBound working_set_size = null; @@ -158,8 +150,8 @@ public class AccordSpec public volatile DurationSpec.IntSecondsBound fast_path_update_delay = null; - public volatile DurationSpec.IntSecondsBound gc_delay = new DurationSpec.IntSecondsBound("5m"); public volatile int shard_durability_target_splits = 16; + public volatile int shard_durability_max_splits = 128; public volatile DurationSpec.IntSecondsBound durability_txnid_lag = new DurationSpec.IntSecondsBound(5); public volatile DurationSpec.IntSecondsBound shard_durability_cycle = new DurationSpec.IntSecondsBound(5, TimeUnit.MINUTES); public volatile DurationSpec.IntSecondsBound global_durability_cycle = new DurationSpec.IntSecondsBound(5, TimeUnit.MINUTES); diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 28a74a39b9..4e77038a8f 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -5372,7 +5372,6 @@ public class DatabaseDescriptor case THREAD_PER_SHARD_SYNC_QUEUE: return conf.accord.queue_shard_count.or(DatabaseDescriptor::getAvailableProcessors); case THREAD_POOL_PER_SHARD: - case THREAD_POOL_PER_SHARD_EXCLUDES_IO: int defaultMax = getAccordQueueSubmissionModel() == AccordSpec.QueueSubmissionModel.SYNC ? 8 : 4; return conf.accord.queue_shard_count.or(Math.min(defaultMax, DatabaseDescriptor.getAvailableProcessors())); } @@ -5393,6 +5392,11 @@ public class DatabaseDescriptor return conf.accord.max_queued_range_loads.or(Math.max(4, getAccordConcurrentOps() / 4)); } + public static int getAccordProgressLogMaxConcurrency() + { + return conf.accord.max_progress_log_concurrency.or(64); + } + public static boolean getAccordCacheShrinkingOn() { return conf.accord.shrink_cache_entries_before_eviction; @@ -5426,16 +5430,16 @@ public class DatabaseDescriptor return bound == null ? -1 : bound.to(TimeUnit.MILLISECONDS); } - public static long getAccordGCDelay(TimeUnit unit) - { - return conf.accord.gc_delay.to(unit); - } - public static int getAccordShardDurabilityTargetSplits() { return conf.accord.shard_durability_target_splits; } + public static int getAccordShardDurabilityMaxSplits() + { + return conf.accord.shard_durability_max_splits; + } + public static long getAccordScheduleDurabilityTxnIdLag(TimeUnit unit) { return conf.accord.durability_txnid_lag.to(unit); diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index a499bc42f2..4cc86d4889 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1374,6 +1374,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean, Memtable.Owner } } cfs.replaceFlushed(memtable, sstables); + memtable.notifyFlushed(); reclaim(memtable); cfs.compactionStrategyManager.compactionLogger.flush(sstables); logger.debug("Flushed to {} ({} sstables, {}), biggest {}, smallest {}", diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java index 0cbd9fe3a8..eb331dda8f 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java @@ -30,6 +30,7 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import com.google.common.collect.Iterables; import com.google.common.util.concurrent.RateLimiter; + import org.apache.cassandra.db.compaction.unified.UnifiedCompactionTask; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; diff --git a/src/java/org/apache/cassandra/db/memtable/AbstractMemtable.java b/src/java/org/apache/cassandra/db/memtable/AbstractMemtable.java index dd2254721b..8983979a91 100644 --- a/src/java/org/apache/cassandra/db/memtable/AbstractMemtable.java +++ b/src/java/org/apache/cassandra/db/memtable/AbstractMemtable.java @@ -18,6 +18,7 @@ package org.apache.cassandra.db.memtable; +import java.util.Collection; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ConcurrentSkipListSet; @@ -25,8 +26,11 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; +import java.util.function.Supplier; import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.ImmutableMap; import org.apache.cassandra.db.RegularAndStaticColumns; import org.apache.cassandra.db.commitlog.CommitLogPosition; @@ -51,6 +55,7 @@ public abstract class AbstractMemtable implements Memtable // The smallest local deletion time for all partitions in this memtable protected AtomicLong minLocalDeletionTime = new AtomicLong(Long.MAX_VALUE); private final long id = nextId.incrementAndGet(); + private Map> onFlush = ImmutableMap.of(); // Note: statsCollector has corresponding statistics to the two above, but starts with an epoch value which is not // correct for their usage. @@ -147,6 +152,35 @@ public abstract class AbstractMemtable implements Memtable return this.flushTransaction.getAndSet(flushTransaction); } + @Override + public synchronized > T ensureFlushListener(Object key, Supplier factory) + { + if (onFlush == null) + return null; + + T listener = (T)onFlush.get(key); + if (null == listener) + { + listener = factory.get(); + onFlush = ImmutableMap.>builder() + .putAll(onFlush) + .put(key, listener) + .build(); + } + return listener; + } + + public void notifyFlushed() + { + Collection> run; + synchronized (this) + { + run = onFlush.values(); + onFlush = null; + } + run.forEach(c -> c.accept(metadata())); + } + protected static class ColumnsCollector { private final HashMap predefined = new HashMap<>(); diff --git a/src/java/org/apache/cassandra/db/memtable/Memtable.java b/src/java/org/apache/cassandra/db/memtable/Memtable.java index dc0b7b3a1d..4d665b8bf9 100644 --- a/src/java/org/apache/cassandra/db/memtable/Memtable.java +++ b/src/java/org/apache/cassandra/db/memtable/Memtable.java @@ -19,6 +19,8 @@ package org.apache.cassandra.db.memtable; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; +import java.util.function.Supplier; import javax.annotation.concurrent.NotThreadSafe; import org.apache.cassandra.db.ColumnFamilyStore; @@ -420,6 +422,10 @@ public interface Memtable extends Comparable, UnfilteredSource return shouldSwitch(reason, metadata()); } + // returns null if already flushed + > T ensureFlushListener(Object key, Supplier factory); + void notifyFlushed(); + /** * Called when the table's metadata is updated. The memtable's metadata reference now points to the new version. * This will not be called if {@link #shouldSwitch } (SCHEMA_CHANGE) returns true, as the memtable will be swapped out diff --git a/src/java/org/apache/cassandra/service/accord/AccordCache.java b/src/java/org/apache/cassandra/service/accord/AccordCache.java index b7340d3a66..839f31f1a6 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCache.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCache.java @@ -52,7 +52,6 @@ import accord.utils.IntrusiveLinkedList; import accord.utils.Invariants; import accord.utils.QuadFunction; import accord.utils.TriFunction; -import accord.utils.async.Cancellable; import org.agrona.collections.Object2ObjectHashMap; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.config.DatabaseDescriptor; @@ -60,6 +59,7 @@ import org.apache.cassandra.exceptions.UnknownTableException; import org.apache.cassandra.io.util.DataInputBuffer; import org.apache.cassandra.metrics.AccordCacheMetrics; import org.apache.cassandra.metrics.CacheAccessMetrics; +import org.apache.cassandra.service.accord.AccordCacheEntry.LoadExecutor; import org.apache.cassandra.service.accord.AccordCacheEntry.Status; import org.apache.cassandra.service.accord.events.CacheEvents; import org.apache.cassandra.service.accord.serializers.Version; @@ -72,6 +72,8 @@ import static accord.utils.Invariants.require; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.EVICTED; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.LOADED; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.MODIFIED; +import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.SAVING; +import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.WAITING_TO_SAVE; /** * Cache for AccordCommand and AccordCommandsForKey, available memory is shared between the two object types. @@ -100,6 +102,7 @@ public class AccordCache implements CacheSize { @Nullable V load(AccordCommandStore commandStore, K key); @Nullable Runnable save(AccordCommandStore commandStore, K key, @Nullable V value, @Nullable Object shrunk); + default boolean canSave(@Nullable V value, @Nullable Object shrunk) { return true; } // a result of null means we can immediately evict, without saving @Nullable V quickShrink(V value); // a result of null means we cannot shrink, and should save/evict as appropriate @@ -138,9 +141,7 @@ public class AccordCache implements CacheSize } private final List> types = new CopyOnWriteArrayList<>(); - private final Function saveExecutor; - private final AccordCacheEntry.OnSaved onSaved; - // TODO (required): monitor this queue and periodically clean up entries, or implement an eviction deadline system + final AccordCacheEntry.SaveExecutor saveExecutor; private final IntrusiveLinkedList> evictQueue = new IntrusiveLinkedList<>(); private final IntrusiveLinkedList> noEvictQueue = new IntrusiveLinkedList<>(); @@ -155,10 +156,9 @@ public class AccordCache implements CacheSize final AccordCacheMetrics metrics; final Stats stats = new Stats(); - public AccordCache(Function saveExecutor, AccordCacheEntry.OnSaved onSaved, long maxSizeInBytes, AccordCacheMetrics metrics) + public AccordCache(AccordCacheEntry.SaveExecutor saveExecutor, long maxSizeInBytes, AccordCacheMetrics metrics) { this.saveExecutor = saveExecutor; - this.onSaved = onSaved; this.maxSizeInBytes = maxSizeInBytes; this.metrics = metrics; } @@ -265,14 +265,21 @@ public class AccordCache implements CacheSize evict(node, true); break; case MODIFIED: - Type parent = node.owner.parent(); - node.save(saveExecutor, parent.adapter, onSaved); + node.save(); boolean evict = node.status() == LOADED; node.unlink(); if (evict) evict(node, true); } } + public void saveWhenReadyExclusive(AccordCacheEntry entry, Runnable onSuccess) + { + if (!entry.isSavingOrWaiting() && !entry.saveWhenReady()) + onSuccess.run(); + else + entry.savingOrWaitingToSave().identity.onSuccess(onSuccess); + } + private void evict(AccordCacheEntry node, boolean updateUnreferenced) { if (logger.isTraceEnabled()) @@ -302,10 +309,10 @@ public class AccordCache implements CacheSize node.evicted(); } - Collection> load(BiFunction loadExecutor, P param, AccordCacheEntry node, AccordCacheEntry.OnLoaded onLoaded) + Collection> load(LoadExecutor loadExecutor, P1 p1, P2 p2, AccordCacheEntry node) { Type parent = node.owner.parent(); - return node.load(loadExecutor, param, parent.adapter, onLoaded).waiters(); + return node.load(loadExecutor, p1, p2).waiters(); } void loaded(AccordCacheEntry node, V value) @@ -329,7 +336,7 @@ public class AccordCache implements CacheSize void saved(AccordCacheEntry node, Object identity, Throwable fail) { - if (node.saved(identity, fail) && node.references() == 0) + if (node.saved(identity, fail) && node.references() == 0 && node.isUnqueued()) evictQueue.addFirst(node); // add to front since we have just saved, so we were eligible for eviction } @@ -539,6 +546,7 @@ public class AccordCache implements CacheSize switch (status) { default: throw new IllegalStateException("Unhandled status " + status); + case WAITING_TO_SAVE: case WAITING_TO_LOAD: case LOADING: case LOADED: @@ -560,7 +568,7 @@ public class AccordCache implements CacheSize return cache.values().stream(); } - Type parent() + final Type parent() { return Type.this; } @@ -568,7 +576,7 @@ public class AccordCache implements CacheSize @Override public Iterator> iterator() { - return stream().iterator(); + return cache.values().iterator(); } void validateLoadEvicted(AccordCacheEntry node) @@ -672,7 +680,6 @@ public class AccordCache implements CacheSize listeners = null; return true; } - } private final Class keyClass; @@ -733,7 +740,7 @@ public class AccordCache implements CacheSize AccordCache.this.stats.misses++; } - AccordCache parent() + final AccordCache parent() { return AccordCache.this; } @@ -768,7 +775,7 @@ public class AccordCache implements CacheSize return ((SettableWrapper)adapter).load; } - Adapter adapter() + final Adapter adapter() { return adapter; } @@ -1120,6 +1127,12 @@ public class AccordCache implements CacheSize return commandStore.saveCommandsForKey(key, value, serialized); } + @Override + public boolean canSave(@Nullable CommandsForKey value, @Nullable Object serialized) + { + return value == null || !value.isLoadingPruned(); + } + @Override public CommandsForKey quickShrink(CommandsForKey value) { diff --git a/src/java/org/apache/cassandra/service/accord/AccordCacheEntry.java b/src/java/org/apache/cassandra/service/accord/AccordCacheEntry.java index 58bf4cdf75..e322a51ee8 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCacheEntry.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCacheEntry.java @@ -20,10 +20,11 @@ package org.apache.cassandra.service.accord; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.List; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.function.BiConsumer; -import java.util.function.BiFunction; -import java.util.function.Function; + +import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; import com.google.common.primitives.Ints; @@ -43,6 +44,7 @@ import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.LOADIN import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.MODIFIED; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.SAVING; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.WAITING_TO_LOAD; +import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.WAITING_TO_SAVE; /** * Global (per CommandStore) state of a cached entity (Command or CommandsForKey). @@ -52,16 +54,29 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode public enum Status { UNINITIALIZED, + + UNUSED1, // spacing to permit easier bit masks + WAITING_TO_LOAD(UNINITIALIZED), LOADING(WAITING_TO_LOAD), + /** * Consumers should never see this state */ FAILED_TO_LOAD(LOADING), + UNUSED2, // spacing to permit easier bit masks + UNUSED3, // spacing to permit easier bit masks + UNUSED4, // spacing to permit easier bit masks + LOADED(true, false, UNINITIALIZED, LOADING), MODIFIED(true, false, LOADED), - SAVING(true, true, MODIFIED), + + UNUSED5, // spacing to permit easier bit masks + UNUSED6, // spacing to permit easier bit masks + + WAITING_TO_SAVE(true, true, MODIFIED), + SAVING(true, true, MODIFIED, WAITING_TO_SAVE), /** * Attempted to save but failed. Shouldn't normally happen unless we have a bug in serialization, @@ -69,7 +84,7 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode */ FAILED_TO_SAVE(true, true, SAVING), - UNUSED, // spacing to permit easier bit masks + UNUSED7, // spacing to permit easier bit masks EVICTED(WAITING_TO_LOAD, LOADING, LOADED, FAILED_TO_LOAD), ; @@ -84,8 +99,11 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode LOADED.permittedFrom |= 1 << MODIFIED.ordinal(); for (Status status : VALUES) { + if (status.name().startsWith("UNUSED")) continue; Invariants.require((status.ordinal() & IS_LOADED) != 0 == status.loaded); Invariants.require(((status.ordinal() & IS_LOADED) != 0 && (status.ordinal() & IS_NESTED) != 0) == status.nested); + Invariants.require(((status.ordinal() & IS_LOADING_OR_WAITING_MASK) == IS_LOADING_OR_WAITING) == (status == LOADING || status == WAITING_TO_LOAD)); + Invariants.require(((status.ordinal() & IS_SAVING_OR_WAITING_MASK) == IS_SAVING_OR_WAITING) == (status == SAVING || status == WAITING_TO_SAVE)); } } @@ -110,10 +128,13 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode static final int STATUS_MASK = 0x0000001F; static final int SHRUNK = 0x00000040; static final int NO_EVICT = 0x00000020; - static final int IS_LOADED = 0x4; - static final int IS_NESTED = 0x2; // only valid to test if already tested NORMAL - static final int IS_LOADING_OR_WAITING_MASK = 0x6; // only valid to test if already tested NORMAL - static final int IS_LOADING_OR_WAITING = 0x2; // only valid to test if already tested NORMAL + static final int IS_NOT_EVICTED = 0xF; + static final int IS_LOADED = 0x8; + static final int IS_NESTED = 0x4; + static final int IS_LOADING_OR_WAITING_MASK = 0x6; + static final int IS_LOADING_OR_WAITING = 0x2; + static final int IS_SAVING_OR_WAITING_MASK = 0xE; + static final int IS_SAVING_OR_WAITING = 0xC; static final long EMPTY_SIZE = ObjectSizes.measure(new AccordCacheEntry<>(null, null)); private final K key; @@ -166,6 +187,11 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode return (status & IS_LOADED) != 0; } + boolean isModified() + { + return (status & IS_NOT_EVICTED) >= MODIFIED.ordinal(); + } + boolean isNested() { Invariants.require(isLoaded()); @@ -187,6 +213,11 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode return (status & IS_LOADING_OR_WAITING_MASK) == IS_LOADING_OR_WAITING; } + boolean isSavingOrWaiting() + { + return (status & IS_SAVING_OR_WAITING_MASK) == IS_SAVING_OR_WAITING; + } + public boolean isComplete() { return !is(LOADING) && !is(SAVING); @@ -288,71 +319,34 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode owner.notifyListeners(notify, this); } - public interface OnLoaded + public interface LoadExecutor { - void onLoaded(AccordCacheEntry state, V value, Throwable fail); + Cancellable load(P1 p1, P2 p2, AccordCacheEntry entry); + } - static OnLoaded immediate() + // functions as both an identity object, and a register of listeners + public static class UniqueSave + { + @Nullable List onSuccess; + void onSuccess(Runnable onSuccess) { - return new OnLoaded() - { - @Override - public void onLoaded(AccordCacheEntry state, V value, Throwable fail) - { - if (fail == null) state.loaded(value); - else state.failedToLoad(); - } - }; + if (this.onSuccess == null) + this.onSuccess = new ArrayList<>(); + this.onSuccess.add(onSuccess); } } - public interface OnSaved + public interface SaveExecutor { - void onSaved(AccordCacheEntry state, Object identity, Throwable fail); - - static OnSaved immediate() - { - return new OnSaved() - { - @Override - public void onSaved(AccordCacheEntry state, Object identity, Throwable fail) - { - state.saved(identity, fail); - } - }; - } + Cancellable save(AccordCacheEntry saving, UniqueSave identity, Runnable save); } - public

Loading load(BiFunction loadExecutor, P param, Adapter adapter, OnLoaded onLoaded) + public Loading load(LoadExecutor loadExecutor, P1 p1, P2 p2) { Invariants.require(is(WAITING_TO_LOAD), "%s", this); WaitingToLoad cur = (WaitingToLoad)state; - Collection> waiters = cur.waiters; - Loading loading = cur.load(loadExecutor.apply(param, new Runnable() - { - @Override - public void run() - { - V result; - try - { - result = adapter.load(owner.commandStore, key); - } - catch (Throwable t) - { - onLoaded.onLoaded(AccordCacheEntry.this, null, t); - throw t; - } - onLoaded.onLoaded(AccordCacheEntry.this, result, null); - } - - @Override - public String toString() - { - return "Loading " + key + " on " + owner.commandStore + " for " + waiters; - } - })); + Loading loading = cur.load(loadExecutor.load(p1, p2, this)); setStatus(LOADING); state = loading; return loading; @@ -409,11 +403,20 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode return; Saving cancel = is(SAVING) ? ((Saving)state) : null; - setStatus(MODIFIED); - state = value; + if (is(WAITING_TO_SAVE)) + { + ((WaitingToSave) state).state = value; + if (owner.parent().adapter().canSave(value, null)) + save(); + } + else + { + setStatus(MODIFIED); + state = value; + } updateSize(owner.parent()); - // TODO (expected): do we want to always cancel in-progress saving? - if (cancel != null) + // TODO (expected): do we want to cancel in-progress saving? + if (cancel != null && cancel.identity.onSuccess == null) cancel.saving.cancel(); } @@ -463,48 +466,48 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode return state == null; } + boolean saveWhenReady() + { + V full = isShrunk() ? null : (V)state; + Object shrunk = isShrunk() ? state : null; + if (owner.parent().adapter().canSave(full, shrunk)) + return save(); + + setStatus(WAITING_TO_SAVE); + UniqueSave identity = new UniqueSave(); + state = new WaitingToSave<>(identity, state); + return true; + } + /** * Submits a save runnable to the specified executor. When the runnable * has completed, the state save will have either completed or failed. */ @VisibleForTesting - void save(Function saveExecutor, Adapter adapter, OnSaved onSaved) + boolean save() { + WaitingToSave waitingToSave = is(WAITING_TO_SAVE) ? (WaitingToSave)state : null; + Object state = waitingToSave == null ? this.state : waitingToSave.state; V full = isShrunk() ? null : (V)state; Object shrunk = isShrunk() ? state : null; - Runnable save = adapter.save(owner.commandStore, key, full, shrunk); + Runnable save = owner.parent().adapter().save(owner.commandStore, key, full, shrunk); + + UniqueSave identity = waitingToSave == null ? new UniqueSave() : waitingToSave.identity; if (null == save) // null mutation -> null Runnable -> no change on disk { setStatus(LOADED); + if (waitingToSave != null) + this.state = state; + if (identity.onSuccess != null) + identity.onSuccess.forEach(Runnable::run); + return false; } else { setStatus(SAVING); - Object identity = new Object(); - Cancellable saving = saveExecutor.apply(new Runnable() - { - @Override - public void run() - { - try - { - save.run(); - } - catch (Throwable t) - { - onSaved.onSaved(AccordCacheEntry.this, identity, t); - throw t; - } - onSaved.onSaved(AccordCacheEntry.this, identity, null); - } - - @Override - public String toString() - { - return "Saving " + key + " for Accord cache eviction"; - } - }); - state = new Saving(saving, identity, state); + Cancellable saving = owner.parent().parent().saveExecutor.save(this, identity, save); + this.state = new Saving(saving, identity, state); + return true; } } @@ -537,9 +540,9 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode setStatus(LOADED); } - public Cancellable saving() + public SavingOrWaitingToSave savingOrWaitingToSave() { - return ((Saving)state).saving; + return (SavingOrWaitingToSave) state; } public AccordCacheEntry evicted() @@ -558,15 +561,16 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode private boolean tryShrink(K key, Adapter adapter) { - Invariants.require(!isNested()); if (isShrunk() || state == null) return false; - Object update = adapter.fullShrink(key, (V)state); - if (update == null || update == state) + Object cur = unwrap(); + Object upd = adapter.fullShrink(key, (V)cur); + if (upd == null || upd == cur) return false; - state = update; + if (isNested()) ((Nested)this.state).state = upd; + else this.state = upd; status |= SHRUNK; return true; } @@ -665,19 +669,36 @@ public class AccordCacheEntry extends IntrusiveLinkedListNode Object state; } - static class Saving extends Nested + static class SavingOrWaitingToSave extends Nested { - final Cancellable saving; - final Object identity; + final UniqueSave identity; - Saving(Cancellable saving, Object identity, Object state) + SavingOrWaitingToSave(UniqueSave identity, Object state) { - this.saving = saving; this.identity = identity; this.state = state; } } + static class Saving extends SavingOrWaitingToSave + { + final Cancellable saving; + + Saving(Cancellable saving, UniqueSave identity, Object state) + { + super(identity, state); + this.saving = saving; + } + } + + static class WaitingToSave extends SavingOrWaitingToSave + { + WaitingToSave(UniqueSave identity, Object state) + { + super(identity, state); + } + } + static class FailedToSave extends Nested { final Throwable cause; diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index e25f7a3827..7345500a91 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -19,8 +19,10 @@ package org.apache.cassandra.service.accord; import java.util.List; +import java.util.Map; import java.util.NavigableMap; import java.util.Objects; +import java.util.TreeMap; import java.util.concurrent.Callable; import java.util.concurrent.Executor; import java.util.concurrent.TimeUnit; @@ -34,14 +36,18 @@ import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import accord.api.Agent; import accord.api.DataStore; import accord.api.Journal; import accord.api.LocalListeners; import accord.api.ProgressLog; import accord.api.RoutingKey; -import accord.impl.AbstractLoader; +import accord.impl.AbstractReplayer; import accord.impl.AbstractSafeCommandStore.CommandStoreCaches; +import accord.impl.progresslog.DefaultProgressLog; import accord.local.Command; import accord.local.CommandStore; import accord.local.CommandStores; @@ -61,6 +67,9 @@ import accord.primitives.TxnId; import accord.utils.Invariants; import accord.utils.async.AsyncChain; import accord.utils.async.AsyncChains; +import accord.utils.async.AsyncResults.CountingResult; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.accord.AccordKeyspace.CommandsForKeyAccessor; import org.apache.cassandra.service.accord.IAccordService.AccordCompactionInfo; @@ -76,6 +85,8 @@ import static accord.utils.Invariants.require; public class AccordCommandStore extends CommandStore { + private static final Logger logger = LoggerFactory.getLogger(AccordCommandStore.class); + // TODO (required): track this via a PhantomReference, so that if we remove a CommandStore without clearing the caches we can be sure to release them public static class Caches { @@ -140,7 +151,7 @@ public class AccordCommandStore extends CommandStore static final AtomicLong nextSafeRedundantBeforeTicket = new AtomicLong(); public final String loggingId; - private final Journal journal; + public final Journal journal; private final RangeSearcher rangeSearcher; private final AccordExecutor sharedExecutor; final AccordExecutor.SequentialExecutor exclusiveExecutor; @@ -152,8 +163,6 @@ public class AccordCommandStore extends CommandStore private AccordSafeCommandStore current; - private final AccordCommandStoreLoader loader; - public AccordCommandStore(int id, NodeCommandStoreService node, Agent agent, @@ -169,6 +178,8 @@ public class AccordCommandStore extends CommandStore this.journal = journal; this.rangeSearcher = RangeSearcher.extractRangeSearcher(journal); this.sharedExecutor = sharedExecutor; + if (this.progressLog() instanceof DefaultProgressLog) + ((DefaultProgressLog)this.progressLog).setMaxConcurrency(DatabaseDescriptor.getAccordProgressLogMaxConcurrency()); final AccordCache.Type.Instance commands; final AccordCache.Type.Instance commandsForKey; @@ -181,7 +192,6 @@ public class AccordCommandStore extends CommandStore this.exclusiveExecutor = sharedExecutor.executor(); this.commandsForRanges = new CommandsForRanges.Manager(this); - this.loader = new AccordCommandStoreLoader(this); maybeLoadRedundantBefore(journal.loadRedundantBefore(id())); maybeLoadBootstrapBeganAt(journal.loadBootstrapBeganAt(id())); @@ -216,8 +226,6 @@ public class AccordCommandStore extends CommandStore @Override public void markShardDurable(SafeCommandStore safeStore, TxnId globalSyncId, Ranges ranges, Status.Durability durability) { - if (durability.compareTo(Status.Durability.UniversalOrInvalidated) >= 0) - dataStore.snapshot(ranges, globalSyncId); super.markShardDurable(safeStore, globalSyncId, ranges, durability); if (durability.compareTo(Status.Durability.UniversalOrInvalidated) >= 0) commandsForRanges.gcBefore(globalSyncId, ranges); @@ -484,9 +492,60 @@ public class AccordCommandStore extends CommandStore return rangeSearcher; } - public AccordCommandStoreLoader loader() + public AccordCommandStoreReplayer replayer() { - return loader; + return new AccordCommandStoreReplayer(this); + } + + static final AtomicLong nextDurabilityLoggingId = new AtomicLong(); + @Override + protected void ensureDurable(Ranges ranges, RedundantBefore onCommandStoreDurable) + { + long reportId = nextDurabilityLoggingId.incrementAndGet(); + logger.info("{} awaiting local metadata durability for {} ({})", this, ranges, reportId); + executor().afterSubmittedAndConsequences(() -> { + logger.debug("{}: saving intersecting keys ({})", this, reportId); + class Ready extends CountingResult implements Runnable + { + public Ready() { super(1); } + @Override public void run() { decrement(); } + } + + Map saving = new TreeMap<>(); + Map excluding = new TreeMap<>(); + Ready ready = new Ready(); + try (ExclusiveCaches caches = lockCaches()) + { + for (AccordCacheEntry e : caches.commandsForKeys()) + { + if (ranges.contains(e.key()) && e.isModified()) + { + ready.increment(); + caches.global().saveWhenReadyExclusive(e, ready); + } + } + } + + ready.begin((success, fail) -> { + if (fail != null) + { + logger.error("{}: failed to ensure durability of {} ({})", this, ranges, reportId, fail); + } + else + { + logger.debug("{}: waiting for CommandsForKey to flush ({})", this, reportId); + ColumnFamilyStore cfs = AccordKeyspace.AccordColumnFamilyStores.commandsForKey; + + AccordDurableOnFlush onFlush = null; + while (onFlush == null) + onFlush = cfs.getCurrentMemtable().ensureFlushListener(AccordDataStore.FlushListenerKey.KEY, AccordDurableOnFlush::new); + + if (!onFlush.add(id, onCommandStoreDurable)) + AccordDurableOnFlush.notify(cfs.metadata(), this, onCommandStoreDurable); + } + }); + ready.decrement(); + }); } @VisibleForTesting @@ -495,19 +554,28 @@ public class AccordCommandStore extends CommandStore super.unsafeUpsertRedundantBefore(addRedundantBefore); } - public static class AccordCommandStoreLoader extends AbstractLoader + public static class AccordCommandStoreReplayer extends AbstractReplayer { private final AccordCommandStore store; - private AccordCommandStoreLoader(AccordCommandStore store) + private AccordCommandStoreReplayer(AccordCommandStore store) { + super(store.unsafeGetRedundantBefore()); this.store = store; } @Override - public AsyncChain load(TxnId txnId) + public AsyncChain replay(TxnId txnId) { + if (!maybeShouldReplay(txnId)) + { + return AsyncChains.success(null); + } + return store.submit(PreLoadContext.contextFor(txnId, "Replay"), safeStore -> { + if (!shouldReplay(txnId, safeStore.unsafeGet(txnId).current().participants())) + return null; + initialiseState(safeStore, txnId); return safeStore.unsafeGet(txnId).current().route(); }); diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java index 3fa4bb0bdc..18b9bccee2 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java @@ -35,7 +35,6 @@ import accord.primitives.Range; import accord.topology.Topology; import accord.utils.RandomSource; import org.apache.cassandra.cache.CacheSize; -import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.AccordSpec.QueueShardModel; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.metrics.AccordCacheMetrics; @@ -50,7 +49,6 @@ import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueSubmi 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; -import static org.apache.cassandra.service.accord.AccordExecutor.constantFactory; public class AccordCommandStores extends CommandStores implements CacheSize { @@ -109,13 +107,10 @@ public class AccordCommandStores extends CommandStores implements CacheSize { case THREAD_PER_SHARD: case THREAD_PER_SHARD_SYNC_QUEUE: - executors[id] = factory.get(id, shardModel == THREAD_PER_SHARD ? RUN_WITHOUT_LOCK : RUN_WITH_LOCK, 1, constant(baseName + ']'), metrics, constantFactory(Stage.READ.executor()), constantFactory(Stage.MUTATION.executor()), constantFactory(Stage.READ.executor()), agent); + executors[id] = factory.get(id, shardModel == THREAD_PER_SHARD ? RUN_WITHOUT_LOCK : RUN_WITH_LOCK, 1, constant(baseName + ']'), metrics, agent); break; case THREAD_POOL_PER_SHARD: - executors[id] = factory.get(id, RUN_WITHOUT_LOCK, threads, num -> baseName + ',' + num + ']', metrics, AccordExecutor::submitIOToSelf, AccordExecutor::submitIOToSelf, AccordExecutor::submitIOToSelf, agent); - break; - case THREAD_POOL_PER_SHARD_EXCLUDES_IO: - executors[id] = factory.get(id, RUN_WITHOUT_LOCK, threads, num -> baseName + ',' + num + ']', metrics, constantFactory(Stage.READ.executor()), constantFactory(Stage.MUTATION.executor()), constantFactory(Stage.READ.executor()), agent); + executors[id] = factory.get(id, RUN_WITHOUT_LOCK, threads, num -> baseName + ',' + num + ']', metrics, agent); break; } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordDataStore.java b/src/java/org/apache/cassandra/service/accord/AccordDataStore.java index 993e713115..507f0f91cd 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordDataStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordDataStore.java @@ -18,44 +18,27 @@ package org.apache.cassandra.service.accord; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.concurrent.TimeUnit; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import accord.api.DataStore; +import accord.local.CommandStore; import accord.local.Node; +import accord.local.RedundantBefore; import accord.local.SafeCommandStore; import accord.primitives.Range; import accord.primitives.Ranges; import accord.primitives.SyncPoint; -import accord.primitives.TxnId; -import accord.utils.async.AsyncResult; -import accord.utils.async.AsyncResults; -import org.agrona.collections.Object2ObjectHashMap; -import org.apache.cassandra.concurrent.ScheduledExecutors; -import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; -import org.apache.cassandra.db.DataRange; -import org.apache.cassandra.db.Keyspace; -import org.apache.cassandra.db.filter.ColumnFilter; -import org.apache.cassandra.db.lifecycle.View; import org.apache.cassandra.db.memtable.Memtable; -import org.apache.cassandra.db.memtable.TrieMemtable; -import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; -import org.apache.cassandra.dht.Token; -import org.apache.cassandra.io.sstable.SSTableReadsListener; +import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; -import org.apache.cassandra.schema.TableMetadata; -import org.apache.cassandra.tcm.ClusterMetadata; -import org.apache.cassandra.utils.concurrent.Future; -import org.apache.cassandra.utils.concurrent.FutureCombiner; - -import static accord.utils.Invariants.require; -import static org.apache.cassandra.db.ColumnFamilyStore.FlushReason.ACCORD_TXN_GC; public class AccordDataStore implements DataStore { + private static final Logger logger = LoggerFactory.getLogger(AccordDataStore.class); + enum FlushListenerKey { KEY } + @Override public FetchResult fetch(Node node, SafeCommandStore safeStore, Ranges ranges, SyncPoint syncPoint, FetchRanges callback) { @@ -64,121 +47,41 @@ public class AccordDataStore implements DataStore return coordinator.result(); } - static class SnapshotBounds + /** + * Ensures data for the intersecting ranges is flushed to sstable before calling back with reportOnSuccess. + * This is used to gate journal cleanup, since we skip the CommitLog for applying to the data table. + */ + public void ensureDurable(CommandStore commandStore, Ranges ranges, RedundantBefore reportOnSuccess) { - final List> ranges = new ArrayList<>(); - long id; - } - - @Override - public AsyncResult snapshot(Ranges ranges, TxnId before) - { - AsyncResults.SettableResult result = new AsyncResults.SettableResult<>(); - // TODO (desired): maintain a list of Accord tables, perhaps in ClusterMetadata? - ClusterMetadata metadata = ClusterMetadata.current(); - Object2ObjectHashMap tables = new Object2ObjectHashMap<>(); + logger.debug("{} awaiting local data durability of {}", commandStore, ranges); + ColumnFamilyStore prev = null; for (Range range : ranges) { - tables.computeIfAbsent(((TokenRange)range).table(), ignore -> new SnapshotBounds()) - .ranges.add(((TokenRange) range).toKeyspaceRange()); - } - - for (Map.Entry e : tables.entrySet()) - { - // TODO (required): is it safe to ignore null table metadata / cfs? - TableMetadata tableMetadata = metadata.schema.getTableMetadata(e.getKey()); - if (tableMetadata == null || !tableMetadata.isAccordEnabled()) - continue; - - ColumnFamilyStore cfs = Keyspace.openAndGetStoreIfExists(tableMetadata); + ColumnFamilyStore cfs; + if (prev != null && prev.metadata().id.equals(range.prefix())) cfs = prev; + else cfs = Schema.instance.getColumnFamilyStoreInstance((TableId) range.prefix()); if (cfs == null) - continue; - - // TODO (required): when we can safely map TxnId.hlc() -> local timestamp, consult Memtable timestamps - Memtable memtable = cfs.getCurrentMemtable(); - e.getValue().id = memtable.getMemtableId(); - } - - ScheduledExecutors.scheduledTasks.schedule(() -> { - List> futures = new ArrayList<>(); - for (Map.Entry e : tables.entrySet()) { - // TODO (required): is it safe to ignore null tableMetadata (or ColumnFamilyStore below)? - TableMetadata tableMetadata = metadata.schema.getTableMetadata(e.getKey()); - if (tableMetadata == null) continue; + // TODO (expected): should we record this as durable? + continue; + } - ColumnFamilyStore cfs = Keyspace.openAndGetStoreIfExists(tableMetadata); - if (cfs == null) continue; + while (true) + { + Memtable memtable = cfs.getCurrentMemtable(); + AccordDurableOnFlush onFlush = memtable.ensureFlushListener(FlushListenerKey.KEY, AccordDurableOnFlush::new); + if (onFlush != null && onFlush.add(commandStore.id(), reportOnSuccess)) + break; - SnapshotBounds bounds = e.getValue(); - View view = cfs.getTracker().getView(); - for (Memtable memtable : view.getAllMemtables()) + if (cfs == prev) { - if (memtable.getMemtableId() > bounds.id) continue; - if (!intersects(cfs, memtable, bounds.ranges)) continue; - - futures.add(cfs.forceFlush(ACCORD_TXN_GC)); + // we must already have a successful notify, so just propagate + AccordDurableOnFlush.notify(cfs.metadata(), commandStore, reportOnSuccess); break; } } - FutureCombiner.allOf(futures).addCallback((objects, throwable) -> { - if (throwable != null) - result.setFailure(throwable); - else - result.setSuccess(null); - }); - }, DatabaseDescriptor.getAccordGCDelay(TimeUnit.MILLISECONDS), TimeUnit.MILLISECONDS); - - return result; - } - - private boolean intersects(ColumnFamilyStore cfs, Memtable memtable, List> tableRanges) - { - boolean intersects = false; - // TrieMemtable doesn't support reverse iteration so can't find the last token - if (memtable instanceof TrieMemtable) - intersects = true; - else - { - Token firstToken = null; - try (UnfilteredPartitionIterator iterator = memtable.partitionIterator(ColumnFilter.all(cfs.metadata()), DataRange.allData(cfs.getPartitioner()), SSTableReadsListener.NOOP_LISTENER)) - { - if (iterator.hasNext()) - firstToken = iterator.next().partitionKey().getToken(); - } - Token lastToken = memtable.lastToken(); - - if (firstToken != null) - { - require(lastToken != null); - if (firstToken.equals(lastToken)) - { - for (org.apache.cassandra.dht.Range tableRange : tableRanges) - { - if (tableRange.contains(firstToken)) - { - intersects = true; - break; - } - } - } - else - { - require(firstToken.compareTo(lastToken) < 0); - org.apache.cassandra.dht.Range memtableRange = new org.apache.cassandra.dht.Range<>(firstToken, lastToken); - for (org.apache.cassandra.dht.Range tableRange : tableRanges) - { - if (tableRange.intersects(memtableRange)) - { - intersects = true; - break; - } - } - } - } + prev = cfs; } - - return intersects; } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordDurableOnFlush.java b/src/java/org/apache/cassandra/service/accord/AccordDurableOnFlush.java new file mode 100644 index 0000000000..96cf75d0cd --- /dev/null +++ b/src/java/org/apache/cassandra/service/accord/AccordDurableOnFlush.java @@ -0,0 +1,76 @@ +/* + * 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.Map; +import java.util.function.Consumer; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import accord.local.CommandStore; +import accord.local.CommandStores; +import accord.local.PreLoadContext; +import accord.local.RedundantBefore; +import org.agrona.collections.Int2ObjectHashMap; +import org.apache.cassandra.schema.TableMetadata; + +class AccordDurableOnFlush implements Consumer +{ + private static final Logger logger = LoggerFactory.getLogger(AccordDurableOnFlush.class); + + private Int2ObjectHashMap commandStores = new Int2ObjectHashMap<>(); + + AccordDurableOnFlush() + { + } + + synchronized boolean add(int commandStoreId, RedundantBefore reportOnFlush) + { + if (commandStores == null) + return false; + commandStores.merge(commandStoreId, reportOnFlush, RedundantBefore::merge); + return true; + } + + @Override + public void accept(TableMetadata metadata) + { + Int2ObjectHashMap notify; + synchronized (this) + { + notify = commandStores; + commandStores = null; + } + CommandStores commandStores = AccordService.instance().node().commandStores(); + for (Map.Entry e : notify.entrySet()) + { + RedundantBefore durable = e.getValue(); + notify(metadata, commandStores.forId(e.getKey()), durable); + } + } + + static void notify(TableMetadata metadata, CommandStore commandStore, RedundantBefore report) + { + logger.debug("Reporting flush of {}/{}; reporting {} to {}", metadata.id, metadata, report, commandStore); + commandStore.execute((PreLoadContext.Empty) () -> "Report Durable", safeStore -> { + safeStore.upsertRedundantBefore(report); + }); + } +} diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index 9710cc34d9..63370cf5f9 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -18,8 +18,10 @@ package org.apache.cassandra.service.accord; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.List; +import java.util.Queue; import java.util.concurrent.Callable; import java.util.concurrent.CancellationException; import java.util.concurrent.TimeUnit; @@ -36,6 +38,9 @@ import javax.annotation.Nullable; import com.google.common.base.Functions; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import accord.api.Agent; import accord.api.AsyncExecutor; import accord.api.RoutingKey; @@ -58,20 +63,24 @@ import accord.utils.async.Cancellable; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.concurrent.DebuggableTask; import org.apache.cassandra.concurrent.DebuggableTask.DebuggableTaskRunner; -import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.concurrent.Shutdownable; import org.apache.cassandra.metrics.AccordCacheMetrics; +import org.apache.cassandra.service.accord.AccordCacheEntry.LoadExecutor; +import org.apache.cassandra.service.accord.AccordCacheEntry.SaveExecutor; +import org.apache.cassandra.service.accord.AccordCacheEntry.UniqueSave; import org.apache.cassandra.utils.MonotonicClock; import org.apache.cassandra.utils.concurrent.AsyncPromise; import org.apache.cassandra.utils.concurrent.Condition; import org.apache.cassandra.utils.concurrent.Future; +import static accord.utils.Invariants.createIllegalState; import static org.apache.cassandra.service.accord.AccordCache.CommandAdapter.COMMAND_ADAPTER; import static org.apache.cassandra.service.accord.AccordCache.CommandsForKeyAdapter.CFK_ADAPTER; import static org.apache.cassandra.service.accord.AccordCache.registerJfrListener; import static org.apache.cassandra.service.accord.AccordCacheEntry.Status.EVICTED; import static org.apache.cassandra.service.accord.AccordTask.State.LOADING; +import static org.apache.cassandra.service.accord.AccordTask.State.RUNNING; import static org.apache.cassandra.service.accord.AccordTask.State.SCANNING_RANGES; import static org.apache.cassandra.service.accord.AccordTask.State.WAITING_TO_LOAD; import static org.apache.cassandra.service.accord.AccordTask.State.WAITING_TO_RUN; @@ -79,32 +88,33 @@ import static org.apache.cassandra.service.accord.AccordTask.State.WAITING_TO_RU /** * NOTE: We assume that NO BLOCKING TASKS are submitted to this executor AND WAITED ON by another task executing on this executor. */ -public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLoaded, AccordCacheEntry.OnSaved, Shutdownable, AsyncExecutor +public abstract class AccordExecutor implements CacheSize, LoadExecutor, Boolean>, SaveExecutor, Shutdownable, AsyncExecutor { + private static final Logger logger = LoggerFactory.getLogger(AccordExecutor.class); public interface AccordExecutorFactory { - AccordExecutor get(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent); + AccordExecutor get(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent); } public enum Mode { RUN_WITH_LOCK, RUN_WITHOUT_LOCK } - public interface ExecutorFunction extends BiFunction {} - public interface ExecutorFunctionFactory extends Function {} - // WARNING: this is a shared object, so close is NOT idempotent public static final class ExclusiveGlobalCaches extends GlobalCaches implements AutoCloseable { final Lock lock; + final AccordExecutor executor; - public ExclusiveGlobalCaches(Lock lock, AccordCache global, AccordCache.Type commands, AccordCache.Type commandsForKey) + public ExclusiveGlobalCaches(AccordExecutor executor, AccordCache global, AccordCache.Type commands, AccordCache.Type commandsForKey) { super(global, commands, commandsForKey); - this.lock = lock; + this.lock = executor.lock; + this.executor = executor; } @Override public void close() { + executor.beforeUnlock(); lock.unlock(); } } @@ -127,21 +137,39 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo final Agent agent; final int executorId; private final AccordCache cache; - private final ExecutorFunction loadExecutor; - private final ExecutorFunction rangeLoadExecutor; private final TaskQueue> scanningRanges = new TaskQueue<>(SCANNING_RANGES); // never queried, just parked here while scanning private final TaskQueue> loading = new TaskQueue<>(LOADING); // never queried, just parked here while loading + private final TaskQueue 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); - private final AccordCacheEntry.OnLoaded onRangeLoaded = this::onRangeLoaded; private final ExclusiveGlobalCaches caches; private List waitingForQuiescence; + private Queue waitingForCompletion; + + private static class WaitForCompletion + { + final int position; + int maybeNotify; + final Runnable run; + + private WaitForCompletion(int position, Runnable run) + { + this.position = position; + this.maybeNotify = position - 1; + this.run = run; + } + + public String toString() + { + return run.toString() + " @" + position; + } + } /** * The maximum total number of loads we can queue at once - this includes loads for range transactions, @@ -159,15 +187,13 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo private int activeLoads, activeRangeLoads; private boolean hasPausedLoading; int tasks; - int running; + int runningThreads; - AccordExecutor(Lock lock, int executorId, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + AccordExecutor(Lock lock, int executorId, AccordCacheMetrics metrics, Agent agent) { this.lock = lock; this.executorId = executorId; - this.cache = new AccordCache(alwaysNullTask(saveExecutor.apply(this)), this, 0, metrics); - this.loadExecutor = loadExecutor.apply(this); - this.rangeLoadExecutor = rangeLoadExecutor.apply(this); + this.cache = new AccordCache(this, 0, metrics); this.agent = agent; final AccordCache.Type commands; @@ -178,7 +204,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo commandsForKey = cache.newType(RoutingKey.class, CFK_ADAPTER); registerJfrListener(executorId, commandsForKey, "CommandsForKey"); - this.caches = new ExclusiveGlobalCaches(lock, cache, commands, commandsForKey); + this.caches = new ExclusiveGlobalCaches(this, cache, commands, commandsForKey); ScheduledExecutors.scheduledFastTasks.scheduleAtFixedRate(() -> { executeDirectlyWithLock(cache::processNoEvictQueue); }, 1L, 1L, TimeUnit.SECONDS); @@ -216,7 +242,10 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo Task pollWaitingToRunExclusive() { updateWaitingToRunExclusive(); - return waitingToRun.poll(); + Task next = waitingToRun.poll(); + if (next != null) + next.addToQueue(running); + return next; } void updateWaitingToRunExclusive() @@ -235,7 +264,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo lock.lock(); try { - if (tasks == 0 && running == 0) + if (tasks == 0 && runningThreads == 0) return; if (waitingForQuiescence == null) @@ -250,13 +279,44 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo condition.awaitThrowUncheckedOnInterrupt(); } - protected void signalQuiescentExclusive() + protected void notifyQuiescentExclusive() { if (waitingForQuiescence != null) { waitingForQuiescence.forEach(Condition::signalAll); waitingForQuiescence = null; } + if (waitingForCompletion != null) + { + logger.warn("{} processed all pending tasks (<{}) but found waiting: {}", this, nextPosition, waitingForCompletion); + waitingForCompletion.forEach(w -> w.run.run()); + waitingForCompletion = null; + } + } + + public void afterSubmittedAndConsequences(Runnable run) + { + lock.lock(); + try + { + if (tasks == 0 && runningThreads == 0) + { + run.run(); + return; + } + + if (waitingForCompletion != null) // escape hatch for some bug that means we lose a notification for a given task's queue position + maybeNotifyWaitingForCompletion(); + if (waitingForCompletion == null) + waitingForCompletion = new ArrayDeque<>(); + + int position = nextPosition; + waitingForCompletion.add(new WaitForCompletion(position, run)); + } + finally + { + lock.unlock(); + } } void maybeUnpauseLoading() @@ -272,6 +332,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } public abstract boolean hasTasks(); + abstract void beforeUnlock(); abstract boolean isOwningThread(); private void enqueueLoadsExclusive() @@ -297,7 +358,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo { default: { - failExclusive(next, new AssertionError("Unexpected state: " + next.toDescription())); + failExclusive(next, createIllegalState("Unexpected state: " + next.toDescription())); break; } case WAITING_TO_SCAN_RANGES: @@ -309,7 +370,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo { ++activeRangeLoads; ++activeLoads; - next.rangeScanner().start(rangeLoadExecutor); + next.rangeScanner().start(this); updateQueue(next); } break; @@ -326,15 +387,11 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } Invariants.require(load != null); - AccordCacheEntry.OnLoaded onLoaded = this; ++activeLoads; if (isForRange) - { ++activeRangeLoads; - onLoaded = onRangeLoaded; - } - for (AccordTask task : cache.load(loadExecutor, next, load, onLoaded)) + for (AccordTask task : cache.load(this, next, isForRange, load)) { if (task == next) continue; if (task.onLoading(load)) @@ -448,28 +505,6 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo task.addToQueue(queue == null ? waitingToRun : queue); } - private Cancellable submitIOExclusive(Task parent, Runnable run) - { - Invariants.require(isOwningThread()); - ++tasks; - IORunnable task = new IORunnable(run); - // TODO (expected): adopt queue position of the submitting task - if (parent == null) assignNewQueuePosition(task); - else assignQueueSubPosition(parent, task); - waitingToRun.append(task); - return task; - } - - private void assignNewQueuePosition(Task task) - { - task.queuePosition = (((long)++nextPosition) & 0xffffffffL) << 31; - } - - private void assignQueueSubPosition(Task parent, Task task) - { - task.queuePosition = parent.queuePosition | (++nextPosition & 0x7fffffff); - } - public SequentialExecutor executor() { return new SequentialExecutor(); @@ -480,11 +515,6 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo return new SequentialExecutor(); } - public void submit(AccordTask operation) - { - submit(AccordExecutor::submitExclusive, Function.identity(), operation); - } - public void cancel(AccordTask task) { Invariants.require(task.commandStore.executor() == this, @@ -493,25 +523,16 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo submit(AccordExecutor::cancelExclusive, CancelAsync::new, task); } - public void onScannedRanges(AccordTask task, Throwable fail) + @Override + public Cancellable load(AccordTask parent, Boolean isForRange, AccordCacheEntry entry) { - submit(AccordExecutor::onScannedRangesExclusive, OnScannedRanges::new, task, fail); - } - - public void onSaved(AccordCacheEntry saved, Object identity, Throwable fail) - { - submit(AccordExecutor::onSavedExclusive, OnSaved::new, saved, identity, fail); + return submitPlainExclusive(parent, newLoad(entry, isForRange)); } @Override - public void onLoaded(AccordCacheEntry loaded, V value, Throwable fail) + public Cancellable save(AccordCacheEntry entry, UniqueSave identity, Runnable save) { - submit(AccordExecutor::onLoadedExclusive, OnLoaded::new, loaded, value, fail, false); - } - - public void onRangeLoaded(AccordCacheEntry loaded, V value, Throwable fail) - { - submit(AccordExecutor::onLoadedExclusive, OnLoaded::new, loaded, value, fail, true); + return submitPlainExclusive(null, new SaveRunnable(entry, identity, save)); } private void submit(BiConsumer sync, Function async, P1 p1) @@ -536,27 +557,103 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo abstract void submit(QuintConsumer sync, QuadFunction async, P1s p1s, P1a p1a, P2 p2, P3 p3, P4 p4); - private void submitExclusive(Plain task) + void submit(AccordTask operation) { - ++tasks; - task.queuePosition = ++nextPosition; - waitingToRun(task, task.executor()); - } - - void submitExclusive(Runnable run) - { - submitExclusive(new PlainRunnable(run)); + submit(AccordExecutor::submitExclusive, Function.identity(), operation); } void submitExclusive(AccordTask task) { + assignQueuePosition(task); ++tasks; - assignNewQueuePosition(task); task.setupExclusive(); updateQueue(task); enqueueLoadsExclusive(); } + void submitExclusive(Runnable runnable) + { + submitPlainExclusive(new PlainRunnable(null, runnable, null)); + } + + private void submitPlainExclusive(Plain task) + { + ++tasks; + assignQueuePosition(task); + waitingToRun(task, task.executor()); + } + + Cancellable submitPlainExclusive(Task parent, AbstractIOTask task) + { + return submitPlainExclusive(parent, new WrappedIOTask(task)); + } + + T submitPlainExclusive(Task parent, T task) + { + Invariants.require(isOwningThread()); + ++tasks; + if (parent != null) inheritQueuePosition(parent, task); + else assignNewQueuePosition(task); + waitingToRun.append(task); + return task; + } + + private void assignQueuePosition(Task task) + { + if (task.queuePosition == 0) + assignNewQueuePosition(task); + } + + private void assignNewQueuePosition(Task task) + { + if (nextPosition == 0) nextPosition++; + task.queuePosition = nextPosition++; + } + + private void inheritQueuePosition(Task parent, Task task) + { + task.queuePosition = parent.queuePosition; + } + + void completeTaskExclusive(Task task, boolean cleanupTask) + { + --tasks; + + // for integration with SequentialExecutor, we must : + // - first take the position so that represents the just-executed task + // - call cleanup to submit any following task on the relevant sub-queue + // - remove the previous task from the running collection only if still present (SequentialExecutor will have removed it) + int position = task.queuePosition; + if (cleanupTask) task.cleanupExclusive(); + if (running.contains(task)) + running.remove(task); + + if (waitingForCompletion != null && waitingForCompletion.peek().maybeNotify - position >= 0) + maybeNotifyWaitingForCompletion(); + } + + private void maybeNotifyWaitingForCompletion() + { + int min = minPosition(waitingToRun.peek(), + minPosition(waitingToLoad.peek(), + minPosition(waitingToLoadRangeTxns.peek(), + minPosition(running.peek(), + minPosition(loading.peek(), + minPosition(scanningRanges.peek(), Integer.MAX_VALUE)))))); + + while (!waitingForCompletion.isEmpty() && waitingForCompletion.peek().position - min <= 0) + waitingForCompletion.poll().run.run(); + if (waitingForCompletion.isEmpty()) + waitingForCompletion = null; + else + waitingForCompletion.peek().maybeNotify = min; + } + + private static int minPosition(@Nullable Task task, int min) + { + return task == null ? min : Integer.min(task.queuePosition, min); + } + private void cancelExclusive(AccordTask task) { switch (task.state()) @@ -572,9 +669,9 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo case WAITING_TO_LOAD: case WAITING_TO_SCAN_RANGES: case WAITING_TO_RUN: - --tasks; task.unqueueIfQueued(); task.cancelExclusive(); + completeTaskExclusive(task, false); break; case FAILING: @@ -587,7 +684,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } } - private void onScannedRangesExclusive(AccordTask task, Throwable fail) + void onScannedRangesExclusive(AccordTask task, Throwable fail) { --activeLoads; --activeRangeLoads; @@ -612,13 +709,12 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo if (task.state().isExecuted()) return; - --tasks; try { task.failExclusive(fail); } catch (Throwable t) { agent.onUncaughtException(t); } finally { task.unqueueIfQueued(); - task.cleanupExclusive(null); + completeTaskExclusive(task, true); } } @@ -664,28 +760,22 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo public Future submit(Runnable run) { - return submit(run, null); - } - - // TODO (expected): offer queue jumping/priorities - private Future submit(Runnable run, @Nullable SequentialExecutor executor) - { - PlainRunnable task = new PlainRunnable(new AsyncPromise<>(), run, executor); + PlainRunnable task = new PlainRunnable(new AsyncPromise<>(), run, null); submit(task); return task.result; } - private Cancellable submit(Plain task) - { - submit(AccordExecutor::submitExclusive, Functions.identity(), task); - return task; - } - public void execute(Runnable command) { submit(new PlainRunnable(null, command, null)); } + private Cancellable submit(Plain task) + { + submit(AccordExecutor::submitPlainExclusive, Functions.identity(), task); + return task; + } + public void executeDirectlyWithLock(Runnable command) { lock.lock(); @@ -695,6 +785,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } finally { + beforeUnlock(); lock.unlock(); } } @@ -748,6 +839,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo protected static abstract class TaskRunner implements DebuggableTaskRunner { + // TODO (desired): this probably doesn't need to be volatile private volatile Task running; private static final AtomicReferenceFieldUpdater runningUpdater = AtomicReferenceFieldUpdater.newUpdater(TaskRunner.class, Task.class, "running"); @@ -758,12 +850,17 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo return running == null ? null : running.debuggable(); } - void setDebuggable(Task debuggable) + Task runningTask() + { + return running; + } + + void setRunning(Task debuggable) { runningUpdater.lazySet(this, debuggable); } - void clearDebuggable() + void clearRunning() { runningUpdater.lazySet(this, null); } @@ -771,7 +868,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo public static abstract class Task extends IntrusivePriorityHeap.Node { - long queuePosition; + int queuePosition; protected Task() { @@ -782,12 +879,12 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo /** * Prepare to run while holding the state cache lock */ - abstract protected void preRunExclusive(@Nullable TaskRunner runner); + abstract protected void preRunExclusive(); /** * Run the command; the state cache lock may or may not be held depending on the executor implementation */ - abstract protected void run(); + abstract protected void runInternal(); /** * Fail the command; the state cache lock may or may not be held depending on the executor implementation */ @@ -796,7 +893,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo /** * Cleanup the command while holding the state cache lock */ - abstract protected void cleanupExclusive(@Nullable TaskRunner runner); + abstract protected void cleanupExclusive(); abstract protected void addToQueue(TaskQueue queue); } @@ -817,13 +914,13 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } @Override - protected void preRunExclusive(@Nullable TaskRunner runner) + protected void preRunExclusive() { - queue.preRunTask(runner); + queue.preRunTask(); } @Override - protected void run() + protected void runInternal() { queue.runTask(); } @@ -835,15 +932,16 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } @Override - protected void cleanupExclusive(@Nullable TaskRunner runner) + protected void cleanupExclusive() { - queue.cleanupTask(runner); + queue.cleanupTask(); } @Override protected void addToQueue(TaskQueue queue) { - throw new UnsupportedOperationException(); + Invariants.require(queue.kind == RUNNING); + queue.append(this); } } @@ -861,10 +959,10 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo this.selfTask = new SequentialQueueTask(this); } - void preRunTask(@Nullable TaskRunner runner) + void preRunTask() { Invariants.require(task != null); - task.preRunExclusive(runner); + task.preRunExclusive(); running = true; } @@ -878,7 +976,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo LockSupport.park(); waiting = null; } - task.run(); + task.runInternal(); } void failTask(Throwable t) @@ -886,14 +984,15 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo task.fail(t); } - void cleanupTask(@Nullable TaskRunner runner) + void cleanupTask() { - task.cleanupExclusive(runner); + task.cleanupExclusive(); running = false; owner = null; task = super.poll(); if (task != null) { + AccordExecutor.this.running.remove(selfTask); selfTask.queuePosition = task.queuePosition; waitingToRun.append(selfTask); } @@ -905,10 +1004,12 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo { if (task != null) { + Invariants.require(running || waitingToRun.contains(selfTask)); super.append(newTask); } else { + Invariants.require(!running && isEmpty()); task = newTask; selfTask.queuePosition = newTask.queuePosition; waitingToRun.append(selfTask); @@ -948,6 +1049,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } } } + Invariants.require(task == null || running || waitingToRun.contains(selfTask)); } public boolean inExecutor() @@ -970,18 +1072,21 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo @Override protected boolean contains(Task contains) { - throw new UnsupportedOperationException(); + return super.contains(contains) || (task == contains && !running); } @Override - public AsyncChain build(Callable task) + public AsyncChain build(Callable call) { + int position = inExecutor() && task != null ? task.queuePosition : 0; return new AsyncChains.Head<>() { @Override protected Cancellable start(BiConsumer callback) { - return submit(new PlainChain<>(task, callback, SequentialExecutor.this)); + PlainChain submit = new PlainChain<>(call, callback, SequentialExecutor.this); + submit.queuePosition = position; + return AccordExecutor.this.submit(submit); } }; } @@ -989,7 +1094,10 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo @Override public void execute(Runnable run) { - submit(new PlainRunnable(null, run, this)); + PlainRunnable submit = new PlainRunnable(null, run, this); + if (inExecutor() && this.task != null) + submit.queuePosition = this.task.queuePosition; + AccordExecutor.this.submit(submit); } @Override @@ -1071,90 +1179,6 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo abstract void submitExclusive(AccordExecutor executor); } - private static class OnLoaded extends SubmitAsync - { - static final int FAIL = 1; - static final int RANGE = 2; - final AccordCacheEntry loaded; - final Object result; - final int flags; - - OnLoaded(AccordCacheEntry loaded, V success, Throwable fail, boolean isForRange) - { - this.loaded = loaded; - int flags = isForRange ? RANGE : 0; - if (fail == null) - { - result = success; - } - else - { - result = fail; - flags |= FAIL; - } - this.flags = flags; - } - - V success() - { - return (flags & FAIL) == 0 ? (V) result : null; - } - - Throwable fail() - { - return (flags & FAIL) == 0 ? null : (Throwable) result; - } - - boolean isForRange() - { - return (flags & RANGE) != 0; - } - - @Override - void submitExclusive(AccordExecutor executor) - { - executor.onLoadedExclusive(loaded, success(), fail(), isForRange()); - } - } - - private static class OnScannedRanges extends SubmitAsync - { - final AccordTask scanned; - final Throwable fail; - - private OnScannedRanges(AccordTask scanned, Throwable fail) - { - this.scanned = scanned; - this.fail = fail; - } - - @Override - void submitExclusive(AccordExecutor executor) - { - executor.onScannedRangesExclusive(scanned, fail); - } - } - - private static class OnSaved extends SubmitAsync - { - final AccordCacheEntry state; - final Object identity; - final Throwable fail; - - private OnSaved(AccordCacheEntry state, Object identity, Throwable fail) - { - this.state = state; - this.identity = identity; - this.fail = fail; - } - - @Override - void submitExclusive(AccordExecutor executor) - { - executor.onSavedExclusive(state, identity, fail); - } - } - private static class CancelAsync extends SubmitAsync { final AccordTask cancel; @@ -1176,64 +1200,35 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo return ignore -> out; } - static ExecutorFunctionFactory constantFactory(ExecutorFunction exec) - { - return ignore -> exec; - } - - static ExecutorFunctionFactory constantFactory(ExecutorPlus exec) - { - return ignore -> wrap(exec); - } - - static ExecutorFunction wrap(ExecutorPlus exec) - { - return (t, r) -> wrap(exec.submit(r)); - } - - static Cancellable wrap(Future f) - { - return () -> f.cancel(false); - } - - public static ExecutorFunction submitIOToSelf(AccordExecutor executor) - { - return executor::submitIOExclusive; - } - - private static Function alwaysNullTask(ExecutorFunction f) - { - return r -> f.apply(null, r); - } - abstract class Plain extends SubmittableTask implements Cancellable { abstract SequentialExecutor executor(); @Override - protected void preRunExclusive(@Nullable TaskRunner runner) {} + protected void preRunExclusive() {} @Override - protected void cleanupExclusive(@Nullable TaskRunner runner) {} + protected void cleanupExclusive() {} @Override protected final void addToQueue(TaskQueue queue) { - Invariants.require(queue.kind == WAITING_TO_RUN); + Invariants.require(queue.kind == WAITING_TO_RUN || queue.kind == RUNNING); queue.append(this); } @Override - public final void cancel() + public void cancel() { executeDirectlyWithLock(() -> { SequentialExecutor executor = executor(); TaskQueue queue = executor == null ? waitingToRun : executor; if (queue.contains(this)) { - --tasks; queue.remove(this); - fail(new CancellationException()); + try { fail(new CancellationException()); } + catch (Throwable t) { agent.onUncaughtException(t); } + completeTaskExclusive(this, true); } }); } @@ -1241,7 +1236,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo @Override final void submitExclusive(AccordExecutor owner) { - owner.submitExclusive(this); + owner.submitPlainExclusive(this); } } @@ -1264,7 +1259,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } @Override - protected void run() + protected void runInternal() { run.run(); if (result != null) @@ -1286,45 +1281,24 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } } - class IORunnable extends Plain implements Cancellable, DebuggableTask + // a task that may be submitted to this executor or another + abstract class IOTask extends Plain implements Cancellable, DebuggableTask { - // expected to implement toString - final Runnable run; final long createdAtNanos = MonotonicClock.Global.approxTime.now(); long startedAtNanos; - IORunnable(Runnable run) + abstract void postRunExclusive(); + + @Override + protected void preRunExclusive() { - this.run = run; + startedAtNanos = MonotonicClock.Global.approxTime.now(); } @Override - protected void run() + protected void cleanupExclusive() { - run.run(); - } - - @Override - protected void preRunExclusive(@Nullable TaskRunner runner) - { - if (runner != null) - { - startedAtNanos = MonotonicClock.Global.approxTime.now(); - runner.setDebuggable(this); - } - } - - @Override - protected void cleanupExclusive(@Nullable TaskRunner runner) - { - if (runner != null) - runner.clearDebuggable(); - } - - @Override - protected void fail(Throwable t) - { - agent.onUncaughtException(t); + postRunExclusive(); } @Override @@ -1344,11 +1318,148 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo { return startedAtNanos; } + } + + static class FailureHolder + { + static final FailureHolder NOT_STARTED = new FailureHolder(new RuntimeException("Not started")); + + final Throwable fail; + + FailureHolder(Throwable fail) + { + this.fail = fail; + } + } + + LoadRunnable newLoad(AccordCacheEntry entry, boolean isForRange) + { + return isForRange ? new LoadRangeRunnable<>(entry) : new LoadRunnable<>(entry); + } + + class LoadRunnable extends IOTask + { + final AccordCacheEntry entry; + Object result = FailureHolder.NOT_STARTED; + + LoadRunnable(AccordCacheEntry entry) + { + this.entry = entry; + } + + boolean isForRange() { return false; } + + void postRunExclusive() + { + if (!(result instanceof FailureHolder)) onLoadedExclusive(entry, (V)result, null, isForRange()); + else onLoadedExclusive(entry, null, ((FailureHolder)result).fail, isForRange()); + } + + @Override + public void runInternal() + { + result = entry.owner.parent().adapter().load(entry.owner.commandStore, entry.key()); + } + + @Override + protected void fail(Throwable t) + { + result = new FailureHolder(t); + } @Override public String description() { - return run.toString(); + return "Loading " + entry; + } + } + + final class LoadRangeRunnable extends LoadRunnable + { + LoadRangeRunnable(AccordCacheEntry entry) { super(entry); } + @Override boolean isForRange() { return true; } + } + + static abstract class AbstractIOTask + { + abstract protected void runInternal(); + abstract protected void postRunExclusive(); + abstract protected void fail(Throwable t); + abstract protected String description(); + } + + class WrappedIOTask extends IOTask + { + final AbstractIOTask wrapped; + + WrappedIOTask(AbstractIOTask wrap) + { + this.wrapped = wrap; + } + + @Override + protected void runInternal() + { + wrapped.runInternal(); + } + + @Override + void postRunExclusive() + { + wrapped.postRunExclusive(); + } + + @Override + public String description() + { + return wrapped.description(); + } + + @Override + protected void fail(Throwable fail) + { + wrapped.fail(fail); + } + } + + private static final Throwable NOT_STARTED = new Throwable(); + class SaveRunnable extends IOTask + { + final AccordCacheEntry entry; + final UniqueSave identity; + final Runnable run; + Throwable failure = NOT_STARTED; + + SaveRunnable(AccordCacheEntry entry, UniqueSave identity, Runnable run) + { + this.entry = entry; + this.identity = identity; + this.run = run; + } + + @Override + void postRunExclusive() + { + onSavedExclusive(entry, identity, failure); + } + + @Override + public void runInternal() + { + run.run(); + failure = null; + } + + @Override + protected void fail(Throwable t) + { + failure = t; + } + + @Override + public String description() + { + return "Save " + entry; } } @@ -1372,7 +1483,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } @Override - protected void run() + protected void runInternal() { T success; try @@ -1435,18 +1546,14 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo } @Override - protected void preRunExclusive(@Nullable TaskRunner runner) + protected void preRunExclusive() { startedAtNanos = MonotonicClock.Global.approxTime.now(); - if (runner != null) - runner.setDebuggable(this); } @Override - protected void cleanupExclusive(@Nullable TaskRunner runner) + protected void cleanupExclusive() { - if (runner != null) - runner.clearDebuggable(); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java index 63fe294439..8348e14ee5 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractLockLoop.java @@ -37,9 +37,9 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor boolean isHeldByExecutor; boolean shutdown; - AccordExecutorAbstractLockLoop(Lock lock, int executorId, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + AccordExecutorAbstractLockLoop(Lock lock, int executorId, AccordCacheMetrics metrics, Agent agent) { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); } abstract void notifyWork(); @@ -81,13 +81,13 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor public boolean hasTasks() { - if (tasks > 0 || !submitted.isEmpty() || running > 0) + if (tasks > 0 || !submitted.isEmpty() || runningThreads > 0) return true; lock.lock(); try { - return tasks > 0 || !submitted.isEmpty() || running > 0; + return tasks > 0 || !submitted.isEmpty() || runningThreads > 0; } finally { @@ -95,6 +95,13 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor } } + @Override + void beforeUnlock() + { + if (!isInLoop()) + notifyIfMoreWorkExclusive(); + } + void updateWaitingToRunExclusive() { drainSubmittedExclusive(); @@ -125,14 +132,14 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor private void pauseExclusive() { isHeldByExecutor = false; - if (--running == 0 && tasks == 0) - signalQuiescentExclusive(); + if (--runningThreads == 0 && tasks == 0) + notifyQuiescentExclusive(); } private void resumeExclusive() { isHeldByExecutor = true; - ++running; + ++runningThreads; } LoopTask task(String name, Mode mode) @@ -161,11 +168,11 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor if (task != null) { - --tasks; + setRunning(task); try { - task.preRunExclusive(this); - task.run(); + task.preRunExclusive(); + task.runInternal(); } catch (Throwable t) { @@ -173,7 +180,8 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor } finally { - task.cleanupExclusive(this); + completeTaskExclusive(task, true); + clearRunning(); } } else @@ -222,7 +230,13 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor lock.lock(); try { - if (task != null) task.cleanupExclusive(this); + if (task != null) + { + Task tmp = task; + task = null; + completeTaskExclusive(tmp, true); + clearRunning(); + } else resumeExclusive(); enterLockExclusive(); @@ -231,6 +245,8 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor task = pollWaitingToRunExclusive(); if (task != null) { + setRunning(task); + task.preRunExclusive(); exitLockExclusive(); break; } @@ -247,8 +263,6 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor awaitExclusive(); resumeExclusive(); } - --tasks; - task.preRunExclusive(this); } catch (Throwable t) { @@ -256,12 +270,17 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor { try { task.fail(t); } catch (Throwable t2) { t.addSuppressed(t2); } - try { task.cleanupExclusive(this); } + try { completeTaskExclusive(task, true); } catch (Throwable t2) { t.addSuppressed(t2); } try { agent.onUncaughtException(t); } catch (Throwable t2) { /* nothing we can sensibly do after already reporting */ } task = null; } + else + { + try { agent.onUncaughtException(t); } + catch (Throwable t2) { /* nothing we can sensibly do after already reporting */ } + } if (isHeldByExecutor) pauseExclusive(); exitLockExclusive(); @@ -274,7 +293,7 @@ abstract class AccordExecutorAbstractLockLoop extends AccordExecutor try { - task.run(); + task.runInternal(); } catch (Throwable t) { diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java index 42cde814cf..e7f01e1cd6 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAbstractSemiSyncSubmit.java @@ -27,9 +27,9 @@ import org.apache.cassandra.metrics.AccordCacheMetrics; abstract class AccordExecutorAbstractSemiSyncSubmit extends AccordExecutorAbstractLockLoop { - AccordExecutorAbstractSemiSyncSubmit(Lock lock, int executorId, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + AccordExecutorAbstractSemiSyncSubmit(Lock lock, int executorId, AccordCacheMetrics metrics, Agent agent) { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); } abstract void awaitExclusive() throws InterruptedException; diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java index 99086259a9..cefdb46b59 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorAsyncSubmit.java @@ -31,14 +31,14 @@ class AccordExecutorAsyncSubmit extends AccordExecutorAbstractSemiSyncSubmit private final AccordExecutorLoops loops; private final LockWithAsyncSignal lock; - public AccordExecutorAsyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + public AccordExecutorAsyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(new LockWithAsyncSignal(), executorId, mode, threads, name, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + this(new LockWithAsyncSignal(), executorId, mode, threads, name, metrics, agent); } - private AccordExecutorAsyncSubmit(LockWithAsyncSignal lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + private AccordExecutorAsyncSubmit(LockWithAsyncSignal lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); this.lock = lock; this.loops = new AccordExecutorLoops(mode, threads, name, this::task); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java index 803839de94..ff9169d08b 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java @@ -85,6 +85,13 @@ class AccordExecutorLoops }; } + AccordExecutor.Task activeTaskForCurrentThread() + { + Thread thread = Thread.currentThread(); + LoopTask task = tasks.get(thread.getId()); + return task == null ? null : task.runningTask(); + } + public Stream active() { return tasks.values().stream(); diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java index bd4b494646..c7677f8424 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSemiSyncSubmit.java @@ -33,14 +33,14 @@ class AccordExecutorSemiSyncSubmit extends AccordExecutorAbstractSemiSyncSubmit private final ReentrantLock lock; private final Condition hasWork; - public AccordExecutorSemiSyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + public AccordExecutorSemiSyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(new ReentrantLock(), executorId, mode, threads, name, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + this(new ReentrantLock(), executorId, mode, threads, name, metrics, agent); } - private AccordExecutorSemiSyncSubmit(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + private AccordExecutorSemiSyncSubmit(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); this.lock = lock; this.hasWork = lock.newCondition(); this.loops = new AccordExecutorLoops(mode, threads, name, this::task); diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java index 7e74e09b32..bbeb165d50 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSimple.java @@ -27,7 +27,6 @@ import accord.utils.Invariants; import accord.utils.QuadFunction; import accord.utils.QuintConsumer; import org.apache.cassandra.concurrent.ExecutorPlus; -import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.metrics.AccordCacheMetrics; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; @@ -36,37 +35,28 @@ class AccordExecutorSimple extends AccordExecutor { final ExecutorPlus executor; final ReentrantLock lock; + private Task active; public AccordExecutorSimple(int executorId, String name, AccordCacheMetrics metrics, Agent agent) { - this(executorId, name, metrics, Stage.READ.executor(), Stage.MUTATION.executor(), Stage.READ.executor(), agent); + this(new ReentrantLock(), executorId, name, metrics, agent); } - public AccordExecutorSimple(int executorId, String name, AccordCacheMetrics metrics, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor, ExecutorPlus rangeLoadExecutor, Agent agent) + public AccordExecutorSimple(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(executorId, name, metrics, wrap(loadExecutor), wrap(saveExecutor), wrap(rangeLoadExecutor), agent); + this(new ReentrantLock(), executorId, mode, threads, name, metrics, agent); } - public AccordExecutorSimple(int executorId, String name, AccordCacheMetrics metrics, ExecutorFunction loadExecutor, ExecutorFunction saveExecutor, ExecutorFunction rangeLoadExecutor, Agent agent) + private AccordExecutorSimple(ReentrantLock lock, int executorId, String name, AccordCacheMetrics metrics, Agent agent) { - this(new ReentrantLock(), executorId, name, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); - } - - private AccordExecutorSimple(ReentrantLock lock, int executorId, String name, AccordCacheMetrics metrics, ExecutorFunction loadExecutor, ExecutorFunction saveExecutor, ExecutorFunction rangeLoadExecutor, Agent agent) - { - super(lock, executorId, metrics, constantFactory(loadExecutor), constantFactory(saveExecutor), constantFactory(rangeLoadExecutor), agent); + super(lock, executorId, metrics, agent); this.lock = lock; this.executor = executorFactory().sequential(name); } - public AccordExecutorSimple(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + public AccordExecutorSimple(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(new ReentrantLock(), executorId, mode, threads, name, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); - } - - public AccordExecutorSimple(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) - { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); Invariants.requireArgument(threads == 1); this.lock = lock; this.executor = executorFactory().sequential(name.apply(0)); @@ -78,35 +68,42 @@ class AccordExecutorSimple extends AccordExecutor return tasks + executor.getActiveTaskCount() + executor.getPendingTaskCount() > 0; } + @Override + void beforeUnlock() + { + if (hasWaitingToRun()) + executor.execute(this::run); + } + protected void run() { lock.lock(); try { - running = 1; + runningThreads = 1; while (true) { Task task = pollWaitingToRunExclusive(); + active = task; if (task == null) { - running = 0; - signalQuiescentExclusive(); + runningThreads = 0; + notifyQuiescentExclusive(); return; } - --tasks; - try { task.preRunExclusive(null); task.run(); } + try { task.preRunExclusive(); task.runInternal(); } catch (Throwable t) { task.fail(t); } - finally { task.cleanupExclusive(null); } + finally + { + completeTaskExclusive(task, true); + active = null; + } } } - catch (Throwable t) - { - throw t; - } finally { - running = 0; + runningThreads = 0; if (hasWaitingToRun()) executor.execute(this::run); lock.unlock(); diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java index 3bb07b5f2e..3e94a21fac 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorSyncSubmit.java @@ -26,7 +26,6 @@ import java.util.function.IntFunction; import accord.api.Agent; import accord.utils.QuadFunction; import accord.utils.QuintConsumer; -import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.metrics.AccordCacheMetrics; class AccordExecutorSyncSubmit extends AccordExecutorAbstractLockLoop @@ -35,24 +34,19 @@ class AccordExecutorSyncSubmit extends AccordExecutorAbstractLockLoop private final ReentrantLock lock; private final Condition hasWork; - public AccordExecutorSyncSubmit(int executorId, Mode mode, String name, AccordCacheMetrics metrics, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor, ExecutorPlus rangeLoadExecutor, Agent agent) + public AccordExecutorSyncSubmit(int executorId, Mode mode, String name, AccordCacheMetrics metrics, Agent agent) { - this(executorId, mode, 1, constant(name), metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + this(executorId, mode, 1, constant(name), metrics, agent); } - public AccordExecutorSyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor, ExecutorPlus rangeLoadExecutor, Agent agent) + public AccordExecutorSyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(executorId, mode, threads, name, metrics, constantFactory(loadExecutor), constantFactory(saveExecutor), constantFactory(rangeLoadExecutor), agent); + this(new ReentrantLock(), executorId, mode, threads, name, metrics, agent); } - public AccordExecutorSyncSubmit(int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) + private AccordExecutorSyncSubmit(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, Agent agent) { - this(new ReentrantLock(), executorId, mode, threads, name, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); - } - - private AccordExecutorSyncSubmit(ReentrantLock lock, int executorId, Mode mode, int threads, IntFunction name, AccordCacheMetrics metrics, ExecutorFunctionFactory loadExecutor, ExecutorFunctionFactory saveExecutor, ExecutorFunctionFactory rangeLoadExecutor, Agent agent) - { - super(lock, executorId, metrics, loadExecutor, saveExecutor, rangeLoadExecutor, agent); + super(lock, executorId, metrics, agent); this.lock = lock; this.hasWork = lock.newCondition(); this.loops = new AccordExecutorLoops(mode, threads, name, this::task); diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournal.java b/src/java/org/apache/cassandra/service/accord/AccordJournal.java index cd64259cd1..3e4ec84e07 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournal.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournal.java @@ -529,14 +529,14 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier class ReplayStream implements Closeable { final CommandStore commandStore; - final Loader loader; + final Replayer replayer; final CloseableIterator> iter; JournalKey prev; public ReplayStream(CommandStore commandStore) { this.commandStore = commandStore; - this.loader = commandStore.loader(); + this.replayer = commandStore.replayer(); // Keys in the index are sorted by command store id, so index iteration will be sequential this.iter = journalTable.keyIterator(new JournalKey(TxnId.NONE, COMMAND_DIFF, commandStore.id()), new JournalKey(TxnId.MAX.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id())); } @@ -569,8 +569,8 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier "duplicate key detected %s == %s", key, prev); prev = key; commandParallelism.acquireThrowUncheckedOnInterrupt(1); - loader.load(txnId) - .map(route -> { + replayer.replay(txnId) + .map(route -> { if (segments != null && route != null) { for (long segment : segments) diff --git a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java index a6e76b5dfa..9646d63c6c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java +++ b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java @@ -332,7 +332,7 @@ public class AccordKeyspace { ByteBuffer bytes; if (serialized instanceof ByteBuffer) bytes = (ByteBuffer) serialized; - else bytes = Serialize.toBytesWithoutKey(commandsForKey); + else bytes = Serialize.toBytesWithoutKey(commandsForKey.maximalPrune()); // TODO (expected): we only need to strip pruned, not prune additional txns return makeUpdate(storeId, key, timestampMicros, bytes); } @@ -486,8 +486,6 @@ public class AccordKeyspace } } - public static final CommandsForKeyAccessor CommandsForKeysAccessor = new CommandsForKeyAccessor(CommandsForKeys); - private static TableMetadata.Builder parse(String name, String description, String cql) { return CreateTableStatement.parse(format(cql, name), ACCORD_KEYSPACE_NAME) @@ -512,16 +510,6 @@ public class AccordKeyspace return TABLES; } - public static void truncateAllCaches() - { - Keyspace ks = Keyspace.open(ACCORD_KEYSPACE_NAME); - for (String table : new String[]{ CommandsForKeys.name }) - { - if (!ks.getColumnFamilyStore(table).isEmpty()) - ks.getColumnFamilyStore(table).truncateBlocking(); - } - } - private static ByteBuffer cellValue(Cell cell) { return cell.accessor().toBuffer(cell.value()); diff --git a/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java index 2639e184d3..5cdff4f11f 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSafeCommandStore.java @@ -134,7 +134,7 @@ public class AccordSafeCommandStore extends AbstractSafeCommandStore { - AccordCommandStore.safeRedundantBeforeUpdater.accumulateAndGet((AccordCommandStore)commandStore, update, SafeRedundantBefore::max); + AccordCommandStore.safeRedundantBeforeUpdater.accumulateAndGet(commandStore, update, SafeRedundantBefore::max); }; Runnable prevOnDone = onDone; onDone = prevOnDone == null ? reportRedundantBefore : () -> { diff --git a/src/java/org/apache/cassandra/service/accord/AccordSafeState.java b/src/java/org/apache/cassandra/service/accord/AccordSafeState.java index 66b3eb1d1d..50ce897c03 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSafeState.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSafeState.java @@ -18,7 +18,6 @@ package org.apache.cassandra.service.accord; import accord.impl.SafeState; -import accord.utils.async.Cancellable; public interface AccordSafeState extends SafeState { @@ -45,11 +44,6 @@ public interface AccordSafeState extends SafeState return global().key(); } - default Cancellable saving() - { - return global().saving(); - } - default Throwable failure() { return global().failure(); diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 1d1928e4a6..3d20f18677 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -146,6 +146,7 @@ import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordCommandStoreShardCount; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordGlobalDurabilityCycle; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordShardDurabilityCycle; +import static org.apache.cassandra.config.DatabaseDescriptor.getAccordShardDurabilityMaxSplits; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordShardDurabilityTargetSplits; import static org.apache.cassandra.config.DatabaseDescriptor.getPartitioner; import static org.apache.cassandra.metrics.ClientRequestsMetricsHolder.accordReadBookkeeping; @@ -257,12 +258,12 @@ public class AccordService implements IAccordService, Shutdownable CommandsForKey.disableLinearizabilityViolationsReporting(); try { - AccordKeyspace.truncateAllCaches(); - + as.node.commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().stop()); as.journal().replay(as.node().commandStores()); logger.info("Waiting for command stores to quiesce."); ((AccordCommandStores)as.node.commandStores()).waitForQuiescense(); as.journal.unsafeSetStarted(); + as.node.commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().start()); } finally { @@ -451,8 +452,9 @@ public class AccordService implements IAccordService, Shutdownable configService.start(); fastPathCoordinator.start(); ClusterMetadataService.instance().log().addListener(fastPathCoordinator); - node.durability().shards().setTargetShardSplits(Ints.checkedCast(getAccordShardDurabilityTargetSplits())); - node.durability().shards().setShardCycleTime(Ints.checkedCast(getAccordShardDurabilityCycle(SECONDS)), SECONDS); + node.durability().shards().reconfigure(Ints.checkedCast(getAccordShardDurabilityTargetSplits()), + Ints.checkedCast(getAccordShardDurabilityMaxSplits()), + Ints.checkedCast(getAccordShardDurabilityCycle(SECONDS)), SECONDS); node.durability().global().setGlobalCycleTime(Ints.checkedCast(getAccordGlobalDurabilityCycle(SECONDS)), SECONDS); node.durability().start(); state = State.STARTED; diff --git a/src/java/org/apache/cassandra/service/accord/AccordTask.java b/src/java/org/apache/cassandra/service/accord/AccordTask.java index 9b1bc3f7fc..1414c19472 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTask.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTask.java @@ -29,7 +29,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; import java.util.function.BiConsumer; -import java.util.function.BiFunction; import java.util.function.Consumer; import java.util.function.Function; import javax.annotation.Nonnull; @@ -63,7 +62,6 @@ import org.apache.cassandra.concurrent.DebuggableTask; 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; @@ -92,7 +90,7 @@ 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.utils.Clock.Global.nanoTime; -public abstract class AccordTask extends SubmittableTask implements Runnable, Function, Cancellable, DebuggableTask +public abstract class AccordTask extends SubmittableTask implements Function, Cancellable, DebuggableTask { private static final Logger logger = LoggerFactory.getLogger(AccordTask.class); private static final NoSpamLogger noSpamLogger = NoSpamLogger.getLogger(logger, 1, TimeUnit.MINUTES); @@ -336,6 +334,7 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, // to be invoked only by the CommandStore owning thread, to take references to objects already in use by the current execution public void presetup(AccordTask parent) { + this.queuePosition = parent.queuePosition; // note we use the caches "unsafely" here deliberately, as we only reference commands we already have references to // so we do not mutate anything, except the atomic counter of references if (parent.commands != null) @@ -470,11 +469,12 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, Map map; switch (entryStatus) { - default: throw new IllegalStateException("Unhandled global state: " + entryStatus); + default: throw new UnhandledEnum(entryStatus); case WAITING_TO_LOAD: case LOADING: map = ensureLoading(); break; + case WAITING_TO_SAVE: case SAVING: case LOADED: case MODIFIED: @@ -632,12 +632,10 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, } @Override - protected void preRunExclusive(@Nullable AccordExecutor.TaskRunner runner) + protected void preRunExclusive() { state(RUNNING); queued = null; - if (runner != null) - runner.setDebuggable(this); if (rangeScanner != null) { commandsForRanges = rangeScanner.finish(commandStore.cachesExclusive()); @@ -650,7 +648,7 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, } @Override - public void run() + public void runInternal() { logger.trace("Running {} with state {}", this, state); AccordSafeCommandStore safeStore = null; @@ -749,10 +747,8 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, } } - protected void cleanupExclusive(@Nullable AccordExecutor.TaskRunner runner) + protected void cleanupExclusive() { - if (runner != null) - runner.clearDebuggable(); releaseResources(commandStore.cachesExclusive()); if (state == FAILING) state(FAILED); @@ -952,7 +948,7 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, boolean scanned; - void runInternal() + protected void runInternal() { for (Range range : ranges) { @@ -984,6 +980,7 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, return; case MODIFIED: + case WAITING_TO_SAVE: case SAVING: case LOADED: case FAILED_TO_SAVE: @@ -1027,7 +1024,7 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, } } - public class RangeTxnScanner implements Runnable + public class RangeTxnScanner extends AccordExecutor.AbstractIOTask { class CommandWatcher implements AccordCache.Listener { @@ -1047,23 +1044,9 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, CommandsForRanges.Loader summaryLoader; boolean scanned; + Throwable failure; - @Override - public void run() - { - try - { - runInternal(); - } - catch (Throwable t) - { - commandStore.executor().onScannedRanges(AccordTask.this, t); - throw t; - } - commandStore.executor().onScannedRanges(AccordTask.this, null); - } - - void runInternal() + protected void runInternal() { summaryLoader.intersects(txnId -> { if (summaries.containsKey(txnId)) @@ -1075,12 +1058,24 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, }); } - public void start(BiFunction executor) + @Override + protected void postRunExclusive() + { + commandStore.executor().onScannedRangesExclusive(AccordTask.this, failure); + } + + @Override + protected void fail(Throwable t) + { + this.failure = t; + } + + public void start(AccordExecutor executor) { Caches caches = commandStore.cachesExclusive(); state(SCANNING_RANGES); startInternal(caches); - executor.apply(AccordTask.this, this); + executor.submitPlainExclusive(AccordTask.this, this); } void startInternal(Caches caches) @@ -1116,10 +1111,16 @@ public abstract class AccordTask extends SubmittableTask implements Runnable, } @Override - public String toString() + public String description() { return "Scanning range intersections for " + AccordTask.this; } + + @Override + public String toString() + { + return description(); + } } @Override diff --git a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java index 9c4b241755..eddafc8398 100644 --- a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java +++ b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java @@ -60,6 +60,7 @@ import accord.utils.SymmetricComparator; import accord.utils.UnhandledEnum; import org.agrona.collections.Object2ObjectHashMap; import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.utils.btree.BTree; import org.apache.cassandra.utils.btree.BTreeSet; import org.apache.cassandra.utils.btree.IntervalBTree; import org.apache.cassandra.utils.concurrent.IntrusiveStack; @@ -91,6 +92,12 @@ public class CommandsForRanges extends TreeMap implements Co { this(range.start(), range.end(), txnId); } + + @Override + public String toString() + { + return super.toString() + ':' + txnId; + } } static class IntervalComparators implements IntervalBTree.IntervalComparators @@ -204,6 +211,12 @@ public class CommandsForRanges extends TreeMap implements Co if (!upd.equals(cur)) pushEdit(new IntervalTreeEdit(txnId, toMap(txnId, upd), cur == null ? null : toMap(txnId, cur))); } + else + { + RangeRoute cur = cachedRangeTxnsById.remove(cmd.txnId()); + if (cur != null) + pushEdit(new IntervalTreeEdit(txnId, null, toMap(txnId, cur))); + } } } } @@ -267,8 +280,8 @@ public class CommandsForRanges extends TreeMap implements Co List update = new ArrayList<>(), remove = new ArrayList<>(); for (IntervalTreeEdit edit : editMap.values()) { - if (edit.update != null) update.addAll(BTreeSet.wrap(edit.update, COMPARATORS.totalOrder())); if (edit.remove != null) remove.addAll(BTreeSet.wrap(edit.remove, COMPARATORS.totalOrder())); + if (edit.update != null) update.addAll(BTreeSet.wrap(edit.update, COMPARATORS.totalOrder())); } if (!remove.isEmpty()) @@ -281,6 +294,26 @@ public class CommandsForRanges extends TreeMap implements Co update.sort(COMPARATORS.totalOrder()); cachedRangeTxnsByRange = IntervalBTree.update(cachedRangeTxnsByRange, IntervalBTree.build(update, COMPARATORS), COMPARATORS); } + + if (Invariants.isParanoid()) + { + try (AccordCommandStore.ExclusiveCaches caches = commandStore.tryLockCaches()) + { + if (caches != null) + { + for (TxnIdInterval i : BTree.iterable(cachedRangeTxnsByRange)) + { + if (caches.commands().getUnsafe(i.txnId) == null) + { + boolean removed = pendingEdits != null && pendingEdits.foldl((edit, interval, r) -> { + return r || (edit.txnId.equals(i.txnId) && BTree.find(edit.remove, COMPARATORS.totalOrder(), i) != null); + }, i, false); + Invariants.require(removed); + } + } + } + } + } } private void postUnlock() @@ -472,9 +505,14 @@ public class CommandsForRanges extends TreeMap implements Co if (isRelevant(i)) { TxnId txnId = i.txnId; - Summary summary = ifRelevant(c.getUnsafe(txnId)); - if (summary != null) - f.accept(summary); + AccordCacheEntry entry = c.getUnsafe(txnId); + Invariants.expect(entry != null, "%s found interval %s but no matching transaction in cache", manager.commandStore, i); + if (entry != null) + { + Summary summary = ifRelevant(entry); + if (summary != null) + f.accept(summary); + } } return c; }, forEach, this, caches.commands()); @@ -529,6 +567,7 @@ public class CommandsForRanges extends TreeMap implements Co case LOADED: case MODIFIED: case SAVING: + case WAITING_TO_SAVE: case FAILED_TO_SAVE: } diff --git a/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java b/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java index ed39bbad8d..0ddc6683a5 100644 --- a/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java +++ b/src/java/org/apache/cassandra/service/accord/api/AccordAgent.java @@ -209,10 +209,11 @@ public class AccordAgent implements Agent return SECONDS.toMicros(1); } + // TODO (expected): I don't think we even need this - just prune each time we have doubled in size @Override public long maxConflictsPruneInterval() { - return 100; + return 1024; } /** 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 717f9232d6..69427e3ef6 100644 --- a/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java +++ b/src/java/org/apache/cassandra/service/accord/interop/AccordInteropApply.java @@ -225,10 +225,10 @@ public class AccordInteropApply extends Apply implements LocalListeners.ComplexL { return "AccordInteropApply{" + "txnId:" + txnId + - ", deps:" + deps + + ", deps:" + deps() + ", executeAt:" + executeAt + - ", writes:" + writes + - ", result:" + result + + ", writes:" + writes() + + ", result:" + result() + '}'; } diff --git a/src/java/org/apache/cassandra/service/accord/serializers/AcceptSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/AcceptSerializers.java index d7603b6981..cd96796142 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/AcceptSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/AcceptSerializers.java @@ -54,7 +54,7 @@ public class AcceptSerializers out.writeByte((accept.kind.ordinal() << 1) | (accept.isPartialAccept ? IS_PARTIAL : 0)); CommandSerializers.ballot.serialize(accept.ballot, out); ExecuteAtSerializer.serialize(accept.txnId, accept.executeAt, out); - DepsSerializers.partialDeps.serialize(accept.partialDeps, out); + DepsSerializers.partialDeps.serialize(accept.partialDeps(), out); } @Override @@ -76,7 +76,7 @@ public class AcceptSerializers return 1 + CommandSerializers.ballot.serializedSize(accept.ballot) + ExecuteAtSerializer.serializedSize(accept.txnId, accept.executeAt) - + DepsSerializers.partialDeps.serializedSize(accept.partialDeps); + + DepsSerializers.partialDeps.serializedSize(accept.partialDeps()); } } 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 8716d87709..e3c009d17e 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/ApplySerializers.java @@ -73,11 +73,11 @@ public class ApplySerializers out.writeUnsignedVInt(apply.maxEpoch - apply.minEpoch); kind.serialize(apply.kind, out); ExecuteAtSerializer.serialize(apply.txnId, apply.executeAt, out); - DepsSerializers.partialDeps.serialize(apply.deps, out); - CommandSerializers.nullablePartialTxn.serialize(apply.txn, out, version); + DepsSerializers.partialDeps.serialize(apply.deps(), out); + CommandSerializers.nullablePartialTxn.serialize(apply.txn(), out, version); KeySerializers.nullableFullRoute.serialize(apply.fullRoute, out); if (apply.txnId.is(Write)) - CommandSerializers.writes.serialize(apply.writes, out, version); + CommandSerializers.writes.serialize(apply.writes(), out, version); out.writeUnsignedVInt32(apply.flags.bits()); } @@ -109,10 +109,10 @@ public class ApplySerializers + TypeSizes.sizeofUnsignedVInt(apply.maxEpoch - apply.minEpoch) + kind.serializedSize(apply.kind) + ExecuteAtSerializer.serializedSize(apply.txnId, apply.executeAt) - + DepsSerializers.partialDeps.serializedSize(apply.deps) - + CommandSerializers.nullablePartialTxn.serializedSize(apply.txn, version) + + DepsSerializers.partialDeps.serializedSize(apply.deps()) + + CommandSerializers.nullablePartialTxn.serializedSize(apply.txn(), version) + KeySerializers.nullableFullRoute.serializedSize(apply.fullRoute) - + (apply.txnId.is(Write) ? CommandSerializers.writes.serializedSize(apply.writes, version) : 0) + + (apply.txnId.is(Write) ? CommandSerializers.writes.serializedSize(apply.writes(), version) : 0) + VIntCoding.sizeOfUnsignedVInt(apply.flags.bits()) ; } diff --git a/src/java/org/apache/cassandra/service/accord/serializers/ReadDataSerializer.java b/src/java/org/apache/cassandra/service/accord/serializers/ReadDataSerializer.java index aba9a3788b..4311fb0efa 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/ReadDataSerializer.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/ReadDataSerializer.java @@ -131,15 +131,15 @@ public class ReadDataSerializer implements IVersionedSerializer { ApplyThenWaitUntilApplied msg = (ApplyThenWaitUntilApplied) read; boolean hasMinEpoch = msg.minEpoch() != read.txnId.epoch(); - boolean hasWrites = msg.writes != null; + boolean hasWrites = msg.writes() != null; out.writeUnsignedVInt32((hasMinEpoch ? ATWUA_HAS_MIN_EPOCH : 0) | (hasWrites ? ATWUA_HAS_WRITES : 0)); if (hasMinEpoch) out.writeVInt(read.txnId.epoch() - msg.minEpoch()); - DepsSerializers.partialDeps.serialize(msg.deps, out); + DepsSerializers.partialDeps.serialize(msg.deps(), out); KeySerializers.fullRoute.serialize(msg.route, out); if (hasWrites) - CommandSerializers.writes.serialize(msg.writes, out, version); + CommandSerializers.writes.serialize(msg.writes(), out, version); break; } } @@ -276,15 +276,15 @@ public class ReadDataSerializer implements IVersionedSerializer { ApplyThenWaitUntilApplied msg = (ApplyThenWaitUntilApplied) read; boolean hasMinEpoch = msg.minEpoch() != read.txnId.epoch(); - boolean hasWrites = msg.writes != null; + boolean hasWrites = msg.writes() != null; size += VIntCoding.computeUnsignedVIntSize((hasMinEpoch ? ATWUA_HAS_MIN_EPOCH : 0) | (hasWrites ? ATWUA_HAS_WRITES : 0)); if (hasMinEpoch) size += VIntCoding.computeVIntSize(read.txnId.epoch() - msg.minEpoch()); - size += DepsSerializers.partialDeps.serializedSize(msg.deps); + size += DepsSerializers.partialDeps.serializedSize(msg.deps()); size += KeySerializers.fullRoute.serializedSize(msg.route); if (hasWrites) - size += CommandSerializers.writes.serializedSize(msg.writes, version); + size += CommandSerializers.writes.serializedSize(msg.writes(), version); break; } } diff --git a/src/java/org/apache/cassandra/utils/JVMStabilityInspector.java b/src/java/org/apache/cassandra/utils/JVMStabilityInspector.java index 2c840d4a7d..e77a231da2 100644 --- a/src/java/org/apache/cassandra/utils/JVMStabilityInspector.java +++ b/src/java/org/apache/cassandra/utils/JVMStabilityInspector.java @@ -33,6 +33,7 @@ import com.google.common.collect.ImmutableSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import ch.qos.logback.classic.LoggerContext; import net.nicoulaj.compilecommand.annotations.Exclude; import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.config.Config; @@ -273,6 +274,8 @@ public final class JVMStabilityInspector if (doExit && killing.compareAndSet(false, true)) { + if (LoggerFactory.getILoggerFactory() instanceof LoggerContext) + ((LoggerContext) LoggerFactory.getILoggerFactory()).stop(); StorageService.instance.removeShutdownHook(); System.exit(100); } diff --git a/src/java/org/apache/cassandra/utils/btree/BTree.java b/src/java/org/apache/cassandra/utils/btree/BTree.java index fac23d47f7..0ef17c1696 100644 --- a/src/java/org/apache/cassandra/utils/btree/BTree.java +++ b/src/java/org/apache/cassandra/utils/btree/BTree.java @@ -2494,7 +2494,7 @@ public class BTree final BranchBuilder ensureParent() { if (parent == null) parent = allocateParent(); - else if (!parent.inUse) initParent(); + if (!parent.inUse) initParent(); return parent; } @@ -3500,12 +3500,14 @@ public class BTree while (branch != null && branch.inUse) { branch.count = 0; + branch.hasRightChild = false; clearBranchBuffer(branch.buffer); if (branch.savedBuffer != null && branch.savedBuffer[0] != null) Arrays.fill(branch.savedBuffer, null); // by definition full, if non-empty branch.inUse = false; branch = branch.parent; } + Invariants.require(branch == null || (branch.count == 0 && !branch.hasRightChild)); } /** @@ -3589,8 +3591,13 @@ public class BTree this.allocated = isSimple(updateF) ? -1 : 0; int leafDepth = BTree.depth(update) - 1; LeafOrBranchBuilder builder = leaf(); + Invariants.require(builder.isEmpty()); for (int i = 0; i < leafDepth; ++i) - builder = builder.ensureParent(); + { + BranchBuilder branch = builder.ensureParent(); + Invariants.require(builder.isEmpty() && !branch.hasRightChild); + builder = branch; + } Insert ik = this.insert.next(); ik = updateRecursive(ik, update, null, builder); @@ -3738,13 +3745,14 @@ public class BTree reset(); pool.offer(this); pool = null; - comparator = null; } void reset() { super.reset(); insert.reset(); + comparator = null; + updateF = null; } } @@ -3788,8 +3796,13 @@ public class BTree { int height = this.update.initToRoot(root); LeafOrBranchBuilder level = leaf(); + Invariants.require(level.isEmpty()); for (int d = 1; d < height; ++d) - level = level.ensureParent(); + { + BranchBuilder branch = level.ensureParent(); + Invariants.require(branch.isEmpty() && !branch.hasRightChild); + level = branch; + } level.setSourceNode(root); return level; } @@ -4111,8 +4124,8 @@ public class BTree void reset() { - update.reset(); super.reset(); + update.reset(); } } @@ -4211,6 +4224,20 @@ public class BTree leaf().copy(unode, prev, usz - prev); return true; } + + @Override + public void close() + { + reset(); + } + + @Override + void reset() + { + super.reset(); + remove = null; + comparator = null; + } } private static abstract class AbstractTransformer extends AbstractSeekingTransformer diff --git a/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java b/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java index b027295562..a8e5e3dcd6 100644 --- a/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java +++ b/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java @@ -378,7 +378,7 @@ public class IntervalBTree @Override BTree.BranchBuilder allocateParent() { - return new IntervalBranchBuilder(this).init(adapter); + return new IntervalBranchBuilder(this); } @Override @@ -442,7 +442,7 @@ public class IntervalBTree @Override BTree.BranchBuilder allocateParent() { - return new IntervalBranchBuilder(this).init(adapter); + return new IntervalBranchBuilder(this); } @Override @@ -479,7 +479,7 @@ public class IntervalBTree @Override BTree.BranchBuilder allocateParent() { - return new IntervalBranchBuilder(this).init(adapter); + return new IntervalBranchBuilder(this); } @Override @@ -492,6 +492,8 @@ public class IntervalBTree @Override public void close() { + // TODO (required): validate this in IntervalBTreeTest + super.close(); adapter.reset(); } } @@ -524,7 +526,7 @@ public class IntervalBTree @Override BTree.BranchBuilder allocateParent() { - return new IntervalBranchBuilder(this).init(adapter); + return new IntervalBranchBuilder(this); } @Override @@ -535,10 +537,9 @@ public class IntervalBTree } @Override - public void close() + void reset() { - reset(); - comparator = null; + super.reset(); adapter.reset(); } } diff --git a/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java b/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java index 07229249c5..f50bb625c5 100644 --- a/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java +++ b/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java @@ -25,6 +25,7 @@ import java.util.function.BiFunction; import java.util.function.Consumer; import java.util.function.Function; +import accord.utils.TriFunction; import net.nicoulaj.compilecommand.annotations.Inline; import org.apache.cassandra.utils.LongAccumulator; @@ -208,6 +209,17 @@ public class IntrusiveStack> implements Iterable forEach((T)this, forEach); } + public V foldl(TriFunction foldl, P param, V accumulator) + { + T list = (T) this; + while (list != null) + { + accumulator = foldl.apply(list, param, accumulator); + list = list.next; + } + return accumulator; + } + protected static > void forEach(T list, Consumer forEach) { forEach(list, Function.identity(), forEach); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java index 04a6dca76d..a5aae94516 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java @@ -35,7 +35,6 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import accord.impl.progresslog.DefaultProgressLogs; import accord.local.Node; import accord.local.PreLoadContext; import accord.local.SafeCommand; @@ -147,7 +146,7 @@ public class AccordIncrementalRepairTest extends AccordTestBase public void tearDown() { for (IInvokableInstance instance : SHARED_CLUSTER) - instance.runOnInstance(() -> DefaultProgressLogs.unsafePauseForTesting(false)); + instance.runOnInstance(() -> AccordService.instance().node().commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().start())); SHARED_CLUSTER.filters().reset(); } @@ -293,7 +292,7 @@ public class AccordIncrementalRepairTest extends AccordTestBase // heal partition and wait for node 1 to see node 3 again for (IInvokableInstance instance : SHARED_CLUSTER) instance.runOnInstance(() -> { - DefaultProgressLogs.unsafePauseForTesting(true); + AccordService.instance().node().commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().stop()); Assert.assertFalse(barrierRecordingService().executedBarriers); }); SHARED_CLUSTER.filters().reset(); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIntegrationTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIntegrationTest.java index 144fcdf311..c2b8f1e486 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIntegrationTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIntegrationTest.java @@ -27,10 +27,10 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import accord.impl.progresslog.DefaultProgressLogs; import org.apache.cassandra.distributed.api.IInvokableInstance; import org.apache.cassandra.distributed.api.IMessageFilters; import org.apache.cassandra.net.Verb; +import org.apache.cassandra.service.accord.AccordService; public class AccordIntegrationTest extends AccordTestBase { @@ -125,6 +125,6 @@ public class AccordIntegrationTest extends AccordTestBase private void pauseSimpleProgressLog() { for (IInvokableInstance instance : SHARED_CLUSTER) - instance.runOnInstance(() -> DefaultProgressLogs.unsafePauseForTesting(true)); + instance.runOnInstance(() -> AccordService.instance().node().commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().stop())); } } 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 66e188df5c..4dadcc6904 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordLoadTest.java @@ -51,6 +51,8 @@ import org.apache.cassandra.distributed.api.IMessage; import org.apache.cassandra.distributed.api.IMessageFilters; import org.apache.cassandra.distributed.shared.DistributedTestBase; import org.apache.cassandra.net.Verb; +import org.apache.cassandra.schema.Schema; +import org.apache.cassandra.service.accord.AccordKeyspace; import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.utils.EstimatedHistogram; @@ -58,6 +60,7 @@ 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; public class AccordLoadTest extends AccordTestBase { @@ -70,10 +73,11 @@ public class AccordLoadTest extends AccordTestBase // AccordTestBase.setupCluster(builder -> builder, 3); AccordTestBase.setupCluster(builder -> builder.withConfig(config -> config .with(Feature.NETWORK, Feature.GOSSIP) - .set("accord.shard_durability_target_splits", "32") + .set("accord.shard_durability_target_splits", "8") + .set("accord.shard_durability_max_splits", "16") .set("accord.shard_durability_cycle", "1m") // .set("accord.ephemeral_read_enabled", "true") - .set("accord.gc_delay", "30s")), 3); + ), 3); } @Ignore @@ -101,7 +105,9 @@ public class AccordLoadTest extends AccordTestBase final int repairInterval = Integer.MAX_VALUE; final int compactionInterval = 20_000; // final int flushInterval = 50_000; - final int flushInterval = 500; + final int journalFlushInterval = 2_000; + final int cfkFlushInterval = 10_000; + final int dataFlushInterval = 10_000; final int compactionPeriodSeconds = 0; int restartInterval = 30_000; final int restartDecay = 2; @@ -114,7 +120,9 @@ public class AccordLoadTest extends AccordTestBase final float readChance = 0.33f; long nextRepairAt = repairInterval; long nextCompactionAt = compactionInterval; - long nextFlushAt = flushInterval; + long nextJournalFlushAt = journalFlushInterval; + long nextDataFlushAt = dataFlushInterval; + long nextCfkFlushAt = cfkFlushInterval; long nextRestartAt = restartInterval; final ExecutorService restartExecutor = Executors.newSingleThreadExecutor(); final BitSet initialised = new BitSet(); @@ -205,9 +213,9 @@ public class AccordLoadTest extends AccordTestBase }); } - if ((nextFlushAt -= batchSize) <= 0) + if ((nextJournalFlushAt -= batchSize) <= 0) { - nextFlushAt += flushInterval; + nextJournalFlushAt += journalFlushInterval; System.out.println("flushing journal..."); cluster.forEach(i -> { try @@ -221,7 +229,43 @@ public class AccordLoadTest extends AccordTestBase { logger.error("", t); } - }); + }); + } + + if ((nextDataFlushAt -= batchSize) <= 0) + { + nextDataFlushAt += dataFlushInterval; + System.out.println("flushing data..."); + cluster.forEach(i -> { + try + { + i.acceptOnInstance(name -> { + Schema.instance.getColumnFamilyStoreInstance(Schema.instance.getTableMetadata(KEYSPACE, name).id).forceFlush(UNIT_TESTS); + }, accordTableName); + } + catch (Throwable t) + { + logger.error("", t); + } + }); + } + + if ((nextCfkFlushAt -= batchSize) <= 0) + { + nextCfkFlushAt += cfkFlushInterval; + System.out.println("flushing data..."); + cluster.forEach(i -> { + try + { + i.acceptOnInstance(name -> { + AccordKeyspace.AccordColumnFamilyStores.commandsForKey.forceFlush(UNIT_TESTS); + }, accordTableName); + } + catch (Throwable t) + { + logger.error("", t); + } + }); } if ((nextRestartAt -= batchSize) <= 0) diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java index 49fb52f100..2236f27cd1 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java @@ -46,7 +46,6 @@ import org.slf4j.LoggerFactory; import accord.api.RoutingKey; import accord.coordinate.Invalidated; -import accord.impl.progresslog.DefaultProgressLogs; import accord.messages.PreAccept; import accord.primitives.KeyRoute; import accord.primitives.Routable.Domain; @@ -165,7 +164,7 @@ public abstract class AccordTestBase extends TestBaseImpl public void tearDown() throws Exception { for (IInvokableInstance instance : SHARED_CLUSTER) - instance.runOnInstance(() -> DefaultProgressLogs.unsafePauseForTesting(false)); + instance.runOnInstance(() -> AccordService.instance().node().commandStores().forEachCommandStore(cs -> cs.unsafeProgressLog().start())); truncateSystemTables(); diff --git a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java index c134eaed43..67bb3c6929 100644 --- a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java +++ b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java @@ -32,6 +32,8 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; +import java.util.function.Supplier; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Iterators; @@ -932,6 +934,17 @@ public class ColumnFamilyStoreTest return false; } + @Override + public > T ensureFlushListener(Object key, Supplier onDurablyFlushed) + { + return null; + } + + @Override + public void notifyFlushed() + { + } + @Override public boolean shouldSwitch(FlushReason reason, TableMetadata latest) { diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java index 3b5292ef55..bc83665a21 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionAccordIteratorsTest.java @@ -90,7 +90,7 @@ import org.apache.cassandra.utils.FBUtilities; import static accord.local.LoadKeys.SYNC; import static accord.local.LoadKeysFor.READ_WRITE; import static accord.local.PreLoadContext.contextFor; -import static accord.local.RedundantStatus.SomeStatus.GC_BEFORE_AND_LOCALLY_APPLIED; +import static accord.local.RedundantStatus.SomeStatus.GC_BEFORE_AND_LOCALLY_DURABLE; import static accord.primitives.Routable.Domain.Range; import static accord.primitives.Timestamp.Flag.HLC_BOUND; import static accord.primitives.Timestamp.Flag.SHARD_BOUND; @@ -230,7 +230,7 @@ public class CompactionAccordIteratorsTest { Ranges ranges = AccordTestUtils.fullRange(AccordTestUtils.keys(table, 42)); txnId = txnId.as(Kind.ExclusiveSyncPoint, Range).addFlag(SHARD_BOUND); - return RedundantBefore.create(ranges, Long.MIN_VALUE, Long.MAX_VALUE, txnId, GC_BEFORE_AND_LOCALLY_APPLIED, LT_TXN_ID.as(Range)); + return RedundantBefore.create(ranges, Long.MIN_VALUE, Long.MAX_VALUE, txnId, GC_BEFORE_AND_LOCALLY_DURABLE, LT_TXN_ID.as(Range)); } enum DurableBeforeType diff --git a/test/unit/org/apache/cassandra/service/accord/AccordCacheTest.java b/test/unit/org/apache/cassandra/service/accord/AccordCacheTest.java index d7a89017d6..c8805ccc25 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordCacheTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordCacheTest.java @@ -24,13 +24,12 @@ import org.junit.Assert; import org.junit.Before; import org.junit.Test; -import accord.utils.async.Cancellable; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.ManualExecutor; import org.apache.cassandra.metrics.AccordCacheMetrics; import org.apache.cassandra.metrics.CacheAccessMetrics; -import org.apache.cassandra.service.accord.AccordCacheEntry.OnSaved; +import org.apache.cassandra.service.accord.AccordCacheEntry.SaveExecutor; import org.apache.cassandra.service.accord.AccordCacheEntry.Status; import static org.apache.cassandra.service.accord.AccordTestUtils.testLoad; @@ -90,12 +89,6 @@ public class AccordCacheTest original = global.getExclusive(); } - @Override - public Cancellable saving() - { - return global.saving(); - } - @Override public Throwable failure() { @@ -172,7 +165,7 @@ public class AccordCacheTest public void testAcquisitionAndRelease() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), 500, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), 500, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); @@ -180,7 +173,7 @@ public class AccordCacheTest SafeString safeString1 = instance.acquire("1"); assertCacheState(cache, 1, 1, emptyNodeSize()); - testLoad(executor, instance, safeString1, "1"); + testLoad(executor, safeString1, "1"); Assert.assertTrue(!cache.evictionQueue().iterator().hasNext()); instance.release(safeString1, null); @@ -190,7 +183,7 @@ public class AccordCacheTest SafeString safeString2 = instance.acquire("2"); assertCacheState(cache, 1, 2, DEFAULT_NODE_SIZE + nodeSize(1)); - testLoad(executor, instance, safeString2, "2"); + testLoad(executor, safeString2, "2"); instance.release(safeString2, null); assertCacheState(cache, 0, 2, nodeSize(1) + nodeSize(1)); @@ -205,7 +198,7 @@ public class AccordCacheTest public void testCachingMetricsWithTwoInstances() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), 500, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), 500, cacheMetrics); AccordCache.Type stringType = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance stringInstance = stringType.newInstance(null); @@ -215,20 +208,20 @@ public class AccordCacheTest AccordCache.Type.Instance intInstance = intType.newInstance(null); SafeString safeString1 = stringInstance.acquire("1"); - testLoad(executor, stringInstance, safeString1, "1"); + testLoad(executor, safeString1, "1"); stringInstance.release(safeString1, null); SafeString safeString2 = stringInstance.acquire("2"); - testLoad(executor, stringInstance, safeString2, "2"); + testLoad(executor, safeString2, "2"); stringInstance.release(safeString2, null); SafeInt safeInt1 = intInstance.acquire(3); - testLoad(executor, intInstance, safeInt1, 3); + testLoad(executor, safeInt1, 3); intInstance.release(safeInt1, null); SafeInt safeInt2 = intInstance.acquire(4); - testLoad(executor, intInstance, safeInt2, 4); + testLoad(executor, safeInt2, 4); intInstance.release(safeInt2, null); SafeInt safeInt3 = intInstance.acquire(5); - testLoad(executor, intInstance, safeInt3, 5); + testLoad(executor, safeInt3, 5); intInstance.release(safeInt3, null); assertCacheState(cache, 0, 5, nodeSize(Integer.BYTES) * 3 + nodeSize(1) * 2); @@ -247,7 +240,7 @@ public class AccordCacheTest public void testRotation() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), DEFAULT_NODE_SIZE * 5, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), DEFAULT_NODE_SIZE * 5, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); assertCacheState(cache, 0, 0, 0); @@ -259,7 +252,7 @@ public class AccordCacheTest SafeString safeString = instance.acquire(Integer.toString(i)); items[i] = safeString; Assert.assertNotNull(safeString); - testLoad(executor, instance, safeString, Integer.toString(i)); + testLoad(executor, safeString, Integer.toString(i)); Assert.assertTrue(instance.isReferenced(safeString.key())); instance.release(safeString, null); } @@ -288,7 +281,7 @@ public class AccordCacheTest public void testEvictionOnAcquire() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), nodeSize(1) * 5, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), nodeSize(1) * 5, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); @@ -299,7 +292,7 @@ public class AccordCacheTest { SafeString safeString = instance.acquire(Integer.toString(i)); items[i] = safeString; - testLoad(executor, instance, safeString, Integer.toString(i)); + testLoad(executor, safeString, Integer.toString(i)); Assert.assertTrue(instance.isReferenced(safeString.key())); instance.release(safeString, null); } @@ -322,7 +315,7 @@ public class AccordCacheTest assertCacheMetrics(cache.metrics, 0, 6, 6, 5); assertCacheMetrics(type.typeMetrics, 0, 6, 6, 5); - testLoad(executor, instance, safeString, "5"); + testLoad(executor, safeString, "5"); instance.release(safeString, null); assertCacheState(cache, 0, 5, nodeSize(1) * 5); Assert.assertSame(items[1].global, cache.head()); @@ -335,7 +328,7 @@ public class AccordCacheTest public void testEvictionOnRelease() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), nodeSize(1) * 4, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), nodeSize(1) * 4, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); @@ -346,7 +339,7 @@ public class AccordCacheTest { SafeString safeString = instance.acquire(Integer.toString(i)); items[i] = safeString; - testLoad(executor, instance, safeString, Integer.toString(i)); + testLoad(executor, safeString, Integer.toString(i)); Assert.assertTrue(instance.isReferenced(safeString.key())); } @@ -375,14 +368,14 @@ public class AccordCacheTest public void testMultiAcquireRelease() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), DEFAULT_NODE_SIZE * 4, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), DEFAULT_NODE_SIZE * 4, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); assertCacheState(cache, 0, 0, 0); SafeString safeString1 = instance.acquire("0"); - testLoad(executor, instance, safeString1, "0"); + testLoad(executor, safeString1, "0"); Assert.assertEquals(Status.LOADED, safeString1.global.status()); assertCacheMetrics(cache.metrics, 0, 1, 1, 1); assertCacheMetrics(type.typeMetrics, 0, 1, 1, 1); @@ -408,14 +401,14 @@ public class AccordCacheTest public void evictionBlockedOnSaving() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), nodeSize(1) * 3 + nodeSize(3), cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), nodeSize(1) * 3 + nodeSize(3), cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); assertCacheState(cache, 0, 0, 0); SafeString item = instance.acquire(Integer.toString(0)); - testLoad(executor, instance, item, Integer.toString(0)); + testLoad(executor, item, Integer.toString(0)); item.set("0*"); Assert.assertTrue(instance.isReferenced(item.key())); instance.release(item, null); @@ -423,7 +416,7 @@ public class AccordCacheTest for (int i=1; i<4; i++) { item = instance.acquire(Integer.toString(i)); - testLoad(executor, instance, item, Integer.toString(i)); + testLoad(executor, item, Integer.toString(i)); Assert.assertTrue(instance.isReferenced(item.key())); instance.release(item, null); } @@ -449,14 +442,14 @@ public class AccordCacheTest public void testUpdates() { ManualExecutor executor = new ManualExecutor(); - AccordCache cache = new AccordCache(wrap(executor), OnSaved.immediate(), 500, cacheMetrics); + AccordCache cache = new AccordCache(saveExecutor(executor), 500, cacheMetrics); AccordCache.Type type = cache.newType(String.class, (s, k) -> k, (s, k, c, o) -> null, Function.identity(), (s, k, v) -> true, String::length, SafeString::new); AccordCache.Type.Instance instance = type.newInstance(null); assertCacheState(cache, 0, 0, 0); SafeString safeString = instance.acquire("1"); - testLoad(executor, instance, safeString, "1"); + testLoad(executor, safeString, "1"); assertCacheState(cache, 1, 1, nodeSize(1)); Assert.assertNull(cache.head()); Assert.assertNull(cache.tail()); @@ -498,8 +491,15 @@ public class AccordCacheTest assertThat(stringInstance1).isNotSameAs(integerInstance1); } - private static Function wrap(ExecutorPlus executor) + private static SaveExecutor saveExecutor(ExecutorPlus executor) { - return r -> AccordExecutor.wrap(executor.submit(r)); + return (saving, identity, save) -> { + executor.submit(() -> { + try { save.run(); } + catch (Throwable t) { saving.saved(identity, t); throw t; } + saving.saved(identity, null); + }); + return null; + }; } } diff --git a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java index 543cd0af70..0cfab8b3d6 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java @@ -78,10 +78,10 @@ import accord.topology.TopologyManager; import accord.utils.SortedArrays.SortedArrayList; import accord.utils.async.AsyncChain; import accord.utils.async.AsyncChains; +import accord.utils.async.Cancellable; import org.apache.cassandra.ServerTestUtils; import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.ManualExecutor; -import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.AccordSpec; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.DurationSpec; @@ -99,6 +99,7 @@ import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.service.accord.AccordCacheEntry.LoadExecutor; import org.apache.cassandra.service.accord.api.AccordAgent; import org.apache.cassandra.service.accord.api.PartitionKey; import org.apache.cassandra.service.accord.serializers.TableMetadatas; @@ -108,6 +109,7 @@ import org.apache.cassandra.service.accord.txn.TxnQuery; import org.apache.cassandra.service.accord.txn.TxnRead; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.concurrent.Condition; +import org.apache.cassandra.utils.concurrent.Future; import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import static accord.primitives.Routable.Domain.Key; @@ -118,7 +120,6 @@ import static accord.primitives.Txn.Kind.Write; import static accord.utils.async.AsyncChains.getUninterruptibly; import static java.lang.String.format; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_LOCK; -import static org.apache.cassandra.service.accord.AccordExecutor.wrap; public class AccordTestUtils { @@ -177,10 +178,32 @@ public class AccordTestUtils }; } - public static void testLoad(ManualExecutor executor, AccordCache.Type.Instance instance, AccordSafeState safeState, V val) + private static LoadExecutor loadExecutor(ExecutorPlus executor) + { + return new LoadExecutor<>() + { + @Override + public Cancellable load(P1 p1, P2 p2, AccordCacheEntry entry) + { + Future future = executor.submit(() -> { + V v; + try { v = entry.owner.parent().adapter().load(entry.owner.commandStore, entry.key()); } + catch (Throwable t) + { + entry.failedToLoad(); + throw t; + } + entry.loaded(v); + }); + return () -> future.cancel(true); + } + }; + } + + public static void testLoad(ManualExecutor executor, AccordSafeState safeState, V val) { Assert.assertEquals(AccordCacheEntry.Status.WAITING_TO_LOAD, safeState.global().status()); - safeState.global().load(wrap(executor), null, instance.parent().adapter(), AccordCacheEntry.OnLoaded.immediate()); + safeState.global().load(loadExecutor(executor), null, null); Assert.assertEquals(AccordCacheEntry.Status.LOADING, safeState.global().status()); executor.runOne(); Assert.assertEquals(AccordCacheEntry.Status.LOADED, safeState.global().status()); @@ -323,9 +346,9 @@ public class AccordTestUtils } public static AccordCommandStore createAccordCommandStore( - Node.Id node, LongSupplier now, Topology topology, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor) + Node.Id node, LongSupplier now, Topology topology) { - AccordExecutor executor = new AccordExecutorSyncSubmit(0, RUN_WITH_LOCK, CommandStore.class.getSimpleName() + '[' + 0 + ']', new AccordCacheMetrics("test"), loadExecutor, saveExecutor, loadExecutor, new AccordAgent()); + AccordExecutor executor = new AccordExecutorSyncSubmit(0, RUN_WITH_LOCK, CommandStore.class.getSimpleName() + '[' + 0 + ']', new AccordCacheMetrics("test"), new AccordAgent()); return createAccordCommandStore(node, now, topology, executor); } @@ -381,28 +404,18 @@ public class AccordTestUtils return result; } - public static AccordCommandStore createAccordCommandStore(Node.Id node, LongSupplier now, Topology topology) - { - return createAccordCommandStore(node, now, topology, Stage.READ.executor(), Stage.MUTATION.executor()); - } - public static AccordCommandStore createAccordCommandStore( - LongSupplier now, String keyspace, String table, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor) + LongSupplier now, String keyspace, String table) { TableMetadata metadata = Schema.instance.getTableMetadata(keyspace, table); TokenRange range = TokenRange.fullRange(metadata.id, metadata.partitioner); Node.Id node = new Id(1); Topology topology = new Topology(1, Shard.create(range, new SortedArrayList<>(new Id[] { node }), Sets.newHashSet(node), Collections.emptySet())); - AccordCommandStore store = createAccordCommandStore(node, now, topology, loadExecutor, saveExecutor); + AccordCommandStore store = createAccordCommandStore(node, now, topology); store.execute((PreLoadContext.Empty)()->"Test", safeStore -> ((AccordCommandStore)safeStore.commandStore()).executor().cacheUnsafe().setCapacity(1 << 20)); return store; } - public static AccordCommandStore createAccordCommandStore(LongSupplier now, String keyspace, String table) - { - return createAccordCommandStore(now, keyspace, table, Stage.READ.executor(), Stage.MUTATION.executor()); - } - public static void execute(AccordCommandStore commandStore, Runnable runnable) { try 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 42e62de3ee..c5abff9533 100644 --- a/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java +++ b/test/unit/org/apache/cassandra/service/accord/serializers/CommandsForKeySerializerTest.java @@ -64,6 +64,7 @@ import accord.local.ICommand; import accord.local.Node; import accord.local.NodeCommandStoreService; import accord.local.PreLoadContext; +import accord.local.RedundantBefore; import accord.local.SafeCommand; import accord.local.SafeCommandStore; import accord.local.SequentialAsyncExecutor; @@ -649,7 +650,9 @@ public class CommandsForKeySerializerTest } @Override public boolean inStore() { return true; } - @Override public Journal.Loader loader() { throw new UnsupportedOperationException(); } + @Override public Journal.Replayer replayer() { throw new UnsupportedOperationException(); } + + @Override protected void ensureDurable(Ranges ranges, RedundantBefore onSuccess) {} @Override public Agent agent() { return this; } @Override public AsyncChain build(PreLoadContext context, Consumer consumer) { return null; } @Override public AsyncChain build(PreLoadContext context, Function apply) { throw new UnsupportedOperationException(); }