diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java index c802a8a..7aa3f39 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java @@ -43,6 +43,7 @@ import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.msg.ReplyInbox; import dev.ltms.fleet.msg.ReplyPushLoop; +import dev.ltms.fleet.peer.MemberRole; import dev.ltms.fleet.peer.PeerLauncher; import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; @@ -50,6 +51,7 @@ import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism; import dev.ltms.fleet.power.IdleSleepGuard; import dev.ltms.fleet.rest.FleetApp; import dev.ltms.fleet.session.GitWorktrees; +import dev.ltms.fleet.session.MemberSession; import dev.ltms.fleet.session.SessionManager; import dev.ltms.fleet.session.SessionReaper; import io.javalin.Javalin; @@ -250,24 +252,44 @@ final class FleetdAssembly { log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values()); } // CB-531/CB-579: discover leads by the tab labels the operator writes, one scanner per - // configured lead's own exact `tab:` label. + // configured lead's own exact `tab:` label. fleetd #669: the same scan also recognises a + // configured collaborator's tab, so one herdr pass answers both. final Supplier> leads; + final Supplier> collaboratorTerminals; var leaders = cfg.fleet().leaders(); - if (!leaders.isEmpty()) { + var collaboratorsConfig = cfg.fleet().collaborators(); + if (!leaders.isEmpty() || !collaboratorsConfig.isEmpty()) { Map tabToName = new LinkedHashMap<>(); leaders.forEach((name, leader) -> { if (leader != null && leader.tab() != null && !leader.tab().isBlank()) { tabToName.put(leader.tab(), name); } }); - int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds(); + Map collaboratorTabToName = new LinkedHashMap<>(); + collaboratorsConfig.forEach((name, collaborator) -> { + if (collaborator != null && collaborator.tab() != null && !collaborator.tab().isBlank()) { + collaboratorTabToName.put(collaborator.tab(), name); + } + }); + // A collaborator-only fleet configures no `leaders:` entry to read a scan interval from + // — FleetConfig.Collaborator carries no scanIntervalSeconds of its own. Falling back to + // FleetConfig.Leader's own compact-constructor default keeps a collaborator-only + // deployment on the same rescan cadence as the default lead cadence, instead of + // inventing a second number for the same kind of scan. + int scanIntervalSeconds = leaders.isEmpty() + ? 10 + : leaders.values().iterator().next().scanIntervalSeconds(); // This must use the lead daemon: scanning member tabs would demote the lead to a worker. - leads = new LeadTabScanner(herdr, tabToName, Set.of(), + LeadTabScanner scanner = new LeadTabScanner(herdr, tabToName, collaboratorTabToName, Set.of(), TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock()); - log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)", - tabToName.keySet(), scanIntervalSeconds); + leads = scanner; + collaboratorTerminals = scanner::collaborators; + log.info("lead/collaborator scan: tabs {} host a lead, tabs {} host a collaborator " + + "(rescan every {}s, shared fleet space)", + tabToName.keySet(), collaboratorTabToName.keySet(), scanIntervalSeconds); } else { leads = () -> leadTerminals; + collaboratorTerminals = Map::of; } leadsRef.set(leads); @@ -452,6 +474,15 @@ final class FleetdAssembly { ConnectionIdentity identity = new ConnectionIdentity( new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); + // fleetd #669 Unit D: a live spawned member resolves as its own role, whatever a tab map + // says about the same terminal — read from the roster meant for a hot path (SessionManager + // javadoc), never rosterResolved(), since resolve() runs on every request. + Function spawnedMemberRole = terminal -> sessions.roster().stream() + .filter(s -> terminal.equals(s.terminalId())) + .map(MemberSession::role) + .findFirst() + .orElse(null); + // CB-501: one resolver behind both entry paths. Worker identity still comes from the // connection and is never token-gated, so enabling token mode cannot lock the fleet out. final CallerResolver callers; @@ -461,11 +492,13 @@ final class FleetdAssembly { throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv() + " is unset or empty — export it before starting fleetd"); } - callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members); + callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members, + spawnedMemberRole, collaboratorTerminals); log.info("auth: token mode (bearer required for non-worker callers, env {})", cfg.auth().tokenEnv()); } else { - callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members); + callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members, + spawnedMemberRole, collaboratorTerminals); log.info("auth: loopback-trust (any loopback non-worker caller is the primary)"); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java index 355d5b4..d59d9e8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java @@ -62,11 +62,11 @@ public final class Authz { } /** - * Stands in for the terminal-to-tab registry a collaborator's {@code SEND} is checked - * against, until one exists: answers no for every target, so a collaborator reaches nothing - * today. Both production gates ({@code FleetMcp#denyFor}, {@code FleetApp#allow}) pass this - * exact instance, so the classifier is defined once and replacing it is a one-line change in - * each. + * The fail-closed classifier: answers no for every target, so a collaborator's {@code SEND} + * is refused unless a caller supplies a real one. {@code CallerResolver#knownLeadOrCollaborator()} + * is the real one, read from the same lead and collaborator maps {@code CallerResolver#resolve} + * consults, so a target that classifier calls known is one {@code resolve} would actually + * resolve as a lead or collaborator. */ public static final Predicate NO_KNOWN_LEAD_OR_COLLABORATOR = target -> false; diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java b/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java index 9728510..4988d08 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java @@ -7,6 +7,7 @@ import java.nio.charset.StandardCharsets; import java.security.MessageDigest; import java.util.Map; import java.util.function.Function; +import java.util.function.Predicate; import java.util.function.Supplier; /** @@ -19,6 +20,13 @@ import java.util.function.Supplier; * *

