From 6470020dc2309bb706b26eb73f808e4434abe42e Mon Sep 17 00:00:00 2001 From: David Capwell Date: Thu, 7 Aug 2025 11:21:09 -0700 Subject: [PATCH] Fix org.apache.cassandra.simulator.test.EpochStressTest so it can run again patch by David Capwell; reviewed by Alex Petrov for CASSANDRA-20812 --- .../simulator/ClusterSimulation.java | 22 +++++- .../cassandra/simulator/systems/Failures.java | 7 ++ .../simulator/systems/SimulatedQuery.java | 25 +++++-- .../simulator/test/EpochStressTest.java | 34 +++++++-- .../simulator/test/SimulationTestBase.java | 69 +++++++++++++++---- 5 files changed, 128 insertions(+), 29 deletions(-) diff --git a/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java b/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java index 4db15c6c30..fb9e3f41ae 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java +++ b/test/simulator/main/org/apache/cassandra/simulator/ClusterSimulation.java @@ -709,7 +709,8 @@ public class ClusterSimulation implements AutoCloseable public final Cluster cluster; private final ListenableFileSystem fs; protected final Map> onUnexpectedShutdown = new TreeMap<>(); - protected final List> onShutdown = new CopyOnWriteArrayList<>(); + protected final List> onPreShutdown = new CopyOnWriteArrayList<>(); + protected final List> onPostShutdown = new CopyOnWriteArrayList<>(); protected final ThreadLocalRandomCheck threadLocalRandomCheck; private final RunnableActionScheduler scheduler; @@ -895,7 +896,7 @@ public class ClusterSimulation implements AutoCloseable } }).withClassTransformer(interceptClasses) .withShutdownExecutor((name, classLoader, shuttingDown, call) -> { - onShutdown.add(call); + onPostShutdown.add(call); return null; }) ).createWithoutStarting(); @@ -919,6 +920,9 @@ public class ClusterSimulation implements AutoCloseable Choices.random(random, builder.consensusChanges), minRf, initialRf, maxRf, null); this.factory = factory; + + // during cluster shutdown ignore all failures as there is no reason to track them + onPreShutdown.add(() -> {simulated.failures.ignoreFailures(); return null;}); } public S simulation() @@ -958,6 +962,18 @@ public class ClusterSimulation implements AutoCloseable } } + for (Callable call : onPreShutdown) + { + try + { + call.call(); + } + catch (Throwable t) + { + fail = Throwables.merge(fail, t); + } + } + try { cluster.close(); @@ -966,7 +982,7 @@ public class ClusterSimulation implements AutoCloseable { fail = Throwables.merge(fail, t); } - for (Callable call : onShutdown) + for (Callable call : onPostShutdown) { try { diff --git a/test/simulator/main/org/apache/cassandra/simulator/systems/Failures.java b/test/simulator/main/org/apache/cassandra/simulator/systems/Failures.java index 2c89d3c28e..863d7861c2 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/systems/Failures.java +++ b/test/simulator/main/org/apache/cassandra/simulator/systems/Failures.java @@ -33,13 +33,20 @@ public class Failures implements Consumer, BufferPool.DebugLeaks, Ref { private final List failures = Collections.synchronizedList(new ArrayList<>()); private volatile boolean hasFailure; + private volatile boolean ignoreFailures = false; public void onFailure(Throwable t) { + if (ignoreFailures) return; failures.add(t); hasFailure = true; } + public void ignoreFailures() + { + ignoreFailures = true; + } + public boolean hasFailure() { return hasFailure; diff --git a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedQuery.java b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedQuery.java index 106cd8c027..d9e33717d9 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedQuery.java +++ b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedQuery.java @@ -18,6 +18,10 @@ package org.apache.cassandra.simulator.systems; +import java.util.function.Predicate; + +import javax.annotation.Nullable; + import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.distributed.api.IInvokableInstance; import org.apache.cassandra.distributed.api.SimpleQueryResult; @@ -25,30 +29,41 @@ import org.apache.cassandra.distributed.impl.Query; public class SimulatedQuery extends SimulatedActionCallable { + private static final Predicate DEFAULT_PREDICATE = i -> true; + + private final Predicate onFailure; + public SimulatedQuery(Object description, SimulatedSystems simulated, IInvokableInstance instance, String query, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, Object... params) { - this(description, Modifiers.NONE, Modifiers.NONE, simulated, instance, query, commitConsistency, serialConsistency, params); + this(description, Modifiers.NONE, Modifiers.NONE, simulated, instance, query, commitConsistency, serialConsistency, null, params); } - public SimulatedQuery(Object description, Modifiers self, Modifiers transitive, SimulatedSystems simulated, IInvokableInstance instance, String query, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, Object[] params) + public SimulatedQuery(Object description, SimulatedSystems simulated, IInvokableInstance instance, String query, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, @Nullable Predicate onFailure, Object... params) + { + this(description, Modifiers.NONE, Modifiers.NONE, simulated, instance, query, commitConsistency, serialConsistency, onFailure, params); + } + + private SimulatedQuery(Object description, Modifiers self, Modifiers transitive, SimulatedSystems simulated, IInvokableInstance instance, String query, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, @Nullable Predicate onFailure, Object[] params) { super(description, self, transitive, simulated, instance, new Query(query, -1, commitConsistency, serialConsistency, params)); + this.onFailure = onFailure == null ? DEFAULT_PREDICATE : onFailure; } public SimulatedQuery(Object description, Modifiers self, Modifiers transitive, SimulatedSystems simulated, IInvokableInstance instance, String query, long timestamp, ConsistencyLevel consistency, Object... params) { - this(description, self, transitive, simulated, instance, query, timestamp, consistency, null, params); + this(description, self, transitive, simulated, instance, query, timestamp, consistency, null, null, params); } - public SimulatedQuery(Object description, Modifiers self, Modifiers transitive, SimulatedSystems simulated, IInvokableInstance instance, String query, long timestamp, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, Object[] params) + private SimulatedQuery(Object description, Modifiers self, Modifiers transitive, SimulatedSystems simulated, IInvokableInstance instance, String query, long timestamp, ConsistencyLevel commitConsistency, ConsistencyLevel serialConsistency, @Nullable Predicate onFailure, Object[] params) { super(description, self, transitive, simulated, instance, new Query(query, timestamp, commitConsistency, serialConsistency, params)); + this.onFailure = onFailure == null ? DEFAULT_PREDICATE : onFailure; } @Override public void accept(SimpleQueryResult success, Throwable failure) { - if (failure != null) + if (failure != null && onFailure.test(failure)) simulated.failures.accept(failure); } } diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/EpochStressTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/EpochStressTest.java index fbc976f938..a76d2ca6a1 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/EpochStressTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/EpochStressTest.java @@ -28,7 +28,6 @@ import java.util.function.Consumer; import java.util.function.Supplier; import com.google.common.util.concurrent.Uninterruptibles; -import org.junit.Ignore; import org.junit.Test; import org.slf4j.Logger; @@ -38,6 +37,7 @@ import accord.api.ConfigurationService; import accord.local.Node; import accord.primitives.Ranges; import accord.topology.TopologyManager; +import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.service.accord.AccordConfigurationService; import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.service.consensus.TransactionalMode; @@ -50,13 +50,14 @@ import org.apache.cassandra.simulator.cluster.ClusterActions; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ClusterMetadataService; import org.apache.cassandra.tcm.membership.NodeId; +import org.apache.cassandra.utils.AssertionUtils; import org.apache.cassandra.utils.Clock; import static org.apache.cassandra.simulator.cluster.ClusterActions.InitialConfiguration.initializeAll; import static org.apache.cassandra.simulator.cluster.ClusterActions.Options.noActions; /** - * In order to run these tests in your IDE, you need to first build a simulator jara + * In order to run these tests in your IDE, you need to first build a simulator jar * * ant simulator-jars * @@ -122,11 +123,20 @@ import static org.apache.cassandra.simulator.cluster.ClusterActions.Options.noAc --add-opens jdk.management/com.sun.management.internal=ALL-UNNAMED --add-opens jdk.management.jfr/jdk.management.jfr=ALL-UNNAMED --add-opens java.desktop/com.sun.beans.introspect=ALL-UNNAMED - */ -@Ignore("Something is currently wrong with SimulationTestBase and multi node tests. CASSANDRA-20744 exposed that simulator was swallowing errors which then causes this test to fail 100% of the time; until the root cause is fixed it doesn't make sense to run this test.") public class EpochStressTest extends SimulationTestBase { + private static boolean allowedExceptions(Throwable t) + { + if (AssertionUtils.isInstanceof(InvalidRequestException.class).matches(t)) + { + String msg = t.getMessage(); + return !msg.contains("policy Retry{remainingMs=0") // this used the "request_timeout" config + && !(msg.contains("policy Retry{") && msg.endsWith(", attempts=12} gave up") ); // this used the cms_retry_delay config + } + return true; + } + @Test public void manyEpochsAndAccordConverges() throws IOException { @@ -147,7 +157,7 @@ public class EpochStressTest extends SimulationTestBase for (int i = 0; i < numEpochs; i++) { int node = random.uniform(1, simulation.cluster.size() + 1); - actions.add(simulation.schemaChange(node, "ALTER TABLE ks.tbl WITH comment = 'step=" + i + "'")); + actions.add(simulation.schemaChange(node, "ALTER TABLE ks.tbl WITH comment = 'step=" + i + '\'', EpochStressTest::allowedExceptions)); } return ActionList.of(actions); }, @@ -160,7 +170,19 @@ public class EpochStressTest extends SimulationTestBase }, config -> config.nodes(3, 3) .dcs(1, 1) - .threadCount(100)); + .threadCount(100), + (rs, config) -> { + // Cluster defaults cause ~100% test failure before reaching 100 epochs, so we cover this case + // for potential visibility violations. With CMS retry configs at 1h, failure rate drops to 10% + // and more seeds complete all 100 epochs. + if (rs.decide(0.3f)) return; // 30% of the time use the cluster defaults, so have a high failure rate for TCM transactions + // Make the retry rate long enough that it should recover most of the time. + // As of this moment this only impacts CMS nodes, so if the CQL is sent to a non-CMS node then request_timeout is used for retry timeout, which makes this fail ~10% of the time + config.set("cms_await_timeout", "3600s") + .set("cms_retry_delay", "3600s") + .set("cms_default_max_retries", Integer.toString(Integer.MAX_VALUE)); + } + ); } private static void validate() diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/SimulationTestBase.java b/test/simulator/test/org/apache/cassandra/simulator/test/SimulationTestBase.java index aedd7b407b..38af502f71 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/SimulationTestBase.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/SimulationTestBase.java @@ -21,6 +21,7 @@ package org.apache.cassandra.simulator.test; import java.io.IOException; import java.util.Arrays; import java.util.Collections; +import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.function.Function; import java.util.function.IntSupplier; @@ -28,6 +29,8 @@ import java.util.function.LongConsumer; import java.util.function.LongSupplier; import java.util.function.Predicate; +import javax.annotation.Nullable; + import com.google.common.collect.Iterators; import org.junit.BeforeClass; @@ -117,11 +120,6 @@ public class SimulationTestBase super(simulated, scheduler, cluster, options); } - protected SimpleSimulation(SimulatedSystems simulated, RunnableActionScheduler scheduler, Cluster cluster, ClusterActions clusterActions) - { - super(simulated, scheduler, cluster, clusterActions); - } - public Action executeQuery(int node, String query, ConsistencyLevel cl, Object... bindings) { return new SimulatedQuery(String.format("Execute query: %s %s %s", query, cl, Arrays.toString(bindings)), @@ -134,13 +132,19 @@ public class SimulationTestBase } public Action schemaChange(int node, String query) + { + return schemaChange(node, query, null); + } + + public Action schemaChange(int node, String query, @Nullable Predicate onFailure) { return new SimulatedQuery(String.format("Schema change: %s", query), simulated, cluster.get(node), query, org.apache.cassandra.distributed.api.ConsistencyLevel.ALL, - null); + null, + onFailure); } protected ActionList initialize() @@ -198,14 +202,17 @@ public class SimulationTestBase protected final Function init; protected final Function test; protected final Function teardown; + protected final BiConsumer configUpdater; DTestClusterSimulationBuilder(Function init, Function test, - Function teardown) + Function teardown, + BiConsumer configUpdater) { this.init = init; this.test = test; this.teardown = teardown; + this.configUpdater = configUpdater; } public ClusterSimulation create(long seed) throws IOException @@ -214,7 +221,7 @@ public class SimulationTestBase random.reset(seed); return new ClusterSimulation<>(random, seed, 1, this, - (c) -> {}, + c -> configUpdater.accept(random, c), (simulated, scheduler, cluster, options) -> new SimpleSimulation(simulated, scheduler, cluster) { protected ActionList initialize() @@ -235,15 +242,46 @@ public class SimulationTestBase } } - public static void simulate(Function init, - Function test, - Function teardown, - Consumer> configure) throws IOException + static void simulate(Function init, + Function test, + Function teardown, + Consumer> configure) throws IOException { - simulate(new DTestClusterSimulationBuilder(init, test, teardown), + simulate(init, test, teardown, configure, (i1, i2) -> {}); + } + + @SuppressWarnings("unused") + static void simulate(long seed, + Function init, + Function test, + Function teardown, + Consumer> configure) throws IOException + { + simulate(seed, init, test, teardown, configure, (i1, i2) -> {}); + } + + static void simulate(Function init, + Function test, + Function teardown, + Consumer> configure, + BiConsumer configUpdater) throws IOException + { + simulate(new DTestClusterSimulationBuilder(init, test, teardown, configUpdater), configure); } + @SuppressWarnings("unused") + static void simulate(long seed, + Function init, + Function test, + Function teardown, + Consumer> configure, + BiConsumer configUpdater) throws IOException + { + simulate(() -> seed, new DTestClusterSimulationBuilder(init, test, teardown, configUpdater), + configure); + } + public static void simulate(ClusterSimulation.Builder factory, Consumer> configure) throws IOException { @@ -272,10 +310,10 @@ public class SimulationTestBase public static void simulate(long seed, ClusterSimulation.SimulationFactory factory, Consumer> configure) throws IOException { - BasicSimulationBuilder builder = new BasicSimulationBuilder() + BasicSimulationBuilder builder = new BasicSimulationBuilder<>() { @Override - Simulation create(SimulatedSystems simulated, RunnableActionScheduler scheduler, Cluster cluster, ClusterActions.Options options) + T create(SimulatedSystems simulated, RunnableActionScheduler scheduler, Cluster cluster, ClusterActions.Options options) { return factory.create(simulated, scheduler, cluster, options); } @@ -439,6 +477,7 @@ public class SimulationTestBase factory.startParked("begin", runnable)); } + @SafeVarargs public static T[] arr(T... arr) { return arr;