From d366cdbcb0d13fff67cbe4bf10a30c2f1e54d226 Mon Sep 17 00:00:00 2001 From: Katherine Date: Thu, 10 Sep 2026 13:46:00 -0400 Subject: [PATCH] Instrument more FoundationDB calls --- .../foundationdb/FaultTolerantDatabase.java | 80 ++++++++++++++++--- .../foundationdb/FoundationDBWarmup.java | 2 +- .../FoundationDbMessagePublisher.java | 6 +- .../FoundationDbMessageStore.java | 18 ++--- .../FoundationDbMessageStream.java | 2 +- .../foundationdb/FoundationDbUtil.java | 32 -------- .../foundationdb/VersionstampClock.java | 6 +- ...learOrphanedFoundationDbQueuesCommand.java | 7 +- .../FaultTolerantDatabaseTest.java | 10 +-- .../FoundationDbMessageStoreTest.java | 20 ++--- .../foundationdb/VersionstampClockTest.java | 2 +- ...OrphanedFoundationDbQueuesCommandTest.java | 4 +- 12 files changed, 110 insertions(+), 79 deletions(-) delete mode 100644 service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbUtil.java diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabase.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabase.java index 7b25877d7..3c19fa34d 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabase.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabase.java @@ -5,8 +5,6 @@ package org.whispersystems.textsecuregcm.storage.foundationdb; -import static org.whispersystems.textsecuregcm.storage.foundationdb.FoundationDbUtil.TRANSACTION_ERRORS_COUNTER; - import com.apple.foundationdb.Database; import com.apple.foundationdb.FDBException; import com.apple.foundationdb.ReadTransaction; @@ -16,10 +14,12 @@ import io.micrometer.core.instrument.Metrics; import jakarta.annotation.Nullable; import java.util.concurrent.CompletableFuture; import java.util.function.Function; +import org.whispersystems.textsecuregcm.metrics.MetricsUtil; import org.whispersystems.textsecuregcm.util.ExceptionUtils; import org.whispersystems.textsecuregcm.util.ResilienceUtil; public class FaultTolerantDatabase { + private static final String TRANSACTION_ERRORS_COUNTER = MetricsUtil.name(FaultTolerantDatabase.class, "transactionErrors"); private final Database database; private final CircuitBreaker circuitBreaker; @@ -33,8 +33,18 @@ public class FaultTolerantDatabase { : ResilienceUtil.getCircuitBreakerRegistry().circuitBreaker(circuitBreakerName); } - public T run(final Function retryable) { - return circuitBreaker.executeSupplier(() -> database.run(retryable)); + public T run(final Function retryable, final Context context) { + try { + return circuitBreaker.executeSupplier(() -> database.run(retryable)); + } catch (final Exception e) { + if (e instanceof final FDBException fdbException) { + Metrics.counter(TRANSACTION_ERRORS_COUNTER, + "context", context.getName(), + "code", String.valueOf(fdbException.getCode()) + ).increment(); + } + throw e; + } } /// Returns a cancellation-safe version of the result from [Database#runAsync(Function)]. Since the final stage @@ -46,7 +56,7 @@ public class FaultTolerantDatabase { /// @param the return type of retryable /// @return a cancellation-safe version of the future returned from [Database#runAsync(Function)] public CompletableFuture runAsync(final Function> retryable, - final FoundationDbUtil.Context context) { + final Context context) { return circuitBreaker.executeCompletionStage(() -> database.runAsync(retryable) .whenComplete((_, throwable) -> { if (throwable != null && ExceptionUtils.unwrap(throwable) instanceof final FDBException fdbException) { @@ -60,13 +70,65 @@ public class FaultTolerantDatabase { .toCompletableFuture(); } - public T read(final Function retryable) { - return circuitBreaker.executeSupplier(() -> database.read(retryable)); + public T read(final Function retryable, final Context context) { + try { + return circuitBreaker.executeSupplier(() -> database.read(retryable)); + } catch (final Exception e) { + if (e instanceof final FDBException fdbException) { + Metrics.counter(TRANSACTION_ERRORS_COUNTER, + "context", context.getName(), + "code", String.valueOf(fdbException.getCode()) + ).increment(); + } + throw e; + } } public CompletableFuture readAsync( - final Function> retryable) { - return circuitBreaker.executeCompletionStage(() -> database.readAsync(retryable)) + final Function> retryable, + final Context context) { + return circuitBreaker.executeCompletionStage(() -> database.readAsync(retryable) + .whenComplete((_, throwable) -> { + if (throwable != null && ExceptionUtils.unwrap(throwable) instanceof final FDBException fdbException) { + Metrics.counter(TRANSACTION_ERRORS_COUNTER, + "context", context.getName(), + "code", String.valueOf(fdbException.getCode()) + ).increment(); + } + })) .toCompletableFuture(); } + + public enum Context { + INSERT_MESSAGE_BATCH("insertMessageBatch"), + GET_MESSAGES_BATCH("getMessagesBatch"), + SET_PRESENCE("setPresence"), + GET_PRESENCE("getPresence"), + CLEAR_PRESENCE("clearPresence"), + CLEAR_ACCOUNT_SUBSPACE("clearAccountSubspace"), + CLEAR_DEVICE_SUBSPACE("clearDeviceSubspace"), + CLEAR_EXPIRED_MESSAGES("clearExpiredMessages"), + CLEAR_EXPIRED_VERSIONSTAMPS("clearExpiredVersionstamps"), + GET_END_OF_QUEUE("getEndOfQueue"), + ESTIMATE_QUEUE_SIZE("estimateQueueSize"), + ESTIMATE_QUEUE_SIZE_AND_RANGE_SPLITS("estimateQueueSizeAndRangeSplits"), + GET_RANGE_SPLITS("getRangeSplits"), + TRIM_QUEUE("trimQueue"), + DELETE_MESSAGE("deleteMessage"), + READ_ACIS("readAcis"), + READ_STATUS("readStatus"), + READ_VERSIONSTAMP("readVersionstamp"), + RECORD_VERSIONSTAMP_AND_TIME("recordVersionstampAndTime"), + TEST("test"); + + private final String name; + + Context(final String name) { + this.name = name; + } + + public String getName() { + return name; + } + } } diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDBWarmup.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDBWarmup.java index eef1d411a..4979c7ae5 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDBWarmup.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDBWarmup.java @@ -51,7 +51,7 @@ public class FoundationDBWarmup implements Managed { final Timer.Sample sample = Timer.start(); while (true) { try { - database.readAsync(transaction -> transaction.get(STATUS_JSON_KEY)).join(); + database.readAsync(transaction -> transaction.get(STATUS_JSON_KEY), FaultTolerantDatabase.Context.READ_STATUS).join(); sample.stop(WARMUP_TIMER); return; } catch (final Exception e) { diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessagePublisher.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessagePublisher.java index 4071be4c4..037c28c0d 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessagePublisher.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessagePublisher.java @@ -338,7 +338,7 @@ class FoundationDbMessagePublisher { return keyValues; }); - }, FoundationDbUtil.Context.GET_MESSAGES_BATCH) + }, FaultTolerantDatabase.Context.GET_MESSAGES_BATCH) .thenApply(keyValues -> { if (keyValues.size() < maxMessages) { transitionStateOnEvent(Event.FETCHED_ALL_AVAILABLE_MESSAGES); @@ -491,7 +491,7 @@ class FoundationDbMessagePublisher { return database.runAsync(transaction -> { transaction.set(presenceKey, FoundationDbMessageStore.getPresenceValue(clock.instant(), streamId)); return CompletableFuture.completedFuture(null); - }, FoundationDbUtil.Context.SET_PRESENCE); + }, FaultTolerantDatabase.Context.SET_PRESENCE); } @VisibleForTesting @@ -504,7 +504,7 @@ class FoundationDbMessagePublisher { if (!isPresenceContested(presenceValue)) { transaction.clear(presenceKey); } - }), FoundationDbUtil.Context.CLEAR_PRESENCE) + }), FaultTolerantDatabase.Context.CLEAR_PRESENCE) .whenComplete((_, throwable) -> { if (throwable != null) { LOGGER.warn("Failed to clear presence on disposal", throwable); diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStore.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStore.java index 7ab2034c5..46cfbaeb3 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStore.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStore.java @@ -304,7 +304,7 @@ public class FoundationDbMessageStore { } return CompletableFuture.completedFuture(Optional.empty()); }); - }, FoundationDbUtil.Context.INSERT_MESSAGE_BATCH) + }, FaultTolerantDatabase.Context.INSERT_MESSAGE_BATCH) .thenCompose(Function.identity()) .thenApply(maybeVersionstamp -> insertFuturesByAci.entrySet().stream() .collect(Collectors.toMap(Map.Entry::getKey, entry -> { @@ -395,7 +395,7 @@ public class FoundationDbMessageStore { return databasesByEpoch[getConfigurationEpoch(versionstamp)][getShardId(versionstamp)].runAsync(transaction -> { transaction.clear(messageKey); return CompletableFuture.completedFuture(null); - }, FoundationDbUtil.Context.DELETE_MESSAGE) + }, FaultTolerantDatabase.Context.DELETE_MESSAGE) .thenRun(() -> { sample.stop(DELETE_MESSAGE_TIMER); DELETE_MESSAGE_COUNTER.increment(); @@ -419,7 +419,7 @@ public class FoundationDbMessageStore { } transaction.clear(messageKey); return Optional.of(value); - }), FoundationDbUtil.Context.DELETE_MESSAGE) + }), FaultTolerantDatabase.Context.DELETE_MESSAGE) .whenComplete((_, _) -> sample.stop(DELETE_MESSAGE_TIMER)) .thenApply(maybeValue -> maybeValue.map(value -> { DELETE_MESSAGE_COUNTER.increment(); @@ -435,14 +435,14 @@ public class FoundationDbMessageStore { doForAllDatabasesWithMessages(aci, database -> database.run(transaction -> { transaction.clear(getAccountSubspace(aci).range()); return null; - })); + }, FaultTolerantDatabase.Context.CLEAR_ACCOUNT_SUBSPACE)); } public void clearAll(final AciServiceIdentifier aci, final byte deviceId) { doForAllDatabasesWithMessages(aci, database -> database.run(transaction -> { transaction.clear(getDeviceSubspace(aci, deviceId).range()); return null; - })); + }, FaultTolerantDatabase.Context.CLEAR_DEVICE_SUBSPACE)); } private void doForAllDatabasesWithMessages(final AciServiceIdentifier aci, final Consumer action) { @@ -545,7 +545,7 @@ public class FoundationDbMessageStore { new Range(queueSubspace.getKey(), queueSubspace.pack(Tuple.from(cutoffVersionstamp)))); } return null; - })); + }, FaultTolerantDatabase.Context.CLEAR_EXPIRED_MESSAGES)); }); } @@ -619,7 +619,7 @@ public class FoundationDbMessageStore { final Range deviceQueueRange = getDeviceQueueSubspace(aci, deviceId).range(); return getDistinctDatabasesForAci(aci) - .flatMap(database -> Mono.fromFuture(() -> database.runAsync(transaction -> transaction.getEstimatedRangeSizeBytes(deviceQueueRange), FoundationDbUtil.Context.ESTIMATE_QUEUE_SIZE))) + .flatMap(database -> Mono.fromFuture(() -> database.runAsync(transaction -> transaction.getEstimatedRangeSizeBytes(deviceQueueRange), FaultTolerantDatabase.Context.ESTIMATE_QUEUE_SIZE))) .reduce(0L, Long::sum); } @@ -654,7 +654,7 @@ public class FoundationDbMessageStore { final CompletableFuture rangeSplitPointsFuture = transaction.getRangeSplitPoints(deviceQueueRange, rangeSplitChunkSize); return estimatedQueueSizeFuture.thenCombine(rangeSplitPointsFuture, Pair::new); - }, FoundationDbUtil.Context.ESTIMATE_QUEUE_SIZE_AND_RANGE_SPLITS); + }, FaultTolerantDatabase.Context.ESTIMATE_QUEUE_SIZE_AND_RANGE_SPLITS); } public Mono trimQueue(final AciServiceIdentifier aci, @@ -743,7 +743,7 @@ public class FoundationDbMessageStore { transaction.options().setRetryLimit(batchPriorityTransactionRetryLimit); transaction.clear(range); return CompletableFuture.completedFuture(null); - }, FoundationDbUtil.Context.TRIM_QUEUE); + }, FaultTolerantDatabase.Context.TRIM_QUEUE); } @VisibleForTesting diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStream.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStream.java index 07ead87c2..bae601829 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStream.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStream.java @@ -205,7 +205,7 @@ public class FoundationDbMessageStream implements MessageStream { /// @return a [KeySelector] for the first key greater than the current greatest key in the device queue. private CompletableFuture> getEndOfQueueKeyExclusive(final FaultTolerantDatabase database) { return database.runAsync(transaction -> - transaction.getRange(deviceQueueSubspace.range(), 1, true, StreamingMode.EXACT).asList(), FoundationDbUtil.Context.GET_END_OF_QUEUE) + transaction.getRange(deviceQueueSubspace.range(), 1, true, StreamingMode.EXACT).asList(), FaultTolerantDatabase.Context.GET_END_OF_QUEUE) .>thenApply(items -> { if (items.isEmpty()) { return Optional.empty(); diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbUtil.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbUtil.java deleted file mode 100644 index f6b3ade5c..000000000 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbUtil.java +++ /dev/null @@ -1,32 +0,0 @@ -package org.whispersystems.textsecuregcm.storage.foundationdb; - -import org.whispersystems.textsecuregcm.metrics.MetricsUtil; - -public class FoundationDbUtil { - - public enum Context { - INSERT_MESSAGE_BATCH("insertMessageBatch"), - GET_MESSAGES_BATCH("getMessagesBatch"), - SET_PRESENCE("setPresence"), - GET_PRESENCE("getPresence"), - CLEAR_PRESENCE("clearPresence"), - GET_END_OF_QUEUE("getEndOfQueue"), - ESTIMATE_QUEUE_SIZE("estimateQueueSize"), - ESTIMATE_QUEUE_SIZE_AND_RANGE_SPLITS("estimateQueueSizeAndRangeSplits"), - GET_RANGE_SPLITS("getRangeSplits"), - TRIM_QUEUE("trimQueue"), - DELETE_MESSAGE("deleteMessage"); - - private final String name; - - Context(final String name) { - this.name = name; - } - - public String getName() { - return name; - } - } - - static final String TRANSACTION_ERRORS_COUNTER = MetricsUtil.name(FoundationDbUtil.class, "transactionErrors"); -} diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClock.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClock.java index d4cb008fe..5c11a7569 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClock.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClock.java @@ -53,7 +53,7 @@ public class VersionstampClock { Tuple.from(Versionstamp.incomplete()).packWithVersionstamp()); return transaction.getVersionstamp(); - }) + }, FaultTolerantDatabase.Context.RECORD_VERSIONSTAMP_AND_TIME) .thenApply(Versionstamp::complete) .join(); } @@ -76,7 +76,7 @@ public class VersionstampClock { } return Optional.empty(); - }); + }, FaultTolerantDatabase.Context.READ_VERSIONSTAMP); } /// Remove any entries from the versionstamp-clock namespace that are strictly older than the given timestamp @@ -86,7 +86,7 @@ public class VersionstampClock { database.run(transaction -> { transaction.clear(SUBSPACE.getKey(), getTimestampKey(oldestRetainedEntryTimestamp)); return null; - }); + }, FaultTolerantDatabase.Context.CLEAR_EXPIRED_VERSIONSTAMPS); } private byte[] getTimestampKey(final Instant timestamp) { diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommand.java b/service/src/main/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommand.java index b1178fe09..0da91cd14 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommand.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommand.java @@ -42,7 +42,6 @@ import org.whispersystems.textsecuregcm.storage.AccountLockManager; import org.whispersystems.textsecuregcm.storage.AccountsManager; import org.whispersystems.textsecuregcm.storage.foundationdb.FaultTolerantDatabase; import org.whispersystems.textsecuregcm.storage.foundationdb.FoundationDbMessageStore; -import org.whispersystems.textsecuregcm.storage.foundationdb.FoundationDbUtil; import org.whispersystems.textsecuregcm.util.ManagedExecutors; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -224,7 +223,7 @@ public class ClearOrphanedFoundationDbQueuesCommand extends AbstractCommandWithD transaction.options().setTimeout(transactionTimeout.toMillis()); transaction.clear(getAccountSubspace(messagesSubspace, aci).range()); return null; - }); + }, FaultTolerantDatabase.Context.CLEAR_ACCOUNT_SUBSPACE); } @VisibleForTesting @@ -247,7 +246,7 @@ public class ClearOrphanedFoundationDbQueuesCommand extends AbstractCommandWithD .thenCompose(rangeSize -> { final long chunkSize = Math.ceilDiv(rangeSize, numChunks); return transaction.getRangeSplitPoints(messagesSubspace.range(), chunkSize); - }), FoundationDbUtil.Context.GET_RANGE_SPLITS) + }), FaultTolerantDatabase.Context.GET_RANGE_SPLITS) .thenApply(result -> splitPointsToRanges(result.getKeys())); } @@ -314,7 +313,7 @@ public class ClearOrphanedFoundationDbQueuesCommand extends AbstractCommandWithD }); }) .thenApply(_ -> acis); - }) + }, FaultTolerantDatabase.Context.READ_ACIS) .thenApply(acis -> new BatchReadResult(acis, cursor.get())) ); } diff --git a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabaseTest.java b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabaseTest.java index 01e19aa33..70e7ea5d2 100644 --- a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabaseTest.java +++ b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FaultTolerantDatabaseTest.java @@ -47,20 +47,20 @@ public class FaultTolerantDatabaseTest { when(database.read(any())).thenThrow(new FDBException("error", 1031)); assertThrows(FDBException.class, () -> ftDatabase.read(t -> - t.get("test_key".getBytes()))); + t.get("test_key".getBytes()), FaultTolerantDatabase.Context.TEST)); assertThrows(FDBException.class, () -> ftDatabase.read(t -> - t.get("test_key".getBytes()))); + t.get("test_key".getBytes()), FaultTolerantDatabase.Context.TEST)); assertThrows(CallNotPermittedException.class, () -> ftDatabase.read(t -> - t.get("test_key".getBytes()))); + t.get("test_key".getBytes()), FaultTolerantDatabase.Context.TEST)); Thread.sleep(1001); assertThrows(FDBException.class, () -> ftDatabase.read(t -> - t.get("test_key".getBytes()))); + t.get("test_key".getBytes()), FaultTolerantDatabase.Context.TEST)); assertThrows(CallNotPermittedException.class, () -> ftDatabase.read(t -> - t.get("test_key".getBytes()))); + t.get("test_key".getBytes()), FaultTolerantDatabase.Context.TEST)); } } diff --git a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStoreTest.java b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStoreTest.java index 5f9717593..27466d103 100644 --- a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStoreTest.java +++ b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/FoundationDbMessageStoreTest.java @@ -576,7 +576,8 @@ class FoundationDbMessageStoreTest { final List keyValues = foundationDbMessageStore.getShardForAci(deletedAccountIdentifier, DEFAULT_EPOCH) .readAsync(transaction -> AsyncUtil.collect( - transaction.getRange(deletedAccountSubspace.range().begin, deletedAccountSubspace.range().end, 1))) + transaction.getRange(deletedAccountSubspace.range().begin, deletedAccountSubspace.range().end, 1)), + FaultTolerantDatabase.Context.TEST) .join(); assertEquals(0, keyValues.size()); @@ -840,7 +841,7 @@ class FoundationDbMessageStoreTest { while (!presenceCleared && Instant.now().isBefore(deadline)) { presenceCleared = foundationDbMessageStore.getShardForAci(aci, DEFAULT_EPOCH).runAsync(transaction -> - transaction.get(FoundationDbMessageStore.getPresenceKey(aci, Device.PRIMARY_ID)), FoundationDbUtil.Context.GET_PRESENCE) + transaction.get(FoundationDbMessageStore.getPresenceKey(aci, Device.PRIMARY_ID)), FaultTolerantDatabase.Context.TEST) .thenApply(Objects::isNull) .join(); } @@ -894,7 +895,7 @@ class FoundationDbMessageStoreTest { foundationDbMessageStore.getShardForAci(aci, DEFAULT_EPOCH).run(transaction -> { transaction.set(FoundationDbMessageStore.getPresenceKey(aci, Device.PRIMARY_ID), presenceValue); return null; - }); + }, FaultTolerantDatabase.Context.TEST); }) .thenCancel() .verify(); @@ -905,7 +906,7 @@ class FoundationDbMessageStoreTest { while (!presenceCleared && Instant.now().isBefore(deadline)) { presenceCleared = foundationDbMessageStore.getShardForAci(aci, DEFAULT_EPOCH).runAsync(transaction -> transaction.get(FoundationDbMessageStore.getPresenceKey(aci, Device.PRIMARY_ID)), - FoundationDbUtil.Context.GET_PRESENCE) + FaultTolerantDatabase.Context.TEST) .thenApply(Objects::isNull) .join(); } @@ -917,7 +918,7 @@ class FoundationDbMessageStoreTest { private boolean isPresent(final AciServiceIdentifier aci, final byte deviceId) { return foundationDbMessageStore.getShardForAci(aci, DEFAULT_EPOCH).run(transaction -> transaction.get(FoundationDbMessageStore.getPresenceKey(aci, deviceId)) - .thenApply(foundationDbMessageStore::isClientPresent)) + .thenApply(foundationDbMessageStore::isClientPresent), FaultTolerantDatabase.Context.TEST) .join(); } @@ -1309,7 +1310,7 @@ class FoundationDbMessageStoreTest { final byte[] key = FoundationDbMessageStore.getDeviceQueueSubspace(aci, deviceId) .pack(Tuple.from(versionstamp)); return transaction.get(key); - }).join(); + }, FaultTolerantDatabase.Context.TEST).join(); } private Optional getMessagesAvailableWatch(final AciServiceIdentifier aci) { @@ -1320,7 +1321,7 @@ class FoundationDbMessageStoreTest { return foundationDbMessageStore.getShardForAci(aci, epoch) .read(transaction -> transaction.get(FoundationDbMessageStore.getMessagesAvailableWatchKey(aci)) .thenApply(value -> value == null ? null : Tuple.fromBytes(value).getVersionstamp(0)) - .thenApply(Optional::ofNullable)) + .thenApply(Optional::ofNullable), FaultTolerantDatabase.Context.TEST) .join(); } @@ -1333,7 +1334,7 @@ class FoundationDbMessageStoreTest { final byte[] presenceValue = FoundationDbMessageStore.getPresenceValue(CLOCK.instant().minusSeconds(secondsBeforeCurrentTime), STREAM_ID); transaction.set(presenceKey, presenceValue); return null; - }); + }, FaultTolerantDatabase.Context.TEST); } private AciServiceIdentifier generateRandomAciForShard(final int shardNumber) { @@ -1356,7 +1357,8 @@ class FoundationDbMessageStoreTest { private List getItemsInDeviceQueue(final AciServiceIdentifier aci, final byte deviceId, final int epoch) { return foundationDbMessageStore.getShardForAci(aci, epoch).readAsync(transaction -> - AsyncUtil.collect(transaction.getRange(FoundationDbMessageStore.getDeviceQueueSubspace(aci, deviceId).range()))) + AsyncUtil.collect(transaction.getRange(FoundationDbMessageStore.getDeviceQueueSubspace(aci, deviceId).range())), + FaultTolerantDatabase.Context.TEST) .join(); } } diff --git a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClockTest.java b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClockTest.java index 20338d3d2..410f9909e 100644 --- a/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClockTest.java +++ b/service/src/test/java/org/whispersystems/textsecuregcm/storage/foundationdb/VersionstampClockTest.java @@ -30,7 +30,7 @@ class VersionstampClockTest { FOUNDATION_DB_EXTENSION.getDatabases()[0].run(transaction -> { transaction.clear(VersionstampClock.SUBSPACE.range()); return null; - }); + }, FaultTolerantDatabase.Context.TEST); } @Test diff --git a/service/src/test/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommandTest.java b/service/src/test/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommandTest.java index 876cfb6a4..17fd3d5e5 100644 --- a/service/src/test/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommandTest.java +++ b/service/src/test/java/org/whispersystems/textsecuregcm/workers/ClearOrphanedFoundationDbQueuesCommandTest.java @@ -136,7 +136,7 @@ class ClearOrphanedFoundationDbQueuesCommandTest { new byte[]{43}); }); return null; - }); + }, FaultTolerantDatabase.Context.TEST); final List fetchedAcis = command.getAcisInShard(database, 2, 3, Duration.ofSeconds(2), 5) .collectList() .block(); @@ -173,7 +173,7 @@ class ClearOrphanedFoundationDbQueuesCommandTest { for (final FaultTolerantDatabase database : FOUNDATION_DB_EXTENSION.getDatabases()) { final List keyValues = database.readAsync(transaction -> - AsyncUtil.collect(transaction.getRange(accountRange, 1))).join(); + AsyncUtil.collect(transaction.getRange(accountRange, 1)), FaultTolerantDatabase.Context.TEST).join(); if (!keyValues.isEmpty()) { return true;