import type { OperationRunnerConfig } from '../types/operations.js'; import type { Principal, RunEventLog, RunRecord, ServerAssistant, ServerStateStore, ThreadBusyPolicy, ThreadRecord } from '../types/server.js'; /** What a run needs to start. */ export interface StartRunOptions { /** The assistant to run. */ assistant: string; /** What it runs on. */ input?: unknown; /** The thread it belongs to. Without one the run is stateless. */ threadId?: string; /** Answers an interrupt instead of starting a new input. */ resume?: unknown; /** Who asked. */ principal?: Principal; /** Replays the existing run when a run with this key was already accepted. */ idempotencyKey?: string; /** Application data recorded on the run. */ metadata?: Record; /** What to do when the thread is already running something. Defaults to the server's policy. */ onBusy?: ThreadBusyPolicy; } /** Options for the run manager. */ export interface RunManagerOptions { /** The assistants that can be run, by id. */ assistants: Record; /** Where threads, runs, and cron jobs are recorded. Defaults to memory. */ state?: ServerStateStore; /** Where run events are kept for resumable streaming. Defaults to memory. */ events?: RunEventLog; /** Configures the operation runner underneath: store, dispatcher, retries, leases, webhooks. */ operations?: OperationRunnerConfig; /** What to do when a thread is already running something. Defaults to `reject`. */ onBusy?: ThreadBusyPolicy; /** How long a run may take before it expires, in milliseconds. */ runTimeoutMs?: number; /** How long `enqueue` waits for the run in flight, in milliseconds. Defaults to 30 seconds. */ queueTimeoutMs?: number; /** Receives errors that must not fail a request, such as a failed event append. */ onError?: (error: unknown, context: { runId?: string; threadId?: string; }) => void; /** Replaces the system clock, for tests. */ now?: () => Date; } /** * Runs assistants, records threads and runs, and keeps the event log a client streams from. * * Every run is a durable operation, so a run that outlives the request that started it, a worker * that dies mid-run, a duplicate submission, and a cancellation from another replica are all handled * by the operation runner rather than by anything here. What this adds is the thread: which run owns * it, what happens when a second arrives, and where the events go. */ export declare class RunManager { private readonly options; private readonly runner; /** Handles of runs this worker is executing, so a cancellation aborts them at once. */ private readonly local; private readonly state; private readonly events; private readonly now; constructor(options: RunManagerOptions); /** The event log, which the event-stream route reads. */ get eventLog(): RunEventLog; /** The id this worker writes into operation leases. */ get workerId(): string; /** The assistants that can be run. */ get assistants(): Record; /** Looks an assistant up, or refuses the request. */ assistant(id: string): ServerAssistant; /** Creates a thread for an assistant. */ createThread(options: { assistant: string; threadId?: string; principal?: Principal; metadata?: Record; }): Promise; /** Reads a thread, refusing one that belongs to another tenant. */ thread(threadId: string, principal?: Principal): Promise; /** Threads of a tenant, newest written first. */ threads(principal?: Principal, limit?: number): Promise; /** Deletes a thread's record. The assistant's own checkpoints are left alone. */ deleteThread(threadId: string, principal?: Principal): Promise; /** The assistant's state for a thread, when the assistant reports one. */ threadState(threadId: string, principal?: Principal): Promise; /** Reads a run, refusing one that belongs to another tenant. */ run(runId: string, principal?: Principal): Promise; /** Runs of a tenant, newest written first, optionally for one thread. */ runs(principal?: Principal, options?: { threadId?: string; limit?: number; }): Promise; /** * Accepts a run and starts it in the background. * * Returns as soon as the run is recorded, so the caller can stream its events or hand back its id. * A thread already running something is handled by the busy policy before anything is submitted. */ start(options: StartRunOptions): Promise; /** Cancels a run, including one another replica is executing. */ cancel(runId: string, principal?: Principal, reason?: string): Promise; /** * Re-runs whatever this worker can claim from the store, after a restart or another worker's * crash. The operation runner decides what is claimable; this only supplies the executor. */ recover(limit?: number): Promise; private execute; /** Applies the busy policy, returning the thread once it is free to run something new. */ private settleBusyThread; /** Waits for a run to finish, for the `enqueue` policy. */ private waitForRun; private releaseThread; private record; private settle; private report; }