fleetd #778: scope TASK_READ to a caller's own ticket; gate push-loop nudges on Authz #785
@@ -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));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user