From 68aea4b15ffbc513f0dc6f99a37bed9a69e93a34 Mon Sep 17 00:00:00 2001 From: Alex Petrov Date: Tue, 8 Apr 2025 14:15:17 +0200 Subject: [PATCH] Fix simulator after rebase * avoid running python3-dependent tasks when running simulator tasks * fix a problem with simulated snitch rebase * add a distinction between short-lived daemon threads and infinite loop ones for cases when we need to simulate user-implemented infinite loops Patch by Alex Petrov; reviewed by Benedict Elliott Smith for CASSANDRA-20542 --- .../cassandra/concurrent/ExecutorFactory.java | 36 +++++++---- .../concurrent/InfiniteLoopExecutor.java | 22 +++---- .../AbstractCommitLogSegmentManager.java | 2 +- .../commitlog/AbstractCommitLogService.java | 2 +- .../org/apache/cassandra/journal/Flusher.java | 2 +- .../org/apache/cassandra/journal/Journal.java | 2 +- .../service/accord/AccordExecutorLoops.java | 5 +- .../apache/cassandra/tcm/log/LocalLog.java | 2 +- .../simulator/ClusterSimulation.java | 3 +- .../paxos/AccordSimulationRunner.java | 7 +++ .../systems/InterceptingExecutorFactory.java | 19 ++++-- .../simulator/systems/SimulatedSnitch.java | 59 +++++++------------ .../test/AccordHarrySimulationTest.java | 6 ++ .../simulator/test/HarrySimulatorTest.java | 17 ++++-- .../simulator/test/HarryValidatingQuery.java | 43 +++++++++++--- .../concurrent/ForwardingExecutorFactory.java | 8 +-- .../concurrent/InfiniteLoopExecutorTest.java | 2 +- .../concurrent/SimulatedExecutorFactory.java | 4 +- .../apache/cassandra/repair/FuzzTestBase.java | 8 +-- 19 files changed, 149 insertions(+), 100 deletions(-) diff --git a/src/java/org/apache/cassandra/concurrent/ExecutorFactory.java b/src/java/org/apache/cassandra/concurrent/ExecutorFactory.java index ec3b9c370b..0b72961e13 100644 --- a/src/java/org/apache/cassandra/concurrent/ExecutorFactory.java +++ b/src/java/org/apache/cassandra/concurrent/ExecutorFactory.java @@ -18,15 +18,15 @@ package org.apache.cassandra.concurrent; -import org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon; import org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts; import org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe; import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.Shared; -import static java.lang.Thread.*; +import static java.lang.Thread.NORM_PRIORITY; +import static java.lang.Thread.UncaughtExceptionHandler; import static org.apache.cassandra.concurrent.ExecutorFactory.SimulatorSemantics.NORMAL; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.UNSYNCHRONIZED; import static org.apache.cassandra.concurrent.NamedThreadFactory.createThread; import static org.apache.cassandra.concurrent.NamedThreadFactory.setupThread; @@ -78,6 +78,17 @@ public interface ExecutorFactory extends ExecutorBuilderFactory.Jmxable interruptHandler; private final Condition isTerminated = newOneTimeCondition(); - public InfiniteLoopExecutor(String name, Task task, Daemon daemon) + public InfiniteLoopExecutor(String name, Task task, SystemThreadTag systemTag) { - this(ExecutorFactory.Global.executorFactory(), name, task, daemon, UNSYNCHRONIZED); + this(ExecutorFactory.Global.executorFactory(), name, task, systemTag, UNSYNCHRONIZED); } - public InfiniteLoopExecutor(ExecutorFactory factory, String name, Task task, Daemon daemon) + public InfiniteLoopExecutor(ExecutorFactory factory, String name, Task task, SystemThreadTag systemTag) { - this(factory, name, task, daemon, UNSYNCHRONIZED); + this(factory, name, task, systemTag, UNSYNCHRONIZED); } - public InfiniteLoopExecutor(ExecutorFactory factory, String name, Task task, Daemon daemon, Interrupts interrupts) + public InfiniteLoopExecutor(ExecutorFactory factory, String name, Task task, SystemThreadTag systemTag, Interrupts interrupts) { this.task = task; - this.thread = factory.startThread(name, this::loop, daemon); + this.thread = factory.startThread(name, this::loop, systemTag, SimulatorThreadTag.INFINITE_LOOP); this.interruptHandler = interrupts == SYNCHRONIZED ? interruptHandler(task) : Thread::interrupt; diff --git a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java index dcd791caf3..3348968881 100644 --- a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java +++ b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java @@ -56,7 +56,7 @@ import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import org.apache.cassandra.utils.concurrent.WaitQueue; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.SYNCHRONIZED; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE; import static org.apache.cassandra.db.commitlog.CommitLogSegment.Allocation; diff --git a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogService.java b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogService.java index cd3eb56105..7bba9c4911 100644 --- a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogService.java +++ b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogService.java @@ -38,7 +38,7 @@ import static java.util.concurrent.TimeUnit.MINUTES; import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.SYNCHRONIZED; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE; import static org.apache.cassandra.concurrent.Interruptible.State.NORMAL; diff --git a/src/java/org/apache/cassandra/journal/Flusher.java b/src/java/org/apache/cassandra/journal/Flusher.java index cd89709497..d48d80171b 100644 --- a/src/java/org/apache/cassandra/journal/Flusher.java +++ b/src/java/org/apache/cassandra/journal/Flusher.java @@ -39,7 +39,7 @@ import static java.lang.String.format; import static java.util.concurrent.TimeUnit.MINUTES; import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.SYNCHRONIZED; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE; import static org.apache.cassandra.concurrent.Interruptible.State.NORMAL; diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java index 98150f1cee..e70e93b9c6 100644 --- a/src/java/org/apache/cassandra/journal/Journal.java +++ b/src/java/org/apache/cassandra/journal/Journal.java @@ -64,7 +64,7 @@ import org.jctools.queues.MpscUnboundedArrayQueue; import static java.lang.String.format; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.SYNCHRONIZED; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE; import static org.apache.cassandra.concurrent.Interruptible.State.NORMAL; diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java index 0e8aab8c8a..d26e6e5abe 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutorLoops.java @@ -24,11 +24,14 @@ import java.util.function.Function; import java.util.function.IntFunction; import accord.utils.Invariants; + import io.netty.util.collection.LongObjectHashMap; import org.apache.cassandra.service.accord.AccordExecutor.Mode; import org.apache.cassandra.utils.concurrent.Condition; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; +import static org.apache.cassandra.concurrent.ExecutorFactory.SimulatorThreadTag.INFINITE_LOOP; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_LOCK; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -46,7 +49,7 @@ class AccordExecutorLoops loops = new LongObjectHashMap<>(threads); for (int i = 0; i < threads; ++i) { - Thread thread = executorFactory().startThread(name.apply(i), wrap(loopFactory.apply(mode))); + Thread thread = executorFactory().startThread(name.apply(i), wrap(loopFactory.apply(mode)), NON_DAEMON, INFINITE_LOOP); Thread conflict = loops.putIfAbsent(thread.getId(), thread); Invariants.require(conflict == null || !conflict.isAlive(), "Allocated two threads with the same threadId!"); } diff --git a/src/java/org/apache/cassandra/tcm/log/LocalLog.java b/src/java/org/apache/cassandra/tcm/log/LocalLog.java index 01b6c22016..a0f50d16ad 100644 --- a/src/java/org/apache/cassandra/tcm/log/LocalLog.java +++ b/src/java/org/apache/cassandra/tcm/log/LocalLog.java @@ -70,7 +70,7 @@ import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.concurrent.Condition; import org.apache.cassandra.utils.concurrent.WaitQueue; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.UNSYNCHRONIZED; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE; import static org.apache.cassandra.tcm.Epoch.EMPTY; diff --git a/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java b/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java index 22ce8581dd..ed21d7079d 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java +++ b/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java @@ -784,7 +784,8 @@ public class ClusterSimulation implements AutoCloseable .set("failure_detector", SimulatedFailureDetector.Instance.class.getName()) .set("commitlog_compression", new ParameterizedClass(LZ4Compressor.class.getName(), emptyMap())) .set("commitlog_sync", "batch") - .set("accord.journal.flush_mode", "BATCH"); + .set("accord.journal.flush_mode", "BATCH") + .set("accord.command_store_shard_count", "4"); // TODO: Add remove() to IInstanceConfig if (config instanceof InstanceConfig) { diff --git a/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java b/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java index b0cd80cf30..e782848ea4 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java +++ b/test/simulator/main/org/apache/cassandra/simulator/paxos/AccordSimulationRunner.java @@ -23,14 +23,20 @@ import java.util.concurrent.atomic.AtomicInteger; import org.junit.BeforeClass; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import io.airlift.airline.Cli; import io.airlift.airline.Command; import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.simulator.SimulationRunner; +import org.apache.cassandra.simulator.SimulatorUtils; import org.apache.cassandra.utils.StorageCompatibilityMode; public class AccordSimulationRunner extends SimulationRunner { + private static Logger logger = LoggerFactory.getLogger(AccordSimulationRunner.class); + @BeforeClass public static void beforeAll() { @@ -86,6 +92,7 @@ public class AccordSimulationRunner extends SimulationRunner */ public static void main(String[] args) throws IOException { + SimulatorUtils.verifyAndlogSimulatorArgs(logger, args); AccordClusterSimulation.Builder builder = new AccordClusterSimulation.Builder(); builder.unique(uniqueNum.getAndIncrement()); diff --git a/test/simulator/main/org/apache/cassandra/simulator/systems/InterceptingExecutorFactory.java b/test/simulator/main/org/apache/cassandra/simulator/systems/InterceptingExecutorFactory.java index c7f4dce8b8..aa22e3ef86 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/systems/InterceptingExecutorFactory.java +++ b/test/simulator/main/org/apache/cassandra/simulator/systems/InterceptingExecutorFactory.java @@ -30,13 +30,13 @@ import java.util.function.Consumer; import com.google.common.annotations.VisibleForTesting; +import accord.utils.UnhandledEnum; import io.netty.util.concurrent.FastThreadLocal; import org.apache.cassandra.concurrent.ExecutorBuilder; import org.apache.cassandra.concurrent.ExecutorBuilderFactory; import org.apache.cassandra.concurrent.ExecutorFactory; import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.InfiniteLoopExecutor; -import org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon; import org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts; import org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe; import org.apache.cassandra.concurrent.Interruptible.Task; @@ -69,6 +69,8 @@ import org.apache.cassandra.utils.WithResources; import org.apache.cassandra.utils.concurrent.RunnableFuture; import static org.apache.cassandra.simulator.systems.SimulatedAction.Kind.INFINITE_LOOP; +import static org.apache.cassandra.simulator.systems.SimulatedAction.Kind.SCHEDULED_DAEMON; +import static org.apache.cassandra.simulator.systems.SimulatedAction.Kind.THREAD; public class InterceptingExecutorFactory implements ExecutorFactory, Closeable { @@ -327,9 +329,18 @@ public class InterceptingExecutorFactory implements ExecutorFactory, Closeable return configurePooled(name, threads).build(); } - public Thread startThread(String name, Runnable runnable, Daemon daemon) + public Thread startThread(String name, Runnable runnable, SystemThreadTag systemTag, SimulatorThreadTag simulatorTag) { - return simulatedExecution.intercept().start(SimulatedAction.Kind.THREAD, factory(name)::newThread, runnable); + SimulatedAction.Kind kind; + switch (simulatorTag) + { + default: throw UnhandledEnum.unknown(simulatorTag); + case INFINITE_LOOP: kind = INFINITE_LOOP; break; + case JOB: kind = THREAD; break; + case DAEMON: kind = SCHEDULED_DAEMON; break; + } + + return simulatedExecution.intercept().start(kind, factory(name)::newThread, runnable); } @VisibleForTesting @@ -341,7 +352,7 @@ public class InterceptingExecutorFactory implements ExecutorFactory, Closeable } @Override - public Interruptible infiniteLoop(String name, Task task, SimulatorSafe simulatorSafe, Daemon daemon, Interrupts interrupts) + public Interruptible infiniteLoop(String name, Task task, SimulatorSafe simulatorSafe, SystemThreadTag systemTag, Interrupts interrupts) { if (simulatorSafe != SimulatorSafe.SAFE) { diff --git a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedSnitch.java b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedSnitch.java index 7692e0ae5c..5549556298 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedSnitch.java +++ b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedSnitch.java @@ -28,49 +28,18 @@ import java.util.stream.IntStream; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.IInstanceConfig; -import org.apache.cassandra.locator.*; +import org.apache.cassandra.locator.Endpoint; +import org.apache.cassandra.locator.IEndpointSnitch; +import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.locator.Replica; +import org.apache.cassandra.locator.ReplicaCollection; import org.apache.cassandra.simulator.cluster.NodeLookup; import org.apache.cassandra.utils.Sortable; public class SimulatedSnitch extends NodeLookup { - private static class SimulatedProximity implements NodeProximity - { - @Override - public > C sortedByProximity(InetAddressAndPort address, C addresses) - { - return addresses.sorted(Comparator.comparingInt(SimulatedSnitch::asInt)); - } - - @Override - public int compareEndpoints(InetAddressAndPort target, Replica r1, Replica r2) - { - return Comparator.comparingInt(SimulatedSnitch::asInt).compare(r1, r2); - } - - @Override - public boolean isWorthMergingForRangeQuery(ReplicaCollection merged, ReplicaCollection l1, ReplicaCollection l2) - { - return false; - } - - @Override - public boolean supportCompareByEndpoint() - { - return true; - } - - @Override - public > Comparator endpointComparator(InetAddressAndPort address, C addresses) - { - return Comparator.comparingInt(SimulatedSnitch::asInt); - } - } - public static class Instance implements IEndpointSnitch { - private final NodeProximity proximity = new SimulatedProximity(); - private static volatile Function LOOKUP_DC; public String getRack(InetAddressAndPort endpoint) @@ -85,12 +54,12 @@ public class SimulatedSnitch extends NodeLookup public > C sortedByProximity(InetAddressAndPort address, C addresses) { - return proximity.sortedByProximity(address, addresses); + return addresses.sorted(Comparator.comparingInt(SimulatedSnitch::asInt)); } public int compareEndpoints(InetAddressAndPort target, Replica r1, Replica r2) { - return proximity.compareEndpoints(target, r1, r2); + return Comparator.comparingInt(SimulatedSnitch::asInt).compare(r1, r2); } public void gossiperStarting() @@ -99,13 +68,25 @@ public class SimulatedSnitch extends NodeLookup public boolean isWorthMergingForRangeQuery(ReplicaCollection merged, ReplicaCollection l1, ReplicaCollection l2) { - return proximity.isWorthMergingForRangeQuery(merged, l1, l2); + return false; } public static void setup(Function lookupDc) { LOOKUP_DC = lookupDc; } + + @Override + public boolean supportCompareByEndpoint() + { + return true; + } + + @Override + public > Comparator endpointComparator(InetAddressAndPort address, C addresses) + { + return Comparator.comparingInt(SimulatedSnitch::asInt); + } } final int[] numInDcs; diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java index 8506fb6157..c4dc9f10fc 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java @@ -24,6 +24,7 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.harry.SchemaSpec; import org.apache.cassandra.harry.gen.Generator; import org.apache.cassandra.harry.gen.SchemaGenerators; @@ -62,6 +63,11 @@ public class AccordHarrySimulationTest extends HarrySimulatorTest return schedulers; } + protected ConsistencyLevel validateQueryConsistency() + { + return ConsistencyLevel.QUORUM; + } + public Generator schemaSpecGen(String keyspace, String prefix) { return SchemaGenerators.schemaSpecGen(keyspace, prefix, 1000, SchemaSpec.optionsBuilder().withTransactionalMode(TransactionalMode.full)); diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java index 14073943d7..9bc9b6e759 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java @@ -288,7 +288,7 @@ public class HarrySimulatorTest work.add(interleave("Start generating", HarrySimulatorTest.generateWrites(rowsPerPhase, simulation, cl))); work.add(work("Validate all data locally", - lazy(() -> validateAllLocal(simulation, simulation.nodeState.ring, rf)))); + lazy(() -> validateAllLocal(simulation, simulation.nodeState.ring, rf, validateQueryConsistency(), simulation.rng)))); return arr(work.toArray(new ActionSchedule.Work[0])); }, @@ -339,8 +339,8 @@ public class HarrySimulatorTest run(() -> simulation.nodeState.decommission(node)))); work.add(work("Check node state", assertNodeState(simulation.simulated, simulation.cluster, node, NodeState.LEFT))); } - work.add(work("Validate data locally", - lazy(() -> validateAllLocal(simulation, simulation.nodeState.ring, rf)))); + work.add(work("Validate data with " + validateQueryConsistency(), + lazy(() -> validateAllLocal(simulation, simulation.nodeState.ring, rf, validateQueryConsistency(), simulation.rng)))); boolean tmp = shouldBootstrap; work.add(work("Output message", run(() -> logger.warn("Finished {} of {} and data validation!\n", tmp ? "bootstrap" : "decommission", node)))); @@ -846,7 +846,7 @@ public class HarrySimulatorTest * Given you have used `generate` methods to generate data with Harry, you can use this method to check whether all * data has been propagated everywhere it should be, be it via streaming, read repairs, or regular writes. */ - public static Action validateAllLocal(HarrySimulation simulation, List owernship, TokenPlacementModel.ReplicationFactor rf) + public static Action validateAllLocal(HarrySimulation simulation, List owernship, TokenPlacementModel.ReplicationFactor rf, ConsistencyLevel consistencyLevel, EntropySource rng) { return new Actions.LambdaAction("Validate", Action.Modifiers.RELIABLE_NO_TIMEOUTS, () -> { @@ -861,12 +861,19 @@ public class HarrySimulatorTest Operations.SelectPartition select = new Operations.SelectPartition(Long.MAX_VALUE, pd); actions.add(new HarryValidatingQuery(simulation, simulation.cluster, rf, owernship, new Visit(Long.MAX_VALUE, new Operations.Operation[]{ select }), - simulation.queryBuilder)); + simulation.queryBuilder, + consistencyLevel, + rng)); } return ActionList.of(actions).setStrictlySequential(); }); } + protected ConsistencyLevel validateQueryConsistency() + { + return ConsistencyLevel.NODE_LOCAL; + } + private static Set visitedPds(HarrySimulation simulation) { Set pds = new HashSet<>(); diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/HarryValidatingQuery.java b/test/simulator/test/org/apache/cassandra/simulator/test/HarryValidatingQuery.java index 72cd85de3b..3e34bf69b1 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/HarryValidatingQuery.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/HarryValidatingQuery.java @@ -25,7 +25,9 @@ import org.slf4j.LoggerFactory; import accord.utils.Invariants; import org.apache.cassandra.distributed.Cluster; +import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.distributed.api.IInstance; +import org.apache.cassandra.harry.gen.EntropySource; import org.apache.cassandra.harry.op.Visit; import org.apache.cassandra.harry.op.Operations; import org.apache.cassandra.harry.execution.CompiledStatement; @@ -39,6 +41,8 @@ import org.apache.cassandra.simulator.systems.InterceptedExecution; import org.apache.cassandra.simulator.systems.InterceptingExecutor; import org.apache.cassandra.simulator.systems.SimulatedAction; +import static org.apache.cassandra.simulator.SimulatorUtils.failWithOOM; + public class HarryValidatingQuery extends SimulatedAction { private static final Logger logger = LoggerFactory.getLogger(HarryValidatingQuery.class); @@ -51,13 +55,17 @@ public class HarryValidatingQuery extends SimulatedAction private final HarrySimulatorTest.HarrySimulation simulation; private final Visit visit; private final QueryBuildingVisitExecutor queryBuilder; + private final ConsistencyLevel consistencyLevel; + private final EntropySource rng; public HarryValidatingQuery(HarrySimulatorTest.HarrySimulation simulation, Cluster cluster, TokenPlacementModel.ReplicationFactor rf, List owernship, Visit visit, - QueryBuildingVisitExecutor queryBuilder) + QueryBuildingVisitExecutor queryBuilder, + ConsistencyLevel consistencyLevel, + EntropySource rng) { super(visit, Modifiers.RELIABLE_NO_TIMEOUTS, Modifiers.RELIABLE_NO_TIMEOUTS, null, simulation.simulated); this.rf = rf; @@ -67,7 +75,8 @@ public class HarryValidatingQuery extends SimulatedAction this.visit = visit; this.queryBuilder = queryBuilder; this.simulation = simulation; - + this.consistencyLevel = consistencyLevel; + this.rng = rng; } protected InterceptedExecution task() @@ -78,14 +87,25 @@ public class HarryValidatingQuery extends SimulatedAction { try { - TokenPlacementModel.ReplicatedRanges ring = rf.replicate(owernship); - Invariants.require(visit.operations.length == 1); - Invariants.require(visit.operations[0] instanceof Operations.SelectStatement); - Operations.SelectStatement select = (Operations.SelectStatement) visit.operations[0]; - for (TokenPlacementModel.Replica replica : ring.replicasFor(token(select.pd))) + if (consistencyLevel == ConsistencyLevel.NODE_LOCAL) { + TokenPlacementModel.ReplicatedRanges ring = rf.replicate(owernship); + Invariants.require(visit.operations.length == 1); + Invariants.require(visit.operations[0] instanceof Operations.SelectStatement); + Operations.SelectStatement select = (Operations.SelectStatement) visit.operations[0]; + for (TokenPlacementModel.Replica replica : ring.replicasFor(token(select.pd))) + { + CompiledStatement compiled = queryBuilder.compile(visit); + Object[][] objects = executeNodeLocal(compiled.cql(), replica.node(), compiled.bindings()); + List actualRows = InJvmDTestVisitExecutor.rowsToResultSet(simulation.schema, select, objects); + simulation.model.validate(select, actualRows); + } + } + else + { + Operations.SelectStatement select = (Operations.SelectStatement) visit.operations[0]; CompiledStatement compiled = queryBuilder.compile(visit); - Object[][] objects = executeNodeLocal(compiled.cql(), replica.node(), compiled.bindings()); + Object[][] objects = execute(compiled.cql(), rng.nextInt(cluster.size()) + 1, compiled.bindings()); List actualRows = InJvmDTestVisitExecutor.rowsToResultSet(simulation.schema, select, objects); simulation.model.validate(select, actualRows); } @@ -93,6 +113,7 @@ public class HarryValidatingQuery extends SimulatedAction catch (Throwable t) { logger.error("Caught an exception while validating", t); + failWithOOM(); throw t; } } @@ -113,4 +134,10 @@ public class HarryValidatingQuery extends SimulatedAction .get(); return instance.executeInternal(statement, bindings); } + + protected Object[][] execute(String statement, int id, Object... bindings) + { + IInstance instance = cluster.get(id); + return instance.coordinator().execute(statement, consistencyLevel, bindings); + } } diff --git a/test/unit/org/apache/cassandra/concurrent/ForwardingExecutorFactory.java b/test/unit/org/apache/cassandra/concurrent/ForwardingExecutorFactory.java index b6bae2f1a4..37b77f3dcf 100644 --- a/test/unit/org/apache/cassandra/concurrent/ForwardingExecutorFactory.java +++ b/test/unit/org/apache/cassandra/concurrent/ForwardingExecutorFactory.java @@ -117,15 +117,15 @@ public class ForwardingExecutorFactory implements ExecutorFactory } @Override - public Thread startThread(String name, Runnable runnable, InfiniteLoopExecutor.Daemon daemon) + public Thread startThread(String name, Runnable runnable, SystemThreadTag systemTag, SimulatorThreadTag simulatorTag) { - return delegate().startThread(name, runnable, daemon); + return delegate().startThread(name, runnable, systemTag, simulatorTag); } @Override - public Interruptible infiniteLoop(String name, Interruptible.Task task, InfiniteLoopExecutor.SimulatorSafe simulatorSafe, InfiniteLoopExecutor.Daemon daemon, InfiniteLoopExecutor.Interrupts interrupts) + public Interruptible infiniteLoop(String name, Interruptible.Task task, InfiniteLoopExecutor.SimulatorSafe simulatorSafe, SystemThreadTag systemTag, InfiniteLoopExecutor.Interrupts interrupts) { - return delegate().infiniteLoop(name, task, simulatorSafe, daemon, interrupts); + return delegate().infiniteLoop(name, task, simulatorSafe, systemTag, interrupts); } @Override diff --git a/test/unit/org/apache/cassandra/concurrent/InfiniteLoopExecutorTest.java b/test/unit/org/apache/cassandra/concurrent/InfiniteLoopExecutorTest.java index 9ec702dd50..29dc7ec8ac 100644 --- a/test/unit/org/apache/cassandra/concurrent/InfiniteLoopExecutorTest.java +++ b/test/unit/org/apache/cassandra/concurrent/InfiniteLoopExecutorTest.java @@ -30,7 +30,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.junit.Assert; import org.junit.Test; -import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.DAEMON; +import static org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.DAEMON; public class InfiniteLoopExecutorTest { diff --git a/test/unit/org/apache/cassandra/concurrent/SimulatedExecutorFactory.java b/test/unit/org/apache/cassandra/concurrent/SimulatedExecutorFactory.java index c7a03dc79c..e4ae9282e7 100644 --- a/test/unit/org/apache/cassandra/concurrent/SimulatedExecutorFactory.java +++ b/test/unit/org/apache/cassandra/concurrent/SimulatedExecutorFactory.java @@ -204,7 +204,7 @@ public class SimulatedExecutorFactory implements ExecutorFactory, Clock } @Override - public Thread startThread(String name, Runnable runnable, InfiniteLoopExecutor.Daemon daemon) + public Thread startThread(String name, Runnable runnable, SystemThreadTag systemTag, SimulatorThreadTag simulatorTag) { throw new UnsupportedOperationException("Thread can't be simualted"); } @@ -213,7 +213,7 @@ public class SimulatedExecutorFactory implements ExecutorFactory, Clock public Interruptible infiniteLoop(String name, Interruptible.Task task, InfiniteLoopExecutor.SimulatorSafe simulatorSafe, - InfiniteLoopExecutor.Daemon daemon, + SystemThreadTag systemTag, InfiniteLoopExecutor.Interrupts interrupts) { var delegate = new UnorderedScheduledExecutorService(); diff --git a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java index ad814f0032..8ffd96795b 100644 --- a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java +++ b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java @@ -243,16 +243,16 @@ public abstract class FuzzTestBase extends CQLTester.InMemory } @Override - public Thread startThread(String name, Runnable runnable, InfiniteLoopExecutor.Daemon daemon) + public Thread startThread(String name, Runnable runnable, SystemThreadTag systemTag, SimulatorThreadTag simulatorTag) { if (shouldMock()) return new Thread(); - return delegate.startThread(name, runnable, daemon); + return delegate.startThread(name, runnable, systemTag, simulatorTag); } @Override - public Interruptible infiniteLoop(String name, Interruptible.Task task, InfiniteLoopExecutor.SimulatorSafe simulatorSafe, InfiniteLoopExecutor.Daemon daemon, InfiniteLoopExecutor.Interrupts interrupts) + public Interruptible infiniteLoop(String name, Interruptible.Task task, InfiniteLoopExecutor.SimulatorSafe simulatorSafe, SystemThreadTag systemTag, InfiniteLoopExecutor.Interrupts interrupts) { - return delegate.infiniteLoop(name, task, simulatorSafe, daemon, interrupts); + return delegate.infiniteLoop(name, task, simulatorSafe, systemTag, interrupts); } @Override