Instrument more FoundationDB calls

This commit is contained in:
Katherine authored and GitHub committed 2026-09-10 13:46:00 -04:00
1 parent 52b32ce2d4
commit d366cdbcb0
12 files changed
+110 -79

No files matched your search

@@ -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> T run(final Function<? super Transaction, T> retryable) {
return circuitBreaker.executeSupplier(() -> database.run(retryable));
public <T> T run(final Function<? super Transaction, T> 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 <T> the return type of retryable
/// @return a cancellation-safe version of the future returned from [Database#runAsync(Function)]
public <T> CompletableFuture<T> runAsync(final Function<? super Transaction, ? extends CompletableFuture<T>> 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> T read(final Function<? super ReadTransaction, T> retryable) {
return circuitBreaker.executeSupplier(() -> database.read(retryable));
public <T> T read(final Function<? super ReadTransaction, T> 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 <T> CompletableFuture<T> readAsync(
final Function<? super ReadTransaction, ? extends CompletableFuture<T>> retryable) {
return circuitBreaker.executeCompletionStage(() -> database.readAsync(retryable))
final Function<? super ReadTransaction, ? extends CompletableFuture<T>> 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;
}
}
}
@@ -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) {
@@ -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);
@@ -304,7 +304,7 @@ public class FoundationDbMessageStore {
}
return CompletableFuture.completedFuture(Optional.<Versionstamp>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<FaultTolerantDatabase> 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<KeyArrayResult> 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<Void> 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
@@ -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<Optional<KeySelector>> 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)
.<Optional<KeySelector>>thenApply(items -> {
if (items.isEmpty()) {
return Optional.empty();
@@ -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");
}
@@ -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) {
@@ -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()))
);
}
@@ -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));
}
}
@@ -576,7 +576,8 @@ class FoundationDbMessageStoreTest {
final List<KeyValue> 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<Versionstamp> 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<KeyValue> 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();
}
}
@@ -30,7 +30,7 @@ class VersionstampClockTest {
FOUNDATION_DB_EXTENSION.getDatabases()[0].run(transaction -> {
transaction.clear(VersionstampClock.SUBSPACE.range());
return null;
});
}, FaultTolerantDatabase.Context.TEST);
}
@Test
@@ -136,7 +136,7 @@ class ClearOrphanedFoundationDbQueuesCommandTest {
new byte[]{43});
});
return null;
});
}, FaultTolerantDatabase.Context.TEST);
final List<AciServiceIdentifier> 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<KeyValue> 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;