fleetd #778: scope TASK_READ to a caller's own ticket; gate push-loop nudges on Authz #785

Open
agent wants to merge 2 commits from worker/778-12988a-4 into main
10 changed files with 377 additions and 83 deletions
@@ -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<MessageService> messagesRef = new AtomicReference<>();
AtomicReference<CallerResolver> callersRef = new AtomicReference<>();
BiPredicate<String, String> 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);
@@ -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<String> 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<String> 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.
*
* <p>Its default classifiers deny every collaborator and every observer, so a caller
* enforcing authorization must use the five-argument form instead.
* <p>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<String> 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<String> knownLeadOrCollaborator,
Predicate<String> 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<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
Predicate<String> knownObserverTarget,
Predicate<String> 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
@@ -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) {
@@ -563,11 +563,13 @@ public final class FleetMcp {
Map<String, Object> 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<String> 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<String> 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.
*
@@ -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
@@ -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<String, String> 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<String, String> 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<PendingTicket> 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());
}
@@ -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<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget,
Predicate<String> 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<String> 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;
@@ -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");
}
}
}
@@ -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);
@@ -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<String, String> 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<String, String> 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<String, String> ticketPollAuthorized) {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics,
ticketPollAuthorized);
}
private static AgentControl agentWithStatus(String status) {
return new AgentControl(new FakeHerdrClient(status));
}