CB-512: wire bridged_push_nudges_total metric increments

This commit is contained in:
2026-08-01 22:56:32 +07:00
parent 22ad24db6c
commit 37a11cd168
3 changed files with 64 additions and 4 deletions
@@ -200,11 +200,12 @@ public final class Bridged {
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-push-").unstarted(r));
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
pushScheduler, maxReminders, backoffMs);
// CB-502: the registry is built before the service so send/reply outcomes are counted at
// their single funnel rather than at each of the two caller-facing surfaces.
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
Metrics metrics = BridgedMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
@@ -3,6 +3,8 @@ package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -33,6 +35,7 @@ public final class ReplyPushLoop {
private final ScheduledExecutorService scheduler;
private final int maxReminders;
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** Track targets that have an active schedule. */
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
@@ -40,12 +43,27 @@ public final class ReplyPushLoop {
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs) {
this(primaryRegistry, agents, inbox, scheduler, maxReminders, backoffMs, null);
}
/** As above, with a metric registry (CB-512) so push outcomes are counted. */
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs, Metrics metrics) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.inbox = inbox;
this.scheduler = scheduler;
this.maxReminders = maxReminders;
this.backoffMs = backoffMs;
this.metrics = metrics;
}
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
private void count(String name, String... labels) {
if (metrics != null) {
metrics.inc(name, labels);
}
}
// --- decision logic (package-private for unit-testing) -------------------------------------
@@ -71,6 +89,7 @@ public final class ReplyPushLoop {
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target);
count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted");
return Action.STOP;
}
var primaryTerminal = primaryRegistry.primaryTerminal().orElseThrow();
@@ -130,6 +149,7 @@ public final class ReplyPushLoop {
agents.send(primaryTerminal, nudge);
log.debug("push: nudge {}/{} sent to primary {} for target {}",
reminderCount + 1, maxReminders, primaryTerminal, target);
count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge primary {} for target {} (reminder {}/{}): {}",
primaryTerminal, target, reminderCount + 1, maxReminders, e.toString());
@@ -5,6 +5,8 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -183,6 +185,39 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@Test
void successfulNudgeIncrementsDelivered() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
loop(1, 50, metrics).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (2 agent.send calls) should have been sent");
// The delivered count is bumped on the scheduler thread right after the send that releases
// the latch — settle briefly so the counter is published before we read it.
Thread.sleep(200);
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent nudge must count as delivered");
}
@Test
void reminderCapIncrementsExhausted() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
assertEquals(0, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered"));
}
// --- helpers -------------------------------------------------------------------------------
private ReplyPushLoop loop() {
@@ -193,6 +228,10 @@ class ReplyPushLoopTest {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs);
}
private ReplyPushLoop loop(int maxReminders, long backoffMs, Metrics metrics) {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics);
}
private static AgentControl agentWithStatus(String status) {
return new AgentControl(new FakeHerdrClient(status));
}