#150: report lead readiness from delivery gate #178
@@ -380,7 +380,8 @@ public final class Fleetd {
|
|||||||
sessions.onTurnFailed(target);
|
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);
|
presence::forget);
|
||||||
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||||
poller.start();
|
poller.start();
|
||||||
@@ -592,7 +593,7 @@ public final class Fleetd {
|
|||||||
}));
|
}));
|
||||||
|
|
||||||
Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
|
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());
|
app.start(cfg.bind().host(), cfg.bind().port());
|
||||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||||
cfg.bind().host(), cfg.bind().port(), socket);
|
cfg.bind().host(), cfg.bind().port(), socket);
|
||||||
|
|||||||
@@ -31,6 +31,7 @@ import java.util.LinkedHashMap;
|
|||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
import java.util.function.Predicate;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -58,7 +59,7 @@ public final class FleetApp {
|
|||||||
private final PeerLauncher workers;
|
private final PeerLauncher workers;
|
||||||
private final SessionManager sessions; // CB-301: authoritative session registry
|
private final SessionManager sessions; // CB-301: authoritative session registry
|
||||||
private final MessageService messages;
|
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 HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
||||||
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
|
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
|
||||||
private final Metrics metrics; // CB-502: null → /metrics not exposed
|
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.
|
* behaviour without each needing an auth fixture.
|
||||||
*/
|
*/
|
||||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||||
MessageService messages, MemberPresence presence,
|
MessageService messages, MemberPresence presence,
|
||||||
HttpServlet mcpServlet) {
|
HttpServlet mcpServlet) {
|
||||||
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
|
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null, presence::isPresent);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -82,13 +83,23 @@ public final class FleetApp {
|
|||||||
* the endpoint
|
* the endpoint
|
||||||
*/
|
*/
|
||||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||||
MessageService messages, MemberPresence presence,
|
MessageService messages, MemberPresence presence,
|
||||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
|
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.herdr = herdr;
|
||||||
this.workers = workers;
|
this.workers = workers;
|
||||||
this.sessions = sessions;
|
this.sessions = sessions;
|
||||||
this.messages = messages;
|
this.messages = messages;
|
||||||
this.presence = presence;
|
this.deliverable = deliverable;
|
||||||
this.mcpServlet = mcpServlet;
|
this.mcpServlet = mcpServlet;
|
||||||
this.auth = auth;
|
this.auth = auth;
|
||||||
this.metrics = metrics;
|
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
|
* 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
|
* <em>readiness</em>: {@code ready} is true when the injector can deliver to the target. A
|
||||||
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
|
* spawned member must connect the bridge MCP first, while a known lead is ready without member
|
||||||
* is also true during boot.
|
* presence. This differs from bare {@code idle}, which is also true during member boot.
|
||||||
*/
|
*/
|
||||||
private void sessionStatus(Context ctx) {
|
private void sessionStatus(Context ctx) {
|
||||||
String id = ctx.pathParam("id");
|
String id = ctx.pathParam("id");
|
||||||
@@ -520,7 +531,7 @@ public final class FleetApp {
|
|||||||
Map<String, Object> body = new LinkedHashMap<>();
|
Map<String, Object> body = new LinkedHashMap<>();
|
||||||
body.put("sessionId", id);
|
body.put("sessionId", id);
|
||||||
body.put("status", messages.status(id).name().toLowerCase());
|
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
|
// 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
|
// status poll — surface the open question and how to answer it, same as fleet_poll's
|
||||||
// Phase.ASKING view.
|
// Phase.ASKING view.
|
||||||
|
|||||||
@@ -32,6 +32,7 @@ import java.util.List;
|
|||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
|
import java.util.function.Predicate;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
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) {
|
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(
|
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||||
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||||
placement, "fleet", "worker: {profile} #{n}", null, null, 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.
|
// it directly so the inbox contract holds for those endpoints.
|
||||||
inbox.own("term_a");
|
inbox.own("term_a");
|
||||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
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);
|
.build().start("127.0.0.1", 0);
|
||||||
return app.port();
|
return app.port();
|
||||||
}
|
}
|
||||||
@@ -484,6 +491,16 @@ class FleetAppTest {
|
|||||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
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
|
* 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}
|
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
|
||||||
|
|||||||
Reference in New Issue
Block a user