From 4dd12083ab27ce7b87133f7cb5941bab3252eb0e Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 10:47:09 +0700 Subject: [PATCH 1/2] fleetd #284: free backend-error capacity --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 16 ++++++++++--- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 3 ++- .../fleet/FleetdBackendErrorSinkTest.java | 23 +++++++++++++++++++ .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 15 ++++++++++++ 4 files changed, 53 insertions(+), 4 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 9730fb4..b9f17fc 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -250,9 +250,7 @@ public final class Fleetd { boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn(); SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()), System::nanoTime, contextCap, clearAfterTurn); - liveCountRef.set(profileName -> (int) sessions.roster().stream() - .filter(s -> profileName.equals(s.profile())) - .count()); + liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName)); // CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled. final SessionReaper reaper; @@ -859,6 +857,18 @@ public final class Fleetd { .orElse(null); } + /** + * Count sessions that occupy a profile's spawn capacity. A {@code BACKEND_ERROR} session stays + * in the roster so {@code fleet_list} can show its failure, but its dead backend cannot use a + * seat or accept another delivery. + */ + static int liveSessionCount(List roster, String profileName) { + return (int) roster.stream() + .filter(session -> profileName.equals(session.profile())) + .filter(session -> session.state() != MemberSession.State.BACKEND_ERROR) + .count(); + } + /** * fleetd #248 / fleetd#201 Unit 5: factory for the production {@link BackendErrorSink} — the * collaborator {@link CompletionResolver} notifies when a pane-scrape classification actually diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index fab8fbc..add8557 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -1195,7 +1195,8 @@ public final class FleetMcp { int live = liveCount.apply(profile); int leadSeatCount = leadSeats.seatsFor().apply(profile); int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile())) - .filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE)) + .filter(s -> s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE + || s.state() == MemberSession.State.BACKEND_ERROR) .filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId()))) .count(); Map row = new LinkedHashMap<>(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java index d1c5086..c578632 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java @@ -226,4 +226,27 @@ class FleetdBackendErrorSinkTest { assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off"); assertEquals(1, leadClient.sendCount()); } + + @Test + @DisplayName("a backend-error session no longer blocks the real maxLoad spawn gate") + void backendErrorSessionDoesNotBlockFreshSpawnAtMaxLoad() { + FakeHerdr herdr = new FakeHerdr(); + FleetConfig.Profile profile = new FleetConfig.Profile("terra", "http://gx00.gw:8000", "coder", + null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", + "w #{n}", null, null, null, null, null, null, null, 1.0f, 1); + Map profiles = Map.of("terra", profile); + ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), profiles, "terra", _ -> "tok"); + AtomicReference sessionsRef = new AtomicReference<>(); + CompositePeerLauncher workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles, + PlacementPolicies.fixed(), name -> Fleetd.liveSessionCount(sessionsRef.get().roster(), name)); + SessionManager sessions = new SessionManager(workers); + sessionsRef.set(sessions); + MemberSession failed = sessions.acquire("terra", null, null, null); + + assertTrue(sessions.onBackendError(failed.terminalId(), "backend exited")); + + MemberSession fresh = sessions.acquire("terra", null, null, null); + assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a backend error"); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 6100420..7f4ea3b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -604,6 +604,21 @@ class FleetMcpTest { assertTrue(out.contains("\"reclaimable\":0"), out); } + @Test + void capacityReportsABackendErrorSessionAsReclaimable() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + MemberSession session = sessions.acquire("ltms-local", null, null, null); + assertTrue(sessions.onBackendError(session.terminalId(), "backend exited")); + + String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), + sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 1, + () -> Set.of("ltms-local"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "")); + + assertTrue(out.contains("\"reclaimable\":1"), out); + } + @Test void inertCapacitySourceOmitsCapacityBlock() { FakeHerdr h = new FakeHerdr(); From 11cbfa79b4ddbc12f46e6c4e9eeca84779fcf7e5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 10:55:21 +0700 Subject: [PATCH 2/2] fleetd #284: free failed-session capacity --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 7 +++-- .../dev/ltms/fleet/config/FleetConfig.java | 5 ++-- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 2 +- .../fleet/FleetdBackendErrorSinkTest.java | 29 +++++++++++++++---- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 10 ++++--- 5 files changed, 37 insertions(+), 16 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index b9f17fc..f63d2d8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -858,14 +858,15 @@ public final class Fleetd { } /** - * Count sessions that occupy a profile's spawn capacity. A {@code BACKEND_ERROR} session stays - * in the roster so {@code fleet_list} can show its failure, but its dead backend cannot use a - * seat or accept another delivery. + * Count sessions that occupy a profile's spawn capacity. A {@code BACKEND_ERROR} or + * {@code FAILED} session stays in the roster so {@code fleet_list} can show its failure, but a + * member that cannot accept another delivery does not use a seat. */ static int liveSessionCount(List roster, String profileName) { return (int) roster.stream() .filter(session -> profileName.equals(session.profile())) .filter(session -> session.state() != MemberSession.State.BACKEND_ERROR) + .filter(session -> session.state() != MemberSession.State.FAILED) .count(); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index a154c89..ef96eb1 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -322,8 +322,9 @@ public record FleetConfig( * statement that does not stop being true just because the profile was * named directly. A negative value has no sane meaning (there is no * "excluded" to degrade to below zero) and is refused at config load - * instead, naming the profile and the key. Live means any session the - * registry still owns (acquired and not yet released), in any state. + * instead, naming the profile and the key. Live means a session that can + * receive another delivery. The roster keeps terminal {@code BACKEND_ERROR} + * and {@code FAILED} sessions for diagnostics, but they do not use capacity. * @param kind which peer launcher spawns this profile: {@code "claude-code"} (default — * the {@link dev.ltms.fleet.member.ClaudeCodeLauncher}) or {@code "opencode"}. * The {@code CompositePeerLauncher} routes {@code spawn}/reap by this value, so diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index add8557..35ee563 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -1196,7 +1196,7 @@ public final class FleetMcp { int leadSeatCount = leadSeats.seatsFor().apply(profile); int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile())) .filter(s -> s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE - || s.state() == MemberSession.State.BACKEND_ERROR) + || s.state() == MemberSession.State.BACKEND_ERROR || s.state() == MemberSession.State.FAILED) .filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId()))) .count(); Map row = new LinkedHashMap<>(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java index c578632..432d862 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdBackendErrorSinkTest.java @@ -230,6 +230,28 @@ class FleetdBackendErrorSinkTest { @Test @DisplayName("a backend-error session no longer blocks the real maxLoad spawn gate") void backendErrorSessionDoesNotBlockFreshSpawnAtMaxLoad() { + SessionManager sessions = capacityLimitedSessions(); + MemberSession failed = sessions.acquire("terra", null, null, null); + + assertTrue(sessions.onBackendError(failed.terminalId(), "backend exited")); + + MemberSession fresh = sessions.acquire("terra", null, null, null); + assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a backend error"); + } + + @Test + @DisplayName("a failed session no longer blocks the real maxLoad spawn gate") + void failedSessionDoesNotBlockFreshSpawnAtMaxLoad() { + SessionManager sessions = capacityLimitedSessions(); + MemberSession failed = sessions.acquire("terra", null, null, null); + + sessions.onTurnFailed(failed.terminalId()); + + MemberSession fresh = sessions.acquire("terra", null, null, null); + assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a failed turn"); + } + + private static SessionManager capacityLimitedSessions() { FakeHerdr herdr = new FakeHerdr(); FleetConfig.Profile profile = new FleetConfig.Profile("terra", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", @@ -242,11 +264,6 @@ class FleetdBackendErrorSinkTest { PlacementPolicies.fixed(), name -> Fleetd.liveSessionCount(sessionsRef.get().roster(), name)); SessionManager sessions = new SessionManager(workers); sessionsRef.set(sessions); - MemberSession failed = sessions.acquire("terra", null, null, null); - - assertTrue(sessions.onBackendError(failed.terminalId(), "backend exited")); - - MemberSession fresh = sessions.acquire("terra", null, null, null); - assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a backend error"); + return sessions; } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 7f4ea3b..f011b5d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -605,18 +605,20 @@ class FleetMcpTest { } @Test - void capacityReportsABackendErrorSessionAsReclaimable() { + void capacityReportsTerminalFailureSessionsAsReclaimable() { FakeHerdr h = new FakeHerdr(); SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); - MemberSession session = sessions.acquire("ltms-local", null, null, null); - assertTrue(sessions.onBackendError(session.terminalId(), "backend exited")); + MemberSession backendError = sessions.acquire("ltms-local", null, null, null); + MemberSession failed = sessions.acquire("ltms-local", null, null, null); + assertTrue(sessions.onBackendError(backendError.terminalId(), "backend exited")); + sessions.onTurnFailed(failed.terminalId()); String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 1, () -> Set.of("ltms-local"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), Map.of(), "")); - assertTrue(out.contains("\"reclaimable\":1"), out); + assertTrue(out.contains("\"reclaimable\":2"), out); } @Test