{"version":3,"file":"assistant.d.ts","sourceRoot":"","sources":["../../../src/harness/execution/assistant.ts"],"names":[],"mappings":"AAAA,OAAO,KAAK,EACX,OAAO,IAAI,SAAS,EACpB,GAAG,EACH,gBAAgB,EAChB,qBAAqB,EACrB,2BAA2B,EAC3B,OAAO,EACP,KAAK,EACL,mBAAmB,EACnB,IAAI,EACJ,MAAM,uBAAuB,CAAC;AAC/B,OAAO,KAAK,EAAE,YAAY,EAAE,aAAa,EAAE,MAAM,gBAAgB,CAAC;AAClE,OAAO,EAAE,KAAK,OAAO,EAAuB,MAAM,eAAe,CAAC;AAClE,OAAO,KAAK,EAAE,uBAAuB,EAAE,MAAM,qBAAqB,CAAC;AACnE,OAAO,KAAK,EAAE,yBAAyB,EAAE,MAAM,aAAa,CAAC;AAG7D,qFAAqF;AACrF,MAAM,WAAW,yBAAyB;IACzC,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,OAAO,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,MAAM,CAAC,CAAC;CACjC;AAED,iEAAiE;AACjE,MAAM,WAAW,uBAAuB;IACvC,KAAK,CACJ,OAAO,EAAE,gBAAgB,EACzB,KAAK,EAAE,OAAO,CAAC,qBAAqB,EAAE;QAAE,IAAI,EAAE,OAAO,CAAA;KAAE,CAAC,EACxD,OAAO,EAAE,OAAO,GACd,IAAI,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;IACxB,MAAM,CAAC,OAAO,EAAE,gBAAgB,EAAE,KAAK,EAAE,qBAAqB,EAAE,OAAO,EAAE,OAAO,GAAG,IAAI,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;IACxG,GAAG,CAAC,OAAO,EAAE,uBAAuB,EAAE,OAAO,EAAE,OAAO,GAAG,IAAI,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;CAC9E;AAED,6EAA6E;AAC7E,MAAM,WAAW,4BAA4B;IAC5C,KAAK,EAAE,KAAK,CAAC,GAAG,CAAC,CAAC;IAClB,YAAY,EAAE,MAAM,CAAC;IACrB,KAAK,CAAC,EAAE,IAAI,EAAE,CAAC;IACf,aAAa,EAAE,aAAa,CAAC;IAC7B,aAAa,EAAE,yBAAyB,CAAC;IACzC,gBAAgB,CAAC,EAAE,CAClB,cAAc,EAAE;QAAE,QAAQ,EAAE,YAAY,EAAE,CAAC;QAAC,YAAY,EAAE,MAAM,CAAA;KAAE,EAClE,OAAO,EAAE,OAAO,KACZ,OAAO,CAAC;QAAE,QAAQ,EAAE,YAAY,EAAE,CAAC;QAAC,YAAY,EAAE,MAAM,CAAA;KAAE,CAAC,CAAC;IACjE,kBAAkB,EAAE,CAAC,QAAQ,EAAE,YAAY,EAAE,EAAE,OAAO,EAAE,OAAO,KAAK,OAAO,EAAE,GAAG,OAAO,CAAC,OAAO,EAAE,CAAC,CAAC;IACnG,aAAa,CAAC,EAAE,CACf,OAAO,EAAE,OAAO,EAChB,KAAK,EAAE,KAAK,CAAC,GAAG,CAAC,EACjB,OAAO,EAAE,OAAO,KACZ,OAAO,GAAG,SAAS,GAAG,OAAO,CAAC,OAAO,GAAG,SAAS,CAAC,CAAC;IACxD,aAAa,CAAC,EAAE,CACf,OAAO,EAAE,uBAAuB,EAChC,QAAQ,EAAE,yBAAyB,EACnC,OAAO,EAAE,OAAO,KACZ,OAAO,CAAC,uBAAuB,CAAC,CAAC;IACtC,OAAO,CACN,SAAS,EAAE,SAAS,EACpB,OAAO,EAAE,mBAAmB,EAC5B,OAAO,EAAE,OAAO,GACd,2BAA2B,GAAG,OAAO,CAAC,2BAA2B,CAAC,CAAC;IACtE,QAAQ,EAAE,uBAAuB,CAAC;CAClC;AAoCD,wBAAsB,sBAAsB,CAC3C,MAAM,EAAE,2BAA2B,EACnC,QAAQ,EAAE,uBAAuB,EACjC,aAAa,EACV,CAAC,CAAC,OAAO,EAAE,uBAAuB,EAAE,OAAO,EAAE,OAAO,KAAK,OAAO,CAAC,uBAAuB,CAAC,CAAC,GAC1F,SAAS,EACZ,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,uBAAuB,CAAC,CA2BlC;AAED,gFAAgF;AAChF,wBAAsB,sBAAsB,CAC3C,QAAQ,EAAE,YAAY,EAAE,EACxB,MAAM,EAAE,4BAA4B,EACpC,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,uBAAuB,CAAC,CAmClC","sourcesContent":["import type {\n\tContext as AiContext,\n\tApi,\n\tAssistantMessage,\n\tAssistantMessageEvent,\n\tAssistantMessageEventStream,\n\tMessage,\n\tModel,\n\tSimpleStreamOptions,\n\tTool,\n} from \"@earendil-works/pi-ai\";\nimport type { AgentMessage, ThinkingLevel } from \"../../types.ts\";\nimport { type Context, getTelemetryContext } from \"../context.ts\";\nimport type { SettledAssistantMessage } from \"../session/types.ts\";\nimport type { AgentHarnessStreamOptions } from \"../types.ts\";\nimport { AbortRequested } from \"./effect-gate.ts\";\n\n/** HTTP response metadata captured before the provider response body is consumed. */\nexport interface AssistantResponseMetadata {\n\tstatus?: number;\n\theaders?: Record<string, string>;\n}\n\n/** Process-local lifecycle observer for one assistant stream. */\nexport interface AssistantStreamObserver {\n\tstart(\n\t\tmessage: AssistantMessage,\n\t\tevent: Extract<AssistantMessageEvent, { type: \"start\" }>,\n\t\tcontext: Context,\n\t): void | Promise<void>;\n\tupdate(message: AssistantMessage, event: AssistantMessageEvent, context: Context): void | Promise<void>;\n\tend(message: SettledAssistantMessage, context: Context): void | Promise<void>;\n}\n\n/** Executable inputs for one already-approved assistant provider request. */\nexport interface HarnessAssistantStreamConfig {\n\tmodel: Model<Api>;\n\tsystemPrompt: string;\n\ttools?: Tool[];\n\tthinkingLevel: ThinkingLevel;\n\tstreamOptions: AgentHarnessStreamOptions;\n\ttransformContext?: (\n\t\trequestContext: { messages: AgentMessage[]; systemPrompt: string },\n\t\tcontext: Context,\n\t) => Promise<{ messages: AgentMessage[]; systemPrompt: string }>;\n\ttoProviderMessages: (messages: AgentMessage[], context: Context) => Message[] | Promise<Message[]>;\n\tbeforePayload?: (\n\t\tpayload: unknown,\n\t\tmodel: Model<Api>,\n\t\tcontext: Context,\n\t) => unknown | undefined | Promise<unknown | undefined>;\n\tafterResponse?: (\n\t\tmessage: SettledAssistantMessage,\n\t\tmetadata: AssistantResponseMetadata,\n\t\tcontext: Context,\n\t) => Promise<SettledAssistantMessage>;\n\trequest(\n\t\taiContext: AiContext,\n\t\toptions: SimpleStreamOptions,\n\t\tcontext: Context,\n\t): AssistantMessageEventStream | Promise<AssistantMessageEventStream>;\n\tobserver: AssistantStreamObserver;\n}\n\nfunction createRequestOptions(\n\tconfig: HarnessAssistantStreamConfig,\n\tcaptureMetadata: (metadata: AssistantResponseMetadata) => void,\n\tcontext: Context,\n): SimpleStreamOptions {\n\tconst options = config.streamOptions;\n\treturn {\n\t\ttransport: options.transport,\n\t\ttimeoutMs: options.timeoutMs,\n\t\tmaxRetries: options.maxRetries,\n\t\tmaxRetryDelayMs: options.maxRetryDelayMs,\n\t\theaders: options.headers,\n\t\tmetadata: options.metadata,\n\t\tcacheRetention: options.cacheRetention,\n\t\tdeferred: options.deferred,\n\t\t...(config.thinkingLevel === \"off\" ? {} : { reasoning: config.thinkingLevel }),\n\t\tsignal: context.abortSignal,\n\t\ttelemetryContext: getTelemetryContext(context),\n\t\tonPayload:\n\t\t\tconfig.beforePayload === undefined\n\t\t\t\t? undefined\n\t\t\t\t: (payload, model) => config.beforePayload?.(payload, model, context),\n\t\tonResponse: (response) => {\n\t\t\tcaptureMetadata({ status: response.status, headers: response.headers });\n\t\t},\n\t};\n}\n\nfunction isUpdateEvent(\n\tevent: AssistantMessageEvent,\n): event is Exclude<AssistantMessageEvent, { type: \"start\" | \"done\" | \"error\" }> {\n\treturn event.type !== \"start\" && event.type !== \"done\" && event.type !== \"error\";\n}\n\nexport async function consumeAssistantStream(\n\tstream: AssistantMessageEventStream,\n\tobserver: AssistantStreamObserver,\n\tafterResponse:\n\t\t| ((message: SettledAssistantMessage, context: Context) => Promise<SettledAssistantMessage>)\n\t\t| undefined,\n\tcontext: Context,\n): Promise<SettledAssistantMessage> {\n\tlet started = false;\n\tfor await (const event of stream) {\n\t\tif (event.type === \"start\") {\n\t\t\tif (started) throw new Error(\"Assistant message stream emitted more than one start event\");\n\t\t\tstarted = true;\n\t\t\tawait observer.start({ ...event.partial }, event, context);\n\t\t} else if (isUpdateEvent(event)) {\n\t\t\tif (!started) throw new Error(`Assistant message stream emitted ${event.type} before start`);\n\t\t\tawait observer.update({ ...event.partial }, event, context);\n\t\t} else if (event.type === \"done\" && !started) {\n\t\t\tthrow new Error(\"Assistant message stream emitted done before start\");\n\t\t}\n\t}\n\n\tconst settled = (await stream.result()) as SettledAssistantMessage;\n\tlet finalMessage = settled;\n\tif (afterResponse !== undefined) {\n\t\ttry {\n\t\t\tfinalMessage = await afterResponse(settled, context);\n\t\t} catch (error) {\n\t\t\tif (!(error instanceof AbortRequested)) throw error;\n\t\t\tawait error.cancellation;\n\t\t}\n\t}\n\tawait observer.end(finalMessage, context);\n\treturn finalMessage;\n}\n\n/** Stream one assistant response without mutating the caller's message list. */\nexport async function streamHarnessAssistant(\n\tmessages: AgentMessage[],\n\tconfig: HarnessAssistantStreamConfig,\n\tcontext: Context,\n): Promise<SettledAssistantMessage> {\n\tlet requestContext = { messages: messages.slice(), systemPrompt: config.systemPrompt };\n\tif (config.transformContext) {\n\t\trequestContext = await config.transformContext(requestContext, context);\n\t}\n\n\tconst providerMessages = await config.toProviderMessages(requestContext.messages, context);\n\tconst aiContext: AiContext = {\n\t\tsystemPrompt: requestContext.systemPrompt,\n\t\tmessages: providerMessages,\n\t\ttools: config.tools,\n\t};\n\n\tlet metadata: AssistantResponseMetadata = {};\n\tconst stream = await config.request(\n\t\taiContext,\n\t\tcreateRequestOptions(\n\t\t\tconfig,\n\t\t\t(nextMetadata) => {\n\t\t\t\tmetadata = nextMetadata;\n\t\t\t},\n\t\t\tcontext,\n\t\t),\n\t\tcontext,\n\t);\n\n\tconst afterResponse = config.afterResponse;\n\treturn consumeAssistantStream(\n\t\tstream,\n\t\tconfig.observer,\n\t\tafterResponse === undefined\n\t\t\t? undefined\n\t\t\t: (message, afterContext) => afterResponse(message, metadata, afterContext),\n\t\tcontext,\n\t);\n}\n"]}