Fix CFK restore after replay:

- Remove from CFK any unapplied transactions we know cannot apply
 - Force notification of waiting commands in CFK on replay
Also fix:
 - Don't truncateWithOutcome if pre bootstrap or stale
 - Fix RangeDeps.without(RangeDeps)
 - Fix InMemoryCommandStore replay bug with clearing DefaultLocalListeners
 - Ensure SaveStatus and executeAt are updated together to prevent corruption via expunge
Improve:
 - Inform home shard that command is decided if we cannot execute, to avoid recovery contention
 - Don't recover sync points on the fast path
 - Don't calculate recovery info for RX (since we don't use it anymore, so no need to do the work)

patch by Benedict; reviewed by Alex Petrov for CASSANDRA-20797
This commit is contained in:
Benedict Elliott Smith 2025-07-25 15:24:30 +01:00
parent 87069e4dec
commit 4102ea13e5
8 changed files with 103 additions and 30 deletions

@ -1 +1 @@
Subproject commit 2f287b6c357d0fecb61ab0ba6910270a311ad775
Subproject commit 05b2f2415c3e7c109158208824fbc804dfe43f71

View File

@ -70,7 +70,7 @@ public class AccordMetrics
public static final String FAST_PATHS = "FastPaths";
public static final String MEDIUM_PATHS = "MediumPaths";
public static final String SLOW_PATHS = "SlowPaths";
public static final String PREEMPTS = "Preempts";
public static final String PREEMPTED = "Preempted";
public static final String TIMEOUTS = "Timeouts";
public static final String INVALIDATIONS = "Invalidations";
public static final String RECOVERY_DELAY = "RecoveryDelay";
@ -144,7 +144,7 @@ public class AccordMetrics
/**
* The number of preempted transactions on this coordinator.
*/
public final Meter preempts;
public final Meter preempted;
/**
* The number of timed out transactions on this coordinator.
@ -195,7 +195,7 @@ public class AccordMetrics
fastPaths = Metrics.meter(coordinator.createMetricName(FAST_PATHS));
mediumPaths = Metrics.meter(coordinator.createMetricName(MEDIUM_PATHS));
slowPaths = Metrics.meter(coordinator.createMetricName(SLOW_PATHS));
preempts = Metrics.meter(coordinator.createMetricName(PREEMPTS));
preempted = Metrics.meter(coordinator.createMetricName(PREEMPTED));
timeouts = Metrics.meter(coordinator.createMetricName(TIMEOUTS));
invalidations = Metrics.meter(coordinator.createMetricName(INVALIDATIONS));
recoveryDelay = Metrics.timer(coordinator.createMetricName(RECOVERY_DELAY));

View File

@ -101,7 +101,7 @@ import org.apache.cassandra.service.accord.serializers.FetchSerializers;
import org.apache.cassandra.service.accord.serializers.GetDurableBeforeSerializers;
import org.apache.cassandra.service.accord.serializers.GetEphmrlReadDepsSerializers;
import org.apache.cassandra.service.accord.serializers.GetMaxConflictSerializers;
import org.apache.cassandra.service.accord.serializers.InformDurableSerializers;
import org.apache.cassandra.service.accord.serializers.InformSerializers;
import org.apache.cassandra.service.accord.serializers.LatestDepsSerializers;
import org.apache.cassandra.service.accord.serializers.PreacceptSerializers;
import org.apache.cassandra.service.accord.serializers.ReadDataSerializer;
@ -321,14 +321,14 @@ public enum Verb
ACCORD_ACCEPT_RSP (122, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AcceptSerializers.reply), AccordService::responseHandlerOrNoop ),
ACCORD_ACCEPT_REQ (123, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AcceptSerializers.request), AccordService::requestHandlerOrNoop, ACCORD_ACCEPT_RSP ),
ACCORD_NOT_ACCEPT_REQ (124, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AcceptSerializers.notAccept), AccordService::requestHandlerOrNoop, ACCORD_ACCEPT_RSP ),
ACCORD_READ_RSP (125, P2, readTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.reply), AccordService::responseHandlerOrNoop ),
ACCORD_READ_REQ (126, P2, readTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_STABLE_THEN_READ_REQ (127, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_READ_RSP (125, P2, readTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.reply), AccordService::responseHandlerOrNoop ),
ACCORD_READ_REQ (126, P2, readTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_STABLE_THEN_READ_REQ (127, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_COMMIT_REQ (128, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(CommitSerializers.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_COMMIT_INVALIDATE_REQ (129, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(CommitSerializers.invalidate), AccordService::requestHandlerOrNoop ),
ACCORD_APPLY_RSP (130, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ApplySerializers.reply), AccordService::responseHandlerOrNoop ),
ACCORD_APPLY_REQ (131, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ApplySerializers.request), AccordService::requestHandlerOrNoop, ACCORD_APPLY_RSP ),
ACCORD_APPLY_AND_WAIT_REQ (132, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP),
ACCORD_APPLY_AND_WAIT_REQ (132, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP),
ACCORD_BEGIN_RECOVER_RSP (133, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(RecoverySerializers.reply), AccordService::responseHandlerOrNoop ),
ACCORD_BEGIN_RECOVER_REQ (134, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(RecoverySerializers.request), AccordService::requestHandlerOrNoop, ACCORD_BEGIN_RECOVER_RSP ),
ACCORD_BEGIN_INVALIDATE_RSP (135, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(BeginInvalidationSerializers.reply), AccordService::responseHandlerOrNoop ),
@ -336,10 +336,11 @@ public enum Verb
ACCORD_AWAIT_RSP (137, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AwaitSerializers.syncReply), AccordService::responseHandlerOrNoop ),
ACCORD_AWAIT_REQ (138, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AwaitSerializers.request), AccordService::requestHandlerOrNoop, ACCORD_AWAIT_RSP ),
ACCORD_AWAIT_ASYNC_RSP_REQ (139, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AwaitSerializers.asyncReply), AccordService::requestHandlerOrNoop ),
ACCORD_WAIT_UNTIL_APPLIED_REQ (140, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_WAIT_UNTIL_APPLIED_REQ (140, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(ReadDataSerializer.request), AccordService::requestHandlerOrNoop, ACCORD_READ_RSP ),
ACCORD_RECOVER_AWAIT_RSP (141, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AwaitSerializers.recoverReply), AccordService::responseHandlerOrNoop ),
ACCORD_RECOVER_AWAIT_REQ (142, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(AwaitSerializers.recoverRequest), AccordService::requestHandlerOrNoop, ACCORD_RECOVER_AWAIT_RSP),
ACCORD_INFORM_DURABLE_REQ (143, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(InformDurableSerializers.request), AccordService::requestHandlerOrNoop, ACCORD_SIMPLE_RSP ),
ACCORD_INFORM_DURABLE_REQ (143, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(InformSerializers.durable), AccordService::requestHandlerOrNoop, ACCORD_SIMPLE_RSP ),
ACCORD_INFORM_DECIDED_REQ (171, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(InformSerializers.decided), AccordService::requestHandlerOrNoop ),
ACCORD_CHECK_STATUS_RSP (144, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(CheckStatusSerializers.reply), AccordService::responseHandlerOrNoop ),
ACCORD_CHECK_STATUS_REQ (145, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(CheckStatusSerializers.request), AccordService::requestHandlerOrNoop, ACCORD_CHECK_STATUS_RSP ),
ACCORD_FETCH_DATA_RSP (146, P2, writeTimeout, IMMEDIATE, () -> accordEmbedded(FetchSerializers.reply), AccordService::responseHandlerOrNoop ),

View File

@ -120,6 +120,7 @@ public class AccordMessageSink implements MessageSink
builder.put(WAIT_UNTIL_APPLIED_REQ, Verb.ACCORD_WAIT_UNTIL_APPLIED_REQ);
builder.put(APPLY_THEN_WAIT_UNTIL_APPLIED_REQ, Verb.ACCORD_APPLY_AND_WAIT_REQ);
builder.put(INFORM_DURABLE_REQ, Verb.ACCORD_INFORM_DURABLE_REQ);
builder.put(INFORM_DECIDED_REQ, Verb.ACCORD_INFORM_DECIDED_REQ);
builder.put(CHECK_STATUS_REQ, Verb.ACCORD_CHECK_STATUS_REQ);
builder.put(CHECK_STATUS_RSP, Verb.ACCORD_CHECK_STATUS_RSP);
builder.put(FETCH_DATA_REQ, Verb.ACCORD_FETCH_DATA_REQ);

View File

@ -21,6 +21,7 @@ package org.apache.cassandra.service.accord.interop;
import accord.coordinate.Timeout;
import accord.local.Node;
import accord.messages.Callback;
import accord.messages.ReadData;
import accord.messages.ReadData.ReadOk;
import accord.messages.ReadData.ReadReply;
import accord.utils.Invariants;
@ -68,7 +69,7 @@ public abstract class AccordInteropReadCallback<T> implements Callback<ReadReply
// it and instead opts to trigger additional repair messages based on time.
interopExecution.sendMaximalCommit(id);
}
else
else if (reply != ReadData.CommitOrReadNack.Waiting)
{
wrapped.onFailure(endpoint, RequestFailure.UNKNOWN);
}

View File

@ -20,6 +20,8 @@ package org.apache.cassandra.service.accord.serializers;
import java.io.IOException;
import accord.api.RoutingKey;
import accord.messages.InformDecided;
import accord.messages.InformDurable;
import accord.primitives.Route;
import accord.primitives.Status;
@ -29,9 +31,34 @@ import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
public class InformDurableSerializers
public class InformSerializers
{
public static final IVersionedSerializer<InformDurable> request = new TxnRequestSerializer<InformDurable>()
public static final IVersionedSerializer<InformDecided> decided = new IVersionedSerializer<>()
{
@Override
public void serialize(InformDecided t, DataOutputPlus out, Version version) throws IOException
{
CommandSerializers.txnId.serialize(t.txnId, out);
KeySerializers.routingKey.serialize(t.homeKey, out);
}
@Override
public InformDecided deserialize(DataInputPlus in, Version version) throws IOException
{
TxnId txnId = CommandSerializers.txnId.deserialize(in);
RoutingKey homeKey = KeySerializers.routingKey.deserialize(in);
return new InformDecided(txnId, homeKey);
}
@Override
public long serializedSize(InformDecided t, Version version)
{
return CommandSerializers.txnId.serializedSize(t.txnId)
+ KeySerializers.routingKey.serializedSize(t.homeKey);
}
};
public static final IVersionedSerializer<InformDurable> durable = new TxnRequestSerializer<>()
{
@Override
public void serializeBody(InformDurable msg, DataOutputPlus out, Version version) throws IOException

View File

@ -116,19 +116,23 @@ public class AccordMetricsTest extends AccordTestBase
countingMetrics0 = getMetrics();
assertCoordinatorMetrics(0, "rw", 0, 0, 0, 0, 0);
SHARED_CLUSTER.coordinator(1).executeWithResult(writeCql(), ConsistencyLevel.ALL, 0, 0, 0, 0);
assertClientMetrics(0, "AccordWrite", 0, 0);
assertClientMetrics(1, "AccordWrite", 0, 0);
assertCoordinatorMetrics(0, "rw", 1, 0, 0, 0, 0);
assertCoordinatorMetrics(1, "rw", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "rw", 1, 1, 1);
assertReplicaMetrics(1, "rw", 1, 1, 1);
assertZeroMetrics("ro");
assertZeroMetrics("ro", "AccordRead");
countingMetrics0 = getMetrics();
SHARED_CLUSTER.coordinator(1).executeWithResult(readCql(), ConsistencyLevel.ALL, 0, 0, 1, 1);
assertClientMetrics(0, "AccordRead", 0, 0);
assertClientMetrics(1, "AccordRead", 0, 0);
assertCoordinatorMetrics(0, "ro", 1, 0, 0, 0, 0);
assertCoordinatorMetrics(1, "ro", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "ro", 0, 1, 1);
assertReplicaMetrics(1, "ro", 0, 1, 1);
assertZeroMetrics("rw");
assertZeroMetrics("rw", "AccordWrite");
}
@Test
@ -167,12 +171,15 @@ public class AccordMetricsTest extends AccordTestBase
Assertions.assertThat(ex).is(AssertionUtils.rootCauseIs(AccordWritePreemptedException.class));
}
assertCoordinatorMetrics(0, "rw", 0, 0, 1, 0, 0);
assertClientMetrics(0, "AccordWrite", 1, 0);
assertClientMetrics(1, "AccordWrite", 0, 0);
// TODO (required): reenable or remove metric
// assertCoordinatorMetrics(0, "rw", 0, 0, 1, 0, 0);
assertCoordinatorMetrics(1, "rw", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "rw", 0, 0, 0);
assertReplicaMetrics(1, "rw", 0, 0, 0);
assertZeroMetrics("ro");
assertZeroMetrics("ro", "AccordRead");
countingMetrics0 = getMetrics();
try
@ -185,12 +192,15 @@ public class AccordMetricsTest extends AccordTestBase
Assertions.assertThat(ex).is(AssertionUtils.rootCauseIs(AccordReadPreemptedException.class));
}
assertCoordinatorMetrics(0, "ro", 0, 0, 1, 0, 0);
assertClientMetrics(0, "AccordRead", 1, 0);
assertClientMetrics(1, "AccordRead", 0, 0);
// TODO (required): reenable or remove metric
// assertCoordinatorMetrics(0, "ro", 0, 0, 1, 0, 0);
assertCoordinatorMetrics(1, "ro", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "ro", 0, 0, 0);
assertReplicaMetrics(1, "ro", 0, 0, 0);
assertZeroMetrics("rw");
assertZeroMetrics("rw", "AccordWrite");
}
finally
{
@ -219,12 +229,15 @@ public class AccordMetricsTest extends AccordTestBase
Assertions.assertThat(ex).is(AssertionUtils.rootCauseIs(ReadTimeoutException.class));
}
assertCoordinatorMetrics(0, "ro", 0, 0, 0, 1, 0);
assertClientMetrics(0, "AccordRead", 0, 1);
assertClientMetrics(1, "AccordRead", 0, 0);
// TODO (required): reenable or remove metric
// assertCoordinatorMetrics(0, "ro", 0, 0, 0, 1, 0);
assertCoordinatorMetrics(1, "ro", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "ro", 0, 0, 0);
assertReplicaMetrics(1, "ro", 0, 0, 0);
assertZeroMetrics("rw");
assertZeroMetrics("rw", "AccordWrite");
countingMetrics0 = getMetrics();
try
@ -237,23 +250,47 @@ public class AccordMetricsTest extends AccordTestBase
Assertions.assertThat(ex).is(AssertionUtils.rootCauseIs(WriteTimeoutException.class));
}
assertCoordinatorMetrics(0, "rw", 0, 0, 0, 1, 0);
assertClientMetrics(0, "AccordWrite", 0, 1);
assertClientMetrics(1, "AccordWrite", 0, 0);
// TODO (required): reenable or remove metric
// assertCoordinatorMetrics(0, "rw", 0, 0, 0, 1, 0);
assertCoordinatorMetrics(1, "rw", 0, 0, 0, 0, 0);
assertReplicaMetrics(0, "rw", 0, 0, 0);
assertReplicaMetrics(1, "rw", 0, 0, 0);
assertZeroMetrics("ro");
assertZeroMetrics("ro","AccordRead");
}
private void assertZeroMetrics(String scope)
private void assertZeroMetrics(String scope, String clientScope)
{
for (int i = 0; i < SHARED_CLUSTER.size(); i++)
{
assertClientMetrics(0, clientScope, 0, 0);
assertCoordinatorMetrics(i, scope, 0, 0, 0, 0, 0);
assertReplicaMetrics(i, scope, 0, 0, 0);
}
}
private void assertClientMetrics(int node, String scope, long preempts, long timeouts)
{
DefaultNameFactory nameFactory = new DefaultNameFactory("ClientRequest", scope);
Map<String, Long> metrics = diff(countingMetrics0).get(node);
logger.info("Metrics for node {} / {}: {}", node, scope, metrics);
Function<String, Long> metric = n -> metrics.get(nameFactory.createMetricName(n).getMetricName());
assertThat(metric.apply("Preempted")).isEqualTo(preempts);
assertThat(metric.apply("Timeouts")).isEqualTo(timeouts);
// Verify that coordinator metrics are published to the appropriate virtual table:
// SimpleQueryResult res = SHARED_CLUSTER.get(node + 1)
// .executeInternalWithResult("SELECT * FROM system_metrics.accord_coordinator_group WHERE scope = ?", scope);
// while (res.hasNext())
// {
// Row metricRow = res.next();
// String name = metricRow.getString("name");
// assertThat(metrics).containsKey(name);
// }
}
private void assertCoordinatorMetrics(int node, String scope, long fastPaths, long slowPaths, long preempts, long timeouts, long recoveries)
{
DefaultNameFactory nameFactory = new DefaultNameFactory(AccordMetrics.ACCORD_COORDINATOR, scope);
@ -262,7 +299,7 @@ public class AccordMetricsTest extends AccordTestBase
Function<String, Long> metric = n -> metrics.get(nameFactory.createMetricName(n).getMetricName());
assertThat(metric.apply(AccordMetrics.FAST_PATHS)).isEqualTo(fastPaths);
assertThat(metric.apply(AccordMetrics.SLOW_PATHS)).isEqualTo(slowPaths);
assertThat(metric.apply(AccordMetrics.PREEMPTS)).isEqualTo(preempts);
assertThat(metric.apply(AccordMetrics.PREEMPTED)).isEqualTo(preempts);
assertThat(metric.apply(AccordMetrics.TIMEOUTS)).isEqualTo(timeouts);
assertThat(metric.apply(AccordMetrics.RECOVERY_DELAY)).isEqualTo(recoveries);
assertThat(metric.apply(AccordMetrics.RECOVERY_TIME)).isEqualTo(recoveries);
@ -317,7 +354,12 @@ public class AccordMetricsTest extends AccordTestBase
{
Map<Integer, Map<String, Long>> metrics = new HashMap<>();
for (int i = 0; i < SHARED_CLUSTER.size(); i++)
metrics.put(i, SHARED_CLUSTER.get(i + 1).metrics().getCounters(name -> name.startsWith("org.apache.cassandra.metrics.Accord")));
{
Map<String, Long> map = SHARED_CLUSTER.get(i + 1).metrics().getCounters(name -> name.startsWith("org.apache.cassandra.metrics.Accord") || (name.startsWith("org.apache.cassandra.metrics.ClientRequest") && (name.endsWith("AccordRead") || name.endsWith("AccordWrite"))));
SHARED_CLUSTER.get(i + 1).metrics().getGauges(name -> name.startsWith("org.apache.cassandra.metrics.AccordReplica"))
.forEach((key, value) -> map.put(key, (Long)value));
metrics.put(i, map);
}
return metrics;
}

View File

@ -21,7 +21,6 @@ package org.apache.cassandra.service.accord;
import java.util.EnumSet;
import java.util.Set;
import com.google.common.collect.Sets;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@ -76,6 +75,7 @@ public class CommandChangeTest
missing.remove(Field.PROMISED);
missing.remove(Field.ACCEPTED);
missing.remove(Field.DURABILITY);
missing.remove(Field.EXECUTE_AT);
assertMissing(flags, missing);
}
@ -84,8 +84,9 @@ public class CommandChangeTest
{
int flags = getFlags(null, Command.NotDefined.uninitialised(TxnId.NONE));
EnumSet<Field> has = EnumSet.of(Field.SAVE_STATUS, Field.PARTICIPANTS, Field.DURABILITY, Field.PROMISED,
Field.ACCEPTED /* this is Zero... which kinda means null... */);
Set<Field> missing = Sets.difference(ALL, has);
Field.ACCEPTED, Field.EXECUTE_AT /* this is Zero... which kinda means null... */);
EnumSet<Field> missing = EnumSet.complementOf(has);
has.remove(Field.EXECUTE_AT); // we serialize executeAt whenever we change SaveStatus, so we expect it to be non-null
assertHas(flags, has);
assertMissing(flags, missing);
}