diff --git a/modules/accord b/modules/accord index 05b2f2415c..20df7e5e5a 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 05b2f2415c3e7c109158208824fbc804dfe43f71 +Subproject commit 20df7e5e5ad7bae442ff6fe4ee716eb7104336c1 diff --git a/src/java/org/apache/cassandra/config/AccordSpec.java b/src/java/org/apache/cassandra/config/AccordSpec.java index c34d3fab49..94ff4e946a 100644 --- a/src/java/org/apache/cassandra/config/AccordSpec.java +++ b/src/java/org/apache/cassandra/config/AccordSpec.java @@ -121,7 +121,9 @@ public class AccordSpec public volatile OptionaldPositiveInt max_queued_loads = OptionaldPositiveInt.UNDEFINED; public volatile OptionaldPositiveInt max_queued_range_loads = OptionaldPositiveInt.UNDEFINED; - public volatile OptionaldPositiveInt max_progress_log_concurrency = OptionaldPositiveInt.UNDEFINED; + + public volatile OptionaldPositiveInt progress_log_concurrency = OptionaldPositiveInt.UNDEFINED; + public DurationSpec.IntMillisecondsBound progress_log_query_fallback_timeout = new DurationSpec.IntMillisecondsBound("1m"); public DataStorageSpec.LongMebibytesBound cache_size = null; public DataStorageSpec.LongMebibytesBound working_set_size = null; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 4f4dcf720b..f38dc0e9d5 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -67,6 +67,7 @@ import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import accord.impl.progresslog.DefaultProgressLog; import com.googlecode.concurrenttrees.common.Iterables; import org.apache.cassandra.audit.AuditLogOptions; import org.apache.cassandra.auth.AllowAllInternodeAuthenticator; @@ -189,6 +190,7 @@ public class DatabaseDescriptor private static final int MAX_NUM_TOKENS = 1536; private static Config conf; + private static DefaultProgressLog.Config accordProgressLogConfig; /** * Request timeouts can not be less than below defined value (see CASSANDRA-9375) @@ -561,6 +563,8 @@ public class DatabaseDescriptor applyGuardrails(); + applyAccordProgressLog(); + applyStartupChecks(); } @@ -1297,6 +1301,20 @@ public class DatabaseDescriptor } } + private static void applyAccordProgressLog() + { + try + { + accordProgressLogConfig = new DefaultProgressLog.Config(); + accordProgressLogConfig.concurrency = conf.accord.progress_log_concurrency.or(32); + accordProgressLogConfig.maxActiveRunTime = conf.accord.progress_log_query_fallback_timeout.toDuration(); + } + catch (IllegalArgumentException e) + { + throw new ConfigurationException("Invalid accord progress log configuration: " + e.getMessage(), e); + } + } + public static StartupChecksOptions getStartupChecksOptions() { return startupChecksOptions; @@ -5394,9 +5412,9 @@ public class DatabaseDescriptor return conf.accord.max_queued_range_loads.or(Math.max(4, getAccordConcurrentOps() / 4)); } - public static int getAccordProgressLogMaxConcurrency() + public static DefaultProgressLog.Config getAccordProgressLogConfig() { - return conf.accord.max_progress_log_concurrency.or(64); + return accordProgressLogConfig; } public static boolean getAccordCacheShrinkingOn() diff --git a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java index ae72d01d36..fa3101741e 100644 --- a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java @@ -32,9 +32,12 @@ import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.function.BiConsumer; import java.util.function.Function; import java.util.stream.Collectors; +import javax.annotation.Nullable; + import com.google.common.collect.BoundType; import com.google.common.collect.Range; import com.google.common.collect.Sets; @@ -43,6 +46,10 @@ import org.slf4j.LoggerFactory; import accord.api.RoutingKey; import accord.api.TraceEventType; +import accord.coordinate.FetchData; +import accord.coordinate.FetchRoute; +import accord.coordinate.MaybeRecover; +import accord.coordinate.RecoverWithRoute; import accord.impl.CommandChange; import accord.impl.progresslog.DefaultProgressLog; import accord.impl.progresslog.TxnStateKind; @@ -50,11 +57,13 @@ import accord.local.Cleanup; import accord.local.Command; import accord.local.CommandStore; import accord.local.CommandStores; +import accord.local.CommandStores.LatentStoreSelector; import accord.local.Commands; import accord.local.DurableBefore; import accord.local.LoadKeys; import accord.local.LoadKeysFor; import accord.local.MaxConflicts; +import accord.local.Node; import accord.local.PreLoadContext; import accord.local.RejectBefore; import accord.local.SafeCommand; @@ -63,14 +72,21 @@ import accord.local.StoreParticipants; import accord.local.cfk.CommandsForKey; import accord.local.cfk.SafeCommandsForKey; import accord.local.durability.ShardDurability; +import accord.primitives.FullRoute; +import accord.primitives.Known; import accord.primitives.Participants; +import accord.primitives.ProgressToken; +import accord.primitives.Route; import accord.primitives.SaveStatus; import accord.primitives.Status; +import accord.primitives.Timestamp; import accord.primitives.TxnId; import accord.utils.Invariants; import accord.utils.UnhandledEnum; import accord.utils.async.AsyncChain; import accord.utils.async.AsyncChains; +import accord.utils.async.AsyncResult; +import accord.utils.async.AsyncResults; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.cql3.statements.schema.CreateTableStatement; import org.apache.cassandra.db.ColumnFamilyStore; @@ -104,6 +120,8 @@ import org.apache.cassandra.service.consensus.migration.TableMigrationState; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.utils.LocalizeString; +import static accord.api.TraceEventType.RECOVER; +import static accord.coordinate.Infer.InvalidIf.NotKnownToBeInvalid; import static accord.local.RedundantStatus.Property.GC_BEFORE; import static accord.local.RedundantStatus.Property.LOCALLY_APPLIED; import static accord.local.RedundantStatus.Property.LOCALLY_DURABLE_TO_COMMAND_STORE; @@ -1201,7 +1219,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace public static final class TxnOpsTable extends AbstractMutableVirtualTable implements AbstractVirtualTable.DataSet { // TODO (expected): test each of these operations - enum Op { ERASE_VESTIGIAL, INVALIDATE, TRY_EXECUTE, FORCE_APPLY, FORCE_UPDATE } + enum Op { ERASE_VESTIGIAL, INVALIDATE, TRY_EXECUTE, FORCE_APPLY, FORCE_UPDATE, RECOVER, FETCH, RESET_PROGRESS_LOG } private TxnOpsTable() { super(parse(VIRTUAL_ACCORD_DEBUG, TXN_OPS, @@ -1209,7 +1227,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace "CREATE TABLE %s (\n" + " txn_id text,\n" + " command_store_id int,\n" + - " op text,\n" + + " op text," + " PRIMARY KEY (txn_id, command_store_id)" + ')', UTF8Type.instance)); } @@ -1279,6 +1297,76 @@ public class AccordDebugKeyspace extends VirtualKeyspace return Commands.applyChain(safeStore, (Command.Executed) command); }); break; + case FETCH: + runWithRoute(txnId, commandStoreId, command -> { + Timestamp executeAt = command.executeAtIfKnown(); + return (route, result) -> fetch(txnId, executeAt, route, result); + }); + break; + case RECOVER: + runWithRoute(txnId, commandStoreId, command -> (route, result) -> { + recover(txnId, route, result); + }); + break; + case RESET_PROGRESS_LOG: + run(txnId, commandStoreId, safeStore -> { + ((DefaultProgressLog)safeStore.progressLog()).requeue(safeStore, TxnStateKind.Waiting, txnId); + ((DefaultProgressLog)safeStore.progressLog()).requeue(safeStore, TxnStateKind.Home, txnId); + return AsyncChains.success(null); + }); + } + } + + private void runWithRoute(TxnId txnId, int commandStoreId, Function, AsyncResult.Settable>> apply) + { + run(txnId, commandStoreId, safeStore -> { + SafeCommand safeCommand = safeStore.unsafeGet(txnId); + Command command = safeCommand.current(); + if (command == null) + throw new InvalidRequestException(txnId + " not known"); + Node node = AccordService.instance().node(); + AsyncResult.Settable result = new AsyncResults.SettableResult<>(); + BiConsumer, AsyncResult.Settable> consumer = apply.apply(command); + if (command.route() == null) + { + FetchRoute.fetchRoute(node, txnId, command.maxContactable(), LatentStoreSelector.standard(), (success, fail) -> { + if (fail != null) result.setFailure(fail); + else consumer.accept(success, result); + }); + } + else + { + consumer.accept(command.route(), result); + } + return result; + }); + } + + private void fetch(TxnId txnId, Timestamp executeAtIfKnown, Route route, AsyncResult.Settable result) + { + Node node = AccordService.instance().node(); + FetchData.fetchSpecific(Known.Apply, node, txnId, executeAtIfKnown, route, route.withHomeKey(), LatentStoreSelector.standard(), (success, fail) -> { + if (fail != null) result.setFailure(fail); + else result.setSuccess(null); + }); + } + + private void recover(TxnId txnId, @Nullable Route route, AsyncResult.Settable result) + { + Node node = AccordService.instance().node(); + if (Route.isFullRoute(route)) + { + RecoverWithRoute.recover(node, node.someSequentialExecutor(), txnId, NotKnownToBeInvalid, (FullRoute) route, null, LatentStoreSelector.standard(), (success, fail) -> { + if (fail != null) result.setFailure(fail); + else result.setSuccess(null); + }, node.agent().trace(txnId, RECOVER)); + } + else + { + MaybeRecover.maybeRecover(node, txnId, NotKnownToBeInvalid, route, ProgressToken.NONE, LatentStoreSelector.standard(), (success, fail) -> { + if (fail != null) result.setFailure(fail); + else result.setSuccess(null); + }); } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index f337f8f738..972a7488e0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -181,8 +181,8 @@ public class AccordCommandStore extends CommandStore this.journal = journal; this.rangeSearcher = RangeSearcher.extractRangeSearcher(journal); this.sharedExecutor = sharedExecutor; - if (this.progressLog() instanceof DefaultProgressLog) - ((DefaultProgressLog)this.progressLog).setMaxConcurrency(DatabaseDescriptor.getAccordProgressLogMaxConcurrency()); + if (this.progressLog instanceof DefaultProgressLog) + ((DefaultProgressLog)this.progressLog).unsafeSetConfig(DatabaseDescriptor.getAccordProgressLogConfig()); final AccordCache.Type.Instance commands; final AccordCache.Type.Instance commandsForKey;