Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a49671ceb9 | |||
| 76672ff016 | |||
| 9379f92c23 | |||
| 21ff63b11d | |||
| 21844b54d7 | |||
| 4769481515 |
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -813,7 +813,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);
|
||||
|
||||
@@ -510,6 +510,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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -530,13 +536,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);
|
||||
}
|
||||
|
||||
|
||||
@@ -61,6 +61,8 @@ public final class SessionManager implements TurnListener {
|
||||
private final LongSupplier nowNanos;
|
||||
private final int contextCap;
|
||||
private final boolean clearAfterTurn;
|
||||
/** Null in production; test seam for the interval before an idle session's conditional release. */
|
||||
private final Consumer<MemberSession> beforeIdleRelease;
|
||||
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
||||
/**
|
||||
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
|
||||
@@ -102,13 +104,23 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
|
||||
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Package-private constructor for a deterministic reap/delivery race test. Production callers
|
||||
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
|
||||
*/
|
||||
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
|
||||
this.launcher = launcher;
|
||||
this.worktrees = worktrees;
|
||||
this.presence = new PresenceFleet(this);
|
||||
this.nowNanos = nowNanos;
|
||||
this.contextCap = contextCap;
|
||||
this.clearAfterTurn = clearAfterTurn;
|
||||
this.beforeIdleRelease = beforeIdleRelease;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -272,10 +284,27 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
}
|
||||
|
||||
/**
|
||||
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
|
||||
* transition replaces the immutable record, so this prevents a reap based on an old READY or
|
||||
* DONE record from stopping a worker that delivery has made BUSY.
|
||||
*/
|
||||
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
return false;
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
}
|
||||
|
||||
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
|
||||
ReleaseCause cause) {
|
||||
// fleetd #209: remove right alongside the registry entry so a released session's handle is
|
||||
// never leaked — but keep the local reference below, so the id can still be resolved for
|
||||
// the ReleaseDetail this teardown notifies with.
|
||||
PeerHandle removedHandle = handles.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
String snapshotRef = null;
|
||||
if (removed != null) {
|
||||
@@ -846,14 +875,18 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
long idleNanos = now - s.lastActivityAtNanos();
|
||||
if (idleNanos > idleTtlNanos) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
// CB-581: one session that fails to release must not abort the whole reaping pass —
|
||||
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
|
||||
try {
|
||||
release(s.paneId());
|
||||
reaped++;
|
||||
if (beforeIdleRelease != null) {
|
||||
beforeIdleRelease.accept(s);
|
||||
}
|
||||
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
reaped++;
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
|
||||
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -415,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");
|
||||
}
|
||||
@@ -733,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");
|
||||
}
|
||||
|
||||
@@ -746,4 +755,5 @@ class FleetAppTest {
|
||||
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
|
||||
assertTrue(herdr.called("tab.close"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -178,14 +178,19 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn) {
|
||||
boolean clearAfterTurn) {
|
||||
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
|
||||
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 workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn);
|
||||
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -670,6 +675,24 @@ class SessionManagerTest {
|
||||
"BUSY session remains");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager[] manager = new SessionManager[1];
|
||||
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
|
||||
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
|
||||
manager[0] = sessions;
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
clock[0] = 11;
|
||||
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
|
||||
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a just-delivered session stays registered and busy");
|
||||
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
|
||||
}
|
||||
|
||||
@Test
|
||||
void doneSessionPastIdleTtlIsReaped() {
|
||||
long[] clock = {0};
|
||||
|
||||
Reference in New Issue
Block a user