{"version":3,"file":"recovery.d.ts","sourceRoot":"","sources":["../../../../src/harness/runtime/drive/recovery.ts"],"names":[],"mappings":"AACA,OAAO,KAAK,EACX,+BAA+B,EAC/B,8BAA8B,EAG9B,MAAM,wBAAwB,CAAC;AAChC,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,YAAY,CAAC;AAEvC,OAAO,KAAK,EAAE,KAAK,EAAE,eAAe,EAAE,MAAM,aAAa,CAAC;AAiC1D,kHAAkH;AAClH,wBAAsB,0BAA0B,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACnF,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,UAAU,EAAE,+BAA+B,GACzC,OAAO,CAAC,eAAe,CAAC,CAoC1B;AAED,uGAAuG;AACvG,wBAAsB,+BAA+B,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACxF,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,MAAM,EAAE,+BAA+B,GAAG,8BAA8B,GACtE,OAAO,CAAC,eAAe,CAAC,CAmC1B","sourcesContent":["import { type AssistantMessage, reduceAssistantMessageFrames } from \"@earendil-works/pi-ai\";\nimport type {\n\tAssistantEffectPendingOperation,\n\tDeferredEffectPendingOperation,\n\tLaneConfiguration,\n\tSettledAssistantMessage,\n} from \"../../session/types.ts\";\nimport type { Lane } from \"../lane.ts\";\nimport { readAssistantFrames } from \"../progress.ts\";\nimport type { Drive, ProcedureResult } from \"../types.ts\";\nimport { publishResponse } from \"./response.ts\";\n\nconst ZERO_USAGE = {\n\tinput: 0,\n\toutput: 0,\n\tcacheRead: 0,\n\tcacheWrite: 0,\n\ttotalTokens: 0,\n\tcost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },\n};\n\nfunction interruptedAssistantMessage(\n\tidentity: LaneConfiguration[\"model\"],\n\tpartial: AssistantMessage | undefined,\n): SettledAssistantMessage {\n\tconst warning =\n\t\t\"Assistant request was interrupted. The preceding content is the latest committed partial; newer live output may be missing and the external outcome is unknown.\";\n\treturn partial === undefined\n\t\t? {\n\t\t\t\trole: \"assistant\",\n\t\t\t\tcontent: [],\n\t\t\t\tapi: \"unknown\",\n\t\t\t\tprovider: identity.provider,\n\t\t\t\tmodel: identity.modelId,\n\t\t\t\tusage: ZERO_USAGE,\n\t\t\t\tstopReason: \"error\",\n\t\t\t\terrorMessage: warning,\n\t\t\t\ttimestamp: Date.now(),\n\t\t\t}\n\t\t: { ...partial, usage: ZERO_USAGE, stopReason: \"error\", errorMessage: warning };\n}\n\n/** Settle an orphaned assistant request from its bounded committed frame prefix without another provider call. */\nexport async function recoverAssistantGeneration<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tgeneration: AssistantEffectPendingOperation,\n): Promise<ProcedureResult> {\n\tconst frames = await lane.continueOperation(\n\t\tgeneration,\n\t\tasync (_state, _current, _meta, reader) => ({\n\t\t\tkind: \"return\",\n\t\t\tresult: await readAssistantFrames(reader, drive.operationId, generation.responseEntryId, drive.context),\n\t\t}),\n\t\tdrive.context,\n\t);\n\tif (frames.kind === \"cancel_requested\") return { kind: \"continue\" };\n\n\tconst message = interruptedAssistantMessage(\n\t\tgeneration.generationContext.configuration.model,\n\t\treduceAssistantMessageFrames(frames.value),\n\t);\n\tawait lane.emitBatch(\n\t\t[\n\t\t\t{\n\t\t\t\ttype: \"message_start\",\n\t\t\t\tlane: lane.name,\n\t\t\t\trunId: drive.operationId,\n\t\t\t\tmessage,\n\t\t\t\trecovery: true,\n\t\t\t},\n\t\t\t{\n\t\t\t\ttype: \"message_end\",\n\t\t\t\tlane: lane.name,\n\t\t\t\trunId: drive.operationId,\n\t\t\t\tmessage,\n\t\t\t\tentryId: generation.responseEntryId,\n\t\t\t\trecovery: true,\n\t\t\t},\n\t\t],\n\t\tdrive.context,\n\t);\n\treturn publishResponse(lane, drive, generation, message, { recovery: true });\n}\n\n/** Synthetically settle one cancelled orphaned assistant or deferred effect under its reserved ids. */\nexport async function recoverCancelledAssistantEffect<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\teffect: AssistantEffectPendingOperation | DeferredEffectPendingOperation,\n): Promise<ProcedureResult> {\n\tconst frames = await lane.settleOperation(\n\t\teffect,\n\t\tasync (_state, _current, _meta, reader) => ({\n\t\t\tkind: \"return\",\n\t\t\tresult: await readAssistantFrames(reader, drive.operationId, effect.responseEntryId, drive.context),\n\t\t}),\n\t\tdrive.context,\n\t);\n\tconst identity =\n\t\teffect.at === \"assistant.effect_pending\"\n\t\t\t? effect.generationContext.configuration.model\n\t\t\t: effect.configuration.model;\n\tconst message = interruptedAssistantMessage(identity, reduceAssistantMessageFrames(frames));\n\tawait lane.emitBatch(\n\t\t[\n\t\t\t{\n\t\t\t\ttype: \"message_start\",\n\t\t\t\tlane: lane.name,\n\t\t\t\trunId: drive.operationId,\n\t\t\t\tmessage,\n\t\t\t\trecovery: true,\n\t\t\t},\n\t\t\t{\n\t\t\t\ttype: \"message_end\",\n\t\t\t\tlane: lane.name,\n\t\t\t\trunId: drive.operationId,\n\t\t\t\tmessage,\n\t\t\t\tentryId: effect.responseEntryId,\n\t\t\t\trecovery: true,\n\t\t\t},\n\t\t],\n\t\tdrive.context,\n\t);\n\treturn publishResponse(lane, drive, effect, message, { recovery: true });\n}\n"]}