From c5e24197bfaf1c83493a4282be7ad289ec763a73 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 28 Aug 2026 05:54:40 +0700 Subject: [PATCH] #150: report lead readiness from delivery gate --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 5 +-- .../java/dev/ltms/fleet/rest/FleetApp.java | 33 ++++++++++++------- .../dev/ltms/fleet/rest/FleetAppTest.java | 19 ++++++++++- 3 files changed, 43 insertions(+), 14 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 9b66efc..1c265eb 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -380,7 +380,8 @@ public final class Fleetd { sessions.onTurnFailed(target); } }; - Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads), + Predicate deliverable = deliverableTo(presence, leads); + Injector injector = new Injector(agents, turnListener, deliverable, presence::forget); StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS); poller.start(); @@ -592,7 +593,7 @@ public final class Fleetd { })); Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(), - callers, metrics).build(); + callers, metrics, deliverable).build(); app.start(cfg.bind().host(), cfg.bind().port()); log.info("fleetd listening on {}:{}, herdr socket {}", cfg.bind().host(), cfg.bind().port(), socket); diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 0b2d247..e634390 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -31,6 +31,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.function.Function; +import java.util.function.Predicate; import java.util.stream.Collectors; /** @@ -58,7 +59,7 @@ public final class FleetApp { private final PeerLauncher workers; private final SessionManager sessions; // CB-301: authoritative session registry private final MessageService messages; - private final MemberPresence presence; // CB-113: which workers are MCP-connected (available) + private final Predicate deliverable; private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable) private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour) private final Metrics metrics; // CB-502: null → /metrics not exposed @@ -70,9 +71,9 @@ public final class FleetApp { * behaviour without each needing an auth fixture. */ public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions, - MessageService messages, MemberPresence presence, - HttpServlet mcpServlet) { - this(herdr, workers, sessions, messages, presence, mcpServlet, null, null); + MessageService messages, MemberPresence presence, + HttpServlet mcpServlet) { + this(herdr, workers, sessions, messages, presence, mcpServlet, null, null, presence::isPresent); } /** @@ -82,13 +83,23 @@ public final class FleetApp { * the endpoint */ public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions, - MessageService messages, MemberPresence presence, - HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) { + MessageService messages, MemberPresence presence, + HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) { + this(herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, presence::isPresent); + } + + /** + * @param deliverable the injector's readiness gate, shared so status reports its real result + */ + public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions, + MessageService messages, MemberPresence presence, + HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, + Predicate deliverable) { this.herdr = herdr; this.workers = workers; this.sessions = sessions; this.messages = messages; - this.presence = presence; + this.deliverable = deliverable; this.mcpServlet = mcpServlet; this.auth = auth; this.metrics = metrics; @@ -507,9 +518,9 @@ public final class FleetApp { /** * Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its - * readiness (CB-113): {@code ready} is true once the worker's Claude has connected the - * bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which - * is also true during boot. + * readiness: {@code ready} is true when the injector can deliver to the target. A + * spawned member must connect the bridge MCP first, while a known lead is ready without member + * presence. This differs from bare {@code idle}, which is also true during member boot. */ private void sessionStatus(Context ctx) { String id = ctx.pathParam("id"); @@ -520,7 +531,7 @@ public final class FleetApp { Map body = new LinkedHashMap<>(); body.put("sessionId", id); body.put("status", messages.status(id).name().toLowerCase()); - body.put("ready", presence.isPresent(id)); + body.put("ready", deliverable.test(id)); // CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a // status poll — surface the open question and how to answer it, same as fleet_poll's // Phase.ASKING view. diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java index d0d098b..aef96bf 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java @@ -32,6 +32,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; +import java.util.function.Predicate; import static org.junit.jupiter.api.Assertions.*; @@ -64,6 +65,11 @@ class FleetAppTest { } private int start(FakeHerdr herdr, String workerBaseUrl, Set allow, String placement, Worktrees worktrees) { + return start(herdr, workerBaseUrl, allow, placement, worktrees, ignored -> false); + } + + private int start(FakeHerdr herdr, String workerBaseUrl, Set allow, String placement, + Worktrees worktrees, Predicate deliverable) { FleetConfig.Profile wcfg = new FleetConfig.Profile( "ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null, placement, "fleet", "worker: {profile} #{n}", null, null, null); @@ -84,7 +90,8 @@ class FleetAppTest { // it directly so the inbox contract holds for those endpoints. inbox.own("term_a"); MessageService messages = new MessageService(agents, injector, rendezvous, inbox); - app = new FleetApp(herdr, workers, sessions, messages, this.presence, null) + app = new FleetApp(herdr, workers, sessions, messages, this.presence, null, + null, null, id -> this.presence.isPresent(id) || deliverable.test(id)) .build().start("127.0.0.1", 0); return app.port(); } @@ -484,6 +491,16 @@ class FleetAppTest { assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean()); } + @Test + void sessionStatusReportsRegisteredLeadAsReady() throws Exception { + Map leads = Map.of("term_lead", "terra"); + int port = start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", + new GitWorktrees(), leads::containsKey); + + JsonNode body = mapper.readTree(req(port, "GET", "/sessions/term_lead/status").body()); + assertTrue(body.get("ready").asBoolean()); + } + /** * CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the * ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask} -- 2.52.0