{"version":3,"file":"stream-fn.d.ts","sourceRoot":"","sources":["../../../src/core/shared-inference/stream-fn.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;GAkBG;AAGH,OAAO,KAAK,EAAE,QAAQ,EAAE,MAAM,+BAA+B,CAAC;AAC9D,OAAO,EAGN,KAAK,OAAO,EAEZ,KAAK,KAAK,EACV,YAAY,EACZ,MAAM,uBAAuB,CAAC;AAC/B,OAAO,KAAK,EAAE,iBAAiB,EAAE,MAAM,0BAA0B,CAAC;AAClE,OAAO,EAAiC,KAAK,4BAA4B,EAAE,MAAM,qBAAqB,CAAC;AACvG,OAAO,KAAK,EAAE,oBAAoB,EAAE,MAAM,cAAc,CAAC;AACzD,OAAO,KAAK,EAAE,wBAAwB,EAAE,MAAM,gBAAgB,CAAC;AAC/D,OAAO,KAAK,EAAE,iBAAiB,EAAE,0BAA0B,EAAE,MAAM,YAAY,CAAC;AAEhF,MAAM,WAAW,0BAA0B;IAC1C,cAAc,EAAE,MAAM,CAAC;IACvB,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,YAAY,CAAC,EAAE,MAAM,CAAC;IACtB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,QAAQ,CAAC,EAAE,iBAAiB,CAAC;IAC7B,UAAU,CAAC,EAAE,0BAA0B,CAAC;CACxC;AAED,MAAM,WAAW,wBAAwB;IACxC,yFAAyF;IACzF,SAAS,CAAC,EAAE,4BAA4B,CAAC;IACzC,mFAAmF;IACnF,SAAS,CAAC,EAAE,wBAAwB,CAAC;IACrC,OAAO,CAAC,EAAE,oBAAoB,CAAC;IAC/B,gDAAgD;IAChD,QAAQ,CAAC,EAAE,OAAO,YAAY,CAAC;IAC/B,wEAAwE;IACxE,cAAc,CAAC,EAAE,CAChB,KAAK,EAAE,KAAK,CAAC,GAAG,CAAC,EACjB,OAAO,EAAE,OAAO,EAChB,OAAO,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,KAC5B,0BAA0B,CAAC;IAChC,6DAA6D;IAC7D,mBAAmB,CAAC,EAAE,CAAC,OAAO,EAAE,OAAO,KAAK,MAAM,GAAG,SAAS,CAAC;IAC/D,2EAA2E;IAC3E,gBAAgB,CAAC,EAAE,CAAC,KAAK,EAAE,KAAK,CAAC,GAAG,CAAC,EAAE,OAAO,EAAE,OAAO,KAAK,OAAO,GAAG,OAAO,CAAC,OAAO,CAAC,CAAC;IACvF,wEAAwE;IACxE,uBAAuB,CAAC,EAAE,MAAM,OAAO,CAAC,IAAI,CAAC,CAAC;IAC9C,0DAA0D;IAC1D,UAAU,CAAC,EAAE,iBAAiB,CAAC;CAC/B;AAqDD,wBAAgB,uBAAuB,CAAC,OAAO,EAAE,wBAAwB,GAAG,QAAQ,CAwNnF","sourcesContent":["/**\n * Scheduled inference stream function (3.0.0 foundation).\n *\n * The provider integration seam. The agent loop already calls a `streamFn`\n * before touching the provider. This wrapper keeps scheduling/admission\n * separate from provider protocol handling:\n *\n *   request model inference\n *          ↓\n *   shared scheduler admission (durable queue + slot lease)\n *          ↓\n *   existing provider stream (streamSimple) — untouched\n *          ↓\n *   llama.cpp\n *\n * Non-shared providers are passed through unchanged. The scheduler is\n * capability/resource driven; OpenRouter/cloud providers are not forced through\n * local-Qwen slot semantics.\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport type { StreamFn } from \"@apholdings/jensen-agent-core\";\nimport {\n\ttype AssistantMessage,\n\ttype AssistantMessageEventStream,\n\ttype Context,\n\tcreateAssistantMessageEventStream,\n\ttype Model,\n\tstreamSimple,\n} from \"@apholdings/jensen-ai\";\nimport type { GovernanceService } from \"../governance/service.js\";\nimport { LocalSchedulerAdmissionClient, type SharedInferenceAdmissionPort } from \"./admission-port.js\";\nimport type { LocalSubagentRuntime } from \"./runtime.js\";\nimport type { SharedInferenceScheduler } from \"./scheduler.js\";\nimport type { InferencePriority, InferenceRequestDependency } from \"./types.js\";\n\nexport interface ScheduledStreamCorrelation {\n\tlogicalAgentId: string;\n\tmissionId?: string;\n\tassignmentId?: string;\n\texecutionId?: string;\n\tpriority?: InferencePriority;\n\tdependency?: InferenceRequestDependency;\n}\n\nexport interface ScheduledStreamFnOptions {\n\t/** Admission authority (local or remote). When omitted, a local scheduler is wrapped. */\n\tadmission?: SharedInferenceAdmissionPort;\n\t/** Backward-compatible local scheduler (wrapped into a local admission client). */\n\tscheduler?: SharedInferenceScheduler;\n\truntime?: LocalSubagentRuntime;\n\t/** Provider delegate (default streamSimple). */\n\tdelegate?: typeof streamSimple;\n\t/** Resolve logical-agent + mission/assignment/execution correlation. */\n\tgetCorrelation?: (\n\t\tmodel: Model<any>,\n\t\tcontext: Context,\n\t\toptions: Record<string, unknown>,\n\t) => ScheduledStreamCorrelation;\n\t/** Metadata-only input-token estimate (never the prompt). */\n\testimateInputTokens?: (context: Context) => number | undefined;\n\t/** Optional host/resource-pressure check used to park before admission. */\n\tresourcePressure?: (model: Model<any>, context: Context) => boolean | Promise<boolean>;\n\t/** Wait for a pressure change; the default yields to the event loop. */\n\twaitForResourcePressure?: () => Promise<void>;\n\t/** Optional Governance admission/accounting authority. */\n\tgovernance?: GovernanceService;\n}\n\nconst ZERO_USAGE: AssistantMessage[\"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 estimateTokensHeuristic(context: Context): number | undefined {\n\tlet chars = context.systemPrompt?.length ?? 0;\n\tfor (const message of context.messages ?? []) {\n\t\tconst content = (message as { content?: unknown }).content;\n\t\tif (typeof content === \"string\") chars += content.length;\n\t\telse if (Array.isArray(content)) {\n\t\t\tfor (const block of content) {\n\t\t\t\tif (\n\t\t\t\t\tblock &&\n\t\t\t\t\ttypeof block === \"object\" &&\n\t\t\t\t\t\"text\" in block &&\n\t\t\t\t\ttypeof (block as { text?: unknown }).text === \"string\"\n\t\t\t\t) {\n\t\t\t\t\tchars += (block as { text: string }).text.length;\n\t\t\t\t}\n\t\t\t}\n\t\t}\n\t}\n\tfor (const tool of context.tools ?? []) {\n\t\tchars += (tool.name?.length ?? 0) + (JSON.stringify(tool.parameters ?? {}).length ?? 0);\n\t}\n\treturn Math.max(1, Math.ceil(chars / 4));\n}\n\nfunction errorStream(model: Model<any>, message: string): AssistantMessageEventStream {\n\tconst stream = createAssistantMessageEventStream();\n\tconst output: AssistantMessage = {\n\t\trole: \"assistant\",\n\t\tcontent: [],\n\t\tapi: model.api,\n\t\tprovider: model.provider,\n\t\tmodel: model.id,\n\t\tusage: { ...ZERO_USAGE },\n\t\tstopReason: \"error\",\n\t\terrorMessage: message,\n\t\ttimestamp: Date.now(),\n\t};\n\tstream.push({ type: \"error\", reason: \"error\", error: output });\n\tstream.end();\n\treturn stream;\n}\n\nexport function createScheduledStreamFn(options: ScheduledStreamFnOptions): StreamFn {\n\tconst delegate = options.delegate ?? streamSimple;\n\tconst admission =\n\t\toptions.admission ?? (options.scheduler ? new LocalSchedulerAdmissionClient(options.scheduler) : undefined);\n\tif (!admission && !options.governance) {\n\t\tthrow new Error(\"createScheduledStreamFn requires an admission port, local scheduler, or Governance service\");\n\t}\n\n\treturn async (model, context, streamOptions) => {\n\t\tconst resource = admission?.resourceFor(model);\n\t\tconst correlation = options.getCorrelation\n\t\t\t? options.getCorrelation(model, context, (streamOptions ?? {}) as Record<string, unknown>)\n\t\t\t: { logicalAgentId: \"unknown\" };\n\t\tconst logicalAgentId = correlation.logicalAgentId;\n\t\tconst inferenceRequestId = `inference_${randomUUID()}`;\n\t\tif (options.governance && correlation.missionId) {\n\t\t\tconst governanceAdmission = await options.governance.admitInference({\n\t\t\t\tmissionId: correlation.missionId,\n\t\t\t\teventId: inferenceRequestId,\n\t\t\t\tprovider: model.provider,\n\t\t\t\tmodel: model.id,\n\t\t\t\tatMs: Date.now(),\n\t\t\t});\n\t\t\tif (!governanceAdmission.allowed)\n\t\t\t\treturn errorStream(model, `Governance admission denied: ${governanceAdmission.reason ?? \"unknown\"}`);\n\t\t}\n\n\t\tif (!resource) {\n\t\t\tconst delegateStream = delegate(model, context, streamOptions);\n\t\t\tif (!options.governance || !correlation.missionId) return delegateStream;\n\t\t\tconst stream = createAssistantMessageEventStream();\n\t\t\tvoid (async () => {\n\t\t\t\ttry {\n\t\t\t\t\tfor await (const event of delegateStream) {\n\t\t\t\t\t\tif (event.type === \"done\") {\n\t\t\t\t\t\t\tawait options.governance?.recordInferenceResult({\n\t\t\t\t\t\t\t\tmissionId: correlation.missionId!,\n\t\t\t\t\t\t\t\teventId: `${inferenceRequestId}:result`,\n\t\t\t\t\t\t\t\tprovider: model.provider,\n\t\t\t\t\t\t\t\tmodel: model.id,\n\t\t\t\t\t\t\t\tinputTokens: event.message.usage?.input,\n\t\t\t\t\t\t\t\toutputTokens: event.message.usage?.output,\n\t\t\t\t\t\t\t\tcostUsd: model.provider.startsWith(\"llamacpp-\") ? undefined : event.message.usage?.cost.total,\n\t\t\t\t\t\t\t\tcostStatus: model.provider.startsWith(\"llamacpp-\")\n\t\t\t\t\t\t\t\t\t? undefined\n\t\t\t\t\t\t\t\t\t: event.message.usage?.cost.total === undefined\n\t\t\t\t\t\t\t\t\t\t? \"UNKNOWN\"\n\t\t\t\t\t\t\t\t\t\t: \"KNOWN\",\n\t\t\t\t\t\t\t\tatMs: Date.now(),\n\t\t\t\t\t\t\t});\n\t\t\t\t\t\t}\n\t\t\t\t\t\tstream.push(event);\n\t\t\t\t\t\tif (event.type === \"done\" || event.type === \"error\") return;\n\t\t\t\t\t}\n\t\t\t\t\tstream.end();\n\t\t\t\t} catch (error) {\n\t\t\t\t\tstream.push({\n\t\t\t\t\t\ttype: \"error\",\n\t\t\t\t\t\treason: \"error\",\n\t\t\t\t\t\terror: {\n\t\t\t\t\t\t\trole: \"assistant\",\n\t\t\t\t\t\t\tcontent: [],\n\t\t\t\t\t\t\tapi: model.api,\n\t\t\t\t\t\t\tprovider: model.provider,\n\t\t\t\t\t\t\tmodel: model.id,\n\t\t\t\t\t\t\tusage: { ...ZERO_USAGE },\n\t\t\t\t\t\t\tstopReason: \"error\",\n\t\t\t\t\t\t\terrorMessage: error instanceof Error ? error.message : String(error),\n\t\t\t\t\t\t\ttimestamp: Date.now(),\n\t\t\t\t\t\t},\n\t\t\t\t\t});\n\t\t\t\t\tstream.end();\n\t\t\t\t}\n\t\t\t})();\n\t\t\treturn stream;\n\t\t}\n\n\t\tif (options.resourcePressure && (await options.resourcePressure(model, context))) {\n\t\t\tawait options.runtime?.park(logicalAgentId, \"resource pressure\");\n\t\t\tawait (options.waitForResourcePressure ?? (() => new Promise((resolve) => setImmediate(resolve)))());\n\t\t\tawait options.runtime?.resume(logicalAgentId);\n\t\t}\n\n\t\tawait options.runtime?.transition(logicalAgentId, \"WAITING_INFERENCE\", {\n\t\t\twaitingReason: \"shared inference admission\",\n\t\t\tpendingInferenceRequestId: inferenceRequestId,\n\t\t});\n\n\t\tconst acquired = await admission!.acquire({\n\t\t\tlogicalAgentId,\n\t\t\tresource,\n\t\t\tmodel,\n\t\t\tinferenceRequestId,\n\t\t\tmissionId: correlation.missionId,\n\t\t\tassignmentId: correlation.assignmentId,\n\t\t\texecutionId: correlation.executionId,\n\t\t\tpriority: correlation.priority,\n\t\t\tdependency: correlation.dependency,\n\t\t\testimatedInputTokens: options.estimateInputTokens\n\t\t\t\t? options.estimateInputTokens(context)\n\t\t\t\t: estimateTokensHeuristic(context),\n\t\t\tmaxOutputTokens: typeof streamOptions?.maxTokens === \"number\" ? streamOptions.maxTokens : undefined,\n\t\t\tsignal: streamOptions?.signal,\n\t\t});\n\n\t\tif (acquired.status !== \"admitted\") {\n\t\t\tconst reason = acquired.status === \"cancelled\" ? acquired.reason : `inference admission ${acquired.status}`;\n\t\t\tawait options.runtime?.transition(logicalAgentId, \"RUNNABLE\", {\n\t\t\t\twaitingReason: undefined,\n\t\t\t\tpendingInferenceRequestId: undefined,\n\t\t\t});\n\t\t\treturn errorStream(model, `Shared inference admission failed: ${reason}`);\n\t\t}\n\n\t\tawait options.runtime?.transition(logicalAgentId, \"RUNNING_INFERENCE\", {\n\t\t\twaitingReason: undefined,\n\t\t\tpendingInferenceRequestId: acquired.admitted.inferenceRequestId,\n\t\t});\n\n\t\tconst delegateStream = delegate(model, context, streamOptions);\n\n\t\t// Wrap the delegate stream so the scheduler release happens BEFORE the\n\t\t// terminal event / `result()` resolves. A fire-and-forget release is lost\n\t\t// when a CLI process exits immediately after the final message.\n\t\tconst stream = createAssistantMessageEventStream();\n\t\tlet released = false;\n\t\tconst syntheticFailure = (errorMessage: string): AssistantMessage => ({\n\t\t\trole: \"assistant\",\n\t\t\tcontent: [],\n\t\t\tapi: model.api,\n\t\t\tprovider: model.provider,\n\t\t\tmodel: model.id,\n\t\t\tusage: { ...ZERO_USAGE },\n\t\t\tstopReason: \"error\",\n\t\t\terrorMessage,\n\t\t\ttimestamp: Date.now(),\n\t\t});\n\n\t\tconst releaseOnce = async (message: AssistantMessage, failed: boolean): Promise<Error | undefined> => {\n\t\t\tif (released) return undefined;\n\t\t\treleased = true;\n\t\t\tlet accountingError: Error | undefined;\n\t\t\tif (options.governance && correlation.missionId) {\n\t\t\t\ttry {\n\t\t\t\t\tawait options.governance.recordInferenceResult({\n\t\t\t\t\t\tmissionId: correlation.missionId,\n\t\t\t\t\t\teventId: `${inferenceRequestId}:result`,\n\t\t\t\t\t\tprovider: model.provider,\n\t\t\t\t\t\tmodel: model.id,\n\t\t\t\t\t\tinputTokens: message.usage?.input,\n\t\t\t\t\t\toutputTokens: message.usage?.output,\n\t\t\t\t\t\tcostUsd: model.provider.startsWith(\"llamacpp-\") ? undefined : message.usage?.cost.total,\n\t\t\t\t\t\tcostStatus: model.provider.startsWith(\"llamacpp-\")\n\t\t\t\t\t\t\t? undefined\n\t\t\t\t\t\t\t: message.usage?.cost.total === undefined\n\t\t\t\t\t\t\t\t? \"UNKNOWN\"\n\t\t\t\t\t\t\t\t: \"KNOWN\",\n\t\t\t\t\t\tatMs: Date.now(),\n\t\t\t\t\t});\n\t\t\t\t} catch (error) {\n\t\t\t\t\taccountingError = error instanceof Error ? error : new Error(String(error));\n\t\t\t\t}\n\t\t\t}\n\t\t\ttry {\n\t\t\t\tawait admission!.release(acquired.admitted, {\n\t\t\t\t\tstate: failed || accountingError ? \"FAILED\" : \"COMPLETED\",\n\t\t\t\t\tusage:\n\t\t\t\t\t\tfailed || accountingError\n\t\t\t\t\t\t\t? undefined\n\t\t\t\t\t\t\t: { input: message.usage?.input, output: message.usage?.output },\n\t\t\t\t\terrorMessage: accountingError?.message ?? (failed ? message.errorMessage : undefined),\n\t\t\t\t});\n\t\t\t\tif (failed || accountingError)\n\t\t\t\t\tawait options.runtime?.transition(logicalAgentId, \"RUNNABLE\", {\n\t\t\t\t\t\twaitingReason: undefined,\n\t\t\t\t\t\tpendingInferenceRequestId: undefined,\n\t\t\t\t\t});\n\t\t\t\telse await options.runtime?.resume(logicalAgentId);\n\t\t\t} catch (error) {\n\t\t\t\treturn error instanceof Error ? error : new Error(String(error));\n\t\t\t}\n\t\t\treturn accountingError;\n\t\t};\n\n\t\tvoid (async () => {\n\t\t\ttry {\n\t\t\t\tfor await (const event of delegateStream) {\n\t\t\t\t\tif (event.type === \"done\") {\n\t\t\t\t\t\tconst releaseError = await releaseOnce(event.message, false);\n\t\t\t\t\t\tif (releaseError) {\n\t\t\t\t\t\t\tstream.push({ type: \"error\", reason: \"error\", error: syntheticFailure(releaseError.message) });\n\t\t\t\t\t\t\tstream.end();\n\t\t\t\t\t\t\treturn;\n\t\t\t\t\t\t}\n\t\t\t\t\t\tstream.push(event);\n\t\t\t\t\t\treturn;\n\t\t\t\t\t}\n\t\t\t\t\tif (event.type === \"error\") {\n\t\t\t\t\t\tawait releaseOnce(event.error, true);\n\t\t\t\t\t\tstream.push(event);\n\t\t\t\t\t\treturn;\n\t\t\t\t\t}\n\t\t\t\t\tstream.push(event);\n\t\t\t\t}\n\t\t\t\t// Delegate ended without a terminal event: release honestly, never\n\t\t\t\t// fabricate a completion.\n\t\t\t\tawait releaseOnce(syntheticFailure(\"provider stream ended without a terminal event\"), true);\n\t\t\t\tstream.end();\n\t\t\t} catch (error) {\n\t\t\t\tawait releaseOnce(syntheticFailure(error instanceof Error ? error.message : String(error)), true);\n\t\t\t\tstream.end();\n\t\t\t}\n\t\t})();\n\n\t\treturn stream;\n\t};\n}\n"]}