This commit is contained in:
@@ -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<Map<String, String>> leads;
|
||||
final Supplier<Map<String, String>> collaboratorTerminals;
|
||||
var leaders = cfg.fleet().leaders();
|
||||
if (!leaders.isEmpty()) {
|
||||
var collaboratorsConfig = cfg.fleet().collaborators();
|
||||
if (!leaders.isEmpty() || !collaboratorsConfig.isEmpty()) {
|
||||
Map<String, String> 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<String, String> 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<String, MemberRole> 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)");
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String> NO_KNOWN_LEAD_OR_COLLABORATOR = target -> false;
|
||||
|
||||
|
||||
@@ -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;
|
||||
*
|
||||
* <p><strong>Resolution order</strong> — connection identity first, token second, nothing third:
|
||||
* <ol>
|
||||
* <li>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.</li>
|
||||
* <li>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.</li>
|
||||
* <li>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 <em>live</em> terminal→slot binding (never a request argument), before
|
||||
* the generic worker fallback.</li>
|
||||
* resolved from the <em>live</em> 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.</li>
|
||||
* <li>A loopback peer PID that maps to an operator-labelled collaborator tab ⇒
|
||||
* {@link Role#COLLABORATOR}, carrying that collaborator's name.</li>
|
||||
* <li>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.</li>
|
||||
@@ -64,6 +74,20 @@ public final class CallerResolver {
|
||||
private final Supplier<Map<String, String>> architectTerminals;
|
||||
private final Function<String, MemberRole> memberSlotRoles;
|
||||
private final Function<String, String> 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<String, MemberRole> 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<Map<String, String>> 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.
|
||||
*
|
||||
* <p>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.
|
||||
* <p>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<Map<String, String>> 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}.
|
||||
*
|
||||
* <p>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<Map<String, String>> leadTerminals,
|
||||
MemberRegistry members,
|
||||
Function<String, MemberRole> spawnedMemberRole,
|
||||
Supplier<Map<String, String>> 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<Map<String, String>> fixed(Map<String, String> leadTerminals) {
|
||||
@@ -160,6 +209,17 @@ public final class CallerResolver {
|
||||
Supplier<Map<String, String>> architectTerminals,
|
||||
Function<String, MemberRole> memberSlotRoles,
|
||||
Function<String, String> memberSlotNames) {
|
||||
this(identity, tokenMode, token, leadTerminals, architectTerminals, memberSlotRoles,
|
||||
memberSlotNames, null, null);
|
||||
}
|
||||
|
||||
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
Supplier<Map<String, String>> architectTerminals,
|
||||
Function<String, MemberRole> memberSlotRoles,
|
||||
Function<String, String> memberSlotNames,
|
||||
Function<String, MemberRole> spawnedMemberRole,
|
||||
Supplier<Map<String, String>> 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}.
|
||||
*
|
||||
* <p>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<String, String> 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<String> 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
|
||||
}
|
||||
|
||||
|
||||
@@ -96,13 +96,19 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
|
||||
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<String, String> tabToName;
|
||||
private final Map<String, Entry> tabToEntry;
|
||||
private final Set<String> excludedWorkspaceLabels;
|
||||
private final long ttlNanos;
|
||||
private final LongSupplier clock;
|
||||
|
||||
private Map<String, String> cached = Map.of();
|
||||
private Map<String, Entry> cached = Map.of();
|
||||
private long scannedAtNanos;
|
||||
private boolean everScanned;
|
||||
|
||||
@@ -126,26 +132,53 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
*/
|
||||
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
|
||||
Set<String> 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.<name>.tab}), matched the same way as
|
||||
* {@code tabToName}
|
||||
*/
|
||||
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
|
||||
Map<String, String> collaboratorTabToName,
|
||||
Set<String> 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<String, String> normalize(Map<String, String> 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<String, Entry> buildTabIndex(Map<String, String> tabToName,
|
||||
Map<String, String> collaboratorTabToName) {
|
||||
Map<String, Entry> out = new LinkedHashMap<>();
|
||||
putNormalized(out, collaboratorTabToName, Kind.COLLABORATOR);
|
||||
putNormalized(out, tabToName, Kind.LEAD);
|
||||
return Collections.unmodifiableMap(out);
|
||||
}
|
||||
|
||||
private static void putNormalized(Map<String, Entry> out, Map<String, String> tabToName, Kind kind) {
|
||||
if (tabToName == null) {
|
||||
return;
|
||||
}
|
||||
Map<String, String> 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<Map<String, String>> {
|
||||
*/
|
||||
@Override
|
||||
public synchronized Map<String, String> 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<String, String> collaborators() {
|
||||
return byKind(refresh(), Kind.COLLABORATOR);
|
||||
}
|
||||
|
||||
private static Map<String, String> byKind(Map<String, Entry> entries, Kind kind) {
|
||||
Map<String, String> 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<String, Entry> refresh() {
|
||||
long now = clock.getAsLong();
|
||||
if (everScanned && now - scannedAtNanos < ttlNanos) {
|
||||
return cached;
|
||||
@@ -165,21 +221,21 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
scannedAtNanos = now;
|
||||
everScanned = true;
|
||||
try {
|
||||
Map<String, String> fresh = scan();
|
||||
Map<String, Entry> 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<String, String> scan() {
|
||||
Map<String, String> nameByTab = new LinkedHashMap<>();
|
||||
private Map<String, Entry> scan() {
|
||||
Map<String, Entry> 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<Map<String, String>> {
|
||||
}
|
||||
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<String> 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<String, String>> {
|
||||
}
|
||||
}
|
||||
|
||||
Map<String, String> byTerminal = new LinkedHashMap<>();
|
||||
Map<String, Entry> byTerminal = new LinkedHashMap<>();
|
||||
Set<String> 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<Map<String, String>> {
|
||||
// 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<Map<String, String>> {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*
|
||||
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToName} — no prefix
|
||||
* <p>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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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<String> 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
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>{@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<FleetConfig.Fleet> 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"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -163,6 +163,13 @@ class LeadTabScannerTest {
|
||||
return new LeadTabScanner(herdr, tabToName, Set.of("fleetd-workers"), TTL, clock::get);
|
||||
}
|
||||
|
||||
private LeadTabScanner scannerWithCollaborators(TopologyHerdr herdr, Map<String, String> tabToName,
|
||||
Map<String, String> collaboratorTabToName,
|
||||
AtomicLong clock) {
|
||||
return new LeadTabScanner(herdr, tabToName, collaboratorTabToName,
|
||||
Set.of("fleetd-workers"), TTL, clock::get);
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyConfiguredTabBecomesALeadNamedByItsEntry() {
|
||||
Map<String, String> 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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<String> toLead = send(port, "POST", "/sessions/term_lead_known/message",
|
||||
"{\"content\":\"hi\",\"wait\":false}", null);
|
||||
assertEquals(202, toLead.statusCode(), toLead.body());
|
||||
|
||||
HttpResponse<String> 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<String, String> leadTerminals,
|
||||
Map<String, String> 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
|
||||
|
||||
Reference in New Issue
Block a user