import type { TaskHandle, TaskScheduler } from '../types/agent/scheduler.js'; import type { TaskId } from '../types/ids/index.js'; import { type Logger } from '../utils/logger.js'; /** * Completions that finished with nobody left to hear them. * * A worker's result reaches the supervisor as the `tool_result` of the * `create_task` that launched it. That works whenever the launching call is * still the live path — but it is not the only way a task ends: * * - the launching tool hit its deadline and the executor returned * *"timed out… it may still be running"* to the model. The worker then * finished normally, holding a result nothing would ever read. * - the task was launched in the background on purpose, so there is no * call waiting on it by design. * * In both cases the completion exists, the gateway remembers it, and the * model is never told. That is the gap this closes: the turn subscribes once, * every settled task lands here, and anything a tool did NOT hand over * inline is drained into the transcript as a notification the next turn can * read. * * The disambiguation is the whole design. An earlier version of the envelope * path was removed (`dc16d58`) because it fired for completions the blocking * tool had ALREADY delivered, so the supervisor saw every result twice — * once correctly as a `tool_result`, once as an orphan envelope. Removing it * fixed the duplicate and left the abandoned case with no channel at all. * Claiming is what tells the two apart: a tool that delivers a completion * says so, and only unclaimed completions become envelopes. * * It attaches through `onTaskCompleted`, which every `TaskScheduler` already * has, so a host gateway needs no change to take part — a host that was * firing completions into a listener set with no listeners now has one. */ export declare class CompletionInbox { private readonly log?; private readonly unheard; private readonly claimed; /** Launched with nothing waiting on it, and not settled yet. */ private readonly outstanding; /** * Tasks THIS turn launched. * * `onTaskCompleted` is a broadcast and `TaskHandle` carries no turn id, so * a gateway shared between two supervisors hands every completion to both * of their inboxes. Measured: with two inboxes on one gateway, the turn * that launched nothing drained the other turn's task and would have been * told "a task you launched has finished" — a claim that was false, over * another turn's worker output, in a transcript whose model then has to * account for it. * * A shared gateway is not an abuse of the API: `SupervisorAgentConfig` * takes one, and a host that owns a gateway naturally reuses it. */ private readonly ours; /** * Owned tasks not yet observed to have settled, in launch order. * * A `Set`'s iteration order is insertion order, so the most recently * launched entry is always last — {@link describeOwnedWork} reads that * order back to front. An id leaves this set the moment its completion is * observed, by whichever of the three paths sees it first: the gateway's * broadcast, an announcement recovered from {@link unowned}, or a launch * that finds the task already terminal. It is never evicted for any other * reason, so a long-running task stays in here — and therefore visible — * for exactly as long as it is actually running, regardless of how many * other tasks this turn launches meanwhile. */ private readonly runningOwned; /** * Owned tasks observed to have settled, oldest first, bounded to * {@link OWNED_WORK_DISPLAY_LIMIT}. * * This bound is display-driven, not a memory concern the way * {@link UNOWNED_BUFFER_LIMIT} is: a settled task's own result already * reached the model once (inline or as a notification), so dropping it * from here loses a convenience, not the only record of it. */ private readonly settledOwned; /** * Announcements that arrived before anyone said whose task it was. * * `gateway.createTask` resolves one microtask before its caller can name * the task, and a worker that finishes inside that window is announced * first — `LocalTaskScheduler` attaches its completion continuation before * it returns the handle, so the ordering is guaranteed to be reachable * rather than merely possible. Dropping an unowned announcement outright * would therefore turn the leak fix into a LOST RESULT for exactly the * fast completions the inbox exists to catch. * * So they wait here, and ownership may be claimed retroactively. What * makes that safe rather than a second leak is the bound: on a gateway * shared with other turns this fills with completions that will never be * claimed, each holding a whole worker result. */ private readonly unowned; private readonly arrivals; private detach?; /** Kept for {@link launched}: the source of truth about a task's state. */ private gateway?; /** * `log` is optional and unresolved until it is actually needed (the * eviction-warning path in `hold()`), same reason as `LocalTaskScheduler`. */ constructor(log?: Logger | undefined); /** * Start listening. * * Returns a detach function; calling `attach` twice is a no-op rather * than a second subscription, because a doubly-attached inbox would * queue every completion twice and reproduce the exact duplicate this * class exists to prevent. */ attach(gateway: TaskScheduler): () => void; /** * Park an announcement nobody has claimed yet, evicting the oldest if the * buffer is full. * * An eviction is logged at WARN, and that is not decoration. If the entry * turned out to be ours, its completion has just been dropped — the * original defect, wearing the cap as a disguise — and the only evidence * would otherwise be an absence, which is precisely the shape of failure * this whole session has been closing. A reader who sees this line knows * where to look. */ private hold; /** * Say that this turn launched the task. * * Required before anything about the task can reach this inbox — see * {@link ours}. Every launch says it, whether or not something is waiting * on the result, because the case the inbox exists for is precisely the * one where the waiter gave up. * * **The late-announcement branches are not defensive.** * `gateway.createTask` resolves one microtask before its caller can say * who owns the task, and a worker that finishes inside that window is * announced first — the same ordering that used to leave a permanent * pending flag. Without recovery the ownership check would turn that race * from a stale flag into a lost result. * * There are two, and the second is the safety net for the first. The * buffer ({@link unowned}) needs nothing from the gateway beyond the * announcement it already made. Asking `getTask` covers the case the * buffer cannot — an announcement evicted under load — but it rests on a * gateway still knowing about a task it has just settled, which is a * property of the implementations here rather than of the `TaskScheduler` * contract, and `getTask`'s own docs now say so. */ launched(taskId: TaskId): void; /** * Move an owned task from {@link runningOwned} to {@link settledOwned}. * * Idempotent: `claim` and the completion listener can both observe the * same settlement (the race either can win), and re-marking a task that * is already in {@link settledOwned} moves it to the most-recently-settled * end rather than duplicating it there. */ private markSettled; /** * Say that a task was launched with nothing waiting on it. * * {@link launched} plus the statement that no call will deliver the * result. Without the second half the inbox can only see completions that * have already happened, and a turn whose supervisor launched a background * worker and then answered would settle while the worker was still going — * throwing away the very result the launch existed to produce. Knowing a * task is outstanding is what lets the loop hold the turn open for it. */ expect(taskId: TaskId): void; /** Whether anything is either waiting to be told or still running. */ get hasPendingWork(): boolean; /** * Wait for the next completion, deadline or abort, whichever comes first. * Aborting releases only this waiter; work and undelivered results remain owned. * * Bounded on purpose. A worker that never finishes must not hold a turn * open forever, and the caller decides how long "long enough" is — the * turn's own budget is the only thing that knows. */ waitForArrival(timeoutMs: number, signal?: AbortSignal): Promise; /** * Say that this completion reached the model as a `tool_result`. * * Idempotent, and safe to call before the completion is announced: the * claim is remembered so a late announcement does not re-queue it. */ claim(taskId: TaskId): void; /** Whether anything is waiting to be told. */ get hasUnheard(): boolean; /** * Bounded, non-consuming scheduler observations. Delivery never * establishes user-facing completion. * * Running tasks — most recently launched first, matching the order the * rest of this projection has always used — fill the {@link * OWNED_WORK_DISPLAY_LIMIT} slots first, so a task still going is named * here for as long as it keeps running, no matter how many other tasks * this turn has since launched. The most recently SETTLED tasks fill * whatever slots running tasks leave over. A settled task bumped out * entirely is not reported missing — its own result already reached the * model once — but a RUNNING task that does not fit is: the preamble * says exactly how many, so the model never mistakes the cap for the * task having ended. */ describeOwnedWork(): string | undefined; /** * Tasks this turn launched that are still running. * * Read when a turn ends, so it can say which work it walked away from. * Nothing here is cancelled by being read — the ids are a statement, and * what to do about them is the host's call. */ get outstandingTaskIds(): readonly TaskId[]; /** * Take every unheard completion, leaving the inbox empty. * * Draining rather than peeking: a notification that stays queued after * being delivered is the duplicate-delivery bug in a different costume. */ drain(): TaskHandle[]; /** * Stop expecting a task that is never going to arrive. * * Cancelling is the case this exists for. `expect` puts a task on the * outstanding list and only a COMPLETION takes it off, so a cancelled * worker left `hasPendingWork` true for the rest of the turn — and every * attempt to settle then paid the full grace period waiting for a result * that had been called off. */ forget(taskId: TaskId): void; /** * Stop listening. Safe to call more than once. * * A turn that ends without this leaves its listener on the gateway * forever. On a gateway the host reuses that is measurable — three * sequential turns left three live subscriptions, each still holding its * turn's handles — and the listener set only grows. Ownership stops a * retained listener from DELIVERING another turn's work; closing is what * stops it existing. */ close(): void; } /** * The message a supervisor reads when a worker it stopped waiting for * finishes. * * It carries the task id, because without one the model cannot say which of * five workers this was, and it carries the output, because a notification * that only says "done" forces exactly the follow-up call this mechanism * exists to remove. Long output is truncated with the task id repeated in * the truncation notice, so the full text stays one `wait_for_task` away and * the model knows which id to ask for — that tool takes a `task_id` and * returns immediately for a task that has already finished, where the * listing takes only a state filter and could not have been followed. */ export declare function formatCompletionNotification(handles: readonly TaskHandle[]): string; //# sourceMappingURL=completion-inbox.d.ts.map