This is an automated email from the ASF dual-hosted git repository.
belliottsmith pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-accord.git
The following commit(s) were added to refs/heads/trunk by this push:
new 56b7710a Accord: Compute Dependencies Incrementally
56b7710a is described below
commit 56b7710afaf3ac1492e50669f8ef5c0985654586
Author: Benedict Elliott Smith <[email protected]>
AuthorDate: Sun Sep 6 21:25:04 2026 +0100
Accord: Compute Dependencies Incrementally
Improve tail latency by using the new incremental execution machinery for
computing dependency sets, so that the command store lock is not held for long
periods by sync points.
patch by Benedict; reviewed by Aleksey Yeschenko for CASSANDRA-21682
---
.../accord/coordinate/AbstractCoordination.java | 6 +-
.../java/accord/coordinate/CollectLatestDeps.java | 9 +-
.../accord/coordinate/CoordinateTransaction.java | 89 +++----
.../java/accord/impl/AbstractSafeCommandStore.java | 162 +++++++-----
.../java/accord/impl/InMemoryCommandStore.java | 24 +-
.../src/main/java/accord/local/CommandStore.java | 4 +-
.../src/main/java/accord/local/CommandStores.java | 6 +-
.../main/java/accord/local/CommandSummaries.java | 69 +++--
.../src/main/java/accord/local/DepsCalculator.java | 279 ++++++++++++++++++---
.../main/java/accord/local/ExecutionContext.java | 48 ++--
.../local/{LoadKeysFor.java => FindKeys.java} | 16 +-
.../src/main/java/accord/local/LoadKeys.java | 14 +-
.../java/accord/local/MapReduceCommandStores.java | 2 +
.../main/java/accord/local/SafeCommandStore.java | 17 +-
.../accord/local/durability/DurabilityLevel.java | 4 +-
.../src/main/java/accord/messages/Accept.java | 195 ++++++++++----
.../src/main/java/accord/messages/Apply.java | 2 +-
.../accord/messages/ApplyThenWaitUntilApplied.java | 8 +
.../java/accord/messages/BeginInvalidation.java | 6 +
.../main/java/accord/messages/BeginRecovery.java | 31 ++-
.../java/accord/messages/GetEphemeralReadDeps.java | 74 ++++--
.../main/java/accord/messages/GetLatestDeps.java | 101 +++++---
.../main/java/accord/messages/GetMaxConflict.java | 6 +
.../main/java/accord/messages/InformDurable.java | 6 +
.../main/java/accord/messages/NoWaitRequest.java | 38 +--
.../java/accord/messages/ParticipantsRequest.java | 2 +-
.../src/main/java/accord/messages/PreAccept.java | 96 ++++---
.../src/main/java/accord/messages/ReadData.java | 9 +-
.../src/main/java/accord/messages/ReplyList.java | 93 +++++++
.../main/java/accord/messages/RouteRequest.java | 4 +-
.../src/main/java/accord/primitives/Deps.java | 28 +--
.../main/java/accord/utils/async/AsyncChain.java | 7 +-
.../main/java/accord/utils/async/AsyncChains.java | 1 -
.../main/java/accord/utils/async/AsyncResults.java | 66 ++++-
.../accord/utils/async/CancellableAsyncResult.java | 23 ++
35 files changed, 1098 insertions(+), 447 deletions(-)
diff --git
a/accord-core/src/main/java/accord/coordinate/AbstractCoordination.java
b/accord-core/src/main/java/accord/coordinate/AbstractCoordination.java
index 0aaec1b0..cee085cf 100644
--- a/accord-core/src/main/java/accord/coordinate/AbstractCoordination.java
+++ b/accord-core/src/main/java/accord/coordinate/AbstractCoordination.java
@@ -461,7 +461,7 @@ public abstract class AbstractCoordination<P extends
Participants<?>, Result, Re
}
enum LocalExecuteState { PENDING, SUCCESS, TIMEOUT }
- abstract class AbstractLocalExecute extends
MapReduceConsumeCommandStores<P, Reply> implements Timeouts.Timeout
+ abstract class AbstractLocalExecute<R> extends
MapReduceConsumeCommandStores<P, R> implements Timeouts.Timeout
{
LocalExecuteState state = PENDING;
Cancellable cancel;
@@ -469,7 +469,7 @@ public abstract class AbstractCoordination<P extends
Participants<?>, Result, Re
abstract long expiresAt();
abstract Cancellable submit();
- abstract void acceptInternal(Reply result, Throwable failure);
+ abstract void acceptInternal(R result, Throwable failure);
protected AbstractLocalExecute()
{
@@ -504,7 +504,7 @@ public abstract class AbstractCoordination<P extends
Participants<?>, Result, Re
}
@Override
- public void accept(Reply result, Throwable failure)
+ public void accept(R result, Throwable failure)
{
done();
executor.executeMaybeImmediately(() -> acceptInternal(result,
failure));
diff --git a/accord-core/src/main/java/accord/coordinate/CollectLatestDeps.java
b/accord-core/src/main/java/accord/coordinate/CollectLatestDeps.java
index f02d2240..6b1329c7 100644
--- a/accord-core/src/main/java/accord/coordinate/CollectLatestDeps.java
+++ b/accord-core/src/main/java/accord/coordinate/CollectLatestDeps.java
@@ -29,7 +29,6 @@ import accord.coordinate.tracking.QuorumTracker;
import accord.local.Node;
import accord.local.Node.Id;
import accord.messages.GetLatestDeps;
-import accord.messages.GetLatestDeps.GetLatestDepsOk;
import accord.messages.GetLatestDeps.GetLatestDepsReply;
import accord.primitives.Ballot;
import accord.primitives.FullRoute;
@@ -46,7 +45,7 @@ import static accord.coordinate.tracking.RequestStatus.Failed;
import static accord.coordinate.tracking.RequestStatus.Success;
import static accord.primitives.Routables.Slice.Minimal;
-public class CollectLatestDeps extends AbstractCoordination<Route<?>,
List<LatestDeps>, GetLatestDepsReply, GetLatestDepsOk>
+public class CollectLatestDeps extends AbstractCoordination<Route<?>,
List<LatestDeps>, GetLatestDepsReply, GetLatestDepsReply>
{
final Timestamp executeAt;
final @Nullable Ballot ballot;
@@ -90,7 +89,7 @@ public class CollectLatestDeps extends
AbstractCoordination<Route<?>, List<Lates
{
if (ok.isOk())
{
- recordOk(fromIndex, (GetLatestDepsOk) ok);
+ recordOk(fromIndex, ok);
if (tracker.recordSuccess(from) == Success)
onQuorum();
}
@@ -111,9 +110,9 @@ public class CollectLatestDeps extends
AbstractCoordination<Route<?>, List<Lates
private void onQuorum()
{
Invariants.require(!isDone());
- SortedListMap<Node.Id, GetLatestDepsOk> oks = finishOks();
+ SortedListMap<Node.Id, GetLatestDepsReply> oks = finishOks();
List<LatestDeps> result = new ArrayList<>(oks.size());
- for (GetLatestDepsOk ok : oks.values())
+ for (GetLatestDepsReply ok : oks.values())
result.add(ok.deps);
finishWithSuccess(result);
}
diff --git
a/accord-core/src/main/java/accord/coordinate/CoordinateTransaction.java
b/accord-core/src/main/java/accord/coordinate/CoordinateTransaction.java
index b4cf9ba4..dc25e6c9 100644
--- a/accord-core/src/main/java/accord/coordinate/CoordinateTransaction.java
+++ b/accord-core/src/main/java/accord/coordinate/CoordinateTransaction.java
@@ -20,33 +20,33 @@ package accord.coordinate;
import java.util.List;
import java.util.function.BiConsumer;
-
import javax.annotation.Nullable;
+import accord.api.ExclusiveAsyncExecutor;
import accord.api.ProtocolModifiers;
import accord.api.Result;
import accord.coordinate.CoordinationAdapter.Adapters;
import accord.coordinate.ExecuteFlag.CoordinationFlags;
import accord.coordinate.ExecuteFlag.ExecuteFlags;
import accord.local.Commands;
-import accord.local.DepsCalculator;
import accord.local.LoadKeys;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
+import accord.local.Node;
import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
-import accord.api.ExclusiveAsyncExecutor;
import accord.local.StoreParticipants;
+import accord.messages.PreAccept.PreAcceptDepsCalculator;
import accord.messages.PreAccept.PreAcceptNack;
-import accord.messages.PreAccept.PreAcceptReply;
-import accord.topology.Topologies;
-import accord.local.Node;
import accord.messages.PreAccept.PreAcceptOk;
+import accord.messages.PreAccept.PreAcceptReply;
+import accord.messages.ReplyList;
import accord.primitives.Ballot;
import accord.primitives.Deps;
import accord.primitives.FullRoute;
import accord.primitives.Timestamp;
import accord.primitives.Txn;
import accord.primitives.TxnId;
+import accord.topology.Topologies;
import accord.utils.SortedListMap;
import accord.utils.async.AsyncChain;
import accord.utils.async.AsyncChains;
@@ -60,8 +60,8 @@ import static accord.messages.Accept.Kind.MEDIUM;
import static accord.messages.Accept.Kind.SLOW;
import static accord.messages.MessageType.StandardMessage.PRE_ACCEPT_REQ;
import static accord.primitives.Timestamp.Flag.REJECTED;
-import static accord.primitives.Timestamp.mergeMaxAndFlags;
import static accord.primitives.Timestamp.Flag.SOFT_REJECT;
+import static accord.primitives.Timestamp.mergeMaxAndFlags;
import static accord.primitives.TxnId.FastPath.PrivilegedCoordinatorWithDeps;
import static accord.topology.SelectShards.LIVE;
import static java.util.concurrent.TimeUnit.MICROSECONDS;
@@ -209,7 +209,7 @@ public class CoordinateTransaction extends
CoordinatePreAccept<Result>
return node.coordinationAdapter(txnId, Standard);
}
- class LocalExecute extends AbstractLocalExecute
+ class LocalExecute extends AbstractLocalExecute<ReplyList<PreAcceptReply>>
{
@Override
long expiresAt()
@@ -224,7 +224,7 @@ public class CoordinateTransaction extends
CoordinatePreAccept<Result>
}
@Override
- public void acceptInternal(PreAcceptReply result, Throwable failure)
+ public void acceptInternal(ReplyList<PreAcceptReply> replies,
Throwable failure)
{
if (failure != null)
{
@@ -236,56 +236,49 @@ public class CoordinateTransaction extends
CoordinatePreAccept<Result>
}
else
{
- if (result.isOk())
- {
- PreAcceptOk ok = (PreAcceptOk) result;
- // TODO (desired): we can probably still process and
record fast path votes from peers, just with different quorum requirements
- boolean hasCoordinatorVote = txnId.equals(ok.witnessedAt);
- if (!hasCoordinatorVote &&
txnId.hasPrivilegedCoordinator()) fastPathEnabled = false;
- Deps deps = hasCoordinatorVote &&
txnId.is(PrivilegedCoordinatorWithDeps) ? ok.deps : null;
- contactNotSelf(deps, hasCoordinatorVote);
- onSuccess(node.id(), ok);
- }
- else
- {
- finishOnFailure(Preempted.preempted(node.agent(), txnId,
scope.homeKey()));
- }
+ ReplyList.invoke(replies, PreAcceptReply::reduce, (success,
fail) -> {
+ if (success.isOk())
+ {
+ PreAcceptOk ok = (PreAcceptOk) success;
+ boolean hasCoordinatorVote =
txnId.equals(ok.witnessedAt);
+ if (!hasCoordinatorVote &&
txnId.hasPrivilegedCoordinator()) fastPathEnabled = false;
+ Deps deps = hasCoordinatorVote &&
txnId.is(PrivilegedCoordinatorWithDeps) ? ok.deps : null;
+ contactNotSelf(deps, hasCoordinatorVote);
+ onSuccess(node.id(), ok);
+ }
+ else
+ {
+ finishOnFailure(Preempted.preempted(node.agent(),
txnId, scope.homeKey()));
+ }
+ });
}
}
@Override
- public PreAcceptReply applyInternal(SafeCommandStore safeStore)
+ public ReplyList<PreAcceptReply> applyInternal(SafeCommandStore
safeStore)
{
long minEpoch = topologies.oldestEpoch();
StoreParticipants participants =
StoreParticipants.update(safeStore, scope, minEpoch, txnId, txnId.epoch());
SafeCommand safeCommand = safeStore.get(txnId, participants);
- Timestamp executeAt;
- Deps deps;
- ExecuteFlags flags;
- try (DepsCalculator calculator = new DepsCalculator(txnId))
- {
- deps = calculator.calculate(safeStore, txnId, participants,
minEpoch, txnId, true);
- if (deps == null)
- return PreAcceptNack.INSTANCE;
-
- boolean hasCoordinatorVote = txnId.hasPrivilegedCoordinator();
- Deps coordinatorDeps = txnId.is(PrivilegedCoordinatorWithDeps)
? deps : null;
- Commands.AcceptOutcome outcome = Commands.preaccept(safeStore,
safeCommand, participants, txnId, txn, coordinatorDeps, hasCoordinatorVote);
- if (outcome != Success)
- return PreAcceptNack.INSTANCE;
-
- executeAt = calculator.executeAt(safeCommand, node);
- flags = calculator.executeFlags(txnId);
- }
+ boolean hasCoordinatorVote = txnId.hasPrivilegedCoordinator();
+ Commands.AcceptOutcome outcome = Commands.preaccept(safeStore,
safeCommand, participants, txnId, txn, null, hasCoordinatorVote);
+ if (outcome != Success)
+ return PreAcceptNack.INSTANCE;
- return new PreAcceptOk(txnId, executeAt, deps, flags);
+ Timestamp witnessedAt = safeCommand.current().executeAt;
+ //noinspection resource
+ PreAcceptDepsCalculator calculator = new
PreAcceptDepsCalculator(txnId, witnessedAt, participants, node);
+ ReplyList<PreAcceptReply> reply = calculator.calculate(safeStore,
minEpoch, true);
+ if (reply == null)
+ reply = PreAcceptNack.INSTANCE;
+ return reply;
}
@Override
- public PreAcceptReply reduce(PreAcceptReply r1, PreAcceptReply r2)
+ public ReplyList<PreAcceptReply> reduce(ReplyList<PreAcceptReply> r1,
ReplyList<PreAcceptReply> r2)
{
- return PreAcceptReply.reduce(r1, r2);
+ return ReplyList.merge(r1, r2);
}
@Override
@@ -295,9 +288,9 @@ public class CoordinateTransaction extends
CoordinatePreAccept<Result>
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
- return LoadKeysFor.READ_WRITE;
+ return FindKeys.CONFLICTS;
}
public ExecutionKind executionKind()
diff --git
a/accord-core/src/main/java/accord/impl/AbstractSafeCommandStore.java
b/accord-core/src/main/java/accord/impl/AbstractSafeCommandStore.java
index 6fd9e368..213b6180 100644
--- a/accord-core/src/main/java/accord/impl/AbstractSafeCommandStore.java
+++ b/accord-core/src/main/java/accord/impl/AbstractSafeCommandStore.java
@@ -22,7 +22,11 @@ import java.util.ArrayList;
import java.util.List;
import java.util.NavigableMap;
+import javax.annotation.Nonnull;
+import javax.annotation.Nullable;
+
import accord.api.RoutingKey;
+import accord.local.ExecutionContext.OverrideKeys;
import accord.local.LoadKeys;
import accord.local.ExecutionContext;
import accord.local.RedundantBefore;
@@ -30,7 +34,7 @@ import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
import accord.local.cfk.SafeCommandsForKey;
import accord.primitives.Ranges;
-import accord.primitives.Routable;
+import accord.primitives.Routable.Domain;
import accord.primitives.RoutingKeys;
import accord.primitives.Timestamp;
import accord.primitives.TxnId;
@@ -67,81 +71,95 @@ extends SafeCommandStore
@Override
public ExecutionContext canExecute(ExecutionContext with)
{
- if (with.isEmpty()) return with;
- if (with.keys().domain() == Routable.Domain.Range)
- return with.isSubsetOf(this.context) ? with : null;
-
- LoadKeys require = with.loadKeys();
- if (require != LoadKeys.NONE)
- {
- ExecutionContext context = context();
- if (!context.loadKeys().satisfiesIfPresent(require))
- return null;
+ Unseekables<?> withKeys = with.keys();
+ if (withKeys.domain() == Domain.Range)
+ return with.isSubsetOf(context) ? with : null;
- if (with.loadKeysFor().compareTo(context.loadKeysFor()) > 0)
- return null;
- }
+ LoadKeys loadKeys = with.loadKeys();
+ if (loadKeys != LoadKeys.NONE &&
with.findKeys().compareTo(context.findKeys()) > 0)
+ return null;
- try (Caches caches = tryGetCaches())
+ Caches caches = null;
+ try
{
- for (TxnId txnId : with.txnIds())
+ TxnId primaryTxnId = with.primaryTxnId();
+ if (primaryTxnId != null)
{
- if (null != getInternal(txnId))
- continue;
-
- if (caches == null)
- return null;
+ if (!isPresent(primaryTxnId))
+ {
+ caches = tryGetCaches();
+ if (ifLoadedInternal(primaryTxnId) == null)
+ return null;
+ }
- C safeCommand = caches.acquireIfLoaded(txnId);
- if (safeCommand == null)
- return null;
+ TxnId additionalTxnId = with.additionalTxnId();
+ if (additionalTxnId != null && !isPresent(additionalTxnId))
+ {
+ if (caches == null)
+ caches = tryGetCaches();
- add(safeCommand, caches);
+ if (ifLoadedInternal(additionalTxnId) == null)
+ return null;
+ }
}
- LoadKeys loadKeys = with.loadKeys();
- if (loadKeys == LoadKeys.NONE)
+ if (loadKeys == LoadKeys.NONE || withKeys.isEmpty())
return with;
List<RoutingKey> unavailable = null;
- Unseekables<?> keys = with.keys();
- if (keys.isEmpty())
- return with;
-
- for (int i = 0 ; i < keys.size() ; ++i)
+ for (int i = 0 ; i < withKeys.size() ; ++i)
{
- RoutingKey key = (RoutingKey) keys.get(i);
- if (null != getInternal(key))
- continue; // already in working set
+ RoutingKey key = (RoutingKey) withKeys.get(i);
+ if (isPresent(key))
+ continue;
+
+ if (unavailable == null && caches == null)
+ caches = tryGetCaches();
+
+ if (ifLoadedInternal(caches, key) != null)
+ continue;
- if (caches != null)
- {
- CFK safeCfk = caches.acquireIfLoaded(key);
- if (safeCfk != null)
- {
- add(safeCfk, caches);
- continue;
- }
- }
if (unavailable == null)
unavailable = new ArrayList<>();
+
unavailable.add(key);
}
if (unavailable == null)
return with;
- if (unavailable.size() == keys.size())
+ if (unavailable.size() == withKeys.size())
return null;
- return ExecutionContext.contextFor(with.primaryTxnId(),
with.additionalTxnId(), keys.without(RoutingKeys.ofSortedUnique(unavailable)),
loadKeys, context.loadKeysFor(), context.reason());
+ return new OverrideKeys(with,
withKeys.without(RoutingKeys.ofSortedUnique(unavailable)));
+ }
+ finally
+ {
+ if (caches != null)
+ caches.close();
}
}
- @Override
- public ExecutionContext context()
+ private boolean isPresent(TxnId txnId)
{
- return context;
+ return getInternal(txnId) != null;
+ }
+
+ private boolean isPresentOrLoadedInternal(@Nonnull Caches caches, TxnId
txnId)
+ {
+ return getInternal(txnId) != null || ifLoadedInternal(caches, txnId)
!= null;
+ }
+
+ private C ifLoadedInternal(@Nullable Caches caches, TxnId txnId)
+ {
+ if (caches == null)
+ return null;
+
+ C command = caches.acquireIfLoaded(txnId);
+ if (command == null)
+ return null;
+
+ return add(command, caches);
}
@Override
@@ -149,33 +167,47 @@ extends SafeCommandStore
{
try (Caches caches = tryGetCaches())
{
- if (caches == null)
- return null;
+ return ifLoadedInternal(caches, txnId);
+ }
+ }
- C command = caches.acquireIfLoaded(txnId);
- if (command == null)
- return null;
+ private boolean isPresentOrLoadedInternal(@Nonnull Caches caches,
RoutingKey key)
+ {
+ return getInternal(key) != null || ifLoadedInternal(caches, key) !=
null;
+ }
- return add(command, caches);
- }
+ private boolean isPresent(RoutingKey key)
+ {
+ return getInternal(key) != null;
+ }
+
+ protected CFK ifLoadedInternal(@Nullable Caches caches, RoutingKey key)
+ {
+ if (caches == null)
+ return null;
+
+ CFK cfk = caches.acquireIfLoaded(key);
+ if (cfk == null)
+ return null;
+
+ return add(cfk, caches);
}
@Override
- protected CFK ifLoadedInternal(RoutingKey txnId)
+ protected CFK ifLoadedInternal(RoutingKey key)
{
try (Caches caches = tryGetCaches())
{
- if (caches == null)
- return null;
-
- CFK cfk = caches.acquireIfLoaded(txnId);
- if (cfk == null)
- return null;
-
- return add(cfk, caches);
+ return ifLoadedInternal(caches, key);
}
}
+ @Override
+ public ExecutionContext context()
+ {
+ return context;
+ }
+
// TODO (expected): cleanup the integration hooks here; they're a bit
byzantine. Also clearly document behaviour.
public void postExecute()
{
diff --git a/accord-core/src/main/java/accord/impl/InMemoryCommandStore.java
b/accord-core/src/main/java/accord/impl/InMemoryCommandStore.java
index 62b1757f..b413c01b 100644
--- a/accord-core/src/main/java/accord/impl/InMemoryCommandStore.java
+++ b/accord-core/src/main/java/accord/impl/InMemoryCommandStore.java
@@ -67,7 +67,7 @@ import accord.local.CommandSummaries;
import accord.local.CommandSummaries.Summary;
import accord.local.CommandSummaries.SummaryLoader;
import accord.local.Commands;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
import accord.local.MaxDecidedRX;
import accord.local.NodeCommandStoreService;
import accord.local.ExecutionContext;
@@ -92,8 +92,8 @@ import static
accord.api.ProtocolModifiers.isRangeEndInclusive;
import static accord.api.ProtocolModifiers.isRangeStartInclusive;
import static accord.local.Cleanup.Input.FULL;
import static accord.local.LoadKeys.NONE;
-import static accord.local.LoadKeysFor.RECOVERY;
-import static accord.local.LoadKeysFor.WRITE;
+import static accord.local.FindKeys.SUPERSEDING;
+import static accord.local.FindKeys.DECLARED;
import static accord.local.RedundantStatus.Coverage.ALL;
import static accord.local.StoreParticipants.Filter.LOAD;
import static accord.primitives.Routable.Domain.Key;
@@ -456,7 +456,7 @@ public abstract class InMemoryCommandStore extends
CommandStore
commandsForKey.put(key, safeCfk);
}
}
- else if (context.loadKeysFor() != WRITE)
+ else if (context.findKeys() != DECLARED)
{
SummaryLoader loader =
SummaryLoader.loader(unsafeGetRedundantBefore(), unsafeGetMaxDecidedRX(),
context);
for (GlobalCommandsForKey global :
this.commandsForKey.values())
@@ -839,7 +839,7 @@ public abstract class InMemoryCommandStore extends
CommandStore
if (commandsForRanges != null)
return commandsForRanges;
- Invariants.require(context.loadKeysFor() != WRITE);
+ Invariants.require(context.findKeys() != DECLARED);
MaxDecidedRX maxDecidedRX = commandStore().unsafeGetMaxDecidedRX();
SummaryLoader loader = cfrLoad != null ? cfrLoad.loader
:
SummaryLoader.loader(redundantBefore(), maxDecidedRX, context);
@@ -852,11 +852,11 @@ public abstract class InMemoryCommandStore extends
CommandStore
return commandsForRanges = () -> loaded;
}
- private boolean visitForKey(Unseekables<?> keysOrRanges,
Predicate<CommandsForKey> forEach)
+ private boolean visitForKey(@Nullable Unseekables<?> keysOrRanges,
Predicate<CommandsForKey> forEach)
{
for (SafeCommandsForKey safeCfk : commandsForKey.values())
{
- if (!keysOrRanges.contains(safeCfk.key()))
+ if (keysOrRanges != null &&
!keysOrRanges.contains(safeCfk.key()))
continue;
if (!forEach.test(safeCfk.current()))
@@ -865,18 +865,18 @@ public abstract class InMemoryCommandStore extends
CommandStore
return true;
}
- private <P1, P2> void visitForKey(Unseekables<?> keysOrRanges,
Timestamp startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visitor,
P1 p1, P2 p2)
+ private <P1, P2> void visitForKey(@Nullable Unseekables<?>
keysOrRanges, Timestamp startedBefore, Kinds testKind, ActiveCommandVisitor<P1,
P2> visitor, P1 p1, P2 p2)
{
visitForKey(keysOrRanges, cfk -> { cfk.visit(startedBefore,
testKind, visitor, p1, p2); return true; });
}
- public boolean visitForKey(Unseekables<?> keysOrRanges, TxnId
testTxnId, Kinds testKind, SupersedingCommandVisitor visit)
+ public boolean visitForKey(@Nullable Unseekables<?> keysOrRanges,
TxnId testTxnId, Kinds testKind, SupersedingCommandVisitor visit)
{
return visitForKey(keysOrRanges, cfk -> cfk.visit(testTxnId,
testKind, visit));
}
@Override
- public <P1, P2> void visit(Unseekables<?> keysOrRanges, Timestamp
startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visitor, P1 p1, P2
p2)
+ public <P1, P2> void visit(@Nullable Unseekables<?> keysOrRanges,
Timestamp startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visitor,
P1 p1, P2 p2)
{
visitForKey(keysOrRanges, startedBefore, testKind, visitor, p1,
p2);
commandsForRanges().visit(keysOrRanges, startedBefore, testKind,
visitor, p1, p2);
@@ -937,14 +937,14 @@ public abstract class InMemoryCommandStore extends
CommandStore
protected CommandsForRangeLoad cfrLoad(ExecutionContext context)
{
- if (context.loadKeysFor() != LoadKeysFor.RECOVERY)
+ if (context.findKeys() != FindKeys.SUPERSEDING)
return null;
SummaryLoader loader =
SummaryLoader.loader(unsafeGetRedundantBefore(), unsafeGetMaxDecidedRX(),
context);
commandsForRanges.populateMinFutureRx(loader);
TreeMap<Timestamp, Summary> loaded = new TreeMap<>();
commandsForRanges.search(loader, null, txnId -> {
- Invariants.require(loader.loadKeysFor() == RECOVERY);
+ Invariants.require(loader.loadKeysFor() == SUPERSEDING);
Command command = commands.get(txnId).value();
Summary summary = loader.ifRelevant(command);
// TODO (expected): prune implied invalidations from index, so no
need to special case
diff --git a/accord-core/src/main/java/accord/local/CommandStore.java
b/accord-core/src/main/java/accord/local/CommandStore.java
index 20d9b4fe..12d9a6d9 100644
--- a/accord-core/src/main/java/accord/local/CommandStore.java
+++ b/accord-core/src/main/java/accord/local/CommandStore.java
@@ -64,6 +64,8 @@ import accord.primitives.Status.Durability.HasOutcome;
import accord.utils.DeterministicIdentitySet;
import accord.utils.Invariants;
import accord.utils.Reduce;
+import accord.utils.SortedArrays;
+import accord.utils.SortedArrays.SortedArrayList;
import accord.utils.UnhandledEnum;
import accord.utils.async.AsyncChain;
import accord.utils.async.AsyncChains;
@@ -525,7 +527,7 @@ public abstract class CommandStore implements
AbstractAsyncExecutor, ExclusiveAs
dataStore.ensureDurable(this, ranges, addOnDataStoreDurable, 0);
ensureDurable(ranges, addOnCommandStoreDurable);
Ranges unavailable = unavailable(txnId, txnIdWithFlags, ranges,
safeStore.ranges(), safeStore.safeToReadAt());
- node.durability().report(new DurabilityResult(new
MinimalSyncPoint(txnId, txnIdWithFlags, route.without(unavailable)), new
DurabilityLevel(Self, NoRemote, null), null));
+ node.durability().report(new DurabilityResult(new
MinimalSyncPoint(txnId, txnIdWithFlags, route.without(unavailable)), new
DurabilityLevel(Self, NoRemote, SortedArrayList.ofSorted(node.id())), null));
}
/**
diff --git a/accord-core/src/main/java/accord/local/CommandStores.java
b/accord-core/src/main/java/accord/local/CommandStores.java
index ef9ce1a4..d0896cbc 100644
--- a/accord-core/src/main/java/accord/local/CommandStores.java
+++ b/accord-core/src/main/java/accord/local/CommandStores.java
@@ -1068,15 +1068,15 @@ public abstract class CommandStores implements
AsyncExecutorFactory
public AsyncChain<Void> forEach(String reason, TxnId txnId,
Participants<?> participants, long minEpoch, long maxEpoch,
Consumer<SafeCommandStore> forEach)
{
- return forEach(reason, txnId, participants, LoadKeys.SYNC,
LoadKeysFor.READ_WRITE, minEpoch, maxEpoch, forEach);
+ return forEach(reason, txnId, participants, LoadKeys.SYNC,
FindKeys.CONFLICTS, minEpoch, maxEpoch, forEach);
}
- public AsyncChain<Void> forEach(String reason, TxnId txnId,
Participants<?> participants, LoadKeys loadKeys, LoadKeysFor loadKeysFor, long
minEpoch, long maxEpoch, Consumer<SafeCommandStore> forEach)
+ public AsyncChain<Void> forEach(String reason, TxnId txnId,
Participants<?> participants, LoadKeys loadKeys, FindKeys findKeys, long
minEpoch, long maxEpoch, Consumer<SafeCommandStore> forEach)
{
return mapReduce(StoreFinder.selector(participants, minEpoch,
maxEpoch), new MapReduceCommandStores<Participants<?>, Void>(participants)
{
@Override public LoadKeys loadKeys() { return loadKeys;}
- @Override public LoadKeysFor loadKeysFor() { return loadKeysFor; }
+ @Override public FindKeys findKeys() { return findKeys; }
@Override public Void reduce(Void o1, Void o2) { return null; }
@Override public TxnId primaryTxnId() { return txnId; }
@Override public String reason() { return reason; }
diff --git a/accord-core/src/main/java/accord/local/CommandSummaries.java
b/accord-core/src/main/java/accord/local/CommandSummaries.java
index d1c1de98..a710f04e 100644
--- a/accord-core/src/main/java/accord/local/CommandSummaries.java
+++ b/accord-core/src/main/java/accord/local/CommandSummaries.java
@@ -53,8 +53,7 @@ import static
accord.local.CommandSummaries.SummaryStatus.ACCEPTED;
import static accord.local.CommandSummaries.SummaryStatus.APPLIED;
import static accord.local.CommandSummaries.SummaryStatus.NOTACCEPTED;
import static accord.local.CommandSummaries.SummaryStatus.PREACCEPTED;
-import static accord.local.LoadKeysFor.RECOVERY;
-import static accord.local.LoadKeysFor.WRITE;
+import static accord.local.FindKeys.DECLARED;
import static accord.local.MaxDecidedRX.forDeps;
import static accord.primitives.Known.KnownDeps.NoDeps;
import static accord.primitives.Routables.Slice.Minimal;
@@ -102,7 +101,6 @@ public interface CommandSummaries
private static final Relevance[] lookup = values();
public static final int ENCODED_MASK = 7;
- public static final int ENCODED_BITS = 3;
private final int encoded;
Relevance(int encoded)
@@ -238,7 +236,7 @@ public interface CommandSummaries
{
public interface Factory<L extends SummaryLoader>
{
- L create(RedundantBefore redundantBefore, @Nullable MaxDecidedRX
maxDecidedRX, TxnId primaryTxnId, Unseekables<?> searchKeysOrRanges, Kinds
testKind, TxnId minTxnId, Timestamp executeAt, LoadKeysFor loadKeysFor);
+ L create(RedundantBefore redundantBefore, @Nullable MaxDecidedRX
maxDecidedRX, TxnId primaryTxnId, Unseekables<?> searchKeysOrRanges, Kinds
testKind, TxnId minTxnId, Timestamp executeAt, FindKeys findKeys);
}
private static final ReducingRangeMap<TxnId> NO_RX = new
ReducingRangeMap<>();
@@ -248,7 +246,7 @@ public interface CommandSummaries
protected final Unseekables<?> searchFor;
// TODO (expected): separate out Kinds we need before/after
primaryTxnId/executeAt
protected final Kinds testKind;
- protected final LoadKeysFor loadKeysFor;
+ protected final FindKeys findKeys;
protected final TxnId primaryTxnId, minTxnId;
protected final DecidedRX decidedRx;
protected final Timestamp primaryExecuteAt;
@@ -259,30 +257,30 @@ public interface CommandSummaries
// TODO (expected): provide executeAt to PreLoadContext so we can more
aggressively filter what we load, esp. by Kind
public static SummaryLoader loader(RedundantBefore redundantBefore,
MaxDecidedRX maxDecidedRX, ExecutionContext context)
{
- return loader(redundantBefore, maxDecidedRX,
context.primaryTxnId(), context.executeAt(), context.loadKeysFor(),
context.keys());
+ return loader(redundantBefore, maxDecidedRX,
context.primaryTxnId(), context.executeAt(), context.findKeys(),
context.keys());
}
- public static SummaryLoader loader(RedundantBefore redundantBefore,
MaxDecidedRX maxDecidedRX, TxnId primaryTxnId, Timestamp executeAt, LoadKeysFor
loadKeysFor, Unseekables<?> keysOrRanges)
+ public static SummaryLoader loader(RedundantBefore redundantBefore,
MaxDecidedRX maxDecidedRX, TxnId primaryTxnId, Timestamp executeAt, FindKeys
findKeys, Unseekables<?> keysOrRanges)
{
- return loader(redundantBefore, maxDecidedRX, primaryTxnId,
executeAt, loadKeysFor, keysOrRanges, SummaryLoader::new);
+ return loader(redundantBefore, maxDecidedRX, primaryTxnId,
executeAt, findKeys, keysOrRanges, SummaryLoader::new);
}
public static <L extends SummaryLoader> L loader(RedundantBefore
redundantBefore, MaxDecidedRX maxDecidedRX, ExecutionContext context,
Factory<L> factory)
{
- return loader(redundantBefore, maxDecidedRX,
context.primaryTxnId(), context.executeAt(), context.loadKeysFor(),
context.keys(), factory);
+ return loader(redundantBefore, maxDecidedRX,
context.primaryTxnId(), context.executeAt(), context.findKeys(),
context.keys(), factory);
}
- public static <L extends SummaryLoader> L loader(RedundantBefore
redundantBefore, MaxDecidedRX maxDecidedRX, TxnId primaryTxnId, Timestamp
executeAt, LoadKeysFor loadKeysFor, Unseekables<?> keysOrRanges, Factory<L>
factory)
+ public static <L extends SummaryLoader> L loader(RedundantBefore
redundantBefore, MaxDecidedRX maxDecidedRX, TxnId primaryTxnId, Timestamp
executeAt, FindKeys findKeys, Unseekables<?> keysOrRanges, Factory<L> factory)
{
Invariants.require(primaryTxnId != null);
TxnId minTxnId = redundantBefore.min(keysOrRanges,
Bounds::gcBefore);
- Kinds kinds = primaryTxnId.witnesses().or(loadKeysFor == RECOVERY
? primaryTxnId.witnessedBy() : Nothing);
+ Kinds kinds = primaryTxnId.witnesses().or(findKeys ==
FindKeys.SUPERSEDING ? primaryTxnId.witnessedBy() : Nothing);
if (!primaryTxnId.is(Txn.Kind.ExclusiveSyncPoint)) // the main
distinction between RX and RV is that RV doesn't filter out decided transactions
maxDecidedRX = null;
- return factory.create(redundantBefore, maxDecidedRX, primaryTxnId,
keysOrRanges, kinds, minTxnId, executeAt, loadKeysFor);
+ return factory.create(redundantBefore, maxDecidedRX, primaryTxnId,
keysOrRanges, kinds, minTxnId, executeAt, findKeys);
}
- public SummaryLoader(RedundantBefore redundantBefore, MaxDecidedRX
maxDecidedRX, TxnId primaryTxnId, Unseekables<?> searchFor, Kinds testKind,
TxnId minTxnId, Timestamp primaryExecuteAt, LoadKeysFor loadKeysFor)
+ public SummaryLoader(RedundantBefore redundantBefore, MaxDecidedRX
maxDecidedRX, TxnId primaryTxnId, Unseekables<?> searchFor, Kinds testKind,
TxnId minTxnId, Timestamp primaryExecuteAt, FindKeys findKeys)
{
this.redundantBefore = redundantBefore;
this.maxDecidedRX = maxDecidedRX;
@@ -291,7 +289,7 @@ public interface CommandSummaries
this.testKind = testKind;
this.minTxnId = minTxnId;
this.primaryExecuteAt = primaryExecuteAt;
- this.loadKeysFor = loadKeysFor;
+ this.findKeys = findKeys;
this.decidedRx = forDeps(maxDecidedRX, searchFor, primaryTxnId);
}
@@ -300,9 +298,9 @@ public interface CommandSummaries
return searchFor;
}
- public LoadKeysFor loadKeysFor()
+ public FindKeys loadKeysFor()
{
- return loadKeysFor;
+ return findKeys;
}
public TxnId primaryTxnId()
@@ -393,7 +391,7 @@ public interface CommandSummaries
private Relevance relevanceInternal(TxnId txnId, @Nullable SaveStatus
saveStatus, @Nullable Durability durability, @Nullable Timestamp executeAt,
@Nullable Unseekables<?> participants, boolean permitNull)
{
- if (loadKeysFor == WRITE || !txnId.is(testKind) || (saveStatus !=
null && saveStatus.compareTo(SaveStatus.TruncatedUnapplied) >= 0))
+ if (findKeys == DECLARED || !txnId.is(testKind) || (saveStatus !=
null && saveStatus.compareTo(SaveStatus.TruncatedUnapplied) >= 0))
return IRRELEVANT;
if (!permitNull)
@@ -414,7 +412,7 @@ public interface CommandSummaries
return IRRELEVANT;
Relevance atLeast = IRRELEVANT;
- if (loadKeysFor == RECOVERY && txnId.witnesses(primaryTxnId))
+ if (findKeys == FindKeys.SUPERSEDING &&
txnId.witnesses(primaryTxnId))
{
if (isIgnorableFutureRx(txnId, participants))
return IRRELEVANT;
@@ -531,7 +529,7 @@ public interface CommandSummaries
public final <P> Summary get(Relevance relevance, TxnId txnId,
Timestamp executeAt, SaveStatus saveStatus, Durability durability,
Participants<?> touches, @Nullable P deps, TriPredicate<P, TxnId,
Unseekables<?>> depTester)
{
touches = touches.intersecting(searchFor, Minimal);
- IsDep isDep = loadKeysFor == RECOVERY ? IsDep.NOT_ELIGIBLE : null;
+ IsDep isDep = findKeys == FindKeys.SUPERSEDING ?
IsDep.NOT_ELIGIBLE : null;
switch (relevance)
{
default: throw new UnhandledEnum(relevance);
@@ -569,15 +567,6 @@ public interface CommandSummaries
}
}
- enum ComputeIsDep
- {
- // don't test deps
- IGNORE,
-
- // calculate but don't filter
- EITHER
- }
-
interface ActiveCommandVisitor<P1, P2>
{
void visit(P1 p1, P2 p2, SummaryStatus status, Durability durability,
Unseekable keyOrRange, TxnId txnId);
@@ -592,20 +581,20 @@ public interface CommandSummaries
boolean visit(Unseekable keyOrRange, TxnId txnId, Timestamp executeAt,
SummaryStatus status, @Nullable IsDep dep, Durability minDurability);
}
- boolean visit(Unseekables<?> keysOrRanges, TxnId testTxnId, Kinds
testKind, SupersedingCommandVisitor visit);
+ boolean visit(@Nullable Unseekables<?> keysOrRanges, TxnId testTxnId,
Kinds testKind, SupersedingCommandVisitor visit);
/**
* Visits keys first in ascending order, with equal keys visiting TxnId is
ascending order.
* Visits range transactions in ascending order by TxnId, then visiting
each Range in ascending order
*/
- <P1, P2> void visit(Unseekables<?> keysOrRanges, Timestamp startedBefore,
Kinds testKind, ActiveCommandVisitor<P1, P2> visit, P1 p1, P2 p2);
+ <P1, P2> void visit(@Nullable Unseekables<?> keysOrRanges, Timestamp
startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visit, P1 p1, P2
p2);
// TODO (expected): ByRangeSnapshot based on IntervalBTree so we can elide
superseded dependencies as we do with CommandsForKey
interface ByTxnIdSnapshot extends CommandSummaries
{
NavigableMap<Timestamp, Summary> byTxnId();
- default boolean visit(Unseekables<?> keysOrRanges,
+ default boolean visit(@Nullable Unseekables<?> keysOrRanges,
TxnId testTxnId,
Kinds testKind,
SupersedingCommandVisitor visit)
@@ -620,15 +609,17 @@ public interface CommandSummaries
if (!value.is(MAYBE_SUPERSEDING))
continue;
- Unseekables<?> participants = value.participants;
- Unseekables<?> intersecting =
participants.overlapping(keysOrRanges);
- if (!intersecting.isEmpty())
+ Unseekables<?> overlapping = value.participants;
+ if (keysOrRanges != null)
+ overlapping = overlapping.overlapping(keysOrRanges);
+
+ if (!overlapping.isEmpty())
{
Timestamp executeAt = value.plainExecuteAt();
Invariants.require(executeAt != null);
SummaryStatus status = value.status();
IsDep dep = value.isDep();
- for (Unseekable participant : intersecting)
+ for (Unseekable participant : overlapping)
{
if (!visit.visit(participant, value.plainTxnId(),
executeAt, status, dep, NotDurable))
return false;
@@ -640,7 +631,7 @@ public interface CommandSummaries
}
@Override
- default <P1, P2> void visit(Unseekables<?> keysOrRanges, Timestamp
startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visit, P1 p1, P2 p2)
+ default <P1, P2> void visit(@Nullable Unseekables<?> keysOrRanges,
Timestamp startedBefore, Kinds testKind, ActiveCommandVisitor<P1, P2> visit, P1
p1, P2 p2)
{
NavigableMap<Timestamp, Summary> map = byTxnId();
for (Summary value : map.headMap(startedBefore, false).values())
@@ -654,7 +645,11 @@ public interface CommandSummaries
if (!value.is(ACTIVE))
continue;
- for (Unseekable keyOrRange :
value.participants.intersecting(keysOrRanges, Minimal))
+ Unseekables<?> participants = value.participants;
+ if (keysOrRanges != null)
+ participants = participants.intersecting(keysOrRanges,
Minimal);
+
+ for (Unseekable keyOrRange : participants)
visit.visit(p1, p2, value.status(), value.durability(),
keyOrRange, value.plainTxnId());
}
}
diff --git a/accord-core/src/main/java/accord/local/DepsCalculator.java
b/accord-core/src/main/java/accord/local/DepsCalculator.java
index 5cfc4db1..6391854e 100644
--- a/accord-core/src/main/java/accord/local/DepsCalculator.java
+++ b/accord-core/src/main/java/accord/local/DepsCalculator.java
@@ -18,11 +18,16 @@
package accord.local;
+import java.util.function.BiConsumer;
+import java.util.function.Consumer;
+import java.util.function.Function;
import javax.annotation.Nullable;
import accord.coordinate.ExecuteFlag.ExecuteFlags;
import accord.local.CommandSummaries.SummaryStatus;
import accord.local.MaxDecidedRX.DecidedRX;
+import accord.messages.Reply;
+import accord.messages.ReplyList;
import accord.primitives.Deps;
import accord.primitives.EpochSupplier;
import accord.primitives.Participants;
@@ -33,14 +38,25 @@ import accord.primitives.TxnId;
import accord.primitives.Unseekable;
import accord.primitives.Unseekables;
import accord.utils.Invariants;
+import accord.utils.UnhandledEnum;
+import accord.utils.async.AsyncChain;
+import accord.utils.async.AsyncResult;
+import accord.utils.async.AsyncResults;
+import accord.utils.async.AsyncResults.AbstractImmediate;
+import accord.utils.async.CancellableAsyncResult;
import static accord.coordinate.ExecuteFlag.HAS_UNIQUE_HLC;
import static accord.coordinate.ExecuteFlag.READY_TO_EXECUTE;
import static accord.local.CommandSummaries.SummaryStatus.APPLIED;
+import static accord.local.DepsCalculator.Initialised.DONE;
+import static accord.local.DepsCalculator.Initialised.INCOMPLETE;
+import static accord.local.DepsCalculator.Initialised.REJECTED;
+import static accord.local.LoadKeys.INCR;
import static accord.primitives.Txn.Kind.EphemeralRead;
import static accord.primitives.Txn.Kind.ExclusiveSyncPoint;
+import static accord.utils.Invariants.illegalState;
-public class DepsCalculator extends Deps.Builder implements
CommandSummaries.ActiveCommandVisitor<TxnId,
DepsCalculator.MinDependencyCalculator>
+public abstract class DepsCalculator extends Deps.Builder implements
CommandSummaries.ActiveCommandVisitor<TxnId,
DepsCalculator.MinDependencyCalculator>, Consumer<SafeCommandStore>,
ExecutionContext
{
public static class MinDependencyCalculator
{
@@ -77,32 +93,182 @@ public class DepsCalculator extends Deps.Builder
implements CommandSummaries.Act
}
}
+ public static final class SynchronousDepsCalculator extends DepsCalculator
+ {
+ public SynchronousDepsCalculator(TxnId txnId, Timestamp executeAt,
Participants<?> touches)
+ {
+ super(txnId, executeAt, touches);
+ }
+
+ public Deps calculate(SafeCommandStore safeStore, long minEpoch,
boolean rejectIfRedundant)
+ {
+ Initialised initialised = initialise(safeStore, minEpoch,
rejectIfRedundant);
+ switch (initialised)
+ {
+ default: throw new UnhandledEnum(initialised);
+ case INCOMPLETE: throw illegalState("Could not calculate
dependencies - must have declared insufficient ExecutionContext (%s vs %s)",
safeStore.context(), this);
+ case REJECTED: return null;
+ case DONE: return deps();
+ }
+ }
+
+ public static Deps calculateDeps(SafeCommandStore safeStore, TxnId
txnId, StoreParticipants participants, long minEpoch, Timestamp executeAt,
boolean rejectIfRedundant)
+ {
+ return calculateDeps(safeStore, txnId, participants.touches(),
minEpoch, executeAt, rejectIfRedundant);
+ }
+
+ public static Deps calculateDeps(SafeCommandStore safeStore, TxnId
txnId, Participants<?> touches, long minEpoch, Timestamp executeAt, boolean
rejectIfRedundant)
+ {
+ try (SynchronousDepsCalculator calculator = new
SynchronousDepsCalculator(txnId, executeAt, touches))
+ {
+ return calculator.calculate(safeStore, minEpoch,
rejectIfRedundant);
+ }
+ }
+ }
+
+ public abstract static class AbstractDepsReply<R extends
AbstractDepsReply<R>> extends AbstractImmediate<R> implements ReplyList<R>,
Reply, CancellableAsyncResult<R>
+ {
+ @Override
+ public boolean isSuccess()
+ {
+ return true;
+ }
+
+ @Override
+ public AsyncResult<R> invoke(BiConsumer<? super R, Throwable> callback)
+ {
+ callback.accept((R)this, null);
+ return this;
+ }
+
+ @Override
+ public int size()
+ {
+ return 1;
+ }
+
+ @Override
+ public CancellableAsyncResult<R> get(int i)
+ {
+ Invariants.requireArgument(i == 0);
+ return this;
+ }
+
+ @Override
+ public void cancel()
+ {
+ }
+
+ @Override
+ public void cancelReplies()
+ {
+ }
+ }
+
+ static class AsyncDepsReply<R extends AbstractDepsReply<R>> extends
AsyncResults.CancellableChain<R> implements ReplyList<R>
+ {
+ public AsyncDepsReply(AsyncChain<R> chain)
+ {
+ super(chain);
+ }
+
+ @Override
+ public int size()
+ {
+ return 1;
+ }
+
+ @Override
+ public CancellableAsyncResult<R> get(int i)
+ {
+ Invariants.require(i == 0);
+ return this;
+ }
+
+ @Override
+ public void cancelReplies()
+ {
+ cancel();
+ }
+ }
+
+ public abstract static class DepsReplyCalculator<R extends
AbstractDepsReply<R>> extends DepsCalculator implements Function<Void, R>
+ {
+ public DepsReplyCalculator(TxnId txnId, Timestamp executeAt,
StoreParticipants participants)
+ {
+ super(txnId, executeAt, participants);
+ }
+
+ public DepsReplyCalculator(TxnId txnId, Timestamp executeAt,
Participants<?> touches)
+ {
+ super(txnId, executeAt, touches);
+ }
+
+ public ReplyList<R> calculate(SafeCommandStore safeStore, long
minEpoch, boolean nullIfRedundant)
+ {
+ Initialised initialised;
+ try
+ {
+ initialised = initialise(safeStore, minEpoch, nullIfRedundant);
+ }
+ catch (Throwable t)
+ {
+ close();
+ throw t;
+ }
+
+ switch (initialised)
+ {
+ default: throw new UnhandledEnum(initialised);
+ case REJECTED:
+ return null;
+ case DONE:
+ return apply(null);
+ case INCOMPLETE:
+ AsyncChain<R> chain =
safeStore.commandStore().continuationChain(this, this).map(this);
+ return new AsyncDepsReply<>(chain);
+ }
+ }
+ }
+
// TODO (expected): we can also track whether we have only single-key
writes that have been Accepted with ballot 0 (or timestamp != t0), or else
Committed[1];
// in this case we can decide immediately if we have a unique hlc as we
don't run the risk of other keys inserting some arbitrary timestamp
// [1] probably unsafe to use Accepted with ballot > 0, as there could be
a timestamp battle, and the timestamp we see might not be the one that gets
decided.
- private final long now;
+ protected final TxnId txnId;
+ protected final Timestamp executeAt;
+ private final Participants<?> touches;
+ private RangeDeps redundant;
private long sumUnappliedAge, maxUnappliedAge;
private int unappliedCount;
private long maxAppliedHlc;
+ private MinDependencyCalculator minDepCalc;
- public DepsCalculator(Timestamp timestamp)
+ public DepsCalculator(TxnId txnId, Timestamp executeAt, StoreParticipants
touches)
+ {
+ this(txnId, executeAt, touches.touches());
+ }
+
+ public DepsCalculator(TxnId txnId, Timestamp executeAt, Participants<?>
touches)
{
super(true);
- this.now = timestamp.hlc();
+ this.txnId = txnId;
+ this.touches = touches;
+ this.executeAt = executeAt.equals(txnId) ? txnId : executeAt;
}
@Override
- public void visit(TxnId self, @Nullable MinDependencyCalculator
minDepCalc, SummaryStatus status, Durability durability, Unseekable keyOrRange,
TxnId depId)
+ public final void visit(TxnId self, @Nullable MinDependencyCalculator
minDepCalc, SummaryStatus status, Durability durability, Unseekable keyOrRange,
TxnId depId)
{
if (minDepCalc != null && !minDepCalc.include(durability, keyOrRange,
depId))
return;
if (self == null || !self.equals(depId))
add(keyOrRange, depId);
+
if (status.compareTo(APPLIED) < 0)
{
unappliedCount += 1;
- long age = Math.max(0, now - depId.hlc());
+ long age = Math.max(0, executeAt.hlc() - depId.hlc());
sumUnappliedAge += age;
if (age > maxUnappliedAge)
maxUnappliedAge = age;
@@ -110,13 +276,13 @@ public class DepsCalculator extends Deps.Builder
implements CommandSummaries.Act
}
@Override
- public void visitMaxAppliedHlc(long maxAppliedHlc)
+ public final void visitMaxAppliedHlc(long maxAppliedHlc)
{
if (maxAppliedHlc > this.maxAppliedHlc)
this.maxAppliedHlc = maxAppliedHlc;
}
- public ExecuteFlags executeFlags(TxnId txnId)
+ public final ExecuteFlags executeFlags()
{
ExecuteFlags flags = ExecuteFlags.none();
if (unappliedCount == 0)
@@ -129,59 +295,110 @@ public class DepsCalculator extends Deps.Builder
implements CommandSummaries.Act
return flags;
}
- public Timestamp executeAt(SafeCommand safeCommand, Node node)
+ public final Deps deps()
{
- Timestamp executeAt = safeCommand.current().executeAtOrTxnId();
+ Deps result = super.build();
+ result = new Deps(result.keyDeps, result.rangeDeps.with(redundant));
+ Invariants.require(!txnId.isVisible() || !result.contains(txnId));
+ return result;
+ }
+
+ public final Timestamp executeAt(Timestamp witnessedAt, Node node)
+ {
+ Timestamp executeAt = witnessedAt;
if (unappliedCount > 0 && node.agent().softReject(unappliedCount,
maxUnappliedAge, sumUnappliedAge))
executeAt = executeAt.addFlag(Timestamp.Flag.SOFT_REJECT);
return executeAt;
}
- public Deps calculate(SafeCommandStore safeStore, TxnId txnId,
StoreParticipants participants, long minEpoch, Timestamp executeAt, boolean
nullIfRedundant)
+ public enum Initialised
{
- return calculate(safeStore, txnId, participants.touches(), minEpoch,
executeAt, nullIfRedundant);
+ REJECTED, DONE, INCOMPLETE
}
- public Deps calculate(SafeCommandStore safeStore, TxnId txnId,
Participants<?> touches, long minEpoch, Timestamp executeAt, boolean
nullIfRedundant)
+ public final Initialised initialise(SafeCommandStore safeStore, long
minEpoch, boolean rejectIfRedundant)
{
- RangeDeps redundant;
try (RangeDeps.BuilderByRange redundantBuilder =
RangeDeps.builderByRange())
{
redundant = safeStore.redundantBefore().collectDeps(touches,
redundantBuilder, EpochSupplier.constant(minEpoch), executeAt)
.build();
}
- if (nullIfRedundant && !txnId.is(EphemeralRead))
+ if (rejectIfRedundant && !txnId.is(EphemeralRead))
{
TxnId maxRedundantBefore = redundant.maxTxnId(null);
if (maxRedundantBefore != null &&
maxRedundantBefore.compareTo(executeAt) >= 0)
{
Invariants.require(maxRedundantBefore.isSyncPoint());
- return null;
+ return REJECTED;
}
}
- // NOTE: ExclusiveSyncPoint *relies* on STARTED_BEFORE to ensure it
reports a dependency on *every* earlier TxnId that may execute (before or after
it).
- MinDependencyCalculator minDepCalc = null;
// the main difference between RX and RV is whether we apply this
filtering
- if (txnId.is(ExclusiveSyncPoint)) minDepCalc = new
MinDependencyCalculator(safeStore.maxDecidedRX(), touches, txnId);
- safeStore.visit(touches, executeAt, txnId.witnesses(), this,
executeAt.equals(txnId) ? null : txnId, minDepCalc);
- Deps result = super.build();
- result = new Deps(result.keyDeps, result.rangeDeps.with(redundant));
- Invariants.require(!txnId.isVisible() || !result.contains(txnId));
- return result;
- }
+ if (txnId.is(ExclusiveSyncPoint))
+ minDepCalc = new MinDependencyCalculator(safeStore.maxDecidedRX(),
touches, txnId);
- public static Deps calculateDeps(SafeCommandStore safeStore, TxnId txnId,
StoreParticipants participants, long minEpoch, Timestamp executeAt, boolean
nullIfRedundant)
- {
- return calculateDeps(safeStore, txnId, participants.touches(),
minEpoch, executeAt, nullIfRedundant);
+ safeStore.visit(touches, executeAt, txnId.witnesses(), this, executeAt
== txnId ? null : txnId, minDepCalc);
+ if (safeStore.context().keys().containsAll(touches))
+ return DONE;
+
+ return INCOMPLETE;
}
- public static Deps calculateDeps(SafeCommandStore safeStore, TxnId txnId,
Participants<?> touches, long minEpoch, Timestamp executeAt, boolean
nullIfRedundant)
+ @Override
+ public void accept(SafeCommandStore safeStore)
{
- try (DepsCalculator calculator = new DepsCalculator(executeAt))
+ try
{
- return calculator.calculate(safeStore, txnId, touches, minEpoch,
executeAt, nullIfRedundant);
+ safeStore.visit(null, executeAt, txnId.witnesses(), this,
executeAt == txnId ? null : txnId, minDepCalc);
}
+ catch (Throwable t)
+ {
+ try { close(); }
+ catch (Throwable t2) { try { t.addSuppressed(t2); } catch
(Throwable ignore) {} }
+ throw t;
+ }
+ }
+
+ @Override
+ public TxnId primaryTxnId()
+ {
+ return txnId;
+ }
+
+ @Override
+ public String reason()
+ {
+ return "Calculate Deps";
+ }
+
+ @Override
+ public boolean abandonPartialSuccess()
+ {
+ return true;
+ }
+
+ @Override
+ public Unseekables<?> keys()
+ {
+ return touches;
+ }
+
+ @Override
+ public LoadKeys loadKeys()
+ {
+ return INCR;
+ }
+
+ @Override
+ public FindKeys findKeys()
+ {
+ return FindKeys.CONFLICTS;
+ }
+
+ @Override
+ public ExecutionSequence executionSequence()
+ {
+ return ExecutionSequence.ATOMIC;
}
}
diff --git a/accord-core/src/main/java/accord/local/ExecutionContext.java
b/accord-core/src/main/java/accord/local/ExecutionContext.java
index 54a500ff..4c51968b 100644
--- a/accord-core/src/main/java/accord/local/ExecutionContext.java
+++ b/accord-core/src/main/java/accord/local/ExecutionContext.java
@@ -18,33 +18,31 @@
package accord.local;
-import accord.local.cfk.CommandsForKey;
+import java.util.AbstractList;
+import java.util.List;
+import java.util.function.Consumer;
+
+import javax.annotation.Nullable;
+
+import net.nicoulaj.compilecommand.annotations.Inline;
+
import accord.primitives.Ranges;
import accord.primitives.Routables.Slice;
import accord.primitives.RoutingKeys;
import accord.primitives.Timestamp;
import accord.primitives.TxnId;
-
import accord.primitives.Unseekables;
import accord.utils.Invariants;
-import net.nicoulaj.compilecommand.annotations.Inline;
-
-import java.util.AbstractList;
-import java.util.List;
-import java.util.function.Consumer;
-import javax.annotation.Nullable;
import static accord.local.LoadKeys.INCR;
import static accord.local.LoadKeys.NONE;
import static accord.local.LoadKeys.SYNC;
-import static accord.local.LoadKeysFor.READ_WRITE;
-import static accord.local.LoadKeysFor.WRITE;
+import static accord.local.FindKeys.CONFLICTS;
+import static accord.local.FindKeys.DECLARED;
/**
- * Lists txnids and keys of commands and commands for key that will be needed
for an operation. Used
- * to ensure the necessary state is in memory for an operation before it
executes.
- *
- * TODO (desired): rename to simply Context, or LoadContext
+ * Tasks declare required data and semantics for their execution.
+ * An INCR or ASYNC execution
*/
public interface ExecutionContext
{
@@ -144,13 +142,13 @@ public interface ExecutionContext
}
/**
- * @return keys of the {@link CommandsForKey} objects that need to be
loaded into memory before this operation is run
+ * @return keys or ranges that this key needs CommandSummaries loaded for
*/
default Unseekables<?> keys() { return RoutingKeys.EMPTY; }
default LoadKeys loadKeys() { return NONE; }
- default LoadKeysFor loadKeysFor() { return WRITE; }
+ default FindKeys findKeys() { return DECLARED; }
/**
* Whether this execution may be retried safely; useful only for INCR
tasks that may partially succeed,
@@ -158,6 +156,11 @@ public interface ExecutionContext
*/
default boolean isIdempotent() { return false; }
+ /**
+ * Whether this execution should be retried if partially executes; useful
only for INCR tasks that may partially succeed.
+ */
+ default boolean abandonPartialSuccess() { return false; }
+
default ExecutionKind executionKind() { return ExecutionKind.OTHER; }
default ExecutionSequence executionSequence() { return
ExecutionSequence.BY_PRIORITY; }
@@ -189,7 +192,7 @@ public interface ExecutionContext
requiredHistory = SYNC;
if (requiredHistory.compareTo(superset.loadKeys()) < 0)
return false;
- if (loadKeysFor().compareTo(superset.loadKeysFor()) > 0)
+ if (findKeys().compareTo(superset.findKeys()) > 0)
return false;
}
@@ -229,8 +232,9 @@ public interface ExecutionContext
@Override default ExecutionSequence executionSequence() { return
wrapped().executionSequence(); }
@Override default ExecutionKind executionKind() { return
wrapped().executionKind(); }
@Override default boolean isIdempotent() { return
wrapped().isIdempotent(); }
+ @Override default boolean abandonPartialSuccess() { return
wrapped().abandonPartialSuccess(); }
@Override default LoadKeys loadKeys() { return wrapped().loadKeys(); }
- @Override default LoadKeysFor loadKeysFor() { return
wrapped().loadKeysFor(); }
+ @Override default FindKeys findKeys() { return wrapped().findKeys(); }
@Override default Timestamp executeAt() { return
wrapped().executeAt(); }
@Override default String reason() { return wrapped().reason(); }
@Override default String describe() { return wrapped().describe(); }
@@ -251,7 +255,7 @@ public interface ExecutionContext
@Override public ExecutionContext wrapped() { return wrapped; }
}
- static ExecutionContext contextFor(@Nullable TxnId primary, @Nullable
TxnId additional, Unseekables<?> keys, LoadKeys loadKeys, LoadKeysFor
loadKeysFor, String reason)
+ static ExecutionContext contextFor(@Nullable TxnId primary, @Nullable
TxnId additional, Unseekables<?> keys, LoadKeys loadKeys, FindKeys findKeys,
String reason)
{
Invariants.require(primary == null ? additional == null :
!primary.equals(additional));
return new ExecutionContext()
@@ -260,7 +264,7 @@ public interface ExecutionContext
@Override public @Nullable TxnId additionalTxnId() { return
additional; }
@Override public Unseekables<?> keys() { return keys; }
@Override public LoadKeys loadKeys() { return loadKeys; }
- @Override public LoadKeysFor loadKeysFor() { return loadKeysFor; }
+ @Override public FindKeys findKeys() { return findKeys; }
@Override public String reason() { return reason; }
@Override public String toString() { return describe(); }
};
@@ -328,7 +332,7 @@ public interface ExecutionContext
@Override public @Nullable TxnId primaryTxnId() { return txnId; }
@Override public Unseekables<?> keys() { return keys; }
@Override public LoadKeys loadKeys() { return SYNC; }
- @Override public LoadKeysFor loadKeysFor() { return READ_WRITE; }
+ @Override public FindKeys findKeys() { return CONFLICTS; }
@Override public ExecutionSequence executionSequence() { return
ExecutionSequence.UNSEQUENCED; }
@Override public String reason() { return reason; }
@Override public String toString() { return describe(); }
@@ -342,7 +346,7 @@ public interface ExecutionContext
@Override public @Nullable TxnId primaryTxnId() { return null; }
@Override public Unseekables<?> keys() { return keys; }
@Override public LoadKeys loadKeys() { return SYNC; }
- @Override public LoadKeysFor loadKeysFor() { return READ_WRITE; }
+ @Override public FindKeys findKeys() { return CONFLICTS; }
@Override public ExecutionSequence executionSequence() { return
ExecutionSequence.UNSEQUENCED; }
@Override public String reason() { return reason; }
@Override public String toString() { return describe(); }
diff --git a/accord-core/src/main/java/accord/local/LoadKeysFor.java
b/accord-core/src/main/java/accord/local/FindKeys.java
similarity index 78%
rename from accord-core/src/main/java/accord/local/LoadKeysFor.java
rename to accord-core/src/main/java/accord/local/FindKeys.java
index dcf14f12..6ee0454c 100644
--- a/accord-core/src/main/java/accord/local/LoadKeysFor.java
+++ b/accord-core/src/main/java/accord/local/FindKeys.java
@@ -20,29 +20,29 @@ package accord.local;
/**
* For operations that need information associated with keys or ranges,
- * this indicates whether the data will be queried, updated, or both.
+ * this indicates what data will be queried.
*/
-public enum LoadKeysFor
+public enum FindKeys
{
/**
- * WRITE covers only updating the relevant key state for the primaryTxnId,
+ * covers only updating the relevant key state for the primaryTxnId,
* that is for key transactions this means updating key summaries, and for
* range transactions this means updating any range summaries.
* Importantly, this does not mean range transactions must be able to
* synchronously (or otherwise) write to all intersecting key summaries.
*/
- WRITE,
+ DECLARED,
/**
- * READ covers all intersecting summaries of relevant key or range
transactions,
+ * CONFLICTS covers all intersecting summaries of relevant key or range
transactions,
* including any commands that should be witnessed by primaryTxnId.
* This means range transactions MUST be able to consult all intersecting
key summaries.
*/
- READ_WRITE,
+ CONFLICTS,
/**
- * RECOVERY is READ_WRITE + summary information of keys/transactions that
should have witnessed
+ * SUPERSEDING is CONFLICTS + summary information of keys/transactions
that should have witnessed
* the primaryTxnId.
*/
- RECOVERY
+ SUPERSEDING
}
diff --git a/accord-core/src/main/java/accord/local/LoadKeys.java
b/accord-core/src/main/java/accord/local/LoadKeys.java
index 8de1d782..48549df3 100644
--- a/accord-core/src/main/java/accord/local/LoadKeys.java
+++ b/accord-core/src/main/java/accord/local/LoadKeys.java
@@ -35,10 +35,15 @@ public enum LoadKeys
ASYNC,
/**
- * Load and process the requested keys incrementally; the operation will
be invoked multiples times
- * as keys are loaded, until all the keys have been processed. If
submitted by an already running execution
- * this task must declare a subset of the keys and txnIds declared by the
originating task.
- * It is not permitted to chain INCR tasks together; INCR may only be
submitted by an ASYNC or SYNC task.
+ * Important Notes:
+ * 1) An INCR task only adopts command summaries (keys or ranges) that
were not loaded by its parent task, so the
+ * parent task MUST process any keys that are available to it.
+ * 2) {@link SafeCommandStore#context()} will report only the KEYS that
have been loaded for a run, but if there
+ * are range summaries to process then the first invocation will
include these as well. To visit all loaded
+ * summaries ensure to pass {@code null} to {@link
SafeCommandStore#visit}.
+ *<p>
+ * Load and process the requested key and range command summaries
incrementally; the operation will be invoked
+ * multiples times as keys are loaded, until all the keys have been
processed.
*/
INCR,
@@ -70,5 +75,4 @@ public enum LoadKeys
return this.compareTo(ifSyncRequireAtLeast) >= 0;
}
}
-
}
diff --git a/accord-core/src/main/java/accord/local/MapReduceCommandStores.java
b/accord-core/src/main/java/accord/local/MapReduceCommandStores.java
index 008cd02e..5ca2ee11 100644
--- a/accord-core/src/main/java/accord/local/MapReduceCommandStores.java
+++ b/accord-core/src/main/java/accord/local/MapReduceCommandStores.java
@@ -60,6 +60,8 @@ public abstract class MapReduceCommandStores<P extends
Participants<?>, O> imple
protected AsyncChain<O> applyAsyncInternal(Ranges ranges, CommandStore
commandStore)
{
+ // TODO (desired): shouldn't need to override the context to supply
the ranges we interact with
+ // (perhaps accept Ranges as another method signature, or else leave
to commandStore to figure out)
return commandStore.chain(slice(ranges, Minimal), this);
}
diff --git a/accord-core/src/main/java/accord/local/SafeCommandStore.java
b/accord-core/src/main/java/accord/local/SafeCommandStore.java
index 0a9c2051..85d268f1 100644
--- a/accord-core/src/main/java/accord/local/SafeCommandStore.java
+++ b/accord-core/src/main/java/accord/local/SafeCommandStore.java
@@ -22,6 +22,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.NavigableMap;
import java.util.function.Consumer;
+
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
@@ -53,8 +54,8 @@ import accord.primitives.Timestamp;
import accord.primitives.TxnId;
import accord.primitives.Unseekables;
import accord.utils.Invariants;
-import accord.utils.Reduce;
import accord.utils.LargeBitSet;
+import accord.utils.Reduce;
import accord.utils.SortedList;
import accord.utils.async.AsyncChain;
import accord.utils.async.AsyncChains;
@@ -63,9 +64,9 @@ import static
accord.local.ExecutionContext.unsequencedIdempotentIncrementalWrit
import static accord.local.LoadKeys.INCR;
import static accord.local.LoadKeys.NONE;
import static accord.local.RedundantStatus.Property.LOCALLY_APPLIED;
-import static accord.local.RedundantStatus.SomeStatus.LOCALLY_WITNESSED_ONLY;
import static accord.local.RedundantStatus.Property.LOCALLY_REDUNDANT;
import static accord.local.RedundantStatus.Property.SHARD_APPLIED;
+import static accord.local.RedundantStatus.SomeStatus.LOCALLY_WITNESSED_ONLY;
import static accord.local.cfk.UpdateUnmanagedMode.REGISTER;
import static accord.primitives.Known.KnownRoute.MaybeRoute;
import static accord.primitives.Routable.Domain.Range;
@@ -293,7 +294,7 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
public final boolean canExecuteWith(ExecutionContext context) { return
canExecute(context) == context; }
/**
- * Attempt to ready the provided PreLoadContext; if this can only be
achieved partially, a new PreLoadContext
+ * Attempt to ready the provided PreLoadContext; if this can only be
achieved partially, a new ExecutionContext
* will be returned containing the readily available data. If nothing is
available, null will be returned.
*/
public abstract @Nullable ExecutionContext canExecute(ExecutionContext
context);
@@ -413,7 +414,7 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
return;
// TODO (expected): we don't want to insert any dependencies for those
we only touch; we just need to record them as decided/applied for execution
- ExecutionContext context = new UpdateManagedContext(next.txnId(),
update);
+ ExecutionContext context = new UpdateManagedContext(next.txnId,
update);
ExecutionContext execute = safeStore.canExecute(context);
if (execute != null)
{
@@ -431,7 +432,7 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
Unseekables<?> asyncKeys =
remainingKeys.without(participants.touches()).intersecting(participants.hasTouched(),
Minimal);
if (!asyncKeys.isEmpty())
{
- ExecutionContext async =
unsequencedIdempotentIncrementalWrite(asyncKeys, "Update CommandsForKey");
+ ExecutionContext async = new
UpdateManagedContext(next.txnId, asyncKeys);
updateManagedCommandsForKeyIncremental(async,
safeStore.commandStore(), forceNotify);
remainingKeys = remainingKeys.without(asyncKeys);
}
@@ -528,6 +529,8 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
}
else
{
+ // TODO (expected): no need to filter keys, as INCR tasks should
subtract parent keys
+ // (but should first synchronise executor behaviour between
accord/C*)
if (execute != null)
context = new UpdateUnmanagedContext(txnId,
keys.without(execute.keys()));
@@ -540,7 +543,7 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
}
}
- static class UpdateManagedContext implements ExecutionContext
+ static final class UpdateManagedContext implements ExecutionContext
{
final TxnId primaryTxnId;
final Unseekables<?> keys;
@@ -560,7 +563,7 @@ public abstract class SafeCommandStore implements
RangesForEpochSupplier, Redund
@Override public String toString() { return describe(); }
}
- static class UpdateUnmanagedContext implements ExecutionContext
+ static final class UpdateUnmanagedContext implements ExecutionContext
{
final TxnId primaryTxnId;
final Unseekables<?> keys;
diff --git
a/accord-core/src/main/java/accord/local/durability/DurabilityLevel.java
b/accord-core/src/main/java/accord/local/durability/DurabilityLevel.java
index 8e79f29b..051ec0ed 100644
--- a/accord-core/src/main/java/accord/local/durability/DurabilityLevel.java
+++ b/accord-core/src/main/java/accord/local/durability/DurabilityLevel.java
@@ -79,8 +79,8 @@ public class DurabilityLevel
{
SyncLocal local = min(a.local, b.local);
SyncRemote remote = min(a.remote, b.remote);
- SortedArrayList<Node.Id> including = union(a.including, b.including);
SortedArrayList<Node.Id> excluding = union(a.excluding, b.excluding);
+ SortedArrayList<Node.Id> including = subtract(union(a.including,
b.including), excluding);
if (including != null && excluding != null)
including = including.without(excluding);
return new DurabilityLevel(local, remote, including, excluding);
@@ -91,7 +91,7 @@ public class DurabilityLevel
SyncLocal local = max(a.local, b.local);
SyncRemote remote = max(a.remote, b.remote);
SortedArrayList<Node.Id> including = union(a.including, b.including);
- SortedArrayList<Node.Id> excluding = subtract(a.excluding,
b.excluding);
+ SortedArrayList<Node.Id> excluding = subtract(union(a.excluding,
b.excluding), including);
if (including != null && excluding != null)
including = including.without(excluding);
return new DurabilityLevel(local, remote, including, excluding);
diff --git a/accord-core/src/main/java/accord/messages/Accept.java
b/accord-core/src/main/java/accord/messages/Accept.java
index 85186062..fff16fed 100644
--- a/accord-core/src/main/java/accord/messages/Accept.java
+++ b/accord-core/src/main/java/accord/messages/Accept.java
@@ -18,6 +18,7 @@
package accord.messages;
+import java.util.function.Function;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
@@ -25,9 +26,10 @@ import accord.coordinate.ExecuteFlag.ExecuteFlags;
import accord.local.Command;
import accord.local.Commands;
import accord.local.Commands.AcceptOutcome;
-import accord.local.DepsCalculator;
+import accord.local.DepsCalculator.AbstractDepsReply;
+import accord.local.DepsCalculator.DepsReplyCalculator;
import accord.local.LoadKeys;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
import accord.local.Node.Id;
import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
@@ -49,11 +51,11 @@ import accord.utils.UnhandledEnum;
import accord.utils.async.Cancellable;
import static
accord.api.ProtocolModifiers.filterDuplicateDependenciesFromAcceptReply;
+import static accord.api.ProtocolModifiers.loadKeysAsyncIfPermitted;
import static
accord.api.ProtocolModifiers.syncPointsTrackUnstableMediumPathDependencies;
import static accord.local.Commands.AcceptOutcome.Redundant;
import static accord.local.Commands.AcceptOutcome.RejectedBallot;
import static accord.local.Commands.AcceptOutcome.Success;
-import static accord.local.LoadKeys.SYNC;
import static accord.messages.MessageType.StandardMessage.ACCEPT_REQ;
import static accord.messages.MessageType.StandardMessage.ACCEPT_RSP;
import static accord.messages.MessageType.StandardMessage.NOT_ACCEPT_REQ;
@@ -61,7 +63,7 @@ import static accord.primitives.Known.KnownDeps.DepsKnown;
// TODO (low priority, efficiency): use different objects for send and
receive, so can be more efficient
// (e.g. serialize without slicing, and
without unnecessary fields)
-public class Accept extends RouteRequest.WithUnsynced<Accept.AcceptReply>
+public class Accept extends
RouteRequest.WithUnsynced<ReplyList<Accept.AcceptReply>>
{
public static class SerializerSupport
{
@@ -140,15 +142,15 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
}
@Override
- public AcceptReply applyInternal(SafeCommandStore safeStore)
+ public ReplyList<AcceptReply> applyInternal(SafeCommandStore safeStore)
{
- PartialDeps partialDeps = this.partialDeps;
+ PartialDeps inputDeps = this.partialDeps;
if (ifDoneExpectCancelled()) // check cancellation after reading
nullable fields
return null; // we can't throw an exception here else we override
any non-exceptional reply informing the reason
StoreParticipants participants = StoreParticipants.update(safeStore,
scope, minEpoch, txnId, txnId.epoch(), executeAt.epoch());
SafeCommand safeCommand = safeStore.get(txnId, participants);
- AcceptOutcome outcome = Commands.accept(safeStore, safeCommand,
participants, txnId, kind, ballot, scope, executeAt, partialDeps);
+ AcceptOutcome outcome = Commands.accept(safeStore, safeCommand,
participants, txnId, kind, ballot, scope, executeAt, inputDeps);
switch (outcome)
{
default: throw new UnhandledEnum(outcome);
@@ -159,13 +161,14 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
boolean notOwner = participants.owns().isEmpty();
Participants<?> hasDeps = null;
- Deps deps = null;
+ Deps foundDeps;
if (command.known().is(DepsKnown) && (isPartialAccept() ||
notOwner))
{
- deps = command.partialDeps().asFullUnsafe();
+ foundDeps = command.partialDeps().asFullUnsafe();
hasDeps = command.participants().stillTouches();
}
+ else foundDeps = null;
Ballot superseding = command.promised();
if (superseding.compareTo(ballot) <= 0)
@@ -179,30 +182,36 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
calculateDeps = calculateDeps();
}
+ final Deps removeDeps = filterDeps() ? inputDeps : null;
+ final Timestamp executeAtIfKnown = command.executeAtIfKnown();
+
if (calculateDeps)
{
Participants<?> calculate = participants.touches();
if (hasDeps != null)
calculate = calculate.without(hasDeps);
+ hasDeps = participants.touches();
if (!calculate.isEmpty())
{
- Deps calculatedDeps =
DepsCalculator.calculateDeps(safeStore, txnId, calculate, minEpoch, executeAt,
true);
- if (calculatedDeps == null)
- return AcceptReply.inThePast(ballot, participants,
command);
-
- deps = deps == null ? calculatedDeps :
calculatedDeps.with(deps);
+ Participants<?> successful = replySuccessful(hasDeps);
+ outcome = replyOutcome(notOwner, outcome,
participants, hasDeps);
+ //noinspection resource
+ TruncatedAcceptDepsCalculator calculator = new
TruncatedAcceptDepsCalculator(txnId, executeAt, calculate, foundDeps,
removeDeps, outcome, superseding, successful, executeAtIfKnown);
+ ReplyList<AcceptReply> reply =
calculator.calculate(safeStore, minEpoch, true);
+ if (reply == null)
+ reply = AcceptReply.inThePast(ballot,
participants, command);
+ return reply;
}
- hasDeps = participants.touches();
}
- Participants<?> successful = isPartialAccept() ? hasDeps :
null;
- if (notOwner && (outcome == Redundant || (hasDeps != null &&
hasDeps.containsAll(participants.touches()))))
- outcome = Success;
+ Deps deps = foundDeps;
+ if (deps != null && removeDeps != null)
+ deps = deps.without(removeDeps);
- if (deps != null && filterDeps())
- deps = deps.without(partialDeps);
- return new AcceptReply(outcome, superseding, successful, deps,
command.executeAtIfKnown());
+ Participants<?> successful = replySuccessful(hasDeps);
+ outcome = replyOutcome(notOwner, outcome, participants,
hasDeps);
+ return new AcceptReply(outcome, superseding, successful, deps,
executeAtIfKnown);
}
case RejectedBallot:
@@ -212,33 +221,34 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
// if we're Retired, participants.owns() is empty, so we're
just fetching deps
// TODO (desired): optimise deps calculation; for some keys we
only need to return the last RX
case Success:
- ExecuteFlags flags;
- Deps deps;
- if (calculateDeps())
- {
- try (DepsCalculator calculator = new
DepsCalculator(executeAt))
- {
- deps = calculator.calculate(safeStore, txnId,
participants, minEpoch, executeAt, true);
- if (deps == null)
- return AcceptReply.inThePast(ballot, participants,
safeCommand.current());
- flags = calculator.executeFlags(txnId);
- }
+ {
+ Participants<?> successful = isPartialAccept() ?
participants.touches() : null;
- Invariants.require(deps.maxTxnId(txnId).epoch() <=
executeAt.epoch());
- if (filterDeps())
- deps = deps.without(partialDeps);
- }
- else
- {
- flags = ExecuteFlags.none();
- deps = Deps.NONE;
- }
+ if (!calculateDeps())
+ return new AcceptReply(successful, Deps.NONE,
ExecuteFlags.none());
- Participants<?> successful = isPartialAccept() ?
participants.touches() : null;
- return new AcceptReply(successful, deps, flags);
+ //noinspection resource
+ NormalAcceptDepsCalculator calculator = new
NormalAcceptDepsCalculator(txnId, executeAt, participants, filterDeps() ?
inputDeps : null, successful);
+ ReplyList<AcceptReply> reply = calculator.calculate(safeStore,
minEpoch, true);
+ if (reply == null)
+ reply = AcceptReply.inThePast(ballot, participants,
safeCommand.current());
+ return reply;
+ }
}
}
+ private Participants<?> replySuccessful(Participants<?> hasDeps)
+ {
+ return isPartialAccept() ? hasDeps : null;
+ }
+
+ private AcceptOutcome replyOutcome(boolean notOwner, AcceptOutcome
outcome, StoreParticipants participants, Participants<?> hasDeps)
+ {
+ if (notOwner && (outcome == Redundant || (hasDeps != null &&
hasDeps.containsAll(participants.touches()))))
+ return Success;
+ return outcome;
+ }
+
private boolean isPartialAccept()
{
return AcceptFlags.isPartial(acceptFlags);
@@ -255,9 +265,9 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
}
@Override
- public AcceptReply reduce(AcceptReply r1, AcceptReply r2)
+ public ReplyList<AcceptReply> reduce(ReplyList<AcceptReply> r1,
ReplyList<AcceptReply> r2)
{
- return AcceptReply.reduce(r1, r2);
+ return ReplyList.merge(r1, r2);
}
@Override
@@ -267,23 +277,24 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
}
@Override
- protected void acceptInternal(AcceptReply reply, Throwable failure)
+ protected void acceptInternal(ReplyList<AcceptReply> replies, Throwable
fail)
{
// finished processing, null out large objects
partialDeps = null;
- super.acceptInternal(reply, failure);
+ if (fail != null || replies == null) acceptReply(null, fail);
+ else ReplyList.invoke(replies, AcceptReply::reduce, this::acceptReply);
}
@Override
public LoadKeys loadKeys()
{
- return SYNC;
+ return loadKeysAsyncIfPermitted(txnId);
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
- return calculateDeps() ? LoadKeysFor.READ_WRITE : LoadKeysFor.WRITE;
+ return calculateDeps() ? FindKeys.CONFLICTS : FindKeys.DECLARED;
}
public ExecutionKind executionKind()
@@ -307,7 +318,79 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
'}';
}
- public static final class AcceptReply implements Reply
+ static class TruncatedAcceptDepsCalculator extends
DepsReplyCalculator<AcceptReply> implements Function<Void, AcceptReply>
+ {
+ final Deps foundDeps;
+ final @Nullable Deps removeDeps;
+ final AcceptOutcome outcome;
+ final Ballot superseding;
+ final Participants<?> successful;
+ final Timestamp executeAtIfKnown;
+
+ public TruncatedAcceptDepsCalculator(TxnId txnId, Timestamp executeAt,
Participants<?> touches, Deps foundDeps, @Nullable Deps removeDeps,
AcceptOutcome outcome, Ballot superseding, Participants<?> successful,
Timestamp executeAtIfKnown)
+ {
+ super(txnId, executeAt, touches);
+ this.foundDeps = foundDeps;
+ this.removeDeps = removeDeps;
+ this.outcome = outcome;
+ this.superseding = superseding;
+ this.successful = successful;
+ this.executeAtIfKnown = executeAtIfKnown;
+ }
+
+ public AcceptReply apply(Void ignore)
+ {
+ try
+ {
+ Deps deps = deps();
+ if (foundDeps != null)
+ deps = deps.with(foundDeps);
+
+ if (deps != null && removeDeps != null)
+ deps = deps.without(removeDeps);
+ return new AcceptReply(outcome, superseding, successful, deps,
executeAtIfKnown);
+ }
+ finally { close(); }
+ }
+ }
+
+ static class NormalAcceptDepsCalculator extends
DepsReplyCalculator<AcceptReply> implements Function<Void, AcceptReply>
+ {
+ final @Nullable Deps removeDeps;
+ final @Nullable Participants<?> successful;
+
+ public NormalAcceptDepsCalculator(TxnId txnId, Timestamp executeAt,
StoreParticipants participants, @Nullable Deps removeDeps, @Nullable
Participants<?> successful)
+ {
+ super(txnId, executeAt, participants);
+ this.removeDeps = removeDeps;
+ this.successful = successful;
+ }
+
+ public AcceptReply apply(Void ignore)
+ {
+ return finish();
+ }
+
+ AcceptReply finish()
+ {
+ try
+ {
+ Deps deps = deps();
+ Invariants.require(deps.maxTxnId(txnId).epoch() <=
executeAt.epoch());
+ if (removeDeps != null)
+ deps = deps.without(removeDeps);
+ ExecuteFlags flags = executeFlags();
+ return new AcceptReply(successful, deps, flags);
+ }
+ finally
+ {
+ close();
+ }
+ }
+ }
+
+
+ public static final class AcceptReply extends
AbstractDepsReply<AcceptReply>
{
public static final AcceptReply SUCCESS = new AcceptReply(Success);
@@ -437,9 +520,9 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
case Success:
return "AcceptOk{deps=" + deps + '}';
case Redundant:
- return "AcceptRedundant(" + supersededBy + ',' +
committedExecuteAt + ")";
+ return "AcceptRedundant(" + supersededBy + ',' +
committedExecuteAt + ')';
case RejectedBallot:
- return "AcceptNack(" + supersededBy + ")";
+ return "AcceptNack(" + supersededBy + ')';
}
}
}
@@ -468,6 +551,12 @@ public class Accept extends
RouteRequest.WithUnsynced<Accept.AcceptReply>
return node.commandStores().mapReduceConsume(txnId.epoch(),
txnId.epoch(), this);
}
+ @Override
+ protected void acceptInternal(AcceptReply reply, Throwable failure)
+ {
+ acceptReply(reply, failure);
+ }
+
@Override
public AcceptReply applyInternal(SafeCommandStore safeStore)
{
diff --git a/accord-core/src/main/java/accord/messages/Apply.java
b/accord-core/src/main/java/accord/messages/Apply.java
index 4bf4dd35..15cc6f50 100644
--- a/accord-core/src/main/java/accord/messages/Apply.java
+++ b/accord-core/src/main/java/accord/messages/Apply.java
@@ -152,7 +152,7 @@ public class Apply extends RouteRequest<ApplyReply>
deps = null;
writes = null;
result = null;
- if (reply != null || failure != null) super.acceptInternal(reply,
failure);
+ if (reply != null || failure != null) acceptReply(reply, failure);
else Invariants.require(isCancelled());
}
diff --git
a/accord-core/src/main/java/accord/messages/ApplyThenWaitUntilApplied.java
b/accord-core/src/main/java/accord/messages/ApplyThenWaitUntilApplied.java
index c85bf903..29297d7f 100644
--- a/accord-core/src/main/java/accord/messages/ApplyThenWaitUntilApplied.java
+++ b/accord-core/src/main/java/accord/messages/ApplyThenWaitUntilApplied.java
@@ -23,6 +23,7 @@ import org.slf4j.LoggerFactory;
import accord.api.Result;
import accord.api.Result.PersistableResult;
+import accord.local.LoadKeys;
import accord.local.Node;
import accord.local.SafeCommandStore;
import accord.local.StoreParticipants;
@@ -41,6 +42,7 @@ import accord.primitives.Writes;
import accord.topology.Topologies;
import accord.utils.UnhandledEnum;
+import static accord.api.ProtocolModifiers.loadKeysAsyncIfPermitted;
import static
accord.messages.MessageType.StandardMessage.APPLY_THEN_WAIT_UNTIL_APPLIED_REQ;
import static accord.messages.RouteRequest.computeScope;
@@ -147,6 +149,12 @@ public class ApplyThenWaitUntilApplied extends
WaitUntilApplied
result = null;
}
+ @Override
+ public LoadKeys loadKeys()
+ {
+ return loadKeysAsyncIfPermitted(txnId);
+ }
+
@Override
public MessageType type()
{
diff --git a/accord-core/src/main/java/accord/messages/BeginInvalidation.java
b/accord-core/src/main/java/accord/messages/BeginInvalidation.java
index fb0f141e..32175cc4 100644
--- a/accord-core/src/main/java/accord/messages/BeginInvalidation.java
+++ b/accord-core/src/main/java/accord/messages/BeginInvalidation.java
@@ -59,6 +59,12 @@ public class BeginInvalidation extends
ParticipantsRequest<Participants<?>, Begi
return node.commandStores().mapReduceConsume(txnId.epoch(),
txnId.epoch(), this);
}
+ @Override
+ protected void acceptInternal(InvalidateReply reply, Throwable failure)
+ {
+ acceptReply(reply, failure);
+ }
+
@Override
public InvalidateReply applyInternal(SafeCommandStore safeStore)
{
diff --git a/accord-core/src/main/java/accord/messages/BeginRecovery.java
b/accord-core/src/main/java/accord/messages/BeginRecovery.java
index 65c2da32..2321d7ea 100644
--- a/accord-core/src/main/java/accord/messages/BeginRecovery.java
+++ b/accord-core/src/main/java/accord/messages/BeginRecovery.java
@@ -19,13 +19,22 @@
package accord.messages;
import java.util.Collection;
+
import javax.annotation.Nullable;
import accord.api.Result;
-import accord.local.*;
-import accord.local.Node.Id;
+import accord.local.Command;
+import accord.local.CommandSummaries;
import accord.local.CommandSummaries.IsDep;
import accord.local.CommandSummaries.SummaryStatus;
+import accord.local.Commands;
+import accord.local.DepsCalculator.SynchronousDepsCalculator;
+import accord.local.LoadKeys;
+import accord.local.FindKeys;
+import accord.local.Node.Id;
+import accord.local.SafeCommand;
+import accord.local.SafeCommandStore;
+import accord.local.StoreParticipants;
import accord.primitives.Ballot;
import accord.primitives.Deps;
import accord.primitives.FullRoute;
@@ -48,9 +57,9 @@ import accord.utils.TinyEnumSet;
import accord.utils.UnhandledEnum;
import accord.utils.async.Cancellable;
+import static accord.local.CommandSummaries.SummaryStatus.ACCEPTED;
import static accord.local.CommandSummaries.SummaryStatus.APPLIED;
import static
accord.local.CommandSummaries.SummaryStatus.NOT_DIRECTLY_WITNESSED;
-import static accord.local.CommandSummaries.SummaryStatus.ACCEPTED;
import static accord.local.CommandSummaries.SummaryStatus.STABLE;
import static accord.messages.BeginRecovery.RecoverReply.Kind.Ok;
import static accord.messages.BeginRecovery.RecoverReply.Kind.Reject;
@@ -133,6 +142,12 @@ public class BeginRecovery extends
RouteRequest.WithUnsynced<BeginRecovery.Recov
return node.commandStores().mapReduceConsume(minEpoch,
executeAtOrTxnIdEpoch, this);
}
+ @Override
+ protected void acceptInternal(RecoverReply reply, Throwable failure)
+ {
+ acceptReply(reply, failure);
+ }
+
@Override
public RecoverReply applyInternal(SafeCommandStore safeStore)
{
@@ -156,7 +171,7 @@ public class BeginRecovery extends
RouteRequest.WithUnsynced<BeginRecovery.Recov
Deps localDeps = null;
if (!command.known().deps().hasCommittedOrDecidedDeps() &&
calculateDeps())
{
- localDeps = DepsCalculator.calculateDeps(safeStore, txnId,
participants, minEpoch, txnId, false);
+ localDeps = SynchronousDepsCalculator.calculateDeps(safeStore,
txnId, participants, minEpoch, txnId, false);
}
if (localDeps != null && coordinatedDeps != null &&
!participants.touches().equals(coordinatedDeps.covering))
{
@@ -273,13 +288,13 @@ public class BeginRecovery extends
RouteRequest.WithUnsynced<BeginRecovery.Recov
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
if (recoverFastPath())
- return LoadKeysFor.RECOVERY;
+ return FindKeys.SUPERSEDING;
if (calculateDeps())
- return LoadKeysFor.READ_WRITE;
- return LoadKeysFor.WRITE;
+ return FindKeys.CONFLICTS;
+ return FindKeys.DECLARED;
}
@Override
diff --git
a/accord-core/src/main/java/accord/messages/GetEphemeralReadDeps.java
b/accord-core/src/main/java/accord/messages/GetEphemeralReadDeps.java
index a3b7caf4..0a7ff975 100644
--- a/accord-core/src/main/java/accord/messages/GetEphemeralReadDeps.java
+++ b/accord-core/src/main/java/accord/messages/GetEphemeralReadDeps.java
@@ -18,13 +18,16 @@
package accord.messages;
+import java.util.function.Function;
+
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
import accord.coordinate.ExecuteFlag.ExecuteFlags;
import accord.local.DepsCalculator;
+import accord.local.DepsCalculator.AbstractDepsReply;
import accord.local.LoadKeys;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
import accord.local.Node.Id;
import accord.local.SafeCommandStore;
import accord.local.StoreParticipants;
@@ -37,7 +40,9 @@ import accord.topology.Topologies;
import accord.utils.Invariants;
import accord.utils.async.Cancellable;
-public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<GetEphemeralReadDeps.GetEphemeralReadDepsOk>
+import static accord.api.ProtocolModifiers.loadKeysAsyncIfPermitted;
+
+public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<ReplyList<GetEphemeralReadDeps.GetEphemeralReadDepsOk>>
{
public static final class SerializationSupport
{
@@ -68,7 +73,14 @@ public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<GetEphemeral
}
@Override
- public GetEphemeralReadDepsOk applyInternal(SafeCommandStore safeStore)
+ protected void acceptInternal(ReplyList<GetEphemeralReadDepsOk> replies,
Throwable failure)
+ {
+ if (failure != null) acceptReply(null, failure);
+ else ReplyList.invoke(replies, GetEphemeralReadDepsOk::reduce,
this::acceptReply);
+ }
+
+ @Override
+ public ReplyList<GetEphemeralReadDepsOk> applyInternal(SafeCommandStore
safeStore)
{
long latestEpoch = Math.max(safeStore.node().epoch(), node.epoch());
@@ -77,23 +89,41 @@ public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<GetEphemeral
if (latestEpoch > executionEpoch &&
(safeStore.ranges().removed(executionEpoch,
latestEpoch).intersects(participants.owns()) ||
node.topology().active().hasReplicationMaybeChanged(participants.owns(),
executionEpoch)))
return new GetEphemeralReadDepsOk(latestEpoch);
- Deps deps;
- ExecuteFlags flags;
- try (DepsCalculator calculator = new DepsCalculator(txnId))
+ //noinspection resource
+ GetEphemeralReadDepsCalculator calculator = new
GetEphemeralReadDepsCalculator(txnId, participants, latestEpoch);
+ return calculator.calculate(safeStore, minEpoch, false);
+ }
+
+ static class GetEphemeralReadDepsCalculator extends
DepsCalculator.DepsReplyCalculator<GetEphemeralReadDepsOk> implements
Function<Void, GetEphemeralReadDepsOk>
+ {
+ final long latestEpoch;
+
+ public GetEphemeralReadDepsCalculator(TxnId txnId, StoreParticipants
participants, long latestEpoch)
{
- deps = calculator.calculate(safeStore, txnId, participants,
minEpoch, Timestamp.MAX, false);
- flags = calculator.executeFlags(txnId);
+ super(txnId, Timestamp.MAX, participants);
+ this.latestEpoch = latestEpoch;
+ }
+
+ public GetEphemeralReadDepsOk apply(Void ignore)
+ {
+ try
+ {
+ Deps deps = deps();
+ ExecuteFlags flags = executeFlags();
+ return new GetEphemeralReadDepsOk(deps, latestEpoch, flags);
+ }
+ finally
+ {
+ close();
+ }
}
- return new GetEphemeralReadDepsOk(deps, latestEpoch, flags);
}
+
@Override
- public GetEphemeralReadDepsOk reduce(GetEphemeralReadDepsOk r1,
GetEphemeralReadDepsOk r2)
+ public ReplyList<GetEphemeralReadDepsOk>
reduce(ReplyList<GetEphemeralReadDepsOk> r1, ReplyList<GetEphemeralReadDepsOk>
r2)
{
- long latestEpoch = Math.max(r1.latestEpoch, r2.latestEpoch);
- if (r1.deps == null || r2.deps == null)
- return new GetEphemeralReadDepsOk(latestEpoch);
- return new GetEphemeralReadDepsOk(r1.deps.with(r2.deps), latestEpoch,
r1.flags.and(r2.flags));
+ return ReplyList.merge(r1, r2);
}
@Override
@@ -114,16 +144,16 @@ public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<GetEphemeral
@Override
public LoadKeys loadKeys()
{
- return LoadKeys.SYNC;
+ return loadKeysAsyncIfPermitted(txnId);
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
- return LoadKeysFor.READ_WRITE;
+ return FindKeys.CONFLICTS;
}
- public static class GetEphemeralReadDepsOk implements Reply
+ public static class GetEphemeralReadDepsOk extends
AbstractDepsReply<GetEphemeralReadDepsOk>
{
public enum Flag { READY_TO_EXECUTE }
@@ -156,6 +186,14 @@ public class GetEphemeralReadDeps extends
RouteRequest.WithUnsynced<GetEphemeral
{
return MessageType.StandardMessage.GET_EPHEMERAL_READ_DEPS_RSP;
}
+
+ private static GetEphemeralReadDepsOk reduce(GetEphemeralReadDepsOk
r1, GetEphemeralReadDepsOk r2)
+ {
+ long latestEpoch = Math.max(r1.latestEpoch, r2.latestEpoch);
+ if (r1.deps == null || r2.deps == null)
+ return new GetEphemeralReadDepsOk(latestEpoch);
+ return new GetEphemeralReadDepsOk(r1.deps.with(r2.deps),
latestEpoch, r1.flags.and(r2.flags));
+ }
}
}
diff --git a/accord-core/src/main/java/accord/messages/GetLatestDeps.java
b/accord-core/src/main/java/accord/messages/GetLatestDeps.java
index e9993e50..42790988 100644
--- a/accord-core/src/main/java/accord/messages/GetLatestDeps.java
+++ b/accord-core/src/main/java/accord/messages/GetLatestDeps.java
@@ -18,19 +18,22 @@
package accord.messages;
+import java.util.function.Function;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
import accord.local.Command;
import accord.local.DepsCalculator;
+import accord.local.DepsCalculator.AbstractDepsReply;
import accord.local.LoadKeys;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
import accord.local.Node.Id;
import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
import accord.local.StoreParticipants;
import accord.primitives.Ballot;
import accord.primitives.Deps;
+import accord.primitives.Known.KnownDeps;
import accord.primitives.LatestDeps;
import accord.primitives.PartialDeps;
import accord.primitives.Route;
@@ -44,7 +47,7 @@ import accord.utils.async.Cancellable;
import static accord.messages.MessageType.StandardMessage.GET_LATEST_DEPS_REQ;
import static accord.messages.MessageType.StandardMessage.GET_LATEST_DEPS_RSP;
-public class GetLatestDeps extends
RouteRequest.WithUnsynced<GetLatestDeps.GetLatestDepsReply>
+public class GetLatestDeps extends
RouteRequest.WithUnsynced<ReplyList<GetLatestDeps.GetLatestDepsReply>>
{
public static final class SerializationSupport
{
@@ -78,7 +81,14 @@ public class GetLatestDeps extends
RouteRequest.WithUnsynced<GetLatestDeps.GetLa
}
@Override
- public GetLatestDepsReply applyInternal(SafeCommandStore safeStore)
+ protected void acceptInternal(ReplyList<GetLatestDepsReply> replies,
Throwable failure)
+ {
+ if (failure != null) acceptReply(null, failure);
+ else ReplyList.invoke(replies, GetLatestDepsReply::reduce,
this::acceptReply);
+ }
+
+ @Override
+ public ReplyList<GetLatestDepsReply> applyInternal(SafeCommandStore
safeStore)
{
StoreParticipants participants = StoreParticipants.read(safeStore,
scope, txnId, minEpoch, executeAt.epoch());
SafeCommand safeCommand = safeStore.get(txnId, participants);
@@ -86,26 +96,28 @@ public class GetLatestDeps extends
RouteRequest.WithUnsynced<GetLatestDeps.GetLa
if (ballot != null)
{
if (command.promised().compareTo(ballot) > 0)
- return GetLatestDepsNack.INSTANCE;
+ return GetLatestDepsReply.NACK;
command = safeCommand.updatePromised(ballot);
}
+
PartialDeps coordinatedDeps = command.partialDeps();
- Deps localDeps = null;
- if (!command.known().deps().hasCommittedOrDecidedDeps() &&
!command.hasBeen(Status.Truncated))
+ KnownDeps knownDeps = command.known().deps();
+ Ballot acceptedOrCommitted = command.acceptedOrCommitted();
+ if (knownDeps.hasCommittedOrDecidedDeps() ||
command.hasBeen(Status.Truncated))
{
- localDeps = DepsCalculator.calculateDeps(safeStore, txnId,
participants, minEpoch, txnId, false);
+ LatestDeps deps = LatestDeps.create(participants.owns(),
knownDeps, acceptedOrCommitted, coordinatedDeps, null);
+ return new GetLatestDepsReply(deps);
}
- LatestDeps deps = LatestDeps.create(participants.owns(),
command.known().deps(), command.acceptedOrCommitted(), coordinatedDeps,
localDeps);
- return new GetLatestDepsOk(deps);
+ //noinspection resource
+ GetLatestDepsCalculator calculator = new
GetLatestDepsCalculator(txnId, participants, knownDeps, acceptedOrCommitted,
coordinatedDeps);
+ return calculator.calculate(safeStore, minEpoch, false);
}
@Override
- public GetLatestDepsReply reduce(GetLatestDepsReply r1, GetLatestDepsReply
r2)
+ public ReplyList<GetLatestDepsReply> reduce(ReplyList<GetLatestDepsReply>
r1, ReplyList<GetLatestDepsReply> r2)
{
- if (!r1.isOk()) return r1;
- if (!r2.isOk()) return r2;
- return new
GetLatestDepsOk(LatestDeps.merge(((GetLatestDepsOk)r1).deps,
((GetLatestDepsOk)r2).deps));
+ return ReplyList.merge(r1, r2);
}
@Override
@@ -131,47 +143,61 @@ public class GetLatestDeps extends
RouteRequest.WithUnsynced<GetLatestDeps.GetLa
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
- return LoadKeysFor.READ_WRITE;
+ return FindKeys.CONFLICTS;
}
- public interface GetLatestDepsReply extends Reply
+ static class GetLatestDepsCalculator extends
DepsCalculator.DepsReplyCalculator<GetLatestDepsReply> implements
Function<Void, GetLatestDepsReply>
{
- boolean isOk();
- }
+ final StoreParticipants participants;
+ final KnownDeps known;
+ final Ballot acceptedOrCommitted;
+ final Deps coordinatedDeps;
- public static final class GetLatestDepsNack implements GetLatestDepsReply
- {
- public static final GetLatestDepsNack INSTANCE = new
GetLatestDepsNack();
- private GetLatestDepsNack(){}
-
- @Override
- public boolean isOk()
+ public GetLatestDepsCalculator(TxnId txnId, StoreParticipants
participants, KnownDeps known, Ballot acceptedOrCommitted, Deps coordinatedDeps)
{
- return false;
+ super(txnId, txnId, participants);
+ this.participants = participants;
+ this.known = known;
+ this.acceptedOrCommitted = acceptedOrCommitted;
+ this.coordinatedDeps = coordinatedDeps;
}
- @Override
- public MessageType type()
+ public GetLatestDepsReply apply(Void ignore)
{
- return GET_LATEST_DEPS_RSP;
+ try
+ {
+ Deps localDeps = deps();
+ LatestDeps deps = LatestDeps.create(participants.owns(),
known, acceptedOrCommitted, coordinatedDeps, localDeps);
+ return new GetLatestDepsReply(deps);
+ }
+ finally
+ {
+ close();
+ }
}
}
- public static class GetLatestDepsOk implements GetLatestDepsReply
+ public static class GetLatestDepsReply extends
AbstractDepsReply<GetLatestDepsReply>
{
+ public static final GetLatestDepsReply NACK = new GetLatestDepsReply();
public final LatestDeps deps;
- public GetLatestDepsOk(@Nonnull LatestDeps deps)
+ public GetLatestDepsReply(@Nonnull LatestDeps deps)
{
this.deps = Invariants.nonNull(deps);
}
+ private GetLatestDepsReply()
+ {
+ this.deps = null;
+ }
+
@Override
public String toString()
{
- return "GetLatestDepsOk{" + deps + '}' ;
+ return "GetLatestDepsReply{" + deps + '}' ;
}
@Override
@@ -180,11 +206,16 @@ public class GetLatestDeps extends
RouteRequest.WithUnsynced<GetLatestDeps.GetLa
return GET_LATEST_DEPS_RSP;
}
- @Override
public boolean isOk()
{
- return true;
+ return deps != null;
}
- }
+ private static GetLatestDepsReply reduce(GetLatestDepsReply r1,
GetLatestDepsReply r2)
+ {
+ if (!r1.isOk()) return r1;
+ if (!r2.isOk()) return r2;
+ return new GetLatestDepsReply(LatestDeps.merge(r1.deps, r2.deps));
+ }
+ }
}
diff --git a/accord-core/src/main/java/accord/messages/GetMaxConflict.java
b/accord-core/src/main/java/accord/messages/GetMaxConflict.java
index c03a2e80..d566d7a0 100644
--- a/accord-core/src/main/java/accord/messages/GetMaxConflict.java
+++ b/accord-core/src/main/java/accord/messages/GetMaxConflict.java
@@ -63,6 +63,12 @@ public class GetMaxConflict extends
RouteRequest.WithUnsynced<GetMaxConflict.Get
return node.commandStores().mapReduceConsume(minEpoch, executionEpoch,
this);
}
+ @Override
+ protected void acceptInternal(GetMaxConflictOk reply, Throwable failure)
+ {
+ acceptReply(reply, failure);
+ }
+
@Override
public GetMaxConflictOk applyInternal(SafeCommandStore safeStore)
{
diff --git a/accord-core/src/main/java/accord/messages/InformDurable.java
b/accord-core/src/main/java/accord/messages/InformDurable.java
index e354a74b..c163150a 100644
--- a/accord-core/src/main/java/accord/messages/InformDurable.java
+++ b/accord-core/src/main/java/accord/messages/InformDurable.java
@@ -143,6 +143,12 @@ public class InformDurable extends RouteRequest<Reply>
implements ExecutionConte
return node.commandStores().mapReduceConsume(minEpoch, maxEpoch, this);
}
+ @Override
+ protected void acceptInternal(Reply reply, Throwable failure)
+ {
+ acceptReply(reply, failure);
+ }
+
@Override
public Reply applyInternal(SafeCommandStore safeStore)
{
diff --git a/accord-core/src/main/java/accord/messages/NoWaitRequest.java
b/accord-core/src/main/java/accord/messages/NoWaitRequest.java
index 10de253c..a5b3b8f1 100644
--- a/accord-core/src/main/java/accord/messages/NoWaitRequest.java
+++ b/accord-core/src/main/java/accord/messages/NoWaitRequest.java
@@ -37,7 +37,7 @@ import accord.utils.async.Cancellable;
import static accord.utils.Invariants.illegalState;
import static java.util.concurrent.TimeUnit.MICROSECONDS;
-public abstract class NoWaitRequest<P extends Participants<?>, R extends
Reply> extends AbstractRequest<P, R> implements Timeouts.Timeout
+public abstract class NoWaitRequest<P extends Participants<?>, R> extends
AbstractRequest<P, R> implements Timeouts.Timeout
{
public static final CancellationException CANCELLATION_EXCEPTION = new
CancellationException();
@@ -127,27 +127,31 @@ public abstract class NoWaitRequest<P extends
Participants<?>, R extends Reply>
acceptInternal(reply, failure);
}
- protected void acceptInternal(R reply, Throwable failure)
+ protected abstract void acceptInternal(R reply, Throwable failure);
+
+ protected void acceptReply(Reply reply, Throwable failure)
{
- if (reply == null && failure == null)
{
- Invariants.require(isCancelled());
- if (!(replyContext instanceof LocalDelivery<?>))
+ if (reply == null && failure == null)
{
- if (tracing() != null)
- tracing().trace(null, "Completed with no reply");
- return; // for now we don't report cancellation/timeout
remotely, and rely on the coordinator's timeouts
+ Invariants.require(isCancelled());
+ if (!(replyContext instanceof LocalDelivery<?>))
+ {
+ if (tracing() != null)
+ tracing().trace(null, "Completed with no reply");
+ return; // for now we don't report cancellation/timeout
remotely, and rely on the coordinator's timeouts
+ }
+ // we must report something for local delivery, as we rely on
this callback instead of registering a separate timeout
+ failure = CANCELLATION_EXCEPTION;
}
- // we must report something for local delivery, as we rely on this
callback instead of registering a separate timeout
- failure = CANCELLATION_EXCEPTION;
- }
- if (failure != null || reply.isFinal())
- {
- Invariants.require(!hasSentFinalReply);
- hasSentFinalReply = true;
+ if (failure != null || reply.isFinal())
+ {
+ Invariants.require(!hasSentFinalReply);
+ hasSentFinalReply = true;
+ }
+ if (failure != null) cancel();
+ node.reply(replyTo, replyContext, reply, failure, tracing());
}
- if (failure != null) cancel();
- node.reply(replyTo, replyContext, reply, failure, tracing());
}
@Override
diff --git a/accord-core/src/main/java/accord/messages/ParticipantsRequest.java
b/accord-core/src/main/java/accord/messages/ParticipantsRequest.java
index 7d3f240e..1d308f5c 100644
--- a/accord-core/src/main/java/accord/messages/ParticipantsRequest.java
+++ b/accord-core/src/main/java/accord/messages/ParticipantsRequest.java
@@ -38,7 +38,7 @@ import accord.utils.async.Cancellable;
import static accord.topology.Shard.Flag.MUST_WITNESS;
import static accord.utils.Invariants.illegalArgument;
-public abstract class ParticipantsRequest<P extends Participants<?>, R extends
Reply> extends NoWaitRequest<P, R>
+public abstract class ParticipantsRequest<P extends Participants<?>, R>
extends NoWaitRequest<P, R>
{
public final long waitForEpoch;
diff --git a/accord-core/src/main/java/accord/messages/PreAccept.java
b/accord-core/src/main/java/accord/messages/PreAccept.java
index 77f3e32d..8c8d3a97 100644
--- a/accord-core/src/main/java/accord/messages/PreAccept.java
+++ b/accord-core/src/main/java/accord/messages/PreAccept.java
@@ -19,6 +19,7 @@
package accord.messages;
import java.util.Objects;
+import java.util.function.Function;
import javax.annotation.Nullable;
import org.slf4j.Logger;
@@ -27,17 +28,19 @@ import org.slf4j.LoggerFactory;
import accord.coordinate.ExecuteFlag.ExecuteFlags;
import accord.local.Command;
import accord.local.Commands;
-import accord.local.DepsCalculator;
+import accord.local.DepsCalculator.AbstractDepsReply;
+import accord.local.DepsCalculator.DepsReplyCalculator;
import accord.local.LoadKeys;
-import accord.local.LoadKeysFor;
+import accord.local.FindKeys;
+import accord.local.Node;
import accord.local.Node.Id;
import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
-import accord.primitives.PartialDeps;
import accord.local.StoreParticipants;
import accord.messages.RouteRequest.WithUnsynced;
import accord.primitives.Deps;
import accord.primitives.FullRoute;
+import accord.primitives.PartialDeps;
import accord.primitives.PartialTxn;
import accord.primitives.Route;
import accord.primitives.Status;
@@ -49,11 +52,12 @@ import accord.utils.Invariants;
import accord.utils.UnhandledEnum;
import accord.utils.async.Cancellable;
+import static accord.api.ProtocolModifiers.loadKeysAsyncIfPermitted;
import static accord.messages.MessageType.StandardMessage.PRE_ACCEPT_REQ;
import static accord.messages.MessageType.StandardMessage.PRE_ACCEPT_RSP;
import static accord.primitives.Timestamp.Flag.REJECTED;
-public class PreAccept extends WithUnsynced<PreAccept.PreAcceptReply>
+public class PreAccept extends
WithUnsynced<ReplyList<PreAccept.PreAcceptReply>>
{
@SuppressWarnings("unused")
private static final Logger logger =
LoggerFactory.getLogger(PreAccept.class);
@@ -98,13 +102,13 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
@Override
public LoadKeys loadKeys()
{
- return LoadKeys.SYNC;
+ return loadKeysAsyncIfPermitted(txnId);
}
@Override
- public LoadKeysFor loadKeysFor()
+ public FindKeys findKeys()
{
- return LoadKeysFor.READ_WRITE;
+ return FindKeys.CONFLICTS;
}
@Override
@@ -120,7 +124,14 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
}
@Override
- public PreAcceptReply applyInternal(SafeCommandStore safeStore)
+ protected void acceptInternal(ReplyList<PreAcceptReply> replies, Throwable
failure)
+ {
+ if (failure != null) acceptReply(null, failure);
+ else ReplyList.invoke(replies, PreAcceptReply::reduce,
this::acceptReply);
+ }
+
+ @Override
+ public ReplyList<PreAcceptReply> applyInternal(SafeCommandStore safeStore)
{
StoreParticipants participants = StoreParticipants.update(safeStore,
route, minEpoch, txnId, acceptEpoch);
SafeCommand safeCommand = safeStore.get(txnId, participants);
@@ -142,25 +153,13 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
return new PreAcceptOk(txnId, command.executeAt(),
Deps.NONE, ExecuteFlags.none());
case Retired:
- Timestamp executeAt;
- ExecuteFlags flags;
- Deps deps;
- try (DepsCalculator calculator = new DepsCalculator(txnId))
- {
- deps = calculator.calculate(safeStore, txnId,
participants, minEpoch, txnId, true);
- if (deps == null)
- return PreAcceptNack.INSTANCE;
- flags = calculator.executeFlags(txnId);
- executeAt = calculator.executeAt(safeCommand, node);
- }
-
- // NOTE: we CANNOT test whether we adopt a future dependency
here because it might be that this command
- // is guaranteed to not reach agreement, but that this replica
is unaware of that fact and has pruned
- // all preceding transactions. In which case we may be able to
adopt a future dependency but won't propose it.
- // We do however prohibit later epochs as dependencies as we
cannot handle those effectively
- // when back-filling for execution of the transaction.
- Invariants.require(deps.maxTxnId(txnId).epoch() <=
txnId.epoch());
- return new PreAcceptOk(txnId, executeAt, deps, flags);
+ Timestamp witnessedAt = command.executeAtOrTxnId(); // if
retired, executeAt may be null
+ //noinspection resource
+ PreAcceptDepsCalculator calculator = new
PreAcceptDepsCalculator(txnId, witnessedAt, participants, node);
+ ReplyList<PreAcceptReply> reply =
calculator.calculate(safeStore, minEpoch, true);
+ if (reply == null)
+ return PreAcceptNack.INSTANCE;
+ return reply;
case Truncated:
case RejectedBallot:
@@ -169,9 +168,9 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
}
@Override
- public PreAcceptReply reduce(PreAcceptReply r1, PreAcceptReply r2)
+ public ReplyList<PreAcceptReply> reduce(ReplyList<PreAcceptReply> r1,
ReplyList<PreAcceptReply> r2)
{
- return PreAcceptReply.reduce(r1, r2);
+ return ReplyList.merge(r1, r2);
}
@Override
@@ -180,7 +179,7 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
return PRE_ACCEPT_REQ;
}
- public static abstract class PreAcceptReply implements Reply
+ public static abstract class PreAcceptReply extends
AbstractDepsReply<PreAcceptReply>
{
@Override
public MessageType type()
@@ -209,6 +208,41 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
}
}
+ public static class PreAcceptDepsCalculator extends
DepsReplyCalculator<PreAcceptReply> implements Function<Void, PreAcceptReply>
+ {
+ final Node node;
+ final Timestamp witnessedAt;
+
+ public PreAcceptDepsCalculator(TxnId txnId, Timestamp witnessedAt,
StoreParticipants participants, Node node)
+ {
+ super(txnId, txnId, participants);
+ this.node = node;
+ this.witnessedAt = Invariants.nonNull(witnessedAt);
+ }
+
+ public PreAcceptOk apply(Void ignore)
+ {
+ try
+ {
+ Deps deps = deps();
+ Timestamp executeAt = executeAt(witnessedAt, node);
+ ExecuteFlags flags = executeFlags();
+
+ // NOTE: we CANNOT test whether we adopt a future dependency
here because it might be that this command
+ // is guaranteed to not reach agreement, but that this replica
is unaware of that fact and has pruned
+ // all preceding transactions. In which case we may be able to
adopt a future dependency but won't propose it.
+ // We do however prohibit later epochs as dependencies as we
cannot handle those effectively
+ // when back-filling for execution of the transaction.
+ Invariants.require(deps.maxTxnId(txnId).epoch() <=
txnId.epoch());
+ return new PreAcceptOk(txnId, executeAt, deps, flags);
+ }
+ finally
+ {
+ close();
+ }
+ }
+ }
+
public static class PreAcceptOk extends PreAcceptReply
{
public final TxnId txnId;
@@ -219,7 +253,7 @@ public class PreAccept extends
WithUnsynced<PreAccept.PreAcceptReply>
public PreAcceptOk(TxnId txnId, Timestamp witnessedAt, Deps deps,
ExecuteFlags flags)
{
this.txnId = txnId;
- this.witnessedAt = witnessedAt;
+ this.witnessedAt = Invariants.nonNull(witnessedAt);
this.deps = deps;
this.flags = flags;
}
diff --git a/accord-core/src/main/java/accord/messages/ReadData.java
b/accord-core/src/main/java/accord/messages/ReadData.java
index fd13a668..b0258ed2 100644
--- a/accord-core/src/main/java/accord/messages/ReadData.java
+++ b/accord-core/src/main/java/accord/messages/ReadData.java
@@ -35,6 +35,7 @@ import accord.local.Command;
import accord.local.Command.Committed;
import accord.local.CommandStore;
import accord.local.CommandStores;
+import accord.local.LoadKeys;
import accord.local.Node;
import accord.local.SafeCommand;
import accord.local.SafeCommandStore;
@@ -737,14 +738,6 @@ public abstract class ReadData extends
AbstractRequest<Participants<?>, ReadData
}
}
- @Override
- public Unseekables<?> keys()
- {
- if (flags.contains(READY_TO_EXECUTE) &&
fastReadsMayBypassCommandsForKey(txnId))
- return RoutingKeys.EMPTY;
- return scope;
- }
-
protected void reply(Ranges unavailable, Data data, long uniqueHlc)
{
if (data != null && !txnId.awaitsOnlyDeps() &&
!data.validateReply(txnId, executeAt, validateHlc()))
diff --git a/accord-core/src/main/java/accord/messages/ReplyList.java
b/accord-core/src/main/java/accord/messages/ReplyList.java
new file mode 100644
index 00000000..452c6c13
--- /dev/null
+++ b/accord-core/src/main/java/accord/messages/ReplyList.java
@@ -0,0 +1,93 @@
+/*
+ * 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 accord.messages;
+
+import java.util.ArrayList;
+import java.util.function.BiConsumer;
+
+import accord.utils.Reduce;
+import accord.utils.async.AsyncResults;
+import accord.utils.async.CancellableAsyncResult;
+
+public interface ReplyList<R extends CancellableAsyncResult<R> & ReplyList<R>>
+{
+ int size();
+ CancellableAsyncResult<R> get(int i);
+ void cancelReplies();
+
+ final class ReplyMultiList<R extends CancellableAsyncResult<R> &
ReplyList<R>> extends ArrayList<CancellableAsyncResult<R>> implements
ReplyList<R>
+ {
+ public ReplyMultiList() {}
+ public ReplyMultiList(int initialCapacity) { super(initialCapacity); }
+
+ @Override
+ public void cancelReplies()
+ {
+ forEach(CancellableAsyncResult::cancel);
+ }
+ }
+
+ // parameters must be EITHER a ReplyMultiList<R> OR a single R (that
itself implements both ReplyList<R> and AsyncResult<R>)
+ static <R extends CancellableAsyncResult<R> & ReplyList<R>> ReplyList<R>
merge(ReplyList<R> r1, ReplyList<R> r2)
+ {
+ if (r1 == null || r2 == null)
+ {
+ if (r1 != null) r1.cancelReplies();
+ if (r2 != null) r2.cancelReplies();
+ return null;
+ }
+
+ ReplyMultiList<R> multi;
+ if (r1 instanceof ReplyMultiList)
+ {
+ multi = (ReplyMultiList<R>) r1;
+ if (r2 instanceof ReplyMultiList)
+ {
+ multi.addAll((ReplyMultiList<R>)r2);
+ return multi;
+ }
+ multi.add((CancellableAsyncResult<R>) r2);
+ }
+ else if (r2 instanceof ReplyMultiList)
+ {
+ multi = (ReplyMultiList<R>) r2;
+ multi.add((CancellableAsyncResult<R>) r1);
+ }
+ else
+ {
+ multi = new ReplyMultiList<>();
+ multi.add((CancellableAsyncResult<R>) r1);
+ multi.add((CancellableAsyncResult<R>) r2);
+ }
+ return multi;
+ }
+
+ static <R extends CancellableAsyncResult<R> & ReplyList<R>> void
invoke(ReplyList<R> replies, Reduce<R, R> reduce, BiConsumer<? super R,
Throwable> callback)
+ {
+ if (replies instanceof ReplyMultiList<?>)
+ {
+ AsyncResults.reduce((ReplyMultiList<R>)replies, reduce)
+ .invoke(callback);
+ }
+ else
+ {
+ ((R)replies).invoke(callback);
+ }
+ }
+}
diff --git a/accord-core/src/main/java/accord/messages/RouteRequest.java
b/accord-core/src/main/java/accord/messages/RouteRequest.java
index 3a070444..21891165 100644
--- a/accord-core/src/main/java/accord/messages/RouteRequest.java
+++ b/accord-core/src/main/java/accord/messages/RouteRequest.java
@@ -27,9 +27,9 @@ import accord.primitives.TxnId;
import accord.topology.Topologies;
import accord.utils.async.Cancellable;
-public abstract class RouteRequest<R extends Reply> extends
ParticipantsRequest<Route<?>, R>
+public abstract class RouteRequest<R> extends ParticipantsRequest<Route<?>, R>
{
- public static abstract class WithUnsynced<R extends Reply> extends
RouteRequest<R>
+ public static abstract class WithUnsynced<R> extends RouteRequest<R>
{
public final long minEpoch; // TODO (low priority, clarity): can this
just always be TxnId.epoch?
diff --git a/accord-core/src/main/java/accord/primitives/Deps.java
b/accord-core/src/main/java/accord/primitives/Deps.java
index 862a2641..27a7b525 100644
--- a/accord-core/src/main/java/accord/primitives/Deps.java
+++ b/accord-core/src/main/java/accord/primitives/Deps.java
@@ -18,9 +18,17 @@
package accord.primitives;
+import java.util.AbstractList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.function.BiFunction;
+import java.util.function.Consumer;
+import java.util.function.Function;
+import java.util.function.Predicate;
+import javax.annotation.Nullable;
+
import accord.api.RoutingKey;
import accord.local.cfk.CommandsForKey;
-import accord.primitives.Routable.Domain;
import accord.utils.IndexedFunction;
import accord.utils.Invariants;
import accord.utils.MergeFewDisjointSortedListsCursor;
@@ -32,19 +40,9 @@ import accord.utils.SortedList.MergeCursor;
import accord.utils.TriFunction;
import accord.utils.UnhandledEnum;
-import java.util.AbstractList;
-import java.util.Arrays;
-import java.util.List;
-import java.util.function.BiFunction;
-import java.util.function.Consumer;
-import java.util.function.Function;
-import java.util.function.Predicate;
-import javax.annotation.Nullable;
-
import static accord.local.cfk.CommandsForKey.managesExecution;
import static accord.primitives.Routables.Slice.Minimal;
import static accord.primitives.Timestamp.Flag.UNSTABLE;
-import static accord.utils.Invariants.illegalState;
/**
* A collection of transaction dependencies, keyed by the key or range on
which they were adopted.
@@ -105,14 +103,6 @@ public class Deps
this.rangeBuilder = buildRangesByTxnId ?
RangeDeps.byTxnIdBuilder() : RangeDeps.builderByRange();
}
- public AbstractBuilder<T> addNormalise(Unseekable keyOrRange, TxnId
txnId)
- {
- if (keyOrRange.domain() == txnId.domain()) add(keyOrRange, txnId);
- else if (keyOrRange.domain() == Domain.Key)
add(keyOrRange.asRange(), txnId);
- else throw illegalState();
- return this;
- }
-
public AbstractBuilder<T> add(Unseekable keyOrRange, TxnId txnId)
{
Invariants.requireArgument(keyOrRange.domain() == txnId.domain(),
"%s is not same domain as %s", keyOrRange, txnId);
diff --git a/accord-core/src/main/java/accord/utils/async/AsyncChain.java
b/accord-core/src/main/java/accord/utils/async/AsyncChain.java
index cda831db..8a22482a 100644
--- a/accord-core/src/main/java/accord/utils/async/AsyncChain.java
+++ b/accord-core/src/main/java/accord/utils/async/AsyncChain.java
@@ -165,6 +165,11 @@ public interface AsyncChain<V>
default AsyncResult<V> beginAsResult()
{
- return AsyncResults.forChain(this);
+ return AsyncResults.begin(this);
+ }
+
+ default CancellableAsyncResult<V> beginAsCancellableResult()
+ {
+ return AsyncResults.beginCancellable(this);
}
}
\ No newline at end of file
diff --git a/accord-core/src/main/java/accord/utils/async/AsyncChains.java
b/accord-core/src/main/java/accord/utils/async/AsyncChains.java
index 5d334433..bca728d9 100644
--- a/accord-core/src/main/java/accord/utils/async/AsyncChains.java
+++ b/accord-core/src/main/java/accord/utils/async/AsyncChains.java
@@ -440,7 +440,6 @@ public class AsyncChains
}
}
-
private static class DetectLeak extends AsyncChains.Head<Void>
{
private final AtomicBoolean called = new AtomicBoolean(false);
diff --git a/accord-core/src/main/java/accord/utils/async/AsyncResults.java
b/accord-core/src/main/java/accord/utils/async/AsyncResults.java
index cc794831..e108b504 100644
--- a/accord-core/src/main/java/accord/utils/async/AsyncResults.java
+++ b/accord-core/src/main/java/accord/utils/async/AsyncResults.java
@@ -209,7 +209,7 @@ public class AsyncResults
}
}
- static class Chain<V> extends AbstractResult<V>
+ public static class Chain<V> extends AbstractResult<V> implements
BiConsumer<V, Throwable>
{
private AsyncChain<V> chain;
public Chain(AsyncChain<V> chain)
@@ -221,8 +221,60 @@ public class AsyncResults
@Override
void setResult(V result, Throwable failure)
{
+ chain = null;
super.setResult(result, failure);
+ }
+
+ @Override
+ public void accept(V success, Throwable fail)
+ {
+ setResult(success, fail);
+ }
+
+ @Override
+ public String toString()
+ {
+ AsyncChain<V> chain = this.chain;
+ if (chain != null)
+ return "Waiting On: " + chain;
+ return super.toString();
+ }
+ }
+
+ public static class CancellableChain<V> extends AbstractResult<V>
implements Cancellable, BiConsumer<V, Throwable>, CancellableAsyncResult<V>
+ {
+ private AsyncChain<V> chain;
+ private Cancellable cancel;
+
+ public CancellableChain(AsyncChain<V> chain)
+ {
+ this.chain = chain;
+ this.cancel = chain.begin(this);
+ }
+
+ @Override
+ void setResult(V result, Throwable failure)
+ {
chain = null;
+ cancel = null;
+ super.setResult(result, failure);
+ }
+
+ @Override
+ public void cancel()
+ {
+ Cancellable cancel = this.cancel;
+ if (cancel == null)
+ return;
+ this.chain = null;
+ this.cancel = null;
+ cancel.cancel();
+ }
+
+ @Override
+ public void accept(V success, Throwable fail)
+ {
+ setResult(success, fail);
}
@Override
@@ -307,7 +359,7 @@ public class AsyncResults
}
}
- static abstract class AbstractImmediate<V> implements AsyncResult<V>
+ public static abstract class AbstractImmediate<V> implements AsyncResult<V>
{
@Override
public AsyncChain<V> chain()
@@ -391,11 +443,19 @@ public class AsyncResults
/**
* Creates an AsyncResult for the given chain. This calls begin on the
supplied chain
*/
- public static <V> AsyncResult<V> forChain(AsyncChain<V> chain)
+ public static <V> AsyncResult<V> begin(AsyncChain<V> chain)
{
return new Chain<>(chain);
}
+ /**
+ * Creates an AsyncResult for the given chain. This calls begin on the
supplied chain
+ */
+ public static <V> CancellableAsyncResult<V> beginCancellable(AsyncChain<V>
chain)
+ {
+ return new CancellableChain<>(chain);
+ }
+
public static <V> AsyncResult<V> success(V value)
{
if (value == null)
diff --git
a/accord-core/src/main/java/accord/utils/async/CancellableAsyncResult.java
b/accord-core/src/main/java/accord/utils/async/CancellableAsyncResult.java
new file mode 100644
index 00000000..1085a10f
--- /dev/null
+++ b/accord-core/src/main/java/accord/utils/async/CancellableAsyncResult.java
@@ -0,0 +1,23 @@
+/*
+ * 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 accord.utils.async;
+
+public interface CancellableAsyncResult<V> extends AsyncResult<V>, Cancellable
+{
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]