{"version":3,"file":"lane.d.ts","sourceRoot":"","sources":["../../../src/harness/runtime/lane.ts"],"names":[],"mappings":"AAAA,OAAO,EACN,KAAK,GAAG,EACR,KAAK,YAAY,EACjB,KAAK,KAAK,EACV,KAAK,MAAM,EAEX,KAAK,KAAK,EACV,MAAM,uBAAuB,CAAC;AAC/B,OAAO,KAAK,EAAE,YAAY,EAAE,aAAa,EAAE,MAAM,gBAAgB,CAAC;AAClE,OAAO,KAAK,EACX,kBAAkB,EAClB,WAAW,EACX,SAAS,EACT,kBAAkB,EAClB,gBAAgB,EAChB,YAAY,EACZ,WAAW,EACX,YAAY,EAEZ,iBAAiB,EACjB,YAAY,EACZ,aAAa,EACb,gBAAgB,EAEhB,wBAAwB,EACxB,gBAAgB,EAChB,WAAW,EACX,iBAAiB,EACjB,YAAY,EACZ,SAAS,EAET,WAAW,EACX,MAAM,qBAAqB,CAAC;AAG7B,OAAO,EAAoB,KAAK,OAAO,EAAE,MAAM,eAAe,CAAC;AAE/D,OAAO,KAAK,EAAE,YAAY,EAAE,MAAM,aAAa,CAAC;AAEhD,OAAO,EAUN,iBAAiB,EAKjB,MAAM,cAAc,CAAC;AAGtB,OAAO,KAAK,EACX,UAAU,EACV,KAAK,EAGL,SAAS,EAIT,aAAa,EACb,qBAAqB,EACrB,cAAc,EAGd,OAAO,EACP,aAAa,EAIb,MAAM,qBAAqB,CAAC;AAoB7B,OAAO,EACN,KAAK,MAAM,EACX,KAAK,uBAAuB,EAC5B,KAAK,EACL,KAAK,WAAW,EAChB,KAAK,SAAS,EACd,KAAK,gBAAgB,EACrB,MAAM,YAAY,CAAC;AAEpB,KAAK,SAAS,GAAG,CAAC,MAAM,EAAE,SAAS,YAAY,EAAE,EAAE,OAAO,EAAE,OAAO,KAAK,OAAO,CAAC,IAAI,CAAC,CAAC;AACtF,KAAK,YAAY,GAAG,CAAC,CAAC,EACrB,QAAQ,EAAE,CAAC,EACX,MAAM,EAAE,CAAC,KAAK,EAAE,YAAY,KAAK,OAAO,EACxC,OAAO,EAAE,OAAO,EAChB,UAAU,EAAE,CAAC,OAAO,EAAE,OAAO,EAAE,YAAY,EAAE,MAAM,IAAI,KAAK,OAAO,CAAC,CAAC,CAAC,KAClE,WAAW,CAAC,CAAC,CAAC,CAAC;AACpB,KAAK,YAAY,GAAG,CAAC,KAAK,EAAE,OAAO,EAAE,OAAO,EAAE,OAAO,KAAK,KAAK,CAAC;AA2GhE,qDAAqD;AACrD,qBAAa,IAAI,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,CAAE,YAAW,SAAS;IAC1E,QAAQ,CAAC,IAAI,EAAE,MAAM,CAAC;IACtB,QAAQ,CAAC,OAAO,EAAE,OAAO,CAAC;IAC1B,QAAQ,CAAC,MAAM,EAAE,MAAM,CAAC;IACxB,QAAQ,CAAC,KAAK,EAAE,YAAY,CAAC;IAC7B,QAAQ,CAAC,SAAS,EAAE,SAAS,CAAC;IAC9B,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAe;IACvC,OAAO,CAAC,QAAQ,CAAC,YAAY,CAAe;IAC5C,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAyB;IAChD,OAAO,CAAC,WAAW,CAAgB;IACnC,OAAO,CAAC,kBAAkB,CAAa;IACvC,OAAO,CAAC,SAAS,CAA4B;IAC7C,qHAAqH;IACrH,WAAW,EAAE,KAAK,GAAG,SAAS,CAAC;IAC/B,iFAAiF;IACjF,KAAK,EAAE,SAAS,CAAC;IACjB,WAAW,EAAE,KAAK,GAAG,SAAS,CAAC;IAE/B,YACC,IAAI,EAAE,MAAM,EACZ,OAAO,EAAE,OAAO,EAChB,MAAM,EAAE,MAAM,EACd,KAAK,EAAE,YAAY,EACnB,KAAK,EAAE,SAAS,EAChB,OAAO,EAAE,YAAY,EACrB,SAAS,EAAE,SAAS,EACpB,YAAY,EAAE,YAAY,EAC1B,UAAU,EAAE,MAAM,MAAM,CAAC,QAAQ,CAAC,EAgBlC;IAEK,QAAQ,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,MAAM,GAAG,IAAI,CAAC,CAGxD;IAEK,SAAS,CAAC,WAAW,EAAE,MAAM,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,qBAAqB,GAAG,SAAS,CAAC,CAGjG;IAED,UAAU,IAAI,MAAM,CAAC,QAAQ,CAAC,CAE7B;IAED,QAAQ,CAAC,QAAQ,EAAE,MAAM,EAAE,kBAAkB,EAAE,MAAM,GAAG,IAAI,EAAE,eAAe,EAAE,MAAM,GAAG,IAAI,GAAG,iBAAiB,CAQ/G;YAmBa,QAAQ;IAqBhB,OAAO,CAAC,OAAO,EACpB,IAAI,EAAE,CAAC,KAAK,EAAE,SAAS,EAAE,MAAM,EAAE,aAAa,KAAK,WAAW,CAAC,OAAO,CAAC,GAAG,OAAO,CAAC,WAAW,CAAC,OAAO,CAAC,CAAC,EACvG,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,OAAO,CAAC,CAkDlB;IAED;;;;OAIG;IACH,eAAe,CAAC,MAAM,SAAS,cAAc,EAAE,OAAO,EACrD,WAAW,EAAE,MAAM,EACnB,IAAI,EAAE,CACL,KAAK,EAAE,SAAS,EAChB,OAAO,EAAE,MAAM,EACf,IAAI,EAAE,aAAa,EACnB,MAAM,EAAE,aAAa,KACjB,gBAAgB,CAAC,OAAO,CAAC,GAAG,OAAO,CAAC,gBAAgB,CAAC,OAAO,CAAC,CAAC,EACnE,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,OAAO,CAAC,CAgDlB;IAED;;;OAGG;IACH,iBAAiB,CAAC,MAAM,SAAS,cAAc,EAAE,OAAO,EACvD,UAAU,EAAE,MAAM,EAClB,IAAI,EAAE,CACL,KAAK,EAAE,SAAS,EAChB,OAAO,EAAE,MAAM,EACf,IAAI,EAAE,aAAa,EACnB,MAAM,EAAE,aAAa,KACjB,gBAAgB,CAAC,OAAO,CAAC,GAAG,OAAO,CAAC,gBAAgB,CAAC,OAAO,CAAC,CAAC,EACnE,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,uBAAuB,CAAC,OAAO,CAAC,CAAC,CAkB3C;IAEK,MAAM,CAAC,OAAO,EAAE,gBAAgB,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,wBAAwB,CAAC,CAgB3F;YAEa,SAAS;IAqLvB,OAAO,CAAC,gBAAgB;YAmFV,gBAAgB;IA+JxB,KAAK,CAAC,OAAO,EAAE,YAAY,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CAkFzE;IAED,gGAAgG;IAC1F,qBAAqB,CAAC,WAAW,EAAE,MAAM,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,kBAAkB,CAAC,CA8F9F;IAED,YAAY,CAAC,WAAW,EAAE,MAAM,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,kBAAkB,CAAC,CAE/E;IAED,gBAAgB,CAAC,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,iBAAiB,CAAC,CAyB7D;IAEK,MAAM,CACX,GAAG,IAAI,EACJ,CAAC,IAAI,EAAE,MAAM,EAAE,MAAM,EAAE,YAAY,EAAE,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,CAAC,GACpE,CAAC,OAAO,EAAE,YAAY,GAAG,YAAY,EAAE,EAAE,OAAO,EAAE,OAAO,CAAC,GAC3D,OAAO,CAAC,SAAS,CAAC,CAQpB;IAED,KAAK,CAAC,IAAI,EAAE,MAAM,EAAE,sBAAsB,EAAE,MAAM,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,SAAS,CAAC,CAKpG;IAED,kBAAkB,CAAC,IAAI,EAAE,MAAM,EAAE,IAAI,EAAE,MAAM,EAAE,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,SAAS,CAAC,CAEjG;YAEa,eAAe;IA0CvB,OAAO,CAAC,OAAO,EAAE;QAAE,kBAAkB,CAAC,EAAE,MAAM,CAAA;KAAE,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,gBAAgB,CAAC,CA8B/G;IAEK,YAAY,CACjB,QAAQ,EAAE,MAAM,GAAG,IAAI,EACvB,OAAO,EAAE,OAAO,CAAC,gBAAgB,EAAE;QAAE,IAAI,EAAE,YAAY,CAAA;KAAE,CAAC,CAAC,SAAS,CAAC,EACrE,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,gBAAgB,CAAC,CA4B3B;YAEa,wBAAwB;YAmBxB,uBAAuB;IA0C/B,MAAM,CAAC,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,YAAY,CAAC,CA4CpD;IAEK,KAAK,CAAC,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CA2ClD;IAED,KAAK,CAAC,OAAO,EAAE,MAAM,GAAG,YAAY,EAAE,MAAM,EAAE,YAAY,EAAE,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CAEhH;IAED,QAAQ,CACP,OAAO,EAAE,MAAM,GAAG,YAAY,EAC9B,MAAM,EAAE,YAAY,EAAE,GAAG,SAAS,EAClC,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,WAAW,CAAC,CAEtB;IAED,OAAO,CAAC,OAAO,EAAE,MAAM,GAAG,YAAY,EAAE,MAAM,EAAE,YAAY,EAAE,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CAElH;YAEa,OAAO;IAoFf,YAAY,CAAC,OAAO,EAAE,MAAM,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,kBAAkB,CAAC,CAiCjF;IAEK,WAAW,CAChB,KAAK,EAAE,KAAK,EACZ,OAAO,EAAE;QAAE,OAAO,CAAC,EAAE,MAAM,CAAC;QAAC,OAAO,CAAC,EAAE,SAAS,CAAA;KAAE,GAAG,SAAS,EAC9D,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,iBAAiB,CAAC,CA6B5B;IAEK,WAAW,CAAC,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAkBjD;IAEK,WAAW,CAAC,QAAQ,EAAE,CAAC,OAAO,EAAE,OAAO,KAAK,IAAI,GAAG,OAAO,CAAC,IAAI,CAAC,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAsCvG;IAEK,QAAQ,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,KAAK,CAAC,GAAG,CAAC,GAAG,SAAS,CAAC,CAGjE;IAED,QAAQ,CAAC,KAAK,EAAE,aAAa,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAc9D;IAEK,gBAAgB,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,aAAa,CAAC,CAGhE;IAED,gBAAgB,CAAC,aAAa,EAAE,aAAa,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAW9E;IAEK,cAAc,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,MAAM,EAAE,CAAC,CAGzD;IAED,cAAc,CAAC,eAAe,EAAE,MAAM,EAAE,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAWzE;IAED,KAAK,CAAC,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,YAAY,CAAC,CAAC,CAqB1D;YAEa,mBAAmB;YA2JnB,gBAAgB;IAiBxB,WAAW,CAAC,KAAK,EAAE,UAAU,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,KAAK,EAAE,CAAC,CAOnF;IAEK,SAAS,CAAC,KAAK,EAAE,UAAU,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,KAAK,GAAG,SAAS,CAAC,CAK3F;IAED,aAAa,CAAC,OAAO,EAAE,YAAY,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,MAAM,CAAC,CAEtE;IAED,iBAAiB,CAAC,UAAU,EAAE,MAAM,EAAE,IAAI,EAAE,SAAS,GAAG,SAAS,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,MAAM,CAAC,CAEpG;IAED,OAAO,CAAC,MAAM;IAqEd,IAAI,CAAC,KAAK,EAAE,KAAK,GAAG,OAAO,CAAC,IAAI,CAAC,CAKhC;IAED,OAAO,CAAC,iBAAiB;IASzB,UAAU,IAAI,IAAI,CAEjB;CACD","sourcesContent":["import {\n\ttype Api,\n\ttype ImageContent,\n\ttype Model,\n\ttype Models,\n\treduceAssistantMessageFrames,\n\ttype Usage,\n} from \"@earendil-works/pi-ai\";\nimport type { AgentMessage, ThinkingLevel } from \"../../types.ts\";\nimport type {\n\tAbortRequestResult,\n\tAbortResult,\n\tAgentLane,\n\tCancelQueuedResult,\n\tCompactionResult,\n\tDriveOptions,\n\tDriveResult,\n\tHarnessEvent,\n\tLaneConfigEventPayload,\n\tLaneExecutionInfo,\n\tLaneSnapshot,\n\tModelIdentity,\n\tNavigationResult,\n\tOperationAdmission,\n\tOperationAdmissionResult,\n\tOperationRequest,\n\tQueueResult,\n\tRecordUsageResult,\n\tResumeResult,\n\tRunResult,\n\tSuspendedRun,\n\tWatchHandle,\n} from \"../agent-harness.ts\";\nimport { type BranchPreparation, prepareBranchEntries } from \"../compaction/branch-summarization.ts\";\nimport { prepareCompaction } from \"../compaction/compaction.ts\";\nimport { awaitWithContext, type Context } from \"../context.ts\";\nimport { toolResultFromMessage } from \"../execution/tools.ts\";\nimport type { HookRegistry } from \"../hooks.ts\";\nimport { formatPromptTemplateInvocation } from \"../prompt-templates.ts\";\nimport {\n\tClosed,\n\tHarnessClosed,\n\tHarnessFault,\n\tInvalidMessage,\n\tInvalidNavigation,\n\tLaneBusy,\n\tNoActiveOperation,\n\tNothingToCompact,\n\tNothingToResume,\n\tOperationMismatch,\n\tResult,\n\tUnknownSkill,\n\tUnknownTarget,\n\tUnknownTemplate,\n} from \"../result.ts\";\nimport { insertEntry, insertUsage } from \"../session/commit.ts\";\nimport { SessionInvariantError, SessionPendingAssistantMessageError } from \"../session/session.ts\";\nimport type {\n\tBranchScan,\n\tEntry,\n\tInboxItem,\n\tInboxItemKind,\n\tJsonValue,\n\tNavigationReadyToCommitOperation,\n\tNewEntry,\n\tOperation,\n\tOperationMeta,\n\tOperationResultRecord,\n\tOperationState,\n\tPendingEntry,\n\tRunSettings,\n\tSession,\n\tSessionReader,\n\tStartingOperation,\n\tSummaryDecidingOperation,\n\tWrite,\n} from \"../session/types.ts\";\nimport {\n\tbranchTip,\n\tdeleteValue,\n\tlaneConfig,\n\tlaneState as laneStateValue,\n\toperationMeta as operationMetaValue,\n\toperationPreparation,\n\toperationResult as operationResultValue,\n\toperationState as operationStateValue,\n\toperationToolArgs,\n\tpendingEntry,\n\tpendingToolOutput,\n\tsetValue,\n} from \"../session/values.ts\";\nimport { formatSkillInvocation } from \"../skills.ts\";\nimport { durableBranchPreparation, durableCompactionPreparation } from \"./drive/structural.ts\";\nimport { driveOperation } from \"./drive.ts\";\nimport { readAssistantFrames } from \"./progress.ts\";\nimport { chainEntries, committedEntryEvents, readLaneQueues } from \"./transcript.ts\";\nimport {\n\ttype Config,\n\ttype ContinueOperationResult,\n\tDrive,\n\ttype LaneCommand,\n\ttype LaneState,\n\ttype OperationCommand,\n} from \"./types.ts\";\n\ntype EmitBatch = (events: readonly HarnessEvent[], context: Context) => Promise<void>;\ntype WatchHandler = <T>(\n\tsnapshot: T,\n\tfilter: (event: HarnessEvent) => boolean,\n\tcontext: Context,\n\tresnapshot: (context: Context, markBoundary: () => void) => Promise<T>,\n) => WatchHandle<T>;\ntype FaultHandler = (cause: unknown, context: Context) => Error;\n\ntype LaneCommandOutcome<TResult> =\n\t| { kind: \"return\"; result: TResult; delivery?: Promise<void> }\n\t| { kind: \"reject\"; error: Error }\n\t| { kind: \"idle_blocked\"; owner: Promise<void>; change: Promise<void> };\n\ntype IdleObservation = { kind: \"idle\" } | { kind: \"wait\"; drive: Drive | undefined; change: Promise<void> };\ntype IdleClaimObservation = { kind: \"claimed\" } | { kind: \"wait\"; drive: Drive | undefined; change: Promise<void> };\n\ntype DriveClaim =\n\t| { kind: \"observe\"; drive: Drive; installed: boolean }\n\t| { kind: \"occupied\"; drive: Drive }\n\t| { kind: \"settled\"; outcome: OperationResultRecord }\n\t| { kind: \"mismatch\"; error: OperationMismatch };\n\nfunction isPromiseLike(value: unknown): value is PromiseLike<unknown> {\n\tif (value === null || (typeof value !== \"object\" && typeof value !== \"function\")) return false;\n\treturn \"then\" in value && typeof value.then === \"function\";\n}\n\nfunction inboxItems(inbox: readonly InboxItem[], kind: InboxItemKind): InboxItem[] {\n\treturn inbox.filter((item) => item.kind === kind);\n}\n\nfunction withoutInboxItems(inbox: readonly InboxItem[], removed: readonly InboxItem[]): InboxItem[] {\n\tconst removedIds = new Set(removed.map((item) => item.entryId));\n\treturn inbox.filter((item) => !removedIds.has(item.entryId));\n}\n\nfunction selectAcceptedInbox(\n\tinbox: readonly InboxItem[],\n\tsteeringMode: \"all\" | \"one-at-a-time\",\n\tfollowUpMode: \"all\" | \"one-at-a-time\",\n): { selected: InboxItem[]; remainder: InboxItem[] } {\n\tlet steerTaken = false;\n\tlet followUpTaken = false;\n\tconst selected: InboxItem[] = [];\n\tconst remainder: InboxItem[] = [];\n\tfor (const item of inbox) {\n\t\tconst eligible =\n\t\t\titem.kind === \"write\" ||\n\t\t\titem.kind === \"nextRun\" ||\n\t\t\t(item.kind === \"steer\" && (steeringMode === \"all\" || !steerTaken)) ||\n\t\t\t(item.kind === \"followUp\" && (followUpMode === \"all\" || !followUpTaken));\n\t\tif (eligible) {\n\t\t\tselected.push(item);\n\t\t\tif (item.kind === \"steer\") steerTaken = true;\n\t\t\tif (item.kind === \"followUp\") followUpTaken = true;\n\t\t} else {\n\t\t\tremainder.push(item);\n\t\t}\n\t}\n\treturn { selected, remainder };\n}\n\nfunction capturedSettings<TContext extends object | undefined>(config: Config<TContext>): RunSettings {\n\treturn {\n\t\tcompaction: config.compaction,\n\t\tsteeringMode: config.steeringMode,\n\t\tfollowUpMode: config.followUpMode,\n\t\ttoolExecution: config.toolExecution,\n\t};\n}\n\nfunction durableLaneState(\n\tstate: LaneState,\n\tcurrentOperationId: string | null,\n\tinbox: InboxItem[] = state.inbox,\n\tlastOperationId: string | null = state.lastOperationId,\n) {\n\treturn { currentOperationId, lastOperationId, inbox };\n}\n\nfunction pendingEntryWrite(entryId: string, pending: PendingEntry): NewEntry {\n\treturn pending.type === \"message\"\n\t\t? { id: entryId, parentId: null, type: \"message\", message: pending.payload }\n\t\t: {\n\t\t\t\tid: entryId,\n\t\t\t\tparentId: null,\n\t\t\t\ttype: \"custom\",\n\t\t\t\tcustomType: pending.customType,\n\t\t\t\t...(pending.payload === undefined ? {} : { data: pending.payload }),\n\t\t\t};\n}\n\nfunction capturedModel(operation: Operation): ModelIdentity | undefined {\n\tconst { state } = operation;\n\tswitch (state.at) {\n\t\tcase \"assistant.ready\":\n\t\tcase \"assistant.effect_pending\":\n\t\tcase \"assistant.retry_wait\":\n\t\t\treturn state.generationContext.configuration.model;\n\t\tcase \"tools\":\n\t\t\treturn state.batch.configuration.model;\n\t\tcase \"deferred.suspended\":\n\t\tcase \"deferred.effect_pending\":\n\t\t\treturn state.configuration.model;\n\t\tcase \"summary.ready\":\n\t\tcase \"summary.effect_pending\":\n\t\tcase \"summary.retry_wait\":\n\t\t\treturn state.summaryContext.configuration.model;\n\t\tdefault:\n\t\t\treturn undefined;\n\t}\n}\n\n/** Runtime implementation of one configured lane. */\nexport class Lane<TContext extends object | undefined> implements AgentLane {\n\treadonly name: string;\n\treadonly session: Session;\n\treadonly models: Models;\n\treadonly hooks: HookRegistry;\n\treadonly emitBatch: EmitBatch;\n\tprivate readonly onFault: FaultHandler;\n\tprivate readonly installWatch: WatchHandler;\n\tprivate readonly config: () => Config<TContext>;\n\tprivate stateChange: Promise<void>;\n\tprivate resolveStateChange: () => void;\n\tprivate idleOwner: Promise<void> | undefined;\n\t/** Package-internal drive owner. Public only because deterministic procedure tests install exact owners directly. */\n\tactiveDrive: Drive | undefined;\n\t/** Authoritative live control projection while this harness owns the Session. */\n\tstate: LaneState;\n\tclosedError: Error | undefined;\n\n\tconstructor(\n\t\tname: string,\n\t\tsession: Session,\n\t\tmodels: Models,\n\t\thooks: HookRegistry,\n\t\tstate: LaneState,\n\t\tonFault: FaultHandler,\n\t\temitBatch: EmitBatch,\n\t\tinstallWatch: WatchHandler,\n\t\treadConfig: () => Config<TContext>,\n\t) {\n\t\tthis.session = session;\n\t\tthis.models = models;\n\t\tthis.hooks = hooks;\n\t\tthis.name = name;\n\t\tthis.state = state;\n\t\tthis.onFault = onFault;\n\t\tthis.emitBatch = emitBatch;\n\t\tthis.installWatch = installWatch;\n\t\tthis.config = readConfig;\n\t\tlet resolveStateChange!: () => void;\n\t\tthis.stateChange = new Promise<void>((resolve) => {\n\t\t\tresolveStateChange = resolve;\n\t\t});\n\t\tthis.resolveStateChange = resolveStateChange;\n\t}\n\n\tasync getTipId(_context: Context): Promise<string | null> {\n\t\tthis.assertOpen();\n\t\treturn this.state.tipId;\n\t}\n\n\tasync getResult(operationId: string, context: Context): Promise<OperationResultRecord | undefined> {\n\t\tthis.assertOpen();\n\t\treturn (await this.session.getValue(operationResultValue(operationId), context))?.value;\n\t}\n\n\treadConfig(): Config<TContext> {\n\t\treturn this.config();\n\t}\n\n\tmismatch(expected: string, currentOperationId: string | null, lastOperationId: string | null): OperationMismatch {\n\t\treturn new OperationMismatch({\n\t\t\tlane: this.name,\n\t\t\texpectedOperationId: expected,\n\t\t\t...(currentOperationId === null ? {} : { currentOperationId }),\n\t\t\t...(lastOperationId === null ? {} : { lastOperationId }),\n\t\t\tmessage: `Operation ${expected} does not own lane ${JSON.stringify(this.name)}`,\n\t\t});\n\t}\n\n\t/**\n\t * Run one effect-free command on this lane's serialized mutation line. Owned `state` is authoritative; the\n\t * planner receives that state plus a read-only reader for bounded payload lookups. Committed values and owned\n\t * state are immutable snapshots: update them by replacement, never in place. State-independent input validation\n\t * belongs before `command()`, while every state-dependent decision belongs inside its planner.\n\t *\n\t * A planner may choose exactly one outcome:\n\t * - `commit` commits once, publishes `next`, then synchronously materializes the caller result from\n\t *   storage-assigned `CommitResult` metadata;\n\t * - `return` returns without a commit, boxed so a promise value is not awaited while holding the Session line;\n\t * - `reject` rejects outside the mutation/fault boundary as an expected caller error without a commit.\n\t *\n\t * Planner, commit, and materialization errors fault the harness before releasing the Session line. Close/fault gates\n\t * are checked both before queueing and when the callback starts: close-first rejects, while a callback admitted\n\t * before close may finish its commit, publish memory, and resolve without another open check. Never invoke providers,\n\t * tools, hooks, timers, event handlers, or wait for task completion here; perform those after `command()` returns.\n\t */\n\tprivate async readLane<TResult>(\n\t\tread: (state: LaneState, reader: SessionReader) => TResult | Promise<TResult>,\n\t\tcontext: Context,\n\t): Promise<TResult> {\n\t\tthis.assertOpen();\n\t\ttry {\n\t\t\treturn await this.session.mutate(async (reader) => {\n\t\t\t\tthis.assertOpen();\n\t\t\t\ttry {\n\t\t\t\t\treturn await read(this.state, reader);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tif (this.closedError !== undefined) throw this.closedError;\n\t\t\t\t\tthrow this.onFault(error, context);\n\t\t\t\t}\n\t\t\t}, context);\n\t\t} catch (error) {\n\t\t\tif (this.closedError !== undefined) throw this.closedError;\n\t\t\tthrow error;\n\t\t}\n\t}\n\n\tasync command<TResult>(\n\t\tplan: (state: LaneState, reader: SessionReader) => LaneCommand<TResult> | Promise<LaneCommand<TResult>>,\n\t\tcontext: Context,\n\t): Promise<TResult> {\n\t\tthis.assertOpen();\n\t\twhile (this.idleOwner !== undefined) {\n\t\t\tawait awaitWithContext(Promise.race([this.idleOwner, this.stateChange]), context);\n\t\t\tthis.assertOpen();\n\t\t}\n\t\tlet outcome: LaneCommandOutcome<TResult>;\n\t\ttry {\n\t\t\toutcome = await this.session.mutate(async (mutator) => {\n\t\t\t\tthis.assertOpen();\n\t\t\t\tif (this.idleOwner !== undefined) {\n\t\t\t\t\treturn { kind: \"idle_blocked\", owner: this.idleOwner, change: this.stateChange };\n\t\t\t\t}\n\t\t\t\ttry {\n\t\t\t\t\tconst decision = await plan(this.state, mutator);\n\t\t\t\t\tswitch (decision.kind) {\n\t\t\t\t\t\tcase \"return\":\n\t\t\t\t\t\t\treturn { kind: \"return\", result: decision.result };\n\t\t\t\t\t\tcase \"reject\":\n\t\t\t\t\t\t\treturn { kind: \"reject\", error: decision.error };\n\t\t\t\t\t\tcase \"commit\": {\n\t\t\t\t\t\t\tconst commit = await mutator.commit(decision.writes, context);\n\t\t\t\t\t\t\tthis.state = decision.next;\n\t\t\t\t\t\t\tthis.signalStateChange();\n\t\t\t\t\t\t\tconst result = decision.materialize(commit);\n\t\t\t\t\t\t\tif (isPromiseLike(result)) {\n\t\t\t\t\t\t\t\tthrow new TypeError(\"Lane command materialize() must be synchronous\");\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t\tconst events = decision.events?.(commit) ?? [];\n\t\t\t\t\t\t\tconst delivery = events.length === 0 ? undefined : this.emitBatch(events, context);\n\t\t\t\t\t\t\treturn { kind: \"return\", result, ...(delivery === undefined ? {} : { delivery }) };\n\t\t\t\t\t\t}\n\t\t\t\t\t}\n\t\t\t\t} catch (error) {\n\t\t\t\t\tif (this.closedError !== undefined) throw this.closedError;\n\t\t\t\t\tthrow this.onFault(error, context);\n\t\t\t\t}\n\t\t\t}, context);\n\t\t} catch (error) {\n\t\t\tif (this.closedError !== undefined) throw this.closedError;\n\t\t\tthrow error;\n\t\t}\n\t\tif (outcome.kind === \"idle_blocked\") {\n\t\t\tawait awaitWithContext(Promise.race([outcome.owner, outcome.change]), context);\n\t\t\tthis.assertOpen();\n\t\t\treturn this.command(plan, context);\n\t\t}\n\t\tif (outcome.kind === \"reject\") throw outcome.error;\n\t\tawait outcome.delivery;\n\t\treturn outcome.result;\n\t}\n\n\t/**\n\t * Run a command against the current operation even after cancellation is requested. Use this to settle admitted\n\t * effects, finish the operation, or update concurrent child state. The capability narrows the planner's state type;\n\t * the Drive continuation remains the sole top-level state writer.\n\t */\n\tsettleOperation<TState extends OperationState, TResult>(\n\t\t_capability: TState,\n\t\tplan: (\n\t\t\tstate: LaneState,\n\t\t\tcurrent: TState,\n\t\t\tmeta: OperationMeta,\n\t\t\treader: SessionReader,\n\t\t) => OperationCommand<TResult> | Promise<OperationCommand<TResult>>,\n\t\tcontext: Context,\n\t): Promise<TResult> {\n\t\treturn this.command(async (state, reader) => {\n\t\t\tconst operation = state.operation!;\n\t\t\tconst decision = await plan(state, operation.state as TState, operation.meta, reader);\n\t\t\tif (decision.kind === \"commit\") {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites: [\n\t\t\t\t\t\t...decision.writes,\n\t\t\t\t\t\tsetValue(operationStateValue(operation.meta.operationId), decision.operationState),\n\t\t\t\t\t\t...(decision.lane?.inbox === undefined\n\t\t\t\t\t\t\t? []\n\t\t\t\t\t\t\t: [\n\t\t\t\t\t\t\t\t\tsetValue(\n\t\t\t\t\t\t\t\t\t\tlaneStateValue(this.name),\n\t\t\t\t\t\t\t\t\t\tdurableLaneState(state, operation.meta.operationId, decision.lane.inbox),\n\t\t\t\t\t\t\t\t\t),\n\t\t\t\t\t\t\t\t]),\n\t\t\t\t\t],\n\t\t\t\t\tnext: {\n\t\t\t\t\t\t...state,\n\t\t\t\t\t\t...decision.lane,\n\t\t\t\t\t\toperation: { meta: operation.meta, state: decision.operationState },\n\t\t\t\t\t},\n\t\t\t\t\tmaterialize: decision.materialize,\n\t\t\t\t\t...(decision.events === undefined ? {} : { events: decision.events }),\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (decision.kind !== \"finish\") return decision;\n\t\t\tconst inbox = decision.lane?.inbox ?? state.inbox;\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\t...decision.writes,\n\t\t\t\t\tsetValue(operationResultValue(operation.meta.operationId), decision.record),\n\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, null, inbox, operation.meta.operationId)),\n\t\t\t\t],\n\t\t\t\tnext: {\n\t\t\t\t\t...state,\n\t\t\t\t\t...decision.lane,\n\t\t\t\t\tinbox,\n\t\t\t\t\tlastOperationId: operation.meta.operationId,\n\t\t\t\t\toperation: null,\n\t\t\t\t},\n\t\t\t\tmaterialize: decision.materialize,\n\t\t\t\t...(decision.events === undefined ? {} : { events: decision.events }),\n\t\t\t};\n\t\t}, context);\n\t}\n\n\t/**\n\t * Run an ordinary operation command only while durable control is running. Use this before starting new hooks,\n\t * effects, or forward progress. Returns `cancel_requested` without invoking the planner once cancellation is requested.\n\t */\n\tcontinueOperation<TState extends OperationState, TResult>(\n\t\tcapability: TState,\n\t\tplan: (\n\t\t\tstate: LaneState,\n\t\t\tcurrent: TState,\n\t\t\tmeta: OperationMeta,\n\t\t\treader: SessionReader,\n\t\t) => OperationCommand<TResult> | Promise<OperationCommand<TResult>>,\n\t\tcontext: Context,\n\t): Promise<ContinueOperationResult<TResult>> {\n\t\treturn this.settleOperation<TState, ContinueOperationResult<TResult>>(\n\t\t\tcapability,\n\t\t\tasync (state, latest, meta, reader) => {\n\t\t\t\tif (latest.control.status === \"cancel_requested\") {\n\t\t\t\t\treturn { kind: \"return\", result: { kind: \"cancel_requested\" } };\n\t\t\t\t}\n\t\t\t\tconst decision = await plan(state, latest, meta, reader);\n\t\t\t\tif (decision.kind === \"return\") {\n\t\t\t\t\treturn { kind: \"return\", result: { kind: \"result\", value: decision.result } };\n\t\t\t\t}\n\t\t\t\treturn {\n\t\t\t\t\t...decision,\n\t\t\t\t\tmaterialize: (commit) => ({ kind: \"result\", value: decision.materialize(commit) }),\n\t\t\t\t};\n\t\t\t},\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync accept(request: OperationRequest, context: Context): Promise<OperationAdmissionResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\tconst startedAt = Date.now();\n\t\tconst operationId = request.operationId ?? this.session.idGenerator.next(startedAt);\n\t\tconst acceptanceConfig = this.readConfig();\n\t\tif (request.kind === \"compaction\") {\n\t\t\treturn this.acceptCompaction(request, operationId, startedAt, acceptanceConfig, context);\n\t\t}\n\t\tif (request.kind === \"navigation\") {\n\t\t\treturn this.acceptNavigation(request, operationId, startedAt, acceptanceConfig, context);\n\t\t}\n\n\t\treturn this.acceptRun(request, operationId, startedAt, acceptanceConfig, context);\n\t}\n\n\tprivate async acceptRun(\n\t\trequest: Extract<OperationRequest, { kind: \"prompt\" | \"skill\" | \"prompt_template\" }>,\n\t\toperationId: string,\n\t\tstartedAt: number,\n\t\tacceptanceConfig: Config<TContext>,\n\t\tcontext: Context,\n\t): Promise<OperationAdmissionResult> {\n\t\tlet messages: AgentMessage[];\n\t\tswitch (request.kind) {\n\t\t\tcase \"prompt\":\n\t\t\t\tif (typeof request.prompt !== \"string\") {\n\t\t\t\t\tmessages = Array.isArray(request.prompt) ? request.prompt : [request.prompt];\n\t\t\t\t} else {\n\t\t\t\t\tconst images = request.images ?? [];\n\t\t\t\t\tmessages =\n\t\t\t\t\t\trequest.prompt.length === 0 && images.length === 0\n\t\t\t\t\t\t\t? []\n\t\t\t\t\t\t\t: [\n\t\t\t\t\t\t\t\t\t{\n\t\t\t\t\t\t\t\t\t\trole: \"user\",\n\t\t\t\t\t\t\t\t\t\tcontent: [\n\t\t\t\t\t\t\t\t\t\t\t...(request.prompt.length === 0\n\t\t\t\t\t\t\t\t\t\t\t\t? []\n\t\t\t\t\t\t\t\t\t\t\t\t: [{ type: \"text\" as const, text: request.prompt }]),\n\t\t\t\t\t\t\t\t\t\t\t...images,\n\t\t\t\t\t\t\t\t\t\t],\n\t\t\t\t\t\t\t\t\t\ttimestamp: startedAt,\n\t\t\t\t\t\t\t\t\t},\n\t\t\t\t\t\t\t\t];\n\t\t\t\t}\n\t\t\t\tbreak;\n\t\t\tcase \"skill\": {\n\t\t\t\tconst skill = acceptanceConfig.resources.skills?.find((candidate) => candidate.name === request.name);\n\t\t\t\tif (skill === undefined) {\n\t\t\t\t\treturn Result.err(new UnknownSkill({ name: request.name, message: `Unknown skill: ${request.name}` }));\n\t\t\t\t}\n\t\t\t\tmessages = [\n\t\t\t\t\t{\n\t\t\t\t\t\trole: \"user\",\n\t\t\t\t\t\tcontent: [{ type: \"text\", text: formatSkillInvocation(skill, request.additionalInstructions) }],\n\t\t\t\t\t\ttimestamp: startedAt,\n\t\t\t\t\t},\n\t\t\t\t];\n\t\t\t\tbreak;\n\t\t\t}\n\t\t\tcase \"prompt_template\": {\n\t\t\t\tconst template = acceptanceConfig.resources.promptTemplates?.find(\n\t\t\t\t\t(candidate) => candidate.name === request.name,\n\t\t\t\t);\n\t\t\t\tif (template === undefined) {\n\t\t\t\t\treturn Result.err(\n\t\t\t\t\t\tnew UnknownTemplate({ name: request.name, message: `Unknown prompt template: ${request.name}` }),\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t\tconst content = formatPromptTemplateInvocation(template, request.args);\n\t\t\t\tmessages =\n\t\t\t\t\tcontent.length === 0\n\t\t\t\t\t\t? []\n\t\t\t\t\t\t: [{ role: \"user\", content: [{ type: \"text\", text: content }], timestamp: startedAt }];\n\t\t\t\tbreak;\n\t\t\t}\n\t\t}\n\n\t\tfor (const message of messages) {\n\t\t\tif (message.role === \"assistant\" && message.stopReason === \"pending\") {\n\t\t\t\treturn Result.err(\n\t\t\t\t\tnew InvalidMessage({\n\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\treason: \"pending_assistant\",\n\t\t\t\t\t\tmessage: \"Cannot accept a pending assistant message\",\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t}\n\t\t}\n\t\tconst prompt = messages.map((message) => ({ id: this.session.idGenerator.next(startedAt), message }));\n\n\t\treturn this.command<OperationAdmissionResult>(async (state, reader) => {\n\t\t\tif (state.operation !== null) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\tnew LaneBusy({\n\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\toperationId: state.operation.meta.operationId,\n\t\t\t\t\t\t\toperationKind: state.operation.meta.intent.kind,\n\t\t\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} already has an active operation`,\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t};\n\t\t\t}\n\t\t\tconst { selected: selectedItems, remainder: inbox } = selectAcceptedInbox(\n\t\t\t\tstate.inbox,\n\t\t\t\tacceptanceConfig.steeringMode,\n\t\t\t\tacceptanceConfig.followUpMode,\n\t\t\t);\n\t\t\tconst captured = await Promise.all(\n\t\t\t\tselectedItems.map(async (item) => {\n\t\t\t\t\tconst stored = await reader.getValue(pendingEntry(item.entryId), context);\n\t\t\t\t\tif (stored === undefined) {\n\t\t\t\t\t\tthrow new SessionInvariantError(`Pending ${item.kind} entry ${item.entryId} is missing its payload`);\n\t\t\t\t\t}\n\t\t\t\t\tif (item.kind !== \"write\" && stored.value.type !== \"message\") {\n\t\t\t\t\t\tthrow new SessionInvariantError(`Pending ${item.kind} entry ${item.entryId} is not a message`);\n\t\t\t\t\t}\n\t\t\t\t\tif (\n\t\t\t\t\t\tstored.value.type === \"message\" &&\n\t\t\t\t\t\tstored.value.payload.role === \"assistant\" &&\n\t\t\t\t\t\tstored.value.payload.stopReason === \"pending\"\n\t\t\t\t\t) {\n\t\t\t\t\t\tthrow new SessionInvariantError(\n\t\t\t\t\t\t\t`Pending ${item.kind} entry ${item.entryId} contains a pending assistant`,\n\t\t\t\t\t\t);\n\t\t\t\t\t}\n\t\t\t\t\treturn { item, pending: stored.value };\n\t\t\t\t}),\n\t\t\t);\n\t\t\tconst hasCapturedConversation = selectedItems.some((item) => item.kind !== \"write\");\n\t\t\tif (prompt.length === 0 && !hasCapturedConversation) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\tnew InvalidMessage({\n\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\treason: \"empty\",\n\t\t\t\t\t\t\tmessage: \"Acceptance must append at least one message\",\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t};\n\t\t\t}\n\n\t\t\tconst entries = chainEntries(state.tipId, [\n\t\t\t\t...captured.map(({ item, pending }) => pendingEntryWrite(item.entryId, pending)),\n\t\t\t\t...prompt.map(({ id, message }) => ({ id, parentId: null, type: \"message\" as const, message })),\n\t\t\t]);\n\t\t\tconst parentId = entries[entries.length - 1]!.id;\n\t\t\tconst meta = {\n\t\t\t\toperationId,\n\t\t\t\tlane: this.name,\n\t\t\t\tsourceTipId: state.tipId,\n\t\t\t\tstartedAt,\n\t\t\t\tintent: { kind: \"run\" as const, promptEntryIds: prompt.map(({ id }) => id) },\n\t\t\t};\n\t\t\tconst operationState: StartingOperation = {\n\t\t\t\tat: \"starting\",\n\t\t\t\tcontrol: { status: \"running\" },\n\t\t\t\tsettings: capturedSettings(acceptanceConfig),\n\t\t\t\tlatestAssistantEntryId: null,\n\t\t\t};\n\t\t\tconst remainingQueues = await readLaneQueues(reader, inbox, context);\n\t\t\tconst next: LaneState = {\n\t\t\t\t...state,\n\t\t\t\ttipId: parentId,\n\t\t\t\tinbox,\n\t\t\t\toperation: { meta, state: operationState },\n\t\t\t};\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\t...entries.map((entry) => insertEntry(entry)),\n\t\t\t\t\t...selectedItems.map((item) => deleteValue(pendingEntry(item.entryId))),\n\t\t\t\t\tsetValue(branchTip(this.name), parentId),\n\t\t\t\t\tsetValue(operationMetaValue(operationId), meta),\n\t\t\t\t\tsetValue(operationStateValue(operationId), operationState),\n\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, operationId, inbox)),\n\t\t\t\t],\n\t\t\t\tnext,\n\t\t\t\tmaterialize: () => Result.ok({ operationId, kind: \"run\", startedAt }),\n\t\t\t\tevents: (commit) => {\n\t\t\t\t\tconst events: HarnessEvent[] = [\n\t\t\t\t\t\t{ type: \"run_start\", runId: operationId, startedAt, lane: this.name },\n\t\t\t\t\t\t...committedEntryEvents(entries, commit, this.name, operationId),\n\t\t\t\t\t];\n\t\t\t\t\tif (selectedItems.length > 0) {\n\t\t\t\t\t\tevents.push({ type: \"queue_update\", queues: remainingQueues, lane: this.name });\n\t\t\t\t\t}\n\t\t\t\t\treturn events;\n\t\t\t\t},\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tprivate acceptCompaction(\n\t\trequest: Extract<OperationRequest, { kind: \"compaction\" }>,\n\t\toperationId: string,\n\t\tstartedAt: number,\n\t\tacceptanceConfig: Config<TContext>,\n\t\tcontext: Context,\n\t): Promise<OperationAdmissionResult> {\n\t\tconst taskId = this.session.idGenerator.next(startedAt);\n\t\treturn this.command<OperationAdmissionResult>(async (state, reader) => {\n\t\t\tif (state.operation !== null) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\tnew LaneBusy({\n\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\toperationId: state.operation.meta.operationId,\n\t\t\t\t\t\t\toperationKind: state.operation.meta.intent.kind,\n\t\t\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} already has an active operation`,\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t};\n\t\t\t}\n\t\t\tconst path =\n\t\t\t\tstate.tipId === null\n\t\t\t\t\t? []\n\t\t\t\t\t: (\n\t\t\t\t\t\t\tawait reader.scanBranch(\n\t\t\t\t\t\t\t\t{ start: state.tipId, stopAtType: \"compaction\", order: \"newestFirst\" },\n\t\t\t\t\t\t\t\tcontext,\n\t\t\t\t\t\t\t)\n\t\t\t\t\t\t).reverse();\n\t\t\tconst prepared = prepareCompaction(path, acceptanceConfig.compaction);\n\t\t\tif (!prepared.ok) throw prepared.error;\n\t\t\tif (prepared.value === undefined) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\tnew NothingToCompact({\n\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} has nothing to compact`,\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t};\n\t\t\t}\n\t\t\tconst meta: OperationMeta = {\n\t\t\t\toperationId,\n\t\t\t\tlane: this.name,\n\t\t\t\tsourceTipId: state.tipId,\n\t\t\t\tstartedAt,\n\t\t\t\tintent: {\n\t\t\t\t\tkind: \"compaction\",\n\t\t\t\t\t...(request.customInstructions === undefined ? {} : { customInstructions: request.customInstructions }),\n\t\t\t\t},\n\t\t\t};\n\t\t\tconst operationState: SummaryDecidingOperation = {\n\t\t\t\tat: \"summary.deciding\",\n\t\t\t\tcontrol: { status: \"running\" },\n\t\t\t\tsettings: capturedSettings(acceptanceConfig),\n\t\t\t\tlatestAssistantEntryId: null,\n\t\t\t\ttask: {\n\t\t\t\t\ttaskId,\n\t\t\t\t\treason: \"manual\",\n\t\t\t\t\t...(request.customInstructions === undefined ? {} : { customInstructions: request.customInstructions }),\n\t\t\t\t\tboundary: { kind: \"finish\" },\n\t\t\t\t},\n\t\t\t};\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\tsetValue(operationPreparation(operationId, taskId), durableCompactionPreparation(prepared.value)),\n\t\t\t\t\tsetValue(operationMetaValue(operationId), meta),\n\t\t\t\t\tsetValue(operationStateValue(operationId), operationState),\n\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, operationId)),\n\t\t\t\t],\n\t\t\t\tnext: { ...state, operation: { meta, state: operationState } },\n\t\t\t\tmaterialize: () => Result.ok({ operationId, kind: \"compaction\", startedAt }),\n\t\t\t\tevents: () => [\n\t\t\t\t\t{ type: \"compaction_start\", lane: this.name, runId: operationId, reason: \"manual\", startedAt },\n\t\t\t\t],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tprivate async acceptNavigation(\n\t\trequest: Extract<OperationRequest, { kind: \"navigation\" }>,\n\t\toperationId: string,\n\t\tstartedAt: number,\n\t\tacceptanceConfig: Config<TContext>,\n\t\tcontext: Context,\n\t): Promise<OperationAdmissionResult> {\n\t\tconst taskId = this.session.idGenerator.next(startedAt);\n\t\tconst targetId = request.targetId;\n\t\tconst options = request.options ?? {};\n\t\tconst summarize = options.summarize ?? false;\n\t\tfor (;;) {\n\t\t\tconst observedTipId = this.state.tipId;\n\t\t\tlet preparation: BranchPreparation | undefined;\n\t\t\tif (summarize && observedTipId !== null && targetId !== null) {\n\t\t\t\tconst target = await this.session.getEntries([targetId], context);\n\t\t\t\tif (target.has(targetId)) {\n\t\t\t\t\tconst [oldPath, targetPath] = await Promise.all([\n\t\t\t\t\t\tthis.session.scanBranch({ start: observedTipId, order: \"newestFirst\" }, context),\n\t\t\t\t\t\tthis.session.scanBranch({ start: targetId, order: \"newestFirst\" }, context),\n\t\t\t\t\t]);\n\t\t\t\t\tconst oldIds = new Set(oldPath.map((entry) => entry.id));\n\t\t\t\t\tconst commonAncestorId = targetPath.find((entry) => oldIds.has(entry.id))?.id ?? null;\n\t\t\t\t\tpreparation = prepareBranchEntries(\n\t\t\t\t\t\toldPath\n\t\t\t\t\t\t\t.slice(\n\t\t\t\t\t\t\t\t0,\n\t\t\t\t\t\t\t\tcommonAncestorId === null\n\t\t\t\t\t\t\t\t\t? oldPath.length\n\t\t\t\t\t\t\t\t\t: oldPath.findIndex((entry) => entry.id === commonAncestorId),\n\t\t\t\t\t\t\t)\n\t\t\t\t\t\t\t.reverse(),\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t}\n\t\t\tconst accepted = await this.command<OperationAdmissionResult | undefined>(async (state, reader) => {\n\t\t\t\tif (state.operation !== null) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tnew LaneBusy({\n\t\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\t\toperationId: state.operation.meta.operationId,\n\t\t\t\t\t\t\t\toperationKind: state.operation.meta.intent.kind,\n\t\t\t\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} already has an active operation`,\n\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (state.tipId !== observedTipId) return { kind: \"return\", result: undefined };\n\t\t\t\tif (targetId === state.tipId) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tnew InvalidNavigation({\n\t\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\t\treason: \"current_tip\",\n\t\t\t\t\t\t\t\tmessage: \"Navigation target must differ from the current tip\",\n\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (targetId === null && options.label !== undefined) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tnew InvalidNavigation({\n\t\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\t\treason: \"root_label\",\n\t\t\t\t\t\t\t\tmessage: \"Root navigation cannot set a label\",\n\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (summarize && (state.tipId === null || targetId === null)) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tnew InvalidNavigation({\n\t\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\t\treason: state.tipId === null ? \"source_root\" : \"target_root\",\n\t\t\t\t\t\t\t\tmessage: \"Summarized navigation requires non-root source and target entries\",\n\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (targetId !== null && !(await reader.getEntries([targetId], context)).has(targetId)) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(new UnknownTarget({ targetId, message: `Unknown target: ${targetId}` })),\n\t\t\t\t\t};\n\t\t\t\t}\n\n\t\t\t\tconst intent: OperationMeta[\"intent\"] = {\n\t\t\t\t\tkind: \"navigation\",\n\t\t\t\t\ttargetId,\n\t\t\t\t\tsummarize,\n\t\t\t\t\t...(options.label === undefined ? {} : { label: options.label }),\n\t\t\t\t\t...(options.customInstructions === undefined ? {} : { customInstructions: options.customInstructions }),\n\t\t\t\t};\n\t\t\t\tconst meta: OperationMeta = {\n\t\t\t\t\toperationId,\n\t\t\t\t\tlane: this.name,\n\t\t\t\t\tsourceTipId: state.tipId,\n\t\t\t\t\tstartedAt,\n\t\t\t\t\tintent,\n\t\t\t\t};\n\t\t\t\tconst operationScope = {\n\t\t\t\t\tcontrol: { status: \"running\" as const },\n\t\t\t\t\tsettings: capturedSettings(acceptanceConfig),\n\t\t\t\t\tlatestAssistantEntryId: null,\n\t\t\t\t};\n\t\t\t\tlet operationState: SummaryDecidingOperation | NavigationReadyToCommitOperation;\n\t\t\t\tconst writes: Write[] = [];\n\t\t\t\tif (summarize) {\n\t\t\t\t\tif (state.tipId === null || targetId === null || preparation === undefined) {\n\t\t\t\t\t\tthrow new SessionInvariantError(\"Validated summarized navigation is missing its preparation\");\n\t\t\t\t\t}\n\t\t\t\t\twrites.push(setValue(operationPreparation(operationId, taskId), durableBranchPreparation(preparation)));\n\t\t\t\t\toperationState = {\n\t\t\t\t\t\t...operationScope,\n\t\t\t\t\t\tat: \"summary.deciding\",\n\t\t\t\t\t\ttask: {\n\t\t\t\t\t\t\ttaskId,\n\t\t\t\t\t\t\t...(options.customInstructions === undefined\n\t\t\t\t\t\t\t\t? {}\n\t\t\t\t\t\t\t\t: { customInstructions: options.customInstructions }),\n\t\t\t\t\t\t\tboundary: {\n\t\t\t\t\t\t\t\tkind: \"commit_navigation\",\n\t\t\t\t\t\t\t\ttargetId,\n\t\t\t\t\t\t\t\t...(options.label === undefined ? {} : { label: options.label }),\n\t\t\t\t\t\t\t},\n\t\t\t\t\t\t},\n\t\t\t\t\t};\n\t\t\t\t} else {\n\t\t\t\t\toperationState = {\n\t\t\t\t\t\t...operationScope,\n\t\t\t\t\t\tat: \"navigation.ready_to_commit\",\n\t\t\t\t\t\ttargetId,\n\t\t\t\t\t\t...(options.label === undefined ? {} : { label: options.label }),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\twrites.push(\n\t\t\t\t\tsetValue(operationMetaValue(operationId), meta),\n\t\t\t\t\tsetValue(operationStateValue(operationId), operationState),\n\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, operationId)),\n\t\t\t\t);\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites,\n\t\t\t\t\tnext: { ...state, operation: { meta, state: operationState } },\n\t\t\t\t\tmaterialize: () => Result.ok({ operationId, kind: \"navigation\", startedAt }),\n\t\t\t\t\tevents: () => [{ type: \"navigation_start\", lane: this.name, runId: operationId, targetId, startedAt }],\n\t\t\t\t};\n\t\t\t}, context);\n\t\t\tif (accepted !== undefined) return accepted;\n\t\t}\n\t}\n\n\tasync drive(options: DriveOptions, context: Context): Promise<DriveResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\n\t\tfor (;;) {\n\t\t\tconst claim = await this.command<DriveClaim>(async (state, reader) => {\n\t\t\t\tconst signal = context.abortSignal;\n\t\t\t\tif (signal?.aborted) {\n\t\t\t\t\tconst reason: unknown = signal.reason;\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"reject\",\n\t\t\t\t\t\terror: reason instanceof Error ? reason : new DOMException(\"The operation was aborted\", \"AbortError\"),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (state.operation?.meta.operationId === options.operationId) {\n\t\t\t\t\tif (this.activeDrive === undefined) {\n\t\t\t\t\t\tconst drive = new Drive(options, context);\n\t\t\t\t\t\tthis.activeDrive = drive;\n\t\t\t\t\t\tthis.signalStateChange();\n\t\t\t\t\t\treturn { kind: \"return\", result: { kind: \"observe\", drive, installed: true } };\n\t\t\t\t\t}\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult:\n\t\t\t\t\t\t\tthis.activeDrive.operationId === options.operationId\n\t\t\t\t\t\t\t\t? { kind: \"observe\", drive: this.activeDrive, installed: false }\n\t\t\t\t\t\t\t\t: { kind: \"occupied\", drive: this.activeDrive },\n\t\t\t\t\t};\n\t\t\t\t}\n\n\t\t\t\tconst stored = await reader.getValue(operationResultValue(options.operationId), context);\n\t\t\t\treturn stored === undefined\n\t\t\t\t\t? {\n\t\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\t\tresult: {\n\t\t\t\t\t\t\t\tkind: \"mismatch\",\n\t\t\t\t\t\t\t\terror: this.mismatch(\n\t\t\t\t\t\t\t\t\toptions.operationId,\n\t\t\t\t\t\t\t\t\tstate.operation?.meta.operationId ?? null,\n\t\t\t\t\t\t\t\t\tstate.lastOperationId,\n\t\t\t\t\t\t\t\t),\n\t\t\t\t\t\t\t},\n\t\t\t\t\t\t}\n\t\t\t\t\t: { kind: \"return\", result: { kind: \"settled\", outcome: stored.value } };\n\t\t\t}, context);\n\n\t\t\tif (claim.kind === \"settled\") return Result.ok({ kind: \"settled\", outcome: claim.outcome });\n\t\t\tif (claim.kind === \"mismatch\") return Result.err(claim.error);\n\t\t\tif (claim.kind === \"occupied\") {\n\t\t\t\tawait awaitWithContext(claim.drive.completion, context);\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (claim.installed) {\n\t\t\t\tvoid driveOperation(this, claim.drive).then(\n\t\t\t\t\t(outcome) => {\n\t\t\t\t\t\tif (this.activeDrive === claim.drive) {\n\t\t\t\t\t\t\tthis.activeDrive = undefined;\n\t\t\t\t\t\t\tthis.signalStateChange();\n\t\t\t\t\t\t}\n\t\t\t\t\t\tclaim.drive.settle(outcome);\n\t\t\t\t\t},\n\t\t\t\t\t(error: unknown) => {\n\t\t\t\t\t\tlet failure: unknown = this.closedError;\n\t\t\t\t\t\tif (failure === undefined) {\n\t\t\t\t\t\t\ttry {\n\t\t\t\t\t\t\t\tfailure = this.onFault(error, claim.drive.context);\n\t\t\t\t\t\t\t} catch (faultError) {\n\t\t\t\t\t\t\t\tfailure = faultError;\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t}\n\t\t\t\t\t\tif (this.activeDrive === claim.drive) {\n\t\t\t\t\t\t\tthis.activeDrive = undefined;\n\t\t\t\t\t\t\tthis.signalStateChange();\n\t\t\t\t\t\t}\n\t\t\t\t\t\tclaim.drive.fail(failure);\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t}\n\t\t\treturn Result.ok(await awaitWithContext(claim.drive.completion, context));\n\t\t}\n\t}\n\n\t/** Package-private durable cancellation primitive. Public exposure remains guarded until M8. */\n\tasync requestOperationAbort(operationId: string, context: Context): Promise<AbortRequestResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\n\t\tconst drive = this.activeDrive?.operationId === operationId ? this.activeDrive : undefined;\n\t\tlet resolveCancellation!: () => void;\n\t\tlet rejectCancellation!: (error: unknown) => void;\n\t\tconst cancellation = new Promise<void>((resolve, reject) => {\n\t\t\tresolveCancellation = resolve;\n\t\t\trejectCancellation = reject;\n\t\t});\n\t\tvoid cancellation.catch(() => {});\n\t\tdrive?.beginAbort(cancellation);\n\t\tlet gateSettled = false;\n\t\tconst settleGate = (signal: boolean): void => {\n\t\t\tif (gateSettled) return;\n\t\t\tgateSettled = true;\n\t\t\tresolveCancellation();\n\t\t\tif (signal) drive?.signalAbort();\n\t\t};\n\n\t\ttry {\n\t\t\tconst result = await this.command<AbortRequestResult>(async (state, reader) => {\n\t\t\t\tconst operation = state.operation;\n\t\t\t\tif (operation?.meta.operationId !== operationId) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tthis.mismatch(operationId, operation?.meta.operationId ?? null, state.lastOperationId),\n\t\t\t\t\t\t),\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tif (operation.state.control.status === \"cancel_requested\") {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.ok({ operationId, newlyRequested: false, steer: [], followUp: [] }),\n\t\t\t\t\t};\n\t\t\t\t}\n\n\t\t\t\tconst removed = state.inbox.filter((item) => item.kind === \"steer\" || item.kind === \"followUp\");\n\t\t\t\tconst payloads = await Promise.all(\n\t\t\t\t\tremoved.map(async (item) => {\n\t\t\t\t\t\tconst stored = await reader.getValue(pendingEntry(item.entryId), context);\n\t\t\t\t\t\tif (stored?.value.type !== \"message\") {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(\n\t\t\t\t\t\t\t\t`Pending ${item.kind} entry ${item.entryId} is missing its message`,\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t}\n\t\t\t\t\t\treturn { item, message: stored.value.payload };\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t\tconst steer = payloads.filter(({ item }) => item.kind === \"steer\").map(({ message }) => message);\n\t\t\t\tconst followUp = payloads.filter(({ item }) => item.kind === \"followUp\").map(({ message }) => message);\n\t\t\t\tconst removedIds = new Set(removed.map((item) => item.entryId));\n\t\t\t\tconst inbox = state.inbox.filter((item) => !removedIds.has(item.entryId));\n\t\t\t\tconst queues = await readLaneQueues(reader, inbox, context);\n\t\t\t\tconst operationState: OperationState = {\n\t\t\t\t\t...operation.state,\n\t\t\t\t\tcontrol: { status: \"cancel_requested\", requestedAt: Date.now() },\n\t\t\t\t};\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites: [\n\t\t\t\t\t\t...removed.map((item) => deleteValue(pendingEntry(item.entryId))),\n\t\t\t\t\t\tsetValue(operationStateValue(operationId), operationState),\n\t\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, operationId, inbox)),\n\t\t\t\t\t],\n\t\t\t\t\tnext: { ...state, inbox, operation: { meta: operation.meta, state: operationState } },\n\t\t\t\t\tmaterialize: () => {\n\t\t\t\t\t\tsettleGate(true);\n\t\t\t\t\t\treturn Result.ok({ operationId, newlyRequested: true, steer, followUp });\n\t\t\t\t\t},\n\t\t\t\t\tevents: () => [\n\t\t\t\t\t\t{\n\t\t\t\t\t\t\ttype: \"operation_abort\",\n\t\t\t\t\t\t\toperationId,\n\t\t\t\t\t\t\tsteer,\n\t\t\t\t\t\t\tfollowUp,\n\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t},\n\t\t\t\t\t\t...(removed.length === 0 ? [] : [{ type: \"queue_update\" as const, queues, lane: this.name }]),\n\t\t\t\t\t],\n\t\t\t\t};\n\t\t\t}, context);\n\t\t\t// A fresh Drive may observe an already-durable marker; pull its gate on the repeat path too. A mismatch\n\t\t\t// can leave only the stale Drive's admission gate in aborting state; its cancellation wait is still released.\n\t\t\tsettleGate(result.ok && result.value.newlyRequested === false);\n\t\t\treturn result;\n\t\t} catch (error) {\n\t\t\tif (!gateSettled) rejectCancellation(error);\n\t\t\tthrow error;\n\t\t}\n\t}\n\n\trequestAbort(operationId: string, context: Context): Promise<AbortRequestResult> {\n\t\treturn this.requestOperationAbort(operationId, context);\n\t}\n\n\tinspectExecution(context: Context): Promise<LaneExecutionInfo> {\n\t\treturn this.readLane((state) => {\n\t\t\tconst operation = state.operation;\n\t\t\tconst captured = operation === null ? undefined : capturedModel(operation);\n\t\t\tconst current =\n\t\t\t\toperation === null\n\t\t\t\t\t? null\n\t\t\t\t\t: {\n\t\t\t\t\t\t\tid: operation.meta.operationId,\n\t\t\t\t\t\t\tkind: operation.meta.intent.kind,\n\t\t\t\t\t\t\tstatus:\n\t\t\t\t\t\t\t\toperation.state.control.status === \"cancel_requested\"\n\t\t\t\t\t\t\t\t\t? (\"aborting\" as const)\n\t\t\t\t\t\t\t\t\t: (\"open\" as const),\n\t\t\t\t\t\t\tstartedAt: operation.meta.startedAt,\n\t\t\t\t\t\t\t...(captured === undefined ? {} : { capturedModel: captured }),\n\t\t\t\t\t\t};\n\t\t\treturn {\n\t\t\t\tlane: this.name,\n\t\t\t\ttipId: state.tipId,\n\t\t\t\tconfiguredModel: state.configuration.model,\n\t\t\t\tcurrent,\n\t\t\t\tlastOperationId: state.lastOperationId,\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tasync prompt(\n\t\t...args:\n\t\t\t| [text: string, images: ImageContent[] | undefined, context: Context]\n\t\t\t| [message: AgentMessage | AgentMessage[], context: Context]\n\t): Promise<RunResult> {\n\t\tif (args.length === 3) {\n\t\t\treturn this.driveRunRequest(\n\t\t\t\t{ kind: \"prompt\", prompt: args[0], ...(args[1] === undefined ? {} : { images: args[1] }) },\n\t\t\t\targs[2],\n\t\t\t);\n\t\t}\n\t\treturn this.driveRunRequest({ kind: \"prompt\", prompt: args[0] }, args[1]);\n\t}\n\n\tskill(name: string, additionalInstructions: string | undefined, context: Context): Promise<RunResult> {\n\t\treturn this.driveRunRequest(\n\t\t\t{ kind: \"skill\", name, ...(additionalInstructions === undefined ? {} : { additionalInstructions }) },\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tpromptFromTemplate(name: string, args: string[] | undefined, context: Context): Promise<RunResult> {\n\t\treturn this.driveRunRequest({ kind: \"prompt_template\", name, ...(args === undefined ? {} : { args }) }, context);\n\t}\n\n\tprivate async driveRunRequest(\n\t\trequest: Extract<OperationRequest, { kind: \"prompt\" | \"skill\" | \"prompt_template\" }>,\n\t\tcontext: Context,\n\t): Promise<RunResult> {\n\t\tconst admission = await this.accept(request, context);\n\t\tif (!admission.ok) {\n\t\t\tswitch (admission.error._tag) {\n\t\t\t\tcase \"LaneBusy\":\n\t\t\t\tcase \"InvalidMessage\":\n\t\t\t\tcase \"UnknownSkill\":\n\t\t\t\tcase \"UnknownTemplate\":\n\t\t\t\tcase \"Closed\":\n\t\t\t\t\treturn Result.err(admission.error);\n\t\t\t\tdefault:\n\t\t\t\t\tthrow this.onFault(\n\t\t\t\t\t\tnew SessionInvariantError(`Run acceptance returned ${admission.error._tag}`),\n\t\t\t\t\t\tcontext,\n\t\t\t\t\t);\n\t\t\t}\n\t\t}\n\t\tconst driven = await this.drive({ operationId: admission.value.operationId, waitForRetry: true }, context);\n\t\tif (!driven.ok) {\n\t\t\tif (driven.error._tag === \"Closed\") return Result.err(driven.error);\n\t\t\tthrow this.onFault(\n\t\t\t\tnew SessionInvariantError(`Accepted run ${admission.value.operationId} no longer matches its lane`),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\tif (driven.value.kind === \"settled\") return Result.ok(driven.value.outcome);\n\t\tif (driven.value.reason === \"deferred\") {\n\t\t\treturn Result.ok({\n\t\t\t\toperationId: admission.value.operationId,\n\t\t\t\tstatus: \"suspended\",\n\t\t\t\tdeferred: driven.value.deferred,\n\t\t\t});\n\t\t}\n\t\tthrow this.onFault(\n\t\t\tnew SessionInvariantError(`Run ${admission.value.operationId} returned an unwaited retry`),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync compact(options: { customInstructions?: string } | undefined, context: Context): Promise<CompactionResult> {\n\t\tconst admission = await this.accept(\n\t\t\t{\n\t\t\t\tkind: \"compaction\",\n\t\t\t\t...(options?.customInstructions === undefined ? {} : { customInstructions: options.customInstructions }),\n\t\t\t},\n\t\t\tcontext,\n\t\t);\n\t\tif (!admission.ok) {\n\t\t\tswitch (admission.error._tag) {\n\t\t\t\tcase \"LaneBusy\":\n\t\t\t\tcase \"NothingToCompact\":\n\t\t\t\tcase \"Closed\":\n\t\t\t\t\treturn Result.err(admission.error);\n\t\t\t\tdefault:\n\t\t\t\t\tthrow this.onFault(\n\t\t\t\t\t\tnew SessionInvariantError(`Compaction acceptance returned ${admission.error._tag}`),\n\t\t\t\t\t\tcontext,\n\t\t\t\t\t);\n\t\t\t}\n\t\t}\n\t\tconst compacted = await this.driveStructuralAdmission(admission.value, context);\n\t\tif (!compacted.ok) return compacted;\n\t\tconst continuation = await this.continueAfterStructural(compacted.value, context);\n\t\treturn continuation.ok\n\t\t\t? Result.ok({\n\t\t\t\t\tcompaction: compacted.value,\n\t\t\t\t\t...(continuation.value === undefined ? {} : { run: continuation.value }),\n\t\t\t\t})\n\t\t\t: continuation;\n\t}\n\n\tasync navigateTree(\n\t\ttargetId: string | null,\n\t\toptions: Extract<OperationRequest, { kind: \"navigation\" }>[\"options\"],\n\t\tcontext: Context,\n\t): Promise<NavigationResult> {\n\t\tconst admission = await this.accept(\n\t\t\t{ kind: \"navigation\", targetId, ...(options === undefined ? {} : { options }) },\n\t\t\tcontext,\n\t\t);\n\t\tif (!admission.ok) {\n\t\t\tswitch (admission.error._tag) {\n\t\t\t\tcase \"LaneBusy\":\n\t\t\t\tcase \"InvalidNavigation\":\n\t\t\t\tcase \"UnknownTarget\":\n\t\t\t\tcase \"Closed\":\n\t\t\t\t\treturn Result.err(admission.error);\n\t\t\t\tdefault:\n\t\t\t\t\tthrow this.onFault(\n\t\t\t\t\t\tnew SessionInvariantError(`Navigation acceptance returned ${admission.error._tag}`),\n\t\t\t\t\t\tcontext,\n\t\t\t\t\t);\n\t\t\t}\n\t\t}\n\t\tconst navigated = await this.driveStructuralAdmission(admission.value, context);\n\t\tif (!navigated.ok) return navigated;\n\t\tconst continuation = await this.continueAfterStructural(navigated.value, context);\n\t\treturn continuation.ok\n\t\t\t? Result.ok({\n\t\t\t\t\tnavigation: navigated.value,\n\t\t\t\t\t...(continuation.value === undefined ? {} : { run: continuation.value }),\n\t\t\t\t})\n\t\t\t: continuation;\n\t}\n\n\tprivate async driveStructuralAdmission(\n\t\tadmission: OperationAdmission,\n\t\tcontext: Context,\n\t): Promise<Result<OperationResultRecord, Closed>> {\n\t\tconst driven = await this.drive({ operationId: admission.operationId, waitForRetry: true }, context);\n\t\tif (!driven.ok) {\n\t\t\tif (driven.error._tag === \"Closed\") return Result.err(driven.error);\n\t\t\tthrow this.onFault(\n\t\t\t\tnew SessionInvariantError(`Accepted ${admission.kind} ${admission.operationId} no longer matches its lane`),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\tif (driven.value.kind === \"settled\") return Result.ok(driven.value.outcome);\n\t\tthrow this.onFault(\n\t\t\tnew SessionInvariantError(`${admission.kind} ${admission.operationId} returned ${driven.value.reason}`),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tprivate async continueAfterStructural(\n\t\trecord: OperationResultRecord,\n\t\tcontext: Context,\n\t): Promise<Result<OperationResultRecord | SuspendedRun | undefined, Closed>> {\n\t\tif (record.status === \"aborted\") return Result.ok(undefined);\n\t\tconst admission = await this.accept({ kind: \"prompt\", prompt: \"\" }, context);\n\t\tif (!admission.ok) {\n\t\t\tswitch (admission.error._tag) {\n\t\t\t\tcase \"InvalidMessage\":\n\t\t\t\tcase \"LaneBusy\":\n\t\t\t\t\treturn Result.ok(undefined);\n\t\t\t\tcase \"Closed\":\n\t\t\t\t\treturn Result.err(admission.error);\n\t\t\t\tdefault:\n\t\t\t\t\tthrow this.onFault(\n\t\t\t\t\t\tnew SessionInvariantError(`Structural continuation acceptance returned ${admission.error._tag}`),\n\t\t\t\t\t\tcontext,\n\t\t\t\t\t);\n\t\t\t}\n\t\t}\n\t\tconst driven = await this.drive({ operationId: admission.value.operationId, waitForRetry: true }, context);\n\t\tif (!driven.ok) {\n\t\t\tif (driven.error._tag === \"Closed\") return Result.err(driven.error);\n\t\t\tthrow this.onFault(\n\t\t\t\tnew SessionInvariantError(`Continuation run ${admission.value.operationId} no longer matches its lane`),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\tif (driven.value.kind === \"settled\") return Result.ok(driven.value.outcome);\n\t\tif (driven.value.reason === \"deferred\") {\n\t\t\treturn Result.ok({\n\t\t\t\toperationId: admission.value.operationId,\n\t\t\t\tstatus: \"suspended\",\n\t\t\t\tdeferred: driven.value.deferred,\n\t\t\t});\n\t\t}\n\t\tthrow this.onFault(\n\t\t\tnew SessionInvariantError(`Continuation run ${admission.value.operationId} returned an unwaited retry`),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync resume(context: Context): Promise<ResumeResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\tconst inspected = await this.command<Result<{ operationId: string }, NothingToResume>>((state) => {\n\t\t\tconst operation = state.operation;\n\t\t\treturn operation === null\n\t\t\t\t? {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: Result.err(\n\t\t\t\t\t\t\tnew NothingToResume({\n\t\t\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} has no active operation to resume`,\n\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t),\n\t\t\t\t\t}\n\t\t\t\t: { kind: \"return\", result: Result.ok({ operationId: operation.meta.operationId }) };\n\t\t}, context);\n\t\tif (!inspected.ok) return inspected;\n\n\t\tconst driven = await this.drive(\n\t\t\t{ operationId: inspected.value.operationId, pollDeferred: true, waitForRetry: true },\n\t\t\tcontext,\n\t\t);\n\t\tif (!driven.ok) {\n\t\t\tif (driven.error._tag === \"Closed\") return Result.err(driven.error);\n\t\t\tthrow this.onFault(\n\t\t\t\tnew SessionInvariantError(`Operation ${inspected.value.operationId} no longer matches its lane`),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\tif (driven.value.kind === \"settled\") return Result.ok(driven.value.outcome);\n\t\tif (driven.value.reason === \"deferred\") {\n\t\t\treturn Result.ok({\n\t\t\t\toperationId: inspected.value.operationId,\n\t\t\t\tstatus: \"suspended\",\n\t\t\t\tdeferred: driven.value.deferred,\n\t\t\t});\n\t\t}\n\t\tthrow this.onFault(\n\t\t\tnew SessionInvariantError(`Operation ${inspected.value.operationId} returned an unwaited retry`),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync abort(context: Context): Promise<AbortResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\tconst operationId = await this.command(\n\t\t\t(state) => ({\n\t\t\t\tkind: \"return\",\n\t\t\t\tresult: state.operation?.meta.operationId,\n\t\t\t}),\n\t\t\tcontext,\n\t\t);\n\t\tif (operationId === undefined) {\n\t\t\treturn Result.err(\n\t\t\t\tnew NoActiveOperation({\n\t\t\t\t\tlane: this.name,\n\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} has no active operation`,\n\t\t\t\t}),\n\t\t\t);\n\t\t}\n\t\tconst requested = await this.requestAbort(operationId, context);\n\t\tif (!requested.ok) {\n\t\t\tif (requested.error._tag === \"Closed\") return Result.err(requested.error);\n\t\t\treturn Result.err(\n\t\t\t\tnew NoActiveOperation({\n\t\t\t\t\tlane: this.name,\n\t\t\t\t\tmessage: `Lane ${JSON.stringify(this.name)} no longer has the inspected operation`,\n\t\t\t\t}),\n\t\t\t);\n\t\t}\n\t\tconst driven = await this.drive({ operationId }, context);\n\t\tif (!driven.ok && driven.error._tag === \"Closed\") return Result.err(driven.error);\n\t\tif (!driven.ok) {\n\t\t\tthrow this.onFault(\n\t\t\t\tnew SessionInvariantError(`Cancelled operation ${operationId} no longer matches its lane`),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\treturn Result.ok({\n\t\t\toperationId,\n\t\t\tsteer: requested.value.steer,\n\t\t\tfollowUp: requested.value.followUp,\n\t\t});\n\t}\n\n\tsteer(message: string | AgentMessage, images: ImageContent[] | undefined, context: Context): Promise<QueueResult> {\n\t\treturn this.enqueue(\"steer\", message, images, context);\n\t}\n\n\tfollowUp(\n\t\tmessage: string | AgentMessage,\n\t\timages: ImageContent[] | undefined,\n\t\tcontext: Context,\n\t): Promise<QueueResult> {\n\t\treturn this.enqueue(\"followUp\", message, images, context);\n\t}\n\n\tnextRun(message: string | AgentMessage, images: ImageContent[] | undefined, context: Context): Promise<QueueResult> {\n\t\treturn this.enqueue(\"nextRun\", message, images, context);\n\t}\n\n\tprivate async enqueue(\n\t\tkind: \"steer\" | \"followUp\" | \"nextRun\",\n\t\tinput: string | AgentMessage,\n\t\timages: ImageContent[] | undefined,\n\t\tcontext: Context,\n\t): Promise<QueueResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\tconst at = Date.now();\n\t\tlet message: AgentMessage;\n\t\tif (typeof input === \"string\") {\n\t\t\tif (input.length === 0 && (images === undefined || images.length === 0)) {\n\t\t\t\treturn Result.err(\n\t\t\t\t\tnew InvalidMessage({\n\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\treason: \"empty\",\n\t\t\t\t\t\tmessage: \"Queued input must contain text or an image\",\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t}\n\t\t\tmessage = {\n\t\t\t\trole: \"user\",\n\t\t\t\tcontent: [...(input.length === 0 ? [] : [{ type: \"text\" as const, text: input }]), ...(images ?? [])],\n\t\t\t\ttimestamp: at,\n\t\t\t};\n\t\t} else {\n\t\t\tif (input.role === \"assistant\" && input.stopReason === \"pending\") {\n\t\t\t\treturn Result.err(\n\t\t\t\t\tnew InvalidMessage({\n\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\treason: \"pending_assistant\",\n\t\t\t\t\t\tmessage: \"Cannot queue a pending assistant message\",\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t}\n\t\t\tif (images !== undefined && images.length > 0 && input.role !== \"user\") {\n\t\t\t\treturn Result.err(\n\t\t\t\t\tnew InvalidMessage({\n\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\treason: \"images_with_non_user\",\n\t\t\t\t\t\tmessage: \"Images can be added only to queued user messages\",\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t}\n\t\t\tmessage =\n\t\t\t\timages === undefined || images.length === 0 || input.role !== \"user\"\n\t\t\t\t\t? structuredClone(input)\n\t\t\t\t\t: {\n\t\t\t\t\t\t\t...structuredClone(input),\n\t\t\t\t\t\t\tcontent: [\n\t\t\t\t\t\t\t\t...(typeof input.content === \"string\"\n\t\t\t\t\t\t\t\t\t? input.content.length === 0\n\t\t\t\t\t\t\t\t\t\t? []\n\t\t\t\t\t\t\t\t\t\t: [{ type: \"text\" as const, text: input.content }]\n\t\t\t\t\t\t\t\t\t: input.content),\n\t\t\t\t\t\t\t\t...images,\n\t\t\t\t\t\t\t],\n\t\t\t\t\t\t};\n\t\t}\n\t\tconst entryId = this.session.idGenerator.next(at);\n\t\treturn this.command<QueueResult>(async (state, reader) => {\n\t\t\tconst inbox = [...state.inbox, { entryId, kind }];\n\t\t\tconst queues = [\n\t\t\t\t...(await readLaneQueues(reader, state.inbox, context)),\n\t\t\t\t{ entryId, kind, type: \"message\" as const, message },\n\t\t\t];\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\tsetValue(pendingEntry(entryId), { type: \"message\", payload: message }),\n\t\t\t\t\tsetValue(\n\t\t\t\t\t\tlaneStateValue(this.name),\n\t\t\t\t\t\tdurableLaneState(state, state.operation?.meta.operationId ?? null, inbox),\n\t\t\t\t\t),\n\t\t\t\t],\n\t\t\t\tnext: { ...state, inbox },\n\t\t\t\tmaterialize: () => Result.ok({ entryId }),\n\t\t\t\tevents: () => [{ type: \"queue_update\", queues, lane: this.name }],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tasync cancelQueued(entryId: string, context: Context): Promise<CancelQueuedResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\treturn this.command<CancelQueuedResult>(async (state, reader) => {\n\t\t\tconst queued = state.inbox.find((item) => item.entryId === entryId);\n\t\t\tif (queued === undefined) {\n\t\t\t\tconst consumed = (await reader.getEntries([entryId], context)).has(entryId);\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult: Result.ok({ kind: consumed ? \"already_consumed\" : \"not_found\" }),\n\t\t\t\t};\n\t\t\t}\n\t\t\tif ((await reader.getValue(pendingEntry(entryId), context)) === undefined) {\n\t\t\t\tthrow new SessionInvariantError(`Queued ${queued.kind} entry ${entryId} is missing its payload`);\n\t\t\t}\n\t\t\tconst inbox = state.inbox.filter((item) => item.entryId !== entryId);\n\t\t\tconst queues = await readLaneQueues(reader, inbox, context);\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\tdeleteValue(pendingEntry(entryId)),\n\t\t\t\t\tsetValue(\n\t\t\t\t\t\tlaneStateValue(this.name),\n\t\t\t\t\t\tdurableLaneState(state, state.operation?.meta.operationId ?? null, inbox),\n\t\t\t\t\t),\n\t\t\t\t],\n\t\t\t\tnext: { ...state, inbox },\n\t\t\t\tmaterialize: () => Result.ok({ kind: \"cancelled\" }),\n\t\t\t\tevents: () => [{ type: \"queue_update\", queues, lane: this.name }],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tasync recordUsage(\n\t\tusage: Usage,\n\t\toptions: { entryId?: string; details?: JsonValue } | undefined,\n\t\tcontext: Context,\n\t): Promise<RecordUsageResult> {\n\t\tif (this.closedError instanceof HarnessClosed) {\n\t\t\treturn Result.err(new Closed({ message: this.closedError.message }));\n\t\t}\n\t\tthis.assertOpen();\n\t\treturn this.command<RecordUsageResult>((state) => {\n\t\t\tconst usageId = this.session.idGenerator.next();\n\t\t\tconst row = {\n\t\t\t\tid: usageId,\n\t\t\t\tusage,\n\t\t\t\tadjustment: true,\n\t\t\t\t...(options?.entryId === undefined ? {} : { entryId: options.entryId }),\n\t\t\t\t...(options?.details === undefined ? {} : { details: options.details }),\n\t\t\t};\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [insertUsage(row)],\n\t\t\t\tnext: state,\n\t\t\t\tmaterialize: () => Result.ok({ usageId }),\n\t\t\t\tevents: (commit) => [\n\t\t\t\t\t{\n\t\t\t\t\t\ttype: \"usage\",\n\t\t\t\t\t\tlane: this.name,\n\t\t\t\t\t\trow: { ...row, seq: commit.seqs[0]! },\n\t\t\t\t\t\ttotals: commit.stats.usage,\n\t\t\t\t\t},\n\t\t\t\t],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tasync waitForIdle(context: Context): Promise<void> {\n\t\tfor (;;) {\n\t\t\tconst observation = await this.command<IdleObservation>(\n\t\t\t\t(state) => ({\n\t\t\t\t\tkind: \"return\",\n\t\t\t\t\tresult:\n\t\t\t\t\t\tstate.operation === null && this.activeDrive === undefined\n\t\t\t\t\t\t\t? { kind: \"idle\" }\n\t\t\t\t\t\t\t: { kind: \"wait\", drive: this.activeDrive, change: this.stateChange },\n\t\t\t\t}),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t\tif (observation.kind === \"idle\") return;\n\t\t\tawait awaitWithContext(\n\t\t\t\tobservation.drive === undefined ? observation.change : observation.drive.completion.then(() => undefined),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t}\n\n\tasync runWhenIdle(callback: (context: Context) => void | Promise<void>, context: Context): Promise<void> {\n\t\tlet owner: Promise<void> | undefined;\n\t\tlet releaseOwner: (() => void) | undefined;\n\t\tfor (;;) {\n\t\t\tconst observation = await this.command<IdleClaimObservation>((state) => {\n\t\t\t\tif (state.operation !== null || this.activeDrive !== undefined) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"return\",\n\t\t\t\t\t\tresult: { kind: \"wait\", drive: this.activeDrive, change: this.stateChange },\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t\tlet release!: () => void;\n\t\t\t\tconst claimed = new Promise<void>((resolve) => {\n\t\t\t\t\trelease = resolve;\n\t\t\t\t});\n\t\t\t\towner = claimed;\n\t\t\t\treleaseOwner = release;\n\t\t\t\tthis.idleOwner = claimed;\n\t\t\t\tthis.signalStateChange();\n\t\t\t\treturn { kind: \"return\", result: { kind: \"claimed\" } };\n\t\t\t}, context);\n\t\t\tif (observation.kind === \"claimed\") break;\n\t\t\tawait awaitWithContext(\n\t\t\t\tobservation.drive === undefined ? observation.change : observation.drive.completion.then(() => undefined),\n\t\t\t\tcontext,\n\t\t\t);\n\t\t}\n\t\tif (owner === undefined || releaseOwner === undefined) {\n\t\t\tthrow this.onFault(new SessionInvariantError(\"Idle callback claim was not published\"), context);\n\t\t}\n\t\ttry {\n\t\t\tthis.assertOpen();\n\t\t\tawait callback(context);\n\t\t} finally {\n\t\t\tif (this.idleOwner === owner) this.idleOwner = undefined;\n\t\t\treleaseOwner();\n\t\t\tthis.signalStateChange();\n\t\t}\n\t}\n\n\tasync getModel(_context: Context): Promise<Model<Api> | undefined> {\n\t\tthis.assertOpen();\n\t\treturn this.models.getModel(this.state.configuration.model.provider, this.state.configuration.model.modelId);\n\t}\n\n\tsetModel(model: ModelIdentity, context: Context): Promise<void> {\n\t\treturn this.setConfiguration(\n\t\t\t(configuration) => ({\n\t\t\t\t...configuration,\n\t\t\t\tmodel: { ...model },\n\t\t\t}),\n\t\t\t(previous, value) => ({\n\t\t\t\ttype: \"config_update\",\n\t\t\t\tproperty: \"model\",\n\t\t\t\tprevious: previous.model,\n\t\t\t\tvalue: value.model,\n\t\t\t}),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync getThinkingLevel(_context: Context): Promise<ThinkingLevel> {\n\t\tthis.assertOpen();\n\t\treturn this.state.configuration.thinkingLevel;\n\t}\n\n\tsetThinkingLevel(thinkingLevel: ThinkingLevel, context: Context): Promise<void> {\n\t\treturn this.setConfiguration(\n\t\t\t(configuration) => ({ ...configuration, thinkingLevel }),\n\t\t\t(previous, value) => ({\n\t\t\t\ttype: \"config_update\",\n\t\t\t\tproperty: \"thinkingLevel\",\n\t\t\t\tprevious: previous.thinkingLevel,\n\t\t\t\tvalue: value.thinkingLevel,\n\t\t\t}),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\tasync getActiveTools(_context: Context): Promise<string[]> {\n\t\tthis.assertOpen();\n\t\treturn this.state.configuration.activeToolNames;\n\t}\n\n\tsetActiveTools(activeToolNames: string[], context: Context): Promise<void> {\n\t\treturn this.setConfiguration(\n\t\t\t(configuration) => ({ ...configuration, activeToolNames }),\n\t\t\t(previous, value) => ({\n\t\t\t\ttype: \"config_update\",\n\t\t\t\tproperty: \"activeTools\",\n\t\t\t\tprevious: previous.activeToolNames,\n\t\t\t\tvalue: value.activeToolNames,\n\t\t\t}),\n\t\t\tcontext,\n\t\t);\n\t}\n\n\twatch(context: Context): Promise<WatchHandle<LaneSnapshot>> {\n\t\treturn this.readLane(async (state, reader) => {\n\t\t\tconst watcher = this.installWatch<LaneSnapshot>(\n\t\t\t\t{} as LaneSnapshot,\n\t\t\t\t(event) => event.type === \"usage\" || !(\"lane\" in event) || event.lane === this.name,\n\t\t\t\tcontext,\n\t\t\t\t(resnapshotContext, markBoundary) =>\n\t\t\t\t\tthis.readLane(async (latest, latestReader) => {\n\t\t\t\t\t\tconst snapshot = await this.captureLaneSnapshot(latest, latestReader, resnapshotContext);\n\t\t\t\t\t\tmarkBoundary();\n\t\t\t\t\t\treturn snapshot;\n\t\t\t\t\t}, resnapshotContext),\n\t\t\t);\n\t\t\ttry {\n\t\t\t\twatcher.snapshot = await this.captureLaneSnapshot(state, reader, context);\n\t\t\t\treturn watcher;\n\t\t\t} catch (error) {\n\t\t\t\twatcher.unsubscribe();\n\t\t\t\tthrow error;\n\t\t\t}\n\t\t}, context);\n\t}\n\n\tprivate async captureLaneSnapshot(state: LaneState, reader: SessionReader, context: Context): Promise<LaneSnapshot> {\n\t\tconst captured = structuredClone(state);\n\t\tconst transcript =\n\t\t\tcaptured.tipId === null\n\t\t\t\t? []\n\t\t\t\t: (\n\t\t\t\t\t\tawait reader.scanBranch(\n\t\t\t\t\t\t\t{ start: captured.tipId, stopAtType: \"compaction\", order: \"newestFirst\" },\n\t\t\t\t\t\t\tcontext,\n\t\t\t\t\t\t)\n\t\t\t\t\t).reverse();\n\t\tconst queues = await readLaneQueues(reader, captured.inbox, context);\n\t\tlet lastResult: OperationResultRecord | undefined;\n\t\tif (captured.lastOperationId !== null) {\n\t\t\tconst stored = await reader.getValue(operationResultValue(captured.lastOperationId), context);\n\t\t\tif (stored === undefined) {\n\t\t\t\tthrow new SessionInvariantError(\n\t\t\t\t\t`Lane ${JSON.stringify(this.name)} is missing result ${captured.lastOperationId}`,\n\t\t\t\t);\n\t\t\t}\n\t\t\tlastResult = stored.value;\n\t\t}\n\t\tconst stats = await reader.getStats(context);\n\t\tconst operation = captured.operation;\n\n\t\tlet operationSnapshot: NonNullable<LaneSnapshot[\"operation\"]> | null = null;\n\t\tif (operation !== null) {\n\t\t\tconst runningTools: NonNullable<LaneSnapshot[\"operation\"]>[\"runningTools\"] = [];\n\t\t\tlet streamingMessage: NonNullable<LaneSnapshot[\"operation\"]>[\"streamingMessage\"];\n\t\t\tlet retry: NonNullable<LaneSnapshot[\"operation\"]>[\"retry\"];\n\t\t\tlet deferred: NonNullable<LaneSnapshot[\"operation\"]>[\"deferred\"];\n\t\t\tconst readStreamingMessage = async (responseEntryId: string) =>\n\t\t\t\treduceAssistantMessageFrames(\n\t\t\t\t\tawait readAssistantFrames(reader, operation.meta.operationId, responseEntryId, context),\n\t\t\t\t);\n\t\t\tconst state = operation.state;\n\t\t\tswitch (state.at) {\n\t\t\t\tcase \"assistant.retry_wait\":\n\t\t\t\t\tretry = {\n\t\t\t\t\t\tattempt: state.nextAttempt,\n\t\t\t\t\t\tmaxAttempts: state.generationContext.retryPolicy.maxAttempts,\n\t\t\t\t\t\tnextAttemptAt: state.notBefore,\n\t\t\t\t\t};\n\t\t\t\t\tbreak;\n\t\t\t\tcase \"assistant.effect_pending\":\n\t\t\t\t\tstreamingMessage = await readStreamingMessage(state.responseEntryId);\n\t\t\t\t\tbreak;\n\t\t\t\tcase \"deferred.suspended\":\n\t\t\t\tcase \"deferred.effect_pending\": {\n\t\t\t\t\tconst source = (await reader.getEntries([state.sourceEntryId], context)).get(state.sourceEntryId);\n\t\t\t\t\tif (\n\t\t\t\t\t\tsource?.type !== \"message\" ||\n\t\t\t\t\t\tsource.message.role !== \"assistant\" ||\n\t\t\t\t\t\tsource.message.deferred === undefined\n\t\t\t\t\t) {\n\t\t\t\t\t\tthrow new SessionInvariantError(\"Deferred source is missing its assistant handle\");\n\t\t\t\t\t}\n\t\t\t\t\tdeferred = { handle: source.message.deferred, poll: state.poll };\n\t\t\t\t\tif (state.at === \"deferred.effect_pending\") {\n\t\t\t\t\t\tstreamingMessage = await readStreamingMessage(state.responseEntryId);\n\t\t\t\t\t}\n\t\t\t\t\tbreak;\n\t\t\t\t}\n\t\t\t\tcase \"tools\": {\n\t\t\t\t\tconst { batch } = state;\n\t\t\t\t\tconst assistant = (await reader.getEntries([batch.assistantEntryId], context)).get(\n\t\t\t\t\t\tbatch.assistantEntryId,\n\t\t\t\t\t);\n\t\t\t\t\tif (assistant?.type !== \"message\" || assistant.message.role !== \"assistant\") {\n\t\t\t\t\t\tthrow new SessionInvariantError(\"Tool batch assistant entry is invalid\");\n\t\t\t\t\t}\n\t\t\t\t\tfor (const call of batch.calls) {\n\t\t\t\t\t\tif (call.status === \"planned\" || call.status === \"completed\") continue;\n\t\t\t\t\t\tconst block = assistant.message.content[call.sourceIndex];\n\t\t\t\t\t\tif (block?.type !== \"toolCall\") {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(\n\t\t\t\t\t\t\t\t`Tool call source index ${call.sourceIndex} does not name a tool-call block`,\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t}\n\t\t\t\t\t\tconst args = await reader.getValue(\n\t\t\t\t\t\t\toperationToolArgs(operation.meta.operationId, batch.turnId, call.sourceIndex),\n\t\t\t\t\t\t\tcontext,\n\t\t\t\t\t\t);\n\t\t\t\t\t\tif (call.status === \"effect_pending\") {\n\t\t\t\t\t\t\tif (args === undefined) {\n\t\t\t\t\t\t\t\tthrow new SessionInvariantError(`Tool call ${block.id} is missing persisted arguments`);\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t\tconst checkpoint = await reader.getValue(\n\t\t\t\t\t\t\t\tpendingToolOutput(operation.meta.operationId, call.resultEntryId),\n\t\t\t\t\t\t\t\tcontext,\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t\trunningTools.push({\n\t\t\t\t\t\t\t\tstatus: \"running\",\n\t\t\t\t\t\t\t\ttoolCallId: block.id,\n\t\t\t\t\t\t\t\ttoolName: block.name,\n\t\t\t\t\t\t\t\targs: args.value,\n\t\t\t\t\t\t\t\t...(checkpoint === undefined ? {} : { result: checkpoint.value }),\n\t\t\t\t\t\t\t});\n\t\t\t\t\t\t\tcontinue;\n\t\t\t\t\t\t}\n\t\t\t\t\t\tconst staged = await reader.getValue(pendingEntry(call.resultEntryId), context);\n\t\t\t\t\t\tif (staged?.value.type !== \"message\" || staged.value.payload.role !== \"toolResult\") {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(`Tool call ${call.resultEntryId} is missing its staged result`);\n\t\t\t\t\t\t}\n\t\t\t\t\t\tif (staged.value.payload.toolCallId !== block.id || staged.value.payload.toolName !== block.name) {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(`Tool call ${call.resultEntryId} has a mismatched staged result`);\n\t\t\t\t\t\t}\n\t\t\t\t\t\trunningTools.push({\n\t\t\t\t\t\t\tstatus: \"settled\",\n\t\t\t\t\t\t\ttoolCallId: block.id,\n\t\t\t\t\t\t\ttoolName: block.name,\n\t\t\t\t\t\t\targs: args?.value ?? block.arguments,\n\t\t\t\t\t\t\tresult: toolResultFromMessage(staged.value.payload, call.terminate),\n\t\t\t\t\t\t\tisError: staged.value.payload.isError,\n\t\t\t\t\t\t});\n\t\t\t\t\t}\n\t\t\t\t\tbreak;\n\t\t\t\t}\n\t\t\t\tcase \"summary.retry_wait\":\n\t\t\t\t\tretry = {\n\t\t\t\t\t\tattempt: state.nextAttempt,\n\t\t\t\t\t\tmaxAttempts: state.summaryContext.retryPolicy.maxAttempts,\n\t\t\t\t\t\tnextAttemptAt: state.notBefore,\n\t\t\t\t\t};\n\t\t\t\t\tbreak;\n\t\t\t\tdefault:\n\t\t\t\t\tbreak;\n\t\t\t}\n\n\t\t\toperationSnapshot = {\n\t\t\t\tid: operation.meta.operationId,\n\t\t\t\tkind: operation.meta.intent.kind,\n\t\t\t\tstartedAt: operation.meta.startedAt,\n\t\t\t\tfromTipId: operation.meta.sourceTipId,\n\t\t\t\tstatus: operation.state.control.status === \"cancel_requested\" ? \"aborting\" : \"open\",\n\t\t\t\t...(retry === undefined ? {} : { retry }),\n\t\t\t\t...(deferred === undefined ? {} : { deferred }),\n\t\t\t\t...(streamingMessage === undefined ? {} : { streamingMessage }),\n\t\t\t\trunningTools,\n\t\t\t};\n\t\t}\n\n\t\treturn structuredClone({\n\t\t\tlane: this.name,\n\t\t\ttranscript,\n\t\t\ttipId: captured.tipId,\n\t\t\t...(lastResult === undefined ? {} : { lastResult }),\n\t\t\tconfiguration: captured.configuration,\n\t\t\tstats,\n\t\t\toperation: operationSnapshot,\n\t\t\tqueues,\n\t\t\tfaulted: this.closedError instanceof HarnessFault,\n\t\t});\n\t}\n\n\tprivate async setConfiguration(\n\t\tupdate: (configuration: LaneState[\"configuration\"]) => LaneState[\"configuration\"],\n\t\tevent: (previous: LaneState[\"configuration\"], value: LaneState[\"configuration\"]) => LaneConfigEventPayload,\n\t\tcontext: Context,\n\t): Promise<void> {\n\t\tawait this.command((state) => {\n\t\t\tconst configuration = update(state.configuration);\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [setValue(laneConfig(this.name), configuration)],\n\t\t\t\tnext: { ...state, configuration },\n\t\t\t\tmaterialize: () => undefined,\n\t\t\t\tevents: () => [{ ...event(state.configuration, configuration), lane: this.name }],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tasync findEntries(query: BranchScan | undefined, context: Context): Promise<Entry[]> {\n\t\tquery ??= {};\n\t\tthis.assertOpen();\n\t\tconst start = query.start ?? this.state.tipId;\n\t\treturn start === null\n\t\t\t? []\n\t\t\t: this.session.scanBranch({ ...query, start, order: query.order ?? \"newestFirst\" }, context);\n\t}\n\n\tasync findEntry(query: BranchScan | undefined, context: Context): Promise<Entry | undefined> {\n\t\tquery ??= {};\n\t\treturn (\n\t\t\tawait this.findEntries({ ...query, limit: query.limit === undefined ? 1 : Math.min(query.limit, 1) }, context)\n\t\t)[0];\n\t}\n\n\tappendMessage(message: AgentMessage, context: Context): Promise<string> {\n\t\treturn this.append({ type: \"message\", payload: message }, context);\n\t}\n\n\tappendCustomEntry(customType: string, data: JsonValue | undefined, context: Context): Promise<string> {\n\t\treturn this.append({ type: \"custom\", customType, ...(data === undefined ? {} : { payload: data }) }, context);\n\t}\n\n\tprivate append(pending: PendingEntry, context: Context): Promise<string> {\n\t\tthis.assertOpen();\n\t\tif (\n\t\t\tpending.type === \"message\" &&\n\t\t\tpending.payload.role === \"assistant\" &&\n\t\t\tpending.payload.stopReason === \"pending\"\n\t\t) {\n\t\t\treturn Promise.reject(new SessionPendingAssistantMessageError());\n\t\t}\n\t\tconst id = this.session.idGenerator.next();\n\t\treturn this.command(async (state, reader) => {\n\t\t\tif (state.operation === null) {\n\t\t\t\tconst queued = inboxItems(state.inbox, \"write\");\n\t\t\t\tconst captured = await Promise.all(\n\t\t\t\t\tqueued.map(async (item) => {\n\t\t\t\t\t\tconst stored = await reader.getValue(pendingEntry(item.entryId), context);\n\t\t\t\t\t\tif (stored === undefined) {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(`Pending write ${item.entryId} is missing its payload`);\n\t\t\t\t\t\t}\n\t\t\t\t\t\treturn pendingEntryWrite(item.entryId, stored.value);\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t\tconst inbox = withoutInboxItems(state.inbox, queued);\n\t\t\t\tconst queues = queued.length === 0 ? undefined : await readLaneQueues(reader, inbox, context);\n\t\t\t\tconst entries = chainEntries(state.tipId, [...captured, pendingEntryWrite(id, pending)]);\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites: [\n\t\t\t\t\t\t...entries.map((entry) => insertEntry(entry)),\n\t\t\t\t\t\t...queued.map((item) => deleteValue(pendingEntry(item.entryId))),\n\t\t\t\t\t\tsetValue(branchTip(this.name), id),\n\t\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, null, inbox)),\n\t\t\t\t\t],\n\t\t\t\t\tnext: { ...state, tipId: id, inbox },\n\t\t\t\t\tmaterialize: () => id,\n\t\t\t\t\tevents: (commit) => [\n\t\t\t\t\t\t...committedEntryEvents(entries, commit, this.name),\n\t\t\t\t\t\t...(queues === undefined ? [] : [{ type: \"queue_update\" as const, queues, lane: this.name }]),\n\t\t\t\t\t],\n\t\t\t\t};\n\t\t\t}\n\n\t\t\tconst operation = state.operation;\n\t\t\tconst inbox = [...state.inbox, { entryId: id, kind: \"write\" as const }];\n\t\t\tconst queues = [\n\t\t\t\t...(await readLaneQueues(reader, state.inbox, context)),\n\t\t\t\tpending.type === \"message\"\n\t\t\t\t\t? { entryId: id, kind: \"write\" as const, type: \"message\" as const, message: pending.payload }\n\t\t\t\t\t: {\n\t\t\t\t\t\t\tentryId: id,\n\t\t\t\t\t\t\tkind: \"write\" as const,\n\t\t\t\t\t\t\ttype: \"custom\" as const,\n\t\t\t\t\t\t\tcustomType: pending.customType,\n\t\t\t\t\t\t\t...(pending.payload === undefined ? {} : { data: pending.payload }),\n\t\t\t\t\t\t},\n\t\t\t];\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites: [\n\t\t\t\t\tsetValue(pendingEntry(id), pending),\n\t\t\t\t\tsetValue(laneStateValue(this.name), durableLaneState(state, operation.meta.operationId, inbox)),\n\t\t\t\t],\n\t\t\t\tnext: { ...state, inbox },\n\t\t\t\tmaterialize: () => id,\n\t\t\t\tevents: () => [{ type: \"queue_update\", queues, lane: this.name }],\n\t\t\t};\n\t\t}, context);\n\t}\n\n\tseal(error: Error): Promise<void> {\n\t\tthis.closedError ??= error;\n\t\tthis.activeDrive?.closeGate(error);\n\t\tthis.signalStateChange();\n\t\treturn this.idleOwner ?? Promise.resolve();\n\t}\n\n\tprivate signalStateChange(): void {\n\t\tthis.resolveStateChange();\n\t\tlet resolveStateChange!: () => void;\n\t\tthis.stateChange = new Promise<void>((resolve) => {\n\t\t\tresolveStateChange = resolve;\n\t\t});\n\t\tthis.resolveStateChange = resolveStateChange;\n\t}\n\n\tassertOpen(): void {\n\t\tif (this.closedError !== undefined) throw this.closedError;\n\t}\n}\n"]}