CB-307 Increment 3: bridge_ack tool — per-msgId ack refinement

- MessageService.ackReply(target, msgId) delegates to inbox.ack
- BridgeMcp registers bridge_ack tool with target/msgId args
- Tests: valid/invalid args, ack surface via BridgeMcp
This commit is contained in:
Dai Ha
2026-07-19 10:21:06 +02:00
parent f756933879
commit 7c252b5f5f
3 changed files with 64 additions and 0 deletions
@@ -109,6 +109,11 @@ public final class BridgeMcp {
Map<String, Object> 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<String, Object> 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 "
@@ -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
@@ -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");