diff --git a/.gitmodules b/.gitmodules index 8a45ad8f54..616dacf610 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,4 +1,4 @@ [submodule "modules/accord"] path = modules/accord url = https://github.com/apache/cassandra-accord.git - branch = trunk \ No newline at end of file + branch = trunk diff --git a/modules/accord b/modules/accord index 0d6157fc33..3aba47c29a 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 0d6157fc33dd16f1768030e66205e72fc1d7e9ac +Subproject commit 3aba47c29a1cb13dba05814fd6b007165193fd7e diff --git a/src/java/org/apache/cassandra/io/AsymmetricUnversionedSerializer.java b/src/java/org/apache/cassandra/io/AsymmetricUnversionedSerializer.java index 570741903f..34b0a185f8 100644 --- a/src/java/org/apache/cassandra/io/AsymmetricUnversionedSerializer.java +++ b/src/java/org/apache/cassandra/io/AsymmetricUnversionedSerializer.java @@ -30,6 +30,11 @@ import org.apache.cassandra.io.util.DataOutputPlus; public interface AsymmetricUnversionedSerializer { void serialize(In t, DataOutputPlus out) throws IOException; + + /** + * Note: it is not guaranteed that this output is compatible with the DataInput/OutputPlus variations, + * as the ByteBuffer has an implied length. + */ default ByteBuffer serialize(In t) throws IOException { int size = Math.toIntExact(serializedSize(t)); @@ -59,6 +64,11 @@ public interface AsymmetricUnversionedSerializer } } Out deserialize(DataInputPlus in) throws IOException; + + /** + * Note: it is not guaranteed to be safe to provide an input created by the DataOutputPlus serializer varation + * as the ByteBuffer has an implied length. + */ default Out deserialize(ByteBuffer buffer) throws IOException { try (DataInputBuffer in = new DataInputBuffer(buffer, true)) @@ -78,5 +88,6 @@ public interface AsymmetricUnversionedSerializer throw new UncheckedIOException(e); } } + long serializedSize(In t); } diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index 52986d73dc..bcd607dc25 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -1097,6 +1097,12 @@ public abstract class SSTableReader extends SSTable implements UnfilteredSource, return new SSTableSimpleScanner(this, getPositionsForBoundsIterator(boundsIterator)); } + public ISSTableScanner getScanner(AbstractBounds bounds) + { + PartitionPositionBounds positionBounds = getPositionsForBounds(bounds); + return new SSTableSimpleScanner(this, positionBounds == null ? Collections.emptyList() : Collections.singletonList(positionBounds)); + } + /** * Create a {@link FileDataInput} for the data file of the sstable represented by this reader. This method returns diff --git a/src/java/org/apache/cassandra/journal/ActiveSegment.java b/src/java/org/apache/cassandra/journal/ActiveSegment.java index 85035f4991..5c93655552 100644 --- a/src/java/org/apache/cassandra/journal/ActiveSegment.java +++ b/src/java/org/apache/cassandra/journal/ActiveSegment.java @@ -30,6 +30,9 @@ import java.util.function.Consumer; import accord.utils.Invariants; import com.codahale.metrics.Timer; import com.google.common.annotations.VisibleForTesting; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.utils.*; import org.apache.cassandra.utils.concurrent.OpOrder; @@ -42,6 +45,8 @@ import static org.apache.cassandra.utils.Simulate.With.MONITORS; @Simulate(with=MONITORS) public final class ActiveSegment extends Segment { + private static final Logger logger = LoggerFactory.getLogger(ActiveSegment.class); + final FileChannel channel; // OpOrder used to order appends wrt flush @@ -197,6 +202,7 @@ public final class ActiveSegment extends Segment private void discard() { + logger.debug("Discarding {}", this); selfRef.ensureReleased(); descriptor.fileFor(Component.DATA).deleteIfExists(); diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java index d07cae1c59..2545f622f1 100644 --- a/src/java/org/apache/cassandra/journal/Journal.java +++ b/src/java/org/apache/cassandra/journal/Journal.java @@ -936,9 +936,9 @@ public class Journal implements Shutdownable /** * Static segment iterator iterates all keys in _static_ segments in order. */ - public StaticSegmentKeyIterator staticSegmentKeyIterator() + public StaticSegmentKeyIterator staticSegmentKeyIterator(K min, K max) { - return new StaticSegmentKeyIterator(); + return new StaticSegmentKeyIterator(min, max); } /** @@ -1011,17 +1011,22 @@ public class Journal implements Shutdownable private final ReferencedSegments segments; private final MergeIterator> iterator; - public StaticSegmentKeyIterator() + public StaticSegmentKeyIterator(K min, K max) { - this.segments = selectAndReference(Segment::isStatic); + this.segments = selectAndReference(s -> s.isStatic() && (min == null || keySupport.compare(s.index().lastId(), min) >= 0) && (max == null || keySupport.compare(s.index().firstId(), max) <= 0)); List> iterators = new ArrayList<>(segments.count()); for (Segment segment : segments.allSorted(true)) { - StaticSegment staticSegment = (StaticSegment) segment; + final StaticSegment staticSegment = (StaticSegment) segment; + final OnDiskIndex.IndexReader iter = staticSegment.index().reader(); + if (min != null) iter.seek(min); + if (max != null) iter.seekEnd(max); + if (!iter.hasNext()) + continue; + iterators.add(new AbstractIterator<>() { - final Iterator iter = staticSegment.index().reader(); final Head head = new Head(staticSegment.descriptor.timestamp); @Override diff --git a/src/java/org/apache/cassandra/journal/OnDiskIndex.java b/src/java/org/apache/cassandra/journal/OnDiskIndex.java index 5fae9255e3..0a2ca8cc95 100644 --- a/src/java/org/apache/cassandra/journal/OnDiskIndex.java +++ b/src/java/org/apache/cassandra/journal/OnDiskIndex.java @@ -252,6 +252,7 @@ final class OnDiskIndex extends Index public class IndexReader extends AbstractIterator { + int lastIdx = entryCount - 1; int idx; K key; int offset; @@ -262,6 +263,20 @@ final class OnDiskIndex extends Index idx = -1; } + public void seek(K key) + { + int i = binarySearch(key); + if (i < 0) i = -1 - i; + idx = i - 1; + } + + public void seekEnd(K key) + { + int i = binarySearch(key); + if (i < 0) i = -2 - i; + lastIdx = i; + } + protected K computeNext() { if (advance()) @@ -284,7 +299,7 @@ final class OnDiskIndex extends Index public boolean advance() { - if (idx >= entryCount - 1) + if (idx >= lastIdx) return false; idx++; diff --git a/src/java/org/apache/cassandra/journal/Segments.java b/src/java/org/apache/cassandra/journal/Segments.java index a2475557a6..bdd447ec6b 100644 --- a/src/java/org/apache/cassandra/journal/Segments.java +++ b/src/java/org/apache/cassandra/journal/Segments.java @@ -220,6 +220,12 @@ class Segments } } + @Override + public String toString() + { + return sorted.toString(); + } + private static final Long2ObjectHashMap EMPTY_MAP = new Long2ObjectHashMap<>(); @SuppressWarnings("unchecked") diff --git a/src/java/org/apache/cassandra/journal/StaticSegment.java b/src/java/org/apache/cassandra/journal/StaticSegment.java index 658224e056..35c987c8a4 100644 --- a/src/java/org/apache/cassandra/journal/StaticSegment.java +++ b/src/java/org/apache/cassandra/journal/StaticSegment.java @@ -164,6 +164,8 @@ public final class StaticSegment extends Segment */ void discard(Journal journal) { + logger.debug("Discarding {}", this); + ((Tidier)selfRef.tidier()).discard = true; close(journal); } diff --git a/src/java/org/apache/cassandra/metrics/AccordMetrics.java b/src/java/org/apache/cassandra/metrics/AccordMetrics.java index f489ec899e..7c6a48d38d 100644 --- a/src/java/org/apache/cassandra/metrics/AccordMetrics.java +++ b/src/java/org/apache/cassandra/metrics/AccordMetrics.java @@ -201,12 +201,14 @@ public class AccordMetrics private AccordMetrics forTransaction(TxnId txnId) { - if (txnId.isWrite()) - return writeMetrics; - else if (txnId.isSomeRead()) - return readMetrics; - else - return null; + if (txnId != null) + { + if (txnId.isWrite()) + return writeMetrics; + else if (txnId.isSomeRead()) + return readMetrics; + } + return null; } @Override diff --git a/src/java/org/apache/cassandra/schema/TableId.java b/src/java/org/apache/cassandra/schema/TableId.java index d2730c7f67..7cebcf2c66 100644 --- a/src/java/org/apache/cassandra/schema/TableId.java +++ b/src/java/org/apache/cassandra/schema/TableId.java @@ -38,6 +38,7 @@ import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.tcm.serialization.MetadataSerializer; import org.apache.cassandra.tcm.serialization.Version; import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.Hex; import org.apache.cassandra.utils.ObjectSizes; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.UUIDGen; @@ -92,6 +93,8 @@ public final class TableId implements Comparable public static TableId fromString(String idString) { + if (idString.startsWith("tid:")) + return new TableId(MAGIC, Hex.parseLong(idString, 4, idString.length())); return new TableId(UUID.fromString(idString)); } @@ -185,6 +188,13 @@ public final class TableId implements Comparable @Override public String toString() + { + if (msb == MAGIC) + return "tid:" + Long.toHexString(lsb); + return asUUID().toString(); + } + + public String toLongString() { return asUUID().toString(); } diff --git a/src/java/org/apache/cassandra/schema/TableMetadata.java b/src/java/org/apache/cassandra/schema/TableMetadata.java index c2909dcb61..011a5051df 100644 --- a/src/java/org/apache/cassandra/schema/TableMetadata.java +++ b/src/java/org/apache/cassandra/schema/TableMetadata.java @@ -1693,7 +1693,7 @@ public class TableMetadata implements SchemaElement { if (withInternals) builder.append("ID = ") - .append(id.toString()) + .append(id.toLongString()) .newLine() .append("AND "); diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 7c1e66857e..f7e8c3004d 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -508,7 +508,7 @@ public class AccordCommandStore extends CommandStore public AsyncChain load(TxnId txnId) { return store.submit(txnId, safeStore -> { - maybeApplyWrites(safeStore, txnId); + initialiseState(safeStore, txnId); return safeStore.unsafeGet(txnId).current().route(); }); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index 1ece0d3c3c..6797b59ec0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -848,6 +848,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo final SequentialQueueTask selfTask; private Task task; private volatile Thread owner, waiting; + private boolean running; SequentialExecutor() { @@ -859,12 +860,13 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo { Invariants.require(task != null); task.preRunExclusive(runner); + running = true; } void runTask() { Thread self = Thread.currentThread(); - if (!ownerUpdater.compareAndSet(this, null, self)) + while (!ownerUpdater.compareAndSet(this, null, self)) { waiting = self; while (owner != null) @@ -882,6 +884,7 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo void cleanupTask(@Nullable TaskRunner runner) { task.cleanupExclusive(runner); + running = false; owner = null; task = super.poll(); if (task != null) @@ -910,19 +913,34 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo @Override protected void remove(Task remove) { + Invariants.require(remove != null); if (remove != task) { super.remove(remove); } - else + else if (!running) { - Invariants.require(waitingToRun.contains(selfTask)); + // cannot overwrite task while it is being executed - this cannot happen for AccordTask + // but can for other tasks that don't track their own state + task = super.poll(); - if (task == null) waitingToRun.remove(selfTask); + if (waitingToRun.contains(selfTask)) + { + if (task == null) waitingToRun.remove(selfTask); + else + { + selfTask.queuePosition = task.queuePosition; + waitingToRun.update(selfTask); + } + } else { - selfTask.queuePosition = task.queuePosition; - waitingToRun.update(selfTask); + Invariants.expect(false, "%s should have been queued to run as it had the task %s pending, that has now been cancelled", this, remove); + if (task != null) + { + selfTask.queuePosition = task.queuePosition; + waitingToRun.append(selfTask); + } } } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournal.java b/src/java/org/apache/cassandra/service/accord/AccordJournal.java index 5ed6983dec..dfb4636c4c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournal.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournal.java @@ -19,12 +19,15 @@ package org.apache.cassandra.service.accord; import java.io.IOException; import java.nio.ByteBuffer; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.Deque; import java.util.Iterator; import java.util.List; import java.util.NavigableMap; +import java.util.Queue; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; @@ -34,6 +37,9 @@ import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import accord.impl.CommandChange; import accord.impl.CommandChange.Field; import accord.local.Cleanup; @@ -46,6 +52,7 @@ import accord.local.Node; import accord.local.RedundantBefore; import accord.primitives.EpochSupplier; import accord.primitives.Ranges; +import accord.primitives.Route; import accord.primitives.SaveStatus; import accord.primitives.Status.Durability; import accord.primitives.Timestamp; @@ -55,6 +62,8 @@ import accord.utils.PersistentField; import accord.utils.UnhandledEnum; import accord.utils.async.AsyncResult; import accord.utils.async.AsyncResults; +import org.agrona.collections.Int2ObjectHashMap; +import org.agrona.collections.IntArrayList; import org.apache.cassandra.concurrent.Shutdownable; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; @@ -63,6 +72,7 @@ import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.io.util.File; +import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.journal.Compactor; import org.apache.cassandra.journal.Journal; import org.apache.cassandra.journal.Params; @@ -70,7 +80,6 @@ import org.apache.cassandra.journal.RecordPointer; import org.apache.cassandra.journal.SegmentCompactor; import org.apache.cassandra.journal.StaticSegment; import org.apache.cassandra.journal.ValueSerializer; -import org.apache.cassandra.service.accord.AccordCommandStore.AccordCommandStoreLoader; import org.apache.cassandra.service.accord.AccordJournalValueSerializers.FlyweightImage; import org.apache.cassandra.service.accord.AccordJournalValueSerializers.IdentityAccumulator; import org.apache.cassandra.service.accord.JournalKey.JournalKeySupport; @@ -81,9 +90,9 @@ import org.apache.cassandra.service.accord.serializers.DepsSerializers; import org.apache.cassandra.service.accord.serializers.ResultSerializers; import org.apache.cassandra.service.accord.serializers.Version; import org.apache.cassandra.service.accord.serializers.WaitingOnSerializer; +import org.apache.cassandra.utils.Closeable; import org.apache.cassandra.utils.CloseableIterator; import org.apache.cassandra.utils.ExecutorUtils; -import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.concurrent.Semaphore; import static accord.impl.CommandChange.Field.CLEANUP; @@ -99,9 +108,12 @@ import static accord.impl.CommandChange.unsetIterable; import static accord.impl.CommandChange.validateFlags; import static accord.local.Cleanup.Input.FULL; import static org.apache.cassandra.service.accord.AccordJournalValueSerializers.DurableBeforeAccumulator; +import static org.apache.cassandra.service.accord.JournalKey.Type.COMMAND_DIFF; +import static org.apache.cassandra.utils.FBUtilities.getAvailableProcessors; public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier, Shutdownable { + private static final Logger logger = LoggerFactory.getLogger(AccordJournal.class); static final ThreadLocal keyCRCBytes = ThreadLocal.withInitial(() -> new byte[JournalKeySupport.TOTAL_SIZE]); @VisibleForTesting @@ -257,7 +269,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier @Override public List debugCommand(int commandStoreId, TxnId txnId) { - JournalKey key = new JournalKey(txnId, JournalKey.Type.COMMAND_DIFF, commandStoreId); + JournalKey key = new JournalKey(txnId, COMMAND_DIFF, commandStoreId); List result = new ArrayList<>(); journalTable.readAll(key, (long segment, int position, JournalKey k, ByteBuffer buffer, int userVersion) -> { Builder builder = new Builder(txnId); @@ -325,7 +337,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier return; } - JournalKey key = new JournalKey(update.txnId, JournalKey.Type.COMMAND_DIFF, commandStoreId); + JournalKey key = new JournalKey(update.txnId, COMMAND_DIFF, commandStoreId); RecordPointer pointer = journal.asyncWrite(key, diff); if (journalTable.shouldIndex(key) && diff.hasParticipants() @@ -403,7 +415,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier private Builder loadDiffs(int commandStoreId, TxnId txnId, Load load) { - JournalKey key = new JournalKey(txnId, JournalKey.Type.COMMAND_DIFF, commandStoreId); + JournalKey key = new JournalKey(txnId, COMMAND_DIFF, commandStoreId); Builder builder = new Builder(txnId, load); journalTable.readAll(key, builder::deserializeNext); return builder; @@ -478,7 +490,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier public void forEach(Consumer consumer) { - try (CloseableIterator> iter = journalTable.keyIterator()) + try (CloseableIterator> iter = journalTable.keyIterator(null, null)) { while (iter.hasNext()) { @@ -492,54 +504,207 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier @Override public void replay(CommandStores commandStores) { - final Semaphore concurrency = Semaphore.newSemaphore(FBUtilities.getAvailableProcessors()); + // TODO (expected): make the parallelisms configurable + // Replay is performed in parallel, where at most X commands can be in flight, accross at most Y commands stores. + // That is, you can limit replay parallelism to 1 command store at a time, but load multiple commands within that data store, + // _or_ have multiple commands being loaded accross multiple data stores. + final Semaphore commandParallelism = Semaphore.newSemaphore(getAvailableProcessors()); + final int commandStoreParallelism = Math.max(Math.max(1, Math.min(getAvailableProcessors(), 4)), getAvailableProcessors() / 4); final AtomicBoolean abort = new AtomicBoolean(); + // TODO (expected): balance work submission by AccordExecutor + final IntArrayList activeCommandStoreIds = new IntArrayList(); + final ReplayQueue pendingCommandStores = new ReplayQueue(commandStores.all()); - try (CloseableIterator> iter = journalTable.keyIterator()) + class ReplayStream implements Closeable { - JournalKey prev = null; - while (iter.hasNext()) - { - if (abort.get()) - break; + final CommandStore commandStore; + final Loader loader; + final CloseableIterator> iter; + JournalKey prev; + public ReplayStream(CommandStore commandStore) + { + this.commandStore = commandStore; + this.loader = commandStore.loader(); + // 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())); + } + + boolean replay() + { JournalKey key; long[] segments; + while (true) { + if (!iter.hasNext()) + { + logger.info("Completed replay of {}", commandStore); + return false; + } + Journal.KeyRefs ref = iter.next(); - key = ref.key(); - if (key.type != JournalKey.Type.COMMAND_DIFF) + if (ref.key().type != COMMAND_DIFF) continue; + key = ref.key(); segments = journalTable.shouldIndex(key) ? ref.copyOfSegments() : null; + break; } - CommandStore commandStore = commandStores.forId(key.commandStoreId); - AccordCommandStoreLoader loader = (AccordCommandStoreLoader) commandStore.loader(); - TxnId txnId = key.id; Invariants.require(prev == null || key.commandStoreId != prev.commandStoreId || key.id.compareTo(prev.id) != 0, "duplicate key detected %s == %s", key, prev); prev = key; - - concurrency.acquireThrowUncheckedOnInterrupt(1); + commandParallelism.acquireThrowUncheckedOnInterrupt(1); loader.load(txnId) .map(route -> { - if (segments != null) + if (segments != null && route != null) { for (long segment : segments) - journalTable.safeNotify(index -> index.update(segment, key.commandStoreId, txnId, route)); + journalTable.safeNotify(index -> index.update(segment, key.commandStoreId, txnId, (Route)route)); } return null; }).begin((success, fail) -> { - concurrency.release(1); + commandParallelism.release(1); if (fail != null && !journal.handleError("Could not replay command " + txnId, fail)) abort.set(true); }); + + return true; + } + + @Override + public void close() + { + iter.close(); } } + + // Replay streams by command store id, can hold at most commandStoreParallelism items + final Int2ObjectHashMap replayStreams = new Int2ObjectHashMap<>(); + try + { + // index of the store we're currently pulling from in the activeCommandStoreIds collection + int cur = 0; + while (!abort.get()) + { + if (cur == activeCommandStoreIds.size()) + { + if (activeCommandStoreIds.size() < commandStoreParallelism && !pendingCommandStores.isEmpty()) + { + CommandStore next = pendingCommandStores.next(); + int id = next.id(); + activeCommandStoreIds.add(id); + replayStreams.put(id, new ReplayStream(next)); + } + else if (activeCommandStoreIds.isEmpty()) break; + else cur = 0; + } + + int id = activeCommandStoreIds.get(cur); + ReplayStream replayStream = replayStreams.get(id); + while (!replayStream.replay()) + { + // Replay complete for this command store; close and replace + replayStreams.remove(id).close(); + if (pendingCommandStores.isEmpty()) + { + // no more pending to submit; remove and continue with the next remaining (if any) + activeCommandStoreIds.removeAt(cur); + if (cur == activeCommandStoreIds.size()) + --cur; + if (cur < 0) + break; + id = activeCommandStoreIds.get(cur); + } + else + { + // replace it with a pending command store, and continue processing + CommandStore next = pendingCommandStores.next(streamId(replayStream.commandStore)); + id = next.id(); + activeCommandStoreIds.set(cur, id); + replayStreams.put(id, new ReplayStream(next)); + } + + replayStream = replayStreams.get(id); + } + + ++cur; + } + } + catch (Throwable t) + { + try { FileUtils.close(replayStreams.values()); } + catch (Throwable t2) { t.addSuppressed(t2); } + throw t; + } + } + + static class ReplayQueue + { + final Int2ObjectHashMap> byExecutor = new Int2ObjectHashMap<>(); + final Deque nextId = new ArrayDeque<>(); + + ReplayQueue(CommandStore[] commandStores) + { + for (CommandStore commandStore : commandStores) + { + byExecutor.computeIfAbsent(streamId(commandStore), ignore -> new ArrayDeque<>()) + .add(commandStore); + } + nextId.addAll(byExecutor.keySet()); + } + + boolean isEmpty() + { + return byExecutor.isEmpty(); + } + + CommandStore next() + { + while (true) + { + if (byExecutor.isEmpty()) + return null; + + Integer id = nextId.poll(); + if (id == null) + { + nextId.addAll(byExecutor.keySet()); + id = nextId.poll(); + } + + Queue queue = byExecutor.get(id); + if (queue != null) + { + CommandStore next = queue.poll(); + if (queue.isEmpty()) + byExecutor.remove(id); + if (next != null) + return next; + } + } + } + + CommandStore next(int streamId) + { + Queue queue = byExecutor.get(streamId); + if (queue == null) + return next(); + + CommandStore next = queue.poll(); + if (queue.isEmpty()) + byExecutor.remove(streamId); + + return next; + } + } + + private static int streamId(CommandStore commandStore) + { + return commandStore instanceof AccordCommandStore ? ((AccordCommandStore) commandStore).executor().executorId() : 1; } public static @Nullable ByteBuffer asSerializedChange(Command before, Command after, Version userVersion) throws IOException @@ -831,7 +996,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier for (Field field = nextSetField(iterable) ; field != null; field = nextSetField(iterable = unsetIterable(field, iterable))) { // Since we are iterating in reverse order, we skip the fields that were - // set by entries writter later (i.e. already read ones). + // set by entries written later (i.e. already read ones). if (isChanged(field, flags) && field != CLEANUP) skip(txnId, field, in, userVersion); else diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java index d04038ffcc..7d4b2041d4 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java @@ -64,6 +64,7 @@ import org.apache.cassandra.db.rows.Row; import org.apache.cassandra.db.rows.Unfiltered; import org.apache.cassandra.db.rows.UnfilteredRowIterator; import org.apache.cassandra.db.rows.UnfilteredRowIterators; +import org.apache.cassandra.dht.Bounds; import org.apache.cassandra.index.Index; import org.apache.cassandra.index.accord.OrderedRouteSerializer; import org.apache.cassandra.index.accord.RouteJournalIndex; @@ -78,6 +79,7 @@ import org.apache.cassandra.journal.KeySupport; import org.apache.cassandra.journal.RecordConsumer; import org.apache.cassandra.schema.ColumnMetadata; import org.apache.cassandra.service.RetryStrategy; +import org.apache.cassandra.service.accord.AccordKeyspace.JournalColumns; import org.apache.cassandra.service.accord.api.TokenKey; import org.apache.cassandra.service.accord.serializers.Version; import org.apache.cassandra.utils.CloseableIterator; @@ -86,6 +88,7 @@ import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.MergeIterator; import static org.apache.cassandra.io.sstable.SSTableReadsListener.NOOP_LISTENER; +import static org.apache.cassandra.service.accord.AccordKeyspace.JournalColumns.getJournalKey; public class AccordJournalTable implements RangeSearcher.Supplier { @@ -375,7 +378,7 @@ public class AccordJournalTable implements RangeSearche public TxnId next() { UnfilteredRowIterator next = partitionIterator.next(); - JournalKey partitionKeyComponents = AccordKeyspace.JournalColumns.getJournalKey(next.partitionKey()); + JournalKey partitionKeyComponents = getJournalKey(next.partitionKey()); Invariants.require(partitionKeyComponents.commandStoreId == storeId, () -> String.format("table index returned a command store other than the exepcted one; expected %d != %d", storeId, partitionKeyComponents.commandStoreId)); return partitionKeyComponents.id; @@ -404,7 +407,7 @@ public class AccordJournalTable implements RangeSearche private void readAllFromTable(K key, TableRecordConsumer onEntry) { - DecoratedKey pk = AccordKeyspace.JournalColumns.decorate(key); + DecoratedKey pk = JournalColumns.decorate(key); try (RefViewFragment view = cfs.selectAndReference(View.select(SSTableSet.LIVE, pk))) { if (view.sstables.isEmpty()) @@ -460,9 +463,9 @@ public class AccordJournalTable implements RangeSearche } @SuppressWarnings("resource") // Auto-closeable iterator will release related resources - public CloseableIterator> keyIterator() + public CloseableIterator> keyIterator(@Nullable K min, @Nullable K max) { - return new JournalAndTableKeyIterator(); + return new JournalAndTableKeyIterator(min, max); } private class TableIterator extends AbstractIterator implements CloseableIterator @@ -470,12 +473,17 @@ public class AccordJournalTable implements RangeSearche private final UnfilteredPartitionIterator mergeIterator; private final RefViewFragment view; - private TableIterator() + private TableIterator(JournalKey min, JournalKey max) { - view = cfs.selectAndReference(v -> v.select(SSTableSet.LIVE)); + Invariants.require((min != null && max != null) || min == max); + view = cfs.selectAndReference(View.select(SSTableSet.LIVE, r -> (max == null || JournalKey.SUPPORT.compare(getJournalKey(r.getFirst()), max) <= 0) + && (min == null || JournalKey.SUPPORT.compare(getJournalKey(r.getLast()), min) >= 0))); List scanners = new ArrayList<>(); for (SSTableReader sstable : view.sstables) - scanners.add(sstable.getScanner()); + { + if (min == null) scanners.add(sstable.getScanner()); + else scanners.add(sstable.getScanner(new Bounds(JournalColumns.decorate(min), JournalColumns.decorate(max)))); + } mergeIterator = view.sstables.isEmpty() ? EmptyIterators.unfilteredPartition(cfs.metadata()) @@ -490,7 +498,7 @@ public class AccordJournalTable implements RangeSearche { try (UnfilteredRowIterator partition = mergeIterator.next()) { - ret = (K) AccordKeyspace.JournalColumns.getJournalKey(partition.partitionKey()); + ret = (K) getJournalKey(partition.partitionKey()); while (partition.hasNext()) partition.next(); } @@ -515,10 +523,10 @@ public class AccordJournalTable implements RangeSearche final TableIterator tableIterator; final Journal.StaticSegmentKeyIterator journalIterator; - private JournalAndTableKeyIterator() + private JournalAndTableKeyIterator(K min, K max) { - this.tableIterator = new TableIterator(); - this.journalIterator = journal.staticSegmentKeyIterator(); + this.tableIterator = new TableIterator(min, max); + this.journalIterator = journal.staticSegmentKeyIterator(min, max); } K prevFromTable = null; diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index a36d01c97f..7a534a41d8 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -37,6 +37,7 @@ import javax.annotation.concurrent.GuardedBy; import com.google.common.annotations.VisibleForTesting; import com.google.common.primitives.Ints; + import org.apache.cassandra.utils.Clock; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -250,12 +251,13 @@ public class AccordService implements IAccordService, Shutdownable public static void replayJournal(AccordService as) { logger.info("Starting journal replay."); + long before = Clock.Global.nanoTime(); CommandsForKey.disableLinearizabilityViolationsReporting(); try { AccordKeyspace.truncateAllCaches(); - as.journal().replay(as.node().commandStores()); + as.journal().replay(as.node().commandStores()); logger.info("Waiting for command stores to quiesce."); ((AccordCommandStores)as.node.commandStores()).waitForQuiescense(); as.journal.unsafeSetStarted(); @@ -265,7 +267,8 @@ public class AccordService implements IAccordService, Shutdownable CommandsForKey.enableLinearizabilityViolationsReporting(); } - logger.info("Finished journal replay."); + long after = Clock.Global.nanoTime(); + logger.info("Finished journal replay. {}ms elapsed", NANOSECONDS.toMillis(after - before)); } public static void shutdownServiceAndWait(long timeout, TimeUnit unit) throws InterruptedException, TimeoutException @@ -773,88 +776,96 @@ public class AccordService implements IAccordService, Shutdownable private AsyncChain loadDebug(TxnId txnId, CommandStore store) { CommandStoreTxnBlockedGraph.Builder state = new CommandStoreTxnBlockedGraph.Builder(store.id()); - return populate(state, store, txnId).map(ignore -> state.build()); + populateAsync(state, store, txnId); + return state; } - private static AsyncChain populate(CommandStoreTxnBlockedGraph.Builder state, CommandStore store, TxnId txnId) + private static void populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TxnId blockedBy) { - AsyncChain> submit = store.submit(txnId, in -> { - AsyncChain chain = populate(state, (AccordSafeCommandStore) in, txnId); - return chain == null ? AsyncChains.success(null) : chain; - }); - return submit.flatMap(Function.identity()); + if (safeStore.ifLoadedAndInitialised(blockedBy) != null) populateSync(state, safeStore, blockedBy); + else populateAsync(state, safeStore.commandStore(), blockedBy); } - private static AsyncChain populate(CommandStoreTxnBlockedGraph.Builder state, CommandStore commandStore, TokenKey blockedBy, TxnId txnId, Timestamp executeAt) + private static void populateAsync(CommandStoreTxnBlockedGraph.Builder state, CommandStore store, TxnId txnId) { - AsyncChain> submit = commandStore.submit(PreLoadContext.contextFor(txnId, RoutingKeys.of(blockedBy.toUnseekable()), KeyHistory.SYNC), in -> { - AsyncChain chain = populate(state, (AccordSafeCommandStore) in, blockedBy, txnId, executeAt); - return chain == null ? AsyncChains.success(null) : chain; + state.asyncTxns.incrementAndGet(); + store.execute(txnId, in -> { + populateSync(state, (AccordSafeCommandStore) in, txnId); + if (0 == state.asyncTxns.decrementAndGet() && 0 == state.asyncKeys.get()) + state.complete(); }); - return submit.flatMap(Function.identity()); } @Nullable - private static AsyncChain populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TxnId txnId) + private static void populateSync(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TxnId txnId) { - SafeCommand safeCommand = safeStore.unsafeGet(txnId); - Invariants.nonNull(safeCommand, "Txn %s is not in the cache", txnId); - if (safeCommand.current() == null || safeCommand.current().saveStatus() == SaveStatus.Uninitialised) - return null; - CommandStoreTxnBlockedGraph.TxnState cmdTxnState = populate(state, safeCommand.current()); - if (cmdTxnState.notBlocked()) - return null; - //TODO (expected): check depth - List> chains = new ArrayList<>(); - for (TxnId blockedBy : cmdTxnState.blockedBy) + try { - if (state.knows(blockedBy)) continue; - // need to fetch the state - if (safeStore.ifLoadedAndInitialised(blockedBy) != null) + if (state.txns.containsKey(txnId)) + return; // could plausibly request same txn twice + + SafeCommand safeCommand = safeStore.unsafeGet(txnId); + Invariants.nonNull(safeCommand, "Txn %s is not in the cache", txnId); + if (safeCommand.current() == null || safeCommand.current().saveStatus() == SaveStatus.Uninitialised) + return; + + CommandStoreTxnBlockedGraph.TxnState cmdTxnState = populateSync(state, safeCommand.current()); + if (cmdTxnState.notBlocked()) + return; + + for (TxnId blockedBy : cmdTxnState.blockedBy) { - AsyncChain chain = populate(state, safeStore, blockedBy); - if (chain != null) - chains.add(chain); + if (!state.knows(blockedBy)) + populate(state, safeStore, blockedBy); } - else + for (TokenKey blockedBy : cmdTxnState.blockedByKey) { - // go fetch it - chains.add(populate(state, safeStore.commandStore(), blockedBy)); + if (!state.keys.containsKey(blockedBy)) + populate(state, safeStore, blockedBy, txnId, safeCommand.current().executeAt()); } } - for (TokenKey blockedBy : cmdTxnState.blockedByKey) + catch (Throwable t) { - if (state.keys.containsKey(blockedBy)) continue; - if (safeStore.ifLoadedAndInitialised(blockedBy) != null) - { - AsyncChain chain = populate(state, safeStore, blockedBy, txnId, safeCommand.current().executeAt()); - if (chain != null) - chains.add(chain); - } - else - { - // go fetch it - chains.add(populate(state, safeStore.commandStore(), blockedBy, txnId, safeCommand.current().executeAt())); - } + state.tryFailure(t); } - if (chains.isEmpty()) - return null; - return AsyncChains.allOf(chains).map(ignore -> null); } - private static AsyncChain populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TokenKey pk, TxnId txnId, Timestamp executeAt) + private static void populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TokenKey blockedBy, TxnId txnId, Timestamp executeAt) { - SafeCommandsForKey commandsForKey = safeStore.ifLoadedAndInitialised(pk); - TxnId blocking = commandsForKey.current().blockedOnTxnId(txnId, executeAt); - if (blocking instanceof CommandsForKey.TxnInfo) - blocking = ((CommandsForKey.TxnInfo) blocking).plainTxnId(); - state.keys.put(pk, blocking); - if (state.txns.containsKey(blocking)) return null; - if (safeStore.ifLoadedAndInitialised(blocking) != null) return populate(state, safeStore, blocking); - return populate(state, safeStore.commandStore(), blocking); + if (safeStore.ifLoadedAndInitialised(txnId) != null && safeStore.ifLoadedAndInitialised(blockedBy) != null) populateSync(state, safeStore, blockedBy, txnId, executeAt); + else populateAsync(state, safeStore.commandStore(), blockedBy, txnId, executeAt); } - private static CommandStoreTxnBlockedGraph.TxnState populate(CommandStoreTxnBlockedGraph.Builder state, Command cmd) + private static void populateAsync(CommandStoreTxnBlockedGraph.Builder state, CommandStore commandStore, TokenKey blockedBy, TxnId txnId, Timestamp executeAt) + { + state.asyncKeys.incrementAndGet(); + commandStore.execute(PreLoadContext.contextFor(txnId, RoutingKeys.of(blockedBy.toUnseekable()), KeyHistory.SYNC), in -> { + populateSync(state, (AccordSafeCommandStore) in, blockedBy, txnId, executeAt); + if (0 == state.asyncKeys.decrementAndGet() && 0 == state.asyncTxns.get()) + state.complete(); + }); + } + + private static void populateSync(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TokenKey pk, TxnId txnId, Timestamp executeAt) + { + try + { + SafeCommandsForKey commandsForKey = safeStore.ifLoadedAndInitialised(pk); + TxnId blocking = commandsForKey.current().blockedOnTxnId(txnId, executeAt); + if (blocking instanceof CommandsForKey.TxnInfo) + blocking = ((CommandsForKey.TxnInfo) blocking).plainTxnId(); + state.keys.put(pk, blocking); + if (state.txns.containsKey(blocking)) + return; + populate(state, safeStore, blocking); + } + catch (Throwable t) + { + state.tryFailure(t); + } + } + + private static CommandStoreTxnBlockedGraph.TxnState populateSync(CommandStoreTxnBlockedGraph.Builder state, Command cmd) { CommandStoreTxnBlockedGraph.Builder.TxnBuilder cmdTxnState = state.txn(cmd.txnId(), cmd.executeAt(), cmd.saveStatus()); if (!cmd.hasBeen(Status.Applied) && cmd.hasBeen(Status.Stable)) diff --git a/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java b/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java index 40f2c3703e..7d553dda7c 100644 --- a/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java +++ b/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java @@ -24,6 +24,7 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.atomic.AtomicInteger; import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; @@ -32,6 +33,7 @@ import com.google.common.collect.ImmutableSet; import accord.primitives.SaveStatus; import accord.primitives.Timestamp; import accord.primitives.TxnId; +import accord.utils.async.AsyncResults; import org.apache.cassandra.service.accord.api.TokenKey; public class CommandStoreTxnBlockedGraph @@ -75,8 +77,9 @@ public class CommandStoreTxnBlockedGraph } } - public static class Builder + public static class Builder extends AsyncResults.SettableResult { + final AtomicInteger asyncTxns = new AtomicInteger(), asyncKeys = new AtomicInteger(); final int storeId; final Map txns = new LinkedHashMap<>(); final Map keys = new LinkedHashMap<>(); @@ -91,6 +94,11 @@ public class CommandStoreTxnBlockedGraph return txns.containsKey(id); } + public void complete() + { + trySuccess(build()); + } + public CommandStoreTxnBlockedGraph build() { return new CommandStoreTxnBlockedGraph(this); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java index 1057afe13e..4a0481aa09 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/CommandSerializers.java @@ -487,6 +487,7 @@ public class CommandSerializers static { Invariants.require(EPOCH_MASK << EPOCH_SHIFT >= 0); + Invariants.require(EPOCH_SHIFT + Integer.bitCount(EPOCH_MASK) < 8); } interface Factory @@ -537,6 +538,7 @@ public class CommandSerializers | encodeLength(hlcLength, HLC_SHIFT, HLC_MIN_LENGTH, HLC_MASK) | encodeLength(flagsLength, FLAGS_SHIFT, FLAGS_MIN_LENGTH, FLAGS_MASK) | encodeLength(nodeLength, NODE_SHIFT, NODE_MIN_LENGTH, NODE_MASK); + Invariants.require(((byte)encodingFlags) >= 0); out.writeByte(encodingFlags); out.writeLeastSignificantBytes(epoch, epochLength); out.writeLeastSignificantBytes(hlc, hlcLength); @@ -602,6 +604,7 @@ public class CommandSerializers | encodeLength(hlcLength, HLC_SHIFT, HLC_MIN_LENGTH, HLC_MASK) | encodeLength(flagsLength, FLAGS_SHIFT, FLAGS_MIN_LENGTH, FLAGS_MASK) | encodeLength(nodeLength, NODE_SHIFT, NODE_MIN_LENGTH, NODE_MASK); + Invariants.require(((byte)encodingFlags) >= 0); int position = offset; position += accessor.putByte(dst, position, (byte)encodingFlags); diff --git a/src/java/org/apache/cassandra/service/accord/serializers/LatestDepsSerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/LatestDepsSerializers.java index 0dda159259..26d95b0494 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/LatestDepsSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/LatestDepsSerializers.java @@ -127,6 +127,7 @@ public class LatestDepsSerializers @Override public void serializeBody(GetLatestDeps msg, DataOutputPlus out, Version version) throws IOException { + CommandSerializers.ballot.serialize(msg.ballot, out); ExecuteAtSerializer.serialize(msg.executeAt, out); } @@ -141,7 +142,8 @@ public class LatestDepsSerializers @Override public long serializedBodySize(GetLatestDeps msg, Version version) { - return ExecuteAtSerializer.serializedSize(msg.executeAt); + return CommandSerializers.ballot.serializedSize(msg.ballot) + + ExecuteAtSerializer.serializedSize(msg.executeAt); } }; diff --git a/src/java/org/apache/cassandra/utils/CloseableIterator.java b/src/java/org/apache/cassandra/utils/CloseableIterator.java index 634629f4be..b9cda68465 100644 --- a/src/java/org/apache/cassandra/utils/CloseableIterator.java +++ b/src/java/org/apache/cassandra/utils/CloseableIterator.java @@ -21,10 +21,8 @@ import java.util.Iterator; import java.util.NoSuchElementException; // so we can instantiate anonymous classes implementing both interfaces -public interface CloseableIterator extends Iterator, AutoCloseable +public interface CloseableIterator extends Iterator, Closeable { - public void close(); - public static CloseableIterator wrap(Iterator iter) { return new CloseableIterator() diff --git a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java index a3ce82e2d9..40b76a7a9b 100644 --- a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java +++ b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java @@ -31,7 +31,6 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.Set; -import java.util.UUID; import java.util.concurrent.Callable; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; @@ -1676,7 +1675,7 @@ public class ClusterUtils public static TableId tableId(Cluster cluster, String ks, String table) { String str = cluster.getFirstRunningInstance().callOnInstance(() -> Schema.instance.getKeyspaceInstance(ks).getColumnFamilyStore(table).getTableId().toString()); - return TableId.fromUUID(UUID.fromString(str)); + return TableId.fromString(str); } public static void awaitAccordEpochReady(Cluster cluster, long epoch) diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java index 9b67dc4453..416fe065fb 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordDropTableBase.java @@ -18,8 +18,6 @@ package org.apache.cassandra.distributed.test.accord; -import java.util.UUID; - import com.google.common.base.Throwables; import accord.api.RoutingKey; @@ -128,7 +126,7 @@ public class AccordDropTableBase extends TestBaseImpl for (IInvokableInstance inst : cluster) { inst.runOnInstance(() -> { - TableId tableId = TableId.fromUUID(UUID.fromString(s)); + TableId tableId = TableId.fromString(s); AccordService accord = (AccordService) AccordService.instance(); PreLoadContext ctx = PreLoadContext.contextFor(Ranges.single(TokenRange.fullRange(tableId, getPartitioner())), KeyHistory.SYNC); CommandStores stores = accord.node().commandStores(); diff --git a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java index cde8789060..bd20c03cb3 100644 --- a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java +++ b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java @@ -332,7 +332,7 @@ public class AccordJournalBurnTest extends BurnTestBase private TreeMap read(CommandStores commandStores) { TreeMap result = new TreeMap<>(JournalKey.SUPPORT::compare); - try (CloseableIterator> iter = journalTable.keyIterator()) + try (CloseableIterator> iter = journalTable.keyIterator(null, null)) { JournalKey prev = null; while (iter.hasNext()) diff --git a/test/unit/org/apache/cassandra/cql3/ViewSchemaTest.java b/test/unit/org/apache/cassandra/cql3/ViewSchemaTest.java index b58777e5a9..7f3e0b6910 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewSchemaTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewSchemaTest.java @@ -863,6 +863,6 @@ public class ViewSchemaTest extends ViewAbstractTest view, keyspace(), base, - mv.metadata().id))); + mv.metadata().id.toLongString()))); } } diff --git a/test/unit/org/apache/cassandra/cql3/statements/DescribeStatementTest.java b/test/unit/org/apache/cassandra/cql3/statements/DescribeStatementTest.java index 5d376dcf07..eb5015da61 100644 --- a/test/unit/org/apache/cassandra/cql3/statements/DescribeStatementTest.java +++ b/test/unit/org/apache/cassandra/cql3/statements/DescribeStatementTest.java @@ -508,7 +508,7 @@ public class DescribeStatementTest extends CQLTester " v2 int,\n" + " v3 int,\n" + " PRIMARY KEY ((pk1, pk2), c)\n" + - ") WITH ID = " + id + "\n" + + ") WITH ID = " + id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (c ASC)\n" + " AND " + tableParametersCql(); @@ -597,7 +597,7 @@ public class DescribeStatementTest extends CQLTester " v1 text,\n" + " v2 int,\n" + " v3 int\n" + - ") WITH ID = " + id + "\n" + + ") WITH ID = " + id.toLongString() + "\n" + " AND " + tableParametersCql(); assertRowsNet(executeDescribeNet("DESCRIBE TABLE " + KEYSPACE_PER_TEST + "." + table + " WITH INTERNALS"), @@ -873,7 +873,7 @@ public class DescribeStatementTest extends CQLTester " v1 int,\n" + " v2 int,\n" + " PRIMARY KEY ((pk1, pk2), ck1, ck2)\n" + - ") WITH ID = " + source.id + "\n" + + ") WITH ID = " + source.id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (ck1 ASC, ck2 ASC)\n" + " AND " + tableParametersCql(); String targetTableCreateStatement = "CREATE TABLE " + KEYSPACE_PER_TEST + "." + targetTable + " (\n" + @@ -885,7 +885,7 @@ public class DescribeStatementTest extends CQLTester " v1 int,\n" + " v2 int,\n" + " PRIMARY KEY ((pk1, pk2), ck1, ck2)\n" + - ") WITH ID = " + target.id + "\n" + + ") WITH ID = " + target.id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (ck1 ASC, ck2 ASC)\n" + " AND " + tableParametersCql(); @@ -932,7 +932,7 @@ public class DescribeStatementTest extends CQLTester " v1 text,\n" + " v2 text MASKED WITH system.mask_inner(1, null),\n" + " PRIMARY KEY ((pk1, pk2), ck1, ck2)\n" + - ") WITH ID = " + tableMetadata.id + "\n" + + ") WITH ID = " + tableMetadata.id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (ck1 ASC, ck2 ASC)\n" + " AND " + tableParametersCql(); diff --git a/test/unit/org/apache/cassandra/cql3/validation/operations/DropRecreateAndRestoreTest.java b/test/unit/org/apache/cassandra/cql3/validation/operations/DropRecreateAndRestoreTest.java index e9ab1cecbd..d85a494406 100644 --- a/test/unit/org/apache/cassandra/cql3/validation/operations/DropRecreateAndRestoreTest.java +++ b/test/unit/org/apache/cassandra/cql3/validation/operations/DropRecreateAndRestoreTest.java @@ -29,7 +29,6 @@ import org.apache.cassandra.exceptions.AlreadyExistsException; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.io.util.FileUtils; -import org.apache.cassandra.schema.TableId; public class DropRecreateAndRestoreTest extends CQLTester { @@ -45,7 +44,7 @@ public class DropRecreateAndRestoreTest extends CQLTester execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 0, 0, 0, timeInMicroSecond1); execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ?", 0, 1, 1, timeInMicroSecond1); - TableId id = currentTableMetadata().id; + String id = currentTableMetadata().id.toLongString(); assertRows(execute("SELECT * FROM %s"), row(0, 0, 0), row(0, 1, 1)); Thread.sleep(5); @@ -87,7 +86,7 @@ public class DropRecreateAndRestoreTest extends CQLTester public void testCreateWithIdDuplicate() throws Throwable { createTable("CREATE TABLE %s (a int, b int, c int, PRIMARY KEY(a, b))"); - TableId id = currentTableMetadata().id; + String id = currentTableMetadata().id.toLongString(); execute(String.format("CREATE TABLE %%s (a int, b int, c int, PRIMARY KEY(a, b)) WITH ID = %s", id)); } @@ -101,7 +100,7 @@ public class DropRecreateAndRestoreTest extends CQLTester public void testAlterWithId() throws Throwable { createTable("CREATE TABLE %s (a int, b int, c int, PRIMARY KEY(a, b))"); - TableId id = currentTableMetadata().id; + String id = currentTableMetadata().id.toLongString(); execute(String.format("ALTER TABLE %%s WITH ID = %s", id)); } } diff --git a/test/unit/org/apache/cassandra/db/SchemaCQLHelperTest.java b/test/unit/org/apache/cassandra/db/SchemaCQLHelperTest.java index da5f6d1853..9f3a0fab36 100644 --- a/test/unit/org/apache/cassandra/db/SchemaCQLHelperTest.java +++ b/test/unit/org/apache/cassandra/db/SchemaCQLHelperTest.java @@ -473,7 +473,7 @@ public class SchemaCQLHelperTest extends CQLTester " reg1 " + typeC+ ",\n" + " reg3 int,\n" + " PRIMARY KEY ((pk1, pk2), ck1, ck2)\n" + - ") WITH ID = " + cfs.metadata.id + "\n" + + ") WITH ID = " + cfs.metadata.id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (ck1 ASC, ck2 DESC)"; assertThat(schema, @@ -529,7 +529,7 @@ public class SchemaCQLHelperTest extends CQLTester " reg3 int,\n" + " reg2 int,\n" + " PRIMARY KEY ((pk1, pk2), ck1, ck2)\n" + - ") WITH ID = " + cfs.metadata.id + "\n" + + ") WITH ID = " + cfs.metadata.id.toLongString() + "\n" + " AND CLUSTERING ORDER BY (ck1 ASC, ck2 DESC)"; assertThat(schema, @@ -568,7 +568,7 @@ public class SchemaCQLHelperTest extends CQLTester " reg1 int,\n" + " reg3 int,\n" + " reg2 int\n" + - ") WITH ID = " + cfs.metadata.id + "\n"; + ") WITH ID = " + cfs.metadata.id.toLongString() + "\n"; assertThat(schema, allOf(startsWith(expected), diff --git a/test/unit/org/apache/cassandra/index/accord/RouteIndexTest.java b/test/unit/org/apache/cassandra/index/accord/RouteIndexTest.java index 3ac830611e..b5d8480728 100644 --- a/test/unit/org/apache/cassandra/index/accord/RouteIndexTest.java +++ b/test/unit/org/apache/cassandra/index/accord/RouteIndexTest.java @@ -121,7 +121,6 @@ public class RouteIndexTest extends CQLTester private static final Gen.IntGen NUM_STORES_GEN = Gens.ints().between(1, MAX_STORES); private static final Gen TOKEN_DISTRIBUTION = Gens.mixedDistribution(MIN_TOKEN, MAX_TOKEN + 1); private static final Gen RANGE_SIZE_DISTRIBUTION = Gens.mixedDistribution(10, (int) (TOKEN_RANGE_SIZE * .01)); - private static final Gen> DOMAIN_DISTRIBUTION = Gens.mixedDistribution(Domain.values()); @BeforeClass public static void setUpClass() @@ -475,7 +474,7 @@ public class RouteIndexTest extends CQLTester tables = Collections.singletonList(tableId); tokenGen = TOKEN_DISTRIBUTION.next(rs); rangeGen = rangeGen(rs, tables); - domainGen = DOMAIN_DISTRIBUTION.next(rs); + domainGen = ignore -> Domain.Range; // we shouldn't be saving/searching key transactions against ranges journalTable = Keyspace.open(ACCORD_KEYSPACE_NAME).getColumnFamilyStore(AccordKeyspace.JOURNAL); for (int i = 0 ; i < numStores ; ++i) diff --git a/test/unit/org/apache/cassandra/io/Serializers.java b/test/unit/org/apache/cassandra/io/Serializers.java index 75d2fd1c46..e9e7f04c7b 100644 --- a/test/unit/org/apache/cassandra/io/Serializers.java +++ b/test/unit/org/apache/cassandra/io/Serializers.java @@ -44,8 +44,25 @@ public class Serializers Assertions.assertThat(read).describedAs("The deserialized output does not match the serialized input; difference %s", new LazyToString(() -> ReflectionUtils.recursiveEquals(read, input).toString())).isEqualTo(input); Assertions.assertThat(buffer.remaining()).describedAs("deserialize did not consume all the serialized input").isEqualTo(0); buffer.flip(); + buffer.mark(); serializer.skip(in); Assertions.assertThat(buffer.remaining()).describedAs("skip did not consume all the serialized input").isEqualTo(0); + boolean testByteBufferMethods; + try + { + testByteBufferMethods = serializer.getClass().getMethod("serialize", Object.class).getDeclaringClass() != AsymmetricUnversionedSerializer.class + || serializer.getClass().getMethod("deserialize", ByteBuffer.class).getDeclaringClass() != AsymmetricUnversionedSerializer.class; + } + catch (NoSuchMethodException e) + { + throw new AssertionError(e); + } + if (testByteBufferMethods) + { + ByteBuffer serialized2 = serializer.serialize(input); + T read2 = serializer.deserialize(serialized2); + Assertions.assertThat(read2).describedAs("The deserialized output does not match the serialized input; difference %s", new LazyToString(() -> ReflectionUtils.recursiveEquals(read2, input).toString())).isEqualTo(input); + } } public static void testSerde(DataOutputBuffer output, ParameterisedUnversionedSerializer serializer, T input, P p) throws IOException diff --git a/test/unit/org/apache/cassandra/service/accord/serializers/LatestDepsSerializerTest.java b/test/unit/org/apache/cassandra/service/accord/serializers/LatestDepsSerializerTest.java index b0e090ff2a..0efc44023d 100644 --- a/test/unit/org/apache/cassandra/service/accord/serializers/LatestDepsSerializerTest.java +++ b/test/unit/org/apache/cassandra/service/accord/serializers/LatestDepsSerializerTest.java @@ -18,18 +18,96 @@ package org.apache.cassandra.service.accord.serializers; +import java.io.IOException; +import java.util.Arrays; + import org.junit.Test; +import accord.api.RoutingKey; +import accord.primitives.Ballot; +import accord.primitives.Deps; +import accord.primitives.Known; import accord.primitives.LatestDeps; +import accord.primitives.Txn; +import accord.utils.AccordGens; +import accord.utils.Gen; +import accord.utils.Gens; +import accord.utils.RandomSource; +import accord.utils.RandomTestRunner; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Token; import org.apache.cassandra.io.Serializers; import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.schema.TableId; +import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.utils.AccordGenerators; +import org.apache.cassandra.utils.CassandraGenerators; + +import static accord.primitives.Routable.Domain.Key; +import static accord.primitives.Routable.Domain.Range; +import static org.apache.cassandra.utils.AccordGenerators.fromQT; public class LatestDepsSerializerTest { + { + DatabaseDescriptor.toolInitialization(); + } + @Test public void emptySerializerTest() throws Throwable { DataOutputBuffer buf = new DataOutputBuffer(); Serializers.testSerde(buf, LatestDepsSerializers.latestDeps, LatestDeps.EMPTY); } + + @Test + public void testN() + { + for (int i = 0 ; i < 10000 ; ++i) + { + Gen partitioners = AccordGenerators.partitioner(); + RandomTestRunner.test().check(rs -> { + try + { + testOne(partitioners.next(rs), rs); + } + catch (IOException e) + { + throw new AssertionError(e); + } + }); + } + } + + private void testOne(IPartitioner partitioner, RandomSource rs) throws IOException + { + DatabaseDescriptor.setPartitionerUnsafe(partitioner); + TableId tableId = fromQT(CassandraGenerators.TABLE_ID_GEN).next(rs); + Gen tokens = fromQT(CassandraGenerators.token(partitioner)); + Gen routingKeys = AccordGenerators.routingKeyGen(ignore -> tableId, tokens, partitioner); + Gen deps = AccordGens.deps(AccordGens.keyDeps(routingKeys, AccordGens.txnIds(Gens.pick(Txn.Kind.values()), ignore -> Key)), + AccordGens.rangeDeps(AccordGenerators.range(partitioner, ignore -> tableId), AccordGens.txnIds(Gens.pick(Txn.Kind.values()), ignore -> Range))); + Gen knownDeps = Gens.pick(Known.KnownDeps.values()); + Gen ballots = AccordGens.ballot(); + int size = 1 + rs.nextInt(7); + RoutingKey[] starts = new RoutingKey[size + 1]; + LatestDeps.LatestEntry[] entries = new LatestDeps.LatestEntry[size]; + for (int i = 0 ; i <= size ; ++i) + starts[i] = routingKeys.next(rs); + Arrays.sort(starts); + for (int i = 0 ; i < size ; ++i) + { + if (rs.nextBoolean()) continue; + entries[i] = new LatestDeps.LatestEntry(knownDeps.next(rs), + rs.nextBoolean() ? rs.nextBoolean() ? Ballot.ZERO : Ballot.MAX : ballots.next(rs), + rs.nextBoolean() ? null : deps.next(rs), + rs.nextBoolean() ? null : deps.next(rs)); + } + LatestDeps latestDeps = LatestDeps.SerializerSupport.create(true, starts, entries); + DataOutputBuffer buf = new DataOutputBuffer(); + Serializers.testSerde(buf, LatestDepsSerializers.latestDeps, latestDeps); + } + + } \ No newline at end of file diff --git a/test/unit/org/apache/cassandra/utils/AccordGenerators.java b/test/unit/org/apache/cassandra/utils/AccordGenerators.java index 1f3d4fb4ac..bd994fc27a 100644 --- a/test/unit/org/apache/cassandra/utils/AccordGenerators.java +++ b/test/unit/org/apache/cassandra/utils/AccordGenerators.java @@ -103,8 +103,8 @@ import static org.apache.cassandra.service.accord.AccordTestUtils.createPartialT public class AccordGenerators { - private static final Gen PARTITIONER_GEN = fromQT(CassandraGenerators.nonLocalPartitioners()); - private static final Gen TABLE_ID_GEN = fromQT(CassandraGenerators.TABLE_ID_GEN); + public static final Gen PARTITIONER_GEN = fromQT(CassandraGenerators.nonLocalPartitioners()); + public static final Gen TABLE_ID_GEN = fromQT(CassandraGenerators.TABLE_ID_GEN); private AccordGenerators() {