Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ce6e6d11c6 |
@@ -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
|
||||
|
||||
@@ -1,4 +0,0 @@
|
||||
{
|
||||
"description": "fleet mod: cross-account session messaging and PTY-free delivery",
|
||||
"modules": ["./register.js"]
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
Reference in New Issue
Block a user