post-rebase fixes, mostly around CASSANDRA-19341 and CASSANDRA-19567

This commit is contained in:
Caleb Rackliffe 2024-05-08 12:24:33 -05:00 committed by David Capwell
parent 980f1963f3
commit 1b3c3f32a4
44 changed files with 598 additions and 200 deletions

2
.gitmodules vendored
View File

@ -1,4 +1,4 @@
[submodule "modules/accord"]
path = modules/accord
url = ../cassandra-accord.git
url = https://github.com/apache/cassandra-accord.git
branch = trunk

@ -1 +1 @@
Subproject commit 202e67358396a1e413e29498bea71047bd586d06
Subproject commit 256b35e27d170db9fcd8024d5678b4f6e9d3a956

View File

@ -18,6 +18,8 @@
package org.apache.cassandra.config;
import com.fasterxml.jackson.annotation.JsonIgnore;
import org.apache.cassandra.journal.Params;
import org.apache.cassandra.service.consensus.TransactionalMode;
public class AccordSpec
@ -71,4 +73,62 @@ public class AccordSpec
public TransactionalMode default_transactional_mode = TransactionalMode.off;
public boolean ephemeralReadEnabled = false;
public boolean state_cache_listener_jfr_enabled = true;
public final JournalSpec journal = new JournalSpec();
public static class JournalSpec implements Params
{
public int segmentSize = 32 << 20;
public FailurePolicy failurePolicy = FailurePolicy.STOP;
public FlushMode flushMode = FlushMode.BATCH;
public DurationSpec.IntMillisecondsBound flushPeriod; // pulls default from 'commitlog_sync_period'
public DurationSpec.IntMillisecondsBound periodicFlushLagBlock = new DurationSpec.IntMillisecondsBound("1500ms");
@Override
public int segmentSize()
{
return segmentSize;
}
@Override
public FailurePolicy failurePolicy()
{
return failurePolicy;
}
@Override
public FlushMode flushMode()
{
return flushMode;
}
@JsonIgnore
@Override
public int flushPeriodMillis()
{
return flushPeriod == null ? DatabaseDescriptor.getCommitLogSyncPeriod()
: flushPeriod.toMilliseconds();
}
@JsonIgnore
@Override
public int periodicFlushLagBlock()
{
return periodicFlushLagBlock.toMilliseconds();
}
/**
* This is required by the journal, but we don't have multiple versions, so block it from showing up, so we don't need to worry about maintaining it
*/
@JsonIgnore
@Override
public int userVersion()
{
/*
* NOTE: when accord journal version gets bumped, expose it via yaml.
* This way operators can force previous version on upgrade, temporarily,
* to allow easier downgrades if something goes wrong.
*/
return 1;
}
}
}

View File

