{"version":3,"file":"deferred.d.ts","sourceRoot":"","sources":["../../../../src/harness/runtime/drive/deferred.ts"],"names":[],"mappings":"AAAA,OAAO,KAAK,EAAO,cAAc,EAAS,MAAM,uBAAuB,CAAC;AACxE,OAAO,EAAE,KAAK,OAAO,EAAwC,MAAM,kBAAkB,CAAC;AAItF,OAAO,EACN,KAAK,8BAA8B,EACnC,KAAK,0BAA0B,EAK/B,KAAK,aAAa,EAElB,MAAM,wBAAwB,CAAC;AAGhC,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,YAAY,CAAC;AACvC,OAAO,KAAK,EAA2B,KAAK,EAAE,eAAe,EAAE,MAAM,aAAa,CAAC;AAGnF,KAAK,YAAY,GAAG,0BAA0B,GAAG,8BAA8B,CAAC;AA0BhF,wBAAsB,wBAAwB,CAC7C,MAAM,EAAE,aAAa,EACrB,QAAQ,EAAE,YAAY,EACtB,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,cAAc,CAAC,CAqBzB;AA2LD,oFAAoF;AACpF,wBAAgB,oBAAoB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACvE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,QAAQ,EAAE,0BAA0B,GAClC,OAAO,CAAC,eAAe,CAAC,CAE1B;AAED,iGAAiG;AACjG,wBAAgB,mBAAmB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACtE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,QAAQ,EAAE,8BAA8B,GACtC,OAAO,CAAC,eAAe,CAAC,CAE1B;AAED,6DAA6D;AAC7D,wBAAgB,WAAW,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EAC9D,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,QAAQ,EAAE,0BAA0B,GAAG,8BAA8B,GACnE,OAAO,CAAC,eAAe,CAAC,CAI1B","sourcesContent":["import type { Api, DeferredHandle, Model } from \"@earendil-works/pi-ai\";\nimport { type Context, getTelemetryContext, withAbortSignal } from \"../../context.ts\";\nimport { consumeAssistantStream } from \"../../execution/assistant.ts\";\nimport { applyStreamOptionsPatch } from \"../../hooks.ts\";\nimport { SessionInvariantError } from \"../../session/session.ts\";\nimport {\n\ttype DeferredEffectPendingOperation,\n\ttype DeferredSuspendedOperation,\n\ttype LaneConfiguration,\n\ttype OperationError,\n\ttype OperationScope,\n\toperationScopeOf,\n\ttype SessionReader,\n\ttype SettledAssistantMessage,\n} from \"../../session/types.ts\";\nimport { deleteList, pendingAssistantFrames } from \"../../session/values.ts\";\nimport type { AgentHarnessStreamOptions } from \"../../types.ts\";\nimport type { Lane } from \"../lane.ts\";\nimport type { ContinueOperationResult, Drive, ProcedureResult } from \"../types.ts\";\nimport { openAssistantResponse, publishConfigurationFailure, publishResponse } from \"./response.ts\";\n\ntype DeferredLeaf = DeferredSuspendedOperation | DeferredEffectPendingOperation;\n/** Deferred effect payload without the run-wide scope fields. */\ntype EffectPendingFields = Omit<DeferredEffectPendingOperation, keyof OperationScope>;\n\ntype PreparedDeferredPoll = {\n\tkind: \"ready\";\n\tsource: DeferredHandle;\n\tmodel: Model<Api>;\n\tpoll: number;\n\tstreamOptions: AgentHarnessStreamOptions;\n};\n\ntype DeferredPreparation =\n\t| PreparedDeferredPoll\n\t| { kind: \"cancel_requested\" }\n\t| { kind: \"waiting\"; source: DeferredHandle }\n\t| { kind: \"configuration_failure\" };\n\nfunction configurationError(identity: LaneConfiguration[\"model\"]): OperationError {\n\treturn {\n\t\tcode: \"model_unavailable\",\n\t\tmessage: \"The configured model is unavailable in this process\",\n\t\tdetails: identity,\n\t};\n}\n\nexport async function readDeferredSourceHandle(\n\treader: SessionReader,\n\tdeferred: DeferredLeaf,\n\tcontext: Context,\n): Promise<DeferredHandle> {\n\tconst source = (await reader.getEntries([deferred.sourceEntryId], context)).get(deferred.sourceEntryId);\n\tif (\n\t\tsource?.type !== \"message\" ||\n\t\tsource.message.role !== \"assistant\" ||\n\t\tsource.message.stopReason !== \"deferred\" ||\n\t\tsource.message.deferred === undefined\n\t) {\n\t\tthrow new SessionInvariantError(`Deferred source ${deferred.sourceEntryId} is missing its assistant handle`);\n\t}\n\tconst handle = source.message.deferred;\n\tconst identity = deferred.configuration.model;\n\tif (\n\t\thandle.id.length === 0 ||\n\t\thandle.provider !== identity.provider ||\n\t\thandle.modelId !== identity.modelId ||\n\t\thandle.api !== source.message.api\n\t) {\n\t\tthrow new SessionInvariantError(`Deferred source ${deferred.sourceEntryId} has an invalid handle`);\n\t}\n\treturn handle;\n}\n\nasync function readSourceHandle<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredLeaf,\n): Promise<ContinueOperationResult<DeferredHandle>> {\n\treturn lane.continueOperation(\n\t\tdeferred,\n\t\tasync (_state, _current, _meta, reader) => ({\n\t\t\tkind: \"return\",\n\t\t\tresult: await readDeferredSourceHandle(reader, deferred, drive.context),\n\t\t}),\n\t\tdrive.context,\n\t);\n}\n\nasync function prepareDeferredPoll<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\texpected: DeferredLeaf,\n): Promise<DeferredPreparation> {\n\tconst source = await readSourceHandle(lane, drive, expected);\n\tif (source.kind === \"cancel_requested\") return source;\n\tif (drive.deferredPermits === 0) return { kind: \"waiting\", source: source.value };\n\n\tconst identity = expected.configuration.model;\n\tconst model = lane.models.getModel(identity.provider, identity.modelId);\n\tif (model === undefined) return { kind: \"configuration_failure\" };\n\tconst baseOptions: AgentHarnessStreamOptions = { ...expected.streamOptions, deferred: false };\n\tconst poll = expected.at === \"deferred.suspended\" ? expected.poll + 1 : expected.poll;\n\tconst beforeRequest = await lane.hooks.runWithGate(\n\t\t\"before_request\",\n\t\t{\n\t\t\tlane: lane.name,\n\t\t\trunId: drive.operationId,\n\t\t\tmodel,\n\t\t\tstep: \"deferred\",\n\t\t\tattempt: poll,\n\t\t\tstreamOptions: baseOptions,\n\t\t},\n\t\tdrive.gate,\n\t\tdrive.context,\n\t);\n\treturn {\n\t\tkind: \"ready\",\n\t\tsource: source.value,\n\t\tmodel,\n\t\tpoll,\n\t\tstreamOptions: {\n\t\t\t...(beforeRequest?.streamOptions === undefined\n\t\t\t\t? baseOptions\n\t\t\t\t: applyStreamOptionsPatch(baseOptions, beforeRequest.streamOptions)),\n\t\t\tdeferred: false,\n\t\t},\n\t};\n}\n\nasync function publishPollIntent<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredLeaf,\n\tprepared: PreparedDeferredPoll,\n\trecovery: boolean,\n): Promise<ContinueOperationResult<DeferredEffectPendingOperation>> {\n\tconst at = Date.now();\n\tconst next: EffectPendingFields = {\n\t\tat: \"deferred.effect_pending\",\n\t\tstepId: deferred.stepId,\n\t\tsourceEntryId: deferred.sourceEntryId,\n\t\tpoll: prepared.poll,\n\t\tresponseEntryId: lane.session.idGenerator.next(at),\n\t\tusageId: lane.session.idGenerator.next(at),\n\t\tconfiguration: deferred.configuration,\n\t\tstreamOptions: deferred.streamOptions,\n\t};\n\treturn lane.continueOperation(\n\t\tdeferred,\n\t\t(_state, current) => {\n\t\t\tconst nextState: DeferredEffectPendingOperation = { ...operationScopeOf(current), ...next };\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites:\n\t\t\t\t\tdeferred.at === \"deferred.effect_pending\"\n\t\t\t\t\t\t? [deleteList(pendingAssistantFrames(drive.operationId, deferred.responseEntryId))]\n\t\t\t\t\t\t: [],\n\t\t\t\toperationState: nextState,\n\t\t\t\tmaterialize: () => {\n\t\t\t\t\tdrive.deferredPermits--;\n\t\t\t\t\treturn nextState;\n\t\t\t\t},\n\t\t\t\tevents: () => [\n\t\t\t\t\t{\n\t\t\t\t\t\ttype: \"run_resume\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\t...(recovery ? { recovery: true as const } : {}),\n\t\t\t\t\t},\n\t\t\t\t\t{\n\t\t\t\t\t\ttype: \"turn_start\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\tturnId: `${next.stepId}:poll:${next.poll}`,\n\t\t\t\t\t\t...(recovery ? { recovery: true as const } : {}),\n\t\t\t\t\t},\n\t\t\t\t],\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n}\n\nasync function performDeferredPoll<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tprepared: PreparedDeferredPoll,\n\tintent: DeferredEffectPendingOperation,\n\trecovery: boolean,\n): Promise<SettledAssistantMessage> {\n\tconst response = openAssistantResponse(lane, drive, intent.responseEntryId, recovery);\n\tlet metadata: { status?: number; headers?: Record<string, string> } = {};\n\tconst admitted = withAbortSignal(drive.gate.signal, drive.context);\n\tconst stream = drive.gate.admit(() =>\n\t\tlane.models.streamDeferred(prepared.model, prepared.source, {\n\t\t\twait: 0,\n\t\t\tsignal: admitted.abortSignal,\n\t\t\ttelemetryContext: getTelemetryContext(admitted),\n\t\t\ttimeoutMs: prepared.streamOptions.timeoutMs,\n\t\t\tmaxRetries: prepared.streamOptions.maxRetries,\n\t\t\tmaxRetryDelayMs: prepared.streamOptions.maxRetryDelayMs,\n\t\t\theaders: prepared.streamOptions.headers,\n\t\t\tonPayload: async (payload, requestModel) => {\n\t\t\t\tconst result = await lane.hooks.runWithGate(\n\t\t\t\t\t\"before_payload\",\n\t\t\t\t\t{ lane: lane.name, runId: drive.operationId, model: requestModel, payload },\n\t\t\t\t\tdrive.gate,\n\t\t\t\t\tdrive.context,\n\t\t\t\t);\n\t\t\t\treturn result?.payload;\n\t\t\t},\n\t\t\tonResponse: (response) => {\n\t\t\t\tmetadata = { status: response.status, headers: response.headers };\n\t\t\t},\n\t\t}),\n\t);\n\n\ttry {\n\t\treturn await consumeAssistantStream(\n\t\t\tstream,\n\t\t\tresponse.observer,\n\t\t\t(message, context) => response.afterResponse(message, metadata, context),\n\t\t\tdrive.context,\n\t\t);\n\t} finally {\n\t\tawait response.close();\n\t}\n}\n\nasync function pollDeferred<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\texpected: DeferredLeaf,\n\trecovery: boolean,\n): Promise<ProcedureResult> {\n\tconst prepared = await prepareDeferredPoll(lane, drive, expected);\n\tif (prepared.kind === \"cancel_requested\") return { kind: \"continue\" };\n\tif (prepared.kind === \"waiting\") {\n\t\treturn {\n\t\t\tkind: \"waiting\",\n\t\t\toutcome: {\n\t\t\t\tkind: \"waiting\",\n\t\t\t\toperationId: drive.operationId,\n\t\t\t\treason: \"deferred\",\n\t\t\t\tdeferred: prepared.source,\n\t\t\t},\n\t\t};\n\t}\n\tif (prepared.kind === \"configuration_failure\") {\n\t\treturn publishConfigurationFailure(lane, drive, expected, configurationError(expected.configuration.model));\n\t}\n\n\tconst intent = await publishPollIntent(lane, drive, expected, prepared, recovery);\n\tif (intent.kind === \"cancel_requested\") return { kind: \"continue\" };\n\tconst response = await performDeferredPoll(lane, drive, prepared, intent.value, recovery);\n\treturn publishResponse(lane, drive, intent.value, response, recovery ? { recovery: true } : {});\n}\n\n/** Poll one durably suspended deferred response when this pass carries a permit. */\nexport function runDeferredSuspended<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredSuspendedOperation,\n): Promise<ProcedureResult> {\n\treturn pollDeferred(lane, drive, deferred, false);\n}\n\n/** Replace one orphaned unknown-outcome poll under fresh ids when this pass carries a permit. */\nexport function recoverDeferredPoll<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredEffectPendingOperation,\n): Promise<ProcedureResult> {\n\treturn pollDeferred(lane, drive, deferred, true);\n}\n\n/** Advance or report the wait for one deferred run phase. */\nexport function runDeferred<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredSuspendedOperation | DeferredEffectPendingOperation,\n): Promise<ProcedureResult> {\n\treturn deferred.at === \"deferred.suspended\"\n\t\t? runDeferredSuspended(lane, drive, deferred)\n\t\t: recoverDeferredPoll(lane, drive, deferred);\n}\n"]}