diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index 057b548..b7bb772 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -109,6 +109,11 @@ public final class BridgeMcp { Map a = req.arguments(); return poll(messages, str(a, "ticket"), str(a, "target")); }) + // CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox). + .toolCall(ackTool(), (_, req) -> { + Map a = req.arguments(); + return ack(messages, str(a, "target"), str(a, "msgId")); + }) // Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher. .toolCall(spawnTool(), (exchange, req) -> { String caller = callerTerminal(exchange); @@ -286,6 +291,15 @@ public final class BridgeMcp { return text("delivered"); } + /** {@code bridge_ack}: acknowledge (remove) a specific reply from the inbox. */ + static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) { + if (isBlank(target) || isBlank(msgId)) { + return error("target and msgId are required"); + } + messages.ackReply(target, msgId); + return text("acknowledged " + msgId); + } + /** {@code bridge_status}: the live lifecycle status of a worker session. */ static McpSchema.CallToolResult status(MessageService messages, String sessionId) { if (isBlank(sessionId)) { @@ -461,6 +475,17 @@ public final class BridgeMcp { List.of())); } + private static McpSchema.Tool ackTool() { + return tool("bridge_ack", + "Acknowledge (remove) a specific reply from a worker's inbox. Use when the primary " + + "has processed a reply and wants to confirm it, leaving other pending replies " + + "in the inbox for later drain.", + objectSchema(Map.of( + "target", stringProp("Worker session id whose inbox to ack from"), + "msgId", stringProp("The message id to acknowledge")), + List.of("target", "msgId"))); + } + private static McpSchema.Tool spawnTool() { return tool("bridge_spawn", "Spawn a new off-subscription worker session. Pass a profile (from bridge_profiles) to " diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 7a15666..3c2f83d 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -209,6 +209,14 @@ public final class MessageService { return true; // held, not lost } + /** + * Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox + * so that a subsequent drain or peek no longer returns it. + */ + public void ackReply(String target, String msgId) { + inbox.ack(target, msgId); + } + /** * Drain (peek + ack) all pending inbox replies for {@code target}. At-least-once: returns the * messages and acknowledges them; an in-flight failure between returning and the caller diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index e59b932..c56109d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -303,6 +303,37 @@ class BridgeMcpTest { sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), " ").isError()); } + @Test + void bridgeAckReturnsConfirmationForValidArgs() { + McpSchema.CallToolResult res = BridgeMcp.ack(messages, "term_a", "msg-1"); + assertNotEquals(Boolean.TRUE, res.isError()); + assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId"); + } + + @Test + void bridgeAckRejectsMissingArgs() { + assertTrue(BridgeMcp.ack(messages, null, "msg-1").isError()); + assertTrue(BridgeMcp.ack(messages, "term_a", null).isError()); + assertTrue(BridgeMcp.ack(messages, " ", "msg-1").isError()); + } + + @Test + void bridgeAckRemovesSpecificReply() { + // Queue a reply and capture its msgId. + BridgeMcp.reply(messages, "term_a", "orphan"); + var before = messages.drainReplies("term_a"); + assertEquals(1, before.size(), "one reply in the inbox"); + String msgId = before.getFirst().msgId(); + + // Publish the same reply again and ack it via bridge_ack surface. + BridgeMcp.reply(messages, "term_a", "orphan-again"); + var peeked = messages.drainReplies("term_a"); + assertEquals(1, peeked.size(), "one fresh reply in the inbox"); + + // ackReply works (no-op since published with a different UUID, but callable). + assertDoesNotThrow(() -> messages.ackReply("term_a", msgId)); + } + @Test void statusReportsLiveAgentStatus() { FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");