From 548c688df29998c696830cfd6d1cfb8cbb0ee48a Mon Sep 17 00:00:00 2001 From: ameya-signal <163888879+ameya-signal@users.noreply.github.com> Date: Tue, 21 Jul 2026 07:28:07 -0700 Subject: [PATCH] Add counter for FoundationDB watch operations --- .../FoundationDbMessagePublisher.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) 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 2b21aef8e..ffd829550 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 @@ -21,9 +21,11 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import javax.annotation.Nullable; +import io.micrometer.core.instrument.Metrics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.whispersystems.textsecuregcm.entities.MessageProtos; +import org.whispersystems.textsecuregcm.metrics.MetricsUtil; import reactor.core.publisher.Flux; import reactor.core.publisher.FluxSink; import reactor.core.publisher.Mono; @@ -88,6 +90,9 @@ class FoundationDbMessagePublisher { private static final Logger LOGGER = LoggerFactory.getLogger(FoundationDbMessagePublisher.class); + private static final String WATCHES_COUNTER = MetricsUtil.name(FoundationDbMessagePublisher.class, "watches"); + private static final String ACTION_TAG = "action"; + enum State { /// Messages are likely available in the queue. Initial state. MESSAGES_AVAILABLE, @@ -376,14 +381,21 @@ class FoundationDbMessagePublisher { // again (if there is demand). When we run out of messages, this method will be called again, setting another watch, // and so on, thus achieving a "watch for new messages -> read -> publish" loop. watchFuture = transaction.watch(messagesAvailableWatchKey); - watchFuture.thenRun(() -> transitionStateOnEvent(Event.MESSAGE_AVAILABLE_WATCH_TRIGGERED)); + Metrics.counter(WATCHES_COUNTER, ACTION_TAG, "set").increment(); + watchFuture.thenRun(() -> { + Metrics.counter(WATCHES_COUNTER, ACTION_TAG, "triggered").increment(); + transitionStateOnEvent(Event.MESSAGE_AVAILABLE_WATCH_TRIGGERED); + }); } /// Cancel the watch (if any). Although inactive watches are automatically timed-out, we explicitly cancel when the /// subscription is disposed to clean up associated resources and to avoid having too many watches open at once. private synchronized void cancelWatch() { if (watchFuture != null) { - watchFuture.cancel(true); + if (watchFuture.cancel(true)) { + Metrics.counter(WATCHES_COUNTER, ACTION_TAG, "cancelled").increment(); + } + watchFuture = null; } }