diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 675ce49..6364365 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index ac61679..8fa0eb7 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -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 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()); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index e2dd5ec..dbbf859 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -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)); }