Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0dae97e6a9 | |||
| 3cbd84eb07 |
@@ -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));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user