{"version":3,"file":"checkpoint.d.ts","sourceRoot":"","sources":["../../../../src/harness/runtime/drive/checkpoint.ts"],"names":[],"mappings":"AAEA,OAAO,EACN,KAAK,mBAAmB,EAIxB,KAAK,iBAAiB,EAEtB,MAAM,wBAAwB,CAAC;AAEhC,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,YAAY,CAAC;AAEvC,OAAO,KAAK,EAAE,KAAK,EAAE,eAAe,EAAE,MAAM,aAAa,CAAC;AAU1D,4DAA4D;AAC5D,wBAAsB,QAAQ,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACjE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,GAAG,EAAE,iBAAiB,GACpB,OAAO,CAAC,eAAe,CAAC,CA+D1B;AAED,gEAAgE;AAChE,wBAAsB,aAAa,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACtE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,GAAG,EAAE,mBAAmB,GACtB,OAAO,CAAC,eAAe,CAAC,CA2F1B","sourcesContent":["import { insertEntry } from \"../../session/commit.ts\";\nimport { SessionInvariantError } from \"../../session/session.ts\";\nimport {\n\ttype CheckpointOperation,\n\ttype MessageEntry,\n\ttype NewEntry,\n\toperationScopeOf,\n\ttype StartingOperation,\n\ttype SummaryDecidingOperation,\n} from \"../../session/types.ts\";\nimport { branchTip, operationPreparation, setValue } from \"../../session/values.ts\";\nimport type { Lane } from \"../lane.ts\";\nimport { chainEntries, committedEntryEvents } from \"../transcript.ts\";\nimport type { Drive, ProcedureResult } from \"../types.ts\";\nimport {\n\tassistantReadyAtBoundary,\n\ttype BoundaryFinishPending,\n\tboundaryPlacementEvents,\n\tfinishRunBoundary,\n\tplanBoundaryInbox,\n} from \"./boundary.ts\";\nimport { prepareCompactionThreshold } from \"./structural.ts\";\n\n/** Consume before_run and commit the initial checkpoint. */\nexport async function startRun<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\trun: StartingOperation,\n): Promise<ProcedureResult> {\n\tconst prompt = await lane.continueOperation(\n\t\trun,\n\t\tasync (_state, _current, meta, reader) => {\n\t\t\tif (meta.intent.kind !== \"run\") throw new SessionInvariantError(\"Run operation has non-run intent\");\n\t\t\tconst entries = await reader.getEntries(meta.intent.promptEntryIds, drive.context);\n\t\t\tconst messages = meta.intent.promptEntryIds.map((id) => {\n\t\t\t\tconst entry = entries.get(id);\n\t\t\t\tif (entry?.type !== \"message\") {\n\t\t\t\t\tthrow new SessionInvariantError(`Run prompt entry ${id} is missing its message`);\n\t\t\t\t}\n\t\t\t\treturn entry.message;\n\t\t\t});\n\t\t\treturn { kind: \"return\", result: messages };\n\t\t},\n\t\tdrive.context,\n\t);\n\tif (prompt.kind === \"cancel_requested\") return { kind: \"continue\" };\n\n\tconst hook = await lane.hooks.runWithGate(\n\t\t\"before_run\",\n\t\t{ lane: lane.name, runId: drive.operationId, prompt: prompt.value, resources: lane.readConfig().resources },\n\t\tdrive.gate,\n\t\tdrive.context,\n\t);\n\tconst injected = hook?.messages ?? [];\n\tfor (const message of injected) {\n\t\tif (message.role === \"assistant\" && message.stopReason === \"pending\") {\n\t\t\tthrow new SessionInvariantError(\"before_run returned a pending assistant message\");\n\t\t}\n\t}\n\tconst reserved = injected.map((message) => ({ id: lane.session.idGenerator.next(), message }));\n\n\tconst result = await lane.continueOperation(\n\t\trun,\n\t\t(state, current) => {\n\t\t\tconst entries: NewEntry<MessageEntry>[] = chainEntries(\n\t\t\t\tstate.tipId,\n\t\t\t\treserved.map(({ id, message }) => ({ id, type: \"message\" as const, message })),\n\t\t\t);\n\t\t\tconst triggerEntryId = entries.at(-1)?.id ?? state.tipId;\n\t\t\tif (triggerEntryId === null) throw new SessionInvariantError(\"Run start has no trigger entry\");\n\t\t\tconst nextState: CheckpointOperation = {\n\t\t\t\t...operationScopeOf(current),\n\t\t\t\tat: \"checkpoint\",\n\t\t\t\tcontinuation: { kind: \"need_assistant\", overflowRecoveryUsed: false },\n\t\t\t\ttriggerEntryId,\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...(entries.length === 0 ? [] : [setValue(branchTip(lane.name), triggerEntryId)]),\n\t\t\t\t],\n\t\t\t\toperationState: nextState,\n\t\t\t\tlane: { tipId: triggerEntryId },\n\t\t\t\tmaterialize: () => ({ kind: \"continue\" }) as const,\n\t\t\t\tevents: (commit) => committedEntryEvents(entries, commit, lane.name, drive.operationId),\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n\treturn result.kind === \"cancel_requested\" ? { kind: \"continue\" } : result.value;\n}\n\n/** Advance one durable run boundary with at most one commit. */\nexport async function runCheckpoint<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\trun: CheckpointOperation,\n): Promise<ProcedureResult> {\n\tconst threshold = await prepareCompactionThreshold(lane, drive, run);\n\tif (threshold.kind === \"cancel_requested\") return { kind: \"continue\" };\n\tconst planned = await lane.continueOperation<CheckpointOperation, ProcedureResult | BoundaryFinishPending>(\n\t\trun,\n\t\tasync (state, current, _meta, reader) => {\n\t\t\tconst placement = await planBoundaryInbox(\n\t\t\t\tlane,\n\t\t\t\tdrive,\n\t\t\t\tstate,\n\t\t\t\tcurrent,\n\t\t\t\treader,\n\t\t\t\tstate.tipId,\n\t\t\t\tthreshold.value === undefined && current.continuation.kind === \"may_finish\",\n\t\t\t);\n\t\t\tif (placement.triggerEntryId !== undefined) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites: placement.writes,\n\t\t\t\t\toperationState: assistantReadyAtBoundary(lane, state, current, placement.triggerEntryId, false),\n\t\t\t\t\tlane: { tipId: placement.tipId, inbox: placement.inbox },\n\t\t\t\t\tmaterialize: () => ({ kind: \"continue\" }) as const,\n\t\t\t\t\tevents: (commit) => boundaryPlacementEvents(placement, commit, 0, lane.name, drive.operationId),\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (threshold.value !== undefined) {\n\t\t\t\tconst structural: SummaryDecidingOperation = {\n\t\t\t\t\t...operationScopeOf(current),\n\t\t\t\t\tat: \"summary.deciding\",\n\t\t\t\t\ttask: {\n\t\t\t\t\t\ttaskId: threshold.value.taskId,\n\t\t\t\t\t\treason: \"threshold\",\n\t\t\t\t\t\tboundary: {\n\t\t\t\t\t\t\tkind: \"resume_checkpoint\",\n\t\t\t\t\t\t\tresumeAfter: { continuation: current.continuation, triggerEntryId: current.triggerEntryId },\n\t\t\t\t\t\t},\n\t\t\t\t\t},\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...placement.writes,\n\t\t\t\t\t\tsetValue(\n\t\t\t\t\t\t\toperationPreparation(drive.operationId, threshold.value.taskId),\n\t\t\t\t\t\t\tthreshold.value.preparation,\n\t\t\t\t\t\t),\n\t\t\t\t\t],\n\t\t\t\t\toperationState: structural,\n\t\t\t\t\tlane: { tipId: placement.tipId, inbox: placement.inbox },\n\t\t\t\t\tmaterialize: () => ({ kind: \"continue\" }) as const,\n\t\t\t\t\tevents: (commit) => [\n\t\t\t\t\t\t...boundaryPlacementEvents(placement, commit, 0, lane.name, drive.operationId),\n\t\t\t\t\t\t{\n\t\t\t\t\t\t\ttype: \"compaction_start\",\n\t\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\t\treason: \"threshold\",\n\t\t\t\t\t\t\tstartedAt: commit.timestamp,\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\tif (current.continuation.kind === \"need_assistant\") {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\twrites: placement.writes,\n\t\t\t\t\toperationState: assistantReadyAtBoundary(\n\t\t\t\t\t\tlane,\n\t\t\t\t\t\tstate,\n\t\t\t\t\t\tcurrent,\n\t\t\t\t\t\tcurrent.triggerEntryId,\n\t\t\t\t\t\tcurrent.continuation.overflowRecoveryUsed,\n\t\t\t\t\t),\n\t\t\t\t\tlane: { tipId: placement.tipId, inbox: placement.inbox },\n\t\t\t\t\tmaterialize: () => ({ kind: \"continue\" }) as const,\n\t\t\t\t\tevents: (commit) => boundaryPlacementEvents(placement, commit, 0, lane.name, drive.operationId),\n\t\t\t\t};\n\t\t\t}\n\t\t\treturn {\n\t\t\t\tkind: \"return\",\n\t\t\t\tresult: { kind: \"finish_pending\", entryIds: placement.entries.map((entry) => entry.id) } as const,\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n\tif (planned.kind === \"cancel_requested\") return { kind: \"continue\" };\n\tif (planned.value.kind !== \"finish_pending\") return planned.value;\n\tif (run.continuation.kind !== \"may_finish\") {\n\t\tthrow new SessionInvariantError(\"Checkpoint finish mediation requires a finish continuation\");\n\t}\n\treturn finishRunBoundary(lane, drive, run, run.continuation, planned.value.entryIds);\n}\n"]}