diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java index 98ed3a64..dee4d20c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java @@ -32,6 +32,13 @@ 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. */ @@ -107,9 +114,10 @@ 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}), 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}, + * {@code INBOX}), 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 @@ -159,7 +167,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 -> caller.ownsSession(targetSession); + case REPLY, ASK, INBOX -> 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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java index 550c8dd7..e7ef5c22 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java @@ -132,6 +132,12 @@ 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 targets = new ConcurrentHashMap<>(); /** Delivery only; completion signalling is a no-op and every target is treated as available. */ @@ -203,6 +209,7 @@ 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 ready, @@ -232,6 +239,7 @@ public final class Injector { this.ready = ready; this.forget = forget; this.nowMillis = nowMillis; + this.paneInbox = new PaneInbox(nowMillis); } private AgentControl agentsFor(String target) { @@ -348,6 +356,10 @@ public final class Injector { 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); @@ -389,6 +401,11 @@ public final class Injector { return cancellationOf(p); } p.state = Pending.State.CANCELLED; + // The cancelled entry may be the one currently offered to a mod-served pane, so take + // the offer back: a cancelled message the caller was told never arrived must not stay + // collectable. The next poll re-offers whatever is at the head by then. + paneInbox.withdrawAll(p.target); + t.inboxOffer = null; if (isQuiescent(t)) { targets.remove(p.target, t); } @@ -470,10 +487,14 @@ public final class Injector { t.injectableSincePickup = 0; t.awaitingCompletion = false; t.turnObserved = false; - } else { + } else if (!t.deliveredViaInbox) { // 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; } } @@ -498,45 +519,77 @@ public final class Injector { if (p != null && ready.test(target)) { t.notReadySincePoll = 0; t.notReadySinceMillis = 0; - // 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; + 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; } - sent = p; - sendError = e; } } else if (p != null) { // fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in @@ -555,6 +608,8 @@ 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 @@ -633,6 +688,8 @@ public final class Injector { 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, @@ -846,6 +903,24 @@ 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. + * + *

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 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). @@ -899,6 +974,11 @@ 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; diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/PaneInbox.java b/fleetd/src/main/java/dev/ltms/fleet/inject/PaneInbox.java new file mode 100644 index 00000000..da8909dc --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/PaneInbox.java @@ -0,0 +1,141 @@ +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. + * + *

