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. */
public interface MemberLifecycle {
/** A slot held before a member process starts. */
record SlotReservation(String slot, String profile) { }
MemberLifecycle NONE = new MemberLifecycle() {
@Override
public MemberRole acquired(MemberRole role, String profile, String terminal) {
@@ -19,6 +22,20 @@ public interface MemberLifecycle {
public void requireSlotFor(MemberRole role, String profile) {
// 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.
* 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
* normal operation this fallback should not happen once {@link #requireSlotFor} has
* refused every unbindable spawn upfront — but a slot can still be lost between that
* check and this call to a concurrent spawn racing for the same slot, so the honest
* answer is still needed here too.
* normal operation this fallback should not happen once a reservation has been bound.
* It remains the honest answer if a caller has no reservation, or if binding a
* reservation unexpectedly fails.
*/
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
*/
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;
/** 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<>();
/** 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. */
public MemberRegistry(FleetConfig.Fleet fleet) {
@@ -167,7 +169,7 @@ public final class MemberRegistry implements MemberLifecycle {
if (existingSlot != null) {
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
}
terminalToSlot.put(terminal, slot);
@@ -275,6 +277,50 @@ public final class MemberRegistry implements MemberLifecycle {
+ " 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. */
@Override
public void released(String terminal) {
@@ -190,49 +190,44 @@ public final class SessionManager implements TurnListener {
if (profile != null && !profile.isBlank()) {
memberLifecycle.requireSlotFor(memberRole, profile);
}
MemberLifecycle.SlotReservation reservation = memberLifecycle.reserve(memberRole, profile);
String launchProfile = reservation == null ? profile : reservation.profile();
if (wt == null) {
// 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
// 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;
boolean bound = false;
try {
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) {
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;
}
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,
sessionName, resumeSessionId);
try {
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,
String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId) {
String sessionName, String resumeSessionId,
MemberLifecycle.SlotReservation reservation) {
String preResolvedProfile = (profile == null || profile.isBlank())
? launcher.defaultProfile() : profile;
// 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();
// 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.
MemberRole actualRole = memberLifecycle.acquired(memberRole, resolvedProfile, handle.terminalId());
MemberRole actualRole = acquired(memberRole, resolvedProfile, handle.terminalId(), reservation);
MemberSession session = new MemberSession(
handle.id(),
handle.terminalId(),
@@ -545,6 +541,18 @@ public final class SessionManager implements TurnListener {
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) {
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.core.read.ListAppender;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.MemberLifecycle;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
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
* carries the profile) before anything spawns, but a profile that DOES carry a slot can still
* 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.
* fleetd #226: a contended slot is refused through the real {@link SessionManager#acquire}
* path before the real launcher can hand an architect charter to a process.
*/
@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
void aSecondArchitectOnAnAlreadyBoundProfileIsHeldAsDevNotArchitectAndWarnsLoudly() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberRegistry members = architectRegistry();
sessions.setMemberLifecycle(bindFailureAfterReservation(members));
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger registryLog = (ch.qos.logback.classic.Logger)
LoggerFactory.getLogger(MemberRegistry.class);
@@ -356,35 +370,75 @@ class SessionManagerTest {
registryLog.addAppender(appender);
registryLog.setLevel(Level.WARN);
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,
"/caller", "term_primary", null);
assertEquals(MemberRole.DEV, session.role(),
"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");
assertEquals(MemberRole.DEV, session.role(), "a failed reservation bind must use the fallback");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.findFirst()
.orElse("no slot-exhaustion WARN logged");
assertTrue(warn.contains("ltms-local"), "the log names the profile: " + warn);
assertTrue(warn.contains(session.terminalId()), "the log names the terminal: " + warn);
assertTrue(warn.contains("ltms-local"), "the WARN names the profile: " + warn);
assertTrue(warn.contains(session.terminalId()), "the WARN names the terminal: " + warn);
} finally {
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
void rosterReflectsAcquiredMinusReleased() {
FakeHerdr herdr = new FakeHerdr();
@@ -651,6 +705,32 @@ class SessionManagerTest {
"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 --------------------------
@Test
@@ -95,10 +95,9 @@ class WorktreeSessionManagerTest {
null, "/caller/proj", null, null);
assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId()));
MemberSession overflow = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null);
assertNull(members.slotForTerminal(overflow.terminalId()),
"a full slot pool must not stop the architect spawn");
assertThrows(IllegalArgumentException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null),
"a full slot pool must stop the spawn before it can receive an architect charter");
}
@Test