diff --git a/modules/accord b/modules/accord index 555337a7d4..68778350bb 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 555337a7d41158f74033818facf94fed6904bf5a +Subproject commit 68778350bb45a1545cbe38af290d7778ffb79454 diff --git a/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java b/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java index 61f64ce3d0..480fcbe951 100644 --- a/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java +++ b/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java @@ -66,6 +66,7 @@ import org.apache.cassandra.utils.concurrent.Refs; import static accord.local.durability.DurabilityService.SyncLocal.Self; import static accord.local.durability.DurabilityService.SyncRemote.NoRemote; import static com.google.common.base.Preconditions.checkNotNull; +import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.config.CassandraRelevantProperties.REPAIR_MUTATION_REPAIR_ROWS_PER_BATCH; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -260,10 +261,11 @@ public class CassandraStreamReceiver implements StreamReceiver { Ranges accordRanges = AccordTopology.toAccordRanges(cfs.getTableId(), ranges); long startedAtNanos = nanoTime(); - long deadlineNanos = startedAtNanos + DatabaseDescriptor.getAccordRangeSyncPointTimeoutNanos(); + long timeoutNanos = DatabaseDescriptor.getAccordRangeSyncPointTimeoutNanos(); + long deadlineNanos = startedAtNanos + timeoutNanos; // TODO (expected): use the source bounds for the streams to avoid waiting unnecessarily long AccordService.getBlocking(accordService.maxConflict(accordRanges) - .flatMap(min -> accordService.sync("[Stream #" + session.planId() + ']', min, accordRanges, null, Self, NoRemote)) + .flatMap(min -> accordService.sync("[Stream #" + session.planId() + ']', min, accordRanges, null, Self, NoRemote, timeoutNanos, NANOSECONDS)) , accordRanges, new LatencyRequestBookkeeping(cfs.metric.accordPostStreamRepair), startedAtNanos, deadlineNanos); } diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 414f7b6d5a..a36d01c97f 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -516,9 +516,9 @@ public class AccordService implements IAccordService, Shutdownable } @Override - public AsyncChain sync(Object requestedBy, Timestamp minBound, Ranges ranges, @Nullable Collection include, DurabilityService.SyncLocal syncLocal, DurabilityService.SyncRemote syncRemote) + public AsyncChain sync(Object requestedBy, Timestamp minBound, Ranges ranges, @Nullable Collection include, DurabilityService.SyncLocal syncLocal, DurabilityService.SyncRemote syncRemote, long timeout, TimeUnit timeoutUnits) { - return node.durability().sync(requestedBy, minBound, ranges, include, syncLocal, syncRemote); + return node.durability().sync(requestedBy, minBound, ranges, include, syncLocal, syncRemote, timeout, timeoutUnits); } @Override @@ -998,10 +998,11 @@ public class AccordService implements IAccordService, Shutdownable if (rangeList.isEmpty()) return; // nothing to see here Ranges ranges = Ranges.of(rangeList.toArray(accord.primitives.Range[]::new)); + long timeout = DatabaseDescriptor.getAccordRepairTimeoutNanos(); long startedAt = nanoTime(); - long deadline = startedAt + DatabaseDescriptor.getAccordRangeSyncPointTimeoutNanos(); + long deadline = startedAt + timeout; // TODO (required): relax this requirement - too expensive - getBlocking(node.durability().sync("Drop Keyspace/Table (Epoch " + epoch + ')', TxnId.minForEpoch(epoch), ranges, Self, All), ranges, new LatencyRequestBookkeeping(null), startedAt, deadline, false); + getBlocking(node.durability().sync("Drop Keyspace/Table (Epoch " + epoch + ')', TxnId.minForEpoch(epoch), ranges, Self, All, DatabaseDescriptor.getAccordRangeSyncPointTimeoutNanos(), NANOSECONDS), ranges, new LatencyRequestBookkeeping(null), startedAt, deadline, false); } public Params journalConfiguration() diff --git a/src/java/org/apache/cassandra/service/accord/IAccordService.java b/src/java/org/apache/cassandra/service/accord/IAccordService.java index adb23a37af..136bd5290e 100644 --- a/src/java/org/apache/cassandra/service/accord/IAccordService.java +++ b/src/java/org/apache/cassandra/service/accord/IAccordService.java @@ -78,7 +78,7 @@ public interface IAccordService IVerbHandler requestHandler(); IVerbHandler responseHandler(); - AsyncChain sync(Object requestedBy, @Nullable Timestamp minBound, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote); + AsyncChain sync(Object requestedBy, @Nullable Timestamp minBound, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote, long timeout, TimeUnit timeoutUnits); AsyncChain sync(@Nullable Timestamp minBound, Keys keys, SyncLocal syncLocal, SyncRemote syncRemote); AsyncChain maxConflict(Ranges ranges); @@ -193,7 +193,7 @@ public interface IAccordService } @Override - public AsyncChain sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote) + public AsyncChain sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote, long timeout, TimeUnit timeoutUnits) { throw new UnsupportedOperationException("No accord transaction should be executed when accord.enabled = false in cassandra.yaml"); } @@ -361,9 +361,9 @@ public interface IAccordService } @Override - public AsyncChain sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote) + public AsyncChain sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, SyncLocal syncLocal, SyncRemote syncRemote, long timeout, TimeUnit timeoutUnits) { - return delegate.sync(requestedBy, onOrAfter, ranges, include, syncLocal, syncRemote); + return delegate.sync(requestedBy, onOrAfter, ranges, include, syncLocal, syncRemote, timeout, timeoutUnits); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/repair/AccordRepair.java b/src/java/org/apache/cassandra/service/accord/repair/AccordRepair.java index d3a68eb20a..59b5abd4c7 100644 --- a/src/java/org/apache/cassandra/service/accord/repair/AccordRepair.java +++ b/src/java/org/apache/cassandra/service/accord/repair/AccordRepair.java @@ -54,6 +54,7 @@ import static accord.local.durability.DurabilityService.SyncRemote.All; import static accord.local.durability.DurabilityService.SyncRemote.Quorum; import static accord.primitives.Timestamp.mergeMax; import static accord.primitives.Timestamp.minForEpoch; +import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.config.DatabaseDescriptor.getAccordRepairTimeoutNanos; /* @@ -149,10 +150,11 @@ public class AccordRepair Ranges ranges = AccordService.intersecting(Ranges.of(range)); waiting = Thread.currentThread(); RequestBookkeeping bookkeeping = new LatencyRequestBookkeeping(latency); + long timeoutNanos = getAccordRepairTimeoutNanos(); AccordService.getBlocking(service.maxConflict(ranges).flatMap(conflict -> { conflict = mergeMax(conflict, minForEpoch(this.minEpoch.getEpoch())); - return service.sync("[repairId #" + repairId + ']', conflict, Ranges.of(range), ids, NoLocal, syncRemote); - }), ranges, bookkeeping, start, start + getAccordRepairTimeoutNanos()); + return service.sync("[repairId #" + repairId + ']', conflict, Ranges.of(range), ids, NoLocal, syncRemote, timeoutNanos, NANOSECONDS); + }), ranges, bookkeeping, start, start + timeoutNanos); waiting = null; if (shouldAbort != null) diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java index 89400fbd55..387d8ce872 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordIncrementalRepairTest.java @@ -89,9 +89,9 @@ public class AccordIncrementalRepairTest extends AccordTestBase } @Override - public AsyncResult sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, DurabilityService.SyncLocal syncLocal, DurabilityService.SyncRemote syncRemote) + public AsyncResult sync(Object requestedBy, @Nullable Timestamp onOrAfter, Ranges ranges, @Nullable Collection include, DurabilityService.SyncLocal syncLocal, DurabilityService.SyncRemote syncRemote, long timeout, TimeUnit timeoutUnits) { - return delegate.sync(requestedBy, onOrAfter, ranges, include, syncLocal, syncRemote).map(v -> { + return delegate.sync(requestedBy, onOrAfter, ranges, include, syncLocal, syncRemote, 10L, TimeUnit.MINUTES).map(v -> { executedBarriers = true; return v; }).beginAsResult(); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/journal/JournalAccessRouteIndexOnStartupRaceTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/journal/JournalAccessRouteIndexOnStartupRaceTest.java index 9ec080cd8d..ba7eff17ec 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/journal/JournalAccessRouteIndexOnStartupRaceTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/journal/JournalAccessRouteIndexOnStartupRaceTest.java @@ -20,6 +20,7 @@ package org.apache.cassandra.distributed.test.accord.journal; import java.io.IOException; import java.util.concurrent.Callable; +import java.util.concurrent.TimeUnit; import org.junit.Test; @@ -89,7 +90,7 @@ public class JournalAccessRouteIndexOnStartupRaceTest extends TestBaseImpl Ranges ranges = Ranges.single(TokenRange.fullRange(metadata.id, metadata.partitioner)); for (int i = 0; i < 10; i++) { - AsyncChains.getBlockingAndRethrow(accord.sync(null, Timestamp.NONE, ranges, null, DurabilityService.SyncLocal.Self, DurabilityService.SyncRemote.Quorum)); + AsyncChains.getBlockingAndRethrow(accord.sync(null, Timestamp.NONE, ranges, null, DurabilityService.SyncLocal.Self, DurabilityService.SyncRemote.Quorum, 10L, TimeUnit.MINUTES)); accord.journal().closeCurrentSegmentForTestingIfNonEmpty(); accord.journal().runCompactorForTesting(); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/journal/StatefulJournalRestartTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/journal/StatefulJournalRestartTest.java index c2cb10f15c..76eeda8878 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/journal/StatefulJournalRestartTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/journal/StatefulJournalRestartTest.java @@ -20,6 +20,7 @@ package org.apache.cassandra.distributed.test.accord.journal; import java.io.IOException; import java.time.Duration; +import java.util.concurrent.TimeUnit; import org.junit.Ignore; import org.junit.Test; @@ -118,7 +119,7 @@ public class StatefulJournalRestartTest extends TestBaseImpl Ranges ranges = Ranges.single(TokenRange.fullRange(metadata.id, metadata.partitioner)); for (int i = 0; i < 10; i++) { - AsyncChains.getBlockingAndRethrow(accord.sync(null, Timestamp.NONE, ranges, null, DurabilityService.SyncLocal.Self, DurabilityService.SyncRemote.Quorum)); + AsyncChains.getBlockingAndRethrow(accord.sync(null, Timestamp.NONE, ranges, null, DurabilityService.SyncLocal.Self, DurabilityService.SyncRemote.Quorum, 10L, TimeUnit.MINUTES)); accord.journal().closeCurrentSegmentForTestingIfNonEmpty(); accord.journal().runCompactorForTesting();