Compare commits

...

12 Commits

Author SHA1 Message Date
Dai Ha 9425a9b696 fleetd #689: pin the answerGatePasses call site via the audit trail
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Failing after 2m3s
A unit test on the extracted helper proves the helper, not the call
site in sendMessage. allow() logs an AuditLog.allowed() entry for
every granted non-READ/METRICS/TASK_READ action, so a granted turnId
request must log both SEND and ANSWER, and a granted plain request
must log SEND alone. Verified this goes red when the call site is
deleted from sendMessage, and restores to a clean diff.
2026-10-03 22:56:04 +02:00
Dai Ha c6430d8edd fleetd #689: check SEND before reading the request body in sendMessage
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m56s
Authorize twice: the coarse SEND grant first, with no body read, then
parse the body, then check ANSWER too when turnId is present. Restores
the pre-#687 ordering (no attacker-controlled body parse before the
gate) while keeping the SEND/ANSWER split #687 introduced.
2026-10-03 22:47:25 +02:00
Dai Ha 2eb2d6112e Merge PR #688: fleetd #675 — pin three unpinned FleetdAssembly constructor args
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m35s
2026-10-03 22:27:32 +02:00
Dai Ha 7c458e8bf2 Merge PR #687: fleetd #669 Unit A — split SEND and READ into their real call shapes (closes #678) 2026-10-03 22:27:32 +02:00
Dai Ha 133f03e428 Merge PR #684: fleetd #683 — decouple the completion-fallback test's two 5s budgets 2026-10-03 22:27:27 +02:00
Dai Ha 804279175d fleetd #675: pin assembly loop timing defaults
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m42s
2026-10-03 22:12:46 +02:00
Dai Ha 9dea289975 fleetd #669 Unit A: split SEND and READ into their real call shapes
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Failing after 2m7s
SEND covered three different call shapes under one action (local
sessionId delivery, the coordId cross-host broker route, and the
turnId answer-a-blocked-worker form). READ covered both roster/
profile/identity observation and ticket-polling/session-status.
Split each into its own Authz.Action — SEND/COORD_SEND/ANSWER and
READ/TASK_READ — with every new action granted to exactly who held
the combined action before, on both the MCP and REST entry paths.

Also fixes fleetd #678's Authz.java comment: READ no longer claims
"the roster carries no secrets" for ticket replies and pending
questions, because those now live under TASK_READ.
2026-10-03 22:12:13 +02:00
Dai Ha 1a397e962e fleetd #683: decouple the completion-fallback test's send budget from its own setup clock
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Failing after 1m53s
completionFallbackResolvesATurnThatNeverCalledFleetReply gave messages.send a 5000 ms budget
that started ticking the instant sendAsync() ran, then raced that same clock against
awaitWaiting()'s own 2000 ms deadline plus several onStatus/readText calls before asserting with
send.get(5, SECONDS). On a loaded machine the setup could eat enough of the 5000 ms that the
production call expired first, returning TIMED_OUT_QUEUED instead of COMPLETED_UNREPLIED.

Give this one test's send a 30 000 ms budget (sendAsync(content, timeoutMillis)) so the setup can
never compete with it; send.get(5, SECONDS) stays the one clock the test depends on. A new
regression test injects a deterministic 5500 ms delay in the same spot and proves the budget is no
longer the binding constraint — reverting it to 5000 ms turns that test red with the same
TIMED_OUT_QUEUED mismatch, confirmed by mutation.
2026-10-03 22:06:26 +02:00
Dai Ha 209e1231ea Merge PR #682: fleetd #651 — raise turnSettleSeconds default to 300, correct the settle javadoc
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m54s
The deferred roll waits for the calling lead's own turn to end. The old 20s bound
was shorter than one ordinary closing message (measured elapsed=20394ms), so a
lead that followed the handover skill's instruction to say its goodbye in the
same turn landed in the refusal branch. 300 matches idleAfterSeconds' existing
written reason: five minutes absorbs a normal pause without stalling.

The javadoc said both waits are for the pane to report an injectable state.
AgentStatus.injectable() accepts BLOCKED; the settle check requires IDLE or DONE.
The wording named the wrong predicate.
2026-10-03 21:43:25 +02:00
Dai Ha a6aeda39e7 Merge PR #681: fleetd #680 — pin the jar-path split, check the plist's jar, widen the locator
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 46s
CI / build (push) Failing after 1m50s
Verified by the lead before merge, each mutation re-run independently:
- three-dot diff: 1 commit, 3 files, nothing dragged in; clean trial merge
- suite exit 0, 134 test functions (5 new); the 3 FAIL and 5 mktemp lines are
  the documented deliberate self-test output, unchanged from baseline
