Compare commits

..

2 Commits

Author SHA1 Message Date
Dai Ha 0dae97e6a9 fleetd #778: never nudge a pane to run a call it may not perform
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 1m55s
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.
2026-10-05 20:09:37 +02:00
Dai Ha 3cbd84eb07 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.
2026-10-05 20:09:24 +02:00
16 changed files with 421 additions and 434 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) {
@@ -137,18 +137,6 @@ public final class AgentControl {
return result.path("read").path("text").asText("");
}
/**
* Read an agent's terminal with its ANSI styling kept, instead of the stripped text {@link
* #read} returns. Needed when a caller must tell apart text the pane draws dim (a placeholder
* hint) from text drawn plain (the operator's own typing).
*
* @param source one of {@code visible|recent|recent_unwrapped|detection}
*/
public String readWithStyling(String target, String source) {
JsonNode result = agentCall("agent.read", target, Map.of("source", source, "strip_ansi", false));
return result.path("read").path("text").asText("");
}
/** Current agent record (status, session UUID, pane). */
public Agent get(String target) {
return Agent.from(agentCall("agent.get", target, Map.of()).get("agent"));
@@ -3,12 +3,9 @@ package dev.ltms.fleet.herdr;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
/**
* Whether an agent pane's input box is clear for a delivery.
@@ -32,22 +29,20 @@ public final class PromptBox {
private static final Logger log = LoggerFactory.getLogger(PromptBox.class);
/**
* herdr {@code agent.read} source, read with its ANSI styling kept. The input box is always
* drawn here, carrying transcript scrollback above it, which is why only the last box line is
* the live one. Styling must survive the read because the pane draws a placeholder hint — the
* pane's own last submitted prompt — in the same spot as unsubmitted text, dimmed; only the
* escape codes tell the two apart.
* herdr {@code agent.read} source. {@code detection} is the region herdr itself uses for status
* detection, so the input box is always drawn in it. It is not a short tail: it carries transcript
* scrollback above the box, including earlier prompts the operator has already submitted, which is
* why only the last box line on it is the live one.
*/
static final String PROBE_SOURCE = "visible";
static final String PROBE_SOURCE = "detection";
/** Consecutive holds for one target before one warning is logged. */
static final int HOLD_WARN_STREAK = 20;
/**
* Input box markers, each matched only as a line's first characters once any leading ANSI
* escape codes are skipped: the caret the current TUI draws, and the bordered box an older one
* drew. A marker further along a line is transcript text, such as a caret inside something the
* operator quoted.
* Input box markers, each matched only as a line's first characters: the caret the current TUI
* draws, and the bordered box an older one drew. A marker further along a line is transcript text,
* such as a caret inside something the operator quoted.
*/
private static final List<String> BOX_MARKERS = List.of("❯", "│ >");
@@ -57,15 +52,6 @@ public final class PromptBox {
/** Block glyphs a terminal capture can leave in an otherwise empty box for the cursor cell. */
private static final String CURSOR_GLYPHS = "█▉▊▋▌▍▎▏";
/** An SGR escape sequence, e.g. {@code ESC[2m} (faint) or {@code ESC[0m} (reset). */
private static final Pattern SGR = Pattern.compile("\u001b\\[([0-9;]*)m");
/** The SGR code that dims text — herdr's placeholder hint is drawn inside a span of this. */
private static final String FAINT_CODE = "2";
/** The SGR code (or an empty code list) that clears every attribute, including faint. */
private static final List<String> RESET_CODES = List.of("", "0");
/** What a box holds: nothing, unsubmitted characters, or a pane this cannot read as a box. */
public enum State { EMPTY, DRAFT, UNREADABLE }
@@ -108,7 +94,7 @@ public final class PromptBox {
private Reading inspect(String target) {
String pane;
try {
pane = agents.readWithStyling(target, PROBE_SOURCE);
pane = agents.read(target, PROBE_SOURCE);
} catch (RuntimeException e) {
log.debug("prompt box read for {} failed, holding delivery: {}", target, e.getMessage());
return new Reading(State.UNREADABLE, 0);
@@ -128,9 +114,10 @@ public final class PromptBox {
* an earlier turn's marker survives; treating that as a live turn would make {@link State#EMPTY}
* unreachable and hold every delivery forever.
*
* <p>Whitespace, a trailing box border and a cursor block count as nothing. A placeholder hint —
* text the pane draws faint, in the same spot as unsubmitted text — also counts as nothing: only
* a character drawn outside a faint span is the operator's own typing.
* <p>Whitespace, a trailing box border and a cursor block count as nothing. Any other character
* counts as the operator's unsubmitted text, including a placeholder hint a future TUI might draw
* there; that direction holds a delivery it could have sent, which {@link #HOLD_WARN_STREAK} makes
* visible.
*/
static Reading classify(String pane) {
if (pane == null || pane.isBlank()) return new Reading(State.UNREADABLE, 0);
@@ -155,77 +142,28 @@ public final class PromptBox {
return found;
}
/**
* Length of the marker prefix — any leading SGR escape codes, then a box marker — this line
* starts with, or {@code 0} if it starts with neither. The colour drawn on the caret itself
* (e.g. an empty box's grey) sits before the marker glyph, so it must be skipped before the
* marker can match.
*/
/** Length of the box marker this line starts with, or {@code 0} if it starts with none. */
private static int markerLength(String line) {
int skip = leadingEscapeLength(line);
for (String marker : BOX_MARKERS) {
if (line.startsWith(marker, skip)) return skip + marker.length();
if (line.startsWith(marker)) return marker.length();
}
return 0;
}
/** Length of the run of SGR escape codes starting at the beginning of {@code line}. */
private static int leadingEscapeLength(String line) {
Matcher m = SGR.matcher(line);
int pos = 0;
while (m.find(pos) && m.start() == pos) pos = m.end();
return pos;
}
private static String firstLine(String text) {
int newline = text.indexOf('\n');
return newline < 0 ? text : text.substring(0, newline);
}
/** One rendered character of a box line, and whether it was drawn inside a faint (dim) span. */
private record Glyph(char c, boolean faint) {
}
/**
* The text the box holds: its own line after the marker, with border, padding, cursor and any
* faint (placeholder-hint) text left out — only a character drawn outside a faint span is the
* operator's own typing.
*/
/** The text the box holds: its own line after the marker, stripped of border, padding and cursor. */
private static String boxContent(String boxLine) {
List<Glyph> glyphs = renderedGlyphs(boxLine.substring(markerLength(boxLine)));
int end = glyphs.size();
while (end > 0 && isBoxPadding(glyphs.get(end - 1).c())) end--;
if (end > 0 && glyphs.get(end - 1).c() == '│') end--;
String line = boxLine.substring(markerLength(boxLine)).stripTrailing();
if (line.endsWith("│")) line = line.substring(0, line.length() - 1);
StringBuilder content = new StringBuilder();
for (int i = 0; i < end; i++) {
Glyph glyph = glyphs.get(i);
if (glyph.faint() || isBoxPadding(glyph.c())) continue;
content.append(glyph.c());
for (char c : line.toCharArray()) {
if (Character.isWhitespace(c) || CURSOR_GLYPHS.indexOf(c) >= 0) continue;
content.append(c);
}
return content.toString();
}
/** Decode {@code text} into its rendered characters, tracking the faint (SGR 2) span each sits in. */
private static List<Glyph> renderedGlyphs(String text) {
List<Glyph> glyphs = new ArrayList<>();
Matcher m = SGR.matcher(text);
boolean faint = false;
int i = 0;
while (i < text.length()) {
if (m.find(i) && m.start() == i) {
String codes = m.group(1);
if (RESET_CODES.contains(codes)) faint = false;
else if (FAINT_CODE.equals(codes)) faint = true;
i = m.end();
continue;
}
glyphs.add(new Glyph(text.charAt(i), faint));
i++;
}
return glyphs;
}
private static boolean isBoxPadding(char c) {
return Character.isWhitespace(c) || Character.isSpaceChar(c) || CURSOR_GLYPHS.indexOf(c) >= 0;
}
}
@@ -94,17 +94,6 @@ public final class Injector {
*/
private static final int READINESS_GRACE_POLLS = 240;
/**
* How many consecutive polls a message may sit queued with no delivery attempt at all before
* it is failed and the queue cleared — covers every reason the head of the queue is never
* reached, including a target that stays busy ({@code working}) or unclassifiable
* ({@code unknown}) for the whole window. Set well above an ordinary turn so a worker
* genuinely mid-task is never cut off, and below a caller's own overall timeout so a target
* that never frees up fails with this specific reason instead of riding out that longer wait
* silently.
*/
private static final int QUEUE_WAIT_GRACE_POLLS = 4800;
/**
* The single source for the injector poll cadence — how often the {@link StatusPoller} drives
* {@link #onStatus} at. {@code Fleetd} passes this to every {@link StatusPoller} it constructs,
@@ -307,16 +296,13 @@ public final class Injector {
final String text;
final TurnToken token;
final CompletableFuture<Void> delivered;
final long enqueuedAtMillis;
volatile State state = State.QUEUED; // written under the owning Target monitor
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered,
long enqueuedAtMillis) {
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered) {
this.target = target;
this.text = text;
this.token = token;
this.delivered = delivered;
this.enqueuedAtMillis = enqueuedAtMillis;
}
String text() {
@@ -343,7 +329,6 @@ public final class Injector {
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
long notReadySinceMillis; // wall-clock time of the FIRST non-ready sample in the current notReadySincePoll streak (fleetd #501); reset alongside it
int queueWaitSincePoll; // consecutive polls the queue has held an undelivered message with no attempt made
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
boolean awaitingPostTurnPickup;
boolean postTurnObserved;
@@ -364,7 +349,7 @@ public final class Injector {
*/
public Delivery enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(target, text, token, delivered, nowMillis.getAsLong());
Pending p = new Pending(target, text, token, delivered);
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
@@ -430,7 +415,6 @@ public final class Injector {
boolean resubmit = false;
boolean startPostTurn = false;
List<Pending> notReady = null; // queued messages failed because the worker never became ready
List<Pending> queueStalled = null; // queued messages failed because the queue never drained
synchronized (t) {
if (status == AgentStatus.WORKING) {
if (t.awaitingPostTurnPickup) {
@@ -622,28 +606,6 @@ public final class Injector {
}
}
// A message still queued and never attempted this poll is bounded on its own, whatever
// the reason the head of the queue was never reached — a target stuck WORKING or
// UNKNOWN for the whole window hits this even though neither branch above ever looks at
// the queue. Any poll that did attempt the head (`sent != null`, success or failure
// alike) counts as progress and resets the streak, even if messages remain behind it.
if (!t.queue.isEmpty() && sent == null) {
if (++t.queueWaitSincePoll >= QUEUE_WAIT_GRACE_POLLS) {
queueStalled = new ArrayList<>(t.queue);
for (Pending pending : queueStalled) {
pending.state = Pending.State.NOT_DELIVERED;
}
log.warn("queue for {} never drained after {} polls (limit={} polls/{}s): "
+ "failing {} queued message(s) that were never attempted",
target, t.queueWaitSincePoll, QUEUE_WAIT_GRACE_POLLS,
QUEUE_WAIT_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, queueStalled.size());
t.queue.clear();
t.queueWaitSincePoll = 0;
}
} else {
t.queueWaitSincePoll = 0;
}
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
// completion awaited), so the map cannot grow without bound across short-lived workers.
if (isQuiescent(t)) {
@@ -714,19 +676,6 @@ public final class Injector {
forget.accept(target);
turnListener.onTurnFailed(target);
}
if (queueStalled != null) {
// The target is not gone — it may still be genuinely busy — so this does not call
// forget.accept: that would clear presence/readiness state for a worker that is
// simply taking a long turn. It still resolves the awaiting send's own waiter via
// onTurnFailed (mirroring notReady above), so a caller learns this specific message
// never reached the pane instead of riding out its own much longer timeout.
RuntimeException cause = new IllegalStateException(
target + " never freed up to receive this message within the queue wait grace");
for (Pending p : queueStalled) {
p.delivered().completeExceptionally(cause);
}
turnListener.onTurnFailed(target, cause.getMessage());
}
if (turnCompleted) {
if (startPostTurn) {
// fleetd #553: the listener call is wrapped so `t.postTurnPending` (set true inside
@@ -863,23 +812,6 @@ public final class Injector {
.collect(Collectors.toSet());
}
/**
* How long the oldest still-queued, never-attempted message for {@code target} has been
* waiting, or {@code null} when nothing is queued (including when the head has already been
* attempted or delivered). A caller uses this to tell a message that genuinely never reached
* the pane apart from one that was delivered and is now simply being worked on.
*/
public Long queuedWaitMillis(String target) {
Target t = targets.get(target);
if (t == null) {
return null;
}
synchronized (t) {
Pending head = t.queue.peek();
return head != null ? nowMillis.getAsLong() - head.enqueuedAtMillis : null;
}
}
/**
* Forget a target whose worker is gone, failing every still-queued message so awaiting callers
* unblock instead of hanging forever. If a message had already been <em>delivered</em> but its
@@ -8,7 +8,6 @@ import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.inject.MemberPresence;
@@ -564,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.
@@ -733,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
@@ -746,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
}
@@ -1107,29 +1128,7 @@ public final class FleetMcp {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted, creator);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket
+ notInjectableWarning(messages, sessionId));
}
/**
* A synchronous status check at accept time: when {@code sessionId} is not currently
* idle/blocked/done, this message is queued rather than reaching the pane right away. Empty
* when the target is injectable or its status could not be read — a best-effort warning, not
* a reason to withhold the accept receipt.
*/
private static String notInjectableWarning(MessageService messages, String sessionId) {
AgentStatus status;
try {
status = messages.status(sessionId);
} catch (RuntimeException e) {
return "";
}
if (status.injectable()) {
return "";
}
return "\n\nWarning: " + sessionId + " is currently " + status.name().toLowerCase()
+ ", not idle/blocked/done — this message is queued, not yet delivered, and will "
+ "wait until the target frees up. Poll fleet_poll to see when it lands.";
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
@@ -1238,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
@@ -1426,7 +1437,7 @@ public final class MessageService {
return new TaskView(ticket, Phase.ASKING, question.text(), null,
"worker is waiting for your answer", question.turnId());
}
return new TaskView(ticket, Phase.PENDING, null, null, pendingDetail(task.target), null);
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
}
// CB-588: the ticket is terminal and being handed to the caller right here — tell the push
// loop it is collected so a later tick's nudge never names a ticket the lead already has.
@@ -1499,21 +1510,6 @@ public final class MessageService {
return task != null && task.completedNanos != null;
}
/**
* Detail text for a {@link Phase#PENDING} poll of a plain (non-asking) delegation.
* Distinguishes a message still sitting in the injector's queue, never delivered, from one
* that already reached the pane and is simply being worked on — so a caller cannot read
* "worker working" as "received" when it was not.
*/
private String pendingDetail(String target) {
Long queuedMillis = injector.queuedWaitMillis(target);
if (queuedMillis != null) {
return "queued, not yet delivered (target is " + liveStatus(target) + "; queued "
+ (queuedMillis / 1000) + "s)";
}
return "worker " + liveStatus(target);
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
private String liveStatus(String target) {
try {
@@ -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");
}
}
}
@@ -254,11 +254,10 @@ public final class FakeHerdr implements HerdrClient {
}
/**
* Override the text {@code agent.read} returns for the {@code detection} and {@code visible}
* sources only — the prompt/footer tail herdr uses for status detection and input-box probing, a
* different region from the transcript the other sources carry. Needed by a test whose subject
* reads the input box, since one {@link #readText} cannot be both a worker's transcript and a
* lead's empty prompt.
* Override the text {@code agent.read} returns for the {@code detection} source only — the
* prompt/footer tail herdr uses for status detection, a different region from the transcript the
* other sources carry. Needed by a test whose subject reads the input box, since one
* {@link #readText} cannot be both a worker's transcript and a lead's empty prompt.
*/
public FakeHerdr detectionText(String text) {
this.detectionText = text;
@@ -414,8 +413,8 @@ public final class FakeHerdr implements HerdrClient {
}
case "agent.read" -> {
Object source = params instanceof Map<?, ?> m ? m.get("source") : null;
boolean probeSource = "detection".equals(source) || "visible".equals(source);
String text = probeSource && detectionText != null ? detectionText : readText;
String text = "detection".equals(source) && detectionText != null
? detectionText : readText;
yield mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", text))));
}
@@ -12,18 +12,6 @@ class PromptBoxTest {
private static final String EMPTY = FakeHerdr.IDLE_PROMPT_CARET;
private static final String DRAFTED = FakeHerdr.DRAFTED_PROMPT_CARET;
/** An empty box: the pane draws a placeholder hint (its last submitted prompt) at the caret, faint. */
private static final String HINT_1 = "❯\u00a0\u001b[0m\u001b[2mstart md2trilium work for #2\u001b[0m";
/** Same shape as {@link #HINT_1}, a different placeholder hint. */
private static final String HINT_2 = "❯\u00a0\u001b[0m\u001b[2myes, push them\u001b[0m";
/** An empty box with no hint: the caret itself is drawn grey, with escape codes before the marker. */
private static final String EMPTY_GREY_CARET = "\u001b[0m\u001b[38;2;153;153;153m❯\u00a0\u001b[0m";
/** An empty box with no styling at all. */
private static final String EMPTY_PLAIN = "❯\u00a0";
// --- pure classification -------------------------------------------------
@Test
@@ -69,7 +57,7 @@ class PromptBoxTest {
void theLastBoxLineOnThePaneIsTheLiveOne() {
assertEquals(PromptBox.State.DRAFT,
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯ typing now").state(),
"the probed region carries scrollback, so earlier prompts sit above the live box");
"the detection region carries scrollback, so earlier prompts sit above the live box");
assertEquals(PromptBox.State.EMPTY,
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯").state());
}
@@ -102,54 +90,18 @@ class PromptBoxTest {
"that marker survives in scrollback, and holding on it would hold every delivery forever");
}
@Test
void aPlaceholderHintReadsAsAnEmptyBox() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(HINT_1),
"the hint is the pane's own last prompt, drawn faint — it is not the operator's typing");
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(HINT_2));
}
@Test
void anEmptyBoxWithAGreyCaretReadsAsEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY_GREY_CARET),
"the caret's own colour sits before the marker and must not stop the marker matching");
}
@Test
void anEmptyBoxWithNoStylingAtAllReadsAsEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY_PLAIN));
}
@Test
void typedTextWithNoStylingIsADraftAndCountsItsCharacters() {
PromptBox.Reading reading = PromptBox.classify("❯\u00a0deploy the thing");
assertEquals(PromptBox.State.DRAFT, reading.state());
assertEquals("deploythething".length(), reading.characters());
}
@Test
void typedTextAfterAFaintHintCountsOnlyTheTextOutsideTheFaintSpan() {
PromptBox.Reading reading =
PromptBox.classify("❯\u00a0\u001b[0m\u001b[2mhint\u001b[0m and typed");
assertEquals(PromptBox.State.DRAFT, reading.state());
assertEquals("andtyped".length(), reading.characters(),
"the faint hint is excluded; only \"and typed\" was drawn plain");
}
// --- the gate ------------------------------------------------------------
@Test
void anEmptyBoxClearsTheGateAndReadsTheVisibleRegionWithStylingKept() {
void anEmptyBoxClearsTheGateAndReadsTheDetectionRegion() {
FakeHerdr herdr = new FakeHerdr().detectionText(EMPTY);
assertTrue(new PromptBox(new AgentControl(herdr)).clearToSubmit("term_a"));
@SuppressWarnings("unchecked")
var params = (java.util.Map<String, Object>) herdr.lastCall("agent.read").params();
assertEquals("visible", params.get("source"),
"the input box is drawn in the visible region, not in transcript scrollback");
assertEquals(false, params.get("strip_ansi"),
"styling must survive the read, or a faint placeholder hint reads as plain typed text");
assertEquals("detection", params.get("source"),
"the input box is drawn in the detection region, not in transcript scrollback");
}
@Test
@@ -559,21 +559,20 @@ class InjectorTest {
void readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants() {
// fleetd #501, defect 2: the old line computed "({}s)" as READINESS_GRACE_POLLS *
// POLL_INTERVAL_MILLIS / 1000 — arithmetic on two constants, never a measurement, and wrong
// in the direction that says everything ran on schedule. This stub clock returns three FIXED
// values (an unused enqueue-time stamp, 1_000ms at the first non-ready sample, 318_412ms at
// the poll that trips the grace) whose last two differ — 317_412ms — which does NOT equal
// 240 * POLL_INTERVAL_MILLIS (=60_000ms). Asserting on that literal, non-derived number is
// what makes this test able to fail if the production code goes back to printing the
// constant-arithmetic value instead of the injected clock's measurement.
long[] readings = {0L, 1_000L, 318_412L};
// in the direction that says everything ran on schedule. This stub clock returns two FIXED
// values (1_000ms at the first non-ready sample, 318_412ms at the poll that trips the grace)
// whose difference — 317_412ms — does NOT equal 240 * POLL_INTERVAL_MILLIS (=60_000ms).
// Asserting on that literal, non-derived number is what makes this test able to fail if the
// production code goes back to printing the constant-arithmetic value instead of the
// injected clock's measurement.
long[] readings = {1_000L, 318_412L};
AtomicInteger call = new AtomicInteger(0);
LongSupplier stubClock = () -> {
int i = call.getAndIncrement();
if (i >= readings.length) {
throw new AssertionError("nowMillis read more times than this fixture expects (" + i
+ "); the readiness-not-ready branch should read the clock exactly three times —"
+ " once to stamp the queued message's own enqueue time, once to stamp the "
+ "first non-ready sample, once at grace expiry");
+ "); the readiness-not-ready branch should read the clock exactly twice — "
+ "once to stamp the first non-ready sample, once at grace expiry");
}
return readings[i];
};
@@ -615,46 +614,6 @@ class InjectorTest {
assertEquals(List.of("task"), sent(), "a worker that connects within the grace is delivered to");
}
// ~20 minutes of a continuously busy target at the 250ms prod poll interval; enough to trip
// the queue-wait grace.
private static final int QUEUE_WAIT_SAMPLES = 4800;
@Test
void failsAQueuedMessageWhoseTargetNeverFreesUp() {
// A target that stays WORKING the whole time never reaches the branch that looks at the
// queue at all, so nothing else bounds this. The message must fail rather than wait
// forever, and the caller's future unblocks through the same turn-failure path a readiness
// timeout uses — but presence must not be touched, since this target is merely busy, not
// gone.
Captor cap = new Captor();
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> true, forgotten::add);
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
for (int i = 0; i < QUEUE_WAIT_SAMPLES; i++) inj.onStatus(T, AgentStatus.WORKING);
assertEquals(List.of(), sent(), "a target that never frees up is never delivered to");
assertTrue(f.isCompletedExceptionally(), "the caller's future fails instead of hanging forever");
assertEquals(List.of(T), cap.failed, "the awaiting send resolves through the turn-failure path");
assertEquals(List.of(), cap.completed, "a never-freed target is a failure, not a completion");
assertEquals(List.of(), forgotten, "the target is busy, not gone — presence must not be cleared");
assertTrue(inj.activeTargets().isEmpty(), "the target is reclaimed, not polled forever");
}
@Test
void aTargetThatFreesUpBeforeTheQueueWaitGraceIsDeliveredNormally() {
// Positive control: a target that is merely busy for a while, then frees up before the
// grace elapses, is still delivered normally — the long-task case this grace must not break.
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP);
inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < 100; i++) inj.onStatus(T, AgentStatus.WORKING); // busy, well under the grace
assertEquals(List.of(), sent());
inj.onStatus(T, AgentStatus.IDLE); // frees up
assertEquals(List.of("task"), sent(), "a target that frees up within the grace is delivered to");
}
@Test
void dropClearsWorkerPresence() {
// CB-114 (finding #1): a vanished worker's readiness must be forgotten so a stale entry cannot
@@ -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);
@@ -1706,43 +1728,6 @@ class MessageServiceTest {
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* A message still sitting in the injector's queue — never attempted, let alone delivered —
* must not poll as the worker working on it. {@code herdr.agentStatus("working")} supplies the
* target's real, independent live status (busy with something else entirely), while no
* {@code injector.onStatus} call is ever made, so the message is never even attempted.
*/
@Test
void aQueuedButNeverInjectedMessageDoesNotPollAsWorking() throws Exception {
herdr.agentStatus("working");
String ticket = messages.sendAsync(T, "task");
awaitWaiting();
MessageService.TaskView view = messages.poll(ticket);
assertEquals(MessageService.Phase.PENDING, view.phase());
assertFalse(view.detail().contains("worker working"),
"a message never injected must not read as the worker working on it: " + view.detail());
assertTrue(view.detail().contains("queued") && view.detail().contains("not yet delivered"),
"must report the message as queued, not delivered: " + view.detail());
}
/**
* Positive control for the test above: once the message is actually delivered, the poll's
* detail is the plain live-status text again.
*/
@Test
void aDeliveredMessageStillPollsAsWorkerWorking() throws Exception {
herdr.agentStatus("working");
String ticket = messages.sendAsync(T, "task");
awaitWaiting();
injectDelivery();
MessageService.TaskView view = messages.poll(ticket);
assertEquals(MessageService.Phase.PENDING, view.phase());
assertEquals("worker working", view.detail(),
"once actually delivered, the detail reports the worker's live status directly");
}
/**
* fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in
* {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the
@@ -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));
}