Compare commits

..

14 Commits

Author SHA1 Message Date
Dai Ha 2c467c2553 fleetd #693, #676: pin lead-tab case-insensitivity, drop stale validator count
CI / shell-tests (pull_request) Failing after 20s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m40s
#693: add a FleetConfigTest case for two fleet.leaders tabs differing only
in case — mutating equalsIgnoreCase to equals left this uncaught before.

#676: rename theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames
to theSweepRunsEveryValidateMethodOnAnUnrelatedClass, since there are
eight validators now, not six, and the count had drifted into the name
and its class-javadoc {@link}.
2026-10-03 22:51:23 +02:00
Dai Ha bfee23acc3 Merge PR #691: fleetd #677 — refuse two leads sharing one exact tab (supersedes #686)
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m37s
2026-10-03 22:45:23 +02:00
Dai Ha 7dec74f1b4 Merge PR #690: fleetd #638 — stop the verdict userinfo mask from crossing / or whitespace (supersedes #685)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m58s
2026-10-03 22:39:42 +02:00
Dai Ha 482598e2a6 fleetd #677: refuse two leads that share one exact tab
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Failing after 2m17s
validateLeadTabPrefixes() compared each lead's tab against the member
tabLabel template, but never against another lead's tab. Two leads
configured with the same exact tab loaded cleanly, even though exact
tab is the only thing lead identity is matched on, so only one of them
could ever be found.

Add an independent pass, case-insensitive, that refuses when two
fleet.leaders entries share one exact tab, with its own exception so
the message stays accurate for this relation.

Decided not to also refuse a lead's tab starting with a sibling's
tabPrefix: tabPrefix plays no role in identity resolution (only the
exact tab does), and the default tabPrefix is "lead:", the same
string used as the conventional lead tab prefix throughout this
codebase's own fixtures (e.g. "lead: opus" / "lead: sol"). Refusing
that case would reject the standard multi-lead setup with no matching
identity hazard.
2026-10-03 22:36:44 +02:00
Dai Ha 5f5d16fbd4 fleetd #638: stop the verdict userinfo mask from crossing / or whitespace
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m57s
mask_verdict_userinfo's character class [^@]* crossed a '/' or a space, so a
verdict line with a URI that has no userinfo plus a later @ (e.g. an email
address in diagnostic prose) had everything between them destroyed. Restrict
the class to [^@/[:space:]]* so the match stops at the end of the URI.

Adds the uncovered-direction test: a URI with no userinfo plus a later @ in
the same line must pass through byte for byte. Also drops the design-rationale
sentence from the helper's comment (now in the PR description).
2026-10-03 22:34:02 +02:00
Dai Ha 7df7985a16 Merge remote-tracking branch 'refs/remotes/pr/686' into worker/677-fix-lead-collision-f69073-12 2026-10-03 22:29:29 +02:00
Dai Ha 724b35b46e Merge remote-tracking branch 'refs/remotes/pr/685' into worker/638-fix-overmask-dbb1bf-11 2026-10-03 22:29:14 +02:00
Dai Ha 2eb2d6112e Merge PR #688: fleetd #675 — pin three unpinned FleetdAssembly constructor args
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m35s
2026-10-03 22:27:32 +02:00
Dai Ha 7c458e8bf2 Merge PR #687: fleetd #669 Unit A — split SEND and READ into their real call shapes (closes #678) 2026-10-03 22:27:32 +02:00
Dai Ha 133f03e428 Merge PR #684: fleetd #683 — decouple the completion-fallback test's two 5s budgets 2026-10-03 22:27:27 +02:00
Dai Ha 804279175d fleetd #675: pin assembly loop timing defaults
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m42s
2026-10-03 22:12:46 +02:00
Dai Ha 9dea289975 fleetd #669 Unit A: split SEND and READ into their real call shapes
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Failing after 2m7s
SEND covered three different call shapes under one action (local
sessionId delivery, the coordId cross-host broker route, and the
turnId answer-a-blocked-worker form). READ covered both roster/
profile/identity observation and ticket-polling/session-status.
Split each into its own Authz.Action — SEND/COORD_SEND/ANSWER and
READ/TASK_READ — with every new action granted to exactly who held
the combined action before, on both the MCP and REST entry paths.

Also fixes fleetd #678's Authz.java comment: READ no longer claims
"the roster carries no secrets" for ticket replies and pending
questions, because those now live under TASK_READ.
2026-10-03 22:12:13 +02:00
Dai Ha cb4a6869b9 fleetd #677: guard exact lead tab collisions
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Failing after 1m56s
2026-10-03 22:10:29 +02:00
Dai Ha 28b45d97e5 fleetd #638: mask userinfo in the daemon verdict line before it prints
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Failing after 2m4s
Apply the userinfo-only rewrite (sed -E 's#://[^@]*@#://<redacted>@#g') to every
path that prints config-edit.sh's $VERDICT_LINE: restore_and_confirm's two
prints, report_outcome's shared local copy (covering its clean/needs-restart/
refused branches), and check_mode's "last verdict in log" line via
last_verdict_line. Deliberately not routed through redact() — that function's
key:value masking does not match this line's prose, and the rest of the line
(e.g. the pattern quoted in a parse-failure refusal) is the detail an operator
needs to fix the refusal.

