fleetd #226: reserve an architect slot before launch, so no member runs on a charter it will not hold
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:
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user