/** * The lifecycle seam between a connector's hooks and the AG-UI emitter. * * **This file exists so the cutover is a SWAP, not a swap plus a new lifecycle.** `TranscriptMirror` * presents a synchronous `adopt(path)` / `flush(path)` pair to the hook relay, because that is what * a hook can call: the relay must reply immediately and cannot await a publish. {@link AguiEmitter} * is the opposite shape: `start()` is async (it resolves the channel, runs the replica preflight, * and settles any frame the WAL left pending) and `pump()` is async. Something has to hold the * async thing behind the sync surface, and if that something is written in the same commit that * deletes the mirror, then the irreversible step also carries an untested lifecycle. * * So it lands first, against the shape the mirror already proved, and the cutover becomes a * substitution of one object for another behind an unchanged call pattern. * * **WHY LAZY.** The emitter cannot be built at construction time. Its {@link DurableSource} is the * session's transcript, whose path the connector does not know until a hook hands it over, and * `start()` reaches the broker — work that must not run for a session that never emits. First * `adopt` is the earliest moment the emitter is both constructible and needed. */ import type { AguiEmitter } from "./agui.js"; /** * Holds at most one {@link AguiEmitter}, started on first adopt. * * `T` is the connector's source record type; the holder never inspects one. It owns four things — * WHEN the emitter starts, that it starts AT MOST ONCE, that everything touching it is serialized, * and carrying the first adopt's opaque context to the connector-owned factory — and deliberately * owns no policy about what that context, an event, or a frame means. */ export declare class AguiEmitterHolder { private readonly startEmitter; private readonly onError; private readonly onRunClosed?; private emitter?; /** The path this holder is BOUND to. Set before `starting`, so a second path is refused even * while the first start is still in flight. */ private boundPath?; /** * Terminal. Once set, this holder never starts, pumps, or reports success again. * * It does not retry, and that is a decision rather than an omission. A retry on the next hook * would re-run that preflight and a WAL recovery against a stream this holder has already * failed to establish itself on, on a timer set by how often the user happens to type. The * emitter's own answer to an uncertain publish is to halt rather than to limp, and a holder that * quietly reconnected underneath it would reintroduce, one layer up, exactly the silence the * emitter refuses. */ private dead?; /** * Terminal, and NOT a failure. Set by {@link close} when the owner releases this holder: it never * starts, pumps or closes a run again, `failure` stays empty, and `onError` is not called. Kept * apart from `dead` for exactly that reason: a release read as an error puts a line in the log * that tells an operator the event plane broke when it was retired on purpose. */ private shut; /** ALL mutation runs on this chain. Hook events arrive concurrently on the control socket, so * without it two flushes could read the source at the same cursor. */ private chain; /** * @param startEmitter Builds and starts the emitter for an adopted path. `context` is the opaque * value supplied with the FIRST adopt (if any); the holder only carries it, while the connector * decides what it means. Injected rather than assembled here: the WAL location, the source, and * the record mapper are all connector decisions, and a holder that made them would be a second * place they are decided. * @param onError Where a failure goes. Required, and not defaulted to a swallow: this class runs * behind a hook that must not throw, so the only way a failure reaches a human is if the caller * is made to say where it goes. * @param onRunClosed Told which run {@link closeRun} closed, so the connector's own mapper can * forget the run it will no longer attribute records to. Optional, and it is NOT the holder * doing the forgetting: which state a record mapper keeps is the connector's business, and a * holder that reached into it would be a second place that decides what a run is. Without it a * mapper that still believes the run is open would emit under a `runId` the published stream has * already closed, and the emitter would refuse the batch. */ constructor(startEmitter: (path: string, context: StartContext | undefined) => Promise>, onError: (e: Error) => void, onRunClosed?: ((runId: string) => void) | undefined); /** True once an emitter is running here. False while a start is still in flight — it reports what * IS, never what is about to be. */ get running(): boolean; /** The failure that killed this holder, if one did. */ get failure(): Error | undefined; /** True once {@link close} has run. Distinct from {@link failure}: released, not broken. */ get closed(): boolean; /** * Release this holder: refuse every later adopt, flush and run-close, then await what is already * queued. * * THE REFUSAL IS TAKEN SYNCHRONOUSLY, before anything is awaited, so a hook arriving while the * settle is still outstanding is refused rather than racing it. Nothing is cancelled: work already * on the chain runs to completion, so a caller that cannot wait unboundedly has to bound this * itself, exactly as it bounds {@link settled}. * * WHAT IT DOES NOT RELEASE, because those lifetimes are not this object's to decide. The * write-ahead log and the principal lock are created by the `startEmitter` its owner injected: the * log outlives this holder whenever a later start could still recover from it, and the lock * outlives it because a lock is per principal while a holder is per thread. A holder that removed * either would be deciding a policy it cannot see. It closes the seam; the owner closes what the * seam was holding open. * * AND IT REFUSES AT THE DOOR RATHER THAN AT THE WORK, which is the distinction that keeps it from * quietly cancelling. The refusal is in {@link enqueue}, so a call arriving after close never gets * onto the chain, while a flush already on it still reads its source and publishes. A refusal * placed on the far side would have turned every abandoned drain into a silently dropped one, and * an abandoned drain is uncancelled by design: its owner stopped waiting, it did not stop. */ close(): Promise; /** The path this holder bound to on first adopt, if it has adopted. */ get path(): string | undefined; /** * Adopt a transcript path, starting the emitter if this is the first one. * * Synchronous and non-throwing by contract, because a hook calls it. The work lands on the chain. */ adopt(path: unknown, context?: StartContext): void; /** * Adopt if necessary, then drain the source into frames. * * Same contract as {@link adopt}: synchronous, non-throwing, work on the chain. */ flush(path: unknown): void; /** * Start at most once, bind the path once. * * Returns `undefined` when there is nothing to run against — a dead holder or a path this holder * cannot take — rather than throwing, so a caller cannot mistake "no emitter" for "pumped". */ /** * Close the open run at a turn boundary the record stream cannot see. * * Same contract as {@link adopt} and {@link flush}: synchronous, non-throwing, work on the chain, * because a lifecycle hook calls it and a hook must not be made to wait or to fail. * * It deliberately does NOT start an emitter. A session that never adopted a transcript has nothing * open and nothing to close, and starting one here would reach the broker on the way OUT of a * turn that published nothing. * * **`error` CLOSES THE RUN WITH `RUN_ERROR` INSTEAD**, for a turn the connector knows FAILED. It * is one parameter on the one close rather than a second method, because `RUN_ERROR` closes a run * by itself and a run must never carry both terminals: one call builds one terminal, and the * caller cannot ask for the other afterwards because the run is no longer open. This holder does * not decide what counts as a failure — that is a claim about a harness, and it belongs at the * connector's own mapping site where the harness's record is in hand. */ closeRun(timestamp: number, error?: { message: string; code?: string; }): void; private ensureStarted; private die; private enqueue; /** Await the queued work. For callers that need a settled point — a shutdown, or a cell. */ settled(): Promise; } //# sourceMappingURL=agui-holder.d.ts.map