diff --git a/accord-core/src/main/java/accord/local/CommandStore.java b/accord-core/src/main/java/accord/local/CommandStore.java index d321bc6794..e656989734 100644 --- a/accord-core/src/main/java/accord/local/CommandStore.java +++ b/accord-core/src/main/java/accord/local/CommandStore.java @@ -1219,7 +1219,7 @@ private static ImmutableSortedMap purgeAndInser return ImmutableSortedMap.copyOf(build); } - private static ImmutableSortedMap purgeHistory(NavigableMap in, Ranges remove) + protected static ImmutableSortedMap purgeHistory(NavigableMap in, Ranges remove) { return ImmutableSortedMap.copyOf(purgeHistoryIterator(in, remove)); } diff --git a/accord-core/src/main/java/accord/local/CommandStores.java b/accord-core/src/main/java/accord/local/CommandStores.java index 9a6e2d29af..ea68113029 100644 --- a/accord-core/src/main/java/accord/local/CommandStores.java +++ b/accord-core/src/main/java/accord/local/CommandStores.java @@ -25,6 +25,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.NavigableMap; import java.util.Objects; import java.util.function.BiConsumer; import java.util.function.BiFunction; @@ -82,6 +83,7 @@ import org.agrona.collections.Int2IntHashMap; import org.agrona.collections.Int2ObjectHashMap; +import static accord.local.CommandStore.purgeHistory; import static accord.topology.EpochReady.done; import static accord.api.DataStore.FetchKind.Sync; import static accord.local.CommandStores.BootstrapRangeAction.BOOTSTRAP_NOT_NEEDED; @@ -1121,4 +1123,22 @@ protected Snapshot current() { return current; } + + public AsyncResult> getInUseRangesAndMarkRetiredRangesUnsafeToRead() + { + List> results = new ArrayList<>(); + Snapshot snapshot = current; + for (ShardHolder shard : snapshot.shards) + results.add(shard.store.submit((PreLoadContext.Empty) () -> "Get not retired ranges and mark retired ranges unsafe to read", + safeCommandStore -> { + Ranges notRetiredRanges = shard.ranges().notRetired(safeCommandStore); + Ranges retired = shard.ranges().all().without(notRetiredRanges); + NavigableMap safeToReadAt = safeCommandStore.safeToReadAt(); + if (safeToReadAt.values().stream().anyMatch(r -> r.intersects(retired))) + safeCommandStore.setSafeToRead(purgeHistory(safeToReadAt, retired)); + return notRetiredRanges; + })); + + return AsyncResults.allOf(results); + } }