Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha ce6e6d11c6 fleetd #782: PromptBox ignores non-breaking space padding and bounds the hold streak
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 2m1s
boxContent now treats Character.isSpaceChar(c) the same as isWhitespace(c), so a
no-break space (U+00A0) padding the box no longer reads as the operator's text.

clearToSubmit now gives up after HOLD_GIVE_UP_STREAK consecutive holds on one
target and lets the delivery through, instead of holding a channel shut forever
when the box never reads EMPTY.
2026-10-05 20:03:52 +02:00
19 changed files with 158 additions and 1775 deletions
@@ -32,13 +32,6 @@ public final class Authz {
REPLY,
/** A worker's mid-turn question to the primary. */
ASK,
/**
* Collect the messages queued for the caller's OWN pane, instead of having them typed into
* its terminal. Grouped with {@link #REPLY} and {@link #ASK} below as an only-as-itself
* action: the pane is always the caller's connection-resolved terminal, never an argument,
* so no caller can collect another pane's mail.
*/
INBOX,
/** Collect held replies from a session's inbox. */
DRAIN,
/** Read-only roster, profile, and identity observation: no ticket, task, or turn state. */
@@ -114,10 +107,9 @@ public final class Authz {
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
*
* @param targetSession the session id in the request path; only consulted for the
* worker-scoped actions ({@code REPLY}, {@code ASK},
* {@code INBOX}), for a collaborator's {@code SEND}, and for an
* observer's {@code SEND}, ignored otherwise, may be
* {@code null}
* 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}
* @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
@@ -167,7 +159,7 @@ public final class Authz {
// collaborator's own pane passes through the same check, so each can answer a funnel
// that delegated to it. An unnamed primary (token/loopback, no pane) owns nothing and
// is still excluded.
case REPLY, ASK, INBOX -> caller.ownsSession(targetSession);
case REPLY, ASK -> caller.ownsSession(targetSession);
// READ is roster, profile, and identity observation — fleet_list, fleet_profiles, and
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
@@ -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.
@@ -21,33 +18,36 @@ import java.util.regex.Pattern;
*
* <p>Only a box that is positively empty clears the gate. A box with content, a pane this cannot
* recognise, and a failed read all hold the delivery, because a held delivery is recoverable and a
* submitted half-line is not. Every caller must therefore be a path that retries.
* submitted half-line is not. A hold is bounded, though: after {@link #HOLD_GIVE_UP_STREAK}
* consecutive holds on one target, the gate gives up and lets the delivery through anyway, because a
* channel that never delivers again is worse than one clobbered line.
*
* <p>A pane that holds for {@link #HOLD_WARN_STREAK} consecutive checks gets one warning, so a box
* that never clears is visible instead of silent. The warning repeats only after the box has cleared
* again.
* again, and giving up logs its own warning the same way.
*/
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;
/** Consecutive holds for one target before the gate stops holding and lets the delivery through. */
static final int HOLD_GIVE_UP_STREAK = 1200;
/**
* 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 +57,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 }
@@ -84,7 +75,9 @@ public final class PromptBox {
/**
* Whether {@code target}'s input box is empty, so a delivery would submit only its own text.
* {@code false} means hold and come back; it never means the delivery failed.
* {@code false} means hold and come back; it never means the delivery failed. After {@link
* #HOLD_GIVE_UP_STREAK} consecutive holds on the same target, this returns {@code true} instead,
* so the delivery goes through rather than holding the channel shut forever.
*/
public boolean clearToSubmit(String target) {
Reading reading = inspect(target);
@@ -93,6 +86,13 @@ public final class PromptBox {
return true;
}
int streak = holdStreaks.merge(target, 1, Integer::sum);
if (streak >= HOLD_GIVE_UP_STREAK) {
holdStreaks.remove(target);
log.warn("prompt box of {} never cleared after {} holds ({}, {} character(s) in the box)"
+ " — giving up and sending this delivery anyway, which may submit whatever is in the box",
target, streak, reading.state(), reading.characters());
return true;
}
if (streak == HOLD_WARN_STREAK) {
log.warn("prompt box of {} has held a delivery {} times in a row ({}, {} character(s) in the box)"
+ " — nothing is lost, delivery resumes once the box is empty",
@@ -108,7 +108,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 +128,9 @@ 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 — including non-breaking variants such as U+00A0 — a trailing box border and a
* cursor block all count as nothing. Any other character counts as the operator's unsubmitted
* text, including a placeholder hint a TUI might draw there.
*/
static Reading classify(String pane) {
if (pane == null || pane.isBlank()) return new Reading(State.UNREADABLE, 0);
@@ -155,77 +155,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) || Character.isSpaceChar(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,
@@ -132,12 +121,6 @@ public final class Injector {
* already uses for the same purpose.
*/
private final LongSupplier nowMillis;
/**
* Mail offered to panes that collect it themselves. Owned here because this is the single
* writer of delivery state, and the offer must be made and taken back under the same target
* monitor that guards the queue the message is still sitting on.
*/
private final PaneInbox paneInbox;
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
/** Delivery only; completion signalling is a no-op and every target is treated as available. */
@@ -209,7 +192,6 @@ public final class Injector {
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
this.paneInbox = new PaneInbox(nowMillis);
}
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
@@ -239,7 +221,6 @@ public final class Injector {
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
this.paneInbox = new PaneInbox(nowMillis);
}
private AgentControl agentsFor(String target) {
@@ -315,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() {
@@ -351,15 +329,10 @@ 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;
int injectableSincePostTurnPickup;
/** The head of {@link #queue} as offered to a mod-served pane, or {@code null}. */
PaneInbox.Entry inboxOffer;
/** Whether the delivery the pickup latch is waiting on was collected rather than typed. */
boolean deliveredViaInbox;
synchronized void add(Pending p) {
queue.add(p);
@@ -376,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
@@ -388,8 +361,7 @@ public final class Injector {
/**
* Cancel this exact queued delivery. The target monitor serializes this operation with
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
* rather than claiming the message remained queued. A message a mod-served pane has already
* collected answers the same way, even though no poll has recorded that delivery yet.
* rather than claiming the message remained queued.
*/
public Cancellation cancel(Delivery delivery) {
Pending p = delivery.pending;
@@ -398,22 +370,7 @@ public final class Injector {
return cancellationOf(p);
}
synchronized (t) {
if (p.state != Pending.State.QUEUED) {
return cancellationOf(p);
}
if (t.queue.peek() == p && t.inboxOffer != null) {
// This exact entry is the one offered to a mod-served pane. The pane takes an
// offer on its own thread, so withdraw first and then read the outcome: a taken
// offer means the pane already holds this text, and the next poll records that
// delivery. Cancelling it would tell the caller nothing arrived while the pane
// acts on it.
paneInbox.withdrawAll(p.target);
if (t.inboxOffer.taken()) {
return Cancellation.DELIVERED;
}
t.inboxOffer = null;
}
if (!t.queue.remove(p)) {
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
return cancellationOf(p);
}
p.state = Pending.State.CANCELLED;
@@ -458,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) {
@@ -498,14 +454,10 @@ public final class Injector {
t.injectableSincePickup = 0;
t.awaitingCompletion = false;
t.turnObserved = false;
} else if (!t.deliveredViaInbox) {
} else {
// Delivered but still idle → the worker hasn't picked it up; the submit
// keystroke likely raced the paste (esp. right as the TUI became ready).
// Re-nudge Enter (CB-113) until the worker starts (WORKING) or the grace ends.
//
// A pane that collected the message submits it itself, and nothing was
// typed into it. Pressing Enter there would submit whatever its operator
// has in the prompt box instead.
resubmit = true;
}
}
@@ -530,77 +482,45 @@ public final class Injector {
if (p != null && ready.test(target)) {
t.notReadySincePoll = 0;
t.notReadySinceMillis = 0;
boolean modServed = paneInbox.isModServed(target);
if (t.inboxOffer != null && !modServed) {
// The pane stopped collecting its mail, so take the offer back and
// fall through to the terminal route below. withdrawAll leaves an
// entry the pane took first alone, so the branch under it still sees
// that as the delivery it is.
paneInbox.withdrawAll(target);
if (!t.inboxOffer.taken()) {
t.inboxOffer = null;
}
}
if (t.inboxOffer == null && modServed) {
t.inboxOffer = paneInbox.offer(target, p.text());
}
if (t.inboxOffer != null) {
// The message stays at the head of the queue until the pane takes
// it: nothing has reached the pane yet, so nothing may be recorded
// as delivered and nothing may be failed.
if (t.inboxOffer.taken()) {
t.inboxOffer = null;
t.queue.poll();
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
t.deliveredViaInbox = true;
sent = p;
}
} else {
// fleetd #551: poll and record BEFORE the irreversible send, not after.
// The entry comes off the queue and its state is set to ATTEMPTED here,
// unconditionally — so a Throwable escaping the send call below (caught or
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
// #546 hazard, since peek() alone would let the next onStatus round re-enter
// this block and send the same text again), and no path can write a
// confident DELIVERED or NOT_DELIVERED before we actually know which one
// happened.
t.queue.poll();
p.state = Pending.State.ATTEMPTED;
try {
agentsFor(target).send(target, p.text());
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
t.deliveredViaInbox = false;
sent = p;
} catch (Throwable e) {
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
// catch does not prove the text never reached the pane. Three of the
// four HerdrException throw sites in HerdrCodec fire only after herdr
// has already replied (so it processed the request), and the fourth (a
// transport IOException) leaves it genuinely unknown whether herdr even
// received the bytes — see #551 comment 16867. The one exception is a
// herdr `*_not_found` error: that family is already read as "definitely
// absent, not merely inconclusive" everywhere else in this codebase
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
// target pane/agent does not exist at all, so nothing could have been
// pasted anywhere — #551 keeps the new state consistent with that
// existing vocabulary rather than inventing a second one.
if (e instanceof HerdrException he && he.code() != null
&& he.code().endsWith("_not_found")) {
p.state = Pending.State.NOT_DELIVERED;
}
sent = p;
sendError = e;
// fleetd #551: poll and record BEFORE the irreversible send, not after.
// The entry comes off the queue and its state is set to ATTEMPTED here,
// unconditionally — so a Throwable escaping the send call below (caught or
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
// #546 hazard, since peek() alone would let the next onStatus round re-enter
// this block and send the same text again), and no path can write a
// confident DELIVERED or NOT_DELIVERED before we actually know which one
// happened.
t.queue.poll();
p.state = Pending.State.ATTEMPTED;
try {
agentsFor(target).send(target, p.text());
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
sent = p;
} catch (Throwable e) {
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
// catch does not prove the text never reached the pane. Three of the
// four HerdrException throw sites in HerdrCodec fire only after herdr
// has already replied (so it processed the request), and the fourth (a
// transport IOException) leaves it genuinely unknown whether herdr even
// received the bytes — see #551 comment 16867. The one exception is a
// herdr `*_not_found` error: that family is already read as "definitely
// absent, not merely inconclusive" everywhere else in this codebase
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
// target pane/agent does not exist at all, so nothing could have been
// pasted anywhere — #551 keeps the new state consistent with that
// existing vocabulary rather than inventing a second one.
if (e instanceof HerdrException he && he.code() != null
&& he.code().endsWith("_not_found")) {
p.state = Pending.State.NOT_DELIVERED;
}
sent = p;
sendError = e;
}
} else if (p != null) {
// fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in
@@ -619,8 +539,6 @@ public final class Injector {
for (Pending pending : notReady) {
pending.state = Pending.State.NOT_DELIVERED;
}
paneInbox.withdrawAll(target);
t.inboxOffer = null;
// fleetd #501: t.notReadySincePoll — the loop's own counter, already in
// scope — is printed here instead of the READINESS_GRACE_POLLS constant.
// On this branch the counter has JUST reached the threshold, so the two
@@ -688,30 +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;
}
paneInbox.withdrawAll(target);
t.inboxOffer = null;
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)) {
@@ -782,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
@@ -914,24 +795,6 @@ public final class Injector {
}
}
/**
* Hand {@code terminal} every message held for it and stamp it as collecting its own mail.
* While that stamp is fresh this injector offers that pane's messages instead of typing them;
* once it goes stale the pane's queued mail takes the terminal route again.
*
* <p>Returns the messages in the order they were queued, and an empty list when there are
* none — an empty collection still counts as collecting, so a pane that polls on a timer stays
* mod-served between messages.
*/
public List<String> collectInbox(String terminal) {
return paneInbox.drain(terminal);
}
/** Whether {@code terminal} has collected its mail recently enough to be offered the next one. */
public boolean isModServed(String terminal) {
return paneInbox.isModServed(terminal);
}
/**
* Targets the poller must keep sampling: those with a queued message, an awaited pickup, or an
* awaited turn completion (so the {@code working → idle} boundary is observed).
@@ -949,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
@@ -985,11 +831,6 @@ public final class Injector {
p.state = Pending.State.NOT_DELIVERED;
}
t.queue.clear();
// The pane is gone, so drop its offered mail and its poll stamp together: a terminal
// id can be reused, and a stale stamp would make the next pane under it look mod-served
// before it has ever collected anything.
paneInbox.forget(target);
t.inboxOffer = null;
hadDeliveredTurn = t.awaitingCompletion;
t.awaitingCompletion = false;
t.awaitingPickup = false;
@@ -1,141 +0,0 @@
package dev.ltms.fleet.inject;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
/**
* Mail held for a pane that collects it itself instead of having it typed into its terminal.
*
* <p>A pane becomes <em>mod-served</em> by calling {@code fleet_inbox}: {@link #drain} stamps the
* pane as polling, and {@link #isModServed} answers {@code true} while that stamp is younger than
* {@link #MOD_SERVED_WINDOW_MILLIS}. Nothing else sets it, so a pane that has never polled is
* never mod-served and its mail takes the terminal route.
*
* <p>An offered entry is removed exactly once, by {@link #drain} or by {@link #withdrawAll}, and
* both run under the owning pane's monitor. So an entry the pane took is never also withdrawn, and
* an entry that was withdrawn can never still be collected — which is what lets the {@link
* Injector} keep one message both offered here and queued for the terminal without risking two
* deliveries of it.
*
* <p>This class holds no queue of its own beyond what is currently offered: the {@link Injector}
* keeps the message on its own queue until the pane takes it, so a pane that stops polling strands
* nothing.
*/
public class PaneInbox {
/**
* How long after a {@link #drain} a pane still counts as mod-served. It must exceed the mod's
* own poll interval by enough that a few missed polls are not read as a pane that stopped,
* while staying short enough that a pane which really stopped falls back to the terminal route
* promptly.
*/
public static final long MOD_SERVED_WINDOW_MILLIS = 15_000;
/** One message held for a pane until that pane collects it. */
public static final class Entry {
private final String text;
private boolean taken; // written under the owning pane's monitor
private Entry(String text) {
this.text = text;
}
/** The message text, as it will be handed to the pane. */
public String text() {
return text;
}
/** Whether the pane has collected this entry. Once {@code true} it never goes back. */
public synchronized boolean taken() {
return taken;
}
private synchronized void markTaken() {
taken = true;
}
}
private final LongSupplier nowMillis;
private final ConcurrentHashMap<String, Long> lastPolledAtMillis = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Deque<Entry>> offered = new ConcurrentHashMap<>();
public PaneInbox() {
this(System::currentTimeMillis);
}
public PaneInbox(LongSupplier nowMillis) {
this.nowMillis = Objects.requireNonNull(nowMillis, "nowMillis");
}
/**
* Whether {@code terminal} collected its mail within {@link #MOD_SERVED_WINDOW_MILLIS}. A
* terminal that has never collected any is never mod-served.
*/
public boolean isModServed(String terminal) {
Long at = terminal == null ? null : lastPolledAtMillis.get(terminal);
return at != null && nowMillis.getAsLong() - at <= MOD_SERVED_WINDOW_MILLIS;
}
/** Hold {@code text} for {@code terminal} to collect, and return the entry holding it. */
public Entry offer(String terminal, String text) {
Entry entry = new Entry(text);
Deque<Entry> queue = offered.computeIfAbsent(terminal, _ -> new ArrayDeque<>());
synchronized (queue) {
queue.add(entry);
}
return entry;
}
/**
* Take back every entry {@code terminal} has not collected. An entry the pane took first stays
* taken — this never un-delivers one.
*/
public void withdrawAll(String terminal) {
Deque<Entry> queue = terminal == null ? null : offered.get(terminal);
if (queue == null) {
return;
}
synchronized (queue) {
queue.removeIf(e -> !e.taken());
}
}
/**
* Collect everything held for {@code terminal}, in the order it was offered, and stamp the pane
* as polling. Each returned entry is marked taken, so the {@link Injector} can tell a message
* the pane really has from one it merely offered.
*/
public List<String> drain(String terminal) {
if (terminal == null || terminal.isBlank()) {
return List.of();
}
lastPolledAtMillis.put(terminal, nowMillis.getAsLong());
Deque<Entry> queue = offered.get(terminal);
if (queue == null) {
return List.of();
}
List<String> collected = new ArrayList<>();
synchronized (queue) {
for (Entry e : queue) {
e.markTaken();
collected.add(e.text());
}
queue.clear();
}
return collected;
}
/** Forget a pane that is gone, so neither its poll stamp nor its offered mail lingers. */
public void forget(String terminal) {
if (terminal == null) {
return;
}
lastPolledAtMillis.remove(terminal);
offered.remove(terminal);
}
}
@@ -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;
@@ -553,16 +552,6 @@ public final class FleetMcp {
if (denied != null) return denied;
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
};
// fleet_inbox: the caller collects the mail queued for its OWN pane. Identity is the
// CONNECTION, never an argument, and the authz check is terminal ownership — the same
// shape as fleet_reply above.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> inboxHandler =
(exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_inbox", req.arguments()), self);
if (denied != null) return denied;
return inbox(messages, self);
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
@@ -668,7 +657,6 @@ public final class FleetMcp {
McpSchema.Tool fleetProfiles = profilesTool();
McpSchema.Tool fleetWhoami = whoamiTool();
McpSchema.Tool fleetHandover = handoverTool();
McpSchema.Tool fleetInbox = inboxTool();
// fleetd #469: the tool schemas above are already named from FleetTool.wireName(), but
// this is the check that a schema was not accidentally dropped, duplicated, or added
@@ -679,7 +667,7 @@ public final class FleetMcp {
Set<String> registeredToolNames = Set.of(fleetSend.name(), fleetReply.name(), fleetAsk.name(),
fleetStatus.name(), fleetPoll.name(), fleetAck.name(), fleetSpawn.name(),
fleetList.name(), fleetStop.name(), fleetProfiles.name(), fleetWhoami.name(),
fleetHandover.name(), fleetInbox.name());
fleetHandover.name());
if (!registeredToolNames.equals(FleetTool.wireNames())) {
throw new IllegalStateException("fleetd #469: registered MCP tools " + registeredToolNames
+ " do not match the canonical tool set " + FleetTool.wireNames()
@@ -701,7 +689,6 @@ public final class FleetMcp {
.toolCall(fleetProfiles, profilesHandler)
.toolCall(fleetWhoami, whoamiHandler)
.toolCall(fleetHandover, handoverHandler)
.toolCall(fleetInbox, inboxHandler)
.build();
this.metrics = metrics;
}
@@ -1119,29 +1106,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. */
@@ -1298,7 +1263,6 @@ public final class FleetMcp {
case SEND -> sendAction(str(arguments, "coordId"), str(arguments, "turnId"));
case REPLY -> Authz.Action.REPLY;
case ASK -> Authz.Action.ASK;
case INBOX -> Authz.Action.INBOX;
case STATUS -> Authz.Action.TASK_READ;
case LIST, PROFILES, WHOAMI -> Authz.Action.READ;
case POLL -> pollAction(str(arguments, "target"), str(arguments, "coordId"));
@@ -1442,33 +1406,6 @@ public final class FleetMcp {
return text(outcome.description());
}
/**
* {@code fleet_inbox}: hand the caller the messages queued for its own pane, so it submits them
* itself instead of having them typed into its terminal. {@code callerTerminal} is resolved
* from the connection; there is deliberately no pane argument, so no caller can collect another
* pane's mail.
*
* <p>Calling this is also what marks the pane as collecting its own mail, for a window the
* injector re-checks before every delivery. A pane that stops calling it stops being served
* this way and its queued mail is typed instead, so an empty answer is a normal result that
* must still be requested on a timer.
*
* <p>Returns a JSON object with {@code count} and {@code messages} (in the order they were
* queued), so a caller can tell "no mail" apart from a failure.
*/
static McpSchema.CallToolResult inbox(MessageService messages, String callerTerminal) {
if (callerTerminal == null) {
return error("fleet_inbox could not identify the calling pane from the connection, so "
+ "there is no inbox to collect");
}
List<String> collected = messages.collectInbox(callerTerminal);
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", callerTerminal);
m.put("count", collected.size());
m.put("messages", collected);
return text(json(m));
}
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) {
if (isBlank(target) || isBlank(msgId)) {
@@ -2945,18 +2882,6 @@ public final class FleetMcp {
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool inboxTool() {
return tool(FleetTool.INBOX.wireName(),
"Collect the messages the fleet has queued for YOUR OWN pane, then act on each one. "
+ "Takes no arguments: the pane is resolved from your connection, so you can "
+ "never read another session's mail. Any role may call it for itself. Call it "
+ "on a timer — each call is also what tells the daemon you collect your own "
+ "mail, so it offers your next message here instead of typing it into your "
+ "terminal; stop calling it and your mail is typed instead. An empty "
+ "'messages' array is the normal answer when nothing is waiting.",
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool handoverTool() {
return tool(FleetTool.HANDOVER.wireName(),
"Replace your OWN lead session once its context is full: write a handover file, "
@@ -44,8 +44,7 @@ public enum FleetTool {
STOP("fleet_stop"),
PROFILES("fleet_profiles"),
WHOAMI("fleet_whoami"),
HANDOVER("fleet_handover"),
INBOX("fleet_inbox");
HANDOVER("fleet_handover");
private final String wireName;
@@ -514,20 +514,6 @@ public final class MessageService {
return false;
}
/**
* Hand {@code session} every message queued for it that it has not collected yet, and record
* that it collects its own mail. While that record is fresh, delivery to that session is
* offered for collection instead of typed into its terminal; once it goes stale, the terminal
* route takes over again with nothing lost.
*
* <p>The messages are returned in the order they were queued, and are removed by this call.
* An empty list is an ordinary answer: a session polling on a timer keeps itself collecting
* between messages.
*/
public List<String> collectInbox(String session) {
return injector.collectInbox(session);
}
/**
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
@@ -1440,7 +1426,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.
@@ -1513,21 +1499,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 {
@@ -47,28 +47,6 @@ class AuthzTest {
+ "relies on this gate refusing it first");
}
/**
* {@code fleet_inbox} collects the mail queued for the caller's own pane, so it is gated on
* terminal ownership and on nothing else: every role may do it for itself, and no role may do
* it for another pane.
*/
@Test
void collectingAnInboxIsOnlyEverForTheCallersOwnPane() {
for (Principal self : new Principal[]{WORKER_A, ARCH_DESIGN, COLLABORATOR, OBSERVER}) {
assertTrue(Authz.permits(self, INBOX, self.terminal()),
self.role() + " must be able to collect the mail for its own pane");
assertFalse(Authz.permits(self, INBOX, "term_someone_else"),
self.role() + " must not be able to collect another pane's mail");
}
// A lead carries a pane too, so it collects its own mail on the same rule.
assertTrue(Authz.permits(Principal.leader("opus", "term_lead", 800), INBOX, "term_lead"));
// The unnamed primary owns no pane, so there is no inbox it could be asking for.
assertFalse(Authz.permits(PRIMARY, INBOX, "term_a"),
"a caller with no pane of its own has no inbox to collect");
assertFalse(Authz.permits(WORKER_A, INBOX, null),
"a missing terminal must never match an owner");
}
@Test
void orchestrationBelongsToThePrimaryAlone() {
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, SEND, DRAIN}) {
@@ -310,17 +288,16 @@ class AuthzTest {
}
/**
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer —
* 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. The exempt set is the
* three only-as-itself actions plus the two open reads. {@code SEND} is excluded here and
* given its own matrix below, since — unlike every action in this loop — its grant is
* 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.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxAndSend() {
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == SEND) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
@@ -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());
}
@@ -103,53 +91,36 @@ class PromptBoxTest {
}
@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));
void aNoBreakSpacePaddedBoxIsEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0),
PromptBox.classify("❯     "));
}
@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");
void anOrdinarySpacePaddedBoxIsStillEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify("❯ "));
}
@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");
void aNoBreakSpaceLeadInDoesNotCountTowardTheDraftsCharacters() {
PromptBox.Reading reading = PromptBox.classify("❯ hello");
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");
assertEquals("hello".length(), reading.characters(),
"the no-break space before the typed text does not count, only the text itself");
}
// --- 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
@@ -177,4 +148,41 @@ class PromptBoxTest {
herdr.detectionText(EMPTY);
assertTrue(box.clearToSubmit("term_a"));
}
@Test
void aBoxThatNeverClearsIsLetThroughOnceTheGiveUpStreakIsReached() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
PromptBox box = new PromptBox(new AgentControl(herdr));
for (int i = 1; i < PromptBox.HOLD_GIVE_UP_STREAK; i++) {
assertFalse(box.clearToSubmit("term_a"), "still holding before the give-up streak is reached");
}
assertTrue(box.clearToSubmit("term_a"), "the gate gives up and lets the delivery through");
}
@Test
void givingUpResetsTheStreakSoTheNextHoldStartsOver() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
PromptBox box = new PromptBox(new AgentControl(herdr));
for (int i = 0; i < PromptBox.HOLD_GIVE_UP_STREAK; i++) {
box.clearToSubmit("term_a");
}
assertFalse(box.clearToSubmit("term_a"), "the streak started over, so one more hold is not another give-up");
}
@Test
void theStreakResetsWhenTheBoxEmptiesBeforeTheGiveUpStreak() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
PromptBox box = new PromptBox(new AgentControl(herdr));
for (int i = 0; i < PromptBox.HOLD_GIVE_UP_STREAK - 1; i++) {
box.clearToSubmit("term_a");
}
herdr.detectionText(EMPTY);
assertTrue(box.clearToSubmit("term_a"));
herdr.detectionText(DRAFTED);
assertFalse(box.clearToSubmit("term_a"), "the streak started over, so one more hold is not a give-up");
}
}
@@ -1,200 +0,0 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.msg.TestTurnTokens;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Which route a message takes: offered to a pane that collects its own mail, or typed into the
* pane's terminal. Driven by feeding {@code onStatus}, so no real polling is involved.
*/
class InjectorModServedDeliveryTest {
/** A pane that collects its own mail. */
private static final String MOD = "term_mod";
/** A pane that does not, used as the control for every "nothing was typed" assertion. */
private static final String PTY = "term_pty";
private final FakeHerdr herdr = new FakeHerdr();
private final AtomicLong clock = new AtomicLong(1_000_000);
private final Injector injector = new Injector(new AgentControl(herdr), TurnListener.NOOP,
_ -> true, _ -> {
}, clock::get);
/** The messages typed into a pane, in order. A collected message must never appear here. */
@SuppressWarnings("unchecked")
private List<String> typed() {
return herdr.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.map(c -> ((Map<String, Object>) c.params()).get("text").toString())
.toList();
}
@Test
void aPaneThatCollectsItsOwnMailIsNeverTypedInto() {
injector.collectInbox(MOD); // the pane says it collects its own mail
CompletableFuture<Void> delivered =
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD)).completion();
injector.onStatus(MOD, AgentStatus.IDLE); // offers it for collection
assertFalse(delivered.isDone(), "an offered message has not reached the pane yet");
assertEquals(List.of("do the task"), injector.collectInbox(MOD), "the pane collects it");
injector.onStatus(MOD, AgentStatus.IDLE); // the next sample records the delivery
assertTrue(delivered.isDone(), "a collected message is a delivered message");
assertEquals(List.of(), typed(), "nothing was typed into a pane that collects its own mail");
// The control: without it, an injector that typed nothing anywhere would pass the line
// above. Same injector, same herdr, a pane that never collected its mail.
injector.enqueue(PTY, "type this", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("type this"), typed(), "control: an ordinary pane is still typed into");
}
@Test
void aPaneThatStopsCollectingHasItsMailTypedInstead() {
injector.collectInbox(MOD);
CompletableFuture<Void> delivered =
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD)).completion();
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of(), typed(), "control: while it is still collecting, nothing is typed");
assertFalse(delivered.isDone(), "control: and nothing is reported delivered either");
// The mod stopped calling fleet_inbox, so the pane leaves the window.
clock.addAndGet(PaneInbox.MOD_SERVED_WINDOW_MILLIS + 1);
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("do the task"), typed(), "the message falls back to the terminal route");
assertTrue(delivered.isDone(), "and is reported delivered once it is typed");
assertEquals(List.of(), injector.collectInbox(MOD),
"a message that was typed must not also still be collectable");
}
@Test
void aMessageOfferedForCollectionIsStillReportedAsNotYetDelivered() {
injector.collectInbox(MOD);
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertFalse(injector.queuedWaitMillis(MOD) == null,
"an offered-but-uncollected message is still waiting, not delivered");
injector.collectInbox(MOD);
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(null, injector.queuedWaitMillis(MOD),
"once collected it is off the queue, the same as a typed message");
}
@Test
void aCollectedMessageIsNotFollowedByAnEnterNudge() {
injector.collectInbox(MOD);
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
injector.collectInbox(MOD);
injector.onStatus(MOD, AgentStatus.IDLE); // records the delivery, arms the pickup latch
injector.onStatus(MOD, AgentStatus.IDLE); // still idle: the typed route nudges Enter here
assertFalse(herdr.called("agent.send_keys"),
"a pane that collects its own mail submits it itself; an Enter there would submit "
+ "whatever its operator is typing");
// The control: the nudge really does fire on the typed route, so the absence above is
// this route's behaviour and not a harness that never nudges at all.
injector.enqueue(PTY, "type this", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
injector.onStatus(PTY, AgentStatus.IDLE);
assertTrue(herdr.called("agent.send_keys"), "control: a typed message is nudged");
}
@Test
void aCancelledMessageStopsBeingCollectable() {
injector.collectInbox(MOD);
Injector.Delivery delivery = injector.enqueue(MOD, "retracted", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE); // offered for collection
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(delivery),
"an offered message has not reached the pane, so it can still be cancelled");
assertEquals(List.of(), injector.collectInbox(MOD),
"a cancelled message the caller was told never arrived must not arrive later");
// The control: an uncancelled message on the same route really is collectable, so the
// empty list above is the cancel working and not the offer never being made.
injector.enqueue(MOD, "kept", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("kept"), injector.collectInbox(MOD), "control: an offer is collectable");
}
@Test
void aMessageThePaneAlreadyCollectedCannotBeCancelled() {
injector.collectInbox(MOD);
// The control first: an offer the pane has not taken really is cancellable, so the
// different answer below is the collection and not a cancel that gave up on this route.
Injector.Delivery untaken = injector.enqueue(MOD, "retracted", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(untaken),
"control: an uncollected offer is still cancellable");
Injector.Delivery taken = injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("do the task"), injector.collectInbox(MOD), "the pane takes the offer");
// No poll has run since the pane took it, so the entry is still at the head and still
// QUEUED: the state alone cannot tell this case from an uncollected offer.
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(taken),
"the pane holds this text and will act on it, so nothing can be cancelled");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(taken.completion().isDone(), "the next poll records the delivery");
assertEquals(List.of(), typed(), "and nothing was typed into the pane");
}
@Test
void cancellingALaterMessageLeavesACollectedOneDeliveredOnce() {
injector.collectInbox(MOD);
Injector.Delivery first = injector.enqueue(MOD, "first", TestTurnTokens.inert(MOD));
Injector.Delivery second = injector.enqueue(MOD, "second", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE); // offers the head
assertEquals(List.of("first"), injector.collectInbox(MOD), "the pane takes the head");
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(second),
"a message behind the collected one was never offered, so it cancels");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(first.completion().isDone(),
"cancelling a later message must not lose the record that the head was taken");
assertEquals(List.of(), injector.collectInbox(MOD),
"and the head must not be offered a second time");
// The control: the same injector still hands a later message over, so the empty
// collection above is this one not being re-offered rather than the route going quiet.
injector.onStatus(MOD, AgentStatus.WORKING); // the pane picks the collected message up
injector.onStatus(MOD, AgentStatus.IDLE); // and that turn ends
injector.enqueue(MOD, "third", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("third"), injector.collectInbox(MOD), "control: a later message is offered");
assertEquals(List.of(), typed(), "nothing took the terminal route");
}
@Test
void aPaneThatNeverCollectedIsTypedIntoFromTheStart() {
injector.enqueue(PTY, "do the task", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("do the task"), typed(), "no poll, no offer: the terminal route applies");
assertFalse(injector.isModServed(PTY));
}
}
@@ -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
@@ -1,93 +0,0 @@
package dev.ltms.fleet.inject;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/** What a pane that collects its own mail may and may not see. */
class PaneInboxTest {
private static final String A = "term_a";
private static final String B = "term_b";
private final AtomicLong clock = new AtomicLong(1_000_000);
private final PaneInbox inbox = new PaneInbox(clock::get);
@Test
void aPaneCollectsItsOwnMailAndLeavesTheNextPanesWhereItIs() {
inbox.offer(A, "for a");
inbox.offer(B, "for b");
assertEquals(List.of("for a"), inbox.drain(A), "a pane sees its own message");
// The control for the assertion below: B's message really is there to be missed, so a
// drain that returned everything would have shown it above.
assertEquals(List.of("for b"), inbox.drain(B), "the other pane's message stayed put");
}
@Test
void aCollectedMessageIsNotHandedOverASecondTime() {
inbox.offer(A, "deliver once");
assertEquals(List.of("deliver once"), inbox.drain(A), "control: the first call hands it over");
assertEquals(List.of(), inbox.drain(A), "a collected message is gone from the inbox");
}
@Test
void messagesComeBackInTheOrderTheyWereOffered() {
inbox.offer(A, "first");
inbox.offer(A, "second");
assertEquals(List.of("first", "second"), inbox.drain(A));
}
@Test
void aPaneIsModServedOnlyWhileItKeepsCollecting() {
assertFalse(inbox.isModServed(A), "a pane that has never collected is not mod-served");
inbox.drain(A);
assertTrue(inbox.isModServed(A), "control: collecting is what makes a pane mod-served");
clock.addAndGet(PaneInbox.MOD_SERVED_WINDOW_MILLIS);
assertTrue(inbox.isModServed(A), "control: the window edge still counts as collecting");
clock.addAndGet(1);
assertFalse(inbox.isModServed(A), "a pane that stopped collecting leaves the window");
}
@Test
void withdrawingTakesBackOnlyWhatThePaneHasNotCollected() {
PaneInbox.Entry collected = inbox.offer(A, "already taken");
inbox.drain(A);
PaneInbox.Entry pending = inbox.offer(A, "not taken yet");
inbox.withdrawAll(A);
assertTrue(collected.taken(), "withdrawing must not un-deliver a collected message");
assertFalse(pending.taken(), "control: the uncollected entry was never handed over");
assertEquals(List.of(), inbox.drain(A), "a withdrawn message is no longer collectable");
}
@Test
void forgettingAPaneDropsBothItsMailAndItsPollRecord() {
inbox.drain(A);
inbox.offer(A, "for a");
assertTrue(inbox.isModServed(A), "control: the pane is mod-served and holding mail");
inbox.forget(A);
assertFalse(inbox.isModServed(A), "a gone pane must not look mod-served to the next one");
assertEquals(List.of(), inbox.drain(A), "a gone pane's mail does not outlive it");
}
@Test
void aMissingTerminalCollectsNothingAndIsNeverModServed() {
assertEquals(List.of(), inbox.drain(null));
assertEquals(List.of(), inbox.drain(" "));
assertFalse(inbox.isModServed(null));
}
}
@@ -1,142 +0,0 @@
package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* A delegation whose message the pane collected itself must reach the same ticket states as one
* typed into the pane.
*/
class MessageServiceInboxDeliveryTest {
/** A pane that collects its own mail. */
private static final String MOD = "term_mod";
/** A pane that does not — the control for every route-specific assertion here. */
private static final String PTY = "term_pty";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final Injector injector = new Injector(agents);
private final InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox);
@BeforeEach
void setUp() {
replyInbox.own(MOD);
replyInbox.own(PTY);
}
/** The messages typed into a pane, in order. A collected message must never appear here. */
@SuppressWarnings("unchecked")
private List<String> typed() {
return herdr.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.map(c -> ((Map<String, Object>) c.params()).get("text").toString())
.toList();
}
/** {@code sendAsync} queues on another thread, so wait for the delivery to exist. */
private void awaitWaiting(String target) throws InterruptedException {
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(target),
"the delegation to " + target + " should have opened its rendezvous waiter");
}
/** A ticket's terminal state is stamped by the delegating thread, so poll until it settles. */
private MessageService.TaskView awaitSettled(String ticket) throws InterruptedException {
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while ((view == null || view.phase() == MessageService.Phase.PENDING)
&& System.currentTimeMillis() < deadline) {
view = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(view, ticket + " must still be a known ticket");
return view;
}
@Test
void aTicketCollectedByThePaneReachesTheSameStatesAsATypedOne() throws Exception {
// The collecting pane announces itself before anything is delegated to it; the control
// pane never calls fleet_inbox at all.
assertEquals(List.of(), messages.collectInbox(MOD), "nothing is waiting yet");
String collectedTicket = messages.sendAsync(MOD, "long task");
awaitWaiting(MOD);
String typedTicket = messages.sendAsync(PTY, "long task");
awaitWaiting(PTY);
injector.onStatus(MOD, AgentStatus.IDLE); // offers the task for collection
injector.onStatus(PTY, AgentStatus.IDLE); // types the task into the pane
assertEquals(List.of("long task"), typed(),
"only the control pane was typed into; the collecting pane was not");
assertTrue(messages.poll(collectedTicket).detail().contains("queued, not yet delivered"),
"control: an offered-but-uncollected message has not reached its pane");
assertFalse(messages.poll(typedTicket).detail().contains("queued, not yet delivered"),
"control: a typed message has reached its pane");
assertEquals(List.of("long task"), messages.collectInbox(MOD), "the pane collects the task");
injector.onStatus(MOD, AgentStatus.IDLE); // records the delivery
assertFalse(messages.poll(collectedTicket).detail().contains("queued, not yet delivered"),
"a collected message is reported delivered, the same as a typed one");
injector.onStatus(MOD, AgentStatus.WORKING);
injector.onStatus(PTY, AgentStatus.WORKING);
// The reply still arrives through fleet_reply, and lands the same way on both routes.
MessageService.ReplyOutcome collectedReply = messages.reply(MOD, "task done");
MessageService.ReplyOutcome typedReply = messages.reply(PTY, "task done");
assertEquals(typedReply, collectedReply, "both routes resolve their delegation the same way");
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, collectedReply);
MessageService.TaskView collectedDone = awaitSettled(collectedTicket);
MessageService.TaskView typedDone = awaitSettled(typedTicket);
assertEquals(MessageService.Phase.DONE, collectedDone.phase(),
"a ticket delivered by collection must not strand at PENDING");
assertEquals(MessageService.Phase.DONE, typedDone.phase(),
"control: the typed route reaches the same terminal state");
assertEquals("task done", collectedDone.reply());
assertEquals(typedDone.replySource(), collectedDone.replySource());
assertEquals(List.of("long task"), typed(),
"the collected delegation ran start to finish without typing into its pane");
}
@Test
void collectingAnInboxReachesOnlyTheCallersOwnMail() throws Exception {
assertEquals(List.of(), messages.collectInbox(MOD));
assertEquals(List.of(), messages.collectInbox(PTY));
messages.sendAsync(MOD, "task for the mod pane");
awaitWaiting(MOD);
messages.sendAsync(PTY, "task for the other pane");
awaitWaiting(PTY);
injector.onStatus(MOD, AgentStatus.IDLE);
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("task for the mod pane"), messages.collectInbox(MOD));
// The control: the other pane's task really was waiting for it, so a collect that
// returned everyone's mail would have shown it on the line above.
assertEquals(List.of("task for the other pane"), messages.collectInbox(PTY));
}
}
@@ -1706,43 +1706,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
-4
View File
@@ -1,4 +0,0 @@
{
"description": "fleet mod: cross-account session messaging and PTY-free delivery",
"modules": ["./register.js"]
}
-261
View File
@@ -1,261 +0,0 @@
// fleet mod — session messaging for the claude-bridge fleet.
//
// $.store is kept under CLAUDE_CONFIG_DIR, so /fleet-peers and /fleet-mail reach
// only sessions that share this session's config dir. fleetd names a caller by
// its pane, not its account, so /fleet-whoami and anything else sent through
// fleetTool reach the fleet from either account.
//
// Delivery uses $.prompt.submit, so nothing is typed into a pane and no prompt
// box is read.
//
// There are two inboxes. The $.store one carries /fleet-mail between sessions
// that share this config dir. fleet_inbox carries what the fleet queued for
// this pane, from any account, and calling it is also what tells fleetd to
// queue here rather than type into the terminal.
const PRESENCE_PREFIX = 'presence:'
const INBOX_PREFIX = 'inbox:'
const PRESENCE_REFRESH_MS = 15_000
const INBOX_POLL_MS = 3_000
// fleetd stops queueing for a pane that goes quiet, so this must stay well
// under the window the daemon allows between calls.
const FLEETD_INBOX_POLL_MS = 3_000
// A session whose presence row is older than this is treated as gone. It must
// exceed PRESENCE_REFRESH_MS by enough that one missed refresh is not a death.
const PRESENCE_STALE_MS = 60_000
const MAX_INBOX = 50
/** The store key holding one session's queued messages. */
function inboxKey(sessionId) {
return INBOX_PREFIX + sessionId
}
/** The store key holding one session's presence row. */
function presenceKey(sessionId) {
return PRESENCE_PREFIX + sessionId
}
/**
* Append one message to a target's inbox.
*
* $.store has no compare-and-swap, so two senders writing in the same instant
* can lose a message. Callers that need delivery confirmed should read the
* inbox back.
*/
async function deliver($, target, message) {
const key = inboxKey(target)
const queued = (await $.store.get(key)) || []
queued.push(message)
// Keep the newest: an unread inbox must not grow without bound.
const kept = queued.slice(-MAX_INBOX)
await $.store.set(key, kept)
return kept.length
}
/** Every session that refreshed its presence row recently, newest first. */
async function livePeers($, now) {
const keys = await $.store.keys()
const rows = []
for (const key of keys) {
if (!key.startsWith(PRESENCE_PREFIX)) continue
const row = await $.store.get(key)
if (!row || typeof row.at !== 'number') continue
if (now - row.at > PRESENCE_STALE_MS) continue
rows.push(row)
}
rows.sort((a, b) => b.at - a.at)
return rows
}
const FLEETD_MCP = 'http://127.0.0.1:8765/mcp'
const MCP_HEADERS = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' }
// The status fleetd answers, with "Session not found", for an Mcp-Session-Id it no longer holds.
const MCP_SESSION_GONE = 404
// The MCP session every call below shares. fleetd keeps a server-side session per initialize and
// drops it only on a DELETE, so one initialize per call would leave a session behind every time.
let mcpSessionId = null
/**
* Open an MCP session on the local fleetd and hold it for later calls.
*
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse.
*/
async function openFleetSession($) {
mcpSessionId = null
const init = await $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers: MCP_HEADERS,
body: JSON.stringify({
jsonrpc: '2.0', id: 1, method: 'initialize',
params: { protocolVersion: '2025-06-18', capabilities: {}, clientInfo: { name: 'fleet-mod', version: '0' } },
}),
})
if (!init.ok) throw new Error('fleetd initialize failed with status ' + init.status)
const opened = init.headers['mcp-session-id']
await $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': opened },
body: JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }),
})
mcpSessionId = opened
}
/** Send one tools/call on the session this mod holds, and return the raw HTTP answer. */
function sendFleetToolCall($, tool, args) {
return $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': mcpSessionId },
body: JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: tool, arguments: args } }),
})
}
/**
* Call one fleet_* tool on the local fleetd, and return its text result.
*
* fleetd names the caller from the TCP connection on every request, not from the MCP session, so
* the answer is about this session's own pane whichever Claude account the session runs on, and
* reusing one session never changes whose call it is.
*/
async function fleetTool($, tool, args) {
if (mcpSessionId === null) await openFleetSession($)
let call = await sendFleetToolCall($, tool, args)
if (call.status === MCP_SESSION_GONE) {
// A daemon restart drops every session it held. Open a new one and retry once.
await openFleetSession($)
call = await sendFleetToolCall($, tool, args)
}
if (!call.ok) throw new Error('fleetd ' + tool + ' failed with status ' + call.status)
// A tool call answers as a server-sent event: the JSON is on the data: line.
const line = call.text.split('\n').find((l) => l.startsWith('data:'))
const body = JSON.parse(line ? line.slice(5) : call.text)
if (body.error) throw new Error(body.error.message)
return body.result.content.map((c) => c.text).join('\n')
}
export function register(on) {
on('session.start', async ($, e, next) => {
const self = await $.session.id()
// Announce before the first refresh is due, or a session shorter than one
// refresh interval never appears to its peers at all.
await $.store.set(presenceKey(self), {
sessionId: self,
cwd: await $.session.cwd(),
at: await $.clock.now(),
})
// Keep announcing: a row that stops being refreshed is how another session
// learns this one is gone.
$.clock.every(PRESENCE_REFRESH_MS, async () => {
const now = await $.clock.now()
await $.store.set(presenceKey(self), {
sessionId: self,
cwd: await $.session.cwd(),
at: now,
})
})
// Collect this session's mail and hand it to Claude. $.prompt.submit waits
// for the session to be idle, so this never lands mid-turn.
$.clock.every(INBOX_POLL_MS, async () => {
const key = inboxKey(self)
const queued = (await $.store.get(key)) || []
if (queued.length === 0) return
await $.store.set(key, [])
for (const message of queued) {
await $.prompt.submit({
text: 'Message from fleet session ' + message.from + ':\n\n' + message.text,
})
}
})
// Collect what the fleet queued for this pane and hand each message to
// Claude. Every call also renews fleetd's record that this pane collects its
// own mail, so an empty answer still has to be asked for.
$.clock.every(FLEETD_INBOX_POLL_MS, async () => {
let collected
try {
collected = JSON.parse(await fleetTool($, 'fleet_inbox', {}))
} catch {
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
// case on a host with no fleet running. The timer survives a throw, so
// this only keeps every tick from writing an error to the debug log.
return
}
for (const text of collected.messages || []) {
await $.prompt.submit({ text: 'Message from the fleet, via fleetd:\n\n' + text })
}
})
await $.command.register({
name: 'fleet-peers',
description: 'List fleet sessions on this machine, including other accounts',
})
await $.command.register({
name: 'fleet-whoami',
description: 'Show who fleetd says this session is',
})
await $.command.register({
name: 'fleet-mail',
description: 'Send a message to a fleet session on this machine',
argumentHint: '<sessionId> <text>',
// Runs even while Claude is working, so a correction is never queued
// behind the turn it is meant to correct.
immediate: true,
})
return next(e)
})
on('command.run', { command: 'fleet-peers' }, async ($) => {
const now = await $.clock.now()
const self = await $.session.id()
const peers = await livePeers($, now)
if (peers.length === 0) return { text: 'No fleet sessions have announced themselves yet.' }
const lines = peers.map((p) => {
const age = Math.round((now - p.at) / 1000)
const mark = p.sessionId === self ? ' (this session)' : ''
return p.sessionId + ' ' + p.cwd + ' seen ' + age + 's ago' + mark
})
return { text: 'Fleet sessions on this machine:\n' + lines.join('\n') }
})
on('command.run', { command: 'fleet-mail' }, async ($, e) => {
const args = (e.args || '').trim()
const split = args.indexOf(' ')
if (split < 1) return { text: 'Usage: /fleet-mail <sessionId> <text>' }
const target = args.slice(0, split)
const text = args.slice(split + 1).trim()
if (text === '') return { text: 'Usage: /fleet-mail <sessionId> <text>' }
const now = await $.clock.now()
const peers = await livePeers($, now)
if (!peers.some((p) => p.sessionId === target)) {
return { text: 'No live fleet session ' + target + '. Run /fleet-peers.' }
}
const self = await $.session.id()
const depth = await deliver($, target, { from: self, text: text, at: now })
return { text: 'Queued for ' + target + ' (' + depth + ' in its inbox).' }
})
on('command.run', { command: 'fleet-whoami' }, async ($) => {
try {
return { text: 'fleetd says: ' + (await fleetTool($, 'fleet_whoami', {})) }
} catch (err) {
return { text: 'fleetd unreachable: ' + err.message }
}
})
// Record what arrives over Claude Code's own channel, so a message delivered
// by the fleet and one delivered by SendMessage can be told apart.
//
// This hook gates delivery, so it must never decide the message's fate. The
// .catch handler passes the message on when the logging above throws.
on('session.receive', async ($, e, next) => {
$.ui.log('fleet: inbound ' + (e.origin && e.origin.kind) + ', ' + String(e.text).length + ' chars')
return next(e)
}).catch(async ($, e, next) => {
if (next.called) return undefined
return next(e)
})
}
-349
View File
@@ -1,349 +0,0 @@
import { expect, mock, test } from 'claude-code/testing'
// The store-backed tests write the row another session would write, because a
// test cannot start a second session.
const PEER = 'peer-session'
const SELF = 'this-session'
// A fixed clock keeps the staleness arithmetic exact.
const NOW = 1_700_000_000_000
/**
* Answer the mods API calls the harness has no implementation for, and stand in
* for Claude Code's own behaviour beneath the mod's gating hooks.
*
* Returns the map backing $.store. The harness puts no `store` namespace on the
* test's own `$`, so a test seeds and inspects the carrier through this map.
*/
function stubEngine(on: any): Map<string, any> {
const store = new Map<string, any>()
on('store.get', (_$: any, e: any) => ({ value: store.get(e.key) }))
on('store.set', (_$: any, e: any) => {
store.set(e.key, e.value)
return { value: undefined }
})
on('store.delete', (_$: any, e: any) => {
store.delete(e.key)
return { value: undefined }
})
on('store.keys', () => ({ value: [...store.keys()] }))
on('clock.now', () => ({ value: NOW }))
on('session.id', () => ({ value: SELF }))
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('ui.log', () => ({ value: undefined }))
on('prompt.submit', () => ({ value: undefined }))
// Claude Code's own delivery, which the mod's receive hook must reach.
on('session.receive', (_$: any, e: any) => e)
return store
}
/** The presence row a session on the other account would write. */
function announce(store: Map<string, any>, sessionId: string, cwd: string, at: number) {
store.set('presence:' + sessionId, { sessionId, cwd, at })
}
test('/fleet-peers lists a session that announced itself', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).toContain(PEER)
expect(answer.text).toContain('/Users/x/work-repo')
})
test('/fleet-peers hides a session whose presence row went stale', async ($, on) => {
const store = stubEngine(on)
// One second past the 60s staleness cut.
announce(store, PEER, '/Users/x/work-repo', NOW - 61_000)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).not.toContain(PEER)
})
test('/fleet-peers keeps a row one second inside the staleness cut', async ($, on) => {
// The positive control for the test above: without this, a bug that hid
// every row would still satisfy that assertion.
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW - 59_000)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).toContain(PEER)
})
test('/fleet-mail queues a message in the target session inbox', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
const answer = await $.command.run({
command: 'fleet-mail',
args: PEER + ' rebasing on main is safe now',
})
expect(answer.text).toContain('Queued for ' + PEER)
const inbox = store.get('inbox:' + PEER)
expect(inbox.length).toBe(1)
expect(inbox[0].text).toBe('rebasing on main is safe now')
expect(inbox[0].from).toBe(SELF)
})
test('a second message appends rather than replacing the first', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
await $.command.run({ command: 'fleet-mail', args: PEER + ' first' })
await $.command.run({ command: 'fleet-mail', args: PEER + ' second' })
const inbox = store.get('inbox:' + PEER)
expect(inbox.length).toBe(2)
expect(inbox[0].text).toBe('first')
expect(inbox[1].text).toBe('second')
})
test('/fleet-mail refuses a target that never announced itself', async ($, on) => {
const store = stubEngine(on)
const answer = await $.command.run({ command: 'fleet-mail', args: 'ghost-session hello' })
expect(answer.text).toContain('No live fleet session ghost-session')
// Nothing may be queued for a session we could not confirm.
expect(store.get('inbox:ghost-session')).toBe(undefined)
})
test('/fleet-mail rejects input with no message text', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
for (const args of ['', PEER, PEER + ' ']) {
const answer = await $.command.run({ command: 'fleet-mail', args })
expect(answer.text).toContain('Usage: /fleet-mail')
}
expect(store.get('inbox:' + PEER)).toBe(undefined)
})
test('an inbound peer message is passed on, not consumed', async ($, on) => {
// A receive hook that withheld a message would silently break Claude Code's
// own channel, so assert this mod stays transparent.
stubEngine(on)
const result = await $.session.receive({
text: 'from the other session',
origin: { kind: 'peer' },
})
expect(result?.consumed).toBe(undefined)
})
/** Answer fleetd's three MCP requests the way the live daemon does. */
function stubFleetd(on: any, toolText: string, seen: any[]) {
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
seen.push({ method: body.method, session: e.init.headers['Mcp-Session-Id'] })
if (body.method === 'initialize') {
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-1' }, text: '{}' } }
}
if (body.method === 'tools/call') {
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text: toolText }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'event: message\ndata: ' + JSON.stringify(result) + '\n' } }
}
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
}
test('/fleet-whoami reads the tool result out of the event stream', async ($, on) => {
stubEngine(on)
const seen: any[] = []
stubFleetd(on, '{"role":"observer","sessionId":"term_x"}', seen)
const answer = await $.command.run({ command: 'fleet-whoami', args: '' })
expect(answer.text).toBe('fleetd says: {"role":"observer","sessionId":"term_x"}')
// The session id from initialize must ride on every later request.
expect(seen.map((r) => r.method)).toEqual(['initialize', 'notifications/initialized', 'tools/call'])
expect(seen[2].session).toBe('sid-1')
})
test('/fleet-whoami reports a daemon that is down instead of throwing', async ($, on) => {
stubEngine(on)
on('http.fetch', () => ({ value: { ok: false, status: 503, headers: {}, text: '' } }))
const answer = await $.command.run({ command: 'fleet-whoami', args: '' })
expect(answer.text).toBe('fleetd unreachable: fleetd initialize failed with status 503')
})
/**
* The session.start hook's own calls, for a test that fires it. Kept apart from stubEngine
* because a timer test drives $.clock through mock.clock(on) instead of a fixed clock.now.
*/
function stubSessionStart(on: any, submitted: string[]): Map<string, any> {
const store = new Map<string, any>()
on('store.get', (_$: any, e: any) => ({ value: store.get(e.key) }))
on('store.set', (_$: any, e: any) => {
store.set(e.key, e.value)
return { value: undefined }
})
on('store.keys', () => ({ value: [...store.keys()] }))
on('session.id', () => ({ value: SELF }))
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('command.register', () => ({ value: undefined }))
on('ui.log', () => ({ value: undefined }))
// The engine skips a prompt.submit hook that answers anything but { text } or { drop }, and
// the mod's callback then throws, so this must hand the text straight back.
on('prompt.submit', (_$: any, e: any) => {
submitted.push(e.text)
return { text: e.text }
})
on('session.start', (_$: any, e: any) => e)
return store
}
/**
* Answer fleetd's MCP requests, with the tool result read fresh on every call.
*
* `fetches` collects every method sent, and `sessions` the Mcp-Session-Id of each tools/call.
* Each initialize hands out the next id, so a reused session and a reopened one differ.
* `sessionGone` makes a tools/call answer the way fleetd answers for a session id it no longer
* holds.
*/
function stubFleetdDynamic(
on: any,
toolText: () => string,
fetches: string[],
sessions: string[] = [],
sessionGone: () => boolean = () => false,
) {
let opened = 0
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
fetches.push(body.method)
if (body.method === 'initialize') {
opened += 1
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-' + opened }, text: '{}' } }
}
if (body.method === 'tools/call') {
sessions.push(e.init.headers['Mcp-Session-Id'])
if (sessionGone()) {
const gone = '{"jsonRpcError":{"code":-32603,"message":"Session not found"}}'
return { value: { ok: false, status: 404, headers: {}, text: gone } }
}
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text: toolText() }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
}
/** How many of `fetches` were an initialize. */
function initializes(fetches: string[]): number {
return fetches.filter((m) => m === 'initialize').length
}
test('the fleetd inbox poll submits each collected message and names the sender', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
let inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), [])
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
// An empty inbox is the ordinary answer and must submit nothing.
await clock.advance(3_000)
expect(submitted.length).toBe(0)
// The control for the line above: the same timer, the same stubs, with mail waiting.
inbox = { sessionId: 'term_self', count: 2, messages: ['rebase is safe now', 'build is green'] }
await clock.advance(3_000)
expect(submitted.length).toBe(2)
expect(submitted[0]).toContain('rebase is safe now')
expect(submitted[0]).toContain('fleetd')
expect(submitted[1]).toContain('build is green')
})
test('a collected message is not submitted a second time', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
// fleetd removes a message when it hands it over, so the next poll answers empty.
let inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), [])
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(submitted.length).toBe(1) // control: the first poll really did deliver it
inbox = { sessionId: 'term_self', count: 0, messages: [] }
await clock.advance(3_000)
expect(submitted.length).toBe(1)
})
test('a fleetd that is down leaves the poll timer running', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
on('http.fetch', (_$: any, e: any) => {
fetches.push(JSON.parse(e.init.body).method)
return { value: { ok: false, status: 503, headers: {}, text: '' } }
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
const afterFirst = fetches.length
expect(afterFirst).toBeGreaterThan(0)
expect(submitted.length).toBe(0)
// A throw out of the timer would stop it, so a second tick that still reaches fleetd is the
// positive control for the assertion above: the timer survived the failure.
await clock.advance(3_000)
expect(fetches.length).toBeGreaterThan(afterFirst)
expect(submitted.length).toBe(0)
})
test('a second inbox poll reuses the first MCP session', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
const inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions)
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
// fleetd holds a session per initialize and drops it only on a DELETE, so one initialize per
// poll would leave one behind every three seconds.
expect(sessions.length).toBe(2) // control: both polls really reached fleetd
expect(initializes(fetches)).toBe(1)
expect(sessions[1]).toBe(sessions[0])
})
test('a session fleetd no longer holds is opened again and the call retried', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
let gone = false
const inbox = { sessionId: 'term_self', count: 1, messages: ['the daemon restarted'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions, () => {
// Only the first call after the flag is set is refused; the retry succeeds.
const refuse = gone
gone = false
return refuse
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
// The control: an ordinary poll opens one session and needs no second one.
await clock.advance(3_000)
expect(initializes(fetches)).toBe(1)
expect(submitted.length).toBe(1)
gone = true
await clock.advance(3_000)
expect(initializes(fetches)).toBe(2)
expect(sessions[sessions.length - 1]).toBe('sid-2')
expect(submitted.length).toBe(2)
})