Compare commits
24 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ef4996a01e | |||
| edbd8d816a | |||
| 9425a9b696 | |||
| d0688c8a60 | |||
| 2c467c2553 | |||
| c6430d8edd | |||
| bfee23acc3 | |||
| 7dec74f1b4 | |||
| 482598e2a6 | |||
| 5f5d16fbd4 | |||
| 7df7985a16 | |||
| 724b35b46e | |||
| 2eb2d6112e | |||
| 7c458e8bf2 | |||
| 133f03e428 | |||
| 804279175d | |||
| 9dea289975 | |||
| cb4a6869b9 | |||
| 28b45d97e5 | |||
| 1a397e962e | |||
| 209e1231ea | |||
| 5d8b9d365c | |||
| a6aeda39e7 | |||
| 31b3c24caa |
@@ -164,6 +164,13 @@ fails.
|
||||
- **`accepted` does not mean your pane has been cleared.** It means every gate passed and the roll
|
||||
is scheduled to run once your current turn ends. Say your goodbye in the same turn — you will not
|
||||
get another one.
|
||||
- **If you are still running after that turn, the roll did not happen.** A roll that works clears
|
||||
you, so surviving your own goodbye is itself the signal that it refused. Check with
|
||||
`fleet_handover{action: "status", token}`, using the token you confirmed. `TURN_NEVER_SETTLED`
|
||||
means your turn ran past `leadRollover.turnSettleSeconds` and **no `/clear` was ever sent**: your
|
||||
context is intact and nothing was lost. Open a fresh request and retry. Never assume the roll
|
||||
succeeded because `confirm` answered `accepted` — by the time it refuses, there is no caller left
|
||||
to tell, so this check is the only thing that closes that gap.
|
||||
- **There is no terminal or session parameter, on purpose.** The pane is always your own, resolved
|
||||
from your connection, so you can only ever roll yourself.
|
||||
- **`operatorConfirmed` is your report of what a human told you.** Do not pass `true` because you
|
||||
|
||||
@@ -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}";
|
||||
|
||||
@@ -1399,11 +1393,13 @@ public record FleetConfig(
|
||||
* called FROM the calling lead's own turn, so its pane is still {@code WORKING} the instant
|
||||
* {@code confirm()} validates every gate and schedules the roll. {@code
|
||||
* dev.ltms.fleet.lead.LeadRollover}'s deferred continuation waits up to this many seconds for
|
||||
* that SAME pane to report an injectable state again — i.e. for the calling turn to actually
|
||||
* end — before it sends {@code /clear} at all. If that wait times out, no {@code /clear} is
|
||||
* that SAME pane to report {@code IDLE} or {@code DONE} — i.e. for the calling turn to actually
|
||||
* end — before it sends {@code /clear} at all. {@code BLOCKED} does not count: that is a live
|
||||
* turn merely paused, not one that has finished. If that wait times out, no {@code /clear} is
|
||||
* ever sent: a lead that never goes idle is still doing real work, and clearing it would
|
||||
* destroy live context. This is a separate wait from {@code clearSettleSeconds} below, which
|
||||
* bounds the SECOND wait, for the pane to re-settle AFTER {@code /clear} has already gone out.
|
||||
* bounds the SECOND wait, for the pane to reach {@code IDLE} or {@code DONE} again AFTER
|
||||
* {@code /clear} has already gone out.
|
||||
*
|
||||
* @param handoverPath required when this block is present — where the handover file a fresh
|
||||
* lead session reads must live. There is no sane non-null default for an
|
||||
@@ -1422,12 +1418,13 @@ public record FleetConfig(
|
||||
* @param maxDocAgeSeconds default 3600 — refuse a handover file whose modified time is older
|
||||
* than this many seconds, so a stale leftover from an earlier rollover
|
||||
* attempt can never be mistaken for a fresh one.
|
||||
* @param turnSettleSeconds default 20 — bound on how long the deferred roll waits for the
|
||||
* CALLING lead's own turn to end (its pane to report injectable again)
|
||||
* before sending {@code /clear} at all. See the paragraph above.
|
||||
* @param turnSettleSeconds default 300 — bound on how long the deferred roll waits for the
|
||||
* CALLING lead's own turn to end (its pane to report {@code IDLE} or
|
||||
* {@code DONE}) before sending {@code /clear} at all. See the paragraph
|
||||
* above.
|
||||
* @param clearSettleSeconds default 20 — bound on how long to wait for the lead's pane to
|
||||
* report an injectable state again after {@code /clear} before giving up. A
|
||||
* roll that times out here never sends {@code bootstrapText}.
|
||||
* report {@code IDLE} or {@code DONE} again after {@code /clear} before
|
||||
* giving up. A roll that times out here never sends {@code bootstrapText}.
|
||||
* @param bootstrapText default a sentence naming the RESOLVED handover path — sent to the
|
||||
* lead's pane once it settles after {@code /clear}, telling the fresh
|
||||
* session where to read the handover and carry on. Left {@code null} here
|
||||
@@ -1444,7 +1441,7 @@ public record FleetConfig(
|
||||
public LeadRollover {
|
||||
requireOperatorConfirm = requireOperatorConfirm == null || requireOperatorConfirm;
|
||||
maxDocAgeSeconds = (maxDocAgeSeconds == null || maxDocAgeSeconds <= 0) ? 3600 : maxDocAgeSeconds;
|
||||
turnSettleSeconds = (turnSettleSeconds == null || turnSettleSeconds <= 0) ? 20 : turnSettleSeconds;
|
||||
turnSettleSeconds = (turnSettleSeconds == null || turnSettleSeconds <= 0) ? 300 : turnSettleSeconds;
|
||||
clearSettleSeconds = (clearSettleSeconds == null || clearSettleSeconds <= 0) ? 20 : clearSettleSeconds;
|
||||
bootstrapText = (bootstrapText == null || bootstrapText.isBlank()) ? null : bootstrapText;
|
||||
}
|
||||
@@ -2682,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()) {
|
||||
@@ -2714,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());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -2797,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,22 +49,52 @@ 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);
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* The second gate for {@code POST /sessions/{id}/message}: checked only when {@code turnId}
|
||||
* is present and non-blank, against {@link Authz.Action#ANSWER}. A request with no {@code
|
||||
* turnId} passes this gate unconditionally, without consulting {@code permit} at all, having
|
||||
* already cleared the coarse {@link Authz.Action#SEND} grant checked ahead of it.
|
||||
*
|
||||
* @param permit reports whether the caller holds the named grant
|
||||
*/
|
||||
static boolean answerGatePasses(String turnId, Predicate<Authz.Action> permit) {
|
||||
return turnId == null || turnId.isBlank() || permit.test(Authz.Action.ANSWER);
|
||||
}
|
||||
|
||||
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
|
||||
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
|
||||
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
|
||||
@@ -250,7 +280,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 +634,37 @@ 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 coarse {@link
|
||||
* Authz.Action#SEND} grant is checked first, before the body is read at all; only once that
|
||||
* passes is the body parsed, and a present {@code turnId} is then checked again against
|
||||
* {@link Authz.Action#ANSWER}. A body that fails to parse is rejected with 400 and reaches
|
||||
* neither {@code messages.answer} nor {@code messages.send}.
|
||||
*/
|
||||
private void sendMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
|
||||
return;
|
||||
}
|
||||
String content;
|
||||
String turnId;
|
||||
long timeout;
|
||||
boolean wait;
|
||||
JsonNode body;
|
||||
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)
|
||||
body = mapper.readTree(ctx.body());
|
||||
} catch (Exception e) {
|
||||
body = null;
|
||||
}
|
||||
if (body == null) {
|
||||
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
|
||||
return;
|
||||
}
|
||||
String turnId = body.path("turnId").asText(null);
|
||||
if (!answerGatePasses(turnId, action -> allow(ctx, action, id))) {
|
||||
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 ────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
@@ -3175,4 +3256,41 @@ class FleetConfigTest {
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertTrue(cfg.models().offIds().isEmpty());
|
||||
}
|
||||
|
||||
// ── fleetd #651: leadRollover.turnSettleSeconds default resolution ─────────────────────────
|
||||
|
||||
@Test
|
||||
void turnSettleSecondsDefaultsTo300WhenUnset(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("bare-rollover.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\nleadRollover: {}\n");
|
||||
|
||||
FleetConfig.LeadRollover rollover = FleetConfig.load(f).leadRollover();
|
||||
assertNotNull(rollover);
|
||||
assertEquals(300, rollover.turnSettleSeconds());
|
||||
}
|
||||
|
||||
@Test
|
||||
void turnSettleSecondsUsesAnExplicitPositiveValue(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("rollover.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
leadRollover:
|
||||
turnSettleSeconds: 45
|
||||
""");
|
||||
|
||||
FleetConfig.LeadRollover rollover = FleetConfig.load(f).leadRollover();
|
||||
assertEquals(45, rollover.turnSettleSeconds());
|
||||
}
|
||||
|
||||
@Test
|
||||
void turnSettleSecondsFallsBackTo300WhenZeroOrNegative(@TempDir Path dir) throws Exception {
|
||||
Path zero = dir.resolve("zero.yaml");
|
||||
Files.writeString(zero, "bind:\n port: 8080\nleadRollover:\n turnSettleSeconds: 0\n");
|
||||
assertEquals(300, FleetConfig.load(zero).leadRollover().turnSettleSeconds());
|
||||
|
||||
Path negative = dir.resolve("negative.yaml");
|
||||
Files.writeString(negative, "bind:\n port: 8080\nleadRollover:\n turnSettleSeconds: -5\n");
|
||||
assertEquals(300, FleetConfig.load(negative).leadRollover().turnSettleSeconds());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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")));
|
||||
|
||||
@@ -60,13 +60,25 @@ class MessageServiceTest {
|
||||
inbox.own(T);
|
||||
}
|
||||
|
||||
/**
|
||||
* A send budget large enough that a test's own setup — {@link #awaitWaiting()} plus whatever
|
||||
* status transitions it drives afterward — can never compete with it for the same clock. A test
|
||||
* that needs {@code send.get(...)}'s own window to be the only timing bound it depends on uses
|
||||
* {@link #sendAsync(String, long)} with this value instead of the default 5000 ms.
|
||||
*/
|
||||
private static final long GENEROUS_SEND_BUDGET_MILLIS = 30_000;
|
||||
|
||||
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
|
||||
private CompletableFuture<MessageService.Reply> sendAsync() {
|
||||
return sendAsync("do the task");
|
||||
}
|
||||
|
||||
private CompletableFuture<MessageService.Reply> sendAsync(String content) {
|
||||
return CompletableFuture.supplyAsync(() -> messages.send(T, content, 5000));
|
||||
return sendAsync(content, 5000);
|
||||
}
|
||||
|
||||
private CompletableFuture<MessageService.Reply> sendAsync(String content, long timeoutMillis) {
|
||||
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis));
|
||||
}
|
||||
|
||||
private void awaitWaiting() throws InterruptedException {
|
||||
@@ -80,7 +92,7 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void completionFallbackResolvesATurnThatNeverCalledFleetReply() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", GENEROUS_SEND_BUDGET_MILLIS);
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt"); // pre-turn pane: no answer yet (baseline reference)
|
||||
@@ -96,6 +108,32 @@ class MessageServiceTest {
|
||||
assertTrue(reply.completed(), "a scraped completion still counts as completed");
|
||||
}
|
||||
|
||||
/**
|
||||
* Pins {@link #GENEROUS_SEND_BUDGET_MILLIS} as the budget {@link
|
||||
* #completionFallbackResolvesATurnThatNeverCalledFleetReply} depends on. A 5500 ms delay between
|
||||
* {@link #awaitWaiting()} and the status transitions that drive completion stands in for a loaded
|
||||
* machine's setup overhead — comfortably past the 5000 ms budget this send no longer uses, and
|
||||
* still well inside this method's own 30 000 ms budget. The only clock this test depends on is
|
||||
* {@code send.get}'s own 10 s window.
|
||||
*/
|
||||
@Test
|
||||
void completionFallbackSurvivesASlowHarnessBecauseItsSendBudgetIsNotTheBindingClock() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", GENEROUS_SEND_BUDGET_MILLIS);
|
||||
awaitWaiting();
|
||||
|
||||
Thread.sleep(5500);
|
||||
|
||||
herdr.readText("$ prompt");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
herdr.readText("BUILD GREEN: 391 files");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
MessageService.Reply reply = send.get(10, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
|
||||
"a slow harness must not be mistaken for a timed-out delivery");
|
||||
}
|
||||
|
||||
@Test
|
||||
void completionFallbackReplacesAnEchoedInjectedBriefWithNoReportOutcome() throws Exception {
|
||||
String brief = "Implement the requested change. ".repeat(20);
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
package dev.ltms.fleet.rest;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.Authz;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
@@ -18,6 +21,7 @@ import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.FakeWorktrees;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -28,10 +32,13 @@ import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
@@ -122,12 +129,60 @@ 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"));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #689: {@code answerGatePasses} is the second, conditional gate behind {@code
|
||||
* sendMessage}'s coarse {@link Authz.Action#SEND} check. With the {@code ANSWER} grant denied,
|
||||
* a {@code turnId}-bearing request is refused while a plain one still passes — and the denied
|
||||
* permit is queried only for the {@code turnId} case, never for the plain one, which is what
|
||||
* proves this is a genuinely separate, conditional check rather than the {@code SEND} check
|
||||
* renamed or an unconditional call whose result is ignored. Flipping only the {@code ANSWER}
|
||||
* grant to allowed then flips only the {@code turnId} shape's outcome.
|
||||
*/
|
||||
@Test
|
||||
void answerGatePassesOnlyWhenTurnIdAbsentOrAnswerGranted() {
|
||||
List<Authz.Action> queried = new ArrayList<>();
|
||||
Predicate<Authz.Action> denyAnswer = action -> {
|
||||
queried.add(action);
|
||||
return false;
|
||||
};
|
||||
|
||||
assertFalse(FleetApp.answerGatePasses("turn-1", denyAnswer),
|
||||
"ANSWER denied ⇒ the turnId shape is refused");
|
||||
assertEquals(List.of(Authz.Action.ANSWER), queried,
|
||||
"the ANSWER grant, specifically, must be the one consulted");
|
||||
|
||||
queried.clear();
|
||||
assertTrue(FleetApp.answerGatePasses(null, denyAnswer),
|
||||
"no turnId ⇒ the plain shape passes even though ANSWER is denied");
|
||||
assertTrue(FleetApp.answerGatePasses(" ", denyAnswer), "a blank turnId is treated as absent");
|
||||
assertEquals(List.of(), queried, "the plain shape must never consult the permit at all");
|
||||
|
||||
assertTrue(FleetApp.answerGatePasses("turn-1", action -> true),
|
||||
"flipping only the ANSWER grant to allowed flips only the turnId shape's outcome");
|
||||
}
|
||||
|
||||
private static Set<String> routesTheServerRegisters() {
|
||||
try {
|
||||
String source = Files.readString(REST_SOURCE).lines()
|
||||
@@ -192,6 +247,107 @@ 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");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #689: a caller refused the coarse {@link Authz.Action#SEND} grant is refused on
|
||||
* {@code SEND} specifically, even on the {@code turnId}-bearing shape that otherwise raises
|
||||
* the check to {@link Authz.Action#ANSWER} — proving {@code turnId} was never read from the
|
||||
* body before the refusal (reading it would have changed which action is named in the 403).
|
||||
* The same caller refused with no body at all gets the identical detail, which could not hold
|
||||
* if the decision depended on anything read from the body. Control: a caller who IS granted
|
||||
* reaches past the gate and the body is used normally.
|
||||
*/
|
||||
@Test
|
||||
void aDeniedCallerIsRefusedOnSendEvenWithATurnIdBodyAndNeverReadsTheBody() throws Exception {
|
||||
int workerPort = start(FakeHerdr.WORKER_PID, false, null); // denied: not primary/architect
|
||||
|
||||
HttpResponse<String> withTurnId = send(workerPort, "POST", "/sessions/term_b/message",
|
||||
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
|
||||
assertEquals(403, withTurnId.statusCode());
|
||||
assertTrue(withTurnId.body().contains("may not SEND"),
|
||||
"the SEND check must be the one that fired, not ANSWER — ANSWER would only be "
|
||||
+ "reachable by having already read turnId out of the body");
|
||||
|
||||
HttpResponse<String> noBody = send(workerPort, "POST", "/sessions/term_b/message", null, null);
|
||||
assertEquals(403, noBody.statusCode());
|
||||
assertTrue(noBody.body().contains("may not SEND"),
|
||||
"refused identically with no body at all — the refusal cannot depend on body content");
|
||||
|
||||
// Control: a primary IS granted SEND, so the same turnId body is read and acted on —
|
||||
// reaching messages.answer, which reports this unknown turnId as a stale one.
|
||||
int primaryPort = start(999_999, false, null);
|
||||
HttpResponse<String> granted = send(primaryPort, "POST", "/sessions/term_b/message",
|
||||
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
|
||||
assertEquals(409, granted.statusCode());
|
||||
assertTrue(granted.body().contains("stale_turn"), "a granted caller's body IS read and acted on");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #689 (ticket comment 18353): the only place {@code sendMessage}'s call to {@code
|
||||
* answerGatePasses} is observable is the audit trail — {@code allow()} logs an {@code
|
||||
* "allowed"} entry for every granted action except {@code READ}/{@code METRICS}/{@code
|
||||
* TASK_READ}, and {@code ANSWER} is none of those. A granted {@code turnId} request must
|
||||
* therefore log both a {@code SEND} and an {@code ANSWER} entry; a granted plain request must
|
||||
* log {@code SEND} alone. A unit test of the extracted helper pins the helper; this pins the
|
||||
* call site — deleting the {@code answerGatePasses} call from {@code sendMessage} leaves the
|
||||
* helper's own test green but turns this one red.
|
||||
*/
|
||||
@Test
|
||||
void aGrantedTurnIdRequestAuditsBothSendAndAnswerButAPlainRequestAuditsSendAlone() throws Exception {
|
||||
int port = start(999_999, false, null); // primary: granted both SEND and ANSWER
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
try (CapturedLog audit = CapturedLog.at("audit", Level.INFO)) {
|
||||
send(port, "POST", "/sessions/term_b/message",
|
||||
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
|
||||
|
||||
List<String> allowed = allowedActions(audit, mapper);
|
||||
assertTrue(allowed.contains("SEND"),
|
||||
"a turnId request must still clear the coarse SEND grant first");
|
||||
assertTrue(allowed.contains("ANSWER"),
|
||||
"a turnId request must ALSO clear the ANSWER grant — this is the call site itself");
|
||||
}
|
||||
|
||||
try (CapturedLog audit = CapturedLog.at("audit", Level.INFO)) {
|
||||
send(port, "POST", "/sessions/term_b/message",
|
||||
"{\"content\":\"hi\",\"timeoutMs\":50}", null);
|
||||
|
||||
List<String> allowed = allowedActions(audit, mapper);
|
||||
assertEquals(List.of("SEND"), allowed,
|
||||
"a plain request must log SEND and nothing else — ANSWER is conditional on "
|
||||
+ "turnId, not something every request happens to log");
|
||||
}
|
||||
}
|
||||
|
||||
private static List<String> allowedActions(CapturedLog audit, ObjectMapper mapper) {
|
||||
return audit.events().stream()
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.map(line -> {
|
||||
try {
|
||||
return mapper.readTree(line);
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("audit line is not valid JSON: " + line, e);
|
||||
}
|
||||
})
|
||||
.filter(n -> "allowed".equals(n.path("outcome").asText()))
|
||||
.map(n -> n.path("action").asText())
|
||||
.toList();
|
||||
}
|
||||
|
||||
// --- token mode ---------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -675,6 +675,26 @@ class FleetAppTest {
|
||||
assertEquals(400, postMessage(port, "{}").statusCode());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #689: a body that fails to parse is rejected with 400 before {@code turnId} is ever
|
||||
* read from it, so it reaches neither {@code messages.answer} (which needs a {@code turnId})
|
||||
* nor {@code messages.send} — confirmed here for {@code send} by the fake agent's idle status,
|
||||
* which would otherwise make an immediate {@code agent.prompt} delivery observable.
|
||||
*/
|
||||
@Test
|
||||
void malformedBodyReturns400AndNeverReachesSendOrAnswer() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // idle ⇒ send would deliver right away if reached
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> res = postMessage(port, "not json at all");
|
||||
assertEquals(400, res.statusCode());
|
||||
JsonNode err = mapper.readTree(res.body());
|
||||
assertEquals("bad_request", err.get("error").asText());
|
||||
assertEquals("body must be JSON", err.get("detail").asText());
|
||||
assertFalse(herdr.called("agent.prompt"),
|
||||
"a malformed body must never reach messages.send's delivery");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionStatusReportsLiveAgentStatus() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("blocked");
|
||||
|
||||
+22
-5
@@ -213,6 +213,17 @@ map_masked_lines() {
|
||||
done < "$file"
|
||||
}
|
||||
|
||||
# Masks every `scheme://user:pass@host` userinfo on one line of text, replacing just that
|
||||
# userinfo with `<redacted>` and leaving the rest of the line untouched, byte for byte. The
|
||||
# pattern stops at the first `/`, whitespace, or `@` reached after `://` — a URI's userinfo
|
||||
# component cannot contain any of those three characters — so a URI with no userinfo, followed
|
||||
# later on the same line by an unrelated `@`, never matches. The `g` flag matters: a line can
|
||||
# carry more than one URI. Shared by every caller that prints a line which may hold a
|
||||
# credentialed URI, so the bound lives in exactly one place.
|
||||
mask_url_userinfo() {
|
||||
printf '%s\n' "$1" | sed -E 's#://[^@/[:space:]]*@#://<redacted>@#g'
|
||||
}
|
||||
|
||||
redact() {
|
||||
local old_file="$1" new_file="$2"
|
||||
local line prefix content indent lead key old_line=0 new_line=0 in_hunk=0
|
||||
@@ -270,7 +281,7 @@ redact() {
|
||||
continue
|
||||
fi
|
||||
fi
|
||||
printf '%s\n' "$line" | sed -E 's#://[^@]*@#://<redacted>@#g'
|
||||
mask_url_userinfo "$line"
|
||||
done
|
||||
[ "$saved_nocasematch" = 1 ] || shopt -u nocasematch
|
||||
}
|
||||
@@ -580,6 +591,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() {
|
||||
mask_url_userinfo "$1"
|
||||
}
|
||||
|
||||
# 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 +616,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 +633,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 +690,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
|
||||
|
||||
@@ -233,6 +233,66 @@ test_redaction_holds() {
|
||||
assert_contains "weight" "$RUN_OUTPUT" "a diff must have been demonstrably printed at all"
|
||||
}
|
||||
|
||||
# redact()'s key-name filter only inspects the KEY, so a diff line whose key does not match
|
||||
# TOKEN|SECRET|PASSWORD|PASSWD|PASSPHRASE|CREDENTIAL|URI|_KEY still reaches the final userinfo
|
||||
# sed even when its VALUE holds a credentialed URI. "note" is not a sensitive key name, so this
|
||||
# line must fall all the way through to that sed, not the earlier whole-value branch. The
|
||||
# trailing prose on both sides of the userinfo is a positive control: it proves the line reached
|
||||
# the userinfo sed (which touches only the userinfo) rather than the earlier branch (which would
|
||||
# have replaced the whole value with a bare "<redacted>" and dropped the prose).
|
||||
test_diff_line_userinfo_is_masked_with_positive_control() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.note=see amqp://alice:wonderland@rabbit.local:5672/vhost for details'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "diff-userinfo-case reload exit code"
|
||||
assert_not_contains "alice:wonderland" "$RUN_OUTPUT" "the userinfo must never reach the output"
|
||||
assert_contains "amqp://<redacted>@rabbit.local:5672/vhost" "$RUN_OUTPUT" \
|
||||
"the userinfo must be MASKED, not deleted — the rest of the value must survive"
|
||||
assert_contains "note:" "$RUN_OUTPUT" "the key name must still reach the output"
|
||||
assert_contains "see " "$RUN_OUTPUT" "prose BEFORE the userinfo must still reach the output"
|
||||
assert_contains "for details" "$RUN_OUTPUT" "prose AFTER the userinfo must still reach the output"
|
||||
}
|
||||
|
||||
# A diff line can hold a URL with no userinfo, followed later on the same line by an unrelated @
|
||||
# (free text in a string value, for example an email address). The line must pass through the
|
||||
# userinfo sed byte for byte: the match must stop at the end of the URL and must not treat the
|
||||
# later @ as a second userinfo delimiter.
|
||||
test_diff_line_uri_without_userinfo_survives_a_later_at_sign() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.note2=see https://docs.local/guide and mail ops@example.com'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "diff-no-userinfo-with-later-at-sign reload exit code"
|
||||
assert_contains "note2: see https://docs.local/guide and mail ops@example.com" "$RUN_OUTPUT" \
|
||||
"a URL with no userinfo plus a later @ on the same line must pass through byte for byte"
|
||||
}
|
||||
|
||||
# Two credentialed URIs on one diff line must both be masked — the g flag matters.
|
||||
test_diff_line_masks_multiple_userinfo_with_g_flag() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.note3=amqp://u1:p1@host1/vhost1 and amqp://u2:p2@host2/vhost2'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "diff-two-userinfo-on-one-line reload exit code"
|
||||
assert_not_contains "u1:p1" "$RUN_OUTPUT" "the first userinfo must never reach the output"
|
||||
assert_not_contains "u2:p2" "$RUN_OUTPUT" "the second userinfo must never reach the output"
|
||||
assert_contains "amqp://<redacted>@host1/vhost1" "$RUN_OUTPUT" "the first URI must be masked"
|
||||
assert_contains "amqp://<redacted>@host2/vhost2" "$RUN_OUTPUT" "the second URI must be masked"
|
||||
}
|
||||
|
||||
# ------------------------------------------------------- acceptance criterion 9: forgotten value
|
||||
# `--set .a.b=` is a plausible typo (the value simply forgotten), and it must be refused outright
|
||||
# rather than silently nulling the field — a null numeric field falls back to its default, which
|
||||
@@ -629,6 +689,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() {
|
||||
@@ -694,6 +813,12 @@ echo "== acceptance criterion 6: the marker works =="
|
||||
test_marker_skips_lines_before_it
|
||||
echo "== acceptance criterion 7 (+13: redaction is proven to have run) =="
|
||||
test_redaction_holds
|
||||
echo "== fleetd #692: a diff line's userinfo is masked, rest of the value survives =="
|
||||
test_diff_line_userinfo_is_masked_with_positive_control
|
||||
echo "== fleetd #692: a diff line's URI with no userinfo survives a later @ in the line =="
|
||||
test_diff_line_uri_without_userinfo_survives_a_later_at_sign
|
||||
echo "== fleetd #692: two userinfo URIs on one diff line are both masked =="
|
||||
test_diff_line_masks_multiple_userinfo_with_g_flag
|
||||
echo "== acceptance criterion 9: a forgotten value refuses and installs nothing =="
|
||||
test_forgotten_value_refuses_and_installs_nothing
|
||||
echo "== acceptance criterion 10: an explicit clear writes a bare null =="
|
||||
@@ -720,6 +845,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 =="
|
||||
|
||||
Reference in New Issue
Block a user