import type { AgentRuntime } from '../providers/contract.js'; import type { WakeQueueService } from '../inbox/wake-queue.service.js'; import type { ItemStopReason, RuntimeWorkerConfig, RuntimeItemContext } from './types.js'; import type { AgentRuntimeHandleSnapshot } from '../../shared/snapshot.js'; import { TeamRunLimiter } from './team-run-limiter.js'; interface AgentRuntimeWorkerOptions extends RuntimeWorkerConfig { agentRuntime: AgentRuntime; idleTimeoutMs?: number; /** Test seam: inject a deferred restart-drain probe. */ isRestartDrainActive?: () => Promise; onItemStarted?: (context: RuntimeItemContext) => Promise; onItemFollowupAppended?: (activeContext: RuntimeItemContext, context: RuntimeItemContext) => Promise; onItemSettled?: (context: RuntimeItemContext) => Promise; pollIntervalMs?: number; queue: WakeQueueService; workerIsAlive?: (workerId: string) => boolean; workerId?: string; } export interface AgentRuntimeWorkerCloseOptions { abortReason?: ItemStopReason; drainActive?: boolean; forceAfterMs?: number; } export declare class AgentRuntimeWorker { private readonly options; private readonly logger; private readonly runLimiter; private readonly workerIsAlive; private readonly workerId; private readonly idleTimeoutMs; private readonly queue; private readonly runtimeBridge; private activeItem?; private activeDrain?; private closing; private intakePaused; private pendingWake; private pollTimer?; private unsubscribeWake?; /** * In-flight pre-drain journal backfill only. Not a lifetime success cache: * every drain re-runs (deduped observes are cheap) so transient top-level * queue/store failures are retried, and pre-journal active wakes are covered * without treating once-only startup as the sole correctness boundary. * Concurrent callers share one in-flight promise. */ private journalBackfillInFlight?; constructor(options: AgentRuntimeWorkerOptions, logger?: Pick, runLimiter?: TeamRunLimiter); drainOnce(): Promise; private drainLoop; /** * Pre-drain migration. Coalesces concurrent callers; always clears after * settle so the next drain retries after a transient top-level failure * (and re-scans active wakes that arrived mid-flight). */ private ensureJournalBackfill; start(): NodeJS.Timeout; isActive(): boolean; isProviderQuiescent(): boolean | undefined; setIntakePaused(paused: boolean): void; waitForProviderQuiescent(signal?: AbortSignal): Promise; health(): AgentRuntimeHandleSnapshot; private tick; close(options?: AgentRuntimeWorkerCloseOptions): Promise; private runOne; private readDrainActive; private takeNextRunnable; /** * @returns `'ran'` when a provider turn ran; `'settled'` when the item was * completed without a turn; `'deferred'` when the item was requeued without * tombstone (e.g. cursor-delivery prepare fail-closed) — stop this drain * cycle to avoid a hot reclaim loop. */ private processClaimedItem; private registerActiveItem; private recordRestartResumeActivity; private recoverCorruptProviderSession; private recordMemoryCoherenceCompleted; private memoryCoherenceDigest; private recordMemoryCoherenceFailed; private releaseActiveItem; private notifyItemStarted; private notifyItemFollowupAppended; private notifySettledItems; private settleAbortedItem; } export {}; //# sourceMappingURL=runtime-worker.d.ts.map