From dfcf3aff4eb40129ef5f98f8d289fa500c777278 Mon Sep 17 00:00:00 2001 From: Alex Petrov Date: Thu, 5 Dec 2024 13:51:58 +0100 Subject: [PATCH] Fix topology replay during bootstrap and startup, decouple Accord from TCM Includes multiple changes, primary ones: * Make HarrySimulatorTest work with Accord * Removed nodes now live in TCM, no need to discover historic epochs in order to find removed nodes * CommandStore <-> RangesForEpochs mappings required for startup are now stored in journal, and CS can be set up _without_ topology replay * Topology replay is fully done via journal (where we store topologies themselves), and topology metadata table (where we store redundant/closed information) * Fixed various bugs related to propagation and staleness * TCM was previously relied on for "fetching" epoch: we can not rely on it as there's no guarantee we will see a consecutive epoch when grabbing Metadata#current * Redundant / closed during replay was set with incorrect ranges in 1 of the code paths * TCM was contacted multiple times for historical epochs, which made startup much longer under some circumstances Patch by Alex Petrov; reviewed by Benedict Elliott Smith for CASSANDRA-20142 --- modules/accord | 2 +- .../apache/cassandra/concurrent/Stage.java | 4 +- .../org/apache/cassandra/config/Config.java | 2 + .../cassandra/config/DatabaseDescriptor.java | 5 + .../cassandra/exceptions/RequestFailure.java | 2 + .../exceptions/RequestFailureReason.java | 1 + .../org/apache/cassandra/journal/Journal.java | 6 +- .../locator/CMSPlacementStrategy.java | 2 +- .../apache/cassandra/net/MessageDelivery.java | 14 +- .../apache/cassandra/net/MessagingUtils.java | 67 +++++++ src/java/org/apache/cassandra/net/Verb.java | 30 ++-- .../cassandra/repair/SharedContext.java | 1 + .../service/accord/AccordCommandStore.java | 2 +- .../service/accord/AccordCommandStores.java | 14 +- .../accord/AccordConfigurationService.java | 167 +++++++++--------- .../service/accord/AccordJournal.java | 32 +++- .../service/accord/AccordJournalTable.java | 4 +- .../accord/AccordJournalValueSerializers.java | 107 ++++++++++- .../service/accord/AccordKeyspace.java | 21 +-- .../accord/AccordSegmentCompactor.java | 3 + .../service/accord/AccordService.java | 159 +++++++---------- .../service/accord/AccordSyncPropagator.java | 34 ++-- .../service/accord/AccordTopology.java | 19 +- .../service/accord/EndpointMapping.java | 27 ++- .../service/accord/FetchMinEpoch.java | 44 +++-- .../service/accord/FetchTopology.java | 140 +++++++++++++++ .../cassandra/service/accord/JournalKey.java | 2 + .../serializers/TopologySerializers.java | 19 +- .../apache/cassandra/tcm/ClusterMetadata.java | 4 +- .../org/apache/cassandra/tcm/Startup.java | 12 +- .../cassandra/tcm/membership/Directory.java | 116 +++++++++++- .../cassandra/tcm/serialization/Version.java | 1 + .../org/apache/cassandra/utils/Backoff.java | 24 +++ .../distributed/shared/ClusterUtils.java | 2 +- .../test/accord/AccordBootstrapTest.java | 14 +- .../fuzz/topology/AccordBootstrapTest.java | 161 +++++++++++++++++ .../service/accord/AccordJournalBurnTest.java | 10 +- .../harry/execution/CQLVisitExecutor.java | 9 +- .../execution/InJvmDTestVisitExecutor.java | 7 +- .../cassandra/harry/gen/SchemaGenerators.java | 8 +- .../org/apache/cassandra/harry/op/Visit.java | 6 + .../simulator/paxos/PaxosSimulation.java | 3 +- .../test/AccordHarrySimulationTest.java | 71 ++++++++ .../simulator/test/HarrySimulatorTest.java | 43 +++-- .../cassandra/net/MessageDeliveryTest.java | 59 ++++--- .../AccordConfigurationServiceTest.java | 109 +++++++----- .../service/accord/AccordTestUtils.java | 27 --- .../service/accord/EpochSyncTest.java | 9 +- .../service/accord/FetchMinEpochTest.java | 28 ++- .../accord/LoggingDiskStateManager.java | 4 +- .../service/accord/MockDiskStateManager.java | 2 +- .../cassandra/utils/SimulatedMiniCluster.java | 2 +- 52 files changed, 1215 insertions(+), 446 deletions(-) create mode 100644 src/java/org/apache/cassandra/net/MessagingUtils.java create mode 100644 src/java/org/apache/cassandra/service/accord/FetchTopology.java create mode 100644 test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java create mode 100644 test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java diff --git a/modules/accord b/modules/accord index f7b9bb8887..c076383eb4 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit f7b9bb8887ed672185f269ebcbc9d11e6aeafca9 +Subproject commit c076383eb432670c4d919e9fd0db76296169ea00 diff --git a/src/java/org/apache/cassandra/concurrent/Stage.java b/src/java/org/apache/cassandra/concurrent/Stage.java index b4f0540d93..557c1ca0a4 100644 --- a/src/java/org/apache/cassandra/concurrent/Stage.java +++ b/src/java/org/apache/cassandra/concurrent/Stage.java @@ -47,7 +47,7 @@ public enum Stage MUTATION (true, "MutationStage", "request", DatabaseDescriptor::getConcurrentWriters, DatabaseDescriptor::setConcurrentWriters, Stage::multiThreadedLowSignalStage), COUNTER_MUTATION (true, "CounterMutationStage", "request", DatabaseDescriptor::getConcurrentCounterWriters, DatabaseDescriptor::setConcurrentCounterWriters, Stage::multiThreadedLowSignalStage), VIEW_MUTATION (true, "ViewMutationStage", "request", DatabaseDescriptor::getConcurrentViewWriters, DatabaseDescriptor::setConcurrentViewWriters, Stage::multiThreadedLowSignalStage), - ACCORD_MIGRATION (false, "AccordMigrationStage", "request", DatabaseDescriptor::getAccordConcurrentOps, DatabaseDescriptor::setConcurrentAccordOps, Stage::multiThreadedLowSignalStage), + ACCORD_MIGRATION (false, "AccordMigrationStage", "request", DatabaseDescriptor::getAccordConcurrentOps, DatabaseDescriptor::setConcurrentAccordOps, Stage::multiThreadedLowSignalStage), GOSSIP (true, "GossipStage", "internal", () -> 1, null, Stage::singleThreadedStage), REQUEST_RESPONSE (false, "RequestResponseStage", "request", FBUtilities::getAvailableProcessors, null, Stage::multiThreadedLowSignalStage), ANTI_ENTROPY (false, "AntiEntropyStage", "internal", () -> 1, null, Stage::singleThreadedStage), @@ -58,7 +58,7 @@ public enum Stage IMMEDIATE (false, "ImmediateStage", "internal", () -> 0, null, Stage::immediateExecutor), PAXOS_REPAIR (false, "PaxosRepairStage", "internal", FBUtilities::getAvailableProcessors, null, Stage::multiThreadedStage), INTERNAL_METADATA (false, "InternalMetadataStage", "internal", FBUtilities::getAvailableProcessors, null, Stage::multiThreadedStage), - FETCH_LOG (false, "MetadataFetchLogStage", "internal", () -> 1, null, Stage::singleThreadedStage), + FETCH_METADATA (false, "MetadataFetchLogStage", "internal", () -> 1, null, Stage::singleThreadedStage), ; public final String jmxName; private final Supplier executorSupplier; diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index 615480518d..62ae907ead 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -151,6 +151,8 @@ public class Config @Replaces(oldName = "write_request_timeout_in_ms", converter = Converters.MILLIS_DURATION_LONG, deprecated = true) public volatile DurationSpec.LongMillisecondsBound write_request_timeout = new DurationSpec.LongMillisecondsBound("2000ms"); + public volatile DurationSpec.LongMillisecondsBound short_rpc_timeout = new DurationSpec.LongMillisecondsBound("1000ms"); + @Replaces(oldName = "counter_write_request_timeout_in_ms", converter = Converters.MILLIS_DURATION_LONG, deprecated = true) public volatile DurationSpec.LongMillisecondsBound counter_write_request_timeout = new DurationSpec.LongMillisecondsBound("5000ms"); diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 25035cb610..51add020cd 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -2507,6 +2507,11 @@ public class DatabaseDescriptor return conf.write_request_timeout.to(unit); } + public static long getShortRpcTimeout(TimeUnit unit) + { + return conf.short_rpc_timeout.to(unit); + } + public static void setWriteRpcTimeout(long timeOutInMillis) { conf.write_request_timeout = new DurationSpec.LongMillisecondsBound(timeOutInMillis); diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailure.java b/src/java/org/apache/cassandra/exceptions/RequestFailure.java index 2b2282731a..312102c012 100644 --- a/src/java/org/apache/cassandra/exceptions/RequestFailure.java +++ b/src/java/org/apache/cassandra/exceptions/RequestFailure.java @@ -54,6 +54,7 @@ public class RequestFailure public static final RequestFailure COORDINATOR_BEHIND = new RequestFailure(RequestFailureReason.COORDINATOR_BEHIND); public static final RequestFailure READ_TOO_MANY_INDEXES = new RequestFailure(RequestFailureReason.READ_TOO_MANY_INDEXES); public static final RequestFailure RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM = new RequestFailure(RequestFailureReason.RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM); + public static final RequestFailure BOOTING = new RequestFailure(RequestFailureReason.RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM); static { @@ -144,6 +145,7 @@ public class RequestFailure case COORDINATOR_BEHIND: return COORDINATOR_BEHIND; case READ_TOO_MANY_INDEXES: return READ_TOO_MANY_INDEXES; case RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM: return RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM; + case BOOTING: return BOOTING; } } diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java index 600cdebf86..22f75c44a3 100644 --- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java +++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java @@ -50,6 +50,7 @@ public enum RequestFailureReason // The following codes have been ported from an external fork, where they were offset explicitly to avoid conflicts. INDEX_BUILD_IN_PROGRESS (503), RETRY_ON_DIFFERENT_TRANSACTION_SYSTEM (504), + BOOTING (505), ; static diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java index 630959cd38..14bafa375d 100644 --- a/src/java/org/apache/cassandra/journal/Journal.java +++ b/src/java/org/apache/cassandra/journal/Journal.java @@ -335,12 +335,12 @@ public class Journal implements Shutdownable return null; } - public void readAll(K id, RecordConsumer consumer) + public void readAll(K id, RecordConsumer consumer, boolean asc) { EntrySerializer.EntryHolder holder = new EntrySerializer.EntryHolder<>(); try (OpOrder.Group group = readOrder.start()) { - for (Segment segment : segments.get().allSorted(false)) + for (Segment segment : segments.get().allSorted(asc)) segment.readAll(id, holder, consumer); } } @@ -359,7 +359,7 @@ public class Journal implements Shutdownable // can only throw if serializer is buggy throw new RuntimeException(e); } - }); + }, false); return res; } diff --git a/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java b/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java index 716b40a34c..e00289aef6 100644 --- a/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java +++ b/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java @@ -83,7 +83,7 @@ public class CMSPlacementStrategy { if (!filter.apply(metadata, peerId)) { - tmpDirectory = tmpDirectory.without(peerId); + tmpDirectory = tmpDirectory.without(metadata.nextEpoch(), peerId); tmpTokenMap = tmpTokenMap.unassignTokens(peerId); } } diff --git a/src/java/org/apache/cassandra/net/MessageDelivery.java b/src/java/org/apache/cassandra/net/MessageDelivery.java index 7c36d73c14..9dc374750f 100644 --- a/src/java/org/apache/cassandra/net/MessageDelivery.java +++ b/src/java/org/apache/cassandra/net/MessageDelivery.java @@ -22,7 +22,6 @@ import java.util.Collection; import java.util.Iterator; import java.util.Set; import java.util.concurrent.TimeUnit; - import javax.annotation.Nullable; import org.slf4j.Logger; @@ -100,6 +99,14 @@ public interface MessageDelivery return promise; } + public default Future> sendWithRetries(Verb verb, REQ request, + Iterator candidates, + RetryPredicate shouldRetry, + RetryErrorMessage errorMessage) + { + return sendWithRetries(Backoff.NO_OP.INSTANCE, ImmediateRetryScheduler.instance, verb, request, candidates, shouldRetry, errorMessage); + } + public default void sendWithRetries(Backoff backoff, RetryScheduler retryThreads, Verb verb, REQ request, Iterator candidates, @@ -127,11 +134,15 @@ public interface MessageDelivery interface RetryPredicate { + static RetryPredicate times(int n) { return (attempt, from, failure) -> attempt < n; } + RetryPredicate ALWAYS_RETRY = (i1, i2, i3) -> true; + RetryPredicate NEVER_RETRY = (i1, i2, i3) -> false; boolean test(int attempt, InetAddressAndPort from, RequestFailure failure); } interface RetryErrorMessage { + RetryErrorMessage EMPTY = (i1, i2, i3, i4) -> null; String apply(int attempt, ResponseFailureReason retryFailure, @Nullable InetAddressAndPort from, @Nullable RequestFailure reason); } @@ -165,6 +176,7 @@ public interface MessageDelivery @Override public void onFailure(InetAddressAndPort from, RequestFailure failure) { + // TODO (required): we already have a separate retry predicate, backoff should not be taken into consideration when retrying if (!backoff.mayRetry(attempt)) { onResult.result(attempt, null, new MaxRetriesException(attempt, errorMessage.apply(attempt, ResponseFailureReason.MaxRetries, from, failure))); diff --git a/src/java/org/apache/cassandra/net/MessagingUtils.java b/src/java/org/apache/cassandra/net/MessagingUtils.java new file mode 100644 index 0000000000..2190eaf3a6 --- /dev/null +++ b/src/java/org/apache/cassandra/net/MessagingUtils.java @@ -0,0 +1,67 @@ +/* + * 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.net; + +import java.util.Iterator; + +import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.repair.SharedContext; + +public class MessagingUtils +{ + /** + * Candidate iterator that would try all endpoints known to be alive first, and then try all endpoints + * in a round-robin manner. + */ + public static Iterator tryAliveFirst(SharedContext context, Iterable peers) + { + return new Iterator<>() + { + boolean firstRun = true; + Iterator iter = peers.iterator(); + boolean isEmpty = !iter.hasNext(); + + public boolean hasNext() + { + return !isEmpty; + } + + public InetAddressAndPort next() + { + // At first, try all alive nodes + if (firstRun) + { + while (iter.hasNext()) + { + InetAddressAndPort candidate = iter.next(); + if (context.failureDetector().isAlive(candidate)) + return candidate; + } + firstRun = false; + } + + // After that, cycle through all nodes + if (!iter.hasNext()) + iter = peers.iterator(); + + return iter.next(); + } + }; + } +} diff --git a/src/java/org/apache/cassandra/net/Verb.java b/src/java/org/apache/cassandra/net/Verb.java index fb2271dd1b..0b8bc8dd71 100644 --- a/src/java/org/apache/cassandra/net/Verb.java +++ b/src/java/org/apache/cassandra/net/Verb.java @@ -81,6 +81,7 @@ import org.apache.cassandra.service.SnapshotVerbHandler; import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.service.accord.AccordSyncPropagator; import org.apache.cassandra.service.accord.AccordSyncPropagator.Notification; +import org.apache.cassandra.service.accord.FetchTopology; import org.apache.cassandra.service.accord.FetchMinEpoch; import org.apache.cassandra.service.accord.interop.AccordInteropApply; import org.apache.cassandra.service.accord.interop.AccordInteropCommit; @@ -144,7 +145,7 @@ import org.apache.cassandra.utils.UUIDSerializer; import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.concurrent.Stage.ANTI_ENTROPY; import static org.apache.cassandra.concurrent.Stage.COUNTER_MUTATION; -import static org.apache.cassandra.concurrent.Stage.FETCH_LOG; +import static org.apache.cassandra.concurrent.Stage.FETCH_METADATA; import static org.apache.cassandra.concurrent.Stage.GOSSIP; import static org.apache.cassandra.concurrent.Stage.IMMEDIATE; import static org.apache.cassandra.concurrent.Stage.INTERNAL_METADATA; @@ -174,6 +175,7 @@ import static org.apache.cassandra.net.VerbTimeouts.repairTimeout; import static org.apache.cassandra.net.VerbTimeouts.repairValidationRspTimeout; import static org.apache.cassandra.net.VerbTimeouts.repairWithBackoffTimeout; import static org.apache.cassandra.net.VerbTimeouts.rpcTimeout; +import static org.apache.cassandra.net.VerbTimeouts.shortTimeout; import static org.apache.cassandra.net.VerbTimeouts.truncateTimeout; import static org.apache.cassandra.net.VerbTimeouts.writeTimeout; import static org.apache.cassandra.tcm.ClusterMetadataService.commitRequestHandler; @@ -288,8 +290,8 @@ public enum Verb // transactional cluster metadata TCM_COMMIT_RSP (801, P0, rpcTimeout, INTERNAL_METADATA, MessageSerializers::commitResultSerializer, RESPONSE_HANDLER ), TCM_COMMIT_REQ (802, P0, rpcTimeout, INTERNAL_METADATA, MessageSerializers::commitSerializer, () -> commitRequestHandler(), TCM_COMMIT_RSP ), - TCM_FETCH_CMS_LOG_RSP (803, P0, rpcTimeout, FETCH_LOG, MessageSerializers::logStateSerializer, RESPONSE_HANDLER ), - TCM_FETCH_CMS_LOG_REQ (804, P0, rpcTimeout, FETCH_LOG, () -> FetchCMSLog.serializer, () -> fetchLogRequestHandler(), TCM_FETCH_CMS_LOG_RSP ), + TCM_FETCH_CMS_LOG_RSP (803, P0, rpcTimeout, FETCH_METADATA, MessageSerializers::logStateSerializer, RESPONSE_HANDLER ), + TCM_FETCH_CMS_LOG_REQ (804, P0, rpcTimeout, FETCH_METADATA, () -> FetchCMSLog.serializer, () -> fetchLogRequestHandler(), TCM_FETCH_CMS_LOG_RSP ), TCM_REPLICATION (805, P0, rpcTimeout, INTERNAL_METADATA, MessageSerializers::logStateSerializer, () -> replicationHandler() ), TCM_NOTIFY_RSP (806, P0, rpcTimeout, INTERNAL_METADATA, () -> Epoch.messageSerializer, RESPONSE_HANDLER ), TCM_NOTIFY_REQ (807, P0, rpcTimeout, INTERNAL_METADATA, MessageSerializers::logStateSerializer, () -> logNotifyHandler(), TCM_NOTIFY_RSP ), @@ -299,10 +301,10 @@ public enum Verb TCM_ABORT_MIG (811, P0, rpcTimeout, INTERNAL_METADATA, () -> CMSInitializationRequest.Initiator.serializer,() -> Election.instance.abortHandler, TCM_INIT_MIG_RSP ), TCM_DISCOVER_RSP (812, P0, rpcTimeout, INTERNAL_METADATA, () -> Discovery.serializer, RESPONSE_HANDLER ), TCM_DISCOVER_REQ (813, P0, rpcTimeout, INTERNAL_METADATA, () -> NoPayload.serializer, () -> Discovery.instance.requestHandler, TCM_DISCOVER_RSP ), - TCM_FETCH_PEER_LOG_RSP (818, P0, rpcTimeout, FETCH_LOG, MessageSerializers::logStateSerializer, RESPONSE_HANDLER ), - TCM_FETCH_PEER_LOG_REQ (819, P0, rpcTimeout, FETCH_LOG, () -> FetchPeerLog.serializer, () -> FetchPeerLog.Handler.instance, TCM_FETCH_PEER_LOG_RSP ), - TCM_RECONSTRUCT_EPOCH_RSP (820, P0, rpcTimeout, FETCH_LOG, MessageSerializers::logStateSerializer, () -> ResponseVerbHandler.instance ), - TCM_RECONSTRUCT_EPOCH_REQ (821, P0, rpcTimeout, FETCH_LOG, () -> ReconstructLogState.serializer, () -> ReconstructLogState.Handler.instance, TCM_FETCH_PEER_LOG_RSP ), + TCM_FETCH_PEER_LOG_RSP (818, P0, rpcTimeout, FETCH_METADATA, MessageSerializers::logStateSerializer, RESPONSE_HANDLER ), + TCM_FETCH_PEER_LOG_REQ (819, P0, rpcTimeout, FETCH_METADATA, () -> FetchPeerLog.serializer, () -> FetchPeerLog.Handler.instance, TCM_FETCH_PEER_LOG_RSP ), + TCM_RECONSTRUCT_EPOCH_RSP (820, P0, rpcTimeout, FETCH_METADATA, MessageSerializers::logStateSerializer, () -> ResponseVerbHandler.instance ), + TCM_RECONSTRUCT_EPOCH_REQ (821, P0, rpcTimeout, FETCH_METADATA, () -> ReconstructLogState.serializer, () -> ReconstructLogState.Handler.instance, TCM_FETCH_PEER_LOG_RSP ), INITIATE_DATA_MOVEMENTS_RSP (814, P1, rpcTimeout, MISC, () -> NoPayload.serializer, RESPONSE_HANDLER ), INITIATE_DATA_MOVEMENTS_REQ (815, P1, rpcTimeout, MISC, () -> DataMovement.serializer, () -> DataMovementVerbHandler.instance, INITIATE_DATA_MOVEMENTS_RSP ), @@ -348,12 +350,12 @@ public enum Verb ACCORD_QUERY_DURABLE_BEFORE_RSP (149, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.reply, AccordService::responseHandlerOrNoop ), ACCORD_QUERY_DURABLE_BEFORE_REQ (150, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.request,AccordService::requestHandlerOrNoop, ACCORD_QUERY_DURABLE_BEFORE_RSP ), - ACCORD_SYNC_NOTIFY_RSP (151, P2, writeTimeout, IMMEDIATE, () -> EnumSerializer.simpleReply, RESPONSE_HANDLER), - ACCORD_SYNC_NOTIFY_REQ (152, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SYNC_NOTIFY_RSP ), + ACCORD_SYNC_NOTIFY_RSP (151, P2, writeTimeout, MISC, () -> EnumSerializer.simpleReply, RESPONSE_HANDLER), + ACCORD_SYNC_NOTIFY_REQ (152, P2, writeTimeout, MISC, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SYNC_NOTIFY_RSP ), ACCORD_APPLY_AND_WAIT_REQ (153, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.readData, AccordService::requestHandlerOrNoop, ACCORD_READ_RSP), - CONSENSUS_KEY_MIGRATION (154, P1, writeTimeout, MUTATION, () -> ConsensusKeyMigrationFinished.serializer,() -> ConsensusKeyMigrationState.consensusKeyMigrationFinishedHandler), + CONSENSUS_KEY_MIGRATION (154, P1, writeTimeout, MISC, () -> ConsensusKeyMigrationFinished.serializer,() -> ConsensusKeyMigrationState.consensusKeyMigrationFinishedHandler), ACCORD_INTEROP_READ_RSP (155, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.replySerializer, AccordService::responseHandlerOrNoop), ACCORD_INTEROP_READ_REQ (156, P2, writeTimeout, IMMEDIATE, () -> AccordInteropRead.requestSerializer, AccordService::requestHandlerOrNoop, ACCORD_INTEROP_READ_RSP), @@ -361,8 +363,11 @@ public enum Verb ACCORD_INTEROP_READ_REPAIR_RSP (158, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.replySerializer, AccordService::responseHandlerOrNoop), ACCORD_INTEROP_READ_REPAIR_REQ (159, P2, writeTimeout, IMMEDIATE, () -> AccordInteropReadRepair.requestSerializer, AccordService::requestHandlerOrNoop, ACCORD_INTEROP_READ_REPAIR_RSP), ACCORD_INTEROP_APPLY_REQ (160, P2, writeTimeout, IMMEDIATE, () -> AccordInteropApply.serializer, AccordService::requestHandlerOrNoop, ACCORD_APPLY_RSP), - ACCORD_FETCH_MIN_EPOCH_RSP (166, P2, writeTimeout, IMMEDIATE, () -> FetchMinEpoch.Response.serializer, RESPONSE_HANDLER), - ACCORD_FETCH_MIN_EPOCH_REQ (165, P2, writeTimeout, IMMEDIATE, () -> FetchMinEpoch.serializer, () -> FetchMinEpoch.handler, ACCORD_FETCH_MIN_EPOCH_RSP), + // TODO (desired): swap verb order to make IDS sequential? + ACCORD_FETCH_MIN_EPOCH_RSP (166, P0, shortTimeout, FETCH_METADATA, () -> FetchMinEpoch.Response.serializer, RESPONSE_HANDLER), + ACCORD_FETCH_MIN_EPOCH_REQ (165, P0, shortTimeout, FETCH_METADATA, () -> FetchMinEpoch.serializer, () -> FetchMinEpoch.handler, ACCORD_FETCH_MIN_EPOCH_RSP), + ACCORD_FETCH_TOPOLOGY_RSP (169, P0, shortTimeout, FETCH_METADATA, () -> FetchTopology.Response.serializer, RESPONSE_HANDLER), + ACCORD_FETCH_TOPOLOGY_REQ (170, P0, shortTimeout, FETCH_METADATA, () -> FetchTopology.serializer, () -> FetchTopology.handler, ACCORD_FETCH_TOPOLOGY_RSP), // generic failure response FAILURE_RSP (99, P0, noTimeout, REQUEST_RESPONSE, () -> RequestFailure.serializer, RESPONSE_HANDLER ), @@ -634,6 +639,7 @@ public enum Verb class VerbTimeouts { static final ToLongFunction rpcTimeout = DatabaseDescriptor::getRpcTimeout; + static final ToLongFunction shortTimeout = DatabaseDescriptor::getShortRpcTimeout; static final ToLongFunction writeTimeout = DatabaseDescriptor::getWriteRpcTimeout; static final ToLongFunction readTimeout = DatabaseDescriptor::getReadRpcTimeout; static final ToLongFunction rangeTimeout = DatabaseDescriptor::getRangeRpcTimeout; diff --git a/src/java/org/apache/cassandra/repair/SharedContext.java b/src/java/org/apache/cassandra/repair/SharedContext.java index 70afa9bc1d..3fdb15c920 100644 --- a/src/java/org/apache/cassandra/repair/SharedContext.java +++ b/src/java/org/apache/cassandra/repair/SharedContext.java @@ -51,6 +51,7 @@ import org.apache.cassandra.utils.TimeUUID; * * See {@link Global#instance} for the main production path */ +// TODO (required, clarity): move under Util since this is a class with shared logic public interface SharedContext { InetAddressAndPort broadcastAddressAndPort(); diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 81ce3ac8a4..6e45b6bbb3 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -187,7 +187,7 @@ public class AccordCommandStore extends CommandStore loadRangesForEpoch(journal.loadRangesForEpoch(id())); } - static Factory factory(AccordJournal journal, IntFunction executorFactory) + static Factory factory(Journal journal, IntFunction executorFactory) { return (id, node, agent, dataStore, progressLogFactory, listenerFactory, rangesForEpoch) -> new AccordCommandStore(id, node, agent, dataStore, progressLogFactory, listenerFactory, rangesForEpoch, journal, executorFactory.apply(id)); diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java index 0a530bf4d3..7188628772 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java @@ -26,6 +26,7 @@ import java.util.concurrent.TimeUnit; import accord.api.Agent; import accord.api.DataStore; +import accord.api.Journal; import accord.api.LocalListeners; import accord.api.ProgressLog; import accord.local.CommandStores; @@ -47,8 +48,8 @@ import org.apache.cassandra.service.accord.api.AccordRoutingKey; import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import static org.apache.cassandra.config.AccordSpec.QueueShardModel.THREAD_PER_SHARD; -import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueSubmissionModel; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueShardCount; +import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueSubmissionModel; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITHOUT_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.constant; @@ -60,15 +61,16 @@ public class AccordCommandStores extends CommandStores implements CacheSize private final CacheSizeMetrics cacheSizeMetrics; private final AccordExecutor[] executors; + private long cacheSize, workingSetSize; private int maxQueuedLoads, maxQueuedRangeLoads; private boolean shrinkingOn; AccordCommandStores(NodeCommandStoreService node, Agent agent, DataStore store, RandomSource random, ShardDistributor shardDistributor, ProgressLog.Factory progressLogFactory, LocalListeners.Factory listenerFactory, - AccordJournal journal, AccordExecutor[] executors) + Journal journal, AccordExecutor[] executors) { - super(node, agent, store, random, shardDistributor, progressLogFactory, listenerFactory, + super(node, agent, store, random, journal, shardDistributor, progressLogFactory, listenerFactory, AccordCommandStore.factory(journal, id -> executors[id % executors.length])); this.executors = executors; this.cacheSizeMetrics = new CacheSizeMetrics(ACCORD_STATE_CACHE, this); @@ -80,9 +82,9 @@ public class AccordCommandStores extends CommandStores implements CacheSize refreshCapacities(); } - static Factory factory(AccordJournal journal) + static Factory factory() { - return (time, agent, store, random, shardDistributor, progressLogFactory, listenerFactory) -> { + return (NodeCommandStoreService time, Agent agent, DataStore store, RandomSource random, Journal journal, ShardDistributor shardDistributor, ProgressLog.Factory progressLogFactory, LocalListeners.Factory listenersFactory) -> { AccordExecutor[] executors = new AccordExecutor[getAccordQueueShardCount()]; AccordExecutorFactory factory; int maxThreads = Integer.MAX_VALUE; @@ -119,7 +121,7 @@ public class AccordCommandStores extends CommandStores implements CacheSize } } - return new AccordCommandStores(time, agent, store, random, shardDistributor, progressLogFactory, listenerFactory, journal, executors); + return new AccordCommandStores(time, agent, store, random, shardDistributor, progressLogFactory, listenersFactory, journal, executors); }; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java index 42ac0325a8..bd1605c6bc 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordConfigurationService.java @@ -18,12 +18,11 @@ package org.apache.cassandra.service.accord; +import java.util.HashSet; +import java.util.Map; import java.util.Objects; -import java.util.OptionalLong; import java.util.Set; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; -import java.util.function.Consumer; import java.util.stream.Collectors; import javax.annotation.Nullable; import javax.annotation.concurrent.GuardedBy; @@ -31,6 +30,7 @@ import javax.annotation.concurrent.GuardedBy; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Sets; +import accord.api.Journal; import accord.impl.AbstractConfigurationService; import accord.local.Node; import accord.primitives.Ranges; @@ -49,25 +49,26 @@ import org.apache.cassandra.gms.IFailureDetector; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.MessageDelivery; import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.repair.SharedContext; import org.apache.cassandra.service.accord.AccordKeyspace.EpochDiskState; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ClusterMetadataService; -import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.listeners.ChangeListener; -import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Simulate; import org.apache.cassandra.utils.concurrent.AsyncPromise; import org.apache.cassandra.utils.concurrent.Future; +import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import static org.apache.cassandra.utils.Simulate.With.MONITORS; // TODO: listen to FailureDetector and rearrange fast path accordingly @Simulate(with=MONITORS) -public class AccordConfigurationService extends AbstractConfigurationService implements ChangeListener, AccordEndpointMapper, AccordSyncPropagator.Listener, Shutdownable +public class AccordConfigurationService extends AbstractConfigurationService implements AccordEndpointMapper, AccordSyncPropagator.Listener, Shutdownable { private final AccordSyncPropagator syncPropagator; private final DiskStateManager diskStateManager; + private final Journal journal; @GuardedBy("this") private EpochDiskState diskState = EpochDiskState.EMPTY; @@ -130,7 +131,11 @@ public class AccordConfigurationService extends AbstractConfigurationService pending, EpochDiskState diskState); EpochDiskState setCompletedLocalSync(long epoch, EpochDiskState diskState); @@ -151,7 +156,7 @@ public class AccordConfigurationService extends AbstractConfigurationService {}); - } - - public synchronized void start(Consumer callback) { Invariants.checkState(state == State.INITIALIZED, "Expected state to be INITIALIZED but was %s", state); state = State.LOADING; + EndpointMapping snapshot = mapping; - //TODO (restart): if there are topologies loaded then there is likely failures if reporting is needed, as mapping is not setup yet - AtomicReference previousRef = new AtomicReference<>(null); - diskState = diskStateManager.loadTopologies(((epoch, metadata, topology, syncStatus, pendingSyncNotify, remoteSyncComplete, closed, redundant) -> { - updateMapping(metadata); - reportTopology(topology, syncStatus == SyncStatus.NOT_STARTED, true); - - Topology previous = previousRef.get(); - if (previous != null) - { - // for all nodes removed, or pending removal, mark them as removed so we don't wait on their replies - Sets.SetView removedNodes = Sets.difference(previous.nodes(), topology.nodes()); - if (!removedNodes.isEmpty()) - onNodesRemoved(topology.epoch(), currentTopology(), removedNodes); - } - previousRef.set(topology); - + diskStateManager.loadLocalTopologyState((epoch, syncStatus, pendingSyncNotify, remoteSyncComplete, closed, redundant) -> { getOrCreateEpochState(epoch).setSyncStatus(syncStatus); - if (syncStatus == SyncStatus.NOTIFYING) - { - // 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 + // 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 + if (Objects.requireNonNull(syncStatus) == SyncStatus.NOTIFYING) syncPropagator.reportSyncComplete(epoch, Sets.filter(pendingSyncNotify, snapshot::containsId), localId); - } remoteSyncComplete.forEach(id -> receiveRemoteSyncComplete(id, epoch)); // TODO (required): 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 receiveClosed(closed, epoch); receiveRedundant(redundant, epoch); - })); + }); state = State.STARTED; - callback.accept(diskState.isEmpty() ? OptionalLong.empty() : OptionalLong.of(diskState.maxEpoch)); - ClusterMetadataService.instance().log().addListener(this); + + // for all nodes removed, or pending removal, mark them as removed, so we don't wait on their replies + Map removedNodes = mapping.removedNodes(); + for (Map.Entry e : removedNodes.entrySet()) + onNodeRemoved(e.getValue(), currentTopology(), e.getKey()); + + ClusterMetadataService.instance().log().addListener(listener); } @Override @@ -272,7 +274,7 @@ public class AccordConfigurationService extends AbstractConfigurationService removedNodes = Sets.difference(previous.nodes(), topology.nodes()); - if (!removedNodes.isEmpty()) - onNodesRemoved(topology.epoch(), previous, removedNodes); + // TODO (desired, efficiency): there should be no need to notify every epoch for every removed node + for (Node.Id removedNode : removedNodes) + { + if (topology.epoch() >= epochs.minEpoch()) + onNodeRemoved(topology.epoch(), previous, removedNode); + } } private static boolean shareShard(Topology current, Node.Id target, Node.Id self) @@ -366,22 +367,17 @@ public class AccordConfigurationService extends AbstractConfigurationService removed) + public void onNodeRemoved(long epoch, Topology current, Node.Id removed) { - if (removed.isEmpty()) return; syncPropagator.onNodesRemoved(removed); + // TODO (now): it seems to be incorrect to mark remote syncs complete if/when node got removed. for (long oldEpoch : nonCompletedEpochsBefore(epoch)) - { - for (Node.Id node : removed) - receiveRemoteSyncCompletePreListenerNotify(node, oldEpoch); - } - listeners.forEach(l -> l.onRemoveNodes(epoch, removed)); + receiveRemoteSyncCompletePreListenerNotify(removed, oldEpoch); - for (Node.Id node : removed) - { - if (shareShard(current, node, localId)) - AccordService.instance().tryMarkRemoved(current, node); - } + listeners.forEach(l -> l.onRemoveNode(epoch, removed)); + + if (shareShard(current, removed, localId)) + AccordService.instance().tryMarkRemoved(current, removed); } private long[] nonCompletedEpochsBefore(long max) @@ -416,35 +412,32 @@ public class AccordConfigurationService extends AbstractConfigurationService reportMetadata(metadata)); } - @Override - public void notifyPostCommit(ClusterMetadata prev, ClusterMetadata next, boolean fromSnapshot) - { - maybeReportMetadata(next); - } - @Override protected void fetchTopologyInternal(long epoch) { - ClusterMetadata metadata = ClusterMetadata.current(); - if (metadata.directory.peerIds().isEmpty()) - return; // just let CMS handle it when it's ready - - // TODO (desired): randomise - NodeId first = metadata.directory.peerIds().first(); - InetAddressAndPort peer = metadata.directory.getNodeAddresses(first).broadcastAddress; - if (FBUtilities.getBroadcastAddressAndPort().equals(peer)) + try { - NodeId second = metadata.directory.peerIds().higher(first); - if (second == null) + Set peers = new HashSet<>(ClusterMetadata.current().directory.allJoinedEndpoints()); + peers.remove(FBUtilities.getBroadcastAddressAndPort()); + if (peers.isEmpty()) return; - - peer = metadata.directory.getNodeAddresses(second).broadcastAddress; + Topology topology; + while ((topology =FetchTopology.fetch(SharedContext.Global.instance, peers, epoch).get()) == null) {} + reportTopology(topology); + } + catch (InterruptedException e) + { + if (currentEpoch() >= epoch) + return; + Thread.currentThread().interrupt(); + throw new UncheckedInterruptedException(e); + } + catch (Throwable e) + { + if (currentEpoch() >= epoch) + return; + throw new RuntimeException(e.getCause()); } - ClusterMetadataService.instance().fetchLogFromPeerOrCMSAsync(metadata, peer, Epoch.create(epoch)) - .addCallback((success, fail) -> { - if (fail != null) - fetchTopologyInternal(epoch); - }); } @Override @@ -569,7 +562,7 @@ public class AccordConfigurationService extends AbstractConfigurationService result) + static ResultStatus of(AsyncResult result) { if (result == null || !result.isDone()) return PENDING; diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournal.java b/src/java/org/apache/cassandra/service/accord/AccordJournal.java index a17c3528f5..f82b0b3e4d 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournal.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournal.java @@ -20,6 +20,7 @@ package org.apache.cassandra.service.accord; import java.io.IOException; import java.nio.ByteBuffer; import java.util.Collections; +import java.util.Iterator; import java.util.NavigableMap; import java.util.Set; import java.util.concurrent.TimeUnit; @@ -240,28 +241,28 @@ public class AccordJournal implements accord.api.Journal, Shutdownable @Override public RedundantBefore loadRedundantBefore(int store) { - IdentityAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.REDUNDANT_BEFORE, store)); + IdentityAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.REDUNDANT_BEFORE, store), false); return accumulator.get(); } @Override public NavigableMap loadBootstrapBeganAt(int store) { - IdentityAccumulator> accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.BOOTSTRAP_BEGAN_AT, store)); + IdentityAccumulator> accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.BOOTSTRAP_BEGAN_AT, store), false); return accumulator.get(); } @Override public NavigableMap loadSafeToRead(int store) { - IdentityAccumulator> accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.SAFE_TO_READ, store)); + IdentityAccumulator> accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.SAFE_TO_READ, store), false); return accumulator.get(); } @Override public CommandStores.RangesForEpoch loadRangesForEpoch(int store) { - IdentityAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.RANGES_FOR_EPOCH, store)); + IdentityAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.RANGES_FOR_EPOCH, store), false); return accumulator.get(); } @@ -282,6 +283,21 @@ public class AccordJournal implements accord.api.Journal, Shutdownable journal.onDurable(pointer, onFlush); } + @Override + public Iterator replayTopologies() + { + AccordJournalValueSerializers.MapAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.TOPOLOGY_UPDATE, -1), false); + return accumulator.get().values().iterator(); + } + + @Override + public void saveTopology(TopologyUpdate topologyUpdate, Runnable onFlush) + { + RecordPointer pointer = appendInternal(new JournalKey(TxnId.NONE, JournalKey.Type.TOPOLOGY_UPDATE, -1), topologyUpdate); + if (onFlush != null) + journal.onDurable(pointer, onFlush); + } + @Override public PersistentField.Persister durableBeforePersister() { @@ -301,7 +317,7 @@ public class AccordJournal implements accord.api.Journal, Shutdownable @Override public DurableBefore load() { - DurableBeforeAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.DURABLE_BEFORE, 0)); + DurableBeforeAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.DURABLE_BEFORE, 0), false); return accumulator.get(); } }; @@ -334,7 +350,7 @@ public class AccordJournal implements accord.api.Journal, Shutdownable { JournalKey key = new JournalKey(txnId, JournalKey.Type.COMMAND_DIFF, commandStoreId); Builder builder = new Builder(txnId, load); - journalTable.readAll(key, builder::deserializeNext); + journalTable.readAll(key, builder::deserializeNext, true); return builder; } @@ -344,13 +360,13 @@ public class AccordJournal implements accord.api.Journal, Shutdownable return loadDiffs(commandStoreId, txnId, Load.ALL); } - private BUILDER readAll(JournalKey key) + private BUILDER readAll(JournalKey key, boolean asc) { BUILDER builder = (BUILDER) key.type.serializer.mergerFor(key); // TODO: this can be further improved to avoid allocating lambdas AccordJournalValueSerializers.FlyweightSerializer serializer = (AccordJournalValueSerializers.FlyweightSerializer) key.type.serializer; // TODO (expected): for those where we store an image, read only the first entry we find in DESC order - journalTable.readAll(key, (in, userVersion) -> serializer.deserialize(key, builder, in, userVersion)); + journalTable.readAll(key, (in, userVersion) -> serializer.deserialize(key, builder, in, userVersion), asc); return builder; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java index b8d6f8ca56..891abcf1df 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java @@ -164,10 +164,10 @@ public class AccordJournalTable *

* When reading from journal segments, skip descriptors that were read from the table. */ - public void readAll(K key, Reader reader) + public void readAll(K key, Reader reader, boolean asc) { JournalAndTableRecordConsumer consumer = new JournalAndTableRecordConsumer(key, reader); - journal.readAll(key, consumer); + journal.readAll(key, consumer, asc); consumer.readTable(); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournalValueSerializers.java b/src/java/org/apache/cassandra/service/accord/AccordJournalValueSerializers.java index a11dfb744b..eb60d901df 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournalValueSerializers.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournalValueSerializers.java @@ -19,21 +19,28 @@ package org.apache.cassandra.service.accord; import java.io.IOException; +import java.util.Map; import java.util.NavigableMap; +import java.util.TreeMap; +import java.util.function.Function; import com.google.common.collect.ImmutableSortedMap; +import accord.api.Journal; import accord.local.DurableBefore; import accord.local.RedundantBefore; import accord.primitives.Ranges; import accord.primitives.Timestamp; import accord.primitives.TxnId; +import accord.topology.Topology; import accord.utils.Invariants; +import org.agrona.collections.Int2ObjectHashMap; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.service.accord.serializers.CommandStoreSerializers; import org.apache.cassandra.service.accord.serializers.KeySerializers; +import org.apache.cassandra.service.accord.serializers.TopologySerializers; import static accord.api.Journal.Load.ALL; import static accord.local.CommandStores.RangesForEpoch; @@ -194,7 +201,8 @@ public class AccordJournalValueSerializers } } - public static class DurableBeforeSerializer implements FlyweightSerializer + public static class DurableBeforeSerializer + implements FlyweightSerializer { public DurableBeforeAccumulator mergerFor(JournalKey journalKey) { @@ -257,7 +265,8 @@ public class AccordJournalValueSerializers } } - public static class SafeToReadSerializer implements FlyweightSerializer, IdentityAccumulator>> + public static class SafeToReadSerializer + implements FlyweightSerializer, IdentityAccumulator>> { @Override public IdentityAccumulator> mergerFor(JournalKey key) @@ -285,9 +294,8 @@ public class AccordJournalValueSerializers } public static class RangesForEpochSerializer - implements FlyweightSerializer> + implements FlyweightSerializer> { - public IdentityAccumulator mergerFor(JournalKey key) { return new IdentityAccumulator<>(null); @@ -311,13 +319,13 @@ public class AccordJournalValueSerializers } @Override - public void reserialize(JournalKey key, IdentityAccumulator from, DataOutputPlus out, int userVersion) throws IOException + public void reserialize(JournalKey key, Accumulator from, DataOutputPlus out, int userVersion) throws IOException { serialize(key, from.get(), out, messagingVersion); } @Override - public void deserialize(JournalKey key, IdentityAccumulator into, DataInputPlus in, int userVersion) throws IOException + public void deserialize(JournalKey key, Accumulator into, DataInputPlus in, int userVersion) throws IOException { int size = in.readUnsignedVInt32(); Ranges[] ranges = new Ranges[size]; @@ -331,4 +339,91 @@ public class AccordJournalValueSerializers into.update(new RangesForEpoch(epochs, ranges)); } } + + public static class MapAccumulator, V> extends Accumulator, V> + { + private final Function getKey; + + public MapAccumulator(Function getKey) + { + super(new TreeMap<>()); + this.getKey = getKey; + } + + @Override + protected NavigableMap accumulate(NavigableMap accumulator, V newValue) + { + V prev = accumulator.put(getKey.apply(newValue), newValue); + Invariants.checkState(prev == null || prev.equals(newValue)); + return accumulator; + } + } + + public static class TopologyUpdateSerializer + implements FlyweightSerializer> + { + private final RangesForEpochSerializer rangesForEpochSerializer = new RangesForEpochSerializer(); + + @Override + public MapAccumulator mergerFor(JournalKey key) + { + return new MapAccumulator<>(topologyUpdate -> topologyUpdate.global.epoch()); + } + + @Override + public void serialize(JournalKey key, Journal.TopologyUpdate from, DataOutputPlus out, int userVersion) throws IOException + { + out.writeInt(1); + serializeOne(key, from, out, userVersion); + } + + private void serializeOne(JournalKey key, Journal.TopologyUpdate from, DataOutputPlus out, int userVersion) throws IOException + { + out.writeInt(from.commandStores.size()); + for (Map.Entry e : from.commandStores.entrySet()) + { + out.writeInt(e.getKey()); + rangesForEpochSerializer.serialize(key, e.getValue(), out, userVersion); + } + TopologySerializers.topology.serialize(from.local, out, userVersion); + TopologySerializers.topology.serialize(from.global, out, userVersion); + } + + @Override + public void reserialize(JournalKey key, MapAccumulator from, DataOutputPlus out, int userVersion) throws IOException + { + out.writeInt(from.accumulated.size()); + for (Journal.TopologyUpdate update : from.accumulated.values()) + serializeOne(key, update, out, userVersion); + } + + @Override + public void deserialize(JournalKey key, MapAccumulator into, DataInputPlus in, int userVersion) throws IOException + { + int size = in.readInt(); + Accumulator acc = new Accumulator<>(null) + { + @Override + protected RangesForEpoch accumulate(RangesForEpoch oldValue, RangesForEpoch newValue) + { + return this.accumulated = newValue; + } + }; + for (int i = 0; i < size; i++) + { + int commandStoresSize = in.readInt(); + Int2ObjectHashMap commandStores = new Int2ObjectHashMap<>(); + for (int j = 0; j < commandStoresSize; j++) + { + acc.update(null); + int commandStoreId = in.readInt(); + rangesForEpochSerializer.deserialize(key, acc, in, userVersion); + commandStores.put(commandStoreId, acc.accumulated); + } + Topology local = TopologySerializers.topology.deserialize(in, userVersion); + Topology global = TopologySerializers.topology.deserialize(in, userVersion); + into.update(new Journal.TopologyUpdate(commandStores, local, global)); + } + } + } } \ No newline at end of file diff --git a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java index bcad001c4c..5dc7d240c5 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java +++ b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java @@ -137,7 +137,6 @@ import org.apache.cassandra.service.accord.api.AccordRoutingKey.TokenKey; import org.apache.cassandra.service.accord.serializers.AccordRoutingKeyByteSource; import org.apache.cassandra.service.accord.serializers.CommandSerializers; import org.apache.cassandra.service.accord.serializers.KeySerializers; -import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.utils.Clock.Global; import org.apache.cassandra.utils.CloseableIterator; import org.apache.cassandra.utils.btree.BTree; @@ -1309,6 +1308,7 @@ public class AccordKeyspace return diskState; } + // TODO (required): unused public static EpochDiskState markRedundant(Ranges ranges, long epoch, EpochDiskState diskState) { diskState = maybeUpdateMaxEpoch(diskState, epoch); @@ -1369,22 +1369,17 @@ public class AccordKeyspace public interface TopologyLoadConsumer { - void load(long epoch, ClusterMetadata metadata, Topology topology, SyncStatus syncStatus, Set pendingSyncNotify, Set remoteSyncComplete, Ranges closed, Ranges redundant); + void load(long epoch, SyncStatus syncStatus, Set pendingSyncNotify, Set remoteSyncComplete, Ranges closed, Ranges redundant); } @VisibleForTesting - public static void loadEpoch(long epoch, ClusterMetadata metadata, TopologyLoadConsumer consumer) throws IOException + public static void loadEpoch(long epoch, TopologyLoadConsumer consumer) throws IOException { - Topology topology = AccordTopology.createAccordTopology(metadata); - - String cql = "SELECT * FROM " + ACCORD_KEYSPACE_NAME + '.' + TOPOLOGIES + ' ' + - "WHERE epoch=?"; - - UntypedResultSet result = executeInternal(cql, epoch); + UntypedResultSet result = executeInternal(String.format("SELECT * FROM %s.%s WHERE epoch=?", ACCORD_KEYSPACE_NAME, TOPOLOGIES), epoch); if (result.isEmpty()) { // topology updates disk state for epoch but doesn't save the topology to the table, so there maybe an epoch we know about, but no fields are present - consumer.load(epoch, metadata, topology, SyncStatus.NOT_STARTED, Collections.emptySet(), Collections.emptySet(), Ranges.EMPTY, Ranges.EMPTY); + consumer.load(epoch, SyncStatus.NOT_STARTED, Collections.emptySet(), Collections.emptySet(), Ranges.EMPTY, Ranges.EMPTY); return; } checkState(!result.isEmpty(), "Nothing found for epoch %d", epoch); @@ -1402,7 +1397,7 @@ public class AccordKeyspace Ranges closed = row.has("closed") ? blobMapToRanges(row.getMap("closed", BytesType.instance, BytesType.instance)) : Ranges.EMPTY; Ranges redundant = row.has("redundant") ? blobMapToRanges(row.getMap("redundant", BytesType.instance, BytesType.instance)) : Ranges.EMPTY; - consumer.load(epoch, metadata, topology, syncStatus, pendingSyncNotify, remoteSyncComplete, closed, redundant); + consumer.load(epoch, syncStatus, pendingSyncNotify, remoteSyncComplete, closed, redundant); } public static EpochDiskState loadTopologies(TopologyLoadConsumer consumer) @@ -1413,8 +1408,8 @@ public class AccordKeyspace if (diskState == null) return EpochDiskState.EMPTY; - for (ClusterMetadata metadata : AccordService.tcmLoadRange(diskState.minEpoch, diskState.maxEpoch)) - loadEpoch(metadata.epoch.getEpoch(), metadata, consumer); + for (long epoch = diskState.minEpoch; epoch < diskState.maxEpoch; epoch++) + loadEpoch(epoch, consumer); return diskState; } diff --git a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java index 709299e6e8..930afee54f 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java @@ -67,6 +67,8 @@ public class AccordSegmentCompactor implements SegmentCompactor reader = segment.keyOrderReader(); if (reader.advance()) readers.add(reader); + else + reader.close(); } // nothing to compact (all segments empty, should never happen, but it is theoretically possible?) - exit early @@ -126,6 +128,7 @@ public class AccordSegmentCompactor implements SegmentCompactor callbacks, DefaultProgressLogs::new, DefaultLocalListeners.Factory::new, - AccordCommandStores.factory(journal), + AccordCommandStores.factory(), new AccordInteropFactory(agent, configService), journal.durableBeforePersister(), - configuration); + configuration, + journal); this.nodeShutdown = toShutdownable(node); this.requestHandler = new AccordVerbHandler<>(node, configService); this.responseHandler = new AccordResponseVerbHandler<>(callbacks, configService); @@ -351,50 +367,36 @@ public class AccordService implements IAccordService, Shutdownable return; journal.start(node); node.load(); - ClusterMetadataService cms = ClusterMetadataService.instance(); - class Ref { List historic = Collections.emptyList();} - Ref ref = new Ref(); - configService.start((optMaxEpoch -> { - List historic = ref.historic = !optMaxEpoch.isEmpty() - ? tcmLoadRange(optMaxEpoch.getAsLong(), Long.MAX_VALUE) - : discoverHistoric(node, cms); - for (ClusterMetadata m : historic) - configService.reportMetadataInternal(m, true); - })); - ClusterMetadata current = cms.metadata(); - if (!ref.historic.isEmpty()) - { - List historic = ref.historic; - long lastHistoric = ref.historic.get(historic.size() - 1).epoch.getEpoch(); - if (lastHistoric + 1 < current.epoch.getEpoch()) - { - // new epochs added while loading... load the deltas - for (ClusterMetadata metadata : tcmLoadRange(lastHistoric + 1, current.epoch.getEpoch())) - { - historic.add(metadata); - configService.reportMetadataInternal(metadata); - } - } - // sync doesn't happen when this node isn't in the epoch - //TODO (now, correctness): sync points use "closed" and not "syncComplete", so need to call TM.epochRedundant or TM.onEpochClosed - // epochRedundant implies all txn that have been proposed for this epoch have been executed "globally" - we don't have this knowlege - // epochClosed implies no "new" txn can be proposed - for (ClusterMetadata m : historic) + ClusterMetadata metadata = ClusterMetadata.current(); + configService.updateMapping(metadata); + + // Load all active topologies, wihout writing them to journal again. No-op on bootstrap. + node.commandStores().restoreShardStateUnsafe(topology -> configService.reportTopology(topology, true, true)); + configService.start(); + + long minEpoch = fetchMinEpoch(); + if (minEpoch >= 0) + { + for (long epoch = minEpoch; epoch <= metadata.epoch.getEpoch(); epoch++) + node.configService().fetchTopologyForEpoch(epoch); + + try { - Topology t = AccordTopology.createAccordTopology(m); - long epoch = t.epoch(); - for (Id id : t.nodes()) - node.onRemoteSyncComplete(id, epoch); - //TODO (correctness): is this true? - node.onEpochClosed(t.ranges(), t.epoch()); - node.onEpochRedundant(t.ranges(), t.epoch()); + epochReady(metadata.epoch).get(DatabaseDescriptor.getTransactionTimeout(MILLISECONDS), MILLISECONDS); + } + catch (InterruptedException e) + { + throw new UncheckedInterruptedException(e); + } + catch (ExecutionException | TimeoutException e) + { + throw new RuntimeException(e); } } - configService.reportMetadataInternal(current); fastPathCoordinator.start(); - cms.log().addListener(fastPathCoordinator); + ClusterMetadataService.instance().log().addListener(fastPathCoordinator); node.durabilityScheduling().setDefaultRetryDelay(Ints.checkedCast(DatabaseDescriptor.getAccordDefaultDurabilityRetryDelay(SECONDS)), SECONDS); node.durabilityScheduling().setMaxRetryDelay(Ints.checkedCast(DatabaseDescriptor.getAccordMaxDurabilityRetryDelay(SECONDS)), SECONDS); node.durabilityScheduling().setTargetShardSplits(Ints.checkedCast(DatabaseDescriptor.getAccordShardDurabilityTargetSplits())); @@ -405,25 +407,23 @@ public class AccordService implements IAccordService, Shutdownable state = State.STARTED; } - private List discoverHistoric(Node node, ClusterMetadataService cms) + /** + * Queries peers to discover min epoch + */ + private long fetchMinEpoch() { - ClusterMetadata current = cms.metadata(); - Topology topology = AccordTopology.createAccordTopology(current); - Ranges localRanges = topology.rangesForNode(node.id()); - if (!localRanges.isEmpty()) // already joined, nothing to see here - return Collections.emptyList(); - + ClusterMetadata metadata = ClusterMetadata.current(); Map> peers = new HashMap<>(); - for (KeyspaceMetadata keyspace : current.schema.getKeyspaces()) + for (KeyspaceMetadata keyspace : metadata.schema.getKeyspaces()) { List tables = keyspace.tables.stream().filter(TableMetadata::requiresAccordSupport).collect(Collectors.toList()); if (tables.isEmpty()) continue; - DataPlacement placement = current.placements.get(keyspace.params.replication); - DataPlacement whenSettled = current.writePlacementAllSettled(keyspace); - Sets.SetView alive = Sets.intersection(whenSettled.writes.byEndpoint().keySet(), placement.writes.byEndpoint().keySet()); + DataPlacement current = metadata.placements.get(keyspace.params.replication); + DataPlacement settled = metadata.writePlacementAllSettled(keyspace); + Sets.SetView alive = Sets.intersection(settled.writes.byEndpoint().keySet(), current.writes.byEndpoint().keySet()); InetAddressAndPort self = FBUtilities.getBroadcastAddressAndPort(); - whenSettled.writes.forEach((range, group) -> { + settled.writes.forEach((range, group) -> { if (group.endpoints().contains(self)) { for (InetAddressAndPort peer : group.endpoints()) @@ -436,35 +436,12 @@ public class AccordService implements IAccordService, Shutdownable }); } if (peers.isEmpty()) - return Collections.emptyList(); + return -1; Long minEpoch = findMinEpoch(SharedContext.Global.instance, peers); if (minEpoch == null) - return Collections.emptyList(); - return tcmLoadRange(minEpoch, current.epoch.getEpoch()); - } - - public static List tcmLoadRange(long min, long max) - { - List afterLoad = reconstruct(min, max); - - if (Invariants.isParanoid()) - Invariants.checkState(afterLoad.get(0).epoch.getEpoch() == min, "Unexpected epoch: expected %d but given %d", min, afterLoad.get(0).epoch.getEpoch()); - while (!afterLoad.isEmpty() && afterLoad.get(0).epoch.getEpoch() < min) - afterLoad.remove(0); - Invariants.checkState(!afterLoad.isEmpty(), "TCM was unable to return the needed epochs: %d -> %d", min, max); - Invariants.checkState(afterLoad.get(0).epoch.getEpoch() == min, "Unexpected epoch: expected %d but given %d", min, afterLoad.get(0).epoch.getEpoch()); - Invariants.checkState(max == Long.MAX_VALUE || afterLoad.get(afterLoad.size() - 1).epoch.getEpoch() == max, "Unexpected epoch: expected %d but given %d", max, afterLoad.get(afterLoad.size() - 1).epoch.getEpoch()); - return afterLoad; - } - - private static List reconstruct(long min, long max) - { - Epoch start = Epoch.create(min); - Epoch end = Epoch.create(max); - Retry.Deadline deadline = Retry.Deadline.wrap(new Retry.ExponentialBackoff(TCMMetrics.instance.fetchLogRetries)); - return ClusterMetadataService.instance().processor() - .reconstruct(start, end, deadline); + return -1; + return minEpoch; } @VisibleForTesting diff --git a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java index 6f1e479f91..b7f74257cb 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java @@ -232,25 +232,33 @@ public class AccordSyncPropagator '}'; } - public synchronized void onNodesRemoved(Set removed) + public void onNodesRemoved(Node.Id removed) { - for (Node.Id node : removed) + long[] toAck; + boolean[] syncCompletedFor; + + synchronized (AccordSyncPropagator.this) { - PendingEpochs pendingEpochs = pending.get(node.id); - if (pendingEpochs == null) continue; - long[] toComplete = new long[pendingEpochs.size()]; + PendingEpochs pendingEpochs = pending.remove(removed.id); + if (pendingEpochs == null) return; + toAck = new long[pendingEpochs.size()]; + syncCompletedFor = new boolean[pendingEpochs.size()]; Long2ObjectHashMap.KeyIterator it = pendingEpochs.keySet().iterator(); for (int i = 0; it.hasNext(); i++) - toComplete[i] = it.nextLong(); - Arrays.sort(toComplete); - for (long epoch : toComplete) - listener.onEndpointAck(node, epoch); - pending.remove(node.id); - for (long epoch : toComplete) { - if (hasSyncCompletedFor(epoch)) - listener.onComplete(epoch); + long epoch = it.nextLong(); + toAck[i] = epoch; + syncCompletedFor[i] = hasSyncCompletedFor(epoch); } + Arrays.sort(toAck); + } + + for (int i = 0; i < toAck.length; i++) + { + long epoch = toAck[i]; + listener.onEndpointAck(removed, epoch); + if (syncCompletedFor[i]) + listener.onComplete(epoch); } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordTopology.java b/src/java/org/apache/cassandra/service/accord/AccordTopology.java index db89c4fda7..fc4e582ea9 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTopology.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTopology.java @@ -271,7 +271,7 @@ public class AccordTopology Directory directory, AccordFastPath accordFastPath, ShardLookup lookup, AccordStaleReplicas staleReplicas) { - List shards = new ArrayList<>(); + List res = new ArrayList<>(); Set unavailable = accordFastPath.unavailableIds(); Map dcMap = createDCMap(directory); @@ -281,12 +281,11 @@ public class AccordTopology if (tables.isEmpty()) continue; List ksShards = KeyspaceShard.forKeyspace(keyspace, placements, directory); - tables.forEach(table -> ksShards.forEach(shard -> shards.addAll(shard.createForTable(table, unavailable, dcMap, lookup)))); + tables.forEach(table -> ksShards.forEach(shard -> res.addAll(shard.createForTable(table, unavailable, dcMap, lookup)))); } - shards.sort((a, b) -> a.range.compare(b.range)); - - return new Topology(epoch.getEpoch(), SortedArrayList.copyUnsorted(staleReplicas.ids(), Id[]::new), shards.toArray(new Shard[0])); + res.sort((a, b) -> a.range.compare(b.range)); + return new Topology(epoch.getEpoch(), SortedArrayList.copyUnsorted(staleReplicas.ids(), Id[]::new), res.toArray(new Shard[0])); } public static Topology createAccordTopology(ClusterMetadata metadata, ShardLookup lookup) @@ -304,16 +303,16 @@ public class AccordTopology return createAccordTopology(metadata, (Topology) null); } - public static EndpointMapping directoryToMapping(EndpointMapping mapping, long epoch, Directory directory) + public static EndpointMapping directoryToMapping(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 (Id id : mapping.differenceIds(builder)) - builder.add(mapping.mappedEndpoint(id), 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 (Directory.RemovedNode removedNode : directory.removedNodes()) + builder.add(removedNode.endpoint, tcmIdToAccord(removedNode.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 916cea43bf..5745c55112 100644 --- a/src/java/org/apache/cassandra/service/accord/EndpointMapping.java +++ b/src/java/org/apache/cassandra/service/accord/EndpointMapping.java @@ -18,12 +18,12 @@ package org.apache.cassandra.service.accord; -import java.util.Set; +import java.util.Map; 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 com.google.common.collect.ImmutableMap; import accord.local.Node; import accord.utils.Invariants; @@ -31,15 +31,18 @@ import org.apache.cassandra.locator.InetAddressAndPort; class EndpointMapping implements AccordEndpointMapper { - public static final EndpointMapping EMPTY = new EndpointMapping(0, ImmutableBiMap.of()); + public static final EndpointMapping EMPTY = new EndpointMapping(0, ImmutableBiMap.of(), ImmutableMap.of()); private final long epoch; private final ImmutableBiMap mapping; + private final ImmutableMap removedNodes; private EndpointMapping(long epoch, - ImmutableBiMap mapping) + ImmutableBiMap mapping, + ImmutableMap removedNodes) { this.epoch = epoch; this.mapping = mapping; + this.removedNodes = removedNodes; } long epoch() @@ -52,9 +55,9 @@ class EndpointMapping implements AccordEndpointMapper return mapping.containsKey(id); } - public Set differenceIds(Builder builder) + public Map removedNodes() { - return Sets.difference(mapping.keySet(), builder.mapping.keySet()); + return removedNodes; } @Override @@ -73,6 +76,7 @@ class EndpointMapping implements AccordEndpointMapper { private final long epoch; private final BiMap mapping = HashBiMap.create(); + private final ImmutableMap.Builder removed = new ImmutableMap.Builder<>(); public Builder(long epoch) { @@ -87,9 +91,18 @@ class EndpointMapping implements AccordEndpointMapper return this; } + public Builder removed(InetAddressAndPort endpoint, Node.Id id, long epoch) + { + Invariants.checkArgument(!mapping.containsKey(id), "Mapping already exists for Node.Id %s", id); + Invariants.checkArgument(!mapping.containsValue(endpoint), "Mapping already exists for %s", endpoint); + mapping.put(id, endpoint); + removed.put(id, epoch); + return this; + } + public EndpointMapping build() { - return new EndpointMapping(epoch, ImmutableBiMap.copyOf(mapping)); + return new EndpointMapping(epoch, ImmutableBiMap.copyOf(mapping), removed.build()); } } diff --git a/src/java/org/apache/cassandra/service/accord/FetchMinEpoch.java b/src/java/org/apache/cassandra/service/accord/FetchMinEpoch.java index 0e40afc6c1..ef670c572e 100644 --- a/src/java/org/apache/cassandra/service/accord/FetchMinEpoch.java +++ b/src/java/org/apache/cassandra/service/accord/FetchMinEpoch.java @@ -25,7 +25,6 @@ import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Set; - import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; @@ -33,12 +32,13 @@ import com.google.common.collect.Iterators; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.TypeSizes; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.IVerbHandler; -import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.MessageDelivery; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.net.Verb; import org.apache.cassandra.repair.SharedContext; @@ -46,11 +46,15 @@ import org.apache.cassandra.utils.Backoff; import org.apache.cassandra.utils.concurrent.Future; import org.apache.cassandra.utils.concurrent.FutureCombiner; +import static org.apache.cassandra.net.MessageDelivery.RetryErrorMessage; +import static org.apache.cassandra.net.MessageDelivery.RetryPredicate; +import static org.apache.cassandra.net.MessageDelivery.logger; + +// TODO (required, efficiency): this can be simplified: we seem to always use "entire range" public class FetchMinEpoch { public static final IVersionedSerializer serializer = new IVersionedSerializer<>() { - @Override public void serialize(FetchMinEpoch t, DataOutputPlus out, int version) throws IOException { @@ -78,15 +82,20 @@ public class FetchMinEpoch return size; } }; - public static final IVerbHandler handler = new IVerbHandler() - { - @Override - public void doVerb(Message message) throws IOException + + public static final IVerbHandler handler = message -> { + if (AccordService.started()) { Long epoch = AccordService.instance().minEpoch(message.payload.ranges); MessagingService.instance().respond(new Response(epoch), message); } + else + { + logger.error("Accord service is not started, resopnding with error to {}", message); + MessagingService.instance().respondWithFailure(RequestFailure.BOOTING, message); + } }; + public final Collection ranges; public FetchMinEpoch(Collection ranges) @@ -122,13 +131,14 @@ public class FetchMinEpoch List> accum = new ArrayList<>(peers.size()); for (Map.Entry> e : peers.entrySet()) accum.add(fetch(context, e.getKey(), e.getValue())); - return FutureCombiner.successfulOf(accum).map(ls -> { + // TODO (required): we are collecting only successes, but we need some threshold + return FutureCombiner.successfulOf(accum).map(epochs -> { Long min = null; - for (Long l : ls) + for (Long epoch : epochs) { - if (l == null) continue; - if (min == null) min = l; - else min = Math.min(min, l); + if (epoch == null) continue; + if (min == null) min = epoch; + else min = Math.min(min, epoch); } return min; }); @@ -138,12 +148,12 @@ public class FetchMinEpoch static Future fetch(SharedContext context, InetAddressAndPort to, Set value) { FetchMinEpoch req = new FetchMinEpoch(value); - Backoff backoff = Backoff.fromConfig(context, DatabaseDescriptor.getAccord().minEpochSyncRetry); - return context.messaging().sendWithRetries(backoff, context.optionalTasks()::schedule, + return context.messaging().sendWithRetries(Backoff.NO_OP.INSTANCE, + MessageDelivery.ImmediateRetryScheduler.instance, Verb.ACCORD_FETCH_MIN_EPOCH_REQ, req, Iterators.cycle(to), - (i1, i2, i3) -> true, - (i1, i2, i3, i4) -> null) + RetryPredicate.times(DatabaseDescriptor.getAccord().minEpochSyncRetry.maxAttempts.value), + RetryErrorMessage.EMPTY) .map(m -> m.payload.minEpoch); } @@ -206,4 +216,4 @@ public class FetchMinEpoch '}'; } } -} +} \ No newline at end of file diff --git a/src/java/org/apache/cassandra/service/accord/FetchTopology.java b/src/java/org/apache/cassandra/service/accord/FetchTopology.java new file mode 100644 index 0000000000..d40d2654ea --- /dev/null +++ b/src/java/org/apache/cassandra/service/accord/FetchTopology.java @@ -0,0 +1,140 @@ +/* + * 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; + +import java.io.IOException; +import java.util.Collection; + +import accord.topology.Topology; +import org.apache.cassandra.io.IVersionedSerializer; +import org.apache.cassandra.io.util.DataInputPlus; +import org.apache.cassandra.io.util.DataOutputPlus; +import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.net.IVerbHandler; +import org.apache.cassandra.net.MessageDelivery; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.net.MessagingUtils; +import org.apache.cassandra.net.Verb; +import org.apache.cassandra.repair.SharedContext; +import org.apache.cassandra.service.accord.serializers.TopologySerializers; +import org.apache.cassandra.utils.concurrent.Future; + +public class FetchTopology +{ + private final long epoch; + + public static final IVersionedSerializer serializer = new IVersionedSerializer<>() + { + @Override + public void serialize(FetchTopology t, DataOutputPlus out, int version) throws IOException + { + out.writeLong(t.epoch); + } + + @Override + public FetchTopology deserialize(DataInputPlus in, int version) throws IOException + { + return new FetchTopology(in.readLong()); + } + + @Override + public long serializedSize(FetchTopology t, int version) + { + return Long.BYTES; + } + }; + + public FetchTopology(long epoch) + { + this.epoch = epoch; + } + + public static class Response + { + private static Response UNKNOWN = new Response(-1, null) { + public String toString() + { + return "UNKNOWN_TOPOLOGY{}"; + } + }; + + // TODO (required): messaging version after version patch + public static final IVersionedSerializer serializer = new IVersionedSerializer<>() + { + @Override + public void serialize(Response t, DataOutputPlus out, int version) throws IOException + { + if (t == UNKNOWN) + { + out.writeLong(-1); + return; + } + out.writeLong(t.epoch); + TopologySerializers.topology.serialize(t.topology, out, version); + } + + @Override + public Response deserialize(DataInputPlus in, int version) throws IOException + { + long epoch = in.readLong(); + if (epoch == -1) + return UNKNOWN; + Topology topology = TopologySerializers.topology.deserialize(in, version); + return new Response(epoch, topology); + } + + @Override + public long serializedSize(Response t, int version) + { + if (t == UNKNOWN) + return Long.BYTES; + + return Long.BYTES + TopologySerializers.topology.serializedSize(t.topology, version); + } + }; + + private final long epoch; + private final Topology topology; + + public Response(long epoch, Topology topology) + { + this.epoch = epoch; + this.topology = topology; + } + } + + public static final IVerbHandler handler = message -> { + long epoch = message.payload.epoch; + Topology topology = AccordService.instance().topology().maybeGlobalForEpoch(epoch); + if (topology == null) + MessagingService.instance().respond(Response.UNKNOWN, message); + else + MessagingService.instance().respond(new Response(epoch, topology), message); + }; + + public static Future fetch(SharedContext context, Collection peers, long epoch) + { + FetchTopology req = new FetchTopology(epoch); + return context.messaging().sendWithRetries(Verb.ACCORD_FETCH_TOPOLOGY_REQ, req, MessagingUtils.tryAliveFirst(SharedContext.Global.instance, peers), + // If the epoch is already discovered, no need to retry + (attempt, from, failure) -> AccordService.instance().currentEpoch() < epoch, + MessageDelivery.RetryErrorMessage.EMPTY) + .map(m -> m.payload.topology); + } +} \ No newline at end of file diff --git a/src/java/org/apache/cassandra/service/accord/JournalKey.java b/src/java/org/apache/cassandra/service/accord/JournalKey.java index ea95629456..68ab0a1c04 100644 --- a/src/java/org/apache/cassandra/service/accord/JournalKey.java +++ b/src/java/org/apache/cassandra/service/accord/JournalKey.java @@ -40,6 +40,7 @@ import org.apache.cassandra.utils.ByteArrayUtil; import static org.apache.cassandra.db.TypeSizes.BYTE_SIZE; import static org.apache.cassandra.db.TypeSizes.INT_SIZE; import static org.apache.cassandra.db.TypeSizes.LONG_SIZE; +import static org.apache.cassandra.service.accord.AccordJournalValueSerializers.*; import static org.apache.cassandra.service.accord.AccordJournalValueSerializers.RangesForEpochSerializer; import static org.apache.cassandra.service.accord.AccordJournalValueSerializers.SafeToReadSerializer; @@ -266,6 +267,7 @@ public final class JournalKey SAFE_TO_READ (3, new SafeToReadSerializer()), BOOTSTRAP_BEGAN_AT (4, new BootstrapBeganAtSerializer()), RANGES_FOR_EPOCH (5, new RangesForEpochSerializer()), + TOPOLOGY_UPDATE (6, new TopologyUpdateSerializer()), ; public final int id; diff --git a/src/java/org/apache/cassandra/service/accord/serializers/TopologySerializers.java b/src/java/org/apache/cassandra/service/accord/serializers/TopologySerializers.java index 4d5fbad524..865fb77b31 100644 --- a/src/java/org/apache/cassandra/service/accord/serializers/TopologySerializers.java +++ b/src/java/org/apache/cassandra/service/accord/serializers/TopologySerializers.java @@ -119,12 +119,23 @@ public class TopologySerializers } }; - public static final IVersionedSerializer shard = new IVersionedSerializer() + public static final IVersionedSerializer shard = new ShardSerializer((IVersionedSerializer) + (IVersionedSerializer) + TokenRange.serializer); + + public static class ShardSerializer implements IVersionedSerializer { + protected IVersionedSerializer range; + + public ShardSerializer(IVersionedSerializer range) + { + this.range = range; + } + @Override public void serialize(Shard shard, DataOutputPlus out, int version) throws IOException { - TokenRange.serializer.serialize((TokenRange) shard.range, out, version); + range.serialize(shard.range, out, version); CollectionSerializers.serializeList(shard.nodes, out, version, nodeId); CollectionSerializers.serializeCollection(shard.fastPathElectorate, out, version, nodeId); CollectionSerializers.serializeCollection(shard.joining, out, version, nodeId); @@ -134,7 +145,7 @@ public class TopologySerializers @Override public Shard deserialize(DataInputPlus in, int version) throws IOException { - Range range = TokenRange.serializer.deserialize(in, version); + Range range = ShardSerializer.this.range.deserialize(in, version); SortedArrayList nodes = CollectionSerializers.deserializeSortedArrayList(in, version, nodeId, Node.Id[]::new); Set fastPathElectorate = CollectionSerializers.deserializeSet(in, version, nodeId); Set joining = CollectionSerializers.deserializeSet(in, version, nodeId); @@ -144,7 +155,7 @@ public class TopologySerializers @Override public long serializedSize(Shard shard, int version) { - long size = TokenRange.serializer.serializedSize((TokenRange) shard.range, version); + long size = range.serializedSize(shard.range, version); size += CollectionSerializers.serializedListSize(shard.nodes, version, nodeId); size += CollectionSerializers.serializedCollectionSize(shard.fastPathElectorate, version, nodeId); size += CollectionSerializers.serializedCollectionSize(shard.joining, version, nodeId); diff --git a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java index 321a69b7d3..d57df4de5f 100644 --- a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java +++ b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java @@ -456,7 +456,7 @@ public class ClusterMetadata public Transformer unregister(NodeId nodeId) { - directory = directory.withoutRackAndDC(nodeId).without(nodeId); + directory = directory.withoutRackAndDC(nodeId).without(epoch, nodeId); if (!tokenMap.tokens(nodeId).isEmpty()) tokenMap = tokenMap.unassignTokens(nodeId); @@ -528,7 +528,7 @@ public class ClusterMetadata Collection transferringTokens = tokenMap.tokens(replaced); tokenMap = tokenMap.unassignTokens(replaced) .assignTokens(replacement, transferringTokens); - directory = directory.without(replaced) + directory = directory.without(epoch, replaced) .withRackAndDC(replacement) .withNodeState(replacement, NodeState.JOINED); diff --git a/src/java/org/apache/cassandra/tcm/Startup.java b/src/java/org/apache/cassandra/tcm/Startup.java index 07d6dd60a3..64ccbf37e8 100644 --- a/src/java/org/apache/cassandra/tcm/Startup.java +++ b/src/java/org/apache/cassandra/tcm/Startup.java @@ -76,6 +76,7 @@ import static org.apache.cassandra.tcm.compatibility.GossipHelper.emptyWithSchem import static org.apache.cassandra.tcm.compatibility.GossipHelper.fromEndpointStates; import static org.apache.cassandra.tcm.membership.NodeState.JOINED; import static org.apache.cassandra.tcm.membership.NodeState.LEFT; +import static org.apache.cassandra.tcm.membership.NodeState.REGISTERED; import static org.apache.cassandra.utils.FBUtilities.getBroadcastAddressAndPort; /** @@ -422,14 +423,9 @@ import static org.apache.cassandra.utils.FBUtilities.getBroadcastAddressAndPort; metadata = ClusterMetadata.current(); NodeState startingstate = metadata.directory.peerState(self); - switch (startingstate) - { - case REGISTERED: - case LEFT: - break; - default: - AccordService.startup(self); - } + if (startingstate != REGISTERED && startingstate != LEFT) + AccordService.startup(self); + switch (startingstate) { case REGISTERED: diff --git a/src/java/org/apache/cassandra/tcm/membership/Directory.java b/src/java/org/apache/cassandra/tcm/membership/Directory.java index e2ba275cda..7dde51fd9f 100644 --- a/src/java/org/apache/cassandra/tcm/membership/Directory.java +++ b/src/java/org/apache/cassandra/tcm/membership/Directory.java @@ -33,6 +33,8 @@ import com.google.common.collect.Sets; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import accord.utils.Invariants; +import accord.utils.btree.BTreeSet; import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; @@ -61,6 +63,7 @@ public class Directory implements MetadataValue private final int nextId; private final Epoch lastModified; private final BTreeBiMap peers; + private final BTreeSet removedNodes; private final BTreeMap locations; public final BTreeMap states; public final BTreeMap versions; @@ -76,6 +79,7 @@ public class Directory implements MetadataValue this(1, Epoch.EMPTY, BTreeBiMap.empty(), + BTreeSet.empty(RemovedNode::compareTo), BTreeMap.empty(), BTreeMap.empty(), BTreeMap.empty(), @@ -88,6 +92,7 @@ public class Directory implements MetadataValue private Directory(int nextId, Epoch lastModified, BTreeBiMap peers, + BTreeSet removedNodes, BTreeMap locations, BTreeMap states, BTreeMap versions, @@ -99,6 +104,7 @@ public class Directory implements MetadataValue this.nextId = nextId; this.lastModified = lastModified; this.peers = peers; + this.removedNodes = removedNodes; this.locations = locations; this.states = states; this.versions = versions; @@ -145,7 +151,7 @@ public class Directory implements MetadataValue @Override public Directory withLastModified(Epoch epoch) { - return new Directory(nextId, epoch, peers, locations, states, versions, hostIds, addresses, endpointsByDC, racksByDC); + return new Directory(nextId, epoch, peers, removedNodes, locations, states, versions, hostIds, addresses, endpointsByDC, racksByDC); } public Directory withNonUpgradedNode(NodeAddresses addresses, @@ -192,6 +198,7 @@ public class Directory implements MetadataValue return new Directory(nextId + 1, lastModified, peers.without(id).with(id, nodeAddresses.broadcastAddress), + removedNodes, locations.withForce(id, location), states.withForce(id, NodeState.REGISTERED), versions.withForce(id, nodeVersion), @@ -203,14 +210,14 @@ public class Directory implements MetadataValue public Directory withNodeState(NodeId id, NodeState state) { - return new Directory(nextId, lastModified, peers, locations, states.withForce(id, state), versions, hostIds, addresses, endpointsByDC, racksByDC); + return new Directory(nextId, lastModified, peers, removedNodes, locations, states.withForce(id, state), versions, hostIds, addresses, endpointsByDC, racksByDC); } public Directory withNodeVersion(NodeId id, NodeVersion version) { if (Objects.equals(versions.get(id), version)) return this; - return new Directory(nextId, lastModified, peers, locations, states, versions.withForce(id, version), hostIds, addresses, endpointsByDC, racksByDC); + return new Directory(nextId, lastModified, peers, removedNodes, locations, states, versions.withForce(id, version), hostIds, addresses, endpointsByDC, racksByDC); } public Directory withNodeAddresses(NodeId id, NodeAddresses nodeAddresses) @@ -232,7 +239,7 @@ public class Directory implements MetadataValue BTreeMap> updatedEndpointsByRack = racksByDC.withForce(location(id).datacenter, rackEP); return new Directory(nextId, lastModified, - peers.withForce(id,nodeAddresses.broadcastAddress), locations, states, versions, hostIds, addresses.withForce(id, nodeAddresses), + peers.withForce(id,nodeAddresses.broadcastAddress), removedNodes, locations, states, versions, hostIds, addresses.withForce(id, nodeAddresses), updatedEndpointsByDC, updatedEndpointsByRack); } @@ -246,7 +253,7 @@ public class Directory implements MetadataValue rackEP = BTreeMultimap.empty(); rackEP = rackEP.with(location.rack, endpoint); - return new Directory(nextId, lastModified, peers, locations, states, versions, hostIds, addresses, + return new Directory(nextId, lastModified, peers, removedNodes, locations, states, versions, hostIds, addresses, endpointsByDC.with(location.datacenter, endpoint), racksByDC.withForce(location.datacenter, rackEP)); } @@ -266,7 +273,7 @@ public class Directory implements MetadataValue newRacksByDC = racksByDC.without(location.datacenter); else newRacksByDC = racksByDC.withForce(location.datacenter, rackEP); - return new Directory(nextId, lastModified, peers, locations, states, versions, hostIds, addresses, + return new Directory(nextId, lastModified, peers, removedNodes, locations, states, versions, hostIds, addresses, endpointsByDC.without(location.datacenter, endpoint), newRacksByDC); } @@ -287,11 +294,27 @@ public class Directory implements MetadataValue if (locations.get(id).equals(location)) return this; - return new Directory(nextId, lastModified, peers, locations.withForce(id, location), states, versions, hostIds, + return new Directory(nextId, lastModified, peers, removedNodes, locations.withForce(id, location), states, versions, hostIds, addresses, endpointsByDC, racksByDC); } - public Directory without(NodeId id) + public Directory removed(Epoch removedIn, NodeId id, InetAddressAndPort addr) + { + Invariants.checkState(!peers.containsKey(id)); + return new Directory(nextId, + lastModified, + peers, + removedNodes.with(new RemovedNode(removedIn, id, addr)), + locations, + states, + versions, + hostIds, + addresses, + endpointsByDC, + racksByDC); + } + + public Directory without(Epoch removedIn, NodeId id) { InetAddressAndPort endpoint = peers.get(id); Location location = locations.get(id); @@ -303,6 +326,7 @@ public class Directory implements MetadataValue return new Directory(nextId, lastModified, peers.without(id), + removedNodes.with(new RemovedNode(removedIn, id, peers.get(id))), locations.without(id), states.without(id), versions.without(id), @@ -319,6 +343,7 @@ public class Directory implements MetadataValue return new Directory(nextId, lastModified, peers.without(id), + removedNodes.with(new RemovedNode(removedIn, id, peers.get(id))), locations.without(id), states.without(id), versions.without(id), @@ -364,6 +389,11 @@ public class Directory implements MetadataValue return peers.keySet(); } + public BTreeSet removedNodes() + { + return removedNodes; + } + public NodeAddresses getNodeAddresses(NodeId id) { return addresses.get(id); @@ -581,6 +611,17 @@ public class Directory implements MetadataValue } } Epoch.serializer.serialize(t.lastModified, out, version); + + if (version.isAtLeast(Version.V7)) + { + out.writeInt(t.removedNodes.size()); + for (RemovedNode removedNode : t.removedNodes) + { + out.writeLong(removedNode.removedIn.getEpoch()); + NodeId.serializer.serialize(removedNode.id, out, version); + InetAddressAndPort.MetadataSerializer.serializer.serialize(removedNode.endpoint, out, version); + } + } } public Directory deserialize(DataInputPlus in, Version version) throws IOException @@ -636,9 +677,23 @@ public class Directory implements MetadataValue else nextId = maxId.id() + 1; } + + if (version.isAtLeast(Version.V7)) + { + int removedNodes = in.readInt(); + for (int i = 0; i < removedNodes; i++) + { + long epoch = in.readLong(); + NodeId nodeId = NodeId.serializer.deserialize(in, version); + InetAddressAndPort addr = InetAddressAndPort.MetadataSerializer.serializer.deserialize(in, version); + newDir.removed(Epoch.create(epoch), nodeId, addr); + } + } + return new Directory(nextId, lastModified, newDir.peers, + newDir.removedNodes, newDir.locations, newDir.states, newDir.versions, @@ -676,6 +731,18 @@ public class Directory implements MetadataValue } } size += Epoch.serializer.serializedSize(t.lastModified, version); + + if (version.isAtLeast(Version.V7)) + { + size += TypeSizes.INT_SIZE; + for (RemovedNode removedNode : t.removedNodes) + { + size += TypeSizes.LONG_SIZE; + size += NodeId.serializer.serializedSize(removedNode.id, version); + size += InetAddressAndPort.MetadataSerializer.serializer.serializedSize(removedNode.endpoint, version); + } + } + return size; } } @@ -797,4 +864,37 @@ public class Directory implements MetadataValue logger.warn("Value for key {} is only present in the right set: {}", k, r.get(k)); } + + + public static class RemovedNode implements Comparable + { + public final Epoch removedIn; + public final NodeId id; + public final InetAddressAndPort endpoint; + + public RemovedNode(Epoch removedIn, NodeId id, InetAddressAndPort endpoint) + { + this.removedIn = removedIn; + this.id = id; + this.endpoint = endpoint; + } + + public boolean equals(Object object) + { + if (this == object) return true; + if (object == null || getClass() != object.getClass()) return false; + RemovedNode that = (RemovedNode) object; + return Objects.equals(removedIn, that.removedIn) && Objects.equals(id, that.id) && Objects.equals(endpoint, that.endpoint); + } + + public int hashCode() + { + return Objects.hash(removedIn, id, endpoint); + } + + public int compareTo(RemovedNode o) + { + return id.compareTo(o.id); + } + } } diff --git a/src/java/org/apache/cassandra/tcm/serialization/Version.java b/src/java/org/apache/cassandra/tcm/serialization/Version.java index f0b56c9917..a10a4ed807 100644 --- a/src/java/org/apache/cassandra/tcm/serialization/Version.java +++ b/src/java/org/apache/cassandra/tcm/serialization/Version.java @@ -56,6 +56,7 @@ public enum Version V6(6), /** * - Accord + * - Track nodes removed */ V7(7), diff --git a/src/java/org/apache/cassandra/utils/Backoff.java b/src/java/org/apache/cassandra/utils/Backoff.java index 79a68fdbc1..cf978a2ee4 100644 --- a/src/java/org/apache/cassandra/utils/Backoff.java +++ b/src/java/org/apache/cassandra/utils/Backoff.java @@ -27,6 +27,7 @@ import org.apache.cassandra.tcm.Retry; public interface Backoff { + // TODO (required): backoff should not have retries boolean mayRetry(int attempt); long computeWaitTime(int attempt); TimeUnit unit(); @@ -85,6 +86,29 @@ public interface Backoff } } + enum NO_OP implements Backoff + { + INSTANCE; + + @Override + public boolean mayRetry(int attempt) + { + return true; + } + + @Override + public long computeWaitTime(int retryCount) + { + return 0; + } + + @Override + public TimeUnit unit() + { + return TimeUnit.NANOSECONDS; + } + } + class ExponentialBackoff implements Backoff { private final int maxAttempts; diff --git a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java index 6f4a3933a9..1279689fd9 100644 --- a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java +++ b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java @@ -653,7 +653,7 @@ public class ClusterUtils } } - public static void waitForCMSToQuiesce(ICluster cluster, int[] cmsNodes) + public static void waitForCMSToQuiesce(ICluster cluster, int... cmsNodes) { // first step; find the largest epoch waitForCMSToQuiesce(cluster, maxEpoch(cluster, cmsNodes)); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java index 57e452009d..29e152a149 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordBootstrapTest.java @@ -183,12 +183,17 @@ public class AccordBootstrapTest extends TestBaseImpl .with(NETWORK, GOSSIP)) .start()) { + cluster.schemaChange("CREATE KEYSPACE ks WITH REPLICATION={'class':'SimpleStrategy', 'replication_factor':2}"); + cluster.schemaChange("CREATE TABLE ks.tbl (k int, c int, v int, primary key(k, c)) WITH transactional_mode='full'"); + long initialMax = maxEpoch(cluster); + for (IInvokableInstance node : cluster) { node.runOnInstance(() -> { Assert.assertEquals(initialMax, ClusterMetadata.current().epoch.getEpoch()); + System.out.println("Awaiting " + initialMax); awaitEpoch(initialMax); AccordConfigurationService configService = service().configurationService(); long minEpoch = configService.minEpoch(); @@ -211,9 +216,6 @@ public class AccordBootstrapTest extends TestBaseImpl node.runOnInstance(StreamListener::register); } - cluster.schemaChange("CREATE KEYSPACE ks WITH REPLICATION={'class':'SimpleStrategy', 'replication_factor':2}"); - cluster.schemaChange("CREATE TABLE ks.tbl (k int, c int, v int, primary key(k, c)) WITH transactional_mode='full'"); - long schemaChangeMax = maxEpoch(cluster); for (IInvokableInstance node : cluster) { @@ -363,6 +365,9 @@ public class AccordBootstrapTest extends TestBaseImpl .with(NETWORK, GOSSIP)) .start()) { + cluster.schemaChange("CREATE KEYSPACE ks WITH REPLICATION={'class':'SimpleStrategy', 'replication_factor':2}"); + cluster.schemaChange("CREATE TABLE ks.tbl (k int, c int, v int, primary key(k, c)) WITH transactional_mode='full'"); + long initialMax = maxEpoch(cluster); long[] tokens = new long[3]; for (int i=0; i<3; i++) @@ -392,9 +397,6 @@ public class AccordBootstrapTest extends TestBaseImpl }); } - cluster.schemaChange("CREATE KEYSPACE ks WITH REPLICATION={'class':'SimpleStrategy', 'replication_factor':2}"); - cluster.schemaChange("CREATE TABLE ks.tbl (k int, c int, v int, primary key(k, c)) WITH transactional_mode='full'"); - long schemaChangeMax = maxEpoch(cluster); for (IInvokableInstance node : cluster) { diff --git a/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java b/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java new file mode 100644 index 0000000000..5f45518263 --- /dev/null +++ b/test/distributed/org/apache/cassandra/fuzz/topology/AccordBootstrapTest.java @@ -0,0 +1,161 @@ +/* + * 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.fuzz.topology; + +import java.util.HashSet; + +import org.junit.Test; + +import org.apache.cassandra.config.CassandraRelevantProperties; +import org.apache.cassandra.distributed.Cluster; +import org.apache.cassandra.distributed.Constants; +import org.apache.cassandra.distributed.api.ConsistencyLevel; +import org.apache.cassandra.distributed.api.Feature; +import org.apache.cassandra.distributed.api.IInstanceConfig; +import org.apache.cassandra.distributed.api.IInvokableInstance; +import org.apache.cassandra.distributed.api.TokenSupplier; +import org.apache.cassandra.distributed.shared.ClusterUtils; +import org.apache.cassandra.distributed.shared.NetworkTopology; +import org.apache.cassandra.distributed.test.log.FuzzTestBase; +import org.apache.cassandra.harry.SchemaSpec; +import org.apache.cassandra.harry.dsl.HistoryBuilder; +import org.apache.cassandra.harry.dsl.HistoryBuilderHelper; +import org.apache.cassandra.harry.dsl.ReplayingHistoryBuilder; +import org.apache.cassandra.harry.execution.InJvmDTestVisitExecutor; +import org.apache.cassandra.harry.execution.QueryBuildingVisitExecutor; +import org.apache.cassandra.harry.gen.Generator; +import org.apache.cassandra.harry.gen.Generators; +import org.apache.cassandra.harry.gen.Generators.TrackingGenerator; +import org.apache.cassandra.harry.gen.SchemaGenerators; +import org.apache.cassandra.service.consensus.TransactionalMode; + +import static org.apache.cassandra.distributed.shared.ClusterUtils.unpauseCommits; +import static org.apache.cassandra.distributed.shared.ClusterUtils.waitForCMSToQuiesce; +import static org.apache.cassandra.harry.checker.TestHelper.withRandom; + +public class AccordBootstrapTest extends FuzzTestBase +{ + private static final int WRITES = 10; + private static final int POPULATION = 1000; + @Test + public void bootstrapFuzzTest() throws Throwable + { + CassandraRelevantProperties.SYSTEM_TRACES_DEFAULT_RF.setInt(3); + IInvokableInstance forShutdown = null; + try (Cluster cluster = builder().withNodes(3) + .withTokenSupplier(TokenSupplier.evenlyDistributedTokens(100)) + .withNodeIdTopology(NetworkTopology.singleDcNetworkTopology(100, "dc0", "rack0")) + .withConfig((config) -> config.with(Feature.NETWORK, Feature.GOSSIP) + .set("write_request_timeout", "2s") + .set("request_timeout", "5s") + .set("concurrent_accord_operations", 2) + .set("progress_barrier_min_consistency_level", "QUORUM") + .set("progress_barrier_default_consistency_level", "QUORUM") + .set("metadata_snapshot_frequency", 5)) + .start()) + { + IInvokableInstance cmsInstance = cluster.get(1); + forShutdown = cmsInstance; + waitForCMSToQuiesce(cluster, cmsInstance); + + HashSet downInstances = new HashSet<>(); + withRandom(rng -> { + Generator schemaGen = SchemaGenerators.trivialSchema(KEYSPACE, "bootstrap_fuzz", POPULATION, + SchemaSpec.optionsBuilder() + .addWriteTimestamps(false) + .withTransactionalMode(TransactionalMode.full) + ); + + SchemaSpec schema = schemaGen.generate(rng); + TrackingGenerator pkGen = Generators.tracking(Generators.int32(0, Math.min(schema.valueGenerators.pkPopulation(), POPULATION))); + Generator ckGen = Generators.int32(0, Math.min(schema.valueGenerators.ckPopulation(), POPULATION)); + HistoryBuilder history = new ReplayingHistoryBuilder(schema.valueGenerators, + hb -> InJvmDTestVisitExecutor.builder() + .consistencyLevel(ConsistencyLevel.QUORUM) + .wrapQueries(QueryBuildingVisitExecutor.WrapQueries.TRANSACTION) + .pageSizeSelector(p -> InJvmDTestVisitExecutor.PageSizeSelector.NO_PAGING) + .nodeSelector(lts -> { + while (true) + { + int pick = rng.nextInt(1, cluster.size() + 1); + if (!downInstances.contains(pick)) + return pick; + + } + }) + .build(schema, hb, cluster)); + + Runnable writeAndValidate = () -> { + for (int i = 0; i < WRITES; i++) + HistoryBuilderHelper.insertRandomData(schema, pkGen, ckGen, rng, history); + + for (int pk : pkGen.generated()) + history.selectPartition(pk); + }; + + history.customThrowing(() -> { + cluster.schemaChange(String.format("CREATE KEYSPACE %s WITH replication = {'class': 'SimpleStrategy', 'replication_factor' : 3};", KEYSPACE)); + cluster.schemaChange(schema.compile()); + waitForCMSToQuiesce(cluster, cmsInstance); + }, "Setup"); + Thread.sleep(1000); + writeAndValidate.run(); + + history.customThrowing(() -> { + IInstanceConfig config = cluster.newInstanceConfig() + .set("auto_bootstrap", true) + .set(Constants.KEY_DTEST_FULL_STARTUP, true); + cluster.bootstrap(config).startup(); + waitForCMSToQuiesce(cluster, cmsInstance); + }, "Start boostrap"); + + writeAndValidate.run(); + + history.customThrowing(() -> { + downInstances.add(2); + ClusterUtils.stopUnchecked(cluster.get(2)); + cluster.get(1).logs().watchFor("/127.0.0.2:.* is now DOWN"); + }, "Shut down node 2"); + + history.customThrowing(() -> { + IInstanceConfig config = cluster.newInstanceConfig() + .set("auto_bootstrap", true) + .set(Constants.KEY_DTEST_FULL_STARTUP, true); + cluster.bootstrap(config).startup(); + waitForCMSToQuiesce(cluster, cmsInstance); + }, "Bootstrap one more"); + + writeAndValidate.run(); + + history.customThrowing(() -> { + cluster.get(2).startup(); + cluster.get(1).logs().watchFor("/127.0.0.2:.* is now UP"); + downInstances.remove(2); + }, "Start up node 2"); + + }); + } + catch (Throwable t) + { + if (forShutdown != null) + unpauseCommits(forShutdown); + throw t; + } + } +} diff --git a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java index 9e41e0add9..caec9ec444 100644 --- a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java +++ b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java @@ -50,6 +50,7 @@ import org.apache.cassandra.service.accord.serializers.CommandSerializers; import org.apache.cassandra.service.accord.serializers.DepsSerializers; import org.apache.cassandra.service.accord.serializers.KeySerializers; import org.apache.cassandra.service.accord.serializers.ResultSerializers; +import org.apache.cassandra.service.accord.serializers.TopologySerializers; import org.apache.cassandra.tools.FieldUtil; import static accord.impl.PrefixedIntHashKey.ranges; @@ -76,14 +77,21 @@ public class AccordJournalBurnTest extends BurnTestBase BurnTestKeySerializers.update, BurnTestKeySerializers.write), CommandSerializers.class); + FieldUtil.transferFields(new DepsSerializers.Impl(BurnTestKeySerializers.range), DepsSerializers.class); + FieldUtil.setInstanceUnsafe(ResultSerializers.class, BurnTestKeySerializers.result, "result"); + + FieldUtil.setInstanceUnsafe(TopologySerializers.class, + new TopologySerializers.ShardSerializer(BurnTestKeySerializers.range), + "shard"); } - private AtomicInteger counter = new AtomicInteger(); + private static final AtomicInteger counter = new AtomicInteger(); + @Before public void beforeTest() throws Throwable { diff --git a/test/harry/main/org/apache/cassandra/harry/execution/CQLVisitExecutor.java b/test/harry/main/org/apache/cassandra/harry/execution/CQLVisitExecutor.java index 1922f0011a..92a9627683 100644 --- a/test/harry/main/org/apache/cassandra/harry/execution/CQLVisitExecutor.java +++ b/test/harry/main/org/apache/cassandra/harry/execution/CQLVisitExecutor.java @@ -90,9 +90,12 @@ public abstract class CQLVisitExecutor public static void replayAfterFailure(Visit visit, CQLVisitExecutor executor, Model.Replay replay) { QueryBuildingVisitExecutor queryBuilder = executor.queryBuilder; - logger.error("Caught an exception at {} while replaying {}\ncluster.schemaChange(\"{}\");\nOperations _for this partition_ up to this visit:", - visit, queryBuilder.compile(visit), - queryBuilder.schema.compile()); + if (!visit.hasCustom) + { + logger.error("Caught an exception at {} while replaying {}\ncluster.schemaChange(\"{}\");\nOperations _for this partition_ up to this visit:", + visit, queryBuilder.compile(visit), + queryBuilder.schema.compile()); + } // Configurable yet hardcoded for a person who is trying to generate repro ResultDumpMode mode = ResultDumpMode.PARTITION; diff --git a/test/harry/main/org/apache/cassandra/harry/execution/InJvmDTestVisitExecutor.java b/test/harry/main/org/apache/cassandra/harry/execution/InJvmDTestVisitExecutor.java index dca2b7001f..4c4d420d8a 100644 --- a/test/harry/main/org/apache/cassandra/harry/execution/InJvmDTestVisitExecutor.java +++ b/test/harry/main/org/apache/cassandra/harry/execution/InJvmDTestVisitExecutor.java @@ -292,9 +292,10 @@ public class InJvmDTestVisitExecutor extends CQLVisitExecutor public interface RetryPolicy { RetryPolicy RETRY_ON_TIMEOUT = (t) -> { - return t.getMessage().contains("timed out") || - AssertionUtils.isInstanceof(RequestTimeoutException.class) - .matches(Throwables.getRootCause(t)); + return t.getMessage() != null && + (t.getMessage().contains("timed out") || + AssertionUtils.isInstanceof(RequestTimeoutException.class) + .matches(Throwables.getRootCause(t))); }; RetryPolicy NO_RETRY = (t) -> false; boolean retry(Throwable t); diff --git a/test/harry/main/org/apache/cassandra/harry/gen/SchemaGenerators.java b/test/harry/main/org/apache/cassandra/harry/gen/SchemaGenerators.java index c0e0c6e569..a3433c5b56 100644 --- a/test/harry/main/org/apache/cassandra/harry/gen/SchemaGenerators.java +++ b/test/harry/main/org/apache/cassandra/harry/gen/SchemaGenerators.java @@ -97,6 +97,11 @@ public class SchemaGenerators } public static Generator trivialSchema(String ks, String table, int population) + { + return trivialSchema(ks, table, population, SchemaSpec.optionsBuilder().build()); + } + + public static Generator trivialSchema(String ks, String table, int population, SchemaSpec.Options options) { return (rng) -> { return new SchemaSpec(rng.next(), @@ -105,7 +110,8 @@ public class SchemaGenerators Arrays.asList(ColumnSpec.pk("pk1", ColumnSpec.int64Type, Generators.int64())), Arrays.asList(ColumnSpec.ck("ck1", ColumnSpec.int64Type, Generators.int64(), false)), Arrays.asList(ColumnSpec.regularColumn("v1", ColumnSpec.int64Type)), - List.of(ColumnSpec.staticColumn("s1", ColumnSpec.int64Type))); + Arrays.asList(ColumnSpec.staticColumn("s1", ColumnSpec.int64Type)), + options); }; } } diff --git a/test/harry/main/org/apache/cassandra/harry/op/Visit.java b/test/harry/main/org/apache/cassandra/harry/op/Visit.java index 6c8c31390c..7d77b72871 100644 --- a/test/harry/main/org/apache/cassandra/harry/op/Visit.java +++ b/test/harry/main/org/apache/cassandra/harry/op/Visit.java @@ -32,6 +32,8 @@ public class Visit public final Set visitedPartitions; public final boolean selectOnly; + public final boolean hasCustom; + public Visit(long lts, Operation[] operations) { Assert.assertTrue(operations.length > 0); @@ -39,8 +41,11 @@ public class Visit this.operations = operations; this.visitedPartitions = new HashSet<>(); boolean selectOnly = true; + boolean hasCustom = false; for (Operation operation : operations) { + if (operation.kind() == Operations.Kind.CUSTOM) + hasCustom = true; if (selectOnly && !(operation instanceof Operations.SelectStatement)) selectOnly = false; @@ -49,6 +54,7 @@ public class Visit } this.selectOnly = selectOnly; + this.hasCustom = hasCustom; } public String toString() diff --git a/test/simulator/main/org/apache/cassandra/simulator/paxos/PaxosSimulation.java b/test/simulator/main/org/apache/cassandra/simulator/paxos/PaxosSimulation.java index 6c726a1504..74ffd4db56 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/paxos/PaxosSimulation.java +++ b/test/simulator/main/org/apache/cassandra/simulator/paxos/PaxosSimulation.java @@ -36,6 +36,7 @@ import com.google.common.base.Throwables; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import accord.coordinate.Invalidated; import org.apache.cassandra.concurrent.ExecutorFactory; import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.distributed.Cluster; @@ -77,7 +78,7 @@ public abstract class PaxosSimulation implements Simulation, ClusterActionListen protected Class[] expectedExceptions() { - return (Class[]) new Class[] { RequestExecutionException.class }; + return (Class[]) new Class[] { RequestExecutionException.class, Invalidated.class }; } abstract class Operation extends SimulatedActionCallable implements BiConsumer diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java new file mode 100644 index 0000000000..75b19b5b2d --- /dev/null +++ b/test/simulator/test/org/apache/cassandra/simulator/test/AccordHarrySimulationTest.java @@ -0,0 +1,71 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.simulator.test; + +import java.util.Arrays; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; + +import org.apache.cassandra.harry.SchemaSpec; +import org.apache.cassandra.harry.gen.Generator; +import org.apache.cassandra.harry.gen.SchemaGenerators; +import org.apache.cassandra.net.Verb; +import org.apache.cassandra.service.consensus.TransactionalMode; +import org.apache.cassandra.simulator.FixedLossNetworkScheduler; +import org.apache.cassandra.simulator.FutureActionScheduler; +import org.apache.cassandra.simulator.RandomSource; +import org.apache.cassandra.simulator.systems.SimulatedTime; +import org.apache.cassandra.simulator.utils.KindOfSequence; + +public class AccordHarrySimulationTest extends HarrySimulatorTest +{ + + @Override + public Map networkSchedulers(int nodes, SimulatedTime time, RandomSource random) + { + + Set extremelyLossy = new HashSet<>(Arrays.asList(Verb.ACCORD_SIMPLE_RSP, Verb.ACCORD_PRE_ACCEPT_RSP, Verb.ACCORD_PRE_ACCEPT_REQ, + Verb.ACCORD_ACCEPT_RSP, Verb.ACCORD_ACCEPT_REQ, Verb.ACCORD_ACCEPT_INVALIDATE_REQ, + Verb.ACCORD_READ_RSP, Verb.ACCORD_READ_REQ, Verb.ACCORD_COMMIT_REQ, + Verb.ACCORD_COMMIT_INVALIDATE_REQ, Verb.ACCORD_APPLY_RSP, Verb.ACCORD_APPLY_REQ, + Verb.ACCORD_BEGIN_RECOVER_RSP, Verb.ACCORD_BEGIN_RECOVER_REQ, Verb.ACCORD_BEGIN_INVALIDATE_RSP)); + + Set somewhatLossy = new HashSet<>(Arrays.asList(Verb.ACCORD_SYNC_NOTIFY_RSP, Verb.ACCORD_SYNC_NOTIFY_REQ, Verb.ACCORD_APPLY_AND_WAIT_REQ, + Verb.ACCORD_FETCH_MIN_EPOCH_RSP, Verb.ACCORD_FETCH_MIN_EPOCH_REQ, Verb.ACCORD_FETCH_TOPOLOGY_RSP, + Verb.ACCORD_FETCH_TOPOLOGY_REQ)); + + Map schedulers = new HashMap<>(); + for (Verb verb : Verb.values()) + { + if (extremelyLossy.contains(verb)) + schedulers.put(verb, new FixedLossNetworkScheduler(nodes, random, time, KindOfSequence.UNIFORM, .15f, .20f)); + else if (somewhatLossy.contains(verb)) + schedulers.put(verb, new FixedLossNetworkScheduler(nodes, random, time, KindOfSequence.UNIFORM, .1f, .15f)); + } + return schedulers; + } + + public Generator schemaSpecGen(String keyspace, String prefix) + { + return SchemaGenerators.schemaSpecGen(keyspace, prefix, 1000, SchemaSpec.optionsBuilder().withTransactionalMode(TransactionalMode.full)); + } + +} diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java index b870c214bb..61fb38fc7e 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/HarrySimulatorTest.java @@ -56,19 +56,19 @@ import org.apache.cassandra.distributed.api.SimpleQueryResult; import org.apache.cassandra.distributed.impl.Query; import org.apache.cassandra.distributed.shared.WithProperties; import org.apache.cassandra.harry.SchemaSpec; -import org.apache.cassandra.harry.op.Visit; -import org.apache.cassandra.harry.op.Operations; -import org.apache.cassandra.harry.gen.OperationsGenerators; import org.apache.cassandra.harry.execution.CompiledStatement; import org.apache.cassandra.harry.execution.DataTracker; import org.apache.cassandra.harry.execution.QueryBuildingVisitExecutor; import org.apache.cassandra.harry.gen.EntropySource; import org.apache.cassandra.harry.gen.Generator; +import org.apache.cassandra.harry.gen.OperationsGenerators; import org.apache.cassandra.harry.gen.SchemaGenerators; import org.apache.cassandra.harry.gen.rng.JdkRandomEntropySource; import org.apache.cassandra.harry.model.Model; import org.apache.cassandra.harry.model.QuiescentChecker; import org.apache.cassandra.harry.model.TokenPlacementModel; +import org.apache.cassandra.harry.op.Operations; +import org.apache.cassandra.harry.op.Visit; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.Verb; import org.apache.cassandra.schema.ReplicationParams; @@ -122,9 +122,9 @@ import static org.apache.cassandra.simulator.cluster.ClusterActions.Options.noAc * * And then run your test using the following settings (omit add-* if you are running on jdk8): * - -Dstorage-config=/Users/ifesdjeen/p/java/cassandra-accord/test/conf + -Dstorage-config=$MODULE_DIR$/test/conf -Djava.awt.headless=true - -javaagent:/Users/ifesdjeen/p/java/cassandra-accord/lib/jamm-0.4.0.jar + -javaagent:$MODULE_DIR$/lib/jamm-0.4.0.jar -ea -Dcassandra.debugrefcount=true -Xss384k @@ -146,15 +146,15 @@ import static org.apache.cassandra.simulator.cluster.ClusterActions.Options.noAc -Dcassandra.test.messagingService.nonGracefulShutdown=true -Dcassandra.use_nix_recursive_delete=true -Dcie-cassandra.disable_schema_drop_log=true - -Dlogback.configurationFile=file:///Users/ifesdjeen/p/java/cassandra-accord/test/conf/logback-simulator.xml + -Dlogback.configurationFile=file://$MODULE_DIR$/test/conf/logback-simulator.xml -Dcassandra.ring_delay_ms=10000 -Dcassandra.tolerate_sstable_size=true -Dcassandra.skip_sync=true -Dcassandra.debugrefcount=false -Dcassandra.test.simulator.determinismcheck=strict -Dcassandra.test.simulator.print_asm=none - -javaagent:/Users/ifesdjeen/p/java/cassandra-accord/build/test/lib/jars/simulator-asm.jar - -Xbootclasspath/a:/Users/ifesdjeen/p/java/cassandra-accord/build/test/lib/jars/simulator-bootstrap.jar + -javaagent:$MODULE_DIR$/build/test/lib/jars/simulator-asm.jar + -Xbootclasspath/a:$MODULE_DIR$/build/test/lib/jars/simulator-bootstrap.jar -XX:ActiveProcessorCount=4 -XX:-TieredCompilation -XX:-BackgroundCompilation @@ -202,7 +202,7 @@ public class HarrySimulatorTest HarrySimulatorTest test = SingleCommand.singleCommand(HarrySimulatorTest.class).parse(args); if (test.helpOption.showHelpIfRequested()) return; - test.harryTest(); + test.testInternal(); System.exit(1); } @@ -211,11 +211,10 @@ public class HarrySimulatorTest { // To rerun a failing test for a given seed, uncomment the below and set the seed // this.seed = ""; - this.seed = "0xdd3bb3793a6b925a"; - harryTest(); + testInternal(); } - private void harryTest() throws Exception + protected void testInternal() throws Exception { int bootstrapNode1 = 4; int bootstrapNode2 = 8; @@ -514,7 +513,7 @@ public class HarrySimulatorTest } } - static class HarrySimulationBuilder extends ClusterSimulation.Builder + class HarrySimulationBuilder extends ClusterSimulation.Builder { protected final Consumer configUpdater; @@ -526,7 +525,7 @@ public class HarrySimulatorTest @Override public Map perVerbFutureActionSchedulers(int nodeCount, SimulatedTime time, RandomSource random) { - return HarrySimulatorTest.networkSchedulers(nodeCount, time, random); + return networkSchedulers(nodeCount, time, random); } @Override @@ -600,7 +599,7 @@ public class HarrySimulatorTest /** * Custom network scheduler for testing TCM. */ - public static Map networkSchedulers(int nodes, SimulatedTime time, RandomSource random) + public Map networkSchedulers(int nodes, SimulatedTime time, RandomSource random) { Set extremelyLossy = new HashSet<>(Arrays.asList(Verb.TCM_ABORT_MIG, Verb.TCM_REPLICATION, Verb.TCM_COMMIT_REQ, Verb.TCM_NOTIFY_REQ, @@ -633,7 +632,7 @@ public class HarrySimulatorTest return schedulers; } - public Action reconfigureCMS(SimulatedSystems simulated, Cluster cluster, int rf, boolean inEachDc) + public static Action reconfigureCMS(SimulatedSystems simulated, Cluster cluster, int rf, boolean inEachDc) { return new SimulatedActionTask("", Action.Modifiers.RELIABLE_NO_TIMEOUTS, Action.Modifiers.RELIABLE_NO_TIMEOUTS, null, simulated, new InterceptedExecution.InterceptedRunnableExecution((InterceptingExecutor) cluster.get(1).executor(), @@ -690,7 +689,7 @@ public class HarrySimulatorTest cluster.get(node).transfer(runnable))); } - public Action decommission(SimulatedSystems simulated, Cluster cluster, int node) + public static Action decommission(SimulatedSystems simulated, Cluster cluster, int node) { IIsolatedExecutor.SerializableRunnable runnable = () -> { try @@ -764,12 +763,14 @@ public class HarrySimulatorTest CompiledStatement compiledStatement = simulation.queryBuilder.compile(visit); DataTracker tracker = simulation.tracker; + int[] joined = simulation.nodeState.joined(); + int coordinator = joined[simulation.rng.nextInt(joined.length)]; RetryingQuery query = new RetryingQuery(compiledStatement.cql(), cl, compiledStatement.bindings()); Action wrapper = new SimulatedActionCallable<>("Query", Action.Modifiers.RELIABLE_NO_TIMEOUTS, Action.Modifiers.RELIABLE_NO_TIMEOUTS, simulation.simulated, - simulation.cluster.get((int) ((lts % simulation.cluster.size()) + 1)), + simulation.cluster.get(coordinator), query) { @Override @@ -780,7 +781,6 @@ public class HarrySimulatorTest public void run() { tracker.begin(visit); - System.out.println("Started visit = " + visit); // we'll be invoked on the node's executor, but we need to ensure the task is loaded on its classloader try { @@ -804,10 +804,7 @@ public class HarrySimulatorTest if (failure != null) simulated.failures.accept(failure); else - { - System.out.println("Finished visit = " + visit); tracker.end(visit); - } } }; @@ -1038,7 +1035,7 @@ public class HarrySimulatorTest return arr; } - public static Generator schemaSpecGen(String keyspace, String prefix) + public Generator schemaSpecGen(String keyspace, String prefix) { return SchemaGenerators.schemaSpecGen(keyspace, prefix, 1000); } diff --git a/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java b/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java index 9d2fa5de21..cd90f9b977 100644 --- a/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java +++ b/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java @@ -24,6 +24,8 @@ import java.util.List; import java.util.Map; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; import com.google.common.collect.Iterators; @@ -47,14 +49,13 @@ import org.apache.cassandra.utils.Backoff; import org.mockito.Mockito; import static accord.utils.Property.qt; +import static org.apache.cassandra.net.MessageDelivery.RetryErrorMessage; +import static org.apache.cassandra.net.MessageDelivery.RetryPredicate; import static org.assertj.core.api.Assertions.assertThat; public class MessageDeliveryTest { private static final InetAddressAndPort ID1 = InetAddressAndPort.getByNameUnchecked("127.0.0.1"); - private static final MessageDelivery.RetryErrorMessage RETRY_ERROR_MESSAGE = (i1, i2, i3, i4) -> null; - private static final MessageDelivery.RetryPredicate ALWAYS_RETRY = (i1, i2, i3) -> true; - private static final MessageDelivery.RetryPredicate ALWAYS_REJECT = (i1, i2, i3) -> false; static { @@ -79,8 +80,8 @@ public class MessageDeliveryTest scheduler::schedule, Verb.ECHO_REQ, NoPayload.noPayload, Iterators.cycle(ID1), - ALWAYS_RETRY, - RETRY_ERROR_MESSAGE); + RetryPredicate.ALWAYS_RETRY, + RetryErrorMessage.EMPTY); assertThat(result).isNotDone(); factory.processAll(); assertThat(result).isDone(); @@ -104,8 +105,8 @@ public class MessageDeliveryTest scheduler::schedule, Verb.ECHO_REQ, NoPayload.noPayload, Iterators.cycle(ID1), - ALWAYS_RETRY, - RETRY_ERROR_MESSAGE); + RetryPredicate.ALWAYS_RETRY, + RetryErrorMessage.EMPTY); assertThat(result).isNotDone(); factory.processAll(); assertThat(result).isDone(); @@ -135,8 +136,8 @@ public class MessageDeliveryTest scheduler::schedule, Verb.ECHO_REQ, NoPayload.noPayload, Iterators.cycle(ID1), - ALWAYS_RETRY, - RETRY_ERROR_MESSAGE); + RetryPredicate.ALWAYS_RETRY, + RetryErrorMessage.EMPTY); assertThat(result).isNotDone(); factory.processAll(); assertThat(result).isDone(); @@ -163,8 +164,8 @@ public class MessageDeliveryTest scheduler::schedule, Verb.ECHO_REQ, NoPayload.noPayload, Iterators.cycle(ID1), - ALWAYS_REJECT, - RETRY_ERROR_MESSAGE); + RetryPredicate.NEVER_RETRY, + RetryErrorMessage.EMPTY); assertThat(result).isNotDone(); factory.processAll(); assertThat(result).isDone(); @@ -177,6 +178,26 @@ public class MessageDeliveryTest }); } + public static FailedResponseException getFailedResponseException(Future> result) throws InterruptedException + { + FailedResponseException ex; + try + { + result.get(1, TimeUnit.MINUTES); + Assert.fail("Should have failed"); + throw new AssertionError("Not Reachable"); + } + catch (ExecutionException e) + { + ex = (FailedResponseException) e.getCause(); + } + catch (TimeoutException e) + { + throw new RuntimeException(e); + } + return ex; + } + private static MessageDelivery simulatedMessages(RandomSource rs, ScheduledExecutorPlus scheduler, List failures, SimulatedMessageDelivery.ActionSupplier actionSupplier) { Map receivers = new HashMap<>(); @@ -191,22 +212,6 @@ public class MessageDeliveryTest return messaging; } - private static FailedResponseException getFailedResponseException(Future> result) throws InterruptedException - { - FailedResponseException ex; - try - { - result.get(); - Assert.fail("Should have failed"); - throw new AssertionError("Not Reachable"); - } - catch (ExecutionException e) - { - ex = (FailedResponseException) e.getCause(); - } - return ex; - } - private static MaxRetriesException getMaxRetriesException(Future> result) throws InterruptedException { MaxRetriesException ex; diff --git a/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java b/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java index 7210d4a097..04a3289985 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordConfigurationServiceTest.java @@ -21,19 +21,49 @@ package org.apache.cassandra.service.accord; import java.net.UnknownHostException; import java.util.ArrayList; import java.util.Collections; +import java.util.Iterator; import java.util.List; import java.util.Optional; import java.util.UUID; import com.google.common.collect.Sets; +import org.junit.Assert; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Test; + +import accord.api.Journal; +import accord.impl.AbstractConfigurationServiceTest; +import accord.impl.TestAgent; +import accord.impl.basic.InMemoryJournal; +import accord.local.Node.Id; +import accord.topology.Topology; +import accord.utils.SortedArrays.SortedArrayList; +import accord.utils.async.AsyncResult; +import org.agrona.collections.Int2ObjectHashMap; +import org.apache.cassandra.SchemaLoader; +import org.apache.cassandra.ServerTestUtils; +import org.apache.cassandra.concurrent.ScheduledExecutors; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.marshal.Int32Type; +import org.apache.cassandra.dht.Murmur3Partitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.Replica; +import org.apache.cassandra.net.ConnectionType; +import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.MessageDelivery; +import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.schema.DistributedSchema; import org.apache.cassandra.schema.KeyspaceMetadata; +import org.apache.cassandra.schema.KeyspaceParams; +import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.schema.Tables; +import org.apache.cassandra.service.accord.AccordKeyspace.EpochDiskState; +import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ValidatingClusterMetadataService; import org.apache.cassandra.tcm.membership.Location; import org.apache.cassandra.tcm.membership.NodeAddresses; @@ -41,31 +71,6 @@ import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.membership.NodeVersion; import org.apache.cassandra.tcm.ownership.DataPlacement; import org.apache.cassandra.tcm.serialization.Version; -import org.junit.Assert; -import org.junit.Before; -import org.junit.BeforeClass; -import org.junit.Test; - -import accord.api.ConfigurationService.EpochReady; -import accord.impl.AbstractConfigurationServiceTest; -import accord.local.Node.Id; -import accord.topology.Topology; -import accord.utils.SortedArrays.SortedArrayList; -import org.apache.cassandra.SchemaLoader; -import org.apache.cassandra.ServerTestUtils; -import org.apache.cassandra.concurrent.ScheduledExecutors; -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.Keyspace; -import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.net.ConnectionType; -import org.apache.cassandra.net.Message; -import org.apache.cassandra.net.MessageDelivery; -import org.apache.cassandra.net.RequestCallback; -import org.apache.cassandra.schema.KeyspaceParams; -import org.apache.cassandra.schema.TableId; -import org.apache.cassandra.service.accord.AccordKeyspace.EpochDiskState; -import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.utils.MockFailureDetector; import org.apache.cassandra.utils.concurrent.Future; @@ -175,7 +180,8 @@ public class AccordConfigurationServiceTest { ValidatingClusterMetadataService cms = ValidatingClusterMetadataService.createAndRegister(Version.MIN_ACCORD_VERSION); - AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks); + + AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks, null); Assert.assertEquals(null, AccordKeyspace.loadEpochDiskState()); service.start(); Assert.assertEquals(null, AccordKeyspace.loadEpochDiskState()); @@ -183,16 +189,14 @@ public class AccordConfigurationServiceTest Topology topology1 = createTopology(cms); service.reportTopology(topology1); - loadEpoch(1, cms.metadata(), (epoch, cm, topology, syncStatus, pendingSync, remoteSync, closed, redundant) -> { - Assert.assertEquals(topology1, topology); + loadEpoch(1, (epoch, syncStatus, pendingSync, remoteSync, closed, redundant) -> { Assert.assertTrue(remoteSync.isEmpty()); }); Assert.assertEquals(EpochDiskState.create(1), service.diskState()); service.receiveRemoteSyncComplete(ID1, 1); service.receiveRemoteSyncComplete(ID2, 1); - loadEpoch(1, cms.metadata(), (epoch, cm, topology, syncStatus, pendingSync, remoteSync, closed, redundant) -> { - Assert.assertEquals(topology1, topology); + loadEpoch(1, (epoch, syncStatus, pendingSync, remoteSync, closed, redundant) -> { Assert.assertEquals(Sets.newHashSet(ID1, ID2), remoteSync); }); } @@ -202,30 +206,42 @@ public class AccordConfigurationServiceTest { ValidatingClusterMetadataService cms = ValidatingClusterMetadataService.createAndRegister(Version.MIN_ACCORD_VERSION); - AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks); + InMemoryJournal journal = new InMemoryJournal(ID1, new TestAgent()); + AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks, journal); + TestListener listener = new TestListener(service, true) { + @Override + public AsyncResult onTopologyUpdate(Topology topology, boolean isLoad, boolean startSync) + { + // Fake journal save + journal.saveTopology(new Journal.TopologyUpdate(new Int2ObjectHashMap<>(), topology, topology), () -> {}); + return super.onTopologyUpdate(topology, isLoad, startSync); + } + }; + service.registerListener(listener); service.start(); Topology topology1 = createTopology(cms); service.updateMapping(mappingForEpoch(cms.metadata().epoch.getEpoch() + 1)); service.reportTopology(topology1); - service.acknowledgeEpoch(EpochReady.done(1), true); service.receiveRemoteSyncComplete(ID1, 1); service.receiveRemoteSyncComplete(ID2, 1); service.receiveRemoteSyncComplete(ID3, 1); Topology topology2 = createTopology(cms); service.reportTopology(topology2); - service.acknowledgeEpoch(EpochReady.done(2), true); service.receiveRemoteSyncComplete(ID1, 2); Topology topology3 = createTopology(cms); service.reportTopology(topology3); - service.acknowledgeEpoch(EpochReady.done(3), true); - AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks); + AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks, journal); loaded.updateMapping(mappingForEpoch(cms.metadata().epoch.getEpoch() + 1)); - AbstractConfigurationServiceTest.TestListener listener = new AbstractConfigurationServiceTest.TestListener(loaded, true); + listener = new AbstractConfigurationServiceTest.TestListener(loaded, true); loaded.registerListener(listener); + Iterator iter = journal.replayTopologies(); + // Simulate journal replay + while (iter.hasNext()) + loaded.reportTopology(iter.next().global); loaded.start(); listener.assertNoTruncates(); @@ -242,9 +258,17 @@ public class AccordConfigurationServiceTest public void truncateTest() { ValidatingClusterMetadataService cms = ValidatingClusterMetadataService.createAndRegister(Version.MIN_ACCORD_VERSION); - - AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks); - TestListener serviceListener = new TestListener(service, true); + InMemoryJournal journal = new InMemoryJournal(ID1, new TestAgent()); + AccordConfigurationService service = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks, journal); + TestListener serviceListener = new TestListener(service, true) { + @Override + public AsyncResult onTopologyUpdate(Topology topology, boolean isLoad, boolean startSync) + { + // Fake journal save + journal.saveTopology(new Journal.TopologyUpdate(new Int2ObjectHashMap<>(), topology, topology), () -> {}); + return super.onTopologyUpdate(topology, isLoad, startSync); + } + }; service.registerListener(serviceListener); service.start(); @@ -258,13 +282,18 @@ public class AccordConfigurationServiceTest Topology topology3 = createTopology(cms); service.reportTopology(topology3); service.truncateTopologiesUntil(3); + journal.truncateTopologiesForTesting(3); Assert.assertEquals(EpochDiskState.create(3), service.diskState()); serviceListener.assertTruncates(3L); - AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks); + AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector(), AccordConfigurationService.SystemTableDiskStateManager.instance, ScheduledExecutors.scheduledTasks, journal); loaded.updateMapping(mappingForEpoch(cms.metadata().epoch.getEpoch() + 1)); TestListener loadListener = new TestListener(loaded, true); loaded.registerListener(loadListener); + Iterator iter = journal.replayTopologies(); + // Simulate journal replay + while (iter.hasNext()) + loaded.reportTopology(iter.next().global); loaded.start(); loadListener.assertTopologiesFor(3L); } diff --git a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java index 2a10715c8b..3f5b59f0dc 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordTestUtils.java @@ -43,7 +43,6 @@ import accord.api.RoutingKey; import accord.api.Timeouts; import accord.impl.DefaultLocalListeners; import accord.impl.DefaultLocalListeners.NotifySink.NoOpNotifySink; -import accord.impl.InMemoryCommandStore; import accord.local.Command; import accord.local.CommandStore; import accord.local.CommandStores; @@ -356,32 +355,6 @@ public class AccordTestUtils } } - public static InMemoryCommandStore.Synchronized createInMemoryCommandStore(LongSupplier now, String keyspace, String table) - { - TableMetadata metadata = Schema.instance.getTableMetadata(keyspace, table); - TokenRange range = TokenRange.fullRange(metadata.id); - Node.Id node = new Id(1); - NodeCommandStoreService time = new NodeCommandStoreService() - { - private ToLongFunction elapsed = TimeService.elapsedWrapperFromNonMonotonicSource(TimeUnit.MICROSECONDS, this::now); - - @Override public Id id() { return node;} - @Override public Timeouts timeouts() { return null; } - @Override public DurableBefore durableBefore() { return DurableBefore.EMPTY; } - @Override public long epoch() {return 1; } - @Override public long now() {return now.getAsLong(); } - @Override public Timestamp uniqueNow() { return uniqueNow(Timestamp.NONE); } - @Override public Timestamp uniqueNow(Timestamp atLeast) { return Timestamp.fromValues(1, now.getAsLong(), node); } - @Override public long elapsed(TimeUnit timeUnit) { return elapsed.applyAsLong(timeUnit); } - }; - - SingleEpochRanges holder = new SingleEpochRanges(Ranges.of(range)); - InMemoryCommandStore.Synchronized result = new InMemoryCommandStore.Synchronized(0, time, new AccordAgent(), - null, null, cs -> null, holder); - holder.set(); - return result; - } - public static AccordCommandStore createAccordCommandStore( Node.Id node, LongSupplier now, Topology topology, ExecutorPlus loadExecutor, ExecutorPlus saveExecutor) { diff --git a/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java b/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java index 10a5b46334..66c8ead4e5 100644 --- a/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java +++ b/test/unit/org/apache/cassandra/service/accord/EpochSyncTest.java @@ -20,7 +20,6 @@ package org.apache.cassandra.service.accord; import java.net.UnknownHostException; import java.util.ArrayList; -import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.EnumMap; @@ -50,6 +49,7 @@ import org.slf4j.LoggerFactory; import accord.api.ConfigurationService; import accord.api.ConfigurationService.EpochReady; +import accord.api.Journal; import accord.api.LocalConfig; import accord.api.Scheduler; import accord.impl.SizeOfIntersectionSorter; @@ -671,7 +671,8 @@ public class EpochSyncTest // TODO (review): Should there be a real scheduler here? Is it possible to adapt the Scheduler interface to scheduler used in this test? this.topology = new TopologyManager(SizeOfIntersectionSorter.SUPPLIER, new TestAgent.RethrowAgent(), id, Scheduler.NEVER_RUN_SCHEDULED, TimeService.ofNonMonotonic(globalExecutor::currentTimeMillis, TimeUnit.MILLISECONDS), LocalConfig.DEFAULT); AccordConfigurationService.DiskStateManager instance = MockDiskStateManager.instance; - config = new AccordConfigurationService(node, messagingService, failureDetector, instance, scheduler); + Journal journal = null; // TODO + config = new AccordConfigurationService(node, messagingService, failureDetector, instance, scheduler, journal); config.registerListener(new ConfigurationService.Listener() { @Override @@ -699,9 +700,9 @@ public class EpochSyncTest } @Override - public void onRemoveNodes(long epoch, Collection removed) + public void onRemoveNode(long epoch, Node.Id removed) { - topology.onRemoveNodes(epoch, removed); + topology.onRemoveNode(epoch, removed); } @Override diff --git a/test/unit/org/apache/cassandra/service/accord/FetchMinEpochTest.java b/test/unit/org/apache/cassandra/service/accord/FetchMinEpochTest.java index d12f132238..7bfe09b15f 100644 --- a/test/unit/org/apache/cassandra/service/accord/FetchMinEpochTest.java +++ b/test/unit/org/apache/cassandra/service/accord/FetchMinEpochTest.java @@ -38,6 +38,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.RetrySpec; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.io.IVersionedSerializers; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.locator.InetAddressAndPort; @@ -118,8 +119,8 @@ public class FetchMinEpochTest assertThat(f).isNotDone(); cluster.processAll(); assertThat(f).isDone(); - MessageDelivery.MaxRetriesException maxRetries = getMaxRetriesException(f); - Assertions.assertThat(maxRetries.attempts).isEqualTo(expectedMaxAttempts); + MessageDelivery.FailedResponseException maxRetries = getFailedResponseException(f); + Assertions.assertThat(maxRetries.failure).isEqualTo(RequestFailure.TIMEOUT); }); } @@ -259,6 +260,29 @@ public class FetchMinEpochTest DatabaseDescriptor.setPartitionerUnsafe(partitioner); } + private static MessageDelivery.FailedResponseException getFailedResponseException(Future f) throws InterruptedException, ExecutionException + { + MessageDelivery.FailedResponseException exception; + try + { + f.get(); + Assert.fail("Future should have failed"); + throw new AssertionError("Unreachable"); + } + catch (ExecutionException e) + { + if (e.getCause() instanceof MessageDelivery.FailedResponseException) + { + exception = (MessageDelivery.FailedResponseException) e.getCause(); + } + else + { + throw e; + } + } + return exception; + } + private static MessageDelivery.MaxRetriesException getMaxRetriesException(Future f) throws InterruptedException, ExecutionException { MessageDelivery.MaxRetriesException maxRetries; diff --git a/test/unit/org/apache/cassandra/service/accord/LoggingDiskStateManager.java b/test/unit/org/apache/cassandra/service/accord/LoggingDiskStateManager.java index 7b8ce0e335..10f9af4f85 100644 --- a/test/unit/org/apache/cassandra/service/accord/LoggingDiskStateManager.java +++ b/test/unit/org/apache/cassandra/service/accord/LoggingDiskStateManager.java @@ -44,9 +44,9 @@ public class LoggingDiskStateManager implements AccordConfigurationService.DiskS } @Override - public AccordKeyspace.EpochDiskState loadTopologies(AccordKeyspace.TopologyLoadConsumer consumer) { + public AccordKeyspace.EpochDiskState loadLocalTopologyState(AccordKeyspace.TopologyLoadConsumer consumer) { logger.info("[node={}] Calling loadTopologies()", self); - return delegate.loadTopologies(consumer); + return delegate.loadLocalTopologyState(consumer); } @Override diff --git a/test/unit/org/apache/cassandra/service/accord/MockDiskStateManager.java b/test/unit/org/apache/cassandra/service/accord/MockDiskStateManager.java index 9e37602634..64bb9af180 100644 --- a/test/unit/org/apache/cassandra/service/accord/MockDiskStateManager.java +++ b/test/unit/org/apache/cassandra/service/accord/MockDiskStateManager.java @@ -29,7 +29,7 @@ public enum MockDiskStateManager implements AccordConfigurationService.DiskState instance; @Override - public AccordKeyspace.EpochDiskState loadTopologies(AccordKeyspace.TopologyLoadConsumer consumer) { + public AccordKeyspace.EpochDiskState loadLocalTopologyState(AccordKeyspace.TopologyLoadConsumer consumer) { return AccordKeyspace.EpochDiskState.EMPTY; } diff --git a/test/unit/org/apache/cassandra/utils/SimulatedMiniCluster.java b/test/unit/org/apache/cassandra/utils/SimulatedMiniCluster.java index c60a633069..ac31fa97a9 100644 --- a/test/unit/org/apache/cassandra/utils/SimulatedMiniCluster.java +++ b/test/unit/org/apache/cassandra/utils/SimulatedMiniCluster.java @@ -337,7 +337,7 @@ public class SimulatedMiniCluster case MIGRATION: case MISC: case TRACING: - case FETCH_LOG: + case FETCH_METADATA: stage.unsafeSetExecutor(orderedExecutor); break; default: