#150: report lead readiness from delivery gate #178
@@ -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}
|
||||
|
||||
Reference in New Issue
Block a user