Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 83f2aea60f | |||
| 76672ff016 | |||
| 9379f92c23 | |||
| 21ff63b11d | |||
| 21844b54d7 | |||
| 4769481515 | |||
| bbbb4c1eb3 | |||
| 38c248e617 | |||
| 9debc0de27 | |||
| f34361b263 | |||
| 85c90d440a | |||
| e97502d550 | |||
| aef14ff46e | |||
| cba516bda4 |
@@ -625,6 +625,19 @@ public final class Fleetd {
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
|
||||
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
|
||||
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
|
||||
// closure in only one of the two places), exactly the shape #284 was.
|
||||
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
|
||||
profile -> {
|
||||
@@ -636,15 +649,9 @@ public final class Fleetd {
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
}),
|
||||
new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine),
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy),
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
|
||||
|
||||
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
|
||||
@@ -715,9 +722,12 @@ public final class Fleetd {
|
||||
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
|
||||
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
|
||||
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
|
||||
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
|
||||
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable,
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials())).build();
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
|
||||
quarantineSource, outageSource).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
|
||||
@@ -262,11 +262,15 @@ public final class CallerResolver {
|
||||
return token.isEmpty() ? null : token;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #305: delegates to {@link ConnectionIdentity#isLoopback}. This used to be a second,
|
||||
* independent copy of the same rule, and the two drifted: this one accepted all of
|
||||
* {@code 127.0.0.0/8}, {@code ConnectionIdentity}'s accepted only {@code 127.0.0.1}. A caller
|
||||
* from {@code 127.0.0.2} therefore had its identity skipped (so it had no terminal) and was
|
||||
* then read as loopback here — which under loopback-trust is the primary. Sharing the inputs
|
||||
* would not have prevented that; only sharing the computation does.
|
||||
*/
|
||||
private static boolean isLoopback(String remoteAddr) {
|
||||
if (remoteAddr == null) {
|
||||
return false;
|
||||
}
|
||||
return remoteAddr.equals("127.0.0.1") || remoteAddr.equals("::1")
|
||||
|| remoteAddr.equals("0:0:0:0:0:0:0:1") || remoteAddr.startsWith("127.");
|
||||
return ConnectionIdentity.isLoopback(remoteAddr);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -157,6 +157,7 @@ public final class Injector {
|
||||
boolean awaitingCompletion; // a delivered message's turn is not yet known-complete
|
||||
boolean turnObserved; // saw a real `working` sample since that delivery (turn ran)
|
||||
int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
|
||||
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
|
||||
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
|
||||
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
|
||||
boolean awaitingPostTurnPickup;
|
||||
@@ -217,10 +218,12 @@ public final class Injector {
|
||||
t.awaitingPickup = false;
|
||||
t.injectableSincePickup = 0;
|
||||
t.unknownSinceTurn = 0;
|
||||
t.unknownSincePostTurn = 0;
|
||||
t.notReadySincePoll = 0;
|
||||
if (t.awaitingCompletion) t.turnObserved = true;
|
||||
} else if (status.injectable()) { // IDLE or BLOCKED
|
||||
t.unknownSinceTurn = 0;
|
||||
t.unknownSincePostTurn = 0;
|
||||
if (t.awaitingPostTurnPickup) {
|
||||
if (++t.injectableSincePostTurnPickup >= PICKUP_GRACE_POLLS) {
|
||||
t.awaitingPostTurnPickup = false;
|
||||
@@ -315,6 +318,24 @@ public final class Injector {
|
||||
t.unknownSinceTurn = 0;
|
||||
turnFailed = true;
|
||||
}
|
||||
// fleetd #306: the same escape for the post-turn housekeeping phase. Four latches
|
||||
// gate delivery (awaitingCompletion, postTurnPending, awaitingPostTurnPickup,
|
||||
// postTurnObserved) and only the first had a way out of a sustained unknown streak —
|
||||
// a gate that closed one direction only. The other two below are released here as
|
||||
// well; postTurnPending needs no escape because it is cleared unconditionally on the
|
||||
// line after the listener call that sets it.
|
||||
//
|
||||
// This does NOT set turnFailed. The delegated turn already completed and its waiter
|
||||
// already resolved — what is outstanding is adapter housekeeping (the `/clear`).
|
||||
// Reporting a turn failure here would drive SessionManager.onFailed on a session
|
||||
// that genuinely finished its work, which is a worse lie than the wedge.
|
||||
if ((t.awaitingPostTurnPickup || t.postTurnObserved)
|
||||
&& ++t.unknownSincePostTurn >= TURN_STALL_GRACE_POLLS) {
|
||||
t.awaitingPostTurnPickup = false;
|
||||
t.postTurnObserved = false;
|
||||
t.injectableSincePostTurnPickup = 0;
|
||||
t.unknownSincePostTurn = 0;
|
||||
}
|
||||
}
|
||||
|
||||
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
|
||||
|
||||
@@ -60,7 +60,30 @@ public final class ConnectionIdentity {
|
||||
return pid > 0 ? cwds.cwdForPid(pid) : null;
|
||||
}
|
||||
|
||||
private static boolean isLoopback(String addr) {
|
||||
return "127.0.0.1".equals(addr) || "::1".equals(addr) || "0:0:0:0:0:0:0:1".equals(addr);
|
||||
/**
|
||||
* Whether {@code addr} is a same-host address, and therefore one whose peer PID is worth
|
||||
* looking up. <strong>This is the one definition of loopback in the daemon</strong> —
|
||||
* {@code CallerResolver} calls it rather than keeping its own, because the two used to differ
|
||||
* and that difference was a privilege escalation (fleetd #305).
|
||||
*
|
||||
* <p>The whole of {@code 127.0.0.0/8} counts, not just {@code 127.0.0.1}. On Linux every
|
||||
* address in that range is bound to {@code lo} by default, so a process can connect to
|
||||
* {@code 127.0.0.1:8765} with a source address of {@code 127.0.0.2} — measured on the Linux
|
||||
* fleet host, where binding that source succeeds.
|
||||
*
|
||||
* <p><strong>Being strict here does not make the daemon safer; it makes it unsafe.</strong>
|
||||
* That reads backwards, so it is worth stating plainly. This predicate does not decide whether
|
||||
* a caller is trusted — it decides whether the caller's identity is <em>resolved at all</em>.
|
||||
* Returning false means {@link #resolve} answers "no terminal", and downstream a caller with no
|
||||
* terminal is treated as the primary under loopback-trust. So every address excluded here is an
|
||||
* address on which a worker silently becomes the lead. Widening a check normally weakens it;
|
||||
* widening this one is what closes the hole.
|
||||
*/
|
||||
public static boolean isLoopback(String addr) {
|
||||
if (addr == null) {
|
||||
return false;
|
||||
}
|
||||
String a = addr.startsWith("::ffff:") ? addr.substring(7) : addr; // IPv4-mapped IPv6
|
||||
return a.startsWith("127.") || "::1".equals(a) || "0:0:0:0:0:0:0:1".equals(a);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.ShuttingDownException;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
@@ -813,7 +814,12 @@ public final class FleetMcp {
|
||||
return error("fleet_reply is for workers only — could not identify the calling worker "
|
||||
+ "from the connection");
|
||||
}
|
||||
if (content == null) {
|
||||
// fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService
|
||||
// .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch
|
||||
// around it — so a whitespace-only fleet_reply would leave here as an uncaught
|
||||
// IllegalArgumentException instead of this clean tool error. Null and whitespace are the
|
||||
// same mistake by the caller and must get the same answer.
|
||||
if (isBlank(content)) {
|
||||
return error("content is required");
|
||||
}
|
||||
messages.reply(callerTerminal, content);
|
||||
@@ -957,6 +963,10 @@ public final class FleetMcp {
|
||||
return text(json(memberView(member)));
|
||||
} catch (GuardException e) {
|
||||
return error("subscription boundary: " + e.getMessage());
|
||||
} catch (ShuttingDownException e) {
|
||||
// fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather
|
||||
// than register a session drainAll will never see again.
|
||||
return error("shutting down: " + e.getMessage());
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
|
||||
// from "profile does not exist" below.
|
||||
@@ -1015,6 +1025,25 @@ public final class FleetMcp {
|
||||
* both maps at once when it is both exhaustion-quarantined AND cooling off.
|
||||
*/
|
||||
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
|
||||
return text(json(profilesView(workers, quarantine, outage)));
|
||||
}
|
||||
|
||||
/**
|
||||
* The body both front doors answer {@code profiles} with: the configured profile names, the
|
||||
* default, and the two independent outage states — {@code quarantined} (the backend reported it
|
||||
* out of capacity) and {@code coolingOff} (the credential threw repeated non-exhaustion backend
|
||||
* errors). Each map is present only when at least one profile is in that state, and a profile
|
||||
* can appear in both at once, because the two checks are separate.
|
||||
*
|
||||
* <p>fleetd #297: extracted so {@code fleet_profiles} and {@code GET /profiles} render from ONE
|
||||
* body builder rather than two copies. Passing both doors the same {@link QuarantineSource} and
|
||||
* {@link OutageSource} instances is necessary but not sufficient: with the loop written out
|
||||
* twice, a later edit to the row shape — a renamed key, an added field — lands on one door and
|
||||
* not the other, and the two then disagree about a live outage. That is exactly what fleetd
|
||||
* #284 was, where one rule computed in two places was widened in only one and a single response
|
||||
* contradicted itself. Shared inputs do not make duplicated computation safe.
|
||||
*/
|
||||
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("profiles", workers.profiles());
|
||||
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
|
||||
@@ -1046,7 +1075,7 @@ public final class FleetMcp {
|
||||
if (!coolingOff.isEmpty()) {
|
||||
result.put("coolingOff", coolingOff);
|
||||
}
|
||||
return text(json(result));
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -697,7 +697,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
if (paneId == null) {
|
||||
throw new IllegalStateException("pane.split returned no pane — cannot start a peer");
|
||||
}
|
||||
Agent peer = startUniquelyNamed(cfg, argv, paneId).agent();
|
||||
Agent peer;
|
||||
try {
|
||||
peer = startUniquelyNamed(cfg, argv, paneId).agent();
|
||||
} catch (RuntimeException e) {
|
||||
// The peer never started — don't leave the pane we just created orphaned.
|
||||
// Best-effort cleanup; never let it mask the real spawn failure.
|
||||
try {
|
||||
stop(paneId);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to close orphaned pane {} after spawn error: {}",
|
||||
paneId, cleanup.getMessage());
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
log.info("{} started pane={} terminal={}", namePrefix, peer.paneId(), peer.terminalId());
|
||||
return peer;
|
||||
}
|
||||
@@ -996,7 +1009,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* that gap: it stops waiting immediately (never burns the rest of the timeout), runs the same
|
||||
* teardown the timeout path below runs, and throws with a message that says the backend exited
|
||||
* rather than that the pane was slow. Any other {@link HerdrException} still propagates
|
||||
* unchanged — this gate does not know how to recover from it.
|
||||
* unchanged — this gate does not interpret or recover from it, but it still closes the pane
|
||||
* it opened before handing the exception to its caller.
|
||||
*/
|
||||
private void waitUntilInjectableOrThrow(String paneId) {
|
||||
long start = nowMillis.getAsLong();
|
||||
@@ -1010,7 +1024,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
if (isAlreadyGone(e)) {
|
||||
failFastOnGoneBackend(paneId, e, nowMillis.getAsLong() - start);
|
||||
}
|
||||
throw e; // any other herdr failure is not ours to interpret — let it propagate
|
||||
// This gate must not interpret an unrelated herdr error, but the caller does not
|
||||
// receive paneId when spawn throws. Close the pane here before propagating e unchanged.
|
||||
try {
|
||||
stop(paneId);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to close orphaned pane {} after readiness-gate error: {}",
|
||||
paneId, cleanup.getMessage());
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
lastStatus = sample.status();
|
||||
if (lastStatus.injectable() || refinedInjectable(paneId, sample)) {
|
||||
|
||||
@@ -408,10 +408,28 @@ public final class MessageService {
|
||||
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
|
||||
* deterministically instead.
|
||||
*
|
||||
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
|
||||
* method must reject a missing/blank reply the same way, so the check lives here rather than in
|
||||
* either caller: {@code FleetMcp.reply} already refuses a {@code null} content before it ever
|
||||
* calls this method (its own required-arg guard), and no test or production call site anywhere
|
||||
* in the codebase relies on replying with empty content — confirmed by searching every call site
|
||||
* of this method before adding the check, not assumed. Without this guard, a REST {@code
|
||||
* POST /sessions/{id}/reply} whose body omits {@code content} (or a client library that maps a
|
||||
* missing field to {@code ""}) used to reach {@link Rendezvous#resolve} with an empty string,
|
||||
* silently completing the lead's blocking wait with nothing — indistinguishable from a worker
|
||||
* that genuinely replied with nothing, which is worse than a loud failure because it destroys the
|
||||
* information that the reply never arrived.
|
||||
*
|
||||
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
|
||||
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve
|
||||
* anything
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
*/
|
||||
public boolean reply(String session, String content) {
|
||||
if (content == null || content.isBlank()) {
|
||||
throw new IllegalArgumentException("content is required");
|
||||
}
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
||||
return true; // a live send took it — unchanged fast path
|
||||
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.fleet.auth.Principal;
|
||||
import dev.ltms.fleet.guard.GuardException;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
@@ -18,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.ShuttingDownException;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
@@ -86,6 +88,12 @@ public final class FleetApp {
|
||||
// absent() (the honest "no policy configured" view) for every constructor that does not wire
|
||||
// a real one, so existing legacy call sites keep building without knowing this field exists.
|
||||
private final Supplier<MemberCredentialPolicyView> memberCredentials;
|
||||
// fleetd #297: the SAME shared sources FleetMcp.profiles/fleet_profiles reads (BackendQuarantine
|
||||
// and BackendOutagePolicy are each one instance for the whole daemon — see Fleetd wiring) so
|
||||
// GET /profiles cannot drift from fleet_profiles about which profile is quarantined/cooling off.
|
||||
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
|
||||
private final FleetMcp.QuarantineSource quarantine;
|
||||
private final FleetMcp.OutageSource outage;
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
/**
|
||||
@@ -146,6 +154,23 @@ public final class FleetApp {
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
|
||||
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
|
||||
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* @param quarantine the SAME {@link FleetMcp.QuarantineSource} instance passed to {@code
|
||||
* FleetMcp} (fleetd #297), so {@code GET /profiles} reports the identical
|
||||
* exhaustion-quarantine facts as {@code fleet_profiles} rather than a second,
|
||||
* independently-computed copy
|
||||
* @param outage the SAME {@link FleetMcp.OutageSource} instance passed to {@code FleetMcp} —
|
||||
* see {@code quarantine}; a SEPARATE check from it, never merged in
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
this.herdr = herdr;
|
||||
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
|
||||
this.workers = workers;
|
||||
@@ -156,6 +181,8 @@ public final class FleetApp {
|
||||
this.auth = auth;
|
||||
this.metrics = metrics;
|
||||
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
|
||||
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
|
||||
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
|
||||
}
|
||||
|
||||
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
|
||||
@@ -338,8 +365,15 @@ public final class FleetApp {
|
||||
if (!allow(ctx, routeAction("GET /agents"), null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
|
||||
try {
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
|
||||
// {error, detail} envelope every other failure path here uses, not escape as a bare
|
||||
// exception and leave Javalin's default handling to respond outside the JSON contract.
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
|
||||
@@ -347,40 +381,57 @@ public final class FleetApp {
|
||||
if (!allow(ctx, routeAction("GET /members"), null)) {
|
||||
return;
|
||||
}
|
||||
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
|
||||
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
|
||||
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
|
||||
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
|
||||
// drops reads this endpoint. Drop the alias once nothing reads it.
|
||||
body.put("members", out);
|
||||
body.put("workers", out);
|
||||
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
|
||||
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
|
||||
Map.of("count", st.count(), "costBytes", st.costBytes())));
|
||||
ctx.status(200).json(body);
|
||||
try {
|
||||
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
|
||||
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
|
||||
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
|
||||
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
|
||||
// drops reads this endpoint. Drop the alias once nothing reads it.
|
||||
body.put("members", out);
|
||||
body.put("workers", out);
|
||||
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
|
||||
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
|
||||
Map.of("count", st.count(), "costBytes", st.costBytes())));
|
||||
ctx.status(200).json(body);
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #297: same reasoning as agents() above — this is the out-of-band roster a lead
|
||||
// falls back to when its MCP mount drops, so it must stay inside the JSON error contract
|
||||
// exactly when herdr is briefly unreachable, not escape as a bare exception.
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** The configured worker profiles and which one a no-argument spawn uses. */
|
||||
/**
|
||||
* The configured worker profiles, which one a no-argument spawn uses, and (fleetd #297) the two
|
||||
* outage states {@code fleet_profiles} already reports: {@code quarantined} (CB-578 stage B —
|
||||
* the backend reported it out of capacity) and {@code coolingOff} (fleetd #201 Unit 5 — the
|
||||
* credential threw repeated non-exhaustion backend errors). Both are read from the SAME shared
|
||||
* {@link FleetMcp.QuarantineSource}/{@link FleetMcp.OutageSource} instances {@code FleetMcp}
|
||||
* reads, never recomputed, so the two doors cannot disagree about which profile is down and why.
|
||||
* Independent checks, so a profile can appear in both maps at once; each map is present only
|
||||
* when at least one profile is in that state.
|
||||
*/
|
||||
private void profiles(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /profiles"), null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
"profiles", workers.profiles(),
|
||||
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
|
||||
// fleetd #297: ONE body builder, shared with fleet_profiles. Handing both doors the same
|
||||
// QuarantineSource/OutageSource instances stops them reading different facts; rendering
|
||||
// through the same method stops them reporting those facts differently. Both are needed.
|
||||
ctx.status(200).json(FleetMcp.profilesView(workers, quarantine, outage));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -451,6 +502,11 @@ public final class FleetApp {
|
||||
ctx.status(201).json(view(member));
|
||||
} catch (GuardException e) {
|
||||
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
|
||||
} catch (ShuttingDownException e) {
|
||||
// fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500,
|
||||
// so this reads the same as PlacementException below: valid request, refused because
|
||||
// of a transient daemon state rather than a bad argument.
|
||||
ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage()));
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
|
||||
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
|
||||
@@ -460,6 +516,12 @@ public final class FleetApp {
|
||||
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
|
||||
} catch (PeerUnreachableException e) {
|
||||
ctx.status(502).json(Map.of("error", "spawn_timeout", "detail", e.getMessage()));
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #304: not every herdr failure on the spawn path is a readiness timeout, so
|
||||
// PeerUnreachableException above does not cover this. Without this catch the exception
|
||||
// escapes to Javalin's default 500, while fleet_spawn reports the same failure as a
|
||||
// clean named error (FleetMcp.spawn) — the #297 one-door-guarded shape.
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -480,13 +542,29 @@ public final class FleetApp {
|
||||
return (s == null || s.isBlank()) ? null : s;
|
||||
}
|
||||
|
||||
/** Tear a worker down by pane id. */
|
||||
/**
|
||||
* Tear a worker down by pane id.
|
||||
*
|
||||
* <p>fleetd #304: the {@code HerdrException} catch is not cosmetic. {@code release} deregisters
|
||||
* the session, notifies the release listener and preserves a dirty worktree <em>before</em> it
|
||||
* calls {@code launcher.stop}, so a throw from that stop arrives after the teardown the caller
|
||||
* asked for has already happened. Letting it escape gave Javalin's default 500, which tells the
|
||||
* caller to retry — and the retry finds nothing in the registry, reaches the same stop, and
|
||||
* throws again, so it can never succeed. {@code herdrError} instead answers 404 ("the pane is
|
||||
* gone, stop retrying") or 502 ("herdr is upstream and broken, a retry may help"), matching what
|
||||
* {@code fleet_stop} reports for the same failure.
|
||||
*/
|
||||
private void stopMember(Context ctx) {
|
||||
String paneId = ctx.pathParam("paneId");
|
||||
if (!allow(ctx, routeAction("DELETE /members/{paneId}"), paneId)) {
|
||||
return;
|
||||
}
|
||||
sessions.release(paneId);
|
||||
try {
|
||||
sessions.release(paneId);
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
return;
|
||||
}
|
||||
ctx.status(204);
|
||||
}
|
||||
|
||||
@@ -634,7 +712,19 @@ public final class FleetApp {
|
||||
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
|
||||
return;
|
||||
}
|
||||
messages.reply(id, content);
|
||||
// fleetd #302: content is required. `.path("content").asText("")` above turns a missing key
|
||||
// into "" rather than throwing, so without this check an empty/blank reply used to reach
|
||||
// messages.reply(...) and silently resolve the lead's waiter — the same class of bug as the
|
||||
// sibling "content is required" guards on sendMessage/askMessage below, except this one wrote
|
||||
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
|
||||
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
|
||||
// {error, detail} envelope this file uses everywhere else.
|
||||
try {
|
||||
messages.reply(id, content);
|
||||
} catch (IllegalArgumentException e) {
|
||||
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
|
||||
}
|
||||
|
||||
|
||||
@@ -21,6 +21,7 @@ import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -71,6 +72,16 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private volatile String fleetRepoRoot;
|
||||
|
||||
/**
|
||||
* fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot
|
||||
* is even taken — so a spawn already in flight sees the refusal as early as a plain flag can
|
||||
* make it. This alone cannot close the race completely: a caller that read {@code false} just
|
||||
* before the flip can still land in the registry after the snapshot. {@link #drainAll}'s
|
||||
* post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired,
|
||||
* see {@link #drainAll}'s javadoc.
|
||||
*/
|
||||
private final AtomicBoolean draining = new AtomicBoolean(false);
|
||||
|
||||
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
|
||||
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
|
||||
@@ -181,6 +192,14 @@ public final class SessionManager implements TurnListener {
|
||||
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
|
||||
String ownerTerminal, WorktreeRequest wt,
|
||||
String sessionName, String resumeSessionId) {
|
||||
// fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn —
|
||||
// so a caller learns the daemon is going down instead of getting a session drainAll will
|
||||
// never see again. Checked here because every other acquire(...) overload delegates to
|
||||
// this one, so this is the single point every spawn path passes through.
|
||||
if (draining.get()) {
|
||||
throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session "
|
||||
+ "the shutdown drain would never see");
|
||||
}
|
||||
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
|
||||
requireResumeCapability(profile, resumeSessionId);
|
||||
// CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only
|
||||
@@ -864,20 +883,60 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
|
||||
/**
|
||||
* Gracefully drain all registered sessions on daemon shutdown. For each session that is
|
||||
* {@code BUSY}, poll up to {@code timeoutNanos} for it to leave {@code BUSY}, then release it
|
||||
* regardless. Non-busy sessions are released immediately. A failure releasing one session is
|
||||
* logged and does not abort the rest.
|
||||
* Gracefully drain all registered sessions on daemon shutdown. Non-busy sessions are released
|
||||
* immediately; a {@code BUSY} one is polled until it leaves {@code BUSY}, then released
|
||||
* regardless. A failure releasing one session is logged and does not abort the rest.
|
||||
*
|
||||
* <p>{@code timeoutNanos} is a budget for the WHOLE drain, not a grace period per session: the
|
||||
* deadline is taken once, before the loop. So the first BUSY session can spend all of it, and a
|
||||
* later BUSY one is then released with no wait at all. That is deliberate. This drain is only
|
||||
* one phase of shutdown — {@code Fleetd} closes the message service, the push loop, the
|
||||
* heartbeat, MCP and the router after it — and the whole sequence has to finish inside
|
||||
* launchd's exit window. A per-session grace would let N busy members drain for N * the
|
||||
* timeout, overrun that window, and get the daemon SIGKILLed part-way through; the members not
|
||||
* yet reached would then get no clean release, no preserved-worktree log, and no snapshot.
|
||||
* Cutting one turn short is the cheaper failure, and it is not silent: an abandoned BUSY
|
||||
* session is preserved, snapshotted, and logged at WARN by {@code logPreservedForShutdown}.
|
||||
*
|
||||
* <p>CB-544: this is a {@link ReleaseCause#SHUTDOWN} release — the worker's pane is stopped
|
||||
* (the process must end) but its worktree is preserved and its path logged. Shutdown is never
|
||||
* a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY}
|
||||
* when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its
|
||||
* kept worktree.
|
||||
*
|
||||
* <p>fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used
|
||||
* to stop a new session from registering after it was taken — {@link #acquire} stayed open for
|
||||
* as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos}
|
||||
* budget. Two things close that window, deliberately paired because neither alone is complete:
|
||||
* {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire}
|
||||
* refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a
|
||||
* plain flag can close; and the sweep below re-reads the registry once the initial snapshot has
|
||||
* fully drained and drains whatever a straggler — a caller that read the flag as {@code false}
|
||||
* a moment before it flipped — still managed to register. The sweep shares the same
|
||||
* {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE
|
||||
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
|
||||
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
|
||||
* {@link #roster()} call.
|
||||
*/
|
||||
void drainAll(long timeoutNanos) {
|
||||
long deadline = System.nanoTime() + timeoutNanos;
|
||||
for (MemberSession s : roster()) {
|
||||
draining.set(true);
|
||||
drainSnapshot(roster(), deadline);
|
||||
List<MemberSession> stragglers = roster();
|
||||
if (!stragglers.isEmpty()) {
|
||||
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
|
||||
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
|
||||
drainSnapshot(stragglers, deadline);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
|
||||
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
|
||||
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
|
||||
*/
|
||||
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
|
||||
for (MemberSession s : snapshot) {
|
||||
try {
|
||||
if (s.state() == MemberSession.State.BUSY) {
|
||||
while (System.nanoTime() < deadline) {
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
/**
|
||||
* Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown
|
||||
* drain has already begun (fleetd #308).
|
||||
*
|
||||
* <p>{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is
|
||||
* in that snapshot. A session registered after the snapshot is invisible to the drain loop: its
|
||||
* pane is left running and its worktree is never preserved, and nothing else ever reclaims
|
||||
* either — the daemon's in-memory registry dies with the process. Refusing the spawn here,
|
||||
* loudly, is what stops that session from ever being created in the first place, rather than
|
||||
* silently handing the caller a session the daemon can no longer manage.
|
||||
*/
|
||||
public final class ShuttingDownException extends RuntimeException {
|
||||
public ShuttingDownException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
@@ -413,4 +413,29 @@ class CallerResolverTest {
|
||||
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, null));
|
||||
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, " "));
|
||||
}
|
||||
@Test
|
||||
void aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary() {
|
||||
// fleetd #305: the escalation. ConnectionIdentity used to accept only 127.0.0.1, so a
|
||||
// worker connecting from 127.0.0.2 resolved to no terminal, and this resolver's own
|
||||
// (wider) loopback check then made it the PRIMARY — granting spawn, stop, send and drain.
|
||||
// Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds there, so the
|
||||
// path is real and not theoretical.
|
||||
CallerResolver r = new CallerResolver(workerIdentity(), false, null);
|
||||
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "127.1.2.3", "::ffff:127.0.0.2"}) {
|
||||
Principal p = r.resolve(src, 55555, null);
|
||||
assertEquals(Role.WORKER, p.role(), "a worker must stay a worker from source " + src);
|
||||
assertEquals("term_a", p.terminal(), "worker terminal from source " + src);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonWorkerOnAnyLoopbackSourceAddressIsStillThePrimary() {
|
||||
// The other direction of the same fix: widening the identity check must not demote a
|
||||
// legitimate same-host primary that happens to connect from another 127.* address.
|
||||
CallerResolver r = new CallerResolver(nonWorkerIdentity(), false, null);
|
||||
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "::ffff:127.0.0.1"}) {
|
||||
assertEquals(Role.PRIMARY, r.resolve(src, 55555, null).role(), "source " + src);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -40,6 +40,7 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private final Map<String, List<String>> extraTabs = new LinkedHashMap<>();
|
||||
private int agentNameTakenFor = 0;
|
||||
private int agentPaneBusyFor = 0;
|
||||
private String agentStartErrorCode = null;
|
||||
private int workerTabPaneCount = 1;
|
||||
private String paneCloseErrorCode = null;
|
||||
private final Map<String, String> paneCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
@@ -84,6 +85,12 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make every {@code agent.start} call fail with this herdr error code. */
|
||||
public FakeHerdr agentStartFailsWith(String code) {
|
||||
this.agentStartErrorCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make the worker tab (w9:t2) report this many panes in {@code tab.list} (default 1). */
|
||||
public FakeHerdr withWorkerTabPaneCount(int n) {
|
||||
this.workerTabPaneCount = n;
|
||||
@@ -311,6 +318,10 @@ public final class FakeHerdr implements HerdrClient {
|
||||
+ required + "`", "invalid_request", null);
|
||||
}
|
||||
}
|
||||
if (agentStartErrorCode != null) {
|
||||
throw new HerdrException("herdr error [" + agentStartErrorCode + "]: agent.start failed",
|
||||
agentStartErrorCode, null);
|
||||
}
|
||||
long starts = calls.stream().filter(c -> c.method().equals("agent.start")).count();
|
||||
if (starts <= agentPaneBusyFor) {
|
||||
throw new HerdrException(
|
||||
|
||||
@@ -297,6 +297,68 @@ class InjectorTest {
|
||||
assertTrue(inj.activeTargets().isEmpty(), "the wedged target is reclaimed, not polled forever");
|
||||
}
|
||||
|
||||
/** A listener whose post-turn housekeeping always starts, as SessionManager's does with clearAfterTurn on. */
|
||||
private static final class PostTurnListener implements TurnListener {
|
||||
@Override public void onTurnComplete(String target) { }
|
||||
@Override public boolean hasPostTurnAction(String target) { return true; }
|
||||
@Override public boolean onTurnCompleteWithPostAction(String target) { return true; }
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWorkerThatWedgesInUnknownAwaitingPostTurnPickupIsReleased() {
|
||||
// fleetd #306: the post-turn phase had no way out of a sustained unknown streak, so the
|
||||
// pickup latch stayed set, the target was polled forever, and every later message to it was
|
||||
// blocked by the delivery gate — while the session still looked healthy.
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // turn starts
|
||||
inj.onStatus(T, AgentStatus.IDLE); // turn completes; housekeeping dispatched
|
||||
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // then wedges
|
||||
|
||||
assertTrue(inj.activeTargets().isEmpty(),
|
||||
"a target wedged awaiting post-turn pickup must be reclaimed, not polled forever");
|
||||
assertEquals(List.of(), cap.failed,
|
||||
"the delegated turn already completed — a stuck /clear must not be reported as a failed turn");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWorkerThatWedgesInUnknownAfterPickingUpTheResetIsReleased() {
|
||||
// The sibling latch. postTurnObserved is set when the reset is seen picked up (WORKING) and
|
||||
// is cleared only on a later injectable sample, so a wedge right after pickup sticks too.
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
inj.onStatus(T, AgentStatus.WORKING);
|
||||
inj.onStatus(T, AgentStatus.IDLE); // turn complete; reset dispatched
|
||||
inj.onStatus(T, AgentStatus.WORKING); // reset picked up -> postTurnObserved
|
||||
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN);
|
||||
|
||||
assertTrue(inj.activeTargets().isEmpty(), "a wedge after reset pickup must also be reclaimed");
|
||||
assertEquals(List.of(), cap.failed, "still not a turn failure");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBriefUnknownDuringPostTurnHousekeepingDoesNotDropTheLatch() {
|
||||
// The other direction: the escape must not fire on a glitch, or the queued next delegation
|
||||
// would overtake housekeeping that is still running.
|
||||
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
|
||||
inj.enqueue(T, "first", TestTurnTokens.inert(T));
|
||||
inj.enqueue(T, "second", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
inj.onStatus(T, AgentStatus.WORKING);
|
||||
inj.onStatus(T, AgentStatus.IDLE); // first completes; reset dispatched
|
||||
for (int i = 0; i < 10; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // well under the grace
|
||||
|
||||
assertFalse(inj.activeTargets().isEmpty(), "a brief glitch must not release the post-turn latch");
|
||||
assertEquals(List.of("first"), sent(), "the queued delegation must not overtake housekeeping");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() {
|
||||
Captor cap = new Captor();
|
||||
|
||||
@@ -20,6 +20,17 @@ class ConnectionIdentityTest {
|
||||
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.1", 55555));
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesWorkerFromAnyLoopbackSourceAddressNotJust127001() {
|
||||
// fleetd #305. On Linux the whole 127.0.0.0/8 is bound to lo, so a worker can connect with
|
||||
// a source address of 127.0.0.2. If identity resolution skips that address the caller has
|
||||
// no terminal, and a caller with no terminal is the primary under loopback-trust — so this
|
||||
// must resolve the worker, not null.
|
||||
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.2", 55555));
|
||||
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.1.2.3", 55555));
|
||||
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("::ffff:127.0.0.2", 55555));
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullForOffHostCaller() {
|
||||
// A non-loopback peer can't be an on-host worker → treat as primary/unknown.
|
||||
|
||||
@@ -326,6 +326,26 @@ class FleetMcpTest {
|
||||
assertEquals("orphan", drained.getFirst().content());
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() {
|
||||
// fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's
|
||||
// handler is a bare BiFunction with no try/catch around it, so if this guard only checked
|
||||
// `== null` (as it did), a whitespace-only reply would leave the handler as an uncaught
|
||||
// IllegalArgumentException instead of a tool error the caller can read. Null and whitespace
|
||||
// are the same caller mistake and must get the same answer — the sibling fleet_send guard
|
||||
// has always used isBlank for exactly this reason.
|
||||
for (String blank : new String[] {null, "", " ", "\n\t"}) {
|
||||
McpSchema.CallToolResult res = assertDoesNotThrow(
|
||||
() -> FleetMcp.reply(messages, "term_a", blank),
|
||||
"blank content must be refused as a tool error, never thrown out of the handler");
|
||||
assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result");
|
||||
assertTrue(textOf(res).contains("content is required"),
|
||||
"the error names the missing argument: " + textOf(res));
|
||||
}
|
||||
assertEquals(0, messages.drainReplies("term_a").size(),
|
||||
"a refused reply must not reach the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgePollWithTargetDrainsReplies() {
|
||||
// A reply with no open send queues it in the inbox.
|
||||
|
||||
@@ -1165,7 +1165,8 @@ class ClaudeCodeLauncherTest {
|
||||
@Test
|
||||
void spawnLetsAnUnrelatedHerdrErrorPropagateUnchanged() {
|
||||
// Fix 1 must only special-case a "*_not_found" answer. Any other herdr failure keeps
|
||||
// propagating as-is — this gate does not know how to recover from it.
|
||||
// propagating as-is — this gate does not know how to recover from it. The pane still needs
|
||||
// closing because spawn throws before it can return the pane id to a caller that could stop it.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
herdr.agentGetFailsWithAfter(0, "internal_error");
|
||||
@@ -1182,8 +1183,27 @@ class ClaudeCodeLauncherTest {
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertEquals("internal_error", ex.code());
|
||||
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"an error this gate does not recognize is not this gate's teardown to run");
|
||||
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"the unchanged error leaves spawn without a pane id, so this gate closes its orphaned pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void panePlacementClosesTheSplitPaneWhenAgentStartFails() {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("internal_error");
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "pane", "fleetd-workers", "w #{n}", null, null, null);
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
|
||||
dev.ltms.fleet.herdr.HerdrException ex = assertThrows(
|
||||
dev.ltms.fleet.herdr.HerdrException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertEquals("internal_error", ex.code(), "agent.start failure propagates unchanged");
|
||||
assertEquals(1, paneCloseCount(herdr, "w1:pSplit"),
|
||||
"the pane split for a peer that never starts is closed instead of left orphaned");
|
||||
}
|
||||
|
||||
// --- fleetd #176 fix 2: corroborated UNKNOWN refinement --------------------------------------
|
||||
|
||||
@@ -19,6 +19,10 @@ import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.Worktrees;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.member.MemberCredentialPolicyView;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
@@ -32,6 +36,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Predicate;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -70,6 +75,19 @@ class FleetAppTest {
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable) {
|
||||
return start(herdr, workerBaseUrl, allow, placement, worktrees, deliverable,
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297: same wiring as above, plus the two SAME shared sources {@code GET /profiles}
|
||||
* must read — lets a test prove the quarantined/coolingOff facts it reports come from a real
|
||||
* {@link dev.ltms.fleet.placement.BackendQuarantine}/{@link
|
||||
* dev.ltms.fleet.placement.BackendOutagePolicy}, exactly like {@code fleet_profiles}'s own tests.
|
||||
*/
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
placement, "fleet", "worker: {profile} #{n}", null, null, null);
|
||||
@@ -90,8 +108,9 @@ 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,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
|
||||
app = new FleetApp(herdr, herdr, workers, sessions, messages, this.presence, null,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id),
|
||||
MemberCredentialPolicyView::absent, quarantine, outage)
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
@@ -164,6 +183,22 @@ class FleetAppTest {
|
||||
assertEquals("idle", agents.get(0).get("status").asText());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 1: {@code workers.list()} reaches herdr, and a transport failure there must
|
||||
* land in the same {@code {error, detail}} envelope every other failure path in this file uses
|
||||
* (see {@code herdrError}), not escape as a bare exception outside the JSON contract.
|
||||
*/
|
||||
@Test
|
||||
void agentsMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
HttpResponse<String> res = req(port, "GET", "/agents");
|
||||
assertEquals(502, res.statusCode(), res.body());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("herdr_error", body.get("error").asText());
|
||||
assertTrue(body.has("detail"), res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -203,6 +238,42 @@ class FleetAppTest {
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
|
||||
assertEquals("ltms-local", body.get("default").asText());
|
||||
assertEquals("ltms-local", body.get("profiles").get(0).asText());
|
||||
assertFalse(body.has("quarantined"), "nothing is quarantined, so the key is omitted: " + body);
|
||||
assertFalse(body.has("coolingOff"), "nothing is cooling off, so the key is omitted: " + body);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 2: {@code GET /profiles} must report the same two outage states {@code
|
||||
* fleet_profiles} does — CB-578 stage B exhaustion quarantine and fleetd #201 Unit 5 cool-off —
|
||||
* reading the SAME shared {@link BackendQuarantine}/{@link BackendOutagePolicy} instances rather
|
||||
* than recomputing them. The two checks are independent, and this profile is deliberately put in
|
||||
* both states at once, matching {@code FleetMcpTest}'s own coverage of that overlap.
|
||||
*/
|
||||
@Test
|
||||
void profilesReportsQuarantineAndCoolingOffFromTheSameSharedSources() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
|
||||
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
|
||||
outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", new GitWorktrees(),
|
||||
ignored -> false, quarantineSource, outageSource);
|
||||
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
|
||||
assertTrue(body.has("quarantined"), body.toString());
|
||||
assertEquals("shared-openai",
|
||||
body.get("quarantined").get("ltms-local").get("credentialId").asText());
|
||||
assertEquals(1800,
|
||||
body.get("quarantined").get("ltms-local").get("quarantinedForSeconds").asLong());
|
||||
assertTrue(body.has("coolingOff"), body.toString());
|
||||
assertEquals("shared-openai",
|
||||
body.get("coolingOff").get("ltms-local").get("credentialId").asText());
|
||||
assertEquals(60, body.get("coolingOff").get("ltms-local").get("coolingOffForSeconds").asLong());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -239,6 +310,23 @@ class FleetAppTest {
|
||||
"liveStatus is unknown when herdr has no matching pane");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 1: same reasoning as {@code agentsMapsAHerdrFailureToTheJsonErrorEnvelope} —
|
||||
* {@code GET /members} is the endpoint's own comment names as "the out-of-band path a lead falls
|
||||
* back to when its MCP mount drops", so it must stay inside the {@code {error, detail}} envelope
|
||||
* exactly when herdr is briefly unreachable.
|
||||
*/
|
||||
@Test
|
||||
void membersMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
HttpResponse<String> res = req(port, "GET", "/members");
|
||||
assertEquals(502, res.statusCode(), res.body());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("herdr_error", body.get("error").asText());
|
||||
assertTrue(body.has("detail"), res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWithACwdParamRootsTheWorkerThere() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -327,7 +415,11 @@ class FleetAppTest {
|
||||
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(99);
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
assertEquals(500, req(port, "POST", "/members").statusCode());
|
||||
// fleetd #304: 502, not Javalin's default 500 — the herdr failure is named, and the body
|
||||
// carries herdr's own message, matching what fleet_spawn reports for the same failure.
|
||||
HttpResponse<String> res = req(port, "POST", "/members");
|
||||
assertEquals(502, res.statusCode());
|
||||
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
|
||||
assertTrue(herdr.called("tab.create"), "a tab was created before the failed start");
|
||||
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"), "orphaned tab must be closed");
|
||||
}
|
||||
@@ -444,6 +536,53 @@ class FleetAppTest {
|
||||
assertEquals("orphan", body.get("replies").get(0).get("content").asText());
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyWithMissingContentIsRejectedAndDoesNotResolveTheWaiter() throws Exception {
|
||||
// fleetd #302: `.path("content").asText("")` used to turn a missing "content" key into an
|
||||
// empty string that reached rendezvous.resolve, silently completing the lead's blocking wait
|
||||
// with nothing. Prove the fix two ways: the bad call is rejected with 400, AND the real send
|
||||
// it would have wrongly resolved is still open afterwards — a real reply completes it.
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle");
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
var send = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||
try { return postMessage(port, "{\"content\":\"review this\",\"timeoutMs\":4000}"); }
|
||||
catch (Exception e) { throw new RuntimeException(e); }
|
||||
});
|
||||
|
||||
Thread.sleep(200); // let the background send open its rendezvous waiter
|
||||
|
||||
HttpResponse<String> badReply = postJson(port, "/sessions/term_a/reply", "{}");
|
||||
assertEquals(400, badReply.statusCode());
|
||||
JsonNode err = mapper.readTree(badReply.body());
|
||||
assertEquals("bad_request", err.get("error").asText());
|
||||
assertTrue(err.has("detail"));
|
||||
|
||||
// The waiter must still be open — a real reply now completes the ORIGINAL send.
|
||||
HttpResponse<String> goodReply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
|
||||
assertEquals(200, goodReply.statusCode());
|
||||
|
||||
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||
assertEquals(200, res.statusCode());
|
||||
assertEquals("LGTM ship it", mapper.readTree(res.body()).get("reply").asText());
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyWithEmptyOrWhitespaceContentIsRejectedSameAsMissing() throws Exception {
|
||||
// fleetd #302 sibling case: present-but-blank content is treated the same as a missing key —
|
||||
// FleetMcp's own required-content guard (fleet_reply's "content is required") makes no
|
||||
// distinction between the two either, so diverging here would be a new asymmetry.
|
||||
int port = startHealthy();
|
||||
|
||||
HttpResponse<String> empty = postJson(port, "/sessions/term_a/reply", "{\"content\":\"\"}");
|
||||
assertEquals(400, empty.statusCode());
|
||||
assertEquals("bad_request", mapper.readTree(empty.body()).get("error").asText());
|
||||
|
||||
HttpResponse<String> whitespace = postJson(port, "/sessions/term_a/reply", "{\"content\":\" \"}");
|
||||
assertEquals(400, whitespace.statusCode());
|
||||
assertEquals("bad_request", mapper.readTree(whitespace.body()).get("error").asText());
|
||||
}
|
||||
|
||||
@Test
|
||||
void drainRepliesReturnsEmptyForNoReplies() throws Exception {
|
||||
int port = startHealthy();
|
||||
@@ -598,7 +737,12 @@ class FleetAppTest {
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
// A genuine teardown failure must surface, not be reported as a successful 204.
|
||||
assertEquals(500, req(port, "DELETE", "/members/w9:pW").statusCode());
|
||||
// fleetd #304: it surfaces as a named 502 rather than Javalin's default 500. The property
|
||||
// this test guards is "not 204" and the herdr detail reaching the caller — a bare 500 gave
|
||||
// the body "Server Error" and said nothing about herdr.
|
||||
HttpResponse<String> res = req(port, "DELETE", "/members/w9:pW");
|
||||
assertEquals(502, res.statusCode());
|
||||
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
|
||||
assertFalse(herdr.called("tab.close"), "tab is not removed when the pane close failed");
|
||||
}
|
||||
|
||||
@@ -611,4 +755,5 @@ class FleetAppTest {
|
||||
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
|
||||
assertTrue(herdr.called("tab.close"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -824,6 +829,186 @@ class SessionManagerTest {
|
||||
.count();
|
||||
}
|
||||
|
||||
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
|
||||
|
||||
@Test
|
||||
void acquireRefusesANewSpawnOnceDrainAllHasStarted() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately,
|
||||
// but the shutdown guard it flips must stay tripped for the life of the process.
|
||||
|
||||
ShuttingDownException e = assertThrows(ShuttingDownException.class,
|
||||
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
|
||||
"a spawn requested after the drain has begun must be refused loudly (invariant 3), "
|
||||
+ "not silently registered into a registry the drain will never revisit");
|
||||
assertNotNull(e.getMessage());
|
||||
assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed");
|
||||
assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close
|
||||
* all of it — a caller that already passed the {@code draining} check before {@code drainAll}
|
||||
* flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not
|
||||
* instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly
|
||||
* that interleaving with a launcher double that blocks the second {@code spawn()} call and the
|
||||
* first {@code stop()} call until released, then proves the post-loop sweep in {@code
|
||||
* drainAll} still finds and tears down the straggler that lands in the registry afterward.
|
||||
*/
|
||||
@Test
|
||||
void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
RaceLauncher race = new RaceLauncher(delegate);
|
||||
SessionManager sessions = new SessionManager(race);
|
||||
|
||||
// Registered normally, before the drain starts — the first spawn call, never blocked.
|
||||
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
|
||||
|
||||
ExecutorService exec = Executors.newFixedThreadPool(2);
|
||||
try {
|
||||
// The straggler's acquire() reads `draining == false` (checked before this call ever
|
||||
// touches the launcher) and then blocks inside its own spawn() — the second spawn call.
|
||||
Future<MemberSession> straggler = exec.submit(() ->
|
||||
sessions.acquire("ltms-local", "/late", "/caller", "ownerLate"));
|
||||
|
||||
assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS),
|
||||
"the straggler must have passed the shutdown guard and reached spawn() before "
|
||||
+ "drainAll ever runs");
|
||||
assertEquals(1, sessions.roster().size(),
|
||||
"the straggler is still inside spawn() — not registered yet");
|
||||
|
||||
// drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts
|
||||
// releasing that snapshot — its first release() call stops `ready`'s pane, which this
|
||||
// launcher double blocks on so the interleaving below is deterministic, not a timing bet.
|
||||
Future<?> drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5)));
|
||||
|
||||
assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS),
|
||||
"drainAll must be stopping the ready session's pane — proof its initial "
|
||||
+ "registry snapshot has already been taken");
|
||||
|
||||
// Only now does the straggler's spawn complete and register — strictly after the
|
||||
// snapshot drainAll's main pass is working from.
|
||||
race.releaseSecondSpawn.countDown();
|
||||
MemberSession registered = straggler.get(5, TimeUnit.SECONDS);
|
||||
|
||||
// Let drainAll finish releasing `ready`; it then re-checks the registry and must find
|
||||
// (and drain) the straggler that just landed in it.
|
||||
race.releaseFirstStop.countDown();
|
||||
drain.get(5, TimeUnit.SECONDS);
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(),
|
||||
"the post-loop sweep must drain the straggler too, not just the initial snapshot");
|
||||
assertNotNull(registered.paneId());
|
||||
long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count();
|
||||
assertEquals(2, paneCloseCalls,
|
||||
"both ready's pane AND the straggler's pane must actually be stopped — a pane "
|
||||
+ "left running is exactly the orphan this ticket is about");
|
||||
} finally {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST
|
||||
* {@code stop()} call until the test releases them — used to force the fleetd #308 race
|
||||
* deterministically instead of betting on real thread-scheduling timing.
|
||||
*/
|
||||
private static final class RaceLauncher implements PeerLauncher {
|
||||
private final PeerLauncher delegate;
|
||||
private final AtomicInteger spawnCalls = new AtomicInteger();
|
||||
private final AtomicInteger stopCalls = new AtomicInteger();
|
||||
final CountDownLatch enteredSecondSpawn = new CountDownLatch(1);
|
||||
final CountDownLatch releaseSecondSpawn = new CountDownLatch(1);
|
||||
final CountDownLatch enteredFirstStop = new CountDownLatch(1);
|
||||
final CountDownLatch releaseFirstStop = new CountDownLatch(1);
|
||||
|
||||
RaceLauncher(PeerLauncher delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
private static void awaitOrFail(CountDownLatch latch) {
|
||||
try {
|
||||
if (!latch.await(5, TimeUnit.SECONDS)) {
|
||||
throw new AssertionError("RaceLauncher latch timed out");
|
||||
}
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new AssertionError("RaceLauncher latch interrupted", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return delegate.capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return delegate.capabilitiesFor(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
if (spawnCalls.incrementAndGet() == 2) {
|
||||
enteredSecondSpawn.countDown();
|
||||
awaitOrFail(releaseSecondSpawn);
|
||||
}
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return delegate.defaultProfile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return delegate.effectiveCwd(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return delegate.parityOverlay(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return delegate.list();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return delegate.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
if (stopCalls.incrementAndGet() == 1) {
|
||||
enteredFirstStop.countDown();
|
||||
awaitOrFail(releaseFirstStop);
|
||||
}
|
||||
delegate.stop(id);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return delegate.clearContext(id);
|
||||
}
|
||||
}
|
||||
|
||||
private static List<String> promptTexts(FakeHerdr herdr) {
|
||||
return herdr.calls.stream()
|
||||
.filter(c -> "agent.prompt".equals(c.method()))
|
||||
|
||||
Reference in New Issue
Block a user