- JAR moved under target/        -> RED  (the core invariant is pinned now)
- JAR moved elsewhere outside it -> PASS (not a refuse-everything test)
- PATTERN narrowed back to run/  -> RED
- plist-jar mismatch die removed -> RED
- comm=java allowlist removed    -> RED (widening PATTERN did not weaken #593)
- mvn clean install from fleetd/: Tests run: 1929, Failures: 0, 172 reports

Measured separately: during a full Maven build no extra java process matches
the widened 'fleetd.jar' pattern, so assert_single_daemon does not false-
positive while a worker builds.
2026-10-03 21:34:14 +02:00
Dai Ha 31b3c24caa fleetd #651: tell a rolling lead that surviving its goodbye means the roll refused
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 2m9s
A lead calls fleet_handover confirm from inside its own turn, and the skill
tells it to write its goodbye in that same turn. The deferred roll then waits
turnSettleSeconds for that pane to reach IDLE or DONE. Measured on 2026-10-02,
an ordinary goodbye turn took 20394ms against a 20s budget, so the roll wrote
TURN_NEVER_SETTLED and sent no /clear. confirm() had already returned accepted,
so there was no caller left to tell, and the lead carried on believing it had
been replaced.

The refusal is the correct branch — clearing a live turn would destroy context.
The gap is that nobody is told to look afterwards.

fleet_handover{action:"status", token} already reports the outcome: FleetMcp
dispatches it to LeadRollover.status, which reads the same outcomes map the
refusal writes into. Nothing new is needed in the daemon for this half.

The signal costs nothing, because a successful roll clears the lead. A lead
that is still running after its goodbye already knows the roll failed. That is
more reliable than warning it to keep the goodbye short, which would make the
feature worst exactly when it matters most.

The live budget on this host was raised to 300s separately, in fleetd.yaml,
which is hot: the config-watcher reloaded at 21:28:00 with no "needs a restart"
clause, four seconds after the edit, with no daemon restart. The code default
is a separate change.
2026-10-03 21:29:21 +02:00
Dai Ha d105da978d fleetd #680: pin the JAR/BUILD_JAR split, check the plist's jar path, and widen the daemon-locator pattern
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Failing after 2m28s
Part 1: asserts JAR and BUILD_JAR as the script sources them (no test assigns
them first), so reverting JAR to a path under target/ now fails the suite.

Part 2: check_jar_path_matches_plist reads the installed launchd plist's
ProgramArguments and refuses when its jar path does not resolve to $JAR,
mirroring check_log_path_matches_plist. Wired into report_supervisor_state's
launchd branch; unaffected when no plist is installed.

Part 3 (added to the ticket after the brief, by comment): PATTERN narrowed to
'fleetd.jar' so running_pid()/assert_single_daemon see a daemon regardless of
which build layout (run/ or target/) its jar sits under. The comm=java
allowlist still excludes a self-matching shell. Brought
.claude/skills/fleets-status/SKILL.md's pgrep pattern into agreement.

Verified: bash scripts/test-redeploy-fleetd.sh exits 0. Mutation both
directions for part 1 (JAR under target/ -> suite fails; JAR elsewhere ->
suite passes), a positive control for part 2 (neutralizing the mismatch
check makes the new test fail), and a positive control for part 3 (narrowing
PATTERN back to run/fleetd.jar makes the new target/-dir test fail). mvn -o
clean install: Tests run: 1929, Failures: 0, Errors: 0, Skipped: 0, BUILD
SUCCESS, 172 surefire report files.
2026-10-03 21:28:59 +02:00
13 changed files with 805 additions and 69 deletions
+1 -1
View File
@@ -59,7 +59,7 @@ as `matches HEAD`, `drift`, or `unknown`; do not turn an unclear timestamp into
Report the process identifier (PID) and uptime too:
```bash
PIDS="$(pgrep -f 'run/fleetd.jar' || true)"
PIDS="$(pgrep -f 'fleetd.jar' || true)"
if [ -z "$PIDS" ]; then
printf '%s\n' 'fleetd: not running'
else
+7
View File
@@ -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.
@@ -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"));
@@ -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");
+47 -6
View File
@@ -81,11 +81,12 @@ MODULE="$REPO/fleetd"
BUILD_JAR="$MODULE/target/fleetd.jar"
JAR="$MODULE/run/fleetd.jar"
OUT="$MODULE/fleetd.out"
# Matches BOTH the absolute form and the relative `java -jar run/fleetd.jar` a hand-start
# produces from inside fleetd/. Anchoring on the absolute path alone was a real bug: the daemon
# restarted correctly and the script still reported "no process appeared", because it launched with
# a relative path and then looked for an absolute one.
PATTERN='run/fleetd.jar'
# Matches a fleetd daemon's command line wherever its jar sits — absolute or relative, under
# run/, under target/, or anywhere else a build or a hand-start might point it. Detecting a
# daemon this script did not start, including one running from a jar outside $JAR's own
# directory, is this pattern's whole job; running_pid()'s `comm = java` allowlist below is what
# keeps that breadth from counting a shell that merely types the pattern as literal text.
PATTERN='fleetd.jar'
HEALTH='http://127.0.0.1:8765/healthz'
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start — fleetd #603: also the pid-
@@ -231,7 +232,7 @@ report_jar_state() {
# launched as `java -jar ...` — a native image, a renamed launcher — `running_pid()` silently
# returns nothing and `assert_single_daemon` stops noticing a second daemon at all. For a guard,
# that false-negative direction is the worse one to be wrong in. This is not a new assumption,
# though: `PATTERN='run/fleetd.jar'` two lines up already assumes the daemon is a jar, which
# though: `PATTERN='fleetd.jar'` two lines up already assumes the daemon is a jar, which
# is only ever run by `java`. If that launch method changes, `PATTERN` stops matching anything
# before this allowlist would ever get the chance to be wrong — the allowlist rides on the same
# assumption that is already load-bearing, it does not add a new one. Whoever changes the launch
@@ -627,6 +628,45 @@ check_log_path_matches_plist() {
ok "log path check: script and plist agree ($resolved_out)"
}
# Reads the launchd plist's ProgramArguments for the argument that follows "-jar", resolves it
# alongside $jar_path, and dies when the two differ. Call it only when the agent is loaded; it
# never touches launchd or the daemon itself.
check_jar_path_matches_plist() {
local jar_path="$1" plist_path="$2"
local plist_args plist_jar resolved_jar resolved_plist_jar
if ! plist_args="$(/usr/libexec/PlistBuddy -c 'Print :ProgramArguments' "$plist_path" 2>/dev/null)"; then
die "launchd agent is loaded but PlistBuddy could not read ProgramArguments from
$plist_path
— cannot verify which jar the supervised daemon launches. Fix the plist before
redeploying supervised."
fi
plist_jar="$(printf '%s\n' "$plist_args" | awk '
{ gsub(/^[ \t]+|[ \t]+$/, "") }
prev == "-jar" { print; exit }
{ prev = $0 }
')"
if [ -z "$plist_jar" ]; then
die "launchd agent is loaded but its ProgramArguments at
$plist_path
do not contain a '-jar <path>' pair — cannot verify which jar the supervised daemon
launches. Fix the plist before redeploying supervised."
fi
resolved_jar="$(cd "$(dirname "$jar_path")" 2>/dev/null && pwd -P)/$(basename "$jar_path")" || true
resolved_plist_jar="$(cd "$(dirname "$plist_jar")" 2>/dev/null && pwd -P)/$(basename "$plist_jar")" || true
if [ -z "$resolved_jar" ] || [ -z "$resolved_plist_jar" ] || [ "$resolved_jar" != "$resolved_plist_jar" ]; then
die "jar path mismatch — this script deploys to
$jar_path (resolved: ${resolved_jar:-<directory does not exist>})
but the loaded plist's ProgramArguments names
$plist_jar (resolved: ${resolved_plist_jar:-<directory does not exist>})
The swap renames the built jar into place, so the old path stops existing after a redeploy;
a launchd-initiated start from this plist (a reboot, or KeepAlive after a crash) would then
run java against a missing file. Reinstall the plist at
$plist_path
so its ProgramArguments names $jar_path before redeploying supervised."
fi
ok "jar path check: script and plist agree ($resolved_jar)"
}
# fleetd #552: the post-restart fresh-log capture, pulled out of the main flow so it is testable by
# sourcing (the same reason systemd_installed/systemd_loaded above guard their OWN mktemp inline
# instead of leaving it bare) even though its only caller sits below the SOURCED guard. By the time
@@ -960,6 +1000,7 @@ report_supervisor_state() {
SUPERVISED=1
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
check_jar_path_matches_plist "$JAR" "$LAUNCHD_PLIST"
;;
systemd)
SUPERVISED=1
+113
View File
@@ -492,6 +492,38 @@ test_running_pid_finds_a_real_java_named_second_process() {
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) whose own argv holds the pattern: before=[$before] after=[$after]"
}
# PATTERN matches a fleetd jar in either build layout, not only the run/ one: a process whose
# argv names a jar under target/ must be found too, the same way the run/ case above is.
test_running_pid_finds_a_real_java_named_process_from_target_dir() {
local before after standin_pid
before="$(running_pid)"
( exec -a java sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' ) &
standin_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$standin_pid" 2>/dev/null || true
wait "$standin_pid" 2>/dev/null || true
printf '%s\n' "$after" | grep -qxF "$standin_pid" \
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) naming a jar under target/: before=[$before] after=[$after]"
}
# Broadening PATTERN to match both build layouts must not also broaden it into matching a
# non-exec'ing shell that merely holds the target/ text as a literal argument, the same
# self-matching shape test_running_pid_excludes_self_matching_wrapper_shell above excludes for
# the run/ text.
test_running_pid_excludes_self_matching_wrapper_shell_naming_target_dir() {
local before after wrapper_pid
before="$(running_pid)"
sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' &
wrapper_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$wrapper_pid" 2>/dev/null || true
wait "$wrapper_pid" 2>/dev/null || true
[ "$after" = "$before" ] \
|| fail "running_pid() counted a self-matching wrapper shell (pid $wrapper_pid, holding 'target/fleetd.jar' as literal text in its own argv, not the daemon): before=[$before] after=[$after]"
}
# fleetd #593 CORRECTION 1, hole 2 — the round-1 filter denied known shell names (sh/bash/zsh/
# dash/ksh) and counted everything else. `ssh`, `perl`, `python3`, `ruby`, `tail` — anything not on
# that list, carrying the pattern in its own argv — was still counted right alongside the real
@@ -702,6 +734,19 @@ test_report_jar_state_both_absent_is_not_a_mismatch() {
fi
}
# Reads JAR and BUILD_JAR exactly as the script sources them, with nothing here assigning
# either first. JAR must resolve outside $MODULE/target/, and JAR must differ from BUILD_JAR:
# the daemon's live path and Maven's own build output are never the same file.
test_jar_and_build_jar_are_sourced_outside_target_and_differ() {
source "$ROOT/scripts/redeploy-fleetd.sh"
case "$JAR" in
"$MODULE"/target/*)
fail "\$JAR must not live under \$MODULE/target/ — got $JAR" ;;
esac
[ "$JAR" != "$BUILD_JAR" ] \
|| fail "\$JAR and \$BUILD_JAR must not be the same path — got $JAR"
}
# fleetd #493/#664 — never build into the path a running process holds. swap_staged_jar is
# exercised directly against real files on disk (not stubs), because the whole point is file
# behavior (does the content move, does the source disappear, does a failure leave both sides
@@ -1305,12 +1350,71 @@ test_run_drain_gate_declined_reply_refuses() {
source "$ROOT/scripts/redeploy-fleetd.sh"
}
# Writes a launchd-plist fixture naming jar_path as the ProgramArguments entry after "-jar", so
# check_jar_path_matches_plist has something real to read back.
write_launchd_plist_fixture() {
local path="$1" jar_path="$2"
cat > "$path" <<PLIST
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE plist PUBLIC "-//Apple//DTD PLIST 1.0//EN" "http://www.apple.com/DTDs/PropertyList-1.0.dtd">
<plist version="1.0">
<dict>
<key>Label</key>
<string>test.fixture</string>
<key>ProgramArguments</key>
<array>
<string>/usr/bin/java</string>
<string>-jar</string>
<string>$jar_path</string>
<string>fleetd.yaml</string>
</array>
</dict>
</plist>
PLIST
}
# Agreeing case: a plist whose ProgramArguments names the same jar, resolved, must proceed and
# say so, never die.
test_check_jar_path_matches_plist_agrees_ok() {
local dir jar plist output rc=0
dir="$TMP/jar-path-agree"; mkdir -p "$dir/run"
jar="$dir/run/fleetd.jar"
plist="$dir/agree.plist"
write_launchd_plist_fixture "$plist" "$jar"
output="$(check_jar_path_matches_plist "$jar" "$plist" 2>&1)" || rc=$?
[ "$rc" -eq 0 ] \
|| fail "check_jar_path_matches_plist must succeed when the plist names the same jar: $output"
printf '%s' "$output" | grep -qF 'jar path check' \
|| fail "check_jar_path_matches_plist did not print the agreement line: $output"
}
# The disagreeing case, and the positive control this check exists for: a plist naming a
# different jar must die, naming both paths.
test_check_jar_path_matches_plist_mismatch_dies() {
local dir jar other_jar plist output rc=0
dir="$TMP/jar-path-mismatch"; mkdir -p "$dir/run" "$dir/target"
jar="$dir/run/fleetd.jar"
other_jar="$dir/target/fleetd.jar"
plist="$dir/mismatch.plist"
write_launchd_plist_fixture "$plist" "$other_jar"
output="$(check_jar_path_matches_plist "$jar" "$plist" 2>&1)" || rc=$?
[ "$rc" -ne 0 ] \
|| fail "check_jar_path_matches_plist must die when the plist names a different jar"
printf '%s' "$output" | grep -qF "$jar" \
|| fail "die message does not name this script's jar path: $output"
printf '%s' "$output" | grep -qF "$other_jar" \
|| fail "die message does not name the plist's jar path: $output"
}
# fleetd #555 item 4 — the report-state dispatch on $SUPERVISOR_KIND. Inverting this used to report
# the wrong supervisor and, for the launchd arm specifically, skip check_log_path_matches_plist.
CHECK_LOG_PATH_CALLED=0
CHECK_JAR_PATH_CALLED=0
stub_check_log_path_recorder() {
CHECK_LOG_PATH_CALLED=0
CHECK_JAR_PATH_CALLED=0
check_log_path_matches_plist() { CHECK_LOG_PATH_CALLED=1; }
check_jar_path_matches_plist() { CHECK_JAR_PATH_CALLED=1; }
}
test_report_supervisor_state_launchd_sets_supervised_and_checks_log_path() {
@@ -1321,6 +1425,8 @@ test_report_supervisor_state_launchd_sets_supervised_and_checks_log_path() {
[ "$SUPERVISED" = 1 ] || fail "report_supervisor_state launchd must set SUPERVISED=1"
[ "$CHECK_LOG_PATH_CALLED" = 1 ] \
|| fail "report_supervisor_state launchd must call check_log_path_matches_plist"
[ "$CHECK_JAR_PATH_CALLED" = 1 ] \
|| fail "report_supervisor_state launchd must call check_jar_path_matches_plist"
source "$ROOT/scripts/redeploy-fleetd.sh"
}
@@ -1332,6 +1438,8 @@ test_report_supervisor_state_systemd_sets_supervised_without_log_path_check() {
[ "$SUPERVISED" = 1 ] || fail "report_supervisor_state systemd must set SUPERVISED=1"
[ "$CHECK_LOG_PATH_CALLED" = 0 ] \
|| fail "report_supervisor_state systemd must NOT call check_log_path_matches_plist"
[ "$CHECK_JAR_PATH_CALLED" = 0 ] \
|| fail "report_supervisor_state systemd must NOT call check_jar_path_matches_plist"
source "$ROOT/scripts/redeploy-fleetd.sh"
}
@@ -2189,6 +2297,8 @@ test_assert_single_daemon_accepts_one_pid
test_assert_single_daemon_rejects_two_pids
test_running_pid_excludes_self_matching_wrapper_shell
test_running_pid_finds_a_real_java_named_second_process
test_running_pid_finds_a_real_java_named_process_from_target_dir
test_running_pid_excludes_self_matching_wrapper_shell_naming_target_dir
test_running_pid_drops_a_pid_whose_comm_is_not_java
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup
test_running_pid_counts_a_pid_whose_comm_is_java
@@ -2200,6 +2310,7 @@ test_jar_id_reports_unhashable_when_no_hasher_on_path
test_report_jar_state_agrees_when_hashes_match
test_report_jar_state_warns_when_hashes_differ
test_report_jar_state_both_absent_is_not_a_mismatch
test_jar_and_build_jar_are_sourced_outside_target_and_differ
test_swap_staged_jar_moves_staged_onto_live
test_swap_staged_jar_dies_without_staged_file
test_swap_staged_jar_dies_when_mv_fails
@@ -2241,6 +2352,8 @@ test_drain_confirmed_false_on_anything_else
test_run_drain_gate_skips_prompt_when_not_required
test_run_drain_gate_confirmed_reply_does_not_refuse
test_run_drain_gate_declined_reply_refuses
test_check_jar_path_matches_plist_agrees_ok
test_check_jar_path_matches_plist_mismatch_dies
test_report_supervisor_state_launchd_sets_supervised_and_checks_log_path
test_report_supervisor_state_systemd_sets_supervised_without_log_path_check
test_report_supervisor_state_none_leaves_supervised_zero