#150: report lead readiness from delivery gate #178

Merged
ltms merged 1 commits from worker/fleetd-150-lead-ready-fc5d35-6 into main 2026-08-28 00:59:44 +02:00
3 changed files with 43 additions and 14 deletions
@@ -380,7 +380,8 @@ public final class Fleetd {
sessions.onTurnFailed(target);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
Predicate<String> 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);
@@ -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<String> 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<String> 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
* <em>readiness</em> (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.
* <em>readiness</em>: {@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<String, Object> 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.
@@ -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<String> allow, String placement, Worktrees worktrees) {
return start(herdr, workerBaseUrl, allow, placement, worktrees, ignored -> false);
}
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> 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<String, String> 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}