A pane becomes mod-served 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. + * + *

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. + * + *

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 lastPolledAtMillis = new ConcurrentHashMap<>(); + private final ConcurrentHashMap> 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 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 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 drain(String terminal) { + if (terminal == null || terminal.isBlank()) { + return List.of(); + } + lastPolledAtMillis.put(terminal, nowMillis.getAsLong()); + Deque queue = offered.get(terminal); + if (queue == null) { + return List.of(); + } + List 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); + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 3a68d3a6..6d3b92c3 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -553,6 +553,16 @@ 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 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 statusHandler = (exchange, req) -> { McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null); @@ -658,6 +668,7 @@ 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 @@ -668,7 +679,7 @@ public final class FleetMcp { Set 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()); + fleetHandover.name(), fleetInbox.name()); if (!registeredToolNames.equals(FleetTool.wireNames())) { throw new IllegalStateException("fleetd #469: registered MCP tools " + registeredToolNames + " do not match the canonical tool set " + FleetTool.wireNames() @@ -690,6 +701,7 @@ public final class FleetMcp { .toolCall(fleetProfiles, profilesHandler) .toolCall(fleetWhoami, whoamiHandler) .toolCall(fleetHandover, handoverHandler) + .toolCall(fleetInbox, inboxHandler) .build(); this.metrics = metrics; } @@ -1286,6 +1298,7 @@ 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")); @@ -1429,6 +1442,33 @@ 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. + * + *

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. + * + *

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 collected = messages.collectInbox(callerTerminal); + Map 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)) { @@ -2905,6 +2945,18 @@ 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, " diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetTool.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetTool.java index e5fa652a..b41cb52f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetTool.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetTool.java @@ -44,7 +44,8 @@ public enum FleetTool { STOP("fleet_stop"), PROFILES("fleet_profiles"), WHOAMI("fleet_whoami"), - HANDOVER("fleet_handover"); + HANDOVER("fleet_handover"), + INBOX("fleet_inbox"); private final String wireName; diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 64d07a96..b19c147d 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -552,6 +552,20 @@ public final class MessageService { * anything * @return which of the three ways (fleetd #365) the reply actually landed — never {@code null} */ + /** + * 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. + * + *

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 collectInbox(String session) { + return injector.collectInbox(session); + } + public ReplyOutcome reply(String session, String content) { if (content == null || content.isBlank()) { throw new IllegalArgumentException("content is required"); diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java index 9d5828a6..6ef4d636 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java @@ -47,6 +47,28 @@ 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}) { @@ -288,16 +310,17 @@ class AuthzTest { } /** - * Every action beyond READ/METRICS/REPLY/ASK/SEND, asserted denied for an observer — + * Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer — * including {@code TASK_READ}, which is the entire point of this role: an unconfigured pane - * must not be able to poll a ticket or read another session's status. {@code SEND} is excluded - * here and given its own matrix below, since — unlike every action in this loop — its grant is + * 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 * conditional on the target, not fixed. */ @Test - void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() { + void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxAndSend() { for (Authz.Action a : Authz.Action.values()) { - if (a == READ || a == METRICS || a == REPLY || a == ASK || a == SEND) { + if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND) { continue; } assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true), diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorModServedDeliveryTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorModServedDeliveryTest.java new file mode 100644 index 00000000..457c1aac --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorModServedDeliveryTest.java @@ -0,0 +1,147 @@ +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 typed() { + return herdr.calls.stream() + .filter(c -> c.method().equals("agent.prompt")) + .map(c -> ((Map) c.params()).get("text").toString()) + .toList(); + } + + @Test + void aPaneThatCollectsItsOwnMailIsNeverTypedInto() { + injector.collectInbox(MOD); // the pane says it collects its own mail + CompletableFuture 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 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 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)); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/PaneInboxTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/PaneInboxTest.java new file mode 100644 index 00000000..b2fbce8b --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/PaneInboxTest.java @@ -0,0 +1,93 @@ +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)); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceInboxDeliveryTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceInboxDeliveryTest.java new file mode 100644 index 00000000..e76d0bb8 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceInboxDeliveryTest.java @@ -0,0 +1,142 @@ +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 typed() { + return herdr.calls.stream() + .filter(c -> c.method().equals("agent.prompt")) + .map(c -> ((Map) 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)); + } +} diff --git a/plugin/hooks/register.js b/plugin/hooks/register.js index 0a031591..714c6e39 100644 --- a/plugin/hooks/register.js +++ b/plugin/hooks/register.js @@ -7,11 +7,19 @@ // // 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 @@ -131,6 +139,24 @@ export function register(on) { } }) + // 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. Throwing here would kill the + // timer and with it every later message. + 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', diff --git a/plugin/tests/fleet-mod.test.ts b/plugin/tests/fleet-mod.test.ts index cb801cb1..719e78b2 100644 --- a/plugin/tests/fleet-mod.test.ts +++ b/plugin/tests/fleet-mod.test.ts @@ -1,4 +1,4 @@ -import { expect, test } from 'claude-code/testing' +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. @@ -164,3 +164,110 @@ test('/fleet-whoami reports a daemon that is down instead of throwing', async ($ 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 { + const store = new Map() + 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 three MCP requests, with the tool result read fresh on every call. */ +function stubFleetdDynamic(on: any, toolText: () => string, fetches: string[]) { + on('http.fetch', (_$: any, e: any) => { + const body = JSON.parse(e.init.body) + fetches.push(body.method) + 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: 'data: ' + JSON.stringify(result) + '\n' } } + } + return { value: { ok: true, status: 202, headers: {}, text: '' } } + }) +} + +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) +})