From 3cbd84eb07afe888966842ecc4097c6c3728f0f4 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 5 Oct 2026 20:09:24 +0200 Subject: [PATCH 1/2] fleetd #778: let a collaborator or observer read a ticket it created TASK_READ was unconditionally closed to anyone but a primary, worker, or architect, so a non-worker peer that used fleet_send(wait:false) could never collect its own async reply. Add a ticket-ownership classifier to Authz.permits, following the same pattern as the existing SEND classifiers, and expose MessageService.isTicketOwnedBy so FleetMcp's fleet_poll handler and FleetApp's GET /tasks/{ticket} can build it. Both call sites also fix a second bug: they were passing a blank/null authorization target instead of the actual ticket id, so even a correct policy could never have been evaluated against it. Tests: AuthzTest gets a TASK_READ ownership matrix for an observer and a collaborator, each paired with a negative control (a ticket created by someone else). MessageServiceTest covers isTicketOwnedBy directly, same pairing. --- .../main/java/dev/ltms/fleet/auth/Authz.java | 76 ++++++++++++++----- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 40 +++++++++- .../dev/ltms/fleet/msg/MessageService.java | 11 +++ .../java/dev/ltms/fleet/rest/FleetApp.java | 31 +++++++- .../java/dev/ltms/fleet/auth/AuthzTest.java | 63 +++++++++++++-- .../ltms/fleet/msg/MessageServiceTest.java | 22 ++++++ 6 files changed, 209 insertions(+), 34 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java index 98ed3a64..982e37ab 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java @@ -36,7 +36,12 @@ public final class Authz { DRAIN, /** Read-only roster, profile, and identity observation: no ticket, task, or turn state. */ READ, - /** Poll a ticket, or read a session's status. */ + /** + * Poll a ticket, or read a session's status. Open unconditionally to a primary, a worker, + * and an architect; open to any other caller only for a ticket it created itself (see + * {@link #permits(Principal, Action, String, Predicate, Predicate, Predicate)}'s + * {@code callerOwnsTicket} parameter). + */ TASK_READ, /** * Read (never ack) this daemon's own held lead-to-lead coordination mail (fleetd #421). @@ -79,28 +84,52 @@ public final class Authz { */ public static final Predicate NO_KNOWN_OBSERVER_TARGET = target -> false; + /** + * The fail-closed classifier for {@code TASK_READ}'s ticket-ownership grant: answers no for + * every ticket, so the grant is refused unless a caller supplies a real one backed by the + * ticket's recorded creator (see {@code MessageService#isTicketOwnedBy}). + */ + public static final Predicate NO_CALLER_OWNS_TICKET = ticket -> false; + /** * Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's * or an observer's {@code SEND} is refused, as if no terminal were a configured lead, - * collaborator, or observer target — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} - * and {@link #NO_KNOWN_OBSERVER_TARGET} give explicitly. Every other action's result is - * identical to the five-argument form's, since none of them consult either classifier. + * collaborator, or observer target, and {@code TASK_READ} for a ticket that caller did not + * create itself is refused too — the same decisions {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}, + * {@link #NO_KNOWN_OBSERVER_TARGET}, and {@link #NO_CALLER_OWNS_TICKET} give explicitly. Every + * other action's result is identical to the six-argument form's, since none of them consult + * any of the three classifiers. * - *

Its default classifiers deny every collaborator and every observer, so a caller - * enforcing authorization must use the five-argument form instead. + *

