{"version":3,"file":"remote-admission-client.d.ts","sourceRoot":"","sources":["../../../src/core/shared-inference/remote-admission-client.ts"],"names":[],"mappings":"AAAA;;;;;;;;;GASG;AAGH,OAAO,KAAK,EAAE,uBAAuB,EAAE,qBAAqB,EAAE,4BAA4B,EAAE,MAAM,qBAAqB,CAAC;AASxH,OAAO,KAAK,EACX,uBAAuB,EACvB,iBAAiB,EAGjB,uBAAuB,EACvB,uBAAuB,EACvB,MAAM,YAAY,CAAC;AAEpB,MAAM,WAAW,4BAA4B;IAC5C,OAAO,EAAE,MAAM,CAAC;IAChB,KAAK,EAAE,MAAM,CAAC;IACd,WAAW,EAAE,MAAM,CAAC;IACpB,SAAS,EAAE,uBAAuB,EAAE,CAAC;IACrC,gDAAgD;IAChD,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,gCAAgC;IAChC,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,gCAAgC;IAChC,SAAS,CAAC,EAAE,OAAO,KAAK,CAAC;CACzB;AAED,qBAAa,8BAA+B,YAAW,4BAA4B;IAClF,QAAQ,CAAC,IAAI,WAAqB;IAElC,OAAO,CAAC,QAAQ,CAAC,QAAQ,CAAS;IAClC,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAS;IAChC,OAAO,CAAC,QAAQ,CAAC,YAAY,CAAS;IACtC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAA4B;IACvD,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAS;IACjC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAS;IACpC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAe;IAE1C,YAAY,OAAO,EAAE,4BAA4B,EAQhD;IAED,WAAW,CAAC,KAAK,EAAE;QAAE,QAAQ,EAAE,MAAM,CAAC;QAAC,EAAE,EAAE,MAAM,CAAA;KAAE,GAAG,uBAAuB,GAAG,SAAS,CAKxF;IAEK,OAAO,CACZ,KAAK,EAAE,uBAAuB,GAAG;QAAE,MAAM,CAAC,EAAE,WAAW,CAAC;QAAC,kBAAkB,CAAC,EAAE,MAAM,CAAA;KAAE,GACpF,OAAO,CAAC,uBAAuB,CAAC,CAuClC;IAEK,OAAO,CAAC,QAAQ,EAAE,iBAAiB,EAAE,OAAO,EAAE,qBAAqB,GAAG,OAAO,CAAC,uBAAuB,CAAC,CAW3G;IAEK,MAAM,CACX,kBAAkB,EAAE,MAAM,GACxB,OAAO,CAAC;QAAE,MAAM,EAAE,WAAW,GAAG,WAAW,GAAG,SAAS,CAAC;QAAC,kBAAkB,EAAE,MAAM,CAAA;KAAE,CAAC,CAUxF;YAMa,iBAAiB;YAwBjB,OAAO;YAwBP,KAAK;YAmBL,IAAI;YAeJ,eAAe;CAe7B;AAED,qFAAqF;AACrF,qBAAa,yBAA0B,SAAQ,KAAK;IACnD,QAAQ,CAAC,kBAAkB,EAAE,MAAM,CAAC;IACpC,YAAY,MAAM,EAAE,MAAM,EAAE,kBAAkB,EAAE,MAAM,EAIrD;CACD","sourcesContent":["/**\n * Remote shared inference admission client (3.0.0 cross-host bridge).\n *\n * Implements `SharedInferenceAdmissionPort` over the structured HTTP protocol to\n * the local admission service. It NEVER re-implements the scheduling algorithm\n * and NEVER falls back to direct provider access: when the scheduler is\n * unreachable or returns an error, `acquire` resolves to a `cancelled` outcome\n * with a `scheduler_unavailable` reason, which the stream seam surfaces as a\n * fail-closed error stream.\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport type { InferenceAdmissionInput, ReleaseInferenceInput, SharedInferenceAdmissionPort } from \"./admission-port.js\";\nimport {\n\tADMISSION_PROTOCOL_VERSION,\n\ttype AdmissionCancelPayload,\n\ttype AdmissionReleasePayload,\n\ttype AdmissionRequestPayload,\n\ttype AdmissionResponse,\n\tparseAdmissionResponse,\n} from \"./admission-protocol.js\";\nimport type {\n\tAcquireInferenceOutcome,\n\tAdmittedInference,\n\tEnqueueInferenceOutcome,\n\tInferenceAdmissionStatus,\n\tReleaseInferenceOutcome,\n\tSharedInferenceResource,\n} from \"./types.js\";\n\nexport interface RemoteAdmissionClientOptions {\n\tbaseUrl: string;\n\ttoken: string;\n\texecutionId: string;\n\tresources: SharedInferenceResource[];\n\t/** Poll cadence while waiting for admission. */\n\tpollMs?: number;\n\t/** Per-request HTTP timeout. */\n\ttimeoutMs?: number;\n\t/** Injectable fetch (tests). */\n\tfetchImpl?: typeof fetch;\n}\n\nexport class RemoteSchedulerAdmissionClient implements SharedInferenceAdmissionPort {\n\treadonly kind = \"remote\" as const;\n\n\tprivate readonly _baseUrl: string;\n\tprivate readonly _token: string;\n\tprivate readonly _executionId: string;\n\tprivate readonly _resources: SharedInferenceResource[];\n\tprivate readonly _pollMs: number;\n\tprivate readonly _timeoutMs: number;\n\tprivate readonly _fetchImpl: typeof fetch;\n\n\tconstructor(options: RemoteAdmissionClientOptions) {\n\t\tthis._baseUrl = options.baseUrl.replace(/\\/+$/, \"\");\n\t\tthis._token = options.token;\n\t\tthis._executionId = options.executionId;\n\t\tthis._resources = options.resources;\n\t\tthis._pollMs = options.pollMs ?? 250;\n\t\tthis._timeoutMs = options.timeoutMs ?? 10_000;\n\t\tthis._fetchImpl = options.fetchImpl ?? fetch;\n\t}\n\n\tresourceFor(model: { provider: string; id: string }): SharedInferenceResource | undefined {\n\t\tfor (const resource of this._resources) {\n\t\t\tif (resource.backend === model.provider && resource.model === model.id) return { ...resource };\n\t\t}\n\t\treturn undefined;\n\t}\n\n\tasync acquire(\n\t\tinput: InferenceAdmissionInput & { signal?: AbortSignal; queueWaitTimeoutMs?: number },\n\t): Promise<AcquireInferenceOutcome> {\n\t\tconst inferenceRequestId = input.inferenceRequestId ?? `inference_${randomUUID()}`;\n\t\tlet first: EnqueueInferenceOutcome;\n\t\ttry {\n\t\t\tfirst = await this._requestAdmission(input, inferenceRequestId);\n\t\t} catch (error) {\n\t\t\treturn {\n\t\t\t\tstatus: \"cancelled\",\n\t\t\t\tinferenceRequestId,\n\t\t\t\treason:\n\t\t\t\t\terror instanceof AdmissionUnavailableError ? error.message : `scheduler_unavailable:${String(error)}`,\n\t\t\t};\n\t\t}\n\t\tif (first.status === \"admitted\") return { status: \"admitted\", admitted: first.admitted };\n\n\t\tconst enqueuedAtMs = Date.now();\n\t\tconst timeoutMs = input.queueWaitTimeoutMs ?? 0;\n\t\twhile (true) {\n\t\t\tif (input.signal?.aborted) {\n\t\t\t\tawait this.cancel(inferenceRequestId);\n\t\t\t\treturn { status: \"cancelled\", inferenceRequestId, reason: \"aborted\" };\n\t\t\t}\n\t\t\tconst waitedMs = Date.now() - enqueuedAtMs;\n\t\t\tif (timeoutMs > 0 && waitedMs >= timeoutMs) {\n\t\t\t\tawait this.cancel(inferenceRequestId);\n\t\t\t\treturn { status: \"queue_timeout\", inferenceRequestId, waitedMs };\n\t\t\t}\n\n\t\t\tconst status = await this._status(inferenceRequestId);\n\t\t\tif (status.status === \"admitted\") return { status: \"admitted\", admitted: status.admitted };\n\t\t\tif (status.status === \"terminal\") {\n\t\t\t\treturn { status: \"cancelled\", inferenceRequestId, reason: `terminal:${status.state}` };\n\t\t\t}\n\t\t\tif (status.status === \"unknown\") {\n\t\t\t\treturn { status: \"cancelled\", inferenceRequestId, reason: \"removed_from_queue\" };\n\t\t\t}\n\n\t\t\tawait new Promise((resolve) => setTimeout(resolve, this._pollMs));\n\t\t}\n\t}\n\n\tasync release(admitted: AdmittedInference, outcome: ReleaseInferenceInput): Promise<ReleaseInferenceOutcome> {\n\t\tconst response = await this._post(\"/v1/release\", {\n\t\t\tprotocolVersion: ADMISSION_PROTOCOL_VERSION,\n\t\t\texecutionId: this._executionId,\n\t\t\tadmitted,\n\t\t\toutcome,\n\t\t} satisfies AdmissionReleasePayload);\n\t\tif (response.kind !== \"release_result\") {\n\t\t\tthrow new Error(`unexpected release response: ${response.kind}`);\n\t\t}\n\t\treturn response.outcome;\n\t}\n\n\tasync cancel(\n\t\tinferenceRequestId: string,\n\t): Promise<{ status: \"cancelled\" | \"not_found\" | \"running\"; inferenceRequestId: string }> {\n\t\tconst response = await this._post(\"/v1/cancel\", {\n\t\t\tprotocolVersion: ADMISSION_PROTOCOL_VERSION,\n\t\t\texecutionId: this._executionId,\n\t\t\tinferenceRequestId,\n\t\t} satisfies AdmissionCancelPayload);\n\t\tif (response.kind !== \"cancel_result\") {\n\t\t\tthrow new Error(`unexpected cancel response: ${response.kind}`);\n\t\t}\n\t\treturn response.outcome;\n\t}\n\n\t// =========================================================================\n\t// Internals\n\t// =========================================================================\n\n\tprivate async _requestAdmission(\n\t\tinput: InferenceAdmissionInput,\n\t\tinferenceRequestId: string,\n\t): Promise<EnqueueInferenceOutcome> {\n\t\tconst response = await this._post(\"/v1/request\", {\n\t\t\tprotocolVersion: ADMISSION_PROTOCOL_VERSION,\n\t\t\texecutionId: this._executionId,\n\t\t\tlogicalAgentId: input.logicalAgentId,\n\t\t\tprovider: input.model.provider,\n\t\t\tmodel: input.model.id,\n\t\t\tinferenceRequestId,\n\t\t\tmissionId: input.missionId,\n\t\t\tassignmentId: input.assignmentId,\n\t\t\tpriority: input.priority,\n\t\t\tdependency: input.dependency,\n\t\t\testimatedInputTokens: input.estimatedInputTokens,\n\t\t\tmaxOutputTokens: input.maxOutputTokens,\n\t\t} satisfies AdmissionRequestPayload);\n\t\tif (response.kind !== \"request_result\") {\n\t\t\tthrow new AdmissionUnavailableError(`unexpected request response: ${response.kind}`, inferenceRequestId);\n\t\t}\n\t\treturn response.outcome;\n\t}\n\n\tprivate async _status(\n\t\tinferenceRequestId: string,\n\t): Promise<InferenceAdmissionStatus | { status: \"cancelled\"; inferenceRequestId: string; reason: string }> {\n\t\ttry {\n\t\t\tconst response = await this._get(\n\t\t\t\t`/v1/status?inferenceRequestId=${encodeURIComponent(inferenceRequestId)}&executionId=${encodeURIComponent(this._executionId)}`,\n\t\t\t);\n\t\t\tif (response.kind !== \"status_result\") {\n\t\t\t\treturn {\n\t\t\t\t\tstatus: \"cancelled\",\n\t\t\t\t\tinferenceRequestId,\n\t\t\t\t\treason: `scheduler_unavailable:unexpected status response ${response.kind}`,\n\t\t\t\t};\n\t\t\t}\n\t\t\treturn response.status;\n\t\t} catch (error) {\n\t\t\treturn {\n\t\t\t\tstatus: \"cancelled\",\n\t\t\t\tinferenceRequestId,\n\t\t\t\treason: `scheduler_unavailable:${error instanceof Error ? error.message : String(error)}`,\n\t\t\t};\n\t\t}\n\t}\n\n\tprivate async _post(path: string, body: unknown): Promise<AdmissionResponse> {\n\t\tconst controller = new AbortController();\n\t\tconst timer = setTimeout(() => controller.abort(), this._timeoutMs);\n\t\ttry {\n\t\t\tconst response = await this._fetchImpl(`${this._baseUrl}${path}`, {\n\t\t\t\tmethod: \"POST\",\n\t\t\t\theaders: {\n\t\t\t\t\t\"content-type\": \"application/json\",\n\t\t\t\t\tauthorization: `Bearer ${this._token}`,\n\t\t\t\t},\n\t\t\t\tbody: JSON.stringify(body),\n\t\t\t\tsignal: controller.signal,\n\t\t\t});\n\t\t\treturn this._decodeResponse(response);\n\t\t} finally {\n\t\t\tclearTimeout(timer);\n\t\t}\n\t}\n\n\tprivate async _get(path: string): Promise<AdmissionResponse> {\n\t\tconst controller = new AbortController();\n\t\tconst timer = setTimeout(() => controller.abort(), this._timeoutMs);\n\t\ttry {\n\t\t\tconst response = await this._fetchImpl(`${this._baseUrl}${path}`, {\n\t\t\t\tmethod: \"GET\",\n\t\t\t\theaders: { authorization: `Bearer ${this._token}` },\n\t\t\t\tsignal: controller.signal,\n\t\t\t});\n\t\t\treturn this._decodeResponse(response);\n\t\t} finally {\n\t\t\tclearTimeout(timer);\n\t\t}\n\t}\n\n\tprivate async _decodeResponse(response: Response): Promise<AdmissionResponse> {\n\t\tconst text = await response.text();\n\t\tlet parsed: unknown;\n\t\ttry {\n\t\t\tparsed = JSON.parse(text);\n\t\t} catch {\n\t\t\tthrow new Error(`admission service returned non-JSON (HTTP ${response.status})`);\n\t\t}\n\t\tconst decoded = parseAdmissionResponse(parsed);\n\t\tif (!response.ok) {\n\t\t\tif (decoded.kind === \"error\") throw new Error(`${decoded.code}: ${decoded.message}`);\n\t\t\tthrow new Error(`admission service HTTP ${response.status}`);\n\t\t}\n\t\treturn decoded;\n\t}\n}\n\n/** Fail-closed marker used to convert transport failure into a cancelled outcome. */\nexport class AdmissionUnavailableError extends Error {\n\treadonly inferenceRequestId: string;\n\tconstructor(reason: string, inferenceRequestId: string) {\n\t\tsuper(reason);\n\t\tthis.name = \"AdmissionUnavailableError\";\n\t\tthis.inferenceRequestId = inferenceRequestId;\n\t}\n}\n"]}