Resolution order — connection identity first, token second, nothing third: *

    + *
  1. A loopback peer PID that maps to a pane this gateway itself spawned ⇒ that member's own + * role: {@link Role#WORKER} for a dev, hunter, or reviewer; {@link Role#ARCHITECT} for an + * architect, but only while the live slot role still confirms it (fleetd #424 — a slot + * revoked from config demotes an already-bound session on its very next request, so the + * roster's own role is never granted on its word alone). No tab map is consulted — a live + * spawned member's identity comes from the registry that spawned it, never from a label a + * pane could also carry.
  2. *
  3. A loopback peer PID that maps to a pane named by {@code leaders:}, by the legacy * {@code primary.terminal} pin, or by an operator-labelled lead tab (CB-307, CB-530, CB-531) * ⇒ {@link Role#PRIMARY}, carrying that lead's @@ -28,8 +36,10 @@ import java.util.function.Supplier; * so two leads can work as peers rather than one being demoted.
  4. *
  5. A loopback peer PID that maps to a pane bound to a CB-548 architect slot ⇒ * {@link Role#ARCHITECT}, carrying the slot name. Just unforgeable as a worker's, and - * resolved from the live terminal→slot binding (never a request argument), before - * the generic worker fallback.
  6. + * resolved from the live terminal→slot binding (never a request argument). This is + * the case the previous step does not catch: a binding with no live spawned-member session. + *
  7. A loopback peer PID that maps to an operator-labelled collaborator tab ⇒ + * {@link Role#COLLABORATOR}, carrying that collaborator's name.
  8. *
  9. A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#WORKER}. This is * unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured * regardless of auth mode, so enabling auth never breaks the fleet.
  10. @@ -64,6 +74,20 @@ public final class CallerResolver { private final Supplier> architectTerminals; private final Function memberSlotRoles; private final Function memberSlotNames; + /** + * terminal_id → the role of the live spawned member occupying it, or {@code null} for a + * terminal no spawned member occupies. Consulted first, ahead of every tab map: a live + * spawned member's identity is its own, whatever a tab map says about the same terminal. + * A function rather than the roster itself, so a resolve on the hot path never scans a list — + * the lookup strategy is the caller's to choose. + */ + private final Function spawnedMemberRole; + /** + * terminal_id → collaborator name; empty when none are configured. A supplier for the same + * reason as {@link #leadTerminals}: a collaborator tab recognised after construction (the tab + * scan discovering a newly-labelled tab) takes effect without a restart. + */ + private final Supplier> collaboratorTerminals; /** Loopback-trust resolver: no token required, historical behaviour. Test-only. */ CallerResolver(ConnectionIdentity identity) { @@ -124,17 +148,42 @@ public final class CallerResolver { /** * Live registry form that can confirm a bound slot is an architect slot. * - *

    This is the only public construction path. It keeps terminal bindings and slot roles in - * the same {@link MemberRegistry}, so a configured architect can resolve as an architect. + *

    It keeps terminal bindings and slot roles in the same {@link MemberRegistry}, so a + * configured architect can resolve as an architect. No spawned-member roster or collaborator + * registry is consulted — equivalent to {@link #withLeadsAndMembers(ConnectionIdentity, + * boolean, String, Supplier, MemberRegistry, Function, Supplier)} with both absent. Kept for + * every caller that has neither to offer, so adding them did not churn every construction site. */ public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity, boolean tokenMode, String token, Supplier> leadTerminals, MemberRegistry members) { + return withLeadsAndMembers(identity, tokenMode, token, leadTerminals, members, null, null); + } + + /** + * Live registry form that also resolves a live spawned member to its own role, and a + * configured collaborator tab to {@link Role#COLLABORATOR}. + * + *

    This is the only public construction path that exercises the full resolution order. + * + * @param spawnedMemberRole terminal_id → the role of the live spawned member occupying + * it, or {@code null} for a terminal no spawned member occupies. + * {@code null} here means no roster is consulted at all (every + * terminal falls through to the tab maps), not that none matches. + * @param collaboratorTerminals terminal_id → collaborator name, live like {@code leadTerminals} + */ + public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity, + boolean tokenMode, String token, + Supplier> leadTerminals, + MemberRegistry members, + Function spawnedMemberRole, + Supplier> collaboratorTerminals) { return new CallerResolver(identity, tokenMode, token, leadTerminals, members == null ? null : members::snapshot, members == null ? null : members::roleForSlot, - members == null ? null : members::nameForSlot); + members == null ? null : members::nameForSlot, + spawnedMemberRole, collaboratorTerminals); } private static Supplier> fixed(Map leadTerminals) { @@ -160,6 +209,17 @@ public final class CallerResolver { Supplier> architectTerminals, Function memberSlotRoles, Function memberSlotNames) { + this(identity, tokenMode, token, leadTerminals, architectTerminals, memberSlotRoles, + memberSlotNames, null, null); + } + + private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token, + Supplier> leadTerminals, + Supplier> architectTerminals, + Function memberSlotRoles, + Function memberSlotNames, + Function spawnedMemberRole, + Supplier> collaboratorTerminals) { if (tokenMode && (token == null || token.isBlank())) { throw new IllegalArgumentException( "auth.mode=token requires a non-empty token; check that the env var named by " @@ -172,6 +232,8 @@ public final class CallerResolver { this.architectTerminals = architectTerminals == null ? Map::of : architectTerminals; this.memberSlotRoles = memberSlotRoles == null ? _ -> null : memberSlotRoles; this.memberSlotNames = memberSlotNames == null ? Function.identity() : memberSlotNames; + this.spawnedMemberRole = spawnedMemberRole == null ? _ -> null : spawnedMemberRole; + this.collaboratorTerminals = collaboratorTerminals == null ? Map::of : collaboratorTerminals; } /** @@ -198,6 +260,27 @@ public final class CallerResolver { return architectTerminals.get(); } + /** + * The currently-recognised collaborator tabs, {@code terminal_id → name}. + * + *

    Read from the same supplier {@link #resolve} consults, for the reason given in + * {@link #leads()}. Live for the same reason as {@link #leads()}. + */ + public Map collaborators() { + return collaboratorTerminals.get(); + } + + /** + * Whether {@code target} names a terminal this resolver would resolve as a lead or a + * collaborator — the classifier a collaborator's {@code SEND} is checked against, read from the + * exact maps {@link #resolve} consults so a target that would resolve as a lead or collaborator + * is never the one a collaborator is refused to reach, or the reverse. + */ + public Predicate knownLeadOrCollaborator() { + return target -> leadTerminals.get().containsKey(target) + || collaboratorTerminals.get().containsKey(target); + } + /** * Resolve the caller of a request. * @@ -208,6 +291,25 @@ public final class CallerResolver { public Principal resolve(String remoteAddr, int remotePort, String authorizationHeader) { ConnectionIdentity.Caller c = identity.resolve(remoteAddr, remotePort); if (c.terminal() != null) { + MemberRole spawnedRole = spawnedMemberRole.apply(c.terminal()); + if (spawnedRole != null) { + // A live spawned member occupies this pane. Its identity is its own, whatever a tab + // map says about the same terminal — checked before every tab map, consulting none + // of them, so a tab label can never override a roster entry for the same terminal. + if (spawnedRole == MemberRole.ARCHITECT) { + // The roster only answers THAT this pane is a live spawned member; config still + // decides WHAT that member's slot grants (fleetd #424). A slot revoked after the + // bind must still demote this session on its very next request, so the roster's + // own ARCHITECT role is confirmed against the live slot role, exactly as the + // architect-slot step below confirms a binding with no live member session. + String slot = architectTerminals.get().get(c.terminal()); + if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) { + return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid()); + } + return Principal.worker(c.terminal(), c.pid()); + } + return Principal.worker(c.terminal(), c.pid()); + } String lead = leadTerminals.get().get(c.terminal()); if (lead != null) { // The config names this pane as a lead's own. The pane mapping is exactly as @@ -221,10 +323,17 @@ public final class CallerResolver { // The config/live binding names this pane as an architect slot's own. Same // unforgeable pane mapping; the live binding, never a request argument, decides. // Check the slot role too: this defence in depth prevents a bad lifecycle bind from - // escalating a dev, hunter or reviewer into an architect. Checked before - // the worker fallback. + // escalating a dev, hunter or reviewer into an architect. This is the case the + // spawned-member step above does not catch: a binding with no live member session. return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid()); } + String collaborator = collaboratorTerminals.get().get(c.terminal()); + if (collaborator != null) { + // An operator-labelled collaborator tab, confirmed live by the same scan that + // confirms a lead tab. Checked last among the tab maps so a pane also matching one + // of the above keeps that stronger role. + return Principal.collaborator(collaborator, c.terminal(), c.pid()); + } return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated } diff --git a/fleetd/src/main/java/dev/ltms/fleet/herdr/LeadTabScanner.java b/fleetd/src/main/java/dev/ltms/fleet/herdr/LeadTabScanner.java index 39eaa88..b86fac9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/herdr/LeadTabScanner.java +++ b/fleetd/src/main/java/dev/ltms/fleet/herdr/LeadTabScanner.java @@ -96,13 +96,19 @@ public final class LeadTabScanner implements Supplier> { private static final Logger log = LoggerFactory.getLogger(LeadTabScanner.class); + /** What a matched tab names: a lead or a collaborator. */ + private enum Kind { LEAD, COLLABORATOR } + + /** One matched tab's name and what it names. */ + private record Entry(String name, Kind kind) {} + private final HerdrClient herdr; - private final Map tabToName; + private final Map tabToEntry; private final Set excludedWorkspaceLabels; private final long ttlNanos; private final LongSupplier clock; - private Map cached = Map.of(); + private Map cached = Map.of(); private long scannedAtNanos; private boolean everScanned; @@ -126,26 +132,53 @@ public final class LeadTabScanner implements Supplier> { */ public LeadTabScanner(HerdrClient herdr, Map tabToName, Set excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) { + this(herdr, tabToName, Map.of(), excludedWorkspaceLabels, ttlNanos, clock); + } + + /** + * As {@link #LeadTabScanner(HerdrClient, Map, Set, long, LongSupplier)}, additionally scanning + * for configured collaborator tabs in the same pass. + * + * @param collaboratorTabToName every configured collaborator's exact tab label → its name + * ({@code fleet.collaborators..tab}), matched the same way as + * {@code tabToName} + */ + public LeadTabScanner(HerdrClient herdr, Map tabToName, + Map collaboratorTabToName, + Set excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) { this.herdr = herdr; - this.tabToName = normalize(tabToName); + this.tabToEntry = buildTabIndex(tabToName, collaboratorTabToName); this.excludedWorkspaceLabels = excludedWorkspaceLabels == null ? Set.of() : Set.copyOf(excludedWorkspaceLabels); this.ttlNanos = ttlNanos; this.clock = clock; } - /** Keys stripped and lower-cased once, so every lookup is a plain map hit. */ - private static Map normalize(Map tabToName) { - if (tabToName == null || tabToName.isEmpty()) { - return Map.of(); + /** + * Keys stripped and lower-cased once, so every lookup is a plain map hit. Leads and + * collaborators merge into a single index, so {@link #scan()} matches both kinds in one pass + * over the tab list; a label naming both a lead and a collaborator takes the lead entry — + * leads are put last, so a colliding key's lead entry is the one that overwrites — since a lead + * can already do everything a collaborator can. Config validation already refuses a lead and a + * collaborator sharing one exact tab, so this ordering is defence in depth, not the control. + */ + private static Map buildTabIndex(Map tabToName, + Map collaboratorTabToName) { + Map out = new LinkedHashMap<>(); + putNormalized(out, collaboratorTabToName, Kind.COLLABORATOR); + putNormalized(out, tabToName, Kind.LEAD); + return Collections.unmodifiableMap(out); + } + + private static void putNormalized(Map out, Map tabToName, Kind kind) { + if (tabToName == null) { + return; } - Map out = new LinkedHashMap<>(); tabToName.forEach((tab, name) -> { if (tab != null && !tab.isBlank() && name != null && !name.isBlank()) { - out.put(tab.strip().toLowerCase(Locale.ROOT), name); + out.put(tab.strip().toLowerCase(Locale.ROOT), new Entry(name, kind)); } }); - return Collections.unmodifiableMap(out); } /** @@ -156,6 +189,29 @@ public final class LeadTabScanner implements Supplier> { */ @Override public synchronized Map get() { + return byKind(refresh(), Kind.LEAD); + } + + /** + * The current {@code terminal_id → collaborator name} map, sharing the same scan and cache as + * {@link #get()} — both kinds are matched in one pass, so this never costs a second herdr call. + */ + public synchronized Map collaborators() { + return byKind(refresh(), Kind.COLLABORATOR); + } + + private static Map byKind(Map entries, Kind kind) { + Map out = new LinkedHashMap<>(); + entries.forEach((terminal, entry) -> { + if (entry.kind() == kind) { + out.put(terminal, entry.name()); + } + }); + return Collections.unmodifiableMap(out); + } + + /** Rescans if the cache has expired, otherwise returns the cached answer. */ + private Map refresh() { long now = clock.getAsLong(); if (everScanned && now - scannedAtNanos < ttlNanos) { return cached; @@ -165,21 +221,21 @@ public final class LeadTabScanner implements Supplier> { scannedAtNanos = now; everScanned = true; try { - Map fresh = scan(); + Map fresh = scan(); if (!fresh.equals(cached)) { - log.info("lead panes: {}", fresh); + log.info("lead/collaborator panes: {}", fresh); } cached = fresh; } catch (HerdrException e) { - log.warn("lead-tab scan failed, keeping the {} lead(s) already known: {}", + log.warn("lead-tab scan failed, keeping the {} entr(y/ies) already known: {}", cached.size(), e.getMessage()); } return cached; } /** One full pass: labelled tabs → live agents in them → those panes' terminals. */ - private Map scan() { - Map nameByTab = new LinkedHashMap<>(); + private Map scan() { + Map entryByTab = new LinkedHashMap<>(); for (JsonNode w : herdr.call("workspace.list").path("workspaces")) { Workspace ws = Workspace.from(w); if (ws.workspaceId() == null || excludedWorkspaceLabels.contains(ws.label())) { @@ -187,21 +243,22 @@ public final class LeadTabScanner implements Supplier> { } for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", ws.workspaceId())).path("tabs")) { Tab tab = Tab.from(t); - String name = leadNameOf(tab.label()); - if (name != null && tab.tabId() != null) { - nameByTab.put(tab.tabId(), name); + Entry entry = entryOf(tab.label()); + if (entry != null && tab.tabId() != null) { + entryByTab.put(tab.tabId(), entry); } } } - if (nameByTab.isEmpty()) { + if (entryByTab.isEmpty()) { gracedTerminals = Set.of(); return Map.of(); } - // fleetd #359: a labelled tab is only a lead when herdr also reports a running agent in - // it — the same liveness signal LeadLauncher.countLeads trusts for the identical purpose. - // Without this, a tab left behind by a session that has since died reads as live forever. + // fleetd #359: a labelled tab is only a lead (or collaborator) when herdr also reports a + // running agent in it — the same liveness signal LeadLauncher.countLeads trusts for the + // identical purpose. Without this, a tab left behind by a session that has since died reads + // as live forever. Set tabsWithAgent = new HashSet<>(); for (JsonNode a : herdr.call("agent.list").path("agents")) { String tabId = a.path("tab_id").asText(null); @@ -210,18 +267,18 @@ public final class LeadTabScanner implements Supplier> { } } - Map byTerminal = new LinkedHashMap<>(); + Map byTerminal = new LinkedHashMap<>(); Set stillGraced = new HashSet<>(); // One pane.list for every tab: panes carry tab_id, so the join is local. for (JsonNode p : herdr.call("pane.list", Map.of()).path("panes")) { String tabId = p.path("tab_id").asText(null); - String name = nameByTab.get(tabId); + Entry entry = entryByTab.get(tabId); String terminal = p.path("terminal_id").asText(null); - if (name == null || terminal == null || terminal.isBlank()) { + if (entry == null || terminal == null || terminal.isBlank()) { continue; } if (tabsWithAgent.contains(tabId)) { - byTerminal.put(terminal, name); + byTerminal.put(terminal, entry); continue; } // No agent reported for this tab, but its tab/pane are still here — this is the @@ -230,7 +287,7 @@ public final class LeadTabScanner implements Supplier> { // reported as live; a terminal we never reported live gets none, so the original #359 // fix (a genuinely dead tab is never reported) is unaffected for the common case. if (cached.containsKey(terminal) && !gracedTerminals.contains(terminal)) { - byTerminal.put(terminal, name); + byTerminal.put(terminal, entry); stillGraced.add(terminal); } } @@ -239,18 +296,19 @@ public final class LeadTabScanner implements Supplier> { } /** - * The lead name a tab label declares, or {@code null} if it names none of the configured leads. + * The entry a tab label declares, or {@code null} if it names neither a configured lead nor a + * configured collaborator. * - *

    Exact match (case-insensitive, ends stripped) against {@link #tabToName} — no prefix + *

    Exact match (case-insensitive, ends stripped) against {@link #tabToEntry} — no prefix * stripping, so an operator's {@code "lead: something-else"} tab is never mistaken for a * configured lead just because it shares a prefix. The match strips a trailing * {@link PendingCloseMarker} first, so a tab {@code LeadLauncher} has flagged as maybe-dead but * not yet closed keeps resolving normally while that reconcile is pending. */ - private String leadNameOf(String label) { + private Entry entryOf(String label) { if (label == null) { return null; } - return tabToName.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT)); + return tabToEntry.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT)); } } 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 7d05cc9..cb3d68d 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -114,6 +114,13 @@ public final class FleetMcp { * dev.ltms.fleet.herdr.PaneLocator} — that a real assembly wired up. See {@link #identity()}. */ private final ConnectionIdentity identity; + /** + * Kept as a field (rather than only captured by the {@code contextExtractor} closure) so + * {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} — the classifier a + * collaborator's {@code SEND} is checked against, built from the same lead and collaborator + * maps {@link #identity}-based resolution reads. + */ + private final CallerResolver callers; private final Metrics metrics; // CB-502: null → auth failures not counted private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; @@ -409,6 +416,7 @@ public final class FleetMcp { this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode") == AuthorizationMode.ENFORCED; this.identity = identity; + this.callers = callers; this.leadChannel = leadChannel; this.peers = peers == null ? List.of() : List.copyOf(peers); this.capacity = capacity; @@ -689,7 +697,7 @@ public final class FleetMcp { if (!authorizationEnforced) { return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518) } - if (Authz.permits(caller, action, target, Authz.NO_KNOWN_LEAD_OR_COLLABORATOR)) { + if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator())) { if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) { AuditLog.allowed(caller, action, target); // reads would drown the trail } diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 9ae88c3..824f5ff 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -269,12 +269,16 @@ public final class FleetApp { /** * The authorization decision behind {@link #allow}, taking the caller directly rather than * pulling it from a servlet {@link Context} — unit-testable without fabricating a live - * request, the same reason {@code FleetMcp#denyFor} is split from {@code FleetMcp#deny}. The - * classifier is the same shared instance {@link #allow} would use, so a test calling this - * exercises the real production gate, not a test-supplied stand-in. + * request, the same reason {@code FleetMcp#denyFor} is split from {@code FleetMcp#deny}. + * + * @param knownLeadOrCollaborator the classifier a collaborator's {@code SEND} is checked + * against; pass {@link #auth}'s own {@code + * knownLeadOrCollaborator()} to exercise the real production + * gate, as {@link #allow} does */ - static boolean permitsFor(Principal caller, Authz.Action action, String target) { - return Authz.permits(caller, action, target, Authz.NO_KNOWN_LEAD_OR_COLLABORATOR); + static boolean permitsFor(Principal caller, Authz.Action action, String target, + Predicate knownLeadOrCollaborator) { + return Authz.permits(caller, action, target, knownLeadOrCollaborator); } /** @@ -290,7 +294,7 @@ public final class FleetApp { return true; // legacy: authorization not enforced } Principal caller = ctx.attribute(CALLER); - if (permitsFor(caller, action, target)) { + if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator())) { if (action != Authz.Action.READ && action != Authz.Action.METRICS && action != Authz.Action.TASK_READ) { AuditLog.allowed(caller, action, target); // reads would drown the trail diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/CallerResolverTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/CallerResolverTest.java index 5d574fe..6546656 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/auth/CallerResolverTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/CallerResolverTest.java @@ -523,4 +523,129 @@ class CallerResolverTest { } } + // ── fleetd #669 Unit D: a live spawned member outranks every tab map ──────────────────────── + + /** + * Criterion 1: a terminal present in BOTH the spawned-member roster AND the lead tab map + * resolves as its member role, not as a lead — the roster is checked first, consulting no tab + * map at all when it matches. + */ + @Test + void aSpawnedMemberWinsOverALeadTabForTheSamePane() { + Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, + () -> Map.of("term_a", "opus-5.0"), new MemberRegistry(null), + t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.WORKER, p.role(), + "a live spawned member's own identity must win over a tab map naming the same pane a lead"); + assertEquals("term_a", p.terminal()); + } + + /** A spawned architect in the roster resolves ARCHITECT, carrying its bound slot's name. */ + @Test + void aSpawnedArchitectInTheRosterResolvesArchitectWithItsSlotName() { + Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + boundMembers("architect:lead-designer", MemberRole.ARCHITECT), + t -> "term_a".equals(t) ? MemberRole.ARCHITECT : null, Map::of) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.ARCHITECT, p.role()); + assertEquals("lead-designer", p.name()); + assertEquals("term_a", p.terminal()); + } + + /** + * fleetd #424 regression: the roster only answers THAT a pane is a live spawned member; config + * still decides WHAT that member's slot grants. A slot revoked after the bind must still demote + * the session on its very next request, exactly as it would for a pane with no roster entry at + * all — the roster's own ARCHITECT role must never be granted on its word alone. + * + *

    {@code bind} refuses an unconfigured slot, so the revoked state can only be reached by + * binding while the slot is configured and then swapping the config out from under it, the way + * a live reload does. + */ + @Test + void aRevokedArchitectSlotDemotesALiveSpawnedArchitectToWorker() { + FleetConfig.Fleet configured = new FleetConfig.Fleet(Map.of(), + Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null); + java.util.concurrent.atomic.AtomicReference live = + new java.util.concurrent.atomic.AtomicReference<>(configured); + MemberRegistry members = MemberRegistry.live(live::get); + assertTrue(members.bind("architect:lead-designer", "term_a")); + + live.set(new FleetConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), null)); // slot revoked + + // Setup controls: the slot is really gone from config, but the occupancy is still there — + // otherwise this test would pass for the wrong reason. + assertNull(members.roleForSlot("architect:lead-designer"), "setup control: the slot must be gone from config"); + assertEquals("architect:lead-designer", members.snapshot().get("term_a"), + "setup control: the binding itself must still be there"); + + Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + members, t -> "term_a".equals(t) ? MemberRole.ARCHITECT : null, Map::of) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.WORKER, p.role(), + "a revoked slot must demote a live spawned architect on its very next request"); + } + + /** Criterion 3: a configured collaborator tab that is not a spawned member resolves COLLABORATOR. */ + @Test + void aConfiguredCollaboratorTabResolvesToCollaboratorCarryingItsName() { + Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + new MemberRegistry(null), t -> null, () -> Map.of("term_a", "ops")) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.COLLABORATOR, p.role()); + assertEquals("ops", p.name()); + assertEquals("term_a", p.terminal()); + } + + /** Regression: an empty collaborator registry leaves every pane exactly as before. */ + @Test + void anEmptyCollaboratorRegistryLeavesEveryPaneAsBefore() { + Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + new MemberRegistry(null), t -> null, Map::of) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.WORKER, p.role()); + assertNull(p.name()); + } + + @Test + void describeNamesTheCollaborator() { + assertEquals("collaborator:ops", Principal.collaborator("ops", "term_a", 1).describe()); + } + + // ── fleetd #669 Unit D: knownLeadOrCollaborator() reads the same maps resolve() does ─────────── + + @Test + void knownLeadOrCollaboratorIsTrueForALeadTerminal() { + CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, + () -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null), t -> null, Map::of); + + assertTrue(r.knownLeadOrCollaborator().test("term_lead")); + assertFalse(r.knownLeadOrCollaborator().test("term_other")); + } + + @Test + void knownLeadOrCollaboratorIsTrueForACollaboratorTerminal() { + CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + new MemberRegistry(null), t -> null, () -> Map.of("term_collab", "ops")); + + assertTrue(r.knownLeadOrCollaborator().test("term_collab")); + assertFalse(r.knownLeadOrCollaborator().test("term_other")); + } + + @Test + void knownLeadOrCollaboratorIsFalseForASpawnedMembersTerminal() { + // The exact scenario a collaborator's SEND must never reach: a live spawned member's own + // terminal, which is neither a configured lead nor a configured collaborator. + CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of, + new MemberRegistry(null), t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of); + + assertFalse(r.knownLeadOrCollaborator().test("term_a")); + } + } diff --git a/fleetd/src/test/java/dev/ltms/fleet/herdr/LeadTabScannerTest.java b/fleetd/src/test/java/dev/ltms/fleet/herdr/LeadTabScannerTest.java index 1c46f19..86d90c6 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/herdr/LeadTabScannerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/herdr/LeadTabScannerTest.java @@ -163,6 +163,13 @@ class LeadTabScannerTest { return new LeadTabScanner(herdr, tabToName, Set.of("fleetd-workers"), TTL, clock::get); } + private LeadTabScanner scannerWithCollaborators(TopologyHerdr herdr, Map tabToName, + Map collaboratorTabToName, + AtomicLong clock) { + return new LeadTabScanner(herdr, tabToName, collaboratorTabToName, + Set.of("fleetd-workers"), TTL, clock::get); + } + @Test void everyConfiguredTabBecomesALeadNamedByItsEntry() { Map leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get(); @@ -464,4 +471,75 @@ class LeadTabScannerTest { assertEquals(afterFirst, herdr.calls, "the failure path must be rate-limited too"); } + + // ── fleetd #669 Unit D: a collaborator tab is matched the same way as a lead tab, one pass ───── + + @Test + void aConfiguredCollaboratorTabIsReportedByCollaboratorsNotByGet() { + TopologyHerdr herdr = new TopologyHerdr() + .workspace("w1", "main") + .tab("w1:t1", "w1", "collab: ops") + .pane("w1:p1", "w1:t1", "term_ops"); + LeadTabScanner s = scannerWithCollaborators(herdr, Map.of(), Map.of("collab: ops", "ops"), + new AtomicLong()); + + assertEquals(Map.of("term_ops", "ops"), s.collaborators(), + "a collaborator tab is matched exactly like a lead tab"); + assertEquals(Map.of(), s.get(), "a collaborator tab must never also appear as a lead"); + } + + /** + * Criterion 4, scanner level: a collaborator tab that is labelled but runs no agent is not + * reported — the same #359 liveness cross-check a lead tab gets. + */ + @Test + void aDeadCollaboratorTabIsNotReported() { + TopologyHerdr herdr = new TopologyHerdr() + .workspace("w1", "main") + .tab("w1:t1", "w1", "collab: ops") + .pane("w1:p1", "w1:t1", "term_ops") + .deadAgent("w1:t1"); + LeadTabScanner s = scannerWithCollaborators(herdr, Map.of(), Map.of("collab: ops", "ops"), + new AtomicLong()); + + assertFalse(s.collaborators().containsKey("term_ops"), + "a dead collaborator tab must never resolve as a live collaborator"); + } + + /** + * Both kinds are matched in a single pass over the same tab list — not a second scanner, not a + * second scan. Proven by herdr call count: scanning one lead tab and one collaborator tab in the + * same instance costs exactly as many calls as scanning two lead tabs in {@link #twoLeads()}. + */ + @Test + void leadsAndCollaboratorsAreMatchedInOnePassOverTheSameScan() { + TopologyHerdr oneOfEach = new TopologyHerdr() + .workspace("w1", "main") + .tab("w1:t1", "w1", "lead: opus-5.0") + .tab("w1:t2", "w1", "collab: ops") + .pane("w1:p1", "w1:t1", "term_opus") + .pane("w1:p2", "w1:t2", "term_ops"); + LeadTabScanner s = scannerWithCollaborators(oneOfEach, Map.of("lead: opus-5.0", "opus-5.0"), + Map.of("collab: ops", "ops"), new AtomicLong()); + + assertEquals(Map.of("term_opus", "opus-5.0"), s.get()); + assertEquals(Map.of("term_ops", "ops"), s.collaborators()); + + TopologyHerdr twoLeadsBaseline = twoLeads(); + scanner(twoLeadsBaseline, twoLeadsConfigured(), new AtomicLong()).get(); + + assertEquals(twoLeadsBaseline.calls, oneOfEach.calls, + "one lead tab + one collaborator tab must cost exactly as many herdr calls as two " + + "lead tabs — proof this is one pass, not a second scan"); + } + + /** Regression: with no collaborators configured, every existing lead-only behaviour is unchanged. */ + @Test + void anEmptyCollaboratorMapLeavesCollaboratorsEmptyAndGetUnaffected() { + LeadTabScanner s = scannerWithCollaborators(twoLeads(), twoLeadsConfigured(), Map.of(), + new AtomicLong()); + + assertEquals(Map.of(), s.collaborators()); + assertEquals(Map.of("term_opus", "opus-5.0", "term_gpt", "gpt-sol-5.6"), s.get()); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java index 47e56fa..aeb3e75 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java @@ -63,6 +63,17 @@ class FleetMcpAuthzTest { /** A fully wired FleetMcp on fakes — constructing it is itself part of what is under test. */ private FleetMcp mcp(boolean enforce) { + ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999); + return mcp(enforce, CallerResolver.withLeadsAndMembers(identity, false, null, + Map::of, new MemberRegistry(null))); + } + + /** + * As {@link #mcp(boolean)}, with an explicit {@link CallerResolver} — so a test can wire known + * leads/collaborators and drive {@code denyFor}'s real {@code knownLeadOrCollaborator()} + * classifier instead of the default empty one. + */ + private FleetMcp mcp(boolean enforce, CallerResolver callers) { FleetConfig.Profile cfg = new FleetConfig.Profile( "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null, "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null); @@ -83,9 +94,7 @@ class FleetMcpAuthzTest { // to be omitted to reach "legacy" is now always real, and AuthorizationMode is the // separate, explicit choice that governs enforcement. mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(), - new PrimaryRegistry(null), - CallerResolver.withLeadsAndMembers(identity, false, null, - Map::of, new MemberRegistry(null)), + new PrimaryRegistry(null), callers, enforce ? FleetMcp.AuthorizationMode.ENFORCED : FleetMcp.AuthorizationMode.UNENFORCED, metrics, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(), @@ -240,6 +249,26 @@ class FleetMcpAuthzTest { assertTrue(denied.isError()); } + /** + * fleetd #669 Unit D: wires a real {@link CallerResolver} with a known lead and a known + * collaborator tab, and leaves a spawned member's own terminal recognised by neither map — so a + * collaborator's SEND reaches both named peers and is refused for the spawned member's terminal, + * over MCP's {@code denyFor}. + */ + @Test + void aCollaboratorMaySendToAKnownLeadOrCollaboratorButNotToASpawnedMembersTerminal() { + ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999); + CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, + () -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null), + t -> null, () -> Map.of("term_collab_known", "ops2")); + FleetMcp m = mcp(true, callers); + + assertNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_lead_known")); + assertNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_collab_known")); + assertNotNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_a"), + "a spawned member's own terminal must stay unreachable, even once the classifier is real"); + } + @Test void theLegacyConstructorLeavesTheGateOpen() { // The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced. diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java index 21f8f8a..735c116 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -205,16 +205,72 @@ class FleetAppAuthTest { /** * {@code permitsFor} is the exact decision {@link FleetApp#allow} makes, passing the real - * production classifier rather than a test-supplied one — no terminal is recognised as a - * configured lead or collaborator, so a collaborator's SEND is refused through the REST gate. + * production classifier rather than a test-supplied one — built from an empty {@link + * CallerResolver}, so no terminal is recognised as a configured lead or collaborator and a + * collaborator's SEND is refused through the REST gate. */ @Test void aCollaboratorMayNotSendOverRestWithTheRealProductionClassifier() { + ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(new FakeHerdr()), _ -> 700L); + CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, + Map::of, new MemberRegistry(null)); Principal collaborator = Principal.collaborator("ops", "term_collab", 700); - assertFalse(FleetApp.permitsFor(collaborator, Authz.Action.SEND, "term_lead"), + assertFalse(FleetApp.permitsFor(collaborator, Authz.Action.SEND, "term_lead", + callers.knownLeadOrCollaborator()), "no terminal is recognised as a lead or collaborator yet"); } + /** + * fleetd #669 Unit D: wires a real {@link CallerResolver} with a known lead and a known + * collaborator tab, and a spawned member's own terminal recognised by neither map. A + * collaborator's SEND reaches the known lead and the known collaborator, and is refused for the + * spawned member's terminal — over the REST route, not just the unit-level classifier, so a + * test covering only MCP cannot leave this route open. + */ + @Test + void aCollaboratorMaySendToAKnownLeadOrCollaboratorButNotToASpawnedMembersTerminalOverRest() throws Exception { + int port = startWithRealClassifier(FakeHerdr.WORKER_PID, + Map.of("term_lead_known", "lead-x"), Map.of("term_a", "ops2")); + + HttpResponse toLead = send(port, "POST", "/sessions/term_lead_known/message", + "{\"content\":\"hi\",\"wait\":false}", null); + assertEquals(202, toLead.statusCode(), toLead.body()); + + HttpResponse toSpawnedMembersTerminal = send(port, "POST", "/sessions/term_worker/message", + "{\"content\":\"hi\",\"wait\":false}", null); + assertEquals(403, toSpawnedMembersTerminal.statusCode(), toSpawnedMembersTerminal.body()); + } + + /** + * As {@link #start}, but with explicit lead/collaborator maps and no spawned-member roster, so + * a test can wire the real {@link CallerResolver#knownLeadOrCollaborator()} classifier instead + * of the default empty one. + */ + private int startWithRealClassifier(long pid, Map leadTerminals, + Map collaboratorTerminals) { + FakeHerdr herdr = new FakeHerdr(); + FleetConfig.Profile wcfg = new FleetConfig.Profile( + "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null, + "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null); + AgentControl agents = new AgentControl(herdr); + ClaudeCodeLauncher workers = new ClaudeCodeLauncher( + agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), + Map.of(wcfg.profile(), wcfg), wcfg.profile(), + k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null); + SessionManager sessions = new SessionManager(workers, new FakeWorktrees()); + Injector injector = new Injector(agents); + MessageService messages = new MessageService(agents, injector, new Rendezvous()); + + ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid); + CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, + () -> leadTerminals, new MemberRegistry(null), t -> null, () -> collaboratorTerminals); + metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox()); + + app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null, + callers, metrics).build().start("127.0.0.1", 0); + return app.port(); + } + // --- loopback-trust: the caller is the primary ------------------------------------------- @Test