Its default classifiers deny every collaborator, every observer, and every ticket a + * primary/worker/architect did not already have unconditionally, so a caller enforcing + * authorization must use the six-argument form instead. */ public static boolean permits(Principal caller, Action action, String targetSession) { - return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_KNOWN_OBSERVER_TARGET); + return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, + NO_KNOWN_OBSERVER_TARGET, NO_CALLER_OWNS_TICKET); } /** * As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's - * {@code SEND}. An observer's {@code SEND} still fails closed ({@link #NO_KNOWN_OBSERVER_TARGET}) — - * a caller enforcing both grants must use the five-argument form. + * {@code SEND}. An observer's {@code SEND} and {@code TASK_READ}'s ticket-ownership grant + * still fail closed — a caller enforcing all three grants must use the six-argument form. */ public static boolean permits(Principal caller, Action action, String targetSession, Predicate knownLeadOrCollaborator) { - return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_KNOWN_OBSERVER_TARGET); + return permits(caller, action, targetSession, knownLeadOrCollaborator, + NO_KNOWN_OBSERVER_TARGET, NO_CALLER_OWNS_TICKET); + } + + /** + * As {@link #permits(Principal, Action, String, Predicate)}, with a real classifier for an + * observer's {@code SEND}. {@code TASK_READ}'s ticket-ownership grant still fails closed — a + * caller enforcing all three grants must use the six-argument form. + */ + public static boolean permits(Principal caller, Action action, String targetSession, + Predicate knownLeadOrCollaborator, + Predicate knownObserverTarget) { + return permits(caller, action, targetSession, knownLeadOrCollaborator, knownObserverTarget, + NO_CALLER_OWNS_TICKET); } /** @@ -108,8 +137,9 @@ public final class Authz { * * @param targetSession the session id in the request path; only consulted for the * worker-scoped actions ({@code REPLY}, {@code ASK}), for a - * collaborator's {@code SEND}, and for an observer's - * {@code SEND}, ignored otherwise, may be {@code null} + * collaborator's {@code SEND}, for an observer's {@code SEND}, + * and (as the ticket id) for {@code TASK_READ}, ignored + * otherwise, may be {@code null} * @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator — * consulted only for a collaborator's {@code SEND}, to confine * it to another named peer and never a spawned member's @@ -118,10 +148,16 @@ public final class Authz { * {@link Role#OBSERVER} — consulted only for an observer's * {@code SEND}, to confine it to another observer pane and never * a lead, a collaborator, or a spawned member + * @param callerOwnsTicket whether {@code targetSession} (read as a ticket id) was + * created by the calling session — consulted only for + * {@code TASK_READ} by a caller that is none of primary, worker, + * or architect, so a collaborator or an observer may read a + * ticket it created itself and no other */ public static boolean permits(Principal caller, Action action, String targetSession, Predicate knownLeadOrCollaborator, - Predicate knownObserverTarget) { + Predicate knownObserverTarget, + Predicate callerOwnsTicket) { if (caller == null || caller.isAnonymous()) { return false; // authenticated as nothing ⇒ authorized for nothing } @@ -169,12 +205,14 @@ public final class Authz { case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect() || caller.isCollaborator() || caller.isObserver(); - // Ticket polling and session status, open to every role READ is open to except a - // collaborator or an observer. MessageService compares a ticket's creator to the - // caller on every read as well, so dropping this gate would not expose another - // session's reply — it would move the refusal later and widen what a caller that - // never orchestrates can probe. - case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect(); + // Ticket polling and session status, open unconditionally to every role READ is open + // to except a collaborator or an observer. A collaborator or an observer gets it too, + // but only for a ticket it created itself: fleet_send(wait:false) is open to both + // roles, so either can hold a ticket nobody else may read, and the ticket already + // records who created it. MessageService makes the same comparison again on every + // read, so this gate narrows WHO gets to ask, not what the answer is. + case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect() + || callerOwnsTicket.test(targetSession); // 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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index e568a7ab..2f274442 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -563,11 +563,13 @@ public final class FleetMcp { Map a = req.arguments(); String target = str(a, "target"); String coordId = str(a, "coordId"); + String ticket = str(a, "ticket"); // The action depends on the ARGUMENTS, not on the tool name -- see pollAction. - McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target); + Authz.Action action = toolAction("fleet_poll", a); + McpSchema.CallToolResult denied = deny(exchange, action, pollAuthzTarget(action, ticket, target), + t -> t != null && messages.isTicketOwnedBy(t, principal(exchange).ownerKey())); if (denied != null) return denied; - return poll(messages, leadChannel, str(a, "ticket"), target, coordId, - principal(exchange).ownerKey()); + return poll(messages, leadChannel, ticket, target, coordId, principal(exchange).ownerKey()); }; // CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox). // Acking removes a reply from the inbox, so it is a drain, not a read. @@ -732,6 +734,15 @@ public final class FleetMcp { return denyFor(principal(exchange), action, target); } + /** + * As {@link #deny(McpSyncServerExchange, Authz.Action, String)}, also threading the classifier + * {@code TASK_READ}'s ticket-ownership grant is checked against. + */ + private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action, + String target, Predicate callerOwnsTicket) { + return denyFor(principal(exchange), action, target, callerOwnsTicket); + } + /** * The policy half of {@link #deny}: everything except pulling the caller out of the MCP * exchange. Kept separate so the authorization decision — the actual control — is unit-testable @@ -745,13 +756,24 @@ public final class FleetMcp { * @return {@code null} when the call may proceed, or the error result to return when it may not */ McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) { + return denyFor(caller, action, target, Authz.NO_CALLER_OWNS_TICKET); + } + + /** + * As {@link #denyFor(Principal, Authz.Action, String)}, also threading the classifier + * {@code TASK_READ}'s ticket-ownership grant is checked against. + * + * @return {@code null} when the call may proceed, or the error result to return when it may not + */ + McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target, + Predicate callerOwnsTicket) { // The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future // tool that calls this directly must not be able to skip the gate by accident. if (!authorizationEnforced) { return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518) } if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(), - callers.sendableObserverTarget())) { + callers.sendableObserverTarget(), callerOwnsTicket)) { if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) { AuditLog.allowed(caller, action, target); // reads would drown the trail } @@ -1215,6 +1237,16 @@ public final class FleetMcp { return isBlank(target) ? Authz.Action.TASK_READ : Authz.Action.DRAIN; } + /** + * The identifier {@code fleet_poll}'s authorization check runs against: the ticket for + * {@link Authz.Action#TASK_READ} (so a ticket-ownership classifier has something to test), + * {@code target} for every other action, exactly as {@link #pollAction} already decided. + * Pulled into its own method for the same testability reason as {@link #pollAction} itself. + */ + static String pollAuthzTarget(Authz.Action action, String ticket, String target) { + return action == Authz.Action.TASK_READ ? ticket : target; + } + /** * Which authorization action a {@code fleet_send} call needs, decided by its arguments. * diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index c927cb31..6f970745 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -1400,6 +1400,17 @@ public final class MessageService { return poll(ticket, INTERNAL_NO_OWNER_CHECK); } + /** + * Whether {@code ticket} was created by the session identified by {@code callerOwner} — the + * same comparison {@link #poll(String, String)} itself makes, exposed so a caller can be + * granted read access to a ticket before it actually polls one. Returns {@code false} for an + * unknown or expired ticket. + */ + public boolean isTicketOwnedBy(String ticket, String callerOwner) { + Task task = tasks.get(ticket); + return task != null && Objects.equals(callerOwner, task.creatorOwner); + } + /** * Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket. * Refuses a {@code callerOwner} that differs from the owner that created the ticket (see diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index c0fa8edb..aa33e61b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -294,6 +294,17 @@ public final class FleetApp { return Authz.permits(caller, action, target, knownLeadOrCollaborator, knownObserverTarget); } + /** + * As {@link #permitsFor(Principal, Authz.Action, String, Predicate, Predicate)}, also + * threading the classifier {@code TASK_READ}'s ticket-ownership grant is checked against. + */ + static boolean permitsFor(Principal caller, Authz.Action action, String target, + Predicate knownLeadOrCollaborator, + Predicate knownObserverTarget, + Predicate callerOwnsTicket) { + return Authz.permits(caller, action, target, knownLeadOrCollaborator, knownObserverTarget, callerOwnsTicket); + } + /** * Gate a handler on the CB-505 authorization table. Returns {@code true} when the request may * proceed; otherwise writes the error response and returns {@code false}. @@ -303,11 +314,20 @@ public final class FleetApp { * yours" (a worker reaching for another worker's session, or for orchestration). */ private boolean allow(Context ctx, Authz.Action action, String target) { + return allow(ctx, action, target, Authz.NO_CALLER_OWNS_TICKET); + } + + /** + * As {@link #allow(Context, Authz.Action, String)}, also threading the classifier + * {@code TASK_READ}'s ticket-ownership grant is checked against. + */ + private boolean allow(Context ctx, Authz.Action action, String target, Predicate callerOwnsTicket) { if (auth == null) { return true; // legacy: authorization not enforced } Principal caller = ctx.attribute(CALLER); - if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget())) { + if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget(), + callerOwnsTicket)) { if (action != Authz.Action.READ && action != Authz.Action.METRICS && action != Authz.Action.TASK_READ) { AuditLog.allowed(caller, action, target); // reads would drown the trail @@ -929,11 +949,14 @@ public final class FleetApp { /** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */ private void taskStatus(Context ctx) { - if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) { + String ticket = ctx.pathParam("ticket"); + Principal caller = ctx.attribute(CALLER); + String ownerKey = caller == null ? null : caller.ownerKey(); + if (!allow(ctx, routeAction("GET /tasks/{ticket}"), ticket, + t -> t != null && messages.isTicketOwnedBy(t, ownerKey))) { return; } - Principal caller = ctx.attribute(CALLER); - MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey()); + MessageService.TaskView v = messages.poll(ticket, ownerKey); if (v == null) { ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)")); return; diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java index 9d5828a6..99e64f8d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java @@ -237,8 +237,11 @@ class AuthzTest { } /** - * Every action denied to a collaborator, asserted denied even when the classifier would - * accept any target — proving none of these is actually gated on the classifier at all. + * Every action denied to a collaborator, asserted denied even when the SEND classifier would + * accept any target — proving none of these is gated on that classifier. {@code TASK_READ} is + * included here with no ticket-ownership classifier supplied, so it falls back to the + * fail-closed default; its conditional grant for a ticket the collaborator actually created is + * tested separately below. */ @Test void aCollaboratorIsDeniedLifecycleCoordinationAndTicketPolling() { @@ -288,11 +291,12 @@ class AuthzTest { } /** - * Every action beyond READ/METRICS/REPLY/ASK/SEND, asserted denied for an observer — - * including {@code TASK_READ}, which is the entire point of this role: an unconfigured pane - * must not be able to poll a ticket or read another session's status. {@code SEND} is excluded - * here and given its own matrix below, since — unlike every action in this loop — its grant is - * conditional on the target, not fixed. + * Every action beyond READ/METRICS/REPLY/ASK/SEND, asserted denied for an observer with no + * ticket-ownership classifier supplied — so {@code TASK_READ} falls back to the fail-closed + * default here, an unconfigured pane reading a ticket it did not create. Its conditional grant + * for a ticket the observer actually created is tested separately below. {@code SEND} is + * excluded here and given its own matrix below, since — unlike every action in this loop — its + * grant is conditional on the target, not fixed. */ @Test void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() { @@ -363,4 +367,49 @@ class AuthzTest { "an observer's SEND must consult the observer-target classifier, never the " + "lead-or-collaborator one"); } + + // ── the TASK_READ ticket-ownership matrix ────────────────────────────────────────────────────── + + /** + * {@code TASK_READ} for an observer is conditional on ticket ownership, exactly like + * {@code SEND} is conditional on the target: flipping only the ownership classifier's answer + * flips only this outcome. Paired with a negative control — the same observer, the same ticket, + * a classifier that reports a different creator — so a test that forgot to vary the classifier + * at all could not pass by accident. + */ + @Test + void anObserverMayTaskReadATicketOnlyWhenItCreatedIt() { + assertTrue(Authz.permits(OBSERVER, TASK_READ, "task-1", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, + Authz.NO_KNOWN_OBSERVER_TARGET, ticket -> true), + "a ticket the observer created must be readable"); + assertFalse(Authz.permits(OBSERVER, TASK_READ, "task-1", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, + Authz.NO_KNOWN_OBSERVER_TARGET, ticket -> false), + "a ticket created by someone else must still be refused to the same observer"); + assertFalse(Authz.permits(OBSERVER, TASK_READ, "task-1"), + "the default classifier owns no ticket, so the three-argument form still refuses"); + } + + /** As above, for a collaborator — the grant is the same mechanism, not an observer special case. */ + @Test + void aCollaboratorMayTaskReadATicketOnlyWhenItCreatedIt() { + assertTrue(Authz.permits(COLLABORATOR, TASK_READ, "task-1", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, + Authz.NO_KNOWN_OBSERVER_TARGET, ticket -> true), + "a ticket the collaborator created must be readable"); + assertFalse(Authz.permits(COLLABORATOR, TASK_READ, "task-1", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, + Authz.NO_KNOWN_OBSERVER_TARGET, ticket -> false), + "a ticket created by someone else must still be refused to the same collaborator"); + } + + /** + * Control: {@code TASK_READ} for a primary, a worker, and an architect does not move on the + * ticket-ownership classifier — it is already unconditional for all three. + */ + @Test + void theTicketOwnershipClassifierNeverGatesTaskReadForPrimaryWorkerOrArchitect() { + for (Principal p : new Principal[]{PRIMARY, WORKER_A, ARCH_DESIGN}) { + assertTrue(Authz.permits(p, TASK_READ, "task-1", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, + Authz.NO_KNOWN_OBSERVER_TARGET, Authz.NO_CALLER_OWNS_TICKET), + p.describe() + " must read a ticket unconditionally, even when the classifier owns nothing"); + } + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index e4884ccf..fb130c14 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1034,6 +1034,28 @@ class MessageServiceTest { assertEquals("forbidden: this ticket was created by a different session", refused.detail()); } + /** + * {@link MessageService#isTicketOwnedBy} is the primitive {@code Authz}'s {@code TASK_READ} + * grant is built on, exposed so a caller can be granted read access before it ever polls. + * Paired with a negative control: the same ticket, a different creator's owner key. + */ + @Test + void isTicketOwnedByMatchesOnlyTheRealCreator() { + Principal observerA = Principal.observer("term_observer_a", 1); + Principal observerB = Principal.observer("term_observer_b", 2); + String ticket = messages.sendAsync(T, "long task", null, observerA); + + assertTrue(messages.isTicketOwnedBy(ticket, observerA.ownerKey()), + "the observer that created the ticket must be reported as its owner"); + assertFalse(messages.isTicketOwnedBy(ticket, observerB.ownerKey()), + "a different observer must not be reported as the owner"); + } + + @Test + void isTicketOwnedByIsFalseForAnUnknownTicket() { + assertFalse(messages.isTicketOwnedBy("no-such-ticket", Principal.observer("term_x", 1).ownerKey())); + } + @Test void architectOwnershipUsesTerminalRatherThanSlot() { Principal oldArchitect = Principal.architect("opus", "term_OLD", 1); -- 2.52.0 From 0dae97e6a9863329e0739bf7645ce9e1433eb62f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 5 Oct 2026 20:09:37 +0200 Subject: [PATCH 2/2] fleetd #778: never nudge a pane to run a call it may not perform ReplyPushLoop nudged an observer pane to run fleet_poll(ticket=...) five times, which Authz refused every time -- the pane was the right one to nudge, but the instruction was one it could never follow. Gate ticket nudges generally: ReplyPushLoop takes a (lead, ticket) -> boolean authorization predicate, consulted at the one place every other lookup in the class reads from (pendingTicketsFor), so a ticket the resolved lead may not poll is invisible to decide/injectNudge/bumpNudgeCounts alike, never named in a nudge, and never spends its own nudge budget. The real predicate, wired in FleetdAssembly, reconstructs the resolved lead's Principal from its bare terminal (CallerResolver.resolveTerminal, extracted from the existing connection-based resolve()) and runs it through the same Authz.permits/MessageService.isTicketOwnedBy check the real fleet_poll call would face. Reply and question nudges need no equivalent gate: Authz already restricts who can delegate to a worker in the first place (SPAWN and SEND-to-a-worker are primary/architect-only), so the lead resolved for those two sources is always one DRAIN/ANSWER already grants unconditionally. Also fixes the nudge-sent/failed log lines, which called any ticket creator a "lead" regardless of its actual role. Test: ReplyPushLoopTest proves a forbidden ticket produces no nudge, paired with a positive control proving the same setup nudges normally once the predicate authorizes the call. --- .../java/dev/ltms/fleet/FleetdAssembly.java | 23 +++- .../dev/ltms/fleet/auth/CallerResolver.java | 105 ++++++++++-------- .../dev/ltms/fleet/msg/ReplyPushLoop.java | 49 +++++++- .../dev/ltms/fleet/msg/ReplyPushLoopTest.java | 40 +++++++ 4 files changed, 168 insertions(+), 49 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java index 04748183..a643382b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java @@ -1,7 +1,9 @@ package dev.ltms.fleet; +import dev.ltms.fleet.auth.Authz; import dev.ltms.fleet.auth.CallerResolver; import dev.ltms.fleet.auth.MemberRegistry; +import dev.ltms.fleet.auth.Principal; import dev.ltms.fleet.config.ConfigRef; import dev.ltms.fleet.config.ConfigWatcher; import dev.ltms.fleet.config.FleetConfig; @@ -68,6 +70,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiPredicate; import java.util.function.Function; import java.util.function.Predicate; import java.util.function.Supplier; @@ -420,8 +423,24 @@ final class FleetdAssembly { // CB-502: the registry is built before the service and the push loop so send/reply outcomes // are counted at their single funnel rather than at each of the two caller-facing surfaces. Metrics metrics = FleetMetrics.create(sessions, replyInbox); + // The authorization check below needs `messages` and `callers`, both built further down — + // the same construction-order cycle `pushLoopRef` above breaks, broken the same way: a + // mutable holder set once each exists, read lazily from inside the predicate built here. + AtomicReference messagesRef = new AtomicReference<>(); + AtomicReference callersRef = new AtomicReference<>(); + BiPredicate ticketPollAuthorized = (lead, ticket) -> { + MessageService messages = messagesRef.get(); + CallerResolver cr = callersRef.get(); + if (messages == null || cr == null) { + return true; + } + Principal leadPrincipal = cr.resolveTerminal(lead); + return Authz.permits(leadPrincipal, Authz.Action.TASK_READ, ticket, + Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_KNOWN_OBSERVER_TARGET, + t -> messages.isTicketOwnedBy(t, leadPrincipal.ownerKey())); + }; var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox, - pushScheduler, maxReminders, backoffMs, metrics); + pushScheduler, maxReminders, backoffMs, metrics, ticketPollAuthorized); // fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda // above at the real push loop, now that it exists. pushLoopRef.set(pushLoop); @@ -459,6 +478,7 @@ final class FleetdAssembly { leadLauncher, config, leads); MessageService messages = new MessageService(router, injector, rendezvous, replyInbox, pushLoop, metrics); + messagesRef.set(messages); // Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller. // FOURTH of the recurring background loops to start (optional). @@ -517,6 +537,7 @@ final class FleetdAssembly { spawnedMemberRole, collaboratorTerminals); log.info("auth: loopback-trust (any loopback non-worker caller is the primary)"); } + callersRef.set(callers); FleetMcp.QuarantineSource quarantineSource = Fleetd.quarantineSource(config, quarantine, liveExhaustedPatterns, quarantineReasonByCredential); diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java b/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java index c8eab80d..7e6c06f8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/CallerResolver.java @@ -295,6 +295,66 @@ public final class CallerResolver { return slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT; } + /** + * Resolve the Principal a bare terminal would carry, with no connection to read a pid from — + * for a caller deciding whether a session it already knows the terminal of, but is not itself + * handling a request from, may still perform some action. Applies the same precedence {@link + * #resolve} applies to a connection's terminal, so a terminal this method calls a collaborator + * or an observer is exactly one {@link #resolve} would resolve the same way. + * + * @param terminal the herdr {@code terminal_id} to resolve, or {@code null} + */ + public Principal resolveTerminal(String terminal) { + return terminal == null ? Principal.anonymous() : resolveForTerminal(terminal, -1); + } + + private Principal resolveForTerminal(String terminal, long pid) { + MemberRole spawnedRole = spawnedMemberRole.apply(terminal); + if (spawnedRole != null) { + // A live spawned member occupies this pane. Its identity is its own, whatever a tab + // map says about the same terminal — checked before every tab map, consulting none + // of them, so a tab label can never override a roster entry for the same terminal. + if (spawnedRole == MemberRole.ARCHITECT) { + // The roster only answers THAT this pane is a live spawned member; config still + // decides WHAT that member's slot grants (fleetd #424). A slot revoked after the + // bind must still demote this session on its very next request, so the roster's + // own ARCHITECT role is confirmed against the live slot role, exactly as the + // architect-slot step below confirms a binding with no live member session. + String slot = architectTerminals.get().get(terminal); + if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) { + return Principal.architect(memberSlotNames.apply(slot), terminal, pid); + } + return Principal.worker(terminal, pid); + } + return Principal.worker(terminal, pid); + } + String lead = leadTerminals.get().get(terminal); + if (lead != null) { + // The config names this pane as a lead's own. The pane mapping is exactly as + // unforgeable as a worker's, so it outranks the token path — no credential needed. + // Checked before the architect registry so a pane named in BOTH is still the lead + // (CB-548 preserves every existing leader behaviour). + return Principal.leader(lead, terminal, pid); + } + String slot = architectTerminals.get().get(terminal); + if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) { + // The config/live binding names this pane as an architect slot's own. Same + // unforgeable pane mapping; the live binding, never a request argument, decides. + // Check the slot role too: this defence in depth prevents a bad lifecycle bind from + // escalating a dev, hunter or reviewer into an architect. This is the case the + // spawned-member step above does not catch: a binding with no live member session. + return Principal.architect(memberSlotNames.apply(slot), terminal, pid); + } + String collaborator = collaboratorTerminals.get().get(terminal); + if (collaborator != null) { + // An operator-labelled collaborator tab, confirmed live by the same scan that + // confirms a lead tab. Checked last among the tab maps so a pane also matching one + // of the above keeps that stronger role. + return Principal.collaborator(collaborator, terminal, pid); + } + return Principal.observer(terminal, pid); // unforgeable; never token-gated + } + /** * Resolve the caller of a request. * @@ -305,50 +365,7 @@ public final class CallerResolver { public Principal resolve(String remoteAddr, int remotePort, String authorizationHeader) { ConnectionIdentity.Caller c = identity.resolve(remoteAddr, remotePort); if (c.terminal() != null) { - MemberRole spawnedRole = spawnedMemberRole.apply(c.terminal()); - if (spawnedRole != null) { - // A live spawned member occupies this pane. Its identity is its own, whatever a tab - // map says about the same terminal — checked before every tab map, consulting none - // of them, so a tab label can never override a roster entry for the same terminal. - if (spawnedRole == MemberRole.ARCHITECT) { - // The roster only answers THAT this pane is a live spawned member; config still - // decides WHAT that member's slot grants (fleetd #424). A slot revoked after the - // bind must still demote this session on its very next request, so the roster's - // own ARCHITECT role is confirmed against the live slot role, exactly as the - // architect-slot step below confirms a binding with no live member session. - String slot = architectTerminals.get().get(c.terminal()); - if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) { - return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid()); - } - return Principal.worker(c.terminal(), c.pid()); - } - return Principal.worker(c.terminal(), c.pid()); - } - String lead = leadTerminals.get().get(c.terminal()); - if (lead != null) { - // The config names this pane as a lead's own. The pane mapping is exactly as - // unforgeable as a worker's, so it outranks the token path — no credential needed. - // Checked before the architect registry so a pane named in BOTH is still the lead - // (CB-548 preserves every existing leader behaviour). - return Principal.leader(lead, c.terminal(), c.pid()); - } - String slot = architectTerminals.get().get(c.terminal()); - if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) { - // The config/live binding names this pane as an architect slot's own. Same - // unforgeable pane mapping; the live binding, never a request argument, decides. - // Check the slot role too: this defence in depth prevents a bad lifecycle bind from - // escalating a dev, hunter or reviewer into an architect. This is the case the - // spawned-member step above does not catch: a binding with no live member session. - return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid()); - } - String collaborator = collaboratorTerminals.get().get(c.terminal()); - if (collaborator != null) { - // An operator-labelled collaborator tab, confirmed live by the same scan that - // confirms a lead tab. Checked last among the tab maps so a pane also matching one - // of the above keeps that stronger role. - return Principal.collaborator(collaborator, c.terminal(), c.pid()); - } - return Principal.observer(c.terminal(), c.pid()); // unforgeable; never token-gated + return resolveForTerminal(c.terminal(), c.pid()); } if (tokenMode) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java index a2778cf2..6cc5f4f2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -20,6 +20,7 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.function.BiPredicate; import java.util.stream.Collectors; /** @@ -93,6 +94,13 @@ public final class ReplyPushLoop { private final int maxReminders; private final long backoffMs; private final Metrics metrics; // CB-512: nullable — no registry in unit tests + /** + * Whether the resolved lead may still poll a given ticket, consulted before a nudge ever names + * it. {@code fleet_send(wait:false)} is reachable by a collaborator or an observer, neither of + * which this loop's own identity ladder checks, so left unguarded a nudge could instruct a pane + * to run a call its authorization would refuse. + */ + private final BiPredicate ticketPollAuthorized; /** * Worker targets with a reply queued, keyed by target. Each entry carries its own nudge @@ -130,6 +138,20 @@ public final class ReplyPushLoop { public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, ScheduledExecutorService scheduler, int maxReminders, long backoffMs, Metrics metrics) { + this(primaryRegistry, agents, inbox, scheduler, maxReminders, backoffMs, metrics, + (lead, ticket) -> true); + } + + /** + * As above, with {@code ticketPollAuthorized} wired to the real authorization check. The two + * constructors above pass a predicate that always answers yes, matching this loop's behaviour + * before the check existed — the real production wiring supplies the one backed by the + * authorization table. + */ + public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, + ScheduledExecutorService scheduler, + int maxReminders, long backoffMs, Metrics metrics, + BiPredicate ticketPollAuthorized) { this.primaryRegistry = primaryRegistry; this.agents = agents; this.promptBox = new PromptBox(agents); @@ -138,6 +160,7 @@ public final class ReplyPushLoop { this.maxReminders = maxReminders; this.backoffMs = backoffMs; this.metrics = metrics; + this.ticketPollAuthorized = ticketPollAuthorized == null ? (lead, ticket) -> true : ticketPollAuthorized; } /** @@ -194,9 +217,27 @@ public final class ReplyPushLoop { private record PendingTicket(String ticket, String lead, boolean failed, int nudgeCount) { } - /** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */ + /** + * Tickets still pending for {@code lead}, snapshotted fresh for one tick — excluding any ticket + * {@link #ticketPollAuthorized} says {@code lead} may not poll. Filtered here, at the one source + * every other lookup in this class reads from, so an unauthorized ticket is invisible to + * {@link #decide}, {@link #injectNudge}, and {@link #bumpNudgeCounts} alike: it is never named + * in a nudge, never spends its nudge budget, and never keeps a schedule alive on its own. + */ private List pendingTicketsFor(String lead) { - return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList(); + return pendingTickets.values().stream() + .filter(t -> lead.equals(t.lead())) + .filter(t -> authorizedToPoll(lead, t.ticket())) + .toList(); + } + + private boolean authorizedToPoll(String lead, String ticket) { + if (ticketPollAuthorized.test(lead, ticket)) { + return true; + } + log.debug("push: ticket {} is pending for {} but it may not poll it, suppressing the nudge", + ticket, lead); + return false; } /** Ticket ids still pending for {@code lead} — a plain snapshot for race comparison. */ @@ -779,7 +820,7 @@ public final class ReplyPushLoop { String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets); try { agents.send(lead, nudge); - log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; " + log.debug("push: nudge sent to {} (reply {}/{}, ticket {}/{}, question {}/{}; " + "{} reply target(s), {} ticket(s), {} question(s))", lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, questionReminderCount + 1, maxReminders, @@ -796,7 +837,7 @@ public final class ReplyPushLoop { } } } catch (RuntimeException e) { - log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}", + log.warn("push: failed to nudge {} (reply {}/{}, ticket {}/{}, question {}/{}): {}", lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, questionReminderCount + 1, maxReminders, e.toString()); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java index 809d871b..48ccbbfc 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java @@ -23,6 +23,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BiPredicate; import static org.junit.jupiter.api.Assertions.*; @@ -483,6 +484,39 @@ class ReplyPushLoopTest { assertEquals(0, rec.sendCount(), "an already-collected ticket must never be nudged"); } + /** + * A nudge must never instruct a target to run a call its own authorization would refuse — the + * push loop is a defensive check, not the source of truth for whether the call is allowed. + */ + @Test + void aTicketIsNotNudgedWhenTheLeadMayNotPollIt() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + BiPredicate neverAuthorized = (lead, ticket) -> false; + + loop(1, 50, null, neverAuthorized).onTicketTerminal("task-1", WORKER, false); + + Thread.sleep(300); // let the scheduled tick run + assertEquals(0, rec.sendCount(), "a ticket the lead may not poll must never be named in a nudge"); + } + + /** + * Positive control for the test above: the same setup, but a predicate that authorizes the + * call, still sends the nudge — proving the suppression above is the predicate's own doing, + * not some other change that silently stops every ticket nudge. + */ + @Test + void aTicketIsNudgedWhenTheLeadMayPollIt() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + BiPredicate alwaysAuthorized = (lead, ticket) -> true; + + loop(1, 50, null, alwaysAuthorized).onTicketTerminal("task-1", WORKER, false); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "an authorized ticket must still be nudged"); + assertEquals(1, rec.sendCount()); + } + @Test void ticketNudgesSendUpToCapThenStop() throws Exception { int cap = 2; @@ -1164,6 +1198,12 @@ class ReplyPushLoopTest { return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics); } + private ReplyPushLoop loop(int maxReminders, long backoffMs, Metrics metrics, + BiPredicate ticketPollAuthorized) { + return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics, + ticketPollAuthorized); + } + private static AgentControl agentWithStatus(String status) { return new AgentControl(new FakeHerdrClient(status)); } -- 2.52.0