diff --git a/.build/checkstyle_suppressions.xml b/.build/checkstyle_suppressions.xml
index ed4d1443f7..230c808c14 100644
--- a/.build/checkstyle_suppressions.xml
+++ b/.build/checkstyle_suppressions.xml
@@ -21,5 +21,4 @@
"https://checkstyle.org/dtds/suppressions_1_1.dtd">
-
diff --git a/build.xml b/build.xml
index 7544c664c1..55614a36a6 100644
--- a/build.xml
+++ b/build.xml
@@ -226,6 +226,24 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/java/org/apache/cassandra/concurrent/InfiniteLoopExecutor.java b/src/java/org/apache/cassandra/concurrent/InfiniteLoopExecutor.java
index ac10a70c30..b576551ac0 100644
--- a/src/java/org/apache/cassandra/concurrent/InfiniteLoopExecutor.java
+++ b/src/java/org/apache/cassandra/concurrent/InfiniteLoopExecutor.java
@@ -52,6 +52,11 @@ public class InfiniteLoopExecutor implements Interruptible
@Shared(scope = Shared.Scope.SIMULATION)
public enum SimulatorSafe { SAFE, UNSAFE }
+ /**
+ * Does this loop always block on some external work provision that is going to be simulator-controlled, or does
+ * it loop periodically? If the latter, it may prevent simulation making progress between phases, and should be
+ * marked as a DAEMON process.
+ */
@Shared(scope = Shared.Scope.SIMULATION)
public enum Daemon { DAEMON, NON_DAEMON }
diff --git a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
index f1f50e589f..c0c739ec82 100644
--- a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
+++ b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
@@ -596,6 +596,8 @@ public enum CassandraRelevantProperties
* can be also done manually for that particular case: {@code flush(SchemaConstants.SCHEMA_KEYSPACE_NAME);}. */
TEST_FLUSH_LOCAL_SCHEMA_CHANGES("cassandra.test.flush_local_schema_changes", "true"),
TEST_HARRY_SWITCH_AFTER("cassandra.test.harry.progression.switch-after", "1"),
+ TEST_HISTORY_VALIDATOR_LOGGING_ENABLED("cassandra.test.history_validator.logging.enabled", "false"),
+ TEST_IGNORE_SIGAR("cassandra.test.ignore_sigar"),
TEST_INTERVAL_TREE_EXPENSIVE_CHECKS("cassandra.test.interval_tree_expensive_checks"),
TEST_INVALID_LEGACY_SSTABLE_ROOT("invalid-legacy-sstable-root"),
TEST_JVM_DTEST_DISABLE_SSL("cassandra.test.disable_ssl"),
diff --git a/src/java/org/apache/cassandra/db/memtable/AbstractAllocatorMemtable.java b/src/java/org/apache/cassandra/db/memtable/AbstractAllocatorMemtable.java
index b431d360ed..2dbe41374f 100644
--- a/src/java/org/apache/cassandra/db/memtable/AbstractAllocatorMemtable.java
+++ b/src/java/org/apache/cassandra/db/memtable/AbstractAllocatorMemtable.java
@@ -220,6 +220,12 @@ public abstract class AbstractAllocatorMemtable extends AbstractMemtableWithComm
if (current instanceof AbstractAllocatorMemtable)
((AbstractAllocatorMemtable) current).flushIfPeriodExpired();
}
+
+ @Override
+ public String toString()
+ {
+ return "Scheduled Flush of " + owner;
+ }
};
ScheduledExecutors.scheduledTasks.scheduleSelfRecurring(runnable, period, TimeUnit.MILLISECONDS);
}
diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java
index 14cc5f5ada..b84b1f25cb 100644
--- a/src/java/org/apache/cassandra/gms/Gossiper.java
+++ b/src/java/org/apache/cassandra/gms/Gossiper.java
@@ -2046,6 +2046,11 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean,
ExecutorUtils.shutdownAndWait(timeout, unit, executor);
}
+ public void shutdownAndWait(long timeout, TimeUnit unit) throws InterruptedException, TimeoutException
+ {
+ ExecutorUtils.shutdownAndWait(timeout, unit, executor);
+ }
+
@Nullable
private String getReleaseVersionString(InetAddressAndPort ep)
{
diff --git a/src/java/org/apache/cassandra/index/IndexStatusManager.java b/src/java/org/apache/cassandra/index/IndexStatusManager.java
index b11ecd1094..0f50a26276 100644
--- a/src/java/org/apache/cassandra/index/IndexStatusManager.java
+++ b/src/java/org/apache/cassandra/index/IndexStatusManager.java
@@ -24,12 +24,13 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
import com.google.common.annotations.VisibleForTesting;
-import org.apache.cassandra.tcm.ClusterMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -46,7 +47,9 @@ import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.locator.Replica;
import org.apache.cassandra.serializers.MarshalException;
import org.apache.cassandra.service.StorageService;
+import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.utils.CassandraVersion;
+import org.apache.cassandra.utils.ExecutorUtils;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.JsonUtils;
@@ -335,4 +338,9 @@ public class IndexStatusManager
{
return keyspace + '.' + index;
}
+
+ public void shutdownAndWait(long interval, TimeUnit unit) throws InterruptedException, TimeoutException
+ {
+ ExecutorUtils.shutdownAndWait(interval, unit, statusPropagationExecutor);
+ }
}
diff --git a/src/java/org/apache/cassandra/journal/ActiveSegment.java b/src/java/org/apache/cassandra/journal/ActiveSegment.java
index 22a3aba766..f16126c157 100644
--- a/src/java/org/apache/cassandra/journal/ActiveSegment.java
+++ b/src/java/org/apache/cassandra/journal/ActiveSegment.java
@@ -33,6 +33,9 @@ import org.apache.cassandra.utils.concurrent.OpOrder;
import org.apache.cassandra.utils.concurrent.Ref;
import org.apache.cassandra.utils.concurrent.WaitQueue;
+import static org.apache.cassandra.utils.Simulate.With.MONITORS;
+
+@Simulate(with=MONITORS)
final class ActiveSegment extends Segment
{
final FileChannel channel;
@@ -247,6 +250,12 @@ final class ActiveSegment extends Segment
* Flush logic; closing and component flushing
*/
+ boolean shouldFlush()
+ {
+ int allocatePosition = this.allocatePosition.get();
+ return lastFlushedOffset < allocatePosition;
+ }
+
/**
* Possibly force a disk flush for this segment file.
* TODO FIXME: calls from outside Flusher + callbacks
diff --git a/src/java/org/apache/cassandra/journal/Flusher.java b/src/java/org/apache/cassandra/journal/Flusher.java
index 436abc5d0e..04411f74c8 100644
--- a/src/java/org/apache/cassandra/journal/Flusher.java
+++ b/src/java/org/apache/cassandra/journal/Flusher.java
@@ -28,6 +28,7 @@ import org.apache.cassandra.concurrent.Interruptible;
import org.apache.cassandra.concurrent.Interruptible.TerminateException;
import org.apache.cassandra.utils.MonotonicClock;
import org.apache.cassandra.utils.NoSpamLogger;
+import org.apache.cassandra.utils.Simulate;
import org.apache.cassandra.utils.concurrent.Semaphore;
import org.apache.cassandra.utils.concurrent.WaitQueue;
@@ -45,6 +46,9 @@ import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis;
import static org.apache.cassandra.utils.Clock.Global.nanoTime;
import static org.apache.cassandra.utils.LocalizeString.toLowerCaseLocalized;
import static org.apache.cassandra.utils.MonotonicClock.Global.preciseTime;
+import static org.apache.cassandra.utils.Simulate.With.GLOBAL_CLOCK;
+import static org.apache.cassandra.utils.Simulate.With.LOCK_SUPPORT;
+import static org.apache.cassandra.utils.Simulate.With.MONITORS;
import static org.apache.cassandra.utils.concurrent.Semaphore.newSemaphore;
import static org.apache.cassandra.utils.concurrent.WaitQueue.newWaitQueue;
@@ -96,6 +100,7 @@ final class Flusher
flushExecutor.shutdown();
}
+ @Simulate(with={MONITORS,GLOBAL_CLOCK,LOCK_SUPPORT})
private class FlushRunnable implements Interruptible.Task
{
private final MonotonicClock clock;
@@ -151,9 +156,17 @@ final class Flusher
if (state == SHUTTING_DOWN)
return;
- long wakeUpAt = startedRunAt + flushPeriodNanos();
- if (wakeUpAt > now)
- haveWork.tryAcquireUntil(1, wakeUpAt);
+ long flushPeriodNanos = flushPeriodNanos();
+ if (flushPeriodNanos <= 0)
+ {
+ haveWork.acquire(1);
+ }
+ else
+ {
+ long wakeUpAt = startedRunAt + flushPeriodNanos;
+ if (wakeUpAt > now)
+ haveWork.tryAcquireUntil(1, wakeUpAt);
+ }
}
private void doFlush()
@@ -168,6 +181,9 @@ final class Flusher
for (ActiveSegment segment : segmentsToFlush)
{
+ if (!segment.shouldFlush())
+ break;
+
syncedSegment = segment.descriptor.timestamp;
syncedOffset = segment.flush();
@@ -202,8 +218,9 @@ final class Flusher
flushCount++;
flushDuration += (finishedFlushAt - startedFlushAt);
- long lag = finishedFlushAt - (startedFlushAt + flushPeriodNanos());
- if (lag <= 0)
+ long flushPeriodNanos = flushPeriodNanos();
+ long lag = finishedFlushAt - (startedFlushAt + flushPeriodNanos);
+ if (flushPeriodNanos <= 0 || lag <= 0)
return;
lagCount++;
@@ -349,7 +366,7 @@ final class Flusher
private long flushPeriodNanos()
{
- return 1_000_000L * params.flushPeriod();
+ return 1_000_000L * params.flushPeriodMillis();
}
private long periodicFlushLagBlockNanos()
diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java
index bb1ada27f7..844f660796 100644
--- a/src/java/org/apache/cassandra/journal/Journal.java
+++ b/src/java/org/apache/cassandra/journal/Journal.java
@@ -50,6 +50,7 @@ import org.apache.cassandra.journal.Segments.ReferencedSegments;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.Crc;
import org.apache.cassandra.utils.JVMStabilityInspector;
+import org.apache.cassandra.utils.Simulate;
import org.apache.cassandra.utils.concurrent.WaitQueue;
import static java.lang.String.format;
@@ -61,6 +62,7 @@ import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe
import static org.apache.cassandra.concurrent.Interruptible.State.NORMAL;
import static org.apache.cassandra.concurrent.Interruptible.State.SHUTTING_DOWN;
import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis;
+import static org.apache.cassandra.utils.Simulate.With.MONITORS;
import static org.apache.cassandra.utils.concurrent.WaitQueue.newWaitQueue;
/**
@@ -77,6 +79,7 @@ import static org.apache.cassandra.utils.concurrent.WaitQueue.newWaitQueue;
* @param the type of keys used to address the records;
must be fixed-size and byte-order comparable
*/
+@Simulate(with=MONITORS)
public class Journal implements Shutdownable
{
private static final Logger logger = LoggerFactory.getLogger(Journal.class);
diff --git a/src/java/org/apache/cassandra/journal/Params.java b/src/java/org/apache/cassandra/journal/Params.java
index f462f450ac..46b382ea27 100644
--- a/src/java/org/apache/cassandra/journal/Params.java
+++ b/src/java/org/apache/cassandra/journal/Params.java
@@ -41,7 +41,7 @@ public interface Params
/**
* @return milliseconds between journal flushes
*/
- int flushPeriod();
+ int flushPeriodMillis();
/**
* @return milliseconds to block writes for while waiting for a slow disk flush to complete
diff --git a/src/java/org/apache/cassandra/metrics/AccordStateCacheMetrics.java b/src/java/org/apache/cassandra/metrics/AccordStateCacheMetrics.java
index fd4308a356..f63fedf282 100644
--- a/src/java/org/apache/cassandra/metrics/AccordStateCacheMetrics.java
+++ b/src/java/org/apache/cassandra/metrics/AccordStateCacheMetrics.java
@@ -32,7 +32,7 @@ public class AccordStateCacheMetrics extends CacheAccessMetrics
public final Histogram objectSize;
- private final Map, CacheAccessMetrics> instanceMetrics = new ConcurrentHashMap<>(2);
+ private final Map instanceMetrics = new ConcurrentHashMap<>(2);
private final String type;
@@ -45,6 +45,8 @@ public class AccordStateCacheMetrics extends CacheAccessMetrics
public CacheAccessMetrics forInstance(Class> klass)
{
- return instanceMetrics.computeIfAbsent(klass, k -> new CacheAccessMetrics(new DefaultNameFactory(TYPE_NAME, String.format("%s-%s", type, k.getSimpleName()))));
+ // cannot make Class> hashCode deterministic, as cannot rewrite - so cannot safely use as Map key if want deterministic simulation
+ // (or we need to create extra hoops to catch this specific case in method rewriting)
+ return instanceMetrics.computeIfAbsent(klass.getSimpleName(), k -> new CacheAccessMetrics(new DefaultNameFactory(TYPE_NAME, String.format("%s-%s", type, k))));
}
}
diff --git a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java
index ad20fea043..31565f8423 100644
--- a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java
+++ b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java
@@ -20,13 +20,12 @@ package org.apache.cassandra.service.accord;
import java.util.Objects;
import java.util.Set;
+import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import javax.annotation.Nullable;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.Sets;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import accord.impl.AbstractConfigurationService;
import accord.local.Node;
@@ -36,6 +35,7 @@ import accord.utils.Invariants;
import accord.utils.async.AsyncResult;
import accord.utils.async.AsyncResults;
import org.apache.cassandra.concurrent.ScheduledExecutors;
+import org.apache.cassandra.concurrent.Shutdownable;
import org.apache.cassandra.concurrent.Stage;
import org.apache.cassandra.gms.FailureDetector;
import org.apache.cassandra.gms.IFailureDetector;
@@ -44,19 +44,23 @@ import org.apache.cassandra.net.MessageDelivery;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.service.accord.AccordKeyspace.EpochDiskState;
import org.apache.cassandra.tcm.ClusterMetadata;
+import org.apache.cassandra.tcm.ClusterMetadataService;
import org.apache.cassandra.tcm.listeners.ChangeListener;
+import org.apache.cassandra.utils.Simulate;
import org.apache.cassandra.utils.concurrent.AsyncPromise;
import org.apache.cassandra.utils.concurrent.Future;
+import static org.apache.cassandra.utils.Simulate.With.MONITORS;
+
// TODO: listen to FailureDetector and rearrange fast path accordingly
-public class AccordConfigurationService extends AbstractConfigurationService implements ChangeListener, AccordEndpointMapper, AccordSyncPropagator.Listener
+@Simulate(with=MONITORS)
+public class AccordConfigurationService extends AbstractConfigurationService implements ChangeListener, AccordEndpointMapper, AccordSyncPropagator.Listener, Shutdownable
{
- private static final Logger logger = LoggerFactory.getLogger(AccordConfigurationService.class);
private final AccordSyncPropagator syncPropagator;
private EpochDiskState diskState = EpochDiskState.EMPTY;
- private enum State { INITIALIZED, LOADING, STARTED }
+ private enum State { INITIALIZED, LOADING, STARTED, SHUTDOWN }
private State state = State.INITIALIZED;
private volatile EndpointMapping mapping = EndpointMapping.EMPTY;
@@ -150,6 +154,35 @@ public class AccordConfigurationService extends AbstractConfigurationService BOOTSTRAP_SUCCESS = ImmediateFuture.success(null);
private final Node node;
@@ -129,6 +135,8 @@ public class AccordService implements IAccordService, Shutdownable
private final AccordJournal journal;
private final AccordVerbHandler extends Request> verbHandler;
private final LocalConfig configuration;
+ @GuardedBy("this")
+ private State state = State.INIT;
private static final IAccordService NOOP_SERVICE = new IAccordService()
{
@@ -308,13 +316,16 @@ public class AccordService implements IAccordService, Shutdownable
}
@Override
- public void startup()
+ public synchronized void startup()
{
+ if (state != State.INIT)
+ return;
journal.start(node);
configService.start();
ClusterMetadataService.instance().log().addListener(configService);
fastPathCoordinator.start();
ClusterMetadataService.instance().log().addListener(fastPathCoordinator);
+ state = State.STARTED;
}
@Override
@@ -526,15 +537,18 @@ public class AccordService implements IAccordService, Shutdownable
}
@Override
- public void shutdown()
+ public synchronized void shutdown()
{
- ExecutorUtils.shutdown(Arrays.asList(scheduler, nodeShutdown, journal));
+ if (state != State.STARTED)
+ return;
+ ExecutorUtils.shutdown(shutdownableSubsystems());
+ state = State.SHUTDOWN;
}
@Override
public Object shutdownNow()
{
- ExecutorUtils.shutdownNow(Arrays.asList(scheduler, nodeShutdown, journal));
+ shutdown();
return null;
}
@@ -543,7 +557,7 @@ public class AccordService implements IAccordService, Shutdownable
{
try
{
- ExecutorUtils.awaitTermination(timeout, units, Arrays.asList(scheduler, nodeShutdown, journal));
+ ExecutorUtils.awaitTermination(timeout, units, shutdownableSubsystems());
return true;
}
catch (TimeoutException e)
@@ -552,11 +566,16 @@ public class AccordService implements IAccordService, Shutdownable
}
}
+ private List shutdownableSubsystems()
+ {
+ return Arrays.asList(scheduler, nodeShutdown, journal, configService);
+ }
+
@VisibleForTesting
@Override
public void shutdownAndWait(long timeout, TimeUnit unit) throws InterruptedException, TimeoutException
{
- scheduler.shutdownNow();
+ shutdown();
ExecutorUtils.shutdownAndWait(timeout, unit, this);
}
diff --git a/src/java/org/apache/cassandra/utils/concurrent/Semaphore.java b/src/java/org/apache/cassandra/utils/concurrent/Semaphore.java
index c9c253f1d5..a0ac316f29 100644
--- a/src/java/org/apache/cassandra/utils/concurrent/Semaphore.java
+++ b/src/java/org/apache/cassandra/utils/concurrent/Semaphore.java
@@ -23,6 +23,7 @@ import java.util.concurrent.TimeUnit;
import org.apache.cassandra.utils.Intercept;
import org.apache.cassandra.utils.Shared;
+import static org.apache.cassandra.utils.Clock.Global.nanoTime;
import static org.apache.cassandra.utils.Shared.Scope.SIMULATION;
@Shared(scope = SIMULATION)
@@ -139,7 +140,7 @@ public interface Semaphore
*/
public boolean tryAcquireUntil(int acquire, long nanoTimeDeadline) throws InterruptedException
{
- long wait = nanoTimeDeadline - System.nanoTime();
+ long wait = nanoTimeDeadline - nanoTime();
return tryAcquire(acquire, Math.max(0, wait), TimeUnit.NANOSECONDS);
}
diff --git a/test/conf/logback-simulator.xml b/test/conf/logback-simulator.xml
index ffa1ffa088..a4c24aab8d 100644
--- a/test/conf/logback-simulator.xml
+++ b/test/conf/logback-simulator.xml
@@ -19,7 +19,8 @@
-
+
+
@@ -38,7 +39,7 @@
- ./build/test/logs/simulator/${run_start}-${run_seed}/${instance_id}/system.log
+ ./build/test/logs/simulator/${run_start}-${run_seed}/cluster-${cluster_id}/${instance_id}/system.log
%-5level [%thread] ${instance_id} %replace(CS:%X{command_store} ){'CS\:\s+', ''}%replace(OP:%X{async_op} ){'OP\:\s+', ''}%date{ISO8601} %msg%n
diff --git a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
index cac40690af..2f81441c19 100644
--- a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
+++ b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
@@ -102,6 +102,7 @@ import org.apache.cassandra.exceptions.StartupException;
import org.apache.cassandra.gms.Gossiper;
import org.apache.cassandra.hints.DTestSerializer;
import org.apache.cassandra.hints.HintsService;
+import org.apache.cassandra.index.IndexStatusManager;
import org.apache.cassandra.index.SecondaryIndexManager;
import org.apache.cassandra.io.IVersionedAsymmetricSerializer;
import org.apache.cassandra.io.sstable.format.SSTableReader;
@@ -628,6 +629,9 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
{
assert config.networkTopology().contains(config.broadcastAddress()) : String.format("Network topology %s doesn't contain the address %s",
config.networkTopology(), config.broadcastAddress());
+ // org.apache.cassandra.distributed.impl.AbstractCluster.startup sets the exception handler for the thread
+ // so extract it to populate ExecutorFactory.Global
+ ExecutorFactory.Global.tryUnsafeSet(new ExecutorFactory.Default(Thread.currentThread().getContextClassLoader(), null, Thread.getDefaultUncaughtExceptionHandler()));
DistributedTestInitialLocationProvider.assign(config.networkTopology());
CassandraDaemon.getInstanceForTesting().activate(false);
// TODO: filters won't work for the messages dispatched during startup
@@ -927,6 +931,11 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
error = parallelRun(error, executor,
() -> Gossiper.instance.stopShutdownAndWait(1L, MINUTES));
}
+ else
+ {
+ error = parallelRun(error, executor,
+ () -> Gossiper.instance.shutdownAndWait(1L, MINUTES));
+ }
error = parallelRun(error, executor, StorageService.instance::disableAutoCompaction);
@@ -970,6 +979,7 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
() -> ActiveRepairService.instance().shutdownNowAndWait(1L, MINUTES),
() -> EpochAwareDebounce.instance.close(),
SnapshotManager.instance::close,
+ () -> IndexStatusManager.instance.shutdownAndWait(1L, MINUTES),
DiskErrorsHandlerService::close
);
diff --git a/test/distributed/org/apache/cassandra/distributed/impl/IsolatedExecutor.java b/test/distributed/org/apache/cassandra/distributed/impl/IsolatedExecutor.java
index 68ff1e71c6..9e84d32df7 100644
--- a/test/distributed/org/apache/cassandra/distributed/impl/IsolatedExecutor.java
+++ b/test/distributed/org/apache/cassandra/distributed/impl/IsolatedExecutor.java
@@ -126,7 +126,7 @@ public class IsolatedExecutor implements IIsolatedExecutor
public Future shutdown()
{
- isolatedExecutor.shutdownNow();
+ isolatedExecutor.shutdown();
return shutdownExecutor.shutdown(name, classLoader, isolatedExecutor, () -> {
// Shutdown logging last - this is not ideal as the logging subsystem is initialized
diff --git a/test/simulator/asm/org/apache/cassandra/simulator/asm/ClassTransformer.java b/test/simulator/asm/org/apache/cassandra/simulator/asm/ClassTransformer.java
index f9bab8eaed..70fa3a6f04 100644
--- a/test/simulator/asm/org/apache/cassandra/simulator/asm/ClassTransformer.java
+++ b/test/simulator/asm/org/apache/cassandra/simulator/asm/ClassTransformer.java
@@ -189,6 +189,10 @@ class ClassTransformer extends ClassVisitor implements MethodWriterSink
{
if (dependentTypes != null)
Utils.visitIfRefType(descriptor, dependentTypes);
+ // org.apache.cassandra.simulator.systems.SimulatedTime.InstanceTime.nanoTime does not change between invokes which causes AbstractQueuedSynchronizer to loop forever,
+ // so need to make the threshold negative to avoid the spin loop.
+ if (className.equals("java/util/concurrent/locks/AbstractQueuedSynchronizer") && name.equals("SPIN_FOR_TIMEOUT_THRESHOLD"))
+ return super.visitField(makePublic(access), name, descriptor, signature, Long.MIN_VALUE);
return super.visitField(makePublic(access), name, descriptor, signature, value);
}
diff --git a/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptAgent.java b/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptAgent.java
index 4cf1546ca8..8774d867d1 100644
--- a/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptAgent.java
+++ b/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptAgent.java
@@ -30,6 +30,7 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.EnumSet;
import java.util.List;
+import java.util.Objects;
import java.util.function.BiFunction;
import java.util.regex.Pattern;
@@ -93,6 +94,9 @@ public class InterceptAgent
if (className.equals("java/lang/Object"))
return transformObject(bytecode);
+ if (className.equals("java/lang/Class"))
+ return transformClass(bytecode);
+
if (className.equals("java/lang/Enum"))
return transformEnum(bytecode);
@@ -103,10 +107,14 @@ public class InterceptAgent
return transformThreadLocalRandom(bytecode);
if (className.startsWith("java/util/concurrent/ConcurrentHashMap"))
- return transformConcurrent(className, bytecode, DETERMINISTIC, NO_PROXY_METHODS);
+ return InterceptAgent.transform(className, bytecode, DETERMINISTIC, NO_PROXY_METHODS);
if (className.startsWith("java/util/concurrent/locks"))
- return transformConcurrent(className, bytecode, SYSTEM_CLOCK, LOCK_SUPPORT, NO_PROXY_METHODS);
+ {
+ if (className.equals("java/util/concurrent/locks/AbstractQueuedSynchronizer"))
+ return InterceptAgent.transformAbstractQueuedSynchronizer(className, bytecode, SYSTEM_CLOCK, LOCK_SUPPORT, NO_PROXY_METHODS);
+ return InterceptAgent.transform(className, bytecode, SYSTEM_CLOCK, LOCK_SUPPORT, NO_PROXY_METHODS);
+ }
return null;
}
@@ -172,6 +180,29 @@ public class InterceptAgent
return transform(bytes, ObjectVisitor::new);
}
+ /**
+ * We don't want Object.toString() to invoke our overridden identityHashCode by virtue of invoking some overridden hashCode()
+ * So we overwrite Object.toString() to replace calls to Object.hashCode() with direct calls to System.identityHashCode()
+ */
+ private static byte[] transformClass(byte[] bytes)
+ {
+ class ClazzVisitor extends ClassVisitor
+ {
+ public ClazzVisitor(int api, ClassVisitor classVisitor)
+ {
+ super(api, classVisitor);
+ }
+
+ @Override
+ public void visitEnd()
+ {
+ new StringHashcode(api).accept(this);
+ super.visitEnd();
+ }
+ }
+ return transform(bytes, ClazzVisitor::new);
+ }
+
/**
* We want Enum to have a deterministic hashCode() so we simply forward calls to ordinal()
*/
@@ -314,7 +345,7 @@ public class InterceptAgent
else
{
MethodVisitor mv = super.visitMethod(access, name, descriptor, signature, exceptions);
- if (determinismCheck && (name.equals("nextSeed") || name.equals("nextSecondarySeed")))
+ if (determinismCheck && (name.equals("nextSeed") || name.equals("nextSecondarySeed") || name.equals("advanceProbe")))
mv = new ThreadLocalRandomCheckTransformer(api, mv);
return mv;
}
@@ -323,7 +354,61 @@ public class InterceptAgent
return transform(bytes, ThreadLocalRandomVisitor::new);
}
- private static byte[] transform(byte[] bytes, BiFunction constructor)
+ /**
+ * We require ThreadLocalRandom to be deterministic, so we modify its initialisation method to invoke a
+ * global deterministic random value generator
+ */
+ private static byte[] transformAbstractQueuedSynchronizer(String className, byte[] bytes, Flag flag, Flag ... flags)
+ {
+ class AbstractQueuedSynchronizerVisitor extends ClassVisitor
+ {
+ private long defaultSpinForTimeoutThreshold = 1000L;
+
+ public AbstractQueuedSynchronizerVisitor(int api, ClassVisitor classVisitor)
+ {
+ super(api, classVisitor);
+ }
+
+ @Override
+ public FieldVisitor visitField(int access, String name, String descriptor, String signature, Object value)
+ {
+ if (name.equals("SPIN_FOR_TIMEOUT_THRESHOLD"))
+ {
+ defaultSpinForTimeoutThreshold = (Long)value;
+ return super.visitField(access, name, descriptor, signature, 0L);
+ }
+
+ return super.visitField(access, name, descriptor, signature, value);
+ }
+
+ @Override
+ public MethodVisitor visitMethod(int access, String name, String descriptor, String signature, String[] exceptions)
+ {
+ /// !!!!! WARNING !!!!!
+ /// THIS IS SUPER BRITTLE BECAUSE rt.jar INLINES GETSTATIC AS LDC
+ // TODO (desired): visit constructor to fetch actual value of constant in case changes in future release -
+ // but this is brittle enough changes upstream will likely need revisiting anyway
+ MethodVisitor mv = super.visitMethod(access, name, descriptor, signature, exceptions);
+ if (!name.equals("doAcquireNanos") && !name.equals("doAcquireSharedNanos"))
+ return mv;
+
+ return new MethodVisitor(api, mv)
+ {
+ @Override
+ public void visitLdcInsn(Object value)
+ {
+ if (Objects.equals(defaultSpinForTimeoutThreshold, value))
+ super.visitLdcInsn(0L);
+ else
+ super.visitLdcInsn(value);
+ }
+ };
+ }
+ }
+ return transform(className, bytes, AbstractQueuedSynchronizerVisitor::new, flag, flags);
+ }
+
+ private static byte[] transform(byte[] bytes, BiFunction constructor)
{
ClassWriter out = new ClassWriter(0);
ClassReader in = new ClassReader(bytes);
@@ -332,7 +417,7 @@ public class InterceptAgent
return out.toByteArray();
}
- private static byte[] transformConcurrent(String className, byte[] bytes, Flag flag, Flag ... flags)
+ private static byte[] transform(String className, byte[] bytes, Flag flag, Flag ... flags)
{
ClassTransformer transformer = new ClassTransformer(BYTECODE_VERSION, className, EnumSet.of(flag, flags), null);
transformer.readAndTransform(bytes);
@@ -340,4 +425,13 @@ public class InterceptAgent
return null;
return transformer.toBytes();
}
+
+ private static byte[] transform(String className, byte[] bytes, BiFunction constructor, Flag flag, Flag ... flags)
+ {
+ ClassReader in = new ClassReader(bytes);
+ ClassTransformer transformer = new ClassTransformer(BYTECODE_VERSION, className, EnumSet.of(flag, flags), null);
+ ClassVisitor extraTransformer = constructor.apply(BYTECODE_VERSION, transformer);
+ in.accept(extraTransformer, 0);
+ return transformer.toBytes();
+ }
}
diff --git a/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptClasses.java b/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptClasses.java
index dd53ce067f..5043012472 100644
--- a/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptClasses.java
+++ b/test/simulator/asm/org/apache/cassandra/simulator/asm/InterceptClasses.java
@@ -62,6 +62,8 @@ public class InterceptClasses implements BiFunction
"|org[/.]apache[/.]cassandra[/.]distributed[/.]impl[/.]DirectStreamingConnectionFactory.*" +
"|org[/.]apache[/.]cassandra[/.]db[/.]commitlog[/.].*" +
"|org[/.]apache[/.]cassandra[/.]service[/.]paxos[/.].*" +
+ "|org[/.]apache[/.]cassandra[/.]service[/.]accord[/.].*" +
+ "|org[/.]apache[/.]cassandra[/.]journal[/.].*" +
"|accord[/.].*"
);
diff --git a/test/simulator/asm/org/apache/cassandra/simulator/asm/MonitorMethodTransformer.java b/test/simulator/asm/org/apache/cassandra/simulator/asm/MonitorMethodTransformer.java
index d9c9c7ad94..a7c21bbba7 100644
--- a/test/simulator/asm/org/apache/cassandra/simulator/asm/MonitorMethodTransformer.java
+++ b/test/simulator/asm/org/apache/cassandra/simulator/asm/MonitorMethodTransformer.java
@@ -122,8 +122,7 @@ class MonitorMethodTransformer extends MethodNode
}
int invokeCode;
- if (isInstanceMethod && (access & Opcodes.ACC_PRIVATE) != 0) invokeCode = Opcodes.INVOKESPECIAL;
- else if (isInstanceMethod) invokeCode = Opcodes.INVOKEVIRTUAL;
+ if (isInstanceMethod) invokeCode = Opcodes.INVOKESPECIAL;
else invokeCode = Opcodes.INVOKESTATIC;
return invokeCode;
}
diff --git a/test/simulator/asm/org/apache/cassandra/simulator/asm/StringHashcode.java b/test/simulator/asm/org/apache/cassandra/simulator/asm/StringHashcode.java
new file mode 100644
index 0000000000..fc3c57f8b5
--- /dev/null
+++ b/test/simulator/asm/org/apache/cassandra/simulator/asm/StringHashcode.java
@@ -0,0 +1,43 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.cassandra.simulator.asm;
+
+import org.objectweb.asm.Opcodes;
+import org.objectweb.asm.tree.InsnNode;
+import org.objectweb.asm.tree.LabelNode;
+import org.objectweb.asm.tree.MethodInsnNode;
+import org.objectweb.asm.tree.MethodNode;
+
+/**
+ * Generate a new hashCode method in the class that invokes a deterministic hashCode generator
+ */
+class StringHashcode extends MethodNode
+{
+ StringHashcode(int api)
+ {
+ super(api, Opcodes.ACC_PUBLIC, "hashCode", "()I", null, null);
+ maxLocals = 1;
+ maxStack = 1;
+ instructions.add(new LabelNode());
+ instructions.add(new MethodInsnNode(Opcodes.INVOKEVIRTUAL, "java/lang/Object", "toString", "()Ljava/lang/String;", false));
+ instructions.add(new LabelNode());
+ instructions.add(new MethodInsnNode(Opcodes.INVOKEVIRTUAL, "java/lang/Object", "hashCode", "(Ljava/lang/Object;)I", false));
+ instructions.add(new InsnNode(Opcodes.IRETURN));
+ }
+}
diff --git a/test/simulator/main/org/apache/cassandra/simulator/ActionSchedule.java b/test/simulator/main/org/apache/cassandra/simulator/ActionSchedule.java
index 427a777abe..39666077ea 100644
--- a/test/simulator/main/org/apache/cassandra/simulator/ActionSchedule.java
+++ b/test/simulator/main/org/apache/cassandra/simulator/ActionSchedule.java
@@ -281,6 +281,12 @@ public class ActionSchedule implements CloseableIterator