@ -3661,6 +3661,11 @@ public class DatabaseDescriptor
return conf.paxos_topology_repair_strict_each_quorum;
}
public static AccordSpec getAccord()
{
return conf.accord;
}
public static AccordSpec.TransactionalRangeMigration getTransactionalRangeMigration()
{
return conf.accord.range_migration;

View File

@ -19,8 +19,8 @@ package org.apache.cassandra.dht;
import java.io.IOException;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
/**
* Versioned serializer where the serialization depends on partitioner.
@ -28,18 +28,8 @@ import org.apache.cassandra.io.util.DataOutputPlus;
* On serialization the partitioner is given by the entity being serialized. To deserialize the partitioner used must
* be known to the calling method.
*/
public interface IPartitionerDependentSerializer<T>
public interface IPartitionerDependentSerializer<T> extends IVersionedSerializer<T>
{
/**
* Serialize the specified type into the specified DataOutputStream instance.
*
* @param t type that needs to be serialized
* @param out DataOutput into which serialization needs to happen.
* @param version protocol version
* @throws java.io.IOException if serialization fails
*/
public void serialize(T t, DataOutputPlus out, int version) throws IOException;
/**
* Deserialize into the specified DataInputStream instance.
* @param in DataInput from which deserialization needs to happen.
@ -51,11 +41,8 @@ public interface IPartitionerDependentSerializer<T>
*/
public T deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException;
/**
* Calculate serialized size of object without actually serializing.
* @param t object to calculate serialized size
* @param version protocol version
* @return serialized size of object t
*/
public long serializedSize(T t, int version);
default T deserialize(DataInputPlus in, int version) throws IOException
{
return deserialize(in, null, version);
}
}

View File

@ -20,6 +20,12 @@ package org.apache.cassandra.dht;
import java.io.IOException;
import java.io.Serializable;
import java.nio.ByteBuffer;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import com.google.common.collect.Sets;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.TypeSizes;
@ -35,6 +41,8 @@ import org.apache.cassandra.utils.bytecomparable.ByteSourceInverse;
public abstract class Token implements RingPosition<Token>, Serializable
{
private static final Logger logger = LoggerFactory.getLogger(Token.class);
private static final long serialVersionUID = 1L;
public static final TokenSerializer serializer = new TokenSerializer();
@ -178,11 +186,17 @@ public abstract class Token implements RingPosition<Token>, Serializable
}
}
public static volatile boolean logPartitioner = false;
public static final Set<Class<? extends IPartitioner>> serializePartitioners = Sets.newSetFromMap(new ConcurrentHashMap<>());
public static final Set<Class<? extends IPartitioner>> deserializePartitioners = Sets.newSetFromMap(new ConcurrentHashMap<>());
public static class CompactTokenSerializer implements IPartitionerDependentSerializer<Token>
{
public void serialize(Token token, DataOutputPlus out, int version) throws IOException
{
IPartitioner p = token.getPartitioner();
if (logPartitioner && serializePartitioners.add(p.getClass()))
logger.debug("Serializing token with partitioner " + p);
if (!p.isFixedLength())
out.writeUnsignedVInt32(p.getTokenFactory().byteSize(token));
p.getTokenFactory().serialize(token, out);
@ -191,6 +205,8 @@ public abstract class Token implements RingPosition<Token>, Serializable
public Token deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException
{
int size = p.isFixedLength() ? p.getMaxTokenSize() : in.readUnsignedVInt32();
if (logPartitioner && deserializePartitioners.add(p.getClass()))
logger.debug("Deserializing token with partitioner " + p);
byte[] bytes = new byte[size];
in.readFully(bytes);
return p.getTokenFactory().fromByteArray(ByteBuffer.wrap(bytes));

View File

@ -30,7 +30,7 @@ import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics;
// Stolen from org.apache.cassandra.index.sai.metrics.AbstractMetrics
public class IndexMetrics
{
private static final String TYPE = "RouteIndex";
public static final String TYPE = "RouteIndex";
private static final String SCOPE = "IndexMetrics";
private final List<CassandraMetricsRegistry.MetricName> tracked = new ArrayList<>();
@ -61,7 +61,7 @@ public class IndexMetrics
{
metricScope += '.' + indexName;
}
metricScope += '.' + SCOPE + '.' + name;
metricScope += '.' + SCOPE;
CassandraMetricsRegistry.MetricName metricName = new CassandraMetricsRegistry.MetricName(DefaultNameFactory.GROUP_NAME,
TYPE, name, metricScope, createMBeanName(name, SCOPE));

View File

@ -32,7 +32,7 @@ public interface IVersionedAsymmetricSerializer<In, Out>
* @param version protocol version
* @throws IOException if serialization fails
*/
public void serialize(In t, DataOutputPlus out, int version) throws IOException;
void serialize(In t, DataOutputPlus out, int version) throws IOException;
/**
* Deserialize into the specified DataInputStream instance.
@ -41,7 +41,7 @@ public interface IVersionedAsymmetricSerializer<In, Out>
* @return the type that was deserialized
* @throws IOException if deserialization fails
*/
public Out deserialize(DataInputPlus in, int version) throws IOException;
Out deserialize(DataInputPlus in, int version) throws IOException;
/**
* Calculate serialized size of object without actually serializing.
@ -49,5 +49,5 @@ public interface IVersionedAsymmetricSerializer<In, Out>
* @param version protocol version
* @return serialized size of object t
*/
public long serializedSize(In t, int version);
long serializedSize(In t, int version);
}

View File

@ -23,8 +23,10 @@ import org.apache.cassandra.metrics.CassandraMetricsRegistry;
import org.apache.cassandra.metrics.DefaultNameFactory;
import org.apache.cassandra.metrics.MetricNameFactory;
final class Metrics<K, V>
public final class Metrics<K, V>
{
public static final String TYPE_NAME = "Journal";
private static final String WAITING_ON_FLUSH = "WaitingOnFlush";
private static final String WAITING_ON_ALLOCATION = "WaitingOnSegmentAllocation";
private static final String WRITTEN_ENTRIES = "WrittenEntries";
@ -49,7 +51,7 @@ final class Metrics<K, V>
Metrics(String name)
{
this.factory = new DefaultNameFactory("Journal", name);
this.factory = new DefaultNameFactory(TYPE_NAME, name);
}
void register(Flusher<K, V> flusher)

View File

@ -114,6 +114,8 @@ public class CassandraMetricsRegistry extends MetricRegistry
// for virtual tables.
metricGroups = ImmutableSet.<String>builder()
.add(AbstractMetrics.TYPE)
.add(AccordMetrics.ACCORD_COORDINATOR)
.add(AccordMetrics.ACCORD_REPLICA)
.add(BatchMetrics.TYPE_NAME)
.add(BufferPoolMetrics.TYPE_NAME)
.add(CIDRAuthorizerMetrics.TYPE_NAME)
@ -130,8 +132,10 @@ public class CassandraMetricsRegistry extends MetricRegistry
.add(DroppedMessageMetrics.TYPE)
.add(HintedHandoffMetrics.TYPE_NAME)
.add(HintsServiceMetrics.TYPE_NAME)
.add(org.apache.cassandra.index.accord.IndexMetrics.TYPE)
.add(InternodeInboundMetrics.TYPE_NAME)
.add(InternodeOutboundMetrics.TYPE_NAME)
.add(org.apache.cassandra.journal.Metrics.TYPE_NAME)
.add(KeyspaceMetrics.TYPE_NAME)
.add(MemtablePool.TYPE_NAME)
.add(MessagingMetrics.TYPE_NAME)

View File

@ -346,15 +346,15 @@ public enum Verb
ACCORD_SYNC_NOTIFY_REQ (151, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SIMPLE_RSP ),
ACCORD_APPLY_AND_WAIT_REQ (152, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.readData, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP),
ACCORD_APPLY_AND_WAIT_REQ (152, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.readData, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP),
CONSENSUS_KEY_MIGRATION (153, P1, writeTimeout, MUTATION, () -> ConsensusKeyMigrationFinished.serializer,() -> ConsensusKeyMigrationState.consensusKeyMigrationFinishedHandler),
ACCORD_INTEROP_READ_RSP (154, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.replySerializer, RESPONSE_HANDLER),
ACCORD_INTEROP_READ_REQ (155, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.requestSerializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_RSP),
ACCORD_INTEROP_COMMIT_REQ (156, P2, writeTimeout, IMMEDIATE, () -> AccordInteropCommit.serializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_RSP),
ACCORD_INTEROP_READ_REQ (155, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.requestSerializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_RSP),
ACCORD_INTEROP_COMMIT_REQ (156, P2, writeTimeout, IMMEDIATE, () -> AccordInteropCommit.serializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_RSP),
ACCORD_INTEROP_READ_REPAIR_RSP (157, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.replySerializer, RESPONSE_HANDLER),
ACCORD_INTEROP_READ_REPAIR_REQ (158, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.requestSerializer, () -> AccordService.instance().verbHandler(), ACCORD_INTEROP_READ_REPAIR_RSP),
ACCORD_INTEROP_READ_REPAIR_REQ (158, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.requestSerializer, AccordService::verbHandlerOrNoop, ACCORD_INTEROP_READ_REPAIR_RSP),
ACCORD_INTEROP_APPLY_REQ (160, P2, writeTimeout, IMMEDIATE, () -> AccordInteropApply.serializer, AccordService::verbHandlerOrNoop, ACCORD_APPLY_RSP),
// generic failure response

View File

@ -24,7 +24,7 @@ import java.util.Objects;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.dht.IPartitionerDependentSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.locator.InetAddressAndPort;
@ -79,7 +79,7 @@ public class SyncResponse extends RepairMessage
return Objects.hash(desc, success, nodes, summaries);
}
public static final IVersionedSerializer<SyncResponse> serializer = new IVersionedSerializer<SyncResponse>()
public static final IPartitionerDependentSerializer<SyncResponse> serializer = new IPartitionerDependentSerializer<SyncResponse>()
{
public void serialize(SyncResponse message, DataOutputPlus out, int version) throws IOException
{
@ -94,7 +94,8 @@ public class SyncResponse extends RepairMessage
}
}
public SyncResponse deserialize(DataInputPlus in, int version) throws IOException
@Override
public SyncResponse deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException
{
RepairJobDesc desc = RepairJobDesc.serializer.deserialize(in, version);
SyncNodePair nodes = SyncNodePair.serializer.deserialize(in, version);
@ -104,7 +105,7 @@ public class SyncResponse extends RepairMessage
List<SessionSummary> summaries = new ArrayList<>(numSummaries);
for (int i=0; i<numSummaries; i++)
{
summaries.add(SessionSummary.serializer.deserialize(in, IPartitioner.global(), version));
summaries.add(SessionSummary.serializer.deserialize(in, partitioner, version));
}
return new SyncResponse(desc, nodes, success, summaries);

View File

@ -42,6 +42,7 @@ import com.google.common.collect.Multimap;
import com.google.common.primitives.Ints;
import com.google.common.primitives.Longs;
import accord.messages.ApplyThenWaitUntilApplied;
import org.agrona.collections.Long2ObjectHashMap;
import org.agrona.collections.LongArrayList;
import org.agrona.collections.ObjectHashSet;
@ -156,50 +157,6 @@ public class AccordJournal implements IJournal, Shutdownable
private static final ThreadLocal<byte[]> keyCRCBytes = ThreadLocal.withInitial(() -> new byte[21]);
static final Params PARAMS = new Params()
{
@Override
public int segmentSize()
{
return 32 << 20;
}
@Override
public FailurePolicy failurePolicy()
{
return FailurePolicy.STOP;
}
@Override
public FlushMode flushMode()
{
return FlushMode.BATCH;
}
@Override
public int flushPeriodMillis()
{
return DatabaseDescriptor.getCommitLogSyncPeriod();
}
@Override
public int periodicFlushLagBlock()
{
return 1500;
}
@Override
public int userVersion()
{
/*
* NOTE: when accord journal version gets bumped, expose it via yaml.
* This way operators can force previous version on upgrade, temporarily,
* to allow easier downgrades if something goes wrong.
*/
return 1;
}
};
private final File directory;
private final Journal<Key, Object> journal;
private final AccordEndpointMapper endpointMapper;
@ -219,10 +176,10 @@ public class AccordJournal implements IJournal, Shutdownable
private final FrameApplicator frameApplicator = new FrameApplicator();
@VisibleForTesting
public AccordJournal(AccordEndpointMapper endpointMapper)
public AccordJournal(AccordEndpointMapper endpointMapper, Params params)
{
this.directory = new File(DatabaseDescriptor.getAccordJournalDirectory());
this.journal = new Journal<>("AccordJournal", directory, PARAMS, new JournalCallbacks(), Key.SUPPORT, RECORD_SERIALIZER);
this.journal = new Journal<>("AccordJournal", directory, params, new JournalCallbacks(), Key.SUPPORT, RECORD_SERIALIZER);
this.endpointMapper = endpointMapper;
}
@ -969,6 +926,22 @@ public class AccordJournal implements IJournal, Shutdownable
}
}
msgTypeToSynonymousTypesMap = ImmutableListMultimap.copyOf(msgTypeToSynonymousTypes);
//TODO (now): enable as this shows we are currently missing a message
// IllegalStateException e = null;
// for (MessageType t : MessageType.values)
// {
// if (!t.hasSideEffects()) continue;
// Type matches = msgTypeToTypeMap.get(t);
// if (matches == null)
// {
// IllegalStateException ise = new IllegalStateException("Missing MessageType " + t);
// if (e == null) e = ise;
// else e.addSuppressed(ise);
// }
// }
// if (e != null)
// throw e;
}
static Type fromId(int id)
@ -1164,7 +1137,7 @@ public class AccordJournal implements IJournal, Shutdownable
while (null != (request = unframedRequests.poll()))
{
long waitForEpoch = request.waitForEpoch;
if (!node.topology().hasEpoch(waitForEpoch))
if (waitForEpoch != 0 && !node.topology().hasEpoch(waitForEpoch))
{
delayedRequests.computeIfAbsent(waitForEpoch, ignore -> new ArrayList<>()).add(request);
if (!waitForEpochs.containsLong(waitForEpoch))
@ -1394,6 +1367,12 @@ public class AccordJournal implements IJournal, Shutdownable
this.txnId = txnId;
}
@Override
public TxnId txnId()
{
return txnId;
}
@Override
public Set<MessageType> test(Set<MessageType> messages)
{
@ -1464,6 +1443,12 @@ public class AccordJournal implements IJournal, Shutdownable
return readMessage(txnId, STABLE_FAST_PATH_REQ, Commit.class);
}
@Override
public Commit stableSlowPath()
{
return readMessage(txnId, STABLE_SLOW_PATH_REQ, Commit.class);
}
@Override
public Commit stableMaximal()
{
@ -1493,6 +1478,18 @@ public class AccordJournal implements IJournal, Shutdownable
{
return readMessage(txnId, PROPAGATE_APPLY_MSG, Propagate.class);
}
@Override
public Propagate propagateOther()
{
return readMessage(txnId, PROPAGATE_OTHER_MSG, Propagate.class);
}
@Override
public ApplyThenWaitUntilApplied applyThenWaitUntilApplied()
{
return readMessage(txnId, APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, ApplyThenWaitUntilApplied.class);
}
}
private final class LoggingMessageProvider implements SerializerSupport.MessageProvider
@ -1506,6 +1503,12 @@ public class AccordJournal implements IJournal, Shutdownable
this.provider = provider;
}
@Override
public TxnId txnId()
{
return txnId;
}
@Override
public Set<MessageType> test(Set<MessageType> messages)
{
@ -1587,6 +1590,15 @@ public class AccordJournal implements IJournal, Shutdownable
return commit;
}
@Override
public Commit stableSlowPath()
{
logger.debug("Fetching {} message for {}", STABLE_SLOW_PATH_REQ, txnId);
Commit commit = provider.stableSlowPath();
logger.debug("Fetched {} message for {}: {}", STABLE_SLOW_PATH_REQ, txnId, commit);
return commit;
}
@Override
public Commit stableMaximal()
{
@ -1631,5 +1643,23 @@ public class AccordJournal implements IJournal, Shutdownable
logger.debug("Fetched {} message for {}: {}", PROPAGATE_APPLY_MSG, txnId, propagate);
return propagate;
}
@Override
public Propagate propagateOther()
{
logger.debug("Fetching {} message for {}", PROPAGATE_OTHER_MSG, txnId);
Propagate propagate = provider.propagateOther();
logger.debug("Fetched {} message for {}: {}", PROPAGATE_OTHER_MSG, txnId, propagate);
return propagate;
}
@Override
public ApplyThenWaitUntilApplied applyThenWaitUntilApplied()
{
logger.debug("Fetching {} message for {}", APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, txnId);
ApplyThenWaitUntilApplied apply = provider.applyThenWaitUntilApplied();
logger.debug("Fetched {} message for {}: {}", APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, txnId, apply);
return apply;
}
}
}

View File

@ -250,6 +250,7 @@ public class AccordKeyspace
+ format("execute_at %s,", TIMESTAMP_TUPLE)
+ format("promised_ballot %s,", TIMESTAMP_TUPLE)
+ format("accepted_ballot %s,", TIMESTAMP_TUPLE)
+ format("execute_atleast %s,", TIMESTAMP_TUPLE)
+ "waiting_on blob,"
+ "listeners set<blob>, "
+ "PRIMARY KEY((store_id, domain, txn_id))"
@ -298,6 +299,7 @@ public class AccordKeyspace
public static final ColumnMetadata execute_at = getColumn(Commands, "execute_at");
static final ColumnMetadata promised_ballot = getColumn(Commands, "promised_ballot");
static final ColumnMetadata accepted_ballot = getColumn(Commands, "accepted_ballot");
static final ColumnMetadata execute_atleast = getColumn(Commands, "execute_atleast");
static final ColumnMetadata waiting_on = getColumn(Commands, "waiting_on");
static final ColumnMetadata listeners = getColumn(Commands, "listeners");
@ -857,6 +859,8 @@ public class AccordKeyspace
addCellIfModified(CommandsColumns.execute_at, Command::executeAt, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command);
addCellIfModified(CommandsColumns.promised_ballot, Command::promised, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command);
addCellIfModified(CommandsColumns.accepted_ballot, Command::acceptedOrCommitted, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command);
if (command.txnId().kind().awaitsOnlyDeps())
addCellIfModified(CommandsColumns.execute_atleast, Command::executesAtLeast, AccordKeyspace::serializeTimestamp, builder, timestampMicros, nowInSeconds, original, command);
if (command.isStable() && !command.isTruncated())
{
@ -1230,11 +1234,12 @@ public class AccordKeyspace
Timestamp executeAt = deserializeExecuteAtOrNull(row);
Ballot promised = deserializePromisedOrNull(row);
Ballot accepted = deserializeAcceptedOrNull(row);
Timestamp executeAtLeast = status.is(Status.Truncated) && txnId.kind().awaitsOnlyDeps() ? deserializeExecuteAtLeastOrNull(row) : null;
WaitingOnProvider waitingOn = deserializeWaitingOn(txnId, row);
MessageProvider messages = commandStore.makeMessageProvider(txnId);
return SerializerSupport.reconstruct(commandStore.unsafeRangesForEpoch(), attrs, status, executeAt, promised, accepted, waitingOn, messages);
return SerializerSupport.reconstruct(commandStore.unsafeRangesForEpoch(), attrs, status, executeAt, executeAtLeast, promised, accepted, waitingOn, messages);
}
catch (Throwable t)
{
@ -1307,6 +1312,11 @@ public class AccordKeyspace
return deserializeTimestampOrNull(row, "execute_at", Timestamp::fromBits);
}
public static Timestamp deserializeExecuteAtLeastOrNull(UntypedResultSet.Row row)
{
return deserializeTimestampOrNull(row, "execute_atleast", Timestamp::fromBits);
}
public static Ballot deserializePromisedOrNull(UntypedResultSet.Row row)
{
return deserializeTimestampOrNull(row.getBlob("promised_ballot"), Ballot::fromBits);

View File

@ -43,6 +43,7 @@ import accord.impl.CoordinateDurabilityScheduling;
import accord.primitives.SyncPoint;
import org.apache.cassandra.config.CassandraRelevantProperties;
import org.apache.cassandra.cql3.statements.RequestValidations;
import org.apache.cassandra.exceptions.RequestExecutionException;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.service.accord.interop.AccordInteropAdapter.AccordInteropFactory;
@ -315,7 +316,7 @@ public class AccordService implements IAccordService, Shutdownable
this.scheduler = new AccordScheduler();
this.dataStore = new AccordDataStore();
this.configuration = new AccordConfiguration(DatabaseDescriptor.getRawConfig());
this.journal = new AccordJournal(configService);
this.journal = new AccordJournal(configService, DatabaseDescriptor.getAccord().journal);
this.node = new Node(localId,
messageSink,
this::handleLocalRequest,
@ -447,7 +448,7 @@ public class AccordService implements IAccordService, Shutdownable
private long doWithRetries(LongSupplier action, int retryAttempts, long initialBackoffMillis, long maxBackoffMillis) throws InterruptedException
{
// Since we could end up having the barrier transaction or the transaction it listens to invalidated
CoordinationFailed existingFailures = null;
RuntimeException existingFailures = null;
Long success = null;
long backoffMillis = 0;
for (int attempt = 0; attempt < retryAttempts; attempt++)
@ -468,7 +469,7 @@ public class AccordService implements IAccordService, Shutdownable
success = action.getAsLong();
break;
}
catch (CoordinationFailed newFailures)
catch (RequestExecutionException | CoordinationFailed newFailures)
{
existingFailures = Throwables.merge(existingFailures, newFailures);
}

View File

@ -215,7 +215,7 @@ public abstract class AsyncOperation<R> extends AsyncChains.Head<R> implements R
commandStore.abortCurrentOperation();
case LOADING:
context.releaseResources(commandStore);
commandStore.executionOrder().unregister(this);
commandStore.executionOrder().unregisterOutOfOrder(this);
case INITIALIZED:
break; // nothing to clean up, call callback
}
@ -239,6 +239,8 @@ public abstract class AsyncOperation<R> extends AsyncChains.Head<R> implements R
default: throw new IllegalStateException("Unexpected state " + state);
case INITIALIZED:
canRun = commandStore.executionOrder().register(this);
if (Invariants.isParanoid())
Invariants.checkState(canRun.booleanValue() == commandStore.executionOrder().canRun(this), "Register of %s returned canRun=%s but canRun returned %s!", this, canRun, !canRun);
state(LOADING);
case LOADING:
if (null == canRun)

View File

@ -21,6 +21,7 @@ import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.IdentityHashMap;
import java.util.List;
import java.util.function.Consumer;
import accord.api.Key;
import accord.api.RoutingKey;
@ -82,31 +83,9 @@ public class ExecutionOrder
}
}
Conflicts remove(AsyncOperation<?> operation)
Conflicts remove(AsyncOperation<?> operation, boolean allowOutOfOrder)
{
if (operationOrQueue instanceof AsyncOperation<?>)
{
Invariants.checkState(operationOrQueue == operation);
rangeQueues.remove(range);
}
else
{
@SuppressWarnings("unchecked")
ArrayDeque<AsyncOperation<?>> queue = (ArrayDeque<AsyncOperation<?>>) operationOrQueue;
AsyncOperation<?> head = queue.poll();
Invariants.checkState(head == operation);
if (queue.isEmpty())
{
rangeQueues.remove(range);
}
else
{
head = queue.peek();
if (canRun(head))
head.onUnblocked();
}
}
unregister("range", range, operationOrQueue, operation, allowOutOfOrder, () -> rangeQueues.remove(range));
return operationToConflicts.remove(operation);
}
@ -182,20 +161,12 @@ public class ExecutionOrder
result.rangeConflicts.add(e.getKey());
}
RangeState state = e.getValue();
Object operationOrQueue = state.operationOrQueue;
if (operationOrQueue instanceof AsyncOperation)
{
ArrayDeque<AsyncOperation<?>> queue = new ArrayDeque<>(4);
queue.add((AsyncOperation<?>) operationOrQueue);
queue.add(operation);
state.operationOrQueue = queue;
}
else
{
@SuppressWarnings("unchecked")
ArrayDeque<AsyncOperation<?>> queue = (ArrayDeque<AsyncOperation<?>>) operationOrQueue;
queue.add(operation);
}
// a single range could conflict with multiple other ranges, so it is possible that the operation
// exists in the queue already due to another range in the txn... simple example is
// keys = (0, 10], (12, 15]
// e.getKey() == (-100, 100]
// in this case the operation would attempt to double add since it has 2 keys that conflict with this single range
register(state.operationOrQueue, operation, q -> state.operationOrQueue = q);
});
if (result.sameRange != null)
{
@ -205,7 +176,7 @@ public class ExecutionOrder
{
rangeQueues.add(range, new RangeState(range, keyConflicts, result.rangeConflicts, operation));
}
return keyConflicts == null && result.rangeConflicts == null;
return keyConflicts == null && result.rangeConflicts == null && result.sameRange == null;
}
/**
@ -221,12 +192,19 @@ public class ExecutionOrder
return true;
}
register(operationOrQueue, operation, q -> queues.put(keyOrTxnId, q));
return false;
}
private void register(Object operationOrQueue, AsyncOperation<?> operation, Consumer<ArrayDeque<AsyncOperation<?>>> onCreateQueue)
{
if (operationOrQueue instanceof AsyncOperation)
{
Invariants.checkState(operationOrQueue != operation, "Attempted to double register operation %s", operation);
ArrayDeque<AsyncOperation<?>> queue = new ArrayDeque<>(4);
queue.add((AsyncOperation<?>) operationOrQueue);
queue.add(operation);
queues.put(keyOrTxnId, queue);
onCreateQueue.accept(queue);
}
else
{
@ -234,23 +212,35 @@ public class ExecutionOrder
ArrayDeque<AsyncOperation<?>> queue = (ArrayDeque<AsyncOperation<?>>) operationOrQueue;
queue.add(operation);
}
return false;
}
/**
* Unregister the operation as being a dependency for its keys and TxnIds, but do so even if it is unable to run now.
*/
void unregisterOutOfOrder(AsyncOperation<?> operation)
{
unregister(operation, true);
}
/**
* Unregister the operation as being a dependency for its keys and TxnIds
*/
void unregister(AsyncOperation<?> operation)
{
unregister(operation, false);
}
private void unregister(AsyncOperation<?> operation, boolean allowOutOfOrder)
{
for (Seekable seekable : operation.keys())
{
switch (seekable.domain())
{
case Key:
unregister(seekable.asKey(), operation);
unregister(seekable.asKey(), operation, allowOutOfOrder);
break;
case Range:
unregister(seekable.asRange(), operation);
unregister(seekable.asRange(), operation, allowOutOfOrder);
break;
default:
throw new AssertionError("Unexpected domain: " + seekable.domain());
@ -259,48 +249,69 @@ public class ExecutionOrder
}
TxnId primaryTxnId = operation.primaryTxnId();
if (null != primaryTxnId)
unregister(primaryTxnId, operation);
unregister(primaryTxnId, operation, allowOutOfOrder);
}
private void unregister(Range range, AsyncOperation<?> operation)
private void unregister(Range range, AsyncOperation<?> operation, boolean allowOutOfOrder)
{
RangeState state = state(range);
Conflicts conflicts = state.remove(operation);
Conflicts conflicts = state.remove(operation, allowOutOfOrder);
if (conflicts.rangeConflicts != null)
conflicts.rangeConflicts.forEach(r -> state(r).remove(operation));
conflicts.rangeConflicts.forEach(r -> state(r).remove(operation, allowOutOfOrder));
if (conflicts.keyConflicts != null)
conflicts.keyConflicts.forEach(k -> unregister(k, operation));
conflicts.keyConflicts.forEach(k -> unregister(k, operation, allowOutOfOrder));
}
/**
* Unregister the operation as being a dependency for key or TxnId
*/
private void unregister(Object keyOrTxnId, AsyncOperation<?> operation)
private void unregister(Object keyOrTxnId, AsyncOperation<?> operation, boolean allowOutOfOrder)
{
Object operationOrQueue = queues.get(keyOrTxnId);
Invariants.nonNull(operationOrQueue);
unregister("Key or TxnId", keyOrTxnId, operationOrQueue, operation, allowOutOfOrder, () -> queues.remove(keyOrTxnId));
}
private void unregister(String name, Object key, Object operationOrQueue, AsyncOperation<?> operation, boolean allowOutOfOrder, Runnable onEmpty)
{
if (operationOrQueue instanceof AsyncOperation<?>)
{
Invariants.checkState(operationOrQueue == operation);
queues.remove(keyOrTxnId);
Invariants.checkState(operationOrQueue == operation, "Only single operation present and was not %s; %s %s", name, key);
onEmpty.run();
}
else
{
@SuppressWarnings("unchecked")
ArrayDeque<AsyncOperation<?>> queue = (ArrayDeque<AsyncOperation<?>>) operationOrQueue;
AsyncOperation<?> head = queue.poll();
Invariants.checkState(head == operation);
if (queue.isEmpty())
if (allowOutOfOrder)
{
queues.remove(keyOrTxnId);
Invariants.checkState(queue.remove(operation), "Operation %s was not found in queue: %s; %s %s", operation, queue, name, key);
}
else
{
head = queue.peek();
if (canRun(head))
head.onUnblocked();
Invariants.checkState(queue.peek() == operation, "Operation %s is not at the top of the queue; %s; %s %s", operation, queue, name, key);
queue.poll();
}
if (queue.isEmpty())
{
onEmpty.run();
}
else
{
AsyncOperation<?> next = queue.peek();
if (next == operation)
{
// a single range could conflict with multiple other ranges, so it is possible that the operation
// exists in the queue already due to another range in the txn... simple example is
// keys = (0, 10], (12, 15]
// e.getKey() == (-100, 100]
// in this case the operation would attempt to double add since it has 2 keys that conflict with this single range
return;
}
if (canRun(next))
next.onUnblocked();
}
}
}
@ -357,7 +368,7 @@ public class ExecutionOrder
private RangeState state(Range range)
{
List<RangeState> list = rangeQueues.get(range);
assert list.size() == 1 : String.format("Expected 1 element but saw list %s", list);
assert list.size() == 1 : String.format("Expected 1 element for range %s but saw list %s", range, list);
return list.get(0);
}

View File

@ -121,7 +121,7 @@ public class ConsensusRequestRouter
ClusterMetadata cm = ClusterMetadata.current();
TableMetadata metadata = cm.schema.getTableMetadata(tableId);
if (metadata == null)
throw new IllegalStateException("Can't route consensus request for nonexistent table %s".format(tableId.toString()));
throw new IllegalStateException(String.format("Can't route consensus request for nonexistent table %s", tableId));
if (!mayWriteThroughAccord(metadata))
return false;

View File

@ -108,14 +108,14 @@ public class SessionSummary
List<StreamSummary> receivingSummaries = new ArrayList<>(numRcvd);
for (int i=0; i<numRcvd; i++)
{
receivingSummaries.add(StreamSummary.serializer.deserialize(in, partitioner, version));
receivingSummaries.add(StreamSummary.serializer.deserialize(in, version));
}
int numSent = in.readInt();
List<StreamSummary> sendingSummaries = new ArrayList<>(numRcvd);
for (int i=0; i<numSent; i++)
{
sendingSummaries.add(StreamSummary.serializer.deserialize(in, partitioner, version));
sendingSummaries.add(StreamSummary.serializer.deserialize(in, version));
}
return new SessionSummary(coordinator, peer, receivingSummaries, sendingSummaries);

View File

@ -22,7 +22,6 @@ import com.google.common.annotations.VisibleForTesting;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.db.guardrails.GuardrailViolatedException;
import org.apache.cassandra.db.guardrails.Guardrails;
import org.apache.cassandra.locator.InetAddressAndPort;
@ -58,7 +57,7 @@ public class StreamDeserializingTask implements Runnable
try
{
StreamMessage message;
while (null != (message = StreamMessage.deserialize(input, IPartitioner.global(), messagingVersion)))
while (null != (message = StreamMessage.deserialize(input, messagingVersion)))
{
// keep-alives don't necessarily need to be tied to a session (they could be arrive before or after
// wrt session lifecycle, due to races), just log that we received the message and carry on

View File

@ -23,16 +23,17 @@ import java.util.List;
import com.google.common.base.Objects;
import com.google.common.collect.ImmutableList;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.dht.IPartitionerDependentSerializer;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.CollectionSerializers;
/**
@ -40,7 +41,7 @@ import org.apache.cassandra.utils.CollectionSerializers;
*/
public class StreamSummary implements Serializable
{
public static final IPartitionerDependentSerializer<StreamSummary> serializer = new StreamSummarySerializer();
public static final IVersionedSerializer<StreamSummary> serializer = new StreamSummarySerializer();
public final TableId tableId;
public final List<Range<Token>> ranges;
@ -86,25 +87,34 @@ public class StreamSummary implements Serializable
return sb.toString();
}
public static class StreamSummarySerializer implements IPartitionerDependentSerializer<StreamSummary>
public static class StreamSummarySerializer implements IVersionedSerializer<StreamSummary>
{
public void serialize(StreamSummary summary, DataOutputPlus out, int version) throws IOException
{
summary.tableId.serialize(out);
out.writeInt(summary.files);
out.writeLong(summary.totalSize);
Token.logPartitioner = true;
if (version >= MessagingService.VERSION_51)
CollectionSerializers.serializeCollection(summary.ranges, out, version, Range.rangeSerializer);
Token.logPartitioner = false;
}
public StreamSummary deserialize(DataInputPlus in, IPartitioner p, int version) throws IOException
public StreamSummary deserialize(DataInputPlus in, int version) throws IOException
{
TableId tableId = TableId.deserialize(in);
int files = in.readInt();
long totalSize = in.readLong();
List<Range<Token>> ranges = ImmutableList.of();
if (version >= MessagingService.VERSION_51)
{
TableMetadata tableMetadata = Schema.instance.getTableMetadata(tableId);
IPartitioner p = tableMetadata != null ? tableMetadata.partitioner : IPartitioner.global();
Token.logPartitioner = true;
ranges = CollectionSerializers.deserializeList(in, p, version, Range.rangeSerializer);
Token.logPartitioner = false;
}
return new StreamSummary(tableId, ranges, files, totalSize);
}

View File

@ -17,7 +17,6 @@
*/
package org.apache.cassandra.streaming.messages;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingDataOutputPlus;
@ -26,7 +25,7 @@ public class CompleteMessage extends StreamMessage
{
public static Serializer<CompleteMessage> serializer = new Serializer<CompleteMessage>()
{
public CompleteMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version)
public CompleteMessage deserialize(DataInputPlus in, int version)
{
return new CompleteMessage();
}

View File

@ -21,7 +21,6 @@ import java.io.IOException;
import java.util.Objects;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.IncomingStream;
import org.apache.cassandra.streaming.StreamManager;
@ -34,7 +33,7 @@ public class IncomingStreamMessage extends StreamMessage
{
public static Serializer<IncomingStreamMessage> serializer = new Serializer<IncomingStreamMessage>()
{
public IncomingStreamMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException
public IncomingStreamMessage deserialize(DataInputPlus input, int version) throws IOException
{
StreamMessageHeader header = StreamMessageHeader.serializer.deserialize(input, version);
StreamSession session = StreamManager.instance.findSession(header.sender, header.planId, header.sessionIndex, header.sendByFollower);

View File

@ -17,7 +17,6 @@
*/
package org.apache.cassandra.streaming.messages;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingDataOutputPlus;
@ -38,7 +37,7 @@ public class KeepAliveMessage extends StreamMessage
public static Serializer<KeepAliveMessage> serializer = new Serializer<KeepAliveMessage>()
{
public KeepAliveMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version)
public KeepAliveMessage deserialize(DataInputPlus in, int version)
{
return new KeepAliveMessage();
}

View File

@ -21,7 +21,6 @@ import java.io.IOException;
import com.google.common.annotations.VisibleForTesting;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.streaming.OutgoingStream;
@ -33,7 +32,7 @@ public class OutgoingStreamMessage extends StreamMessage
{
public static Serializer<OutgoingStreamMessage> serializer = new Serializer<OutgoingStreamMessage>()
{
public OutgoingStreamMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version)
public OutgoingStreamMessage deserialize(DataInputPlus in, int version)
{
throw new UnsupportedOperationException("Not allowed to call deserialize on an outgoing stream");
}

View File

@ -20,7 +20,6 @@ package org.apache.cassandra.streaming.messages;
import java.io.IOException;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingDataOutputPlus;
@ -34,7 +33,7 @@ public class PrepareAckMessage extends StreamMessage
//nop
}
public PrepareAckMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException
public PrepareAckMessage deserialize(DataInputPlus in, int version) throws IOException
{
return new PrepareAckMessage();
}

View File

@ -22,7 +22,6 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.Collection;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamSummary;
@ -39,12 +38,12 @@ public class PrepareSynAckMessage extends StreamMessage
StreamSummary.serializer.serialize(summary, out, version);
}
public PrepareSynAckMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException
public PrepareSynAckMessage deserialize(DataInputPlus input, int version) throws IOException
{
PrepareSynAckMessage message = new PrepareSynAckMessage();
int numSummaries = input.readInt();
for (int i = 0; i < numSummaries; i++)
message.summaries.add(StreamSummary.serializer.deserialize(input, partitioner, version));
message.summaries.add(StreamSummary.serializer.deserialize(input, version));
return message;
}

View File

@ -21,7 +21,6 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.Collection;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamRequest;
import org.apache.cassandra.streaming.StreamSession;
@ -32,7 +31,7 @@ public class PrepareSynMessage extends StreamMessage
{
public static Serializer<PrepareSynMessage> serializer = new Serializer<PrepareSynMessage>()
{
public PrepareSynMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException
public PrepareSynMessage deserialize(DataInputPlus input, int version) throws IOException
{
PrepareSynMessage message = new PrepareSynMessage();
// requests
@ -42,7 +41,7 @@ public class PrepareSynMessage extends StreamMessage
// summaries
int numSummaries = input.readInt();
for (int i = 0; i < numSummaries; i++)
message.summaries.add(StreamSummary.serializer.deserialize(input, partitioner, version));
message.summaries.add(StreamSummary.serializer.deserialize(input, version));
return message;
}

View File

@ -19,7 +19,6 @@ package org.apache.cassandra.streaming.messages;
import java.io.IOException;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.streaming.StreamSession;
@ -29,7 +28,7 @@ public class ReceivedMessage extends StreamMessage
{
public static Serializer<ReceivedMessage> serializer = new Serializer<ReceivedMessage>()
{
public ReceivedMessage deserialize(DataInputPlus input, IPartitioner partitioner, int version) throws IOException
public ReceivedMessage deserialize(DataInputPlus input, int version) throws IOException
{
return new ReceivedMessage(TableId.deserialize(input), input.readInt());
}

View File

@ -17,7 +17,6 @@
*/
package org.apache.cassandra.streaming.messages;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingDataOutputPlus;
@ -26,7 +25,7 @@ public class SessionFailedMessage extends StreamMessage
{
public static Serializer<SessionFailedMessage> serializer = new Serializer<SessionFailedMessage>()
{
public SessionFailedMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version)
public SessionFailedMessage deserialize(DataInputPlus in, int version)
{
return new SessionFailedMessage();
}

View File

@ -20,7 +20,6 @@ package org.apache.cassandra.streaming.messages;
import java.io.IOException;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.streaming.PreviewKind;
@ -94,7 +93,7 @@ public class StreamInitMessage extends StreamMessage
out.writeInt(message.previewKind.getSerializationVal());
}
public StreamInitMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException
public StreamInitMessage deserialize(DataInputPlus in, int version) throws IOException
{
InetAddressAndPort from = inetAddressAndPortSerializer.deserialize(in, version);
int sessionIndex = in.readInt();

View File

@ -21,7 +21,6 @@ import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingChannel;
@ -45,16 +44,16 @@ public abstract class StreamMessage
return 1 + message.type.outSerializer.serializedSize(message, version);
}
public static StreamMessage deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException
public static StreamMessage deserialize(DataInputPlus in, int version) throws IOException
{
Type type = Type.lookupById(in.readByte());
return type.inSerializer.deserialize(in, partitioner, version);
return type.inSerializer.deserialize(in, version);
}
/** StreamMessage serializer */
public static interface Serializer<V extends StreamMessage>
{
V deserialize(DataInputPlus in, IPartitioner partitioner, int version) throws IOException;
V deserialize(DataInputPlus in, int version) throws IOException;
void serialize(V message, StreamingDataOutputPlus out, int version, StreamSession session) throws IOException;
long serializedSize(V message, int version) throws IOException;
}

View File

@ -28,6 +28,7 @@ import java.util.stream.Collectors;
import org.junit.Test;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.distributed.Cluster;
import org.apache.cassandra.distributed.api.IInvokableInstance;
import org.apache.cassandra.distributed.shared.ClusterUtils;
@ -63,6 +64,7 @@ public class RepairMetadataKeyspaceTest extends TestBaseImpl
IInvokableInstance toRepair = cluster.get(3);
stopUnchecked(toRepair);
DatabaseDescriptor.clientInitialization();
String targetDir = DistributedMetadataLogKeyspace.TABLE_NAME + '-' + DistributedMetadataLogKeyspace.LOG_TABLE_ID.toHexString();
for (File datadir : getDataDirectories(toRepair))
{

View File

@ -26,6 +26,7 @@ import javax.annotation.Nullable;
import com.google.common.collect.ImmutableMap;
import accord.topology.TopologyUtils;
import org.apache.cassandra.config.AccordSpec;
import org.apache.cassandra.schema.*;
import org.junit.Ignore;
import org.junit.Test;
@ -166,7 +167,7 @@ public class AccordJournalSimulationTest extends SimulationTestBase
}
}
private static final ExecutorPlus executor = ExecutorFactory.Global.executorFactory().pooled("name", 10);
private static final AccordJournal journal = new AccordJournal(null);
private static final AccordJournal journal = new AccordJournal(null, new AccordSpec.JournalSpec());
private static final int events = 100;
private static final CountDownLatch eventsWritten = CountDownLatch.newCountDownLatch(events);
private static final CountDownLatch eventsDurable = CountDownLatch.newCountDownLatch(events);

View File

@ -78,6 +78,7 @@ public class DatabaseDescriptorRefTest
"org.apache.cassandra.auth.INetworkAuthorizer",
"org.apache.cassandra.auth.IRoleManager",
"org.apache.cassandra.config.AccordSpec",
"org.apache.cassandra.config.AccordSpec$JournalSpec",
"org.apache.cassandra.config.AccordSpec$TransactionalRangeMigration",
"org.apache.cassandra.config.CassandraRelevantProperties",
"org.apache.cassandra.config.CassandraRelevantProperties$PropertyConverter",
@ -275,6 +276,9 @@ public class DatabaseDescriptorRefTest
"org.apache.cassandra.io.util.PathUtils$IOToLongFunction",
"org.apache.cassandra.io.util.RebufferingInputStream",
"org.apache.cassandra.io.util.SpinningDiskOptimizationStrategy",
"org.apache.cassandra.journal.Params",
"org.apache.cassandra.journal.Params$FailurePolicy",
"org.apache.cassandra.journal.Params$FlushMode",
"org.apache.cassandra.locator.Endpoint",
"org.apache.cassandra.locator.IEndpointSnitch",
"org.apache.cassandra.locator.InetAddressAndPort",

View File

@ -91,6 +91,7 @@ import org.apache.cassandra.service.accord.api.PartitionKey;
import org.apache.cassandra.service.accord.serializers.CommandsForKeySerializer;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.Pair;
import org.assertj.core.api.Assertions;
import static accord.impl.TimestampsForKey.NO_LAST_EXECUTED_HLC;
import static accord.local.KeyHistory.COMMANDS;
@ -369,9 +370,14 @@ public class CompactionAccordIteratorsTest
assertEquals(1, Iterators.size(partition.unfilteredIterator()));
ByteBuffer[] partitionKeyComponents = CommandRows.splitPartitionKey(partition.partitionKey());
Row row = (Row)partition.unfilteredIterator().next();
assertEquals(commands.metadata().regularColumns().size(), row.columnCount());
// execute_atleast is null, so when we read from the scanner the column won't be present in the partition
Assertions.assertThat(new ArrayList<>(row.columns())).isEqualTo(commands.metadata().regularColumns().stream().filter(c -> !c.name.toString().equals("execute_atleast")).collect(Collectors.toList()));
for (ColumnMetadata cm : commands.metadata().regularColumns())
{
if (cm.name.toString().equals("execute_atleast")) continue;
assertNotNull(row.getColumnData(cm));
}
assertEquals(TXN_ID, CommandRows.getTxnId(partitionKeyComponents));
assertEquals(SaveStatus.Applied, AccordKeyspace.CommandRows.getStatus(row));
};

View File

@ -25,6 +25,8 @@ import java.util.List;
import java.util.UUID;
import com.google.common.collect.Lists;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.dht.*;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@ -32,10 +34,7 @@ import org.junit.Test;
import org.apache.cassandra.CassandraTestBase;
import org.apache.cassandra.CassandraTestBase.DDDaemonInitialization;
import org.apache.cassandra.CassandraTestBase.UseMurmur3Partitioner;
import org.apache.cassandra.dht.Murmur3Partitioner;
import org.apache.cassandra.dht.Murmur3Partitioner.LongToken;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputBuffer;
@ -116,7 +115,19 @@ public class RepairMessageSerializationsTest extends CassandraTestBase
buf.flip();
DataInputPlus in = new DataInputBuffer(buf, false);
T deserialized = serializer.deserialize(in, PROTOCOL_VERSION);
T deserialized = null;
if (serializer instanceof IPartitionerDependentSerializer)
{
IPartitionerDependentSerializer<T> pds = (IPartitionerDependentSerializer<T>) serializer;
deserialized = pds.deserialize(in, DatabaseDescriptor.getPartitioner(), PROTOCOL_VERSION);
}
else
{
deserialized = serializer.deserialize(in, PROTOCOL_VERSION);
}
Assert.assertEquals(msg, deserialized);
Assert.assertEquals(msg.hashCode(), deserialized.hashCode());
return deserialized;

View File

@ -23,7 +23,6 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import com.google.common.collect.Lists;
import org.junit.AfterClass;
@ -60,6 +59,7 @@ import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.streaming.PreviewKind;
import org.apache.cassandra.streaming.SessionSummary;
import org.apache.cassandra.streaming.StreamSummary;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.utils.Clock;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.MerkleTrees;
@ -70,6 +70,7 @@ import static java.util.Collections.emptyList;
public class SerializationsTest extends AbstractSerializationsTester
{
private static PartitionerSwitcher partitionerSwitcher;
private static TableId TABLE_ID;
private static TimeUUID RANDOM_UUID;
private static Range<Token> FULL_RANGE;
private static RepairJobDesc DESC;
@ -84,6 +85,7 @@ public class SerializationsTest extends AbstractSerializationsTester
ClusterMetadataTestHelper.setInstanceForTest();
SchemaTestUtil.addOrUpdateKeyspace(KeyspaceMetadata.create("Keyspace1", KeyspaceParams.simple(3)));
SchemaTestUtil.announceNewTable(TableMetadata.minimal("Keyspace1", "Standard1"));
TABLE_ID = ClusterMetadata.current().schema.getKeyspaceMetadata("Keyspace1").getTableOrViewNullable("Standard1").id();
RANDOM_UUID = TimeUUID.fromString("743325d0-4c4b-11ec-8a88-2d67081686db");
FULL_RANGE = new Range<>(Util.testPartitioner().getMinimumToken(), Util.testPartitioner().getMinimumToken());
DESC = new RepairJobDesc(RANDOM_UUID, RANDOM_UUID, "Keyspace1", "Standard1", Arrays.asList(FULL_RANGE));
@ -223,8 +225,8 @@ public class SerializationsTest extends AbstractSerializationsTester
// sync success
List<SessionSummary> summaries = new ArrayList<>();
summaries.add(new SessionSummary(src, dest,
Lists.newArrayList(new StreamSummary(TableId.fromUUID(UUID.randomUUID()), emptyList(), 5, 100)),
Lists.newArrayList(new StreamSummary(TableId.fromUUID(UUID.randomUUID()), emptyList(), 500, 10))
Lists.newArrayList(new StreamSummary(TABLE_ID, emptyList(), 5, 100)),
Lists.newArrayList(new StreamSummary(TABLE_ID, emptyList(), 500, 10))
));
SyncResponse success = new SyncResponse(DESC, src, dest, true, summaries);
// sync fail

View File

@ -76,6 +76,7 @@ import org.apache.cassandra.concurrent.ExecutorPlus;
import org.apache.cassandra.concurrent.ImmediateExecutor;
import org.apache.cassandra.concurrent.ManualExecutor;
import org.apache.cassandra.concurrent.Stage;
import org.apache.cassandra.config.AccordSpec;
import org.apache.cassandra.cql3.QueryOptions;
import org.apache.cassandra.cql3.QueryProcessor;
import org.apache.cassandra.cql3.statements.TransactionStatement;
@ -399,7 +400,7 @@ public class AccordTestUtils
public long unix(TimeUnit timeUnit) { return NodeTimeService.unixWrapper(TimeUnit.MICROSECONDS, this::now).applyAsLong(timeUnit); }
};
AccordJournal journal = new AccordJournal(null);
AccordJournal journal = new AccordJournal(null, new AccordSpec.JournalSpec());
journal.start(null);
SingleEpochRanges holder = new SingleEpochRanges(topology.rangesForNode(node));

View File

@ -28,6 +28,7 @@ import com.google.common.collect.Sets;
import accord.local.SerializerSupport;
import accord.messages.Accept;
import accord.messages.Apply;
import accord.messages.ApplyThenWaitUntilApplied;
import accord.messages.BeginRecovery;
import accord.messages.Commit;
import accord.messages.Message;
@ -43,15 +44,18 @@ import org.apache.cassandra.service.accord.AccordJournal.Type;
import static accord.messages.MessageType.ACCEPT_REQ;
import static accord.messages.MessageType.APPLY_MAXIMAL_REQ;
import static accord.messages.MessageType.APPLY_MINIMAL_REQ;
import static accord.messages.MessageType.APPLY_THEN_WAIT_UNTIL_APPLIED_REQ;
import static accord.messages.MessageType.BEGIN_RECOVER_REQ;
import static accord.messages.MessageType.COMMIT_MAXIMAL_REQ;
import static accord.messages.MessageType.COMMIT_SLOW_PATH_REQ;
import static accord.messages.MessageType.PRE_ACCEPT_REQ;
import static accord.messages.MessageType.PROPAGATE_APPLY_MSG;
import static accord.messages.MessageType.PROPAGATE_OTHER_MSG;
import static accord.messages.MessageType.PROPAGATE_PRE_ACCEPT_MSG;
import static accord.messages.MessageType.PROPAGATE_STABLE_MSG;
import static accord.messages.MessageType.STABLE_FAST_PATH_REQ;
import static accord.messages.MessageType.STABLE_MAXIMAL_REQ;
import static accord.messages.MessageType.STABLE_SLOW_PATH_REQ;
public class MockJournal implements IJournal
{
@ -61,6 +65,12 @@ public class MockJournal implements IJournal
{
return new SerializerSupport.MessageProvider()
{
@Override
public TxnId txnId()
{
return txnId;
}
@Override
public Set<MessageType> test(Set<MessageType> messages)
{
@ -146,6 +156,12 @@ public class MockJournal implements IJournal
return get(STABLE_FAST_PATH_REQ);
}
@Override
public Commit stableSlowPath()
{
return get(STABLE_SLOW_PATH_REQ);
}
@Override
public Commit stableMaximal()
{
@ -175,6 +191,18 @@ public class MockJournal implements IJournal
{
return get(PROPAGATE_APPLY_MSG);
}
@Override
public Propagate propagateOther()
{
return get(PROPAGATE_OTHER_MSG);
}
@Override
public ApplyThenWaitUntilApplied applyThenWaitUntilApplied()
{
return get(APPLY_THEN_WAIT_UNTIL_APPLIED_REQ);
}
};
}

View File

@ -25,6 +25,7 @@ import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.function.BooleanSupplier;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.function.ToLongFunction;
import accord.impl.SizeOfIntersectionSorter;
@ -76,7 +77,7 @@ import static org.apache.cassandra.db.ColumnFamilyStore.FlushReason.UNIT_TESTS;
import static org.apache.cassandra.schema.SchemaConstants.ACCORD_KEYSPACE_NAME;
import static org.apache.cassandra.utils.AccordGenerators.fromQT;
class SimulatedAccordCommandStore implements AutoCloseable
public class SimulatedAccordCommandStore implements AutoCloseable
{
private final List<Throwable> failures = new ArrayList<>();
private final SimulatedExecutorFactory globalExecutor;
@ -90,8 +91,9 @@ class SimulatedAccordCommandStore implements AutoCloseable
public final MockJournal journal;
public final ScheduledExecutorPlus unorderedScheduled;
public final List<String> evictions = new ArrayList<>();
public Predicate<Throwable> ignoreExceptions = ignore -> false;
SimulatedAccordCommandStore(RandomSource rs)
public SimulatedAccordCommandStore(RandomSource rs)
{
globalExecutor = new SimulatedExecutorFactory(accord.utilsfork.RandomSource.wrap(rs).fork(), fromQT(Generators.TIMESTAMP_GEN.map(java.sql.Timestamp::getTime)).mapToLong(TimeUnit.MILLISECONDS::toNanos).next(rs), failures::add);
this.unorderedScheduled = globalExecutor.scheduled("ignored");
@ -151,6 +153,13 @@ class SimulatedAccordCommandStore implements AutoCloseable
{
return false;
}
@Override
public void onUncaughtException(Throwable t)
{
if (ignoreExceptions.test(t)) return;
super.onUncaughtException(t);
}
},
null,
ignore -> AccordTestUtils.NOOP_PROGRESS_LOG,

View File

@ -0,0 +1,207 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.cassandra.service.accord.async;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.BooleanSupplier;
import org.junit.Before;
import org.junit.Test;
import accord.api.Key;
import accord.impl.basic.SimulatedFault;
import accord.local.PreLoadContext;
import accord.local.SafeCommandStore;
import accord.primitives.Keys;
import accord.primitives.Range;
import accord.primitives.Ranges;
import accord.primitives.Seekables;
import accord.utils.Gen;
import accord.utils.Gens;
import accord.utils.RandomSource;
import org.apache.cassandra.dht.Murmur3Partitioner;
import org.apache.cassandra.dht.Murmur3Partitioner.LongToken;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.service.accord.AccordCommandStore;
import org.apache.cassandra.service.accord.AccordKeyspace;
import org.apache.cassandra.service.accord.SimulatedAccordCommandStore;
import org.apache.cassandra.service.accord.SimulatedAccordCommandStoreTestBase;
import org.apache.cassandra.service.accord.TokenRange;
import org.apache.cassandra.service.accord.api.AccordRoutingKey.TokenKey;
import org.apache.cassandra.service.accord.api.PartitionKey;
import org.assertj.core.api.Assertions;
import static accord.utils.Property.qt;
public class SimulatedAsyncOperationTest extends SimulatedAccordCommandStoreTestBase
{
@Before
public void precondition()
{
Assertions.assertThat(intTbl.partitioner).isEqualTo(Murmur3Partitioner.instance);
}
@Test
public void happyPath()
{
qt().withExamples(100).check(rs -> test(rs, 100, intTbl, ignore -> Action.SUCCESS));
}
@Test
public void fuzz()
{
Gen<Action> actionGen = Gens.enums().allWithWeights(Action.class, 10, 1, 1);
qt().withExamples(100).check(rs -> test(rs, 100, intTbl, actionGen));
}
private static void test(RandomSource rs, int numSamples, TableMetadata tbl, Gen<Action> actionGen) throws Exception
{
AccordKeyspace.unsafeClear();
int numKeys = rs.nextInt(20, 1000);
long minToken = 0;
long maxToken = numKeys;
Gen<Key> keyGen = Gens.longs().between(minToken + 1, maxToken).map(t -> new PartitionKey(tbl.id, tbl.partitioner.decorateKey(LongToken.keyForToken(t))));
Gen<Keys> keysGen = Gens.lists(keyGen).unique().ofSizeBetween(1, 10).map(l -> Keys.of(l));
Gen<Ranges> rangesGen = Gens.lists(rangeInsideRange(tbl.id, minToken, maxToken)).uniqueBestEffort().ofSizeBetween(1, 10).map(l -> Ranges.of(l.toArray(Range[]::new)));
Gen<Seekables<?, ?>> seekablesGen = Gens.oneOf(keysGen, rangesGen);
try (var instance = new SimulatedAccordCommandStore(rs))
{
instance.ignoreExceptions = t -> t instanceof SimulatedFault;
Counter counter = new Counter();
for (int i = 0; i < numSamples; i++)
{
PreLoadContext ctx = PreLoadContext.contextFor(seekablesGen.next(rs));
operation(instance, ctx, actionGen.next(rs), rs::nextBoolean).begin((ignore, failure) -> {
counter.counter++;
if (failure != null && !(failure instanceof SimulatedFault)) throw new AssertionError("Unexpected error", failure);
});
}
instance.processAll();
Assertions.assertThat(counter.counter).isEqualTo(numSamples);
}
}
private static Gen<Range> rangeInsideRange(TableId tableId, long minToken, long maxToken)
{
if (minToken + 1 == maxToken)
{
// only one range is possible...
return Gens.constant(range(tableId, minToken, maxToken));
}
return rs -> {
long a = rs.nextLong(minToken, maxToken + 1);
long b = rs.nextLong(minToken, maxToken + 1);
while (a == b)
b = rs.nextLong(minToken, maxToken + 1);
if (a > b)
{
long tmp = a;
a = b;
b = tmp;
}
return range(tableId, a, b);
};
}
private static TokenRange range(TableId tableId, long start, long end)
{
return new TokenRange(new TokenKey(tableId, new LongToken(start)), new TokenKey(tableId, new LongToken(end)));
}
private enum Action {SUCCESS, FAILURE, LOAD_FAILURE}
private static AsyncOperation<Void> operation(SimulatedAccordCommandStore instance, PreLoadContext ctx, Action action, BooleanSupplier delay)
{
return new SimulatedOperation(instance.store, ctx, action == Action.FAILURE ? SimulatedOperation.Action.FAILURE : SimulatedOperation.Action.SUCCESS)
{
@Override
AsyncLoader createAsyncLoader(AccordCommandStore commandStore, PreLoadContext preLoadContext)
{
return new SimulatedLoader(action == SimulatedAsyncOperationTest.Action.LOAD_FAILURE ? SimulatedLoader.Action.FAILURE : SimulatedLoader.Action.SUCCESS, delay.getAsBoolean(), instance.unorderedScheduled);
}
};
}
private static class Counter
{
int counter = 0;
}
private static class SimulatedOperation extends AsyncOperation<Void>
{
enum Action { SUCCESS, FAILURE}
private final Action action;
public SimulatedOperation(AccordCommandStore commandStore, PreLoadContext preLoadContext, Action action)
{
super(commandStore, preLoadContext);
this.action = action;
}
@Override
public Void apply(SafeCommandStore safe)
{
if (action == Action.FAILURE)
throw new SimulatedFault("Operation failed for keys " + keys());
return null;
}
}
private static class SimulatedLoader extends AsyncLoader
{
enum Action { SUCCESS, FAILURE}
private final Action action;
private boolean delay;
private final ScheduledExecutorService executor;
SimulatedLoader(Action action, boolean delay, ScheduledExecutorService executor)
{
super(null, null, null, null);
this.action = action;
this.delay = delay;
this.executor = executor;
}
@Override
public boolean load(AsyncOperation.Context context, BiConsumer<Object, Throwable> callback)
{
if (delay)
{
executor.schedule(() -> {
callback.accept(null, action == Action.FAILURE ? new SimulatedFault("Failure loading " + context) : null);
}, 1, TimeUnit.SECONDS);
delay = false;
return false;
}
if (action == Action.FAILURE)
throw new SimulatedFault("Failure loading " + context);
return true;
}
}
}

View File

@ -30,7 +30,6 @@ import org.junit.Test;
import io.netty.buffer.ByteBuf;
import io.netty.channel.embedded.EmbeddedChannel;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputBuffer;
@ -56,7 +55,6 @@ import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID;
public class StreamingInboundHandlerTest
{
private NettyStreamingChannel streamingChannel;
private EmbeddedChannel channel;
private ByteBuf buf;
@ -125,7 +123,7 @@ public class StreamingInboundHandlerTest
temp.flip();
DataInputPlus in = new DataInputBuffer(temp, false);
// session not found
IncomingStreamMessage.serializer.deserialize(in, IPartitioner.global(), MessagingService.current_version);
IncomingStreamMessage.serializer.deserialize(in, MessagingService.current_version);
}
@Test

View File

@ -20,6 +20,7 @@ package org.apache.cassandra.tcm.listeners;
import java.util.Random;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.Test;
@ -50,7 +51,7 @@ import static org.junit.Assert.assertNull;
public class MetadataSnapshotListenerTest
{
private static final Logger logger = LoggerFactory.getLogger(MetadataSnapshotListenerTest.class);
private IPartitioner partitioner = Murmur3Partitioner.instance;
private final IPartitioner partitioner = Murmur3Partitioner.instance;
private Random r;
@BeforeClass
@ -59,6 +60,7 @@ public class MetadataSnapshotListenerTest
// Set this so that we don't attempt to sort the random placements as this depends on a populated
// TokenMap. This is a temporary element of ClusterMetadata, at least in the current form
CassandraRelevantProperties.TCM_SORT_REPLICA_GROUPS.setBoolean(false);
DatabaseDescriptor.daemonInitialization();
}
@Before