diff --git a/modules/accord b/modules/accord index 7c15f3a620..2ad55e03c4 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 7c15f3a6203939bc6cb398e538df1ca3557cbe03 +Subproject commit 2ad55e03c43ce074cdf5e36cfa14cb4278c2dc0f diff --git a/src/java/org/apache/cassandra/net/Verb.java b/src/java/org/apache/cassandra/net/Verb.java index be2d86ad59..dc8e6d4157 100644 --- a/src/java/org/apache/cassandra/net/Verb.java +++ b/src/java/org/apache/cassandra/net/Verb.java @@ -272,36 +272,36 @@ public enum Verb // accord ACCORD_SIMPLE_RSP (119, P2, writeTimeout, REQUEST_RESPONSE, () -> EnumSerializer.simpleReply, RESPONSE_HANDLER ), ACCORD_PRE_ACCEPT_RSP (121, P2, writeTimeout, REQUEST_RESPONSE, () -> PreacceptSerializers.reply, RESPONSE_HANDLER ), - ACCORD_PRE_ACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_PRE_ACCEPT_RSP ), + ACCORD_PRE_ACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_PRE_ACCEPT_RSP ), ACCORD_ACCEPT_RSP (124, P2, writeTimeout, REQUEST_RESPONSE, () -> AcceptSerializers.reply, RESPONSE_HANDLER ), - ACCORD_ACCEPT_REQ (122, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ), - ACCORD_ACCEPT_INVALIDATE_REQ (123, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.invalidate, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ), + ACCORD_ACCEPT_REQ (122, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_ACCEPT_RSP ), + ACCORD_ACCEPT_INVALIDATE_REQ (123, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.invalidate, AccordService::verbHandlerOrNoop, ACCORD_ACCEPT_RSP ), ACCORD_READ_RSP (126, P2, writeTimeout, REQUEST_RESPONSE, () -> ReadDataSerializers.reply, RESPONSE_HANDLER ), - ACCORD_READ_REQ (125, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ), - ACCORD_COMMIT_REQ (127, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ), - ACCORD_COMMIT_INVALIDATE_REQ (128, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.invalidate, () -> AccordService.instance().verbHandler() ), + ACCORD_READ_REQ (125, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ), + ACCORD_COMMIT_REQ (127, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ), + ACCORD_COMMIT_INVALIDATE_REQ (128, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.invalidate, AccordService::verbHandlerOrNoop ), ACCORD_APPLY_RSP (130, P2, writeTimeout, REQUEST_RESPONSE, () -> ApplySerializers.reply, RESPONSE_HANDLER ), - ACCORD_APPLY_REQ (129, P2, writeTimeout, IMMEDIATE, () -> ApplySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_APPLY_RSP ), + ACCORD_APPLY_REQ (129, P2, writeTimeout, IMMEDIATE, () -> ApplySerializers.request, AccordService::verbHandlerOrNoop, ACCORD_APPLY_RSP ), ACCORD_BEGIN_RECOVER_RSP (132, P2, writeTimeout, REQUEST_RESPONSE, () -> RecoverySerializers.reply, RESPONSE_HANDLER ), - ACCORD_BEGIN_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_RECOVER_RSP ), + ACCORD_BEGIN_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, AccordService::verbHandlerOrNoop, ACCORD_BEGIN_RECOVER_RSP ), ACCORD_BEGIN_INVALIDATE_RSP (134, P2, writeTimeout, REQUEST_RESPONSE, () -> BeginInvalidationSerializers.reply, RESPONSE_HANDLER ), - ACCORD_BEGIN_INVALIDATE_REQ (133, P2, writeTimeout, IMMEDIATE, () -> BeginInvalidationSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_INVALIDATE_RSP ), + ACCORD_BEGIN_INVALIDATE_REQ (133, P2, writeTimeout, IMMEDIATE, () -> BeginInvalidationSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_BEGIN_INVALIDATE_RSP ), ACCORD_WAIT_ON_COMMIT_RSP (136, P2, writeTimeout, REQUEST_RESPONSE, () -> WaitOnCommitSerializer.reply, RESPONSE_HANDLER ), - ACCORD_WAIT_ON_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, () -> AccordService.instance().verbHandler(), ACCORD_WAIT_ON_COMMIT_RSP ), - ACCORD_WAIT_ON_APPLY_REQ (137, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.waitOnApply, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ), - ACCORD_INFORM_OF_TXN_REQ (138, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ), - ACCORD_INFORM_HOME_DURABLE_REQ (139, P2, writeTimeout, IMMEDIATE, () -> InformHomeDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ), - ACCORD_INFORM_DURABLE_REQ (140, P2, writeTimeout, IMMEDIATE, () -> InformDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ), + ACCORD_WAIT_ON_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, AccordService::verbHandlerOrNoop, ACCORD_WAIT_ON_COMMIT_RSP ), + ACCORD_WAIT_ON_APPLY_REQ (137, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.waitOnApply, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ), + ACCORD_INFORM_OF_TXN_REQ (138, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ), + ACCORD_INFORM_HOME_DURABLE_REQ (139, P2, writeTimeout, IMMEDIATE, () -> InformHomeDurableSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ), + ACCORD_INFORM_DURABLE_REQ (140, P2, writeTimeout, IMMEDIATE, () -> InformDurableSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ), ACCORD_CHECK_STATUS_RSP (142, P2, writeTimeout, REQUEST_RESPONSE, () -> CheckStatusSerializers.reply, RESPONSE_HANDLER ), - ACCORD_CHECK_STATUS_REQ (141, P2, writeTimeout, IMMEDIATE, () -> CheckStatusSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_CHECK_STATUS_RSP ), + ACCORD_CHECK_STATUS_REQ (141, P2, writeTimeout, IMMEDIATE, () -> CheckStatusSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_CHECK_STATUS_RSP ), ACCORD_GET_DEPS_RSP (144, P2, writeTimeout, REQUEST_RESPONSE, () -> GetDepsSerializers.reply, RESPONSE_HANDLER ), - ACCORD_GET_DEPS_REQ (143, P2, writeTimeout, IMMEDIATE, () -> GetDepsSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_GET_DEPS_RSP ), + ACCORD_GET_DEPS_REQ (143, P2, writeTimeout, IMMEDIATE, () -> GetDepsSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_GET_DEPS_RSP ), ACCORD_FETCH_DATA_RSP (146, P2, repairTimeout,REQUEST_RESPONSE, () -> FetchSerializers.reply, RESPONSE_HANDLER ), - ACCORD_FETCH_DATA_REQ (145, P2, repairTimeout,IMMEDIATE, () -> FetchSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_FETCH_DATA_RSP ), - ACCORD_SET_SHARD_DURABLE_REQ (147, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.shardDurable, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ), - ACCORD_SET_GLOBALLY_DURABLE_REQ (148, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.globallyDurable,() -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ), + ACCORD_FETCH_DATA_REQ (145, P2, repairTimeout,IMMEDIATE, () -> FetchSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_FETCH_DATA_RSP ), + ACCORD_SET_SHARD_DURABLE_REQ (147, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.shardDurable, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ), + ACCORD_SET_GLOBALLY_DURABLE_REQ (148, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.globallyDurable,AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ), ACCORD_QUERY_DURABLE_BEFORE_RSP (150, P2, writeTimeout, REQUEST_RESPONSE, () -> QueryDurableBeforeSerializers.reply, RESPONSE_HANDLER ), - ACCORD_QUERY_DURABLE_BEFORE_REQ (149, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.request,() -> AccordService.instance().verbHandler(), ACCORD_QUERY_DURABLE_BEFORE_RSP), + ACCORD_QUERY_DURABLE_BEFORE_REQ (149, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.request,AccordService::verbHandlerOrNoop, ACCORD_QUERY_DURABLE_BEFORE_RSP), ACCORD_SYNC_NOTIFY_REQ (151, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SIMPLE_RSP ), diff --git a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java index df77cae975..610147e341 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java @@ -24,6 +24,7 @@ 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; @@ -132,13 +133,17 @@ public class AccordConfigurationService extends AbstractConfigurationService { if (topology != null) reportTopology(topology, syncStatus == SyncStatus.NOT_STARTED); getOrCreateEpochState(epoch).setSyncStatus(syncStatus); if (syncStatus == SyncStatus.NOTIFYING) - syncPropagator.reportSyncComplete(epoch, pendingSyncNotify, localId); + { + // TODO (expected, correctness): since this is loading old topologies, might see nodes no longer present (host replacement, decom, shrink, etc.); attempt to remove unknown nodes + syncPropagator.reportSyncComplete(epoch, Sets.filter(pendingSyncNotify, snapshot::containsId), localId); + } remoteSyncComplete.forEach(id -> receiveRemoteSyncComplete(id, epoch)); // TODO (now): disk doesn't get updated until we see our own notification, so there is an edge case where this instance notified others and fails in the middle, but Apply was already sent! This could leave partial closed/redudant accross the cluster @@ -157,7 +162,7 @@ public class AccordConfigurationService extends AbstractConfigurationService notify = topology.nodes().stream().filter(i -> !localId.equals(i)).collect(Collectors.toSet()); diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 9c4a6b5a18..cf1bd2923d 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -163,12 +163,24 @@ public class AccordService implements IAccordService, Shutdownable } }; - private static Node.Id localId = null; + private static volatile Node.Id localId = null; + private static class Handle { public static final AccordService instance = new AccordService(); } + public static boolean isSetup() + { + return localId != null; + } + + public static IVerbHandler verbHandlerOrNoop() + { + if (!isSetup()) return ignore -> {}; + return instance().verbHandler(); + } + public static void startup(NodeId tcmId) { localId = AccordTopologyUtils.tcmIdToAccord(tcmId); diff --git a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java index 2af9c9472b..3e215a7e68 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java @@ -59,7 +59,11 @@ import static org.apache.cassandra.utils.CollectionSerializers.newListSerializer */ public class AccordSyncPropagator { - public static final IVerbHandler> verbHandler = message -> AccordService.instance().receive(message); + public static final IVerbHandler> verbHandler = message -> { + if (!AccordService.isSetup()) + return; + AccordService.instance().receive(message); + }; interface Listener { diff --git a/src/java/org/apache/cassandra/service/accord/AccordTopologyUtils.java b/src/java/org/apache/cassandra/service/accord/AccordTopologyUtils.java index f234738402..a41e9732d9 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTopologyUtils.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTopologyUtils.java @@ -131,11 +131,16 @@ public class AccordTopologyUtils return new Topology(epoch.getEpoch(), shards.toArray(new Shard[0])); } - public static EndpointMapping directoryToMapping(long epoch, Directory directory) + public static EndpointMapping directoryToMapping(EndpointMapping mapping, long epoch, Directory directory) { EndpointMapping.Builder builder = EndpointMapping.builder(epoch); for (NodeId id : directory.peerIds()) builder.add(directory.endpoint(id), tcmIdToAccord(id)); + + // There are cases where nodes are removed from the cluster (host replacement, decom, etc.), but inflight events may still be happening; + // keep the ids around so pending events do not fail with a mapping error + for (Node.Id id : mapping.differenceIds(builder)) + builder.add(mapping.mappedEndpoint(id), id); return builder.build(); } diff --git a/src/java/org/apache/cassandra/service/accord/EndpointMapping.java b/src/java/org/apache/cassandra/service/accord/EndpointMapping.java index 4a746dc9a2..0c964d3204 100644 --- a/src/java/org/apache/cassandra/service/accord/EndpointMapping.java +++ b/src/java/org/apache/cassandra/service/accord/EndpointMapping.java @@ -18,9 +18,12 @@ package org.apache.cassandra.service.accord; +import java.util.Set; + import com.google.common.collect.BiMap; import com.google.common.collect.HashBiMap; import com.google.common.collect.ImmutableBiMap; +import com.google.common.collect.Sets; import accord.local.Node; import accord.utils.Invariants; @@ -44,6 +47,16 @@ class EndpointMapping implements AccordEndpointMapper return epoch; } + public boolean containsId(Node.Id id) + { + return mapping.containsKey(id); + } + + public Set differenceIds(Builder builder) + { + return Sets.difference(mapping.keySet(), builder.mapping.keySet()); + } + @Override public Node.Id mappedId(InetAddressAndPort endpoint) { diff --git a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java index 4b8b3fbb14..a2d62eb113 100644 --- a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java +++ b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java @@ -981,7 +981,10 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance () -> SharedExecutorPool.SHARED.shutdownAndWait(1L, MINUTES) ); - error = parallelRun(error, executor, () -> AccordService.instance().shutdownAndWait(1l, MINUTES)); + error = parallelRun(error, executor, () -> { + if (!AccordService.isSetup()) return; + AccordService.instance().shutdownAndWait(1l, MINUTES); + }); // CommitLog must shut down after Stage, or threads from the latter may attempt to use the former. // (ex. A Mutation stage thread may attempt to add a mutation to the CommitLog.) diff --git a/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java b/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java index 333bca572b..3cf9da57c1 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java @@ -207,19 +207,19 @@ public class AccordConfigurationServiceTest Topology topology1 = new Topology(1, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, ID_SET)); service.updateMapping(mappingForEpoch(ClusterMetadata.current().epoch.getEpoch() + 1)); service.reportTopology(topology1); - service.acknowledgeEpoch(EpochReady.done(1)); + service.acknowledgeEpoch(EpochReady.done(1), true); service.receiveRemoteSyncComplete(ID1, 1); service.receiveRemoteSyncComplete(ID2, 1); service.receiveRemoteSyncComplete(ID3, 1); Topology topology2 = new Topology(2, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, of(ID1, ID2))); service.reportTopology(topology2); - service.acknowledgeEpoch(EpochReady.done(2)); + service.acknowledgeEpoch(EpochReady.done(2), true); service.receiveRemoteSyncComplete(ID1, 2); Topology topology3 = new Topology(3, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, of(ID1, ID2))); service.reportTopology(topology3); - service.acknowledgeEpoch(EpochReady.done(3)); + service.acknowledgeEpoch(EpochReady.done(3), true); AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector()); loaded.updateMapping(mappingForEpoch(ClusterMetadata.current().epoch.getEpoch() + 1)); diff --git a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java index 4159e6364c..ab6e2790ce 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java @@ -427,7 +427,7 @@ public class AccordSyncPropagatorTest } @Override - protected void localSyncComplete(Topology topology) + protected void localSyncComplete(Topology topology, boolean startSync) { Set notify = topology.nodes().stream().filter(i -> !localId.equals(i)).collect(Collectors.toSet()); instances.get(localId).propagator.reportSyncComplete(topology.epoch(), notify, localId);