fleetd #226: reserve an architect slot before launch, so no member runs on a charter it will not hold
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m31s

Verified by the lead before merge. Read the full production diff: locking is consistent (every MemberRegistry method uses synchronized(terminalToSlot)), and the refusal happens before the launcher starts a process. Proved the restored fallback test is real by removing the DEV fallback from SessionManager and re-running it: it failed with "expected: <DEV> but was: <ARCHITECT>", then reverted. Independent build: BUILD SUCCESS, 1092 tests, 0 failures, 0 skipped.
This commit was merged in pull request #228.
This commit is contained in:
2026-09-02 02:54:52 +02:00
5 changed files with 227 additions and 69 deletions
@@ -5,6 +5,9 @@ import dev.ltms.fleet.peer.MemberRole;
/** Optional session lifecycle hook for live member-slot bindings. */ /** Optional session lifecycle hook for live member-slot bindings. */
public interface MemberLifecycle { public interface MemberLifecycle {
/** A slot held before a member process starts. */
record SlotReservation(String slot, String profile) { }
MemberLifecycle NONE = new MemberLifecycle() { MemberLifecycle NONE = new MemberLifecycle() {
@Override @Override
public MemberRole acquired(MemberRole role, String profile, String terminal) { public MemberRole acquired(MemberRole role, String profile, String terminal) {
@@ -19,6 +22,20 @@ public interface MemberLifecycle {
public void requireSlotFor(MemberRole role, String profile) { public void requireSlotFor(MemberRole role, String profile) {
// no registry configured — nothing to validate against, so nothing is refused // no registry configured — nothing to validate against, so nothing is refused
} }
@Override
public SlotReservation reserve(MemberRole role, String profile) {
return null;
}
@Override
public boolean bind(SlotReservation reservation, String terminal) {
return false;
}
@Override
public void release(SlotReservation reservation) {
}
}; };
/** /**
@@ -29,10 +46,9 @@ public interface MemberLifecycle {
* role — never {@code role} — when a slot-bound role (architect) could not be bound. * role — never {@code role} — when a slot-bound role (architect) could not be bound.
* Callers must record THIS value on the session, never the requested {@code role}, so * Callers must record THIS value on the session, never the requested {@code role}, so
* a later roster read never reports a role the session does not hold (CB-619). In * a later roster read never reports a role the session does not hold (CB-619). In
* normal operation this fallback should not happen once {@link #requireSlotFor} has * normal operation this fallback should not happen once a reservation has been bound.
* refused every unbindable spawn upfront — but a slot can still be lost between that * It remains the honest answer if a caller has no reservation, or if binding a
* check and this call to a concurrent spawn racing for the same slot, so the honest * reservation unexpectedly fails.
* answer is still needed here too.
*/ */
MemberRole acquired(MemberRole role, String profile, String terminal); MemberRole acquired(MemberRole role, String profile, String terminal);
@@ -46,4 +62,13 @@ public interface MemberLifecycle {
* @throws IllegalArgumentException naming the role, the profile, and the pools that do carry it * @throws IllegalArgumentException naming the role, the profile, and the pools that do carry it
*/ */
void requireSlotFor(MemberRole role, String profile); void requireSlotFor(MemberRole role, String profile);
/** Reserve a matching slot before launch, or refuse before a charter can be delivered. */
SlotReservation reserve(MemberRole role, String profile);
/** Convert a reservation into a live terminal binding. */
boolean bind(SlotReservation reservation, String terminal);
/** Return an unbound reservation after a failed launch. */
void release(SlotReservation reservation);
} }
@@ -56,8 +56,10 @@ public final class MemberRegistry implements MemberLifecycle {
} }
private final Map<String, Entry> slots; private final Map<String, Entry> slots;
/** Live {@code terminal_id → qualified slot key}; guarded by {@code this}. */ /** Live {@code terminal_id → qualified slot key}; guarded by {@code terminalToSlot}. */
private final Map<String, String> terminalToSlot = new HashMap<>(); private final Map<String, String> terminalToSlot = new HashMap<>();
/** Slot keys held between reservation and the terminal binding. Guarded by terminalToSlot. */
private final java.util.Set<String> reservedSlots = new java.util.HashSet<>();
/** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */ /** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */
public MemberRegistry(FleetConfig.Fleet fleet) { public MemberRegistry(FleetConfig.Fleet fleet) {
@@ -167,7 +169,7 @@ public final class MemberRegistry implements MemberLifecycle {
if (existingSlot != null) { if (existingSlot != null) {
return slot.equals(existingSlot); // already this slot (idempotent) or a different one return slot.equals(existingSlot); // already this slot (idempotent) or a different one
} }
if (terminalToSlot.containsValue(slot)) { if (terminalToSlot.containsValue(slot) || reservedSlots.contains(slot)) {
return false; // slot already hosts a terminal — no second one return false; // slot already hosts a terminal — no second one
} }
terminalToSlot.put(terminal, slot); terminalToSlot.put(terminal, slot);
@@ -275,6 +277,50 @@ public final class MemberRegistry implements MemberLifecycle {
+ " on one of those profiles instead"); + " on one of those profiles instead");
} }
@Override
public SlotReservation reserve(MemberRole role, String profile) {
if (role != MemberRole.ARCHITECT) {
return null;
}
synchronized (terminalToSlot) {
for (Entry entry : slotsFor(MemberRole.ARCHITECT).values()) {
if ((profile == null || profile.isBlank() || Objects.equals(profile, entry.profile()))
&& !terminalToSlot.containsValue(entry.key()) && reservedSlots.add(entry.key())) {
return new SlotReservation(entry.key(), entry.profile());
}
}
}
throw new IllegalArgumentException("no free architect slot for profile '" + profile
+ "' — every matching slot is already bound or reserved");
}
@Override
public boolean bind(SlotReservation reservation, String terminal) {
if (reservation == null || terminal == null || terminal.isBlank()) {
return false;
}
synchronized (terminalToSlot) {
if (!reservedSlots.remove(reservation.slot())) {
return false;
}
if (!isSlot(reservation.slot()) || terminalToSlot.containsKey(terminal)
|| terminalToSlot.containsValue(reservation.slot())) {
return false;
}
terminalToSlot.put(terminal, reservation.slot());
return true;
}
}
@Override
public void release(SlotReservation reservation) {
if (reservation != null) {
synchronized (terminalToSlot) {
reservedSlots.remove(reservation.slot());
}
}
}
/** Unbind a released terminal using the compare-safe registry operation. */ /** Unbind a released terminal using the compare-safe registry operation. */
@Override @Override
public void released(String terminal) { public void released(String terminal) {
@@ -190,49 +190,44 @@ public final class SessionManager implements TurnListener {
if (profile != null && !profile.isBlank()) { if (profile != null && !profile.isBlank()) {
memberLifecycle.requireSlotFor(memberRole, profile); memberLifecycle.requireSlotFor(memberRole, profile);
} }
MemberLifecycle.SlotReservation reservation = memberLifecycle.reserve(memberRole, profile);
String launchProfile = reservation == null ? profile : reservation.profile();
if (wt == null) { if (wt == null) {
// CB-557: the role must ride on the SpawnRequest, not stay a local. The launcher needs it // CB-557: the role must ride on the SpawnRequest, not stay a local. The launcher needs it
// to pick the profile out of that role's pool and to label the tab; a role kept only on // to pick the profile out of that role's pool and to label the tab; a role kept only on
// the MemberSession is recorded after the spawn it was supposed to steer. // the MemberSession is recorded after the spawn it was supposed to steer.
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, memberRole); SpawnRequest req = new SpawnRequest(launchProfile, requestedCwd, callerCwd, sessionName, resumeSessionId, memberRole);
PeerHandle handle; PeerHandle handle;
boolean bound = false;
try { try {
handle = launcher.spawn(req); handle = launcher.spawn(req);
String resolvedProfile = resolveProfile(handle, launchProfile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
long now = nowNanos.getAsLong();
MemberRole actualRole = acquired(memberRole, resolvedProfile, handle.terminalId(), reservation);
bound = reservation == null || actualRole == MemberRole.ARCHITECT;
MemberSession session = new MemberSession(
handle.id(), handle.terminalId(), resolvedProfile, actualRole, cwd, ownerTerminal, now, now, 0,
MemberSession.State.SPAWNING, null, null, handle.charterReceipt(), handle.agentSessionId());
registry.put(handle.id(), session);
handles.put(handle.id(), handle);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
notifyAcquired(session.terminalId());
return session;
} catch (RuntimeException e) { } catch (RuntimeException e) {
log.warn("spawn failed for profile={} role={}: {}", profile, memberRole, e.getMessage()); if (!bound) memberLifecycle.release(reservation);
log.warn("spawn failed for profile={} role={}: {}", launchProfile, memberRole, e.getMessage());
throw e; throw e;
} }
String resolvedProfile = resolveProfile(handle, profile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
long now = nowNanos.getAsLong();
// CB-619: bind (or fail to bind) BEFORE the session is recorded, and store whatever role
// this call actually returns — never the requested memberRole — so the session's role,
// what GET /members and fleet_list report, is never a lie about what this terminal holds.
MemberRole actualRole = memberLifecycle.acquired(memberRole, resolvedProfile, handle.terminalId());
MemberSession session = new MemberSession(
handle.id(),
handle.terminalId(),
resolvedProfile,
actualRole,
cwd,
ownerTerminal,
now,
now,
0,
MemberSession.State.SPAWNING,
null,
null,
handle.charterReceipt(),
handle.agentSessionId());
registry.put(handle.id(), session);
handles.put(handle.id(), handle);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
notifyAcquired(session.terminalId());
return session;
} }
return acquireWithWorktree(profile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt, try {
sessionName, resumeSessionId); return acquireWithWorktree(launchProfile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt,
sessionName, resumeSessionId, reservation);
} catch (RuntimeException e) {
memberLifecycle.release(reservation);
throw e;
}
} }
/** /**
@@ -478,7 +473,8 @@ public final class SessionManager implements TurnListener {
private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd, private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt, String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId) { String sessionName, String resumeSessionId,
MemberLifecycle.SlotReservation reservation) {
String preResolvedProfile = (profile == null || profile.isBlank()) String preResolvedProfile = (profile == null || profile.isBlank())
? launcher.defaultProfile() : profile; ? launcher.defaultProfile() : profile;
// CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller → // CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller →
@@ -521,7 +517,7 @@ public final class SessionManager implements TurnListener {
long now = nowNanos.getAsLong(); long now = nowNanos.getAsLong();
// CB-619: see the no-worktree path above — bind before recording, and store the returned // CB-619: see the no-worktree path above — bind before recording, and store the returned
// actual role, so this session's role is never a lie about what it actually holds. // actual role, so this session's role is never a lie about what it actually holds.
MemberRole actualRole = memberLifecycle.acquired(memberRole, resolvedProfile, handle.terminalId()); MemberRole actualRole = acquired(memberRole, resolvedProfile, handle.terminalId(), reservation);
MemberSession session = new MemberSession( MemberSession session = new MemberSession(
handle.id(), handle.id(),
handle.terminalId(), handle.terminalId(),
@@ -545,6 +541,18 @@ public final class SessionManager implements TurnListener {
return session; return session;
} }
private MemberRole acquired(MemberRole role, String profile, String terminal,
MemberLifecycle.SlotReservation reservation) {
if (reservation == null) {
return memberLifecycle.acquired(role, profile, terminal);
}
if (memberLifecycle.bind(reservation, terminal)) {
return role;
}
memberLifecycle.release(reservation);
return memberLifecycle.acquired(role, profile, terminal);
}
private String slug(String raw) { private String slug(String raw) {
return raw == null ? "ticket" : raw.toLowerCase().replaceAll("[^a-z0-9]+", "-").replaceAll("^-+|-+$", ""); return raw == null ? "ticket" : raw.toLowerCase().replaceAll("[^a-z0-9]+", "-").replaceAll("^-+|-+$", "");
} }
@@ -5,6 +5,7 @@ import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender; import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.auth.MemberRegistry; import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.MemberLifecycle;
import dev.ltms.fleet.config.FleetConfig; import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentControl;
@@ -335,18 +336,31 @@ class SessionManagerTest {
} }
/** /**
* CB-619 / fleetd #123: {@code requireSlotFor} closes the config-gap case (no slot at all * fleetd #226: a contended slot is refused through the real {@link SessionManager#acquire}
* carries the profile) before anything spawns, but a profile that DOES carry a slot can still * path before the real launcher can hand an architect charter to a process.
* lose the bind to a concurrent spawn racing for the same slot. This drives that residual case
* through the REAL path — {@link SessionManager#acquire} against the real {@link
* dev.ltms.fleet.member.ClaudeCodeLauncher} and {@link FakeHerdr} — never {@link
* dev.ltms.fleet.auth.MemberRegistry#bind} directly for the session under test (only the
* precondition uses it, to occupy the slot before the real spawn happens). The session that
* loses the race must be held as a plain {@code dev}, never left claiming {@code architect} in
* the roster, and the daemon log must say so at WARN.
*/ */
@Test
void aContendedArchitectSlotRefusesBeforeTheLauncherStartsAMember() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberRegistry members = architectRegistry();
assertTrue(members.bind("architect:opus", "term_already_bound"),
"precondition: occupy the sole architect slot before the real spawn under test");
sessions.setMemberLifecycle(members);
assertThrows(IllegalArgumentException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller", "term_primary", null));
assertFalse(herdr.called("agent.start"), "a refused slot must never reach the launcher");
assertTrue(sessions.roster().isEmpty(), "no member exists to receive the wrong charter");
}
@Test @Test
void aSecondArchitectOnAnAlreadyBoundProfileIsHeldAsDevNotArchitectAndWarnsLoudly() { void aSecondArchitectOnAnAlreadyBoundProfileIsHeldAsDevNotArchitectAndWarnsLoudly() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberRegistry members = architectRegistry();
sessions.setMemberLifecycle(bindFailureAfterReservation(members));
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory(); LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger registryLog = (ch.qos.logback.classic.Logger) ch.qos.logback.classic.Logger registryLog = (ch.qos.logback.classic.Logger)
LoggerFactory.getLogger(MemberRegistry.class); LoggerFactory.getLogger(MemberRegistry.class);
@@ -356,35 +370,75 @@ class SessionManagerTest {
registryLog.addAppender(appender); registryLog.addAppender(appender);
registryLog.setLevel(Level.WARN); registryLog.setLevel(Level.WARN);
try { try {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberRegistry members = new MemberRegistry(
new FleetConfig.Fleet(Map.of(), Map.of("opus", new FleetConfig.Slot("ltms-local")),
Map.of(), Map.of(), null));
assertTrue(members.bind("architect:opus", "term_already_bound"),
"precondition: occupy the sole architect slot before the real spawn under test");
sessions.setMemberLifecycle(members);
MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null, MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
"/caller", "term_primary", null); "/caller", "term_primary", null);
assertEquals(MemberRole.DEV, session.role(), assertEquals(MemberRole.DEV, session.role(), "a failed reservation bind must use the fallback");
"the slot is taken, so this session must be held as a plain member, never a lie");
assertEquals("dev", SessionManager.rosterView(session, null).get("role"),
"the roster must report what this session actually holds, not what it asked for");
String warn = appender.list.stream() String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN)) .filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage) .map(ILoggingEvent::getFormattedMessage)
.findFirst() .findFirst()
.orElse("no slot-exhaustion WARN logged"); .orElse("no slot-exhaustion WARN logged");
assertTrue(warn.contains("ltms-local"), "the log names the profile: " + warn); assertTrue(warn.contains("ltms-local"), "the WARN names the profile: " + warn);
assertTrue(warn.contains(session.terminalId()), "the log names the terminal: " + warn); assertTrue(warn.contains(session.terminalId()), "the WARN names the terminal: " + warn);
} finally { } finally {
registryLog.detachAppender(appender); registryLog.detachAppender(appender);
} }
} }
@Test
void reservedArchitectSlotBindsTheLaunchedMember() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
sessions.setMemberLifecycle(architectRegistry());
MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
"/caller", "term_primary", null);
assertEquals(MemberRole.ARCHITECT, session.role());
}
private static MemberRegistry architectRegistry() {
return new MemberRegistry(new FleetConfig.Fleet(Map.of(), Map.of("opus", new FleetConfig.Slot("ltms-local")),
Map.of(), Map.of(), null));
}
private static MemberLifecycle bindFailureAfterReservation(MemberRegistry members) {
return new MemberLifecycle() {
@Override
public MemberRole acquired(MemberRole role, String profile, String terminal) {
return members.acquired(role, profile, terminal);
}
@Override
public void released(String terminal) {
members.released(terminal);
}
@Override
public void requireSlotFor(MemberRole role, String profile) {
members.requireSlotFor(role, profile);
}
@Override
public SlotReservation reserve(MemberRole role, String profile) {
return members.reserve(role, profile);
}
@Override
public boolean bind(SlotReservation reservation, String terminal) {
members.release(reservation);
assertTrue(members.bind(reservation.slot(), "term_racer"), "the racer takes the released slot");
return false;
}
@Override
public void release(SlotReservation reservation) {
members.release(reservation);
}
};
}
@Test @Test
void rosterReflectsAcquiredMinusReleased() { void rosterReflectsAcquiredMinusReleased() {
FakeHerdr herdr = new FakeHerdr(); FakeHerdr herdr = new FakeHerdr();
@@ -651,6 +705,32 @@ class SessionManagerTest {
"no session is registered when spawn times out (roster empty)"); "no session is registered when spawn times out (roster empty)");
} }
@Test
void failedArchitectLaunchReleasesItsReservationForTheNextLaunch() {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown");
long[] clock = {0};
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,
1, () -> clock[0], () -> clock[0] += 10);
SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0);
sessions.setMemberLifecycle(architectRegistry());
assertThrows(PeerUnreachableException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller", "term_primary", null));
herdr.agentStatus("idle");
MemberSession retry = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
"/caller", "term_primary", null);
assertEquals(MemberRole.ARCHITECT, retry.role(),
"the failed launch returned its reservation instead of silently shrinking the slot pool");
}
// --- CB-516: release must notify, so a blocked send can be failed -------------------------- // --- CB-516: release must notify, so a blocked send can be failed --------------------------
@Test @Test
@@ -95,10 +95,9 @@ class WorktreeSessionManagerTest {
null, "/caller/proj", null, null); null, "/caller/proj", null, null);
assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId())); assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId()));
MemberSession overflow = sessions.acquire("ltms-local", MemberRole.ARCHITECT, assertThrows(IllegalArgumentException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null); null, "/caller/proj", null, null),
assertNull(members.slotForTerminal(overflow.terminalId()), "a full slot pool must stop the spawn before it can receive an architect charter");
"a full slot pool must not stop the architect spawn");
} }
@Test @Test