No current refusal message echoes a URI, token or password, so this is a guard
against a future validator doing so, not a fix for an observed leak.
2026-10-03 22:06:54 +02:00
12 changed files with 719 additions and 204 deletions
@@ -20,16 +20,22 @@ public final class Authz {
SPAWN,
/** Tear a worker peer down. */
STOP,
/** Deliver a turn to a session (or answer a worker's question). */
/** Deliver a turn to a local session, addressed by {@code sessionId}. */
SEND,
/** Resolve a worker's blocked question and resume its turn, addressed by {@code turnId}. */
ANSWER,
/** Address a peer lead on another daemon over the coordination broker, by {@code coordId}. */
COORD_SEND,
/** A worker's terminal reply for its own turn. */
REPLY,
/** A worker's mid-turn question to the primary. */
ASK,
/** Collect held replies from a session's inbox. */
DRAIN,
/** Read-only observation: status, roster, profiles, task polling. */
/** Read-only roster, profile, and identity observation: no ticket, task, or turn state. */
READ,
/** Poll a ticket, or read a session's status. */
TASK_READ,
/**
* Read (never ack) this daemon's own held lead-to-lead coordination mail (fleetd #421).
*
@@ -71,11 +77,19 @@ public final class Authz {
// escalating into the orchestrator role.
case SPAWN, STOP, DRAIN, HANDOVER -> caller.isPrimary();
// Delivering a turn is open to the primary and the architect: an architect delegates
// to workers (that is the role's point) but still has no lifecycle rights. A worker is
// excluded — sending would be it escalating.
// Delivering a turn to a local session is open to the primary and the architect: an
// architect delegates to workers (that is the role's point) but still has no lifecycle
// rights. A worker is excluded — sending would be it escalating.
case SEND -> caller.isPrimary() || caller.isArchitect();
// Same grant as SEND. Resolving a worker's blocked question is part of delegating to
// it, not a separate capability.
case ANSWER -> caller.isPrimary() || caller.isArchitect();
// Same grant as SEND. This leaves the daemon over the coordination broker rather than
// addressing a local session, but the caller who may do one may do the other.
case COORD_SEND -> caller.isPrimary() || caller.isArchitect();
// The load-bearing rule: a caller acts only as the pane it occupies. CB-532 widened who
// that can be — a lead answering another lead is replying for its OWN terminal, which
// this already permits — while the rule itself is unchanged, and is what stops anyone
@@ -84,10 +98,17 @@ public final class Authz {
// unnamed primary (token/loopback, no pane) owns nothing and is still excluded.
case REPLY, ASK -> caller.ownsSession(targetSession);
// Observation is open to every authenticated role: a worker legitimately polls its own
// status, and the roster carries no secrets.
// READ is roster, profile, and identity observation — fleet_list, fleet_profiles, and
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
// other session's turn state. Those live under TASK_READ. METRICS is the separate
// Prometheus scrape. Both stay open to every authenticated role.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// Ticket polling and session status, open to every authenticated role the same as READ.
// Unlike READ, a holder may poll a ticket it did not create, or read another session's
// pending question and the turnId that answers it.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
// is coordination between leads, not observation of the roster.
@@ -1129,11 +1129,8 @@ public record FleetConfig(
* {@code tab} can never be discovered, launched or not
* @param instances how many of this lead should be live (default 1). The daemon
* launches only the shortfall, so a restart adopts rather than doubles
* @param tabPrefix no longer used to find a lead's tab — {@code tab} is matched
* exactly. Its only remaining job is the startup collision guard
* ({@link #validateLeadTabPrefixes()}), which still uses it to refuse
* a worker {@code tabLabel} template that could be misread as a lead.
* Default {@code "lead:"}
* @param tabPrefix lead-tab naming convention checked against member labels. Lead
* identity uses {@code tab}. Default {@code "lead:"}
* @param scanIntervalSeconds how long a tab scan is cached before herdr is asked again; also the
* worst case before a newly-labelled tab is recognised. Default 10
* @param kind which agent runs there ({@code claude}, {@code opencode}, …)
@@ -1236,10 +1233,7 @@ public record FleetConfig(
String tabLabel) {
/**
* Role first, so the tab bar reads as the fleet and so the label shares a namespace with a
* lead's {@code tabPrefix}. Because {@code {role}} comes from a closed enum, a generated
* member label can never begin with {@code "lead:"} — the clash that
* {@link #validateLeadTabPrefixes()} used to have to check for is unrepresentable here.
* Role first, so the tab bar identifies the member's fleet role.
*/
public static final String DEFAULT_TAB_LABEL = "{role}: {profile} #{n}";
@@ -2685,28 +2679,14 @@ public record FleetConfig(
}
/**
* Reject a lead-scan convention that a worker tab would also satisfy (CB-531).
* Reject a member tab-label template that could render as a configured lead tab or match a
* lead-tab naming convention, and reject two {@code fleet.leaders} entries that share one exact
* tab.
*
* <p>The scan reads a tab label and concludes "a lead lives here". fleetd also <em>writes</em>
* tab labels — every member gets one rendered into its tab. Choose a lead {@code tabPrefix} that
* a member template matches and the daemon starts labelling its own members as leads, promoting
* the entire fleet to {@link dev.ltms.fleet.auth.Role#PRIMARY} with no message and no diff.
* {@link #validatePanePlacementAgainstLeadTabs()} is the check that stops a pane-placed member
* from landing inside a lead's tab in the first place; this check is a second, independent
* guard that catches the hazard even when every profile places members correctly, by refusing
* a label that a scan would still misread as a lead.
*
* <p>CB-557 shrank this check rather than removing it. The default template is
* {@code "{role}: {profile} #{n}"} and {@code {role}} comes from a closed enum, so a
* <em>generated</em> label can no longer collide by construction. What remains checkable is what
* an operator still writes by hand: the {@code fleet.tabLabel} template and any per-profile
* {@code tabLabel} override.
*
* <p>Fatal rather than a warning, unlike {@link #warnUnknownTopLevelKeys}: an unknown key means
* a feature does nothing, while this means a feature does the opposite of what it says.
*
* @throws IllegalStateException when the fleet template or any profile's {@code tabLabel}
* override starts with a configured lead prefix
* @throws IllegalStateException when the fleet template or a profile {@code tabLabel} override
* can render as a configured lead tab or match a lead-tab prefix,
* or when two {@code fleet.leaders} entries carry the same exact
* {@code tab} (case-insensitively)
*/
public void validateLeadTabPrefixes() {
if (fleet == null || fleet.leaders().isEmpty()) {
@@ -2717,28 +2697,90 @@ public record FleetConfig(
if (leader == null) {
return;
}
String tab = leader.tab();
String prefix = leader.tabPrefix();
// The fleet-wide template is checked once per prefix: it labels every member that has no
// override, so one bad template promotes the entire fleet, not one profile.
if (startsWithIgnoreCase(fleet.tabLabel(), prefix)) {
if (templateCanRenderAs(fleet.tabLabel(), tab)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" can render as the tab of "
+ "lead '" + leadName + "' (\"" + tab + "\")");
} else if (startsWithIgnoreCase(fleet.tabLabel(), prefix)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" starts with the tabPrefix of "
+ "lead '" + leadName + "' (\"" + prefix + "\")");
}
profiles().entrySet().stream()
.filter(e -> startsWithIgnoreCase(e.getValue().tabLabel(), prefix))
.map(Map.Entry::getKey)
.sorted()
.forEach(p -> bad.add("profile '" + p + "' overrides tabLabel with \""
+ profiles().get(p).tabLabel() + "\", which starts with the tabPrefix of "
+ "lead '" + leadName + "' (\"" + prefix + "\")"));
.forEach(p -> {
String label = profiles().get(p).tabLabel();
if (templateCanRenderAs(label, tab)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which can render as the tab of lead '" + leadName
+ "' (\"" + tab + "\")");
} else if (startsWithIgnoreCase(label, prefix)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which starts with the tabPrefix of lead '" + leadName
+ "' (\"" + prefix + "\")");
}
});
});
if (bad.isEmpty()) {
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". Every member labelled that way would be read back as a lead and granted "
+ "spawn/stop/send on the whole fleet. Change one of the two so member tabs "
+ "and lead tabs cannot be confused.");
}
List<String> collisions = new ArrayList<>();
List<String> leadNames = fleet.leaders().keySet().stream().sorted().toList();
for (int i = 0; i < leadNames.size(); i++) {
String nameA = leadNames.get(i);
Leader a = fleet.leaders().get(nameA);
if (a == null || a.tab() == null || a.tab().isBlank()) {
continue;
}
for (int j = i + 1; j < leadNames.size(); j++) {
String nameB = leadNames.get(j);
Leader b = fleet.leaders().get(nameB);
if (b == null || b.tab() == null || b.tab().isBlank()) {
continue;
}
if (a.tab().equalsIgnoreCase(b.tab())) {
collisions.add("lead '" + nameA + "' and lead '" + nameB + "' both use tab \""
+ a.tab() + "\"");
}
}
}
if (collisions.isEmpty()) {
return;
}
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". Every member labelled that way would be read back as a lead and granted "
+ "spawn/stop/send on the whole fleet. Change one of the two so member tabs and "
+ "lead tabs cannot be confused.");
throw new IllegalStateException("refusing to start: " + String.join("; ", collisions)
+ ". Tab identity is matched exactly, so only one of two leads sharing a tab can "
+ "ever be found — the other is silently unreachable. Give each lead its own "
+ "exact tab.");
}
private static boolean templateCanRenderAs(String template, String tab) {
if (template == null || template.isBlank() || tab == null || tab.isBlank()) {
return false;
}
var placeholders = Pattern.compile("\\{(?:role|profile|model|n)}").matcher(template);
StringBuilder expression = new StringBuilder("^");
int literalStart = 0;
while (placeholders.find()) {
expression.append(Pattern.quote(template.substring(literalStart, placeholders.start())));
expression.append(".*");
literalStart = placeholders.end();
}
expression.append(Pattern.quote(template.substring(literalStart))).append("$");
return Pattern.compile(expression.toString(), Pattern.CASE_INSENSITIVE).matcher(tab).matches();
}
/** Case-insensitive prefix test that tolerates a null or blank label. */
private static boolean startsWithIgnoreCase(String label, String prefix) {
if (label == null || prefix == null || prefix.isBlank()) {
return false;
}
String stripped = label.strip();
return stripped.regionMatches(true, 0, prefix, 0, prefix.length());
}
/**
@@ -2800,15 +2842,6 @@ public record FleetConfig(
}
}
/** Case-insensitive prefix test that tolerates a null/blank label. */
private static boolean startsWithIgnoreCase(String label, String prefix) {
if (label == null || prefix == null || prefix.isBlank()) {
return false;
}
String stripped = label.strip();
return stripped.regionMatches(true, 0, prefix, 0, prefix.length());
}
/**
* Reject a subscription profile whose {@code env:} block tries to reseat the Anthropic binding
* (CB-542).
@@ -690,7 +690,7 @@ public final class FleetMcp {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target)) {
if (action != Authz.Action.READ) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
return null;
@@ -1035,31 +1035,21 @@ public final class FleetMcp {
}
/**
* Which authorization action a {@code fleet_poll} call needs, decided by its arguments
* (fleetd #272, widened by fleetd #421).
* Which authorization action a {@code fleet_poll} call needs, decided by its arguments.
*
* <p>{@code fleet_poll} is now <strong>three operations behind one tool name</strong>. With
* <p>{@code fleet_poll} is <strong>three operations behind one tool name</strong>. With
* {@code ticket} it observes an async delegation and changes nothing, which is a {@link
* Authz.Action#READ}. With {@code target} it calls {@link MessageService#drainReplies} on that
* session -- the replies are removed from the inbox and a second call returns nothing -- so it
* is a {@link Authz.Action#DRAIN}, the same gate {@code fleet_ack} already uses for removing a
* single message, and the same one the REST path uses at {@code FleetApp.drainReplies}. With
* {@code coordId} it reads (never acks) this daemon's own held lead-to-lead mail, which is a
* {@link Authz.Action#COORD_READ} -- <strong>not</strong> {@code READ}, even though nothing is
* consumed: {@code READ}'s grant is open to every authenticated role on the premise that the
* roster carries no secrets, and a lead-to-lead body is not the roster. Mapping a non-destructive
* peer-mail read to {@code READ} would let any worker read every peer lead's mail in full.
*
* <p>Before this method existed (fleetd #272) the handler passed a constant {@code READ} for
* both of the original branches. {@code READ} is open to every authenticated role, so any
* worker could read a peer's id out of {@code fleet_list} and destroy the replies that peer had
* queued for the primary. The gate failed open, and it did so because the required action is a
* function of the arguments while the handler chose it before looking at them.
* Authz.Action#TASK_READ}. With {@code target} it calls {@link MessageService#drainReplies} on
* that session -- the replies are removed from the inbox and a second call returns nothing --
* so it is a {@link Authz.Action#DRAIN}, the same gate {@code fleet_ack} already uses for
* removing a single message, and the same one the REST path uses at
* {@code FleetApp.drainReplies}. With {@code coordId} it reads (never acks) this daemon's own
* held lead-to-lead mail, which is a {@link Authz.Action#COORD_READ} -- <strong>not</strong>
* {@code TASK_READ} or {@code READ}: a lead-to-lead body is a different inbox from either, and
* folding it into either would let any worker or architect read every peer lead's mail in full.
*
* <p>The choice lives in this method, and not inline in the handler, so that a test can assert
* the mapping the handler actually uses. {@code FleetMcpAuthzTest} already checked every
* {@link Authz.Action} against every {@link Role} and passed throughout -- it tested the policy
* table, which was correct, while the defect was in which action the caller handed it.
* the mapping the handler actually uses.
*
* <p>Checked first, and exclusively of {@code target}: a call naming {@code coordId} is reading
* a different inbox entirely (this daemon's own lead channel, never a worker's), so it takes
@@ -1072,7 +1062,30 @@ public final class FleetMcp {
if (!isBlank(coordId)) {
return Authz.Action.COORD_READ;
}
return isBlank(target) ? Authz.Action.READ : Authz.Action.DRAIN;
return isBlank(target) ? Authz.Action.TASK_READ : Authz.Action.DRAIN;
}
/**
* Which authorization action a {@code fleet_send} call needs, decided by its arguments.
*
* <p>{@code fleet_send} is three call shapes behind one tool name, mirroring {@link
* #pollAction}. With {@code coordId} it addresses a peer lead on another daemon over the
* coordination broker, which is {@link Authz.Action#COORD_SEND}. With {@code turnId} it
* resolves a worker's blocked {@code fleet_ask} and resumes that turn, which is {@link
* Authz.Action#ANSWER}. Otherwise it delivers to a local session by {@code sessionId}, which is
* the plain {@link Authz.Action#SEND}.
*
* <p>Checked in the same order the handler branches: {@code coordId} first and exclusively of
* {@code turnId}, matching {@link #sendToLead}'s own mutual-exclusion check.
*
* @param coordId the {@code coordId} argument of the call, or {@code null}/blank when absent
* @param turnId the {@code turnId} argument of the call, or {@code null}/blank when absent
*/
static Authz.Action sendAction(String coordId, String turnId) {
if (!isBlank(coordId)) {
return Authz.Action.COORD_SEND;
}
return isBlank(turnId) ? Authz.Action.SEND : Authz.Action.ANSWER;
}
/**
@@ -1097,10 +1110,11 @@ public final class FleetMcp {
*/
private static Authz.Action authzAction(FleetTool tool, Map<String, Object> arguments) {
return switch (tool) {
case SEND -> Authz.Action.SEND;
case SEND -> sendAction(str(arguments, "coordId"), str(arguments, "turnId"));
case REPLY -> Authz.Action.REPLY;
case ASK -> Authz.Action.ASK;
case STATUS, LIST, PROFILES, WHOAMI -> Authz.Action.READ;
case STATUS -> Authz.Action.TASK_READ;
case LIST, PROFILES, WHOAMI -> Authz.Action.READ;
case POLL -> pollAction(str(arguments, "target"), str(arguments, "coordId"));
case ACK -> Authz.Action.DRAIN;
case SPAWN -> Authz.Action.SPAWN;
@@ -49,18 +49,36 @@ import java.util.stream.Collectors;
*/
public final class FleetApp {
/** The authorization action the matching route handler hands to {@link #allow}. */
/**
* The authorization action the matching route handler hands to {@link #allow}, for a route
* whose action does not depend on the request body.
*/
static Authz.Action routeAction(String route) {
return routeAction(route, null);
}
/**
* As above, plus the one route whose action depends on the body: {@code POST
* /sessions/{id}/message} carries a {@code turnId} (the answer-a-blocked-worker shape) or not
* (a plain delivery), mirroring {@code FleetMcp#sendAction}'s split of the same two call
* shapes over MCP. {@code turnId} is ignored by every other route.
*
* @param turnId the request body's {@code turnId}, or {@code null}/blank when absent or not
* applicable to this route
*/
static Authz.Action routeAction(String route, String turnId) {
return switch (route) {
case "GET /metrics" -> Authz.Action.METRICS;
case "POST /members" -> Authz.Action.SPAWN;
case "DELETE /members/{paneId}" -> Authz.Action.STOP;
case "POST /sessions/{id}/message" -> Authz.Action.SEND;
case "POST /sessions/{id}/message" -> turnId == null || turnId.isBlank()
? Authz.Action.SEND : Authz.Action.ANSWER;
case "POST /sessions/{id}/reply" -> Authz.Action.REPLY;
case "GET /sessions/{id}/replies" -> Authz.Action.DRAIN;
case "POST /sessions/{id}/ask" -> Authz.Action.ASK;
case "GET /sessions", "GET /agents", "GET /members", "GET /profiles",
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}" -> Authz.Action.READ;
"GET /member-credentials" -> Authz.Action.READ;
case "GET /sessions/{id}/status", "GET /tasks/{ticket}" -> Authz.Action.TASK_READ;
default -> throw new IllegalArgumentException("route has no authorization gate: " + route);
};
}
@@ -250,7 +268,8 @@ public final class FleetApp {
}
Principal caller = ctx.attribute(CALLER);
if (Authz.permits(caller, action, target)) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
return true;
@@ -603,26 +622,33 @@ public final class FleetApp {
* status-gated injector and block until the worker returns a structured {@code fleet_reply}.
* Times out with a typed 202 (working / queued / busy) rather than an error — the message may
* still land.
*
* <p>Two call shapes share this route, exactly as {@code fleet_send} does over MCP (see
* {@code FleetMcp#sendAction}): a plain delivery to {@code id}, and -- when the body carries
* {@code turnId} -- resolving a worker's blocked question. The body is parsed before the
* authorization check so the right one of {@link Authz.Action#SEND}/{@link Authz.Action#ANSWER}
* reaches the gate; a body that fails to parse is treated as the plain shape for that check
* alone, and is rejected afterward exactly as before.
*/
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
JsonNode body;
try {
body = mapper.readTree(ctx.body());
} catch (Exception e) {
body = null;
}
String turnId = body == null ? null : body.path("turnId").asText(null);
if (!allow(ctx, routeAction("POST /sessions/{id}/message", turnId), id)) {
return;
}
String content;
String turnId;
long timeout;
boolean wait;
try {
JsonNode body = mapper.readTree(ctx.body());
content = body.path("content").asText("");
turnId = body.path("turnId").asText(null);
timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
wait = body.path("wait").asBoolean(true); // default: block for the reply (CB-104)
} catch (Exception e) {
if (body == null) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
String content = body.path("content").asText("");
long timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
boolean wait = body.path("wait").asBoolean(true); // default: block for the reply (CB-104)
if (content.isBlank()) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
return;
@@ -0,0 +1,190 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
/**
* Asserts that the assembled loops use the production reminder, coordination, and delivery timing
* defaults when no {@code primary:} block configures the reply-push values.
*/
class FleetdAssemblyTimingDefaultsTest {
private static final class FakeLeadChannel implements LeadChannelHandle {
@Override
public void publish(String toCoordId, LeadMessage message) {
}
@Override
public List<LeadMessage> peek() {
return List.of();
}
@Override
public void ack(String msgId) {
}
@Override
public String selfCoordId() {
return "test-lead";
}
@Override
public boolean heldDurable() {
return true;
}
@Override
public MailboxState inspect(String coordId) {
return MailboxState.unknown(coordId);
}
@Override
public void close() {
}
}
private static final class TestResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> new FakeLeadChannel();
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private TestResourcePorts ports;
@AfterEach
void tearDown() {
if (ports != null && ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
coordinator:
uri: "amqp://fake-lead-broker/vh"
selfId: "test-lead"
""");
return FleetConfig.load(file);
}
private FleetdRuntime assemble(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ports = new TestResourcePorts();
return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
}
private static long longField(Object target, String name) throws Exception {
Field field = target.getClass().getDeclaredField(name);
field.setAccessible(true);
return field.getLong(target);
}
@Test
void productionBootPathUsesTheExpectedLoopTimingDefaults(@TempDir Path dir) throws Exception {
FleetdRuntime runtime = assemble(dir);
ReplyPushLoop pushLoop = runtime.pushLoop();
assertEquals(5, longField(pushLoop, "maxReminders"),
"without primary:, ReplyPushLoop must stop after five reminder attempts");
assertEquals(15_000L, longField(pushLoop, "backoffMs"),
"without primary:, ReplyPushLoop must wait fifteen seconds before the next reminder");
LeadCoordLoop leadCoordLoop = runtime.leadCoordLoop();
assertNotNull(leadCoordLoop, "control: coordinator: must build LeadCoordLoop");
assertEquals(3_000L, longField(leadCoordLoop, "intervalMs"),
"LeadCoordLoop must poll for peer-lead mail every three seconds");
StatusPoller poller = runtime.poller();
assertEquals(Injector.POLL_INTERVAL_MILLIS, longField(poller, "intervalMillis"),
"StatusPoller must use Injector's delivery poll interval");
}
}
@@ -38,6 +38,37 @@ class AuthzTest {
}
}
/**
* {@code fleet_send} is three call shapes behind one action name until {@code
* FleetMcp#sendAction} picks one: a plain local {@link Authz.Action#SEND}, the {@code coordId}
* route ({@link Authz.Action#COORD_SEND}), and the {@code turnId} answer form ({@link
* Authz.Action#ANSWER}). All three carry the same grant as the undivided action did — a worker
* is excluded from every one, exactly as it was excluded from the one combined action before.
*/
@Test
void theThreeSendShapesCarryTheSameGrantAsTheOldUndividedAction() {
for (Authz.Action a : new Authz.Action[]{SEND, COORD_SEND, ANSWER}) {
assertTrue(Authz.permits(PRIMARY, a, "term_a"), "the primary may " + a);
assertTrue(Authz.permits(ARCH_DESIGN, a, "term_a"), "an architect may " + a);
assertFalse(Authz.permits(WORKER_A, a, "term_a"),
"a worker performing " + a + " would be escalating into the orchestrator role");
assertFalse(Authz.permits(ANON, a, "term_a"));
}
}
/**
* {@code fleet_poll{ticket}} and {@code fleet_status} are {@link Authz.Action#TASK_READ}, split
* out of the roster-only {@link Authz.Action#READ} (fleetd #678). The grant is unchanged from
* what the undivided {@code READ} action gave every one of these callers.
*/
@Test
void taskReadCarriesTheSameGrantReadDidBeforeTheSplit() {
assertTrue(Authz.permits(PRIMARY, TASK_READ, null));
assertTrue(Authz.permits(WORKER_A, TASK_READ, null));
assertTrue(Authz.permits(ARCH_DESIGN, TASK_READ, null));
assertFalse(Authz.permits(ANON, TASK_READ, null));
}
@Test
void aWorkerMayReplyAndAskOnlyAsItself() {
assertTrue(Authz.permits(WORKER_A, REPLY, "term_a"));
@@ -676,12 +676,8 @@ class FleetConfigTest {
assertEquals(5, hb.quietNudgeCap());
}
/**
* The hazard the guard exists for: fleetd writes worker tab labels and reads lead tab labels.
* Overlap the two and every worker it spawns is read back as a lead.
*/
@Test
void aLeadPrefixThatAProfileTabLabelOverrideAlsoMatchesRefusesToStart(@TempDir Path dir)
void aProfileTabLabelOverrideMatchingALeadTabRefusesToStart(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("collide.yaml");
Files.writeString(f, """
@@ -689,32 +685,33 @@ class FleetConfigTest {
port: 8080
profiles:
gx10:
tabLabel: "lead: {profile} #{n}"
tabLabel: "alpha"
fleet:
leaders:
opus:
tab: "lead: opus"
tabPrefix: "lead:"
tab: "alpha"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
assertTrue(e.getMessage().contains("alpha"), "the message must name the offending label");
}
/** A bad fleet-wide template promotes every member, not one profile — so it is checked too. */
@Test
void aFleetTabLabelThatMatchesALeadPrefixRefusesToStart(@TempDir Path dir) throws Exception {
void aFleetTabLabelTemplateThatCanRenderAsALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collide-template.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
pha: {}
fleet:
tabLabel: "lead: {role} {profile}"
tabLabel: "al{profile}"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""");
FleetConfig cfg = FleetConfig.load(f);
@@ -723,12 +720,28 @@ class FleetConfigTest {
assertTrue(e.getMessage().contains("fleet.tabLabel"));
}
/**
* The point of making role the label's first field: {@code {role}} comes from a closed enum, so
* a generated label cannot begin with {@code "lead:"} however the fleet is configured.
*/
@Test
void theDefaultTabLabelCannotCollideWithTheDefaultLeadPrefix(@TempDir Path dir) throws Exception {
void anExactFleetTabLabelCollisionRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("exact-tab-collision.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
tabLabel: "alpha"
leaders:
alpha:
tab: "alpha"
""");
IllegalStateException e = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(f).validateAll());
assertTrue(e.getMessage().contains("fleet.tabLabel"),
"the message must name the offending label");
assertTrue(e.getMessage().contains("alpha"), "the message must name the colliding lead tab");
}
@Test
void aFleetTabLabelTemplateThatCannotRenderAsALeadTabIsAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("ok.yaml");
Files.writeString(f, """
bind:
@@ -737,17 +750,13 @@ class FleetConfigTest {
gx10:
baseUrl: http://gx00.gw:8000
fleet:
tabLabel: "worker-{profile}"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""");
assertDoesNotThrow(() -> FleetConfig.load(f).validateLeadTabPrefixes());
for (MemberRole role : MemberRole.values()) {
assertFalse(FleetConfig.Fleet.DEFAULT_TAB_LABEL
.replace("{role}", role.wireName()).startsWith("lead:"),
"no role renders a label that reads as a lead");
}
assertDoesNotThrow(() -> FleetConfig.load(f).validateAll());
}
@Test
@@ -765,6 +774,78 @@ class FleetConfigTest {
"a label that collides with a convention nobody reads is not a problem");
}
/**
* fleetd #677: identity is matched on a lead's exact {@code tab} alone, so two leads sharing
* one tab means only one of them is ever found — the guard must catch this independently of
* the member-template checks above.
*/
@Test
void twoLeadsSharingTheSameExactTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "shared tab"
sonnet:
tab: "shared tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/**
* fleetd #693: the guard matches tabs case-insensitively, because
* {@code LeadTabScanner} keys its tab map on a lowercased label — two tabs differing only in
* case collide there too, and the guard must catch that independently of the exact-match case
* above.
*/
@Test
void twoLeadsSharingTheSameTabInDifferentCaseRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab-case.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "Shared Tab"
sonnet:
tab: "shared tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/** Control for {@link #twoLeadsSharingTheSameExactTabRefusesToStart}: distinct tabs load cleanly. */
@Test
void twoLeadsWithDistinctExactTabsAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("distinct-tabs.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "opus tab"
sonnet:
tab: "sonnet tab"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
// ── validatePanePlacementAgainstLeadTabs ────────────────────────────────────────────────────
/**
@@ -17,63 +17,21 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The gap this class exists to close: mutation testing on the fleetd ticket "central allow-list
* of usable models" found that although {@link FleetConfig#validateModels()}'s own logic was well
* pinned, nothing proved either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()})
* still invoked it — deleting the call site left the full suite green (1478/0/0/0). A follow-up
* measurement (same technique — remove one call site, run the suite, not read the code) found the
* SAME gap for all five of {@link FleetConfig}'s other validators at startup, and for four of the
* six inside {@link ConfigRef#reload()}. This is a class of gap, not one line's mistake: every one
* of those thirteen tests called the validator itself directly, never the real caller that was
* supposed to.
* Tests the reflective validator sweep and {@link FleetConfig#validateAll()} reachability.
*
* <p>The fix replaces the six individual {@code cfg.validateXxx()} calls at each of the two real
* call sites with one {@link FleetConfig#validateAll()}, which reaches every validator by
* reflection rather than by a hand-maintained list of names. A hand-maintained list of six names
* would have exactly the defect it replaces: the seventh validator someone adds next month has no
* reason to be added to it, and nothing would say so. This class proves TWO separate claims, and
* keeps them separate on purpose:
*
* <ol>
* <li>{@link #theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames()} and its neighbours
* prove the reflective sweep itself ({@link FleetConfig#invokeAllValidators}) is a general
* mechanism — it runs whatever public, no-arg, void {@code validateXxx()} methods a class
* happens to declare today, including a class with more of them than {@link FleetConfig}
* has right now. This is the proof that a future, real seventh validator on {@link
* FleetConfig} would be swept automatically, without needing to add a real (unwanted)
* seventh validator just to exercise the claim.</li>
* <li>{@link #validateAllReachesEveryOneOfTodaysRealValidators()} proves {@link
* FleetConfig#validateAll()} itself is wired to that same generic mechanism and genuinely
* reaches every one of today's real validators — reusing the exact minimal failing
* configurations {@code FleetConfigTest} already established for each one directly, plus a
* dedicated fixture for {@link FleetConfig#validateLeadRollover()}, which no other test
* drives through {@code validateAll()} — so a single call to {@code validateAll()} is shown
* to reproduce every one of those failures.</li>
* </ol>
*
* <p>Together with the direct-{@code Fleetd.main}-invocation tests in {@code
* FleetdStartupValidationTest} (which prove the real startup call site still calls {@code
* validateAll()}) and the {@code ConfigRefTest} reload tests (which prove the same for {@link
* ConfigRef#reload()}), removing {@code cfg.validateAll();} from either real call site now fails
* a test in this module.
*
* <p><b>What is NOT pinned, measured rather than assumed.</b> Reverting {@link
* FleetConfig#validateAll()} to a hardcoded list of today's method calls leaves the whole
* suite green. Nothing ties {@code validateAll()} to
* the generic sweep — claim 1 proves {@link FleetConfig#invokeAllValidators} is generic, and claim
* 2 proves {@code validateAll()} reaches today's validators, and a hardcoded list satisfies both. So the
* reflective sweep is a convenience, not the guarantee. The guarantee is {@link
* #fleetConfigDeclaresExactlyTheseValidatorsToday()}: it fails the moment any validator is added
* or removed, which forces whoever changes the set to look at this file.
* <p>{@link #theSweepRunsEveryValidateMethodOnAnUnrelatedClass()} and its neighbours
* prove that {@link FleetConfig#invokeAllValidators} runs each public, no-arg, void
* {@code validateXxx()} method on its target. {@link #fleetConfigDeclaresExactlyTheseValidatorsToday()}
* is the canary for the validator set. {@link #validateAllReachesEveryOneOfTodaysRealValidators()}
* is the reachability check for that set.
*/
class FleetConfigValidateAllTest {
// ── Claim 1: the reflective sweep is a general mechanism, not six names in disguise ──────────
// ── Claim 1: the reflective sweep is a general mechanism ─────────────────────────────────────
/**
* A throwaway fixture class, unrelated to {@link FleetConfig} in every way except shape: three
* public, no-arg, void methods named {@code validateXxx}. Proves the sweep works on ANY class
* with this shape, not on something special-cased to {@link FleetConfig}.
* Fixture with public, no-arg, void methods named {@code validateXxx}. It proves the sweep uses
* the target's method shape rather than special handling for {@link FleetConfig}.
*/
static class ThreeValidators {
final List<String> ran = new ArrayList<>();
@@ -92,7 +50,7 @@ class FleetConfigValidateAllTest {
}
@Test
void theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames() {
void theSweepRunsEveryValidateMethodOnAnUnrelatedClass() {
ThreeValidators target = new ThreeValidators();
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateGamma"), target.ran,
@@ -102,12 +60,8 @@ class FleetConfigValidateAllTest {
}
/**
* The core of the "self-maintaining" requirement: the exact same class shape as {@link
* ThreeValidators}, plus one more method — standing in for "a developer adds a validator next
* month". Nothing about the sweep changes to pick it up; the new method is invoked purely
* because it exists and matches the shape. This is what makes adding a seventh real validator
* to {@link FleetConfig} safe without touching {@link FleetConfig#validateAll()} or either
* call site — there is no "wire it in" step left to forget.
* Fixture with an added valid method. It proves the sweep reaches a method because it matches
* the validator shape.
*/
static class FourValidators {
final List<String> ran = new ArrayList<>();
@@ -135,7 +89,7 @@ class FleetConfigValidateAllTest {
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateDelta", "validateGamma"),
sorted(target.ran),
"the fourth method must be reached automatically — proving a class can grow the "
"the added method must be reached automatically — proving a class can grow the "
+ "set of things it validates with no change to the sweep itself");
}
@@ -297,16 +251,16 @@ class FleetConfigValidateAllTest {
port: 8765
""", "auth.mode: token");
// validateLeadTabPrefixes: a fleet-wide tabLabel that starts with a lead's own tabPrefix.
// validateLeadTabPrefixes: a fleet-wide tabLabel that equals a lead tab.
assertValidateAllRefuses(dir, "lead-tab-prefixes.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
tabLabel: "lead: {role} {profile}"
tabLabel: "alpha"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""", "fleet.tabLabel");
// validateSubscriptionProfiles: subscription: true with env: reseating ANTHROPIC_BASE_URL.
@@ -156,6 +156,45 @@ class FleetMcpAuthzTest {
}
}
/**
* fleetd #669 Unit A: {@code SEND} is split into three actions ({@link Authz.Action#SEND},
* {@link Authz.Action#COORD_SEND}, {@link Authz.Action#ANSWER}), each carrying the same grant
* the one undivided action gave. An architect holds all three, exactly as it held the one.
*/
@Test
void anArchitectMayUseAllThreeSendShapesOverMcp() {
FleetMcp m = mcp(true);
for (Authz.Action a : new Authz.Action[]{Authz.Action.SEND, Authz.Action.COORD_SEND,
Authz.Action.ANSWER}) {
assertNull(m.denyFor(ARCH_DESIGN, a, "term_a"),
a + " carries the same grant the undivided SEND action gave an architect");
}
}
/** The other half of the same split: a worker is excluded from all three, as it was from one. */
@Test
void aWorkerMayNotUseAnySendShapeOverMcp() {
FleetMcp m = mcp(true);
for (Authz.Action a : new Authz.Action[]{Authz.Action.SEND, Authz.Action.COORD_SEND,
Authz.Action.ANSWER}) {
McpSchema.CallToolResult denied = m.denyFor(WORKER_A, a, "term_a");
assertNotNull(denied, a + " must stay refused to a worker");
assertTrue(denied.isError(), "a refusal is returned as an MCP tool error");
}
}
/**
* fleetd #669 Unit A / #678: {@code TASK_READ} (ticket polling, session status) is split out of
* the roster-only {@code READ}, carrying forward the grant the undivided action gave. A worker
* still has both — it never gained or lost anything by the split.
*/
@Test
void aWorkerKeepsBothReadActionsAfterTheSplit() {
FleetMcp m = mcp(true);
assertNull(m.denyFor(WORKER_A, Authz.Action.READ, null));
assertNull(m.denyFor(WORKER_A, Authz.Action.TASK_READ, null));
}
@Test
void anArchitectMayReplyAndAskOnlyAsItsOwnPaneOverMcp() {
FleetMcp m = mcp(true);
@@ -297,35 +336,53 @@ class FleetMcpAuthzTest {
* the whole time the defect was live -- the table was right, the action fed to it was wrong.
*/
@Test
void pollingByTargetIsADrainAndPollingByTicketIsARead() {
void pollingByTargetIsADrainAndPollingByTicketIsATaskRead() {
assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b", null),
"poll by target removes the replies — that is a drain, not an observation");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, null),
"poll by ticket changes nothing");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(" ", null),
assertEquals(Authz.Action.TASK_READ, FleetMcp.pollAction(null, null),
"poll by ticket changes nothing, but is not the roster-only READ action");
assertEquals(Authz.Action.TASK_READ, FleetMcp.pollAction(" ", null),
"a blank target is an absent target");
}
/**
* fleetd #421: a coordId branch is a THIRD operation behind fleet_poll's one name, and it must
* map to {@link Authz.Action#COORD_READ} — never {@link Authz.Action#READ}, even though this
* branch also consumes nothing. READ's grant is open to every authenticated role on the premise
* that the roster carries no secrets; a lead-to-lead body is not the roster, so folding this
* branch into READ would let any worker read every peer lead's mail in full. coordId also takes
* priority over target when both happen to be present — it addresses a different inbox entirely.
* map to {@link Authz.Action#COORD_READ} — never {@link Authz.Action#READ} or {@link
* Authz.Action#TASK_READ}, even though this branch also consumes nothing. A lead-to-lead body
* is not the roster and not a ticket/status read, so folding this branch into either would let
* any worker or architect read every peer lead's mail in full. coordId also takes priority over
* target when both happen to be present — it addresses a different inbox entirely.
*/
@Test
void pollingByCoordIdIsACoordReadNeverAPlainRead() {
void pollingByCoordIdIsACoordReadNeverAPlainOrTaskRead() {
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(null, "mac-opus"),
"reading held peer mail must not be mapped to the everyone-readable READ action");
"reading held peer mail must not be mapped to a widely-readable action");
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(" ", "mac-opus"),
"a blank target must not fall through to READ/DRAIN when coordId is present");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, " "),
assertEquals(Authz.Action.TASK_READ, FleetMcp.pollAction(null, " "),
"a blank coordId is an absent coordId, same as target/ticket");
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction("term_b", "mac-opus"),
"coordId takes priority over target — this is a different inbox, not a drain");
}
/**
* {@code fleet_send} is three call shapes behind one tool name, exactly as {@code fleet_poll}
* is (fleetd #669 Unit A). {@link FleetMcp#sendAction} picks the action from the arguments, not
* the handler, for the same reason {@link FleetMcp#pollAction} does: a test can assert the
* mapping the handler actually uses.
*/
@Test
void sendMapsToThreeDifferentActionsByItsArguments() {
assertEquals(Authz.Action.SEND, FleetMcp.sendAction(null, null),
"a plain delivery, with neither coordId nor turnId, is a local SEND");
assertEquals(Authz.Action.COORD_SEND, FleetMcp.sendAction("mac-opus", null),
"coordId addresses a peer lead over the coordination broker");
assertEquals(Authz.Action.ANSWER, FleetMcp.sendAction(null, "turn-1"),
"turnId resolves a worker's blocked question");
assertEquals(Authz.Action.COORD_SEND, FleetMcp.sendAction("mac-opus", "turn-1"),
"coordId takes priority over turnId, mirroring sendToLead's own mutual-exclusion check");
}
@Test
void everyRegisteredToolHasItsHandlerActionPinned() {
// fleetd #469: this used to scrape FleetMcp.java's tool("…") calls for the registered set —
@@ -344,16 +401,22 @@ class FleetMcpAuthzTest {
() -> tool + " is registered but has no pinned authorization action"));
assertEquals(Authz.Action.SEND, FleetMcp.toolAction("fleet_send", Map.of()));
assertEquals(Authz.Action.SEND,
FleetMcp.toolAction("fleet_send", Map.of("sessionId", "term_a", "content", "hi")));
assertEquals(Authz.Action.COORD_SEND,
FleetMcp.toolAction("fleet_send", Map.of("coordId", "mac-opus", "content", "hi")));
assertEquals(Authz.Action.ANSWER,
FleetMcp.toolAction("fleet_send", Map.of("turnId", "turn-1", "content", "hi")));
assertEquals(Authz.Action.REPLY, FleetMcp.toolAction("fleet_reply", Map.of()));
assertEquals(Authz.Action.ASK, FleetMcp.toolAction("fleet_ask", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_status", Map.of()));
assertEquals(Authz.Action.TASK_READ, FleetMcp.toolAction("fleet_status", Map.of()));
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_ack", Map.of()));
assertEquals(Authz.Action.SPAWN, FleetMcp.toolAction("fleet_spawn", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_list", Map.of()));
assertEquals(Authz.Action.STOP, FleetMcp.toolAction("fleet_stop", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_profiles", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_whoami", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task")));
assertEquals(Authz.Action.TASK_READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task")));
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_poll", Map.of("target", "term_b")));
assertEquals(Authz.Action.COORD_READ,
FleetMcp.toolAction("fleet_poll", Map.of("coordId", "mac-opus")));
@@ -122,12 +122,28 @@ class FleetAppAuthTest {
assertEquals(Authz.Action.DRAIN, FleetApp.routeAction("GET /sessions/{id}/replies"));
assertEquals(Authz.Action.ASK, FleetApp.routeAction("POST /sessions/{id}/ask"));
for (String route : Set.of("GET /sessions", "GET /agents", "GET /members", "GET /profiles",
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}")) {
"GET /member-credentials")) {
assertEquals(Authz.Action.READ, FleetApp.routeAction(route), route);
}
for (String route : Set.of("GET /sessions/{id}/status", "GET /tasks/{ticket}")) {
assertEquals(Authz.Action.TASK_READ, FleetApp.routeAction(route), route);
}
assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz"));
}
/**
* fleetd #669 Unit A: {@code POST /sessions/{id}/message} is two call shapes behind one route,
* mirroring {@code fleet_send}'s MCP-side split into {@link Authz.Action#SEND} and {@link
* Authz.Action#ANSWER} ({@code FleetMcp#sendAction}). The route never carries a {@code coordId}
* shape — that peer-lead route is MCP-only — so only these two apply here.
*/
@Test
void theMessageRouteIsASendWithNoTurnIdAndAnAnswerWithOne() {
assertEquals(Authz.Action.SEND, FleetApp.routeAction("POST /sessions/{id}/message", null));
assertEquals(Authz.Action.SEND, FleetApp.routeAction("POST /sessions/{id}/message", " "));
assertEquals(Authz.Action.ANSWER, FleetApp.routeAction("POST /sessions/{id}/message", "turn-1"));
}
private static Set<String> routesTheServerRegisters() {
try {
String source = Files.readString(REST_SOURCE).lines()
@@ -192,6 +208,21 @@ class FleetAppAuthTest {
"draining an inbox is the primary's collection step");
}
/**
* fleetd #669 Unit A: the {@code turnId} shape of {@code POST /sessions/{id}/message} maps to
* {@link Authz.Action#ANSWER}, not the plain {@link Authz.Action#SEND} the test above drives —
* a worker must stay refused on this shape too, exactly as it was refused on the one undivided
* action before the split.
*/
@Test
void aWorkerMayNotAnswerAnotherSessionsBlockedQuestionOverRest() throws Exception {
int port = start(FakeHerdr.WORKER_PID, false, null);
assertEquals(403, send(port, "POST", "/sessions/term_b/message",
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null).statusCode(),
"resolving another session's blocked question would be a worker escalating too");
}
// --- token mode ---------------------------------------------------------------------------
@Test
+10 -4
View File
@@ -580,6 +580,12 @@ install_candidate() {
# -------------------------------------------------------------------------------- the report path
#
# Masks basic-auth userinfo (scheme://user:pass@host) in a daemon verdict line before it reaches
# the terminal.
mask_verdict_userinfo() {
printf '%s\n' "$1" | sed -E 's#://[^@/[:space:]]*@#://<redacted>@#g'
}
# Prints the literal command the operator (or a test) can run to restore the backup by hand — the
# absolute path to THIS script plus the overrides actually in force, so it works from any cwd.
restore_command_line() {
@@ -599,8 +605,8 @@ restore_and_confirm() {
ok "restored from $backup"
if wait_for_verdict "$LOG" "$mark2" "$WAIT_SECONDS"; then
case "$VERDICT_KIND" in
refused) warn "the RESTORE was also refused by the daemon: $VERDICT_LINE" ;;
*) ok "restore confirmed: $VERDICT_LINE" ;;
refused) warn "the RESTORE was also refused by the daemon: $(mask_verdict_userinfo "$VERDICT_LINE")" ;;
*) ok "restore confirmed: $(mask_verdict_userinfo "$VERDICT_LINE")" ;;
esac
else
warn "the restore is on disk, but no confirming verdict line appeared within ${WAIT_SECONDS}s"
@@ -616,7 +622,7 @@ report_outcome() {
say "waiting for the daemon's verdict (up to ${WAIT_SECONDS}s)"
if wait_for_verdict "$LOG" "$mark" "$WAIT_SECONDS"; then
kind="$VERDICT_KIND"; line="$VERDICT_LINE"
kind="$VERDICT_KIND"; line="$(mask_verdict_userinfo "$VERDICT_LINE")"
else
kind="none"
fi
@@ -673,7 +679,7 @@ check_mode() {
local verdict
verdict="$(last_verdict_line "$LOG")"
if [ -n "$verdict" ]; then
ok "last verdict in log: $verdict"
ok "last verdict in log: $(mask_verdict_userinfo "$verdict")"
else
warn "no reload verdict line found in $LOG"
fi
+65
View File
@@ -629,6 +629,65 @@ test_refusal_shape_from_parse_failure_wording_is_recognised() {
assert_equals 4 "$RUN_RC" "the parse-failure refusal shape must also exit 4, not be read as silence"
}
# A verdict line carrying a credentialed URI has its userinfo masked, with a positive control
# proving the rest of the line still reaches the output unchanged.
test_verdict_userinfo_is_masked_with_positive_control() {
local dir
dir="$(new_fixture)"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf 'config reload from %s refused, keeping the running config: refusing to start: malformed pattern — profiles.local.errorPattern ("amqp://user:hunter2@host/vhost"): Unclosed character class near index 8\n' \
"$dir/fleetd.yaml" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "refusal-with-userinfo exit code"
assert_not_contains "user:hunter2" "$RUN_OUTPUT" "the userinfo must never reach the output"
assert_contains "amqp://<redacted>@host/vhost" "$RUN_OUTPUT" \
"the userinfo must be MASKED, not deleted — the rest of the quoted value must survive"
# Positive control: the diagnostic prose on both sides of the userinfo must still reach the
# output. Without this, a mutant that drops the whole verdict line would pass identically.
assert_contains "malformed pattern" "$RUN_OUTPUT" "prose BEFORE the userinfo must still reach the output"
assert_contains "Unclosed character class near index 8" "$RUN_OUTPUT" \
"prose AFTER the userinfo must still reach the output"
}
# An ordinary refusal line quotes the offending pattern, not a credential, and must survive byte
# for byte: the rewrite is scoped to userinfo only, and the quoted pattern is the detail an
# operator needs to fix the refusal.
test_ordinary_refusal_line_passes_through_unchanged() {
local dir real_line
dir="$(new_fixture)"
real_line="config reload from $dir/fleetd.yaml refused, keeping the running config: refusing to start: malformed pattern — profiles.local.errorPattern (\"[unclosed\"): Unclosed character class near index 8"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf '%s\n' "$real_line" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "ordinary refusal exit code"
assert_contains "$real_line" "$RUN_OUTPUT" \
"an ordinary refusal with no userinfo must pass through byte for byte, unchanged"
}
# A verdict line can hold a URI with NO userinfo and a later, unrelated @ further on in the same
# line (an email address in diagnostic prose, for example). The rewrite must stop at the end of
# the URI and must not treat the later @ as a second userinfo delimiter.
test_uri_without_userinfo_survives_a_later_at_sign() {
local dir real_line
dir="$(new_fixture)"
real_line="config reload from $dir/fleetd.yaml refused, keeping the running config: broker.uri amqp://broker.local/vhost unreachable, contact ops@example.com"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf '%s\n' "$real_line" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "no-userinfo-with-later-at-sign exit code"
assert_contains "$real_line" "$RUN_OUTPUT" \
"a URI with no userinfo plus a later @ in the same line must pass through byte for byte"
}
# --set runs yq over the whole candidate. It warns when that changes more lines than the requested
# pairs, but a simple file with only the intended changed line must stay quiet.
new_fixture_reformat_sensitive() {
@@ -720,6 +779,12 @@ echo "== extra: --check is read-only and always exits 0 =="
test_check_is_read_only_and_exits_zero
echo "== extra: the parse-failure refusal shape is also recognised =="
test_refusal_shape_from_parse_failure_wording_is_recognised
echo "== verdict-redaction criteria 2+3: verdict userinfo is masked, rest of line survives =="
test_verdict_userinfo_is_masked_with_positive_control
echo "== verdict-redaction criterion 4: an ordinary refusal passes through unchanged =="
test_ordinary_refusal_line_passes_through_unchanged
echo "== fleetd #638: a URI with no userinfo survives a later @ in the same line =="
test_uri_without_userinfo_survives_a_later_at_sign
echo "== acceptance criterion 17: --set warns about yq formatting churn =="
test_set_warns_when_yq_reformats_extra_lines
echo "== acceptance criterion 18: --set stays quiet without formatting churn =="