{"version":3,"file":"scheduler.d.ts","sourceRoot":"","sources":["../../../src/core/shared-inference/scheduler.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;GAgBG;AAIH,OAAO,KAAK,EAAE,mBAAmB,EAA2B,MAAM,sBAAsB,CAAC;AACzF,OAAO,KAAK,EACX,uBAAuB,EACvB,iBAAiB,EACjB,uBAAuB,EACvB,wBAAwB,EACxB,iBAAiB,EACjB,uBAAuB,EACvB,0BAA0B,EAC1B,qBAAqB,EAGrB,uBAAuB,EACvB,qBAAqB,EACrB,uBAAuB,EAEvB,8BAA8B,EAC9B,MAAM,YAAY,CAAC;AAMpB,MAAM,WAAW,+BAA+B;IAC/C,KAAK,EAAE,mBAAmB,CAAC;IAC3B,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,KAAK,CAAC,EAAE,CAAC,EAAE,EAAE,MAAM,KAAK,OAAO,CAAC,IAAI,CAAC,CAAC;IACtC,gEAAgE;IAChE,eAAe,CAAC,EAAE,MAAM,CAAC;IACzB,cAAc,CAAC,EAAE,MAAM,MAAM,CAAC;IAC9B,gBAAgB,CAAC,EAAE,MAAM,MAAM,CAAC;IAChC,0EAA0E;IAC1E,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,8DAA8D;IAC9D,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,oDAAoD;IACpD,kBAAkB,CAAC,EAAE,MAAM,CAAC;IAC5B,2CAA2C;IAC3C,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,iBAAiB,CAAC,EAAE,MAAM,CAAC;IAC3B,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,aAAa,CAAC,EAAE,MAAM,CAAC;IACvB,eAAe,CAAC,EAAE,MAAM,CAAC;IACzB,WAAW,CAAC,EAAE,MAAM,CAAC;CACrB;AAwBD,wBAAgB,6BAA6B,CAAC,KAAK,EAAE,qBAAqB,EAAE,GAAG,EAAE,MAAM,GAAG,OAAO,CAEhG;AAMD,qBAAa,wBAAwB;IACpC,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAsB;IAC7C,OAAO,CAAC,QAAQ,CAAC,UAAU,CAA8C;IACzE,OAAO,CAAC,QAAQ,CAAC,IAAI,CAAe;IACpC,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAgC;IACvD,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAS;IAC1C,OAAO,CAAC,QAAQ,CAAC,eAAe,CAAe;IAC/C,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAAe;IACjD,OAAO,CAAC,QAAQ,CAAC,QAAQ,CAAS;IAClC,OAAO,CAAC,QAAQ,CAAC,WAAW,CAAS;IACrC,OAAO,CAAC,QAAQ,CAAC,mBAAmB,CAAS;IAC7C,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAAS;IAC3C,OAAO,CAAC,QAAQ,CAAC,kBAAkB,CAAS;IAC5C,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAAS;IAC3C,OAAO,CAAC,QAAQ,CAAC,cAAc,CAAS;IACxC,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAS;IAC1C,OAAO,CAAC,QAAQ,CAAC,YAAY,CAAS;IACtC,OAAO,CAAC,qBAAqB,CAAC,CAAS;IACvC,OAAO,CAAC,qBAAqB,CAAK;IAElC,YAAY,OAAO,EAAE,+BAA+B,EAgBnD;IAED,IAAI,OAAO,IAAI,MAAM,CAEpB;IAED,IAAI,KAAK,IAAI,mBAAmB,CAE/B;IAMD,4EAA4E;IACtE,gBAAgB,CAAC,QAAQ,EAAE,uBAAuB,GAAG,OAAO,CAAC,IAAI,CAAC,CA0BvE;IAED,mFAAmF;IACnF,WAAW,CAAC,KAAK,EAAE;QAAE,QAAQ,EAAE,MAAM,CAAC;QAAC,EAAE,EAAE,MAAM,CAAA;KAAE,GAAG,uBAAuB,GAAG,SAAS,CAKxF;IAED,aAAa,IAAI,uBAAuB,EAAE,CAIzC;IAMD;;;;;;OAMG;IACG,OAAO,CAAC,KAAK,EAAE;QACpB,cAAc,EAAE,MAAM,CAAC;QACvB,QAAQ,EAAE,uBAAuB,CAAC;QAClC,KAAK,EAAE;YAAE,QAAQ,EAAE,MAAM,CAAC;YAAC,EAAE,EAAE,MAAM,CAAA;SAAE,CAAC;QACxC,kBAAkB,CAAC,EAAE,MAAM,CAAC;QAC5B,SAAS,CAAC,EAAE,MAAM,CAAC;QACnB,YAAY,CAAC,EAAE,MAAM,CAAC;QACtB,WAAW,CAAC,EAAE,MAAM,CAAC;QACrB,QAAQ,CAAC,EAAE,iBAAiB,CAAC;QAC7B,UAAU,CAAC,EAAE,0BAA0B,CAAC;QACxC,oBAAoB,CAAC,EAAE,MAAM,CAAC;QAC9B,eAAe,CAAC,EAAE,MAAM,CAAC;KACzB,GAAG,OAAO,CAAC,uBAAuB,CAAC,CASnC;IAED,+EAA+E;IACzE,eAAe,CACpB,kBAAkB,EAAE,MAAM,EAC1B,OAAO,GAAE;QAAE,OAAO,CAAC,EAAE,MAAM,CAAA;KAAO,GAChC,OAAO,CAAC,wBAAwB,CAAC,CAkBnC;IAED,uFAAuF;IACjF,KAAK,CAAC,QAAQ,EAAE,iBAAiB,EAAE,OAAO,GAAE;QAAE,GAAG,CAAC,EAAE,MAAM,CAAA;KAAO,GAAG,OAAO,CAAC,qBAAqB,CAAC,CAmCvG;IAED;;;;;OAKG;IACG,OAAO,CAAC,KAAK,EAAE;QACpB,cAAc,EAAE,MAAM,CAAC;QACvB,QAAQ,EAAE,uBAAuB,CAAC;QAClC,KAAK,EAAE;YAAE,QAAQ,EAAE,MAAM,CAAC;YAAC,EAAE,EAAE,MAAM,CAAA;SAAE,CAAC;QACxC,kBAAkB,CAAC,EAAE,MAAM,CAAC;QAC5B,SAAS,CAAC,EAAE,MAAM,CAAC;QACnB,YAAY,CAAC,EAAE,MAAM,CAAC;QACtB,WAAW,CAAC,EAAE,MAAM,CAAC;QACrB,QAAQ,CAAC,EAAE,iBAAiB,CAAC;QAC7B,UAAU,CAAC,EAAE,0BAA0B,CAAC;QACxC,oBAAoB,CAAC,EAAE,MAAM,CAAC;QAC9B,eAAe,CAAC,EAAE,MAAM,CAAC;QACzB,MAAM,CAAC,EAAE,WAAW,CAAC;QACrB,kBAAkB,CAAC,EAAE,MAAM,CAAC;KAC5B,GAAG,OAAO,CAAC,uBAAuB,CAAC,CAgCnC;IAED,OAAO,CAAC,aAAa;YAoCP,eAAe;IA4D7B,kFAAkF;IAC5E,OAAO,CACZ,QAAQ,EAAE,iBAAiB,EAC3B,OAAO,EAAE;QACR,KAAK,EAAE,WAAW,GAAG,QAAQ,GAAG,WAAW,GAAG,aAAa,CAAC;QAC5D,KAAK,CAAC,EAAE;YAAE,KAAK,CAAC,EAAE,MAAM,CAAC;YAAC,MAAM,CAAC,EAAE,MAAM,CAAA;SAAE,CAAC;QAC5C,YAAY,CAAC,EAAE,MAAM,CAAC;KACtB,GACC,OAAO,CAAC,uBAAuB,CAAC,CAsElC;IAED,0GAA0G;IACpG,MAAM,CACX,kBAAkB,EAAE,MAAM,GACxB,OAAO,CAAC;QAAE,MAAM,EAAE,WAAW,GAAG,WAAW,GAAG,SAAS,CAAC;QAAC,kBAAkB,EAAE,MAAM,CAAA;KAAE,CAAC,CA6CxF;IAMD,uGAAuG;IACjG,OAAO,CAAC,OAAO,GAAE;QAAE,GAAG,CAAC,EAAE,MAAM,CAAA;KAAO,GAAG,OAAO,CAAC,uBAAuB,CAAC,CAgF9E;IAMK,MAAM,IAAI,OAAO,CAAC,8BAA8B,CAAC,CAsFtD;IAED,4FAA4F;IAC5F,OAAO,CAAC,qBAAqB;IAkB7B,OAAO,CAAC,WAAW;IAUnB,OAAO,CAAC,aAAa;IAKrB,OAAO,CAAC,UAAU;IA8BlB,wEAAwE;IACxE,OAAO,CAAC,cAAc;IAuCtB,wCAAwC;IACxC,OAAO,CAAC,eAAe;IAKvB,OAAO,CAAC,YAAY;IAIpB,OAAO,CAAC,QAAQ;IAQhB,OAAO,CAAC,UAAU;IAmBlB,OAAO,CAAC,cAAc;CAgBtB","sourcesContent":["/**\n * Shared Inference Scheduler (3.0.0 foundation).\n *\n * Central admission authority for physical inference slots. It controls WHEN a\n * model request may run; it is NOT the model provider and never reimplements\n * provider protocol/stream handling. Scheduling/admission is kept separate from\n * provider protocol handling.\n *\n * Cross-process coordination: the scheduler mutates a durable per-resource\n * ledger (see InferenceQueueStore). Two OS processes sharing the same ledger\n * cannot each believe they own all physical slots. Slot acquisition is a\n * temporary fenced lease; it is never durable agent ownership.\n *\n * Determinism: with the same queued requests, priorities, dependency metadata,\n * controlled clock, and capacity, admission order is deterministic\n * (effective-priority desc, then enqueue time, then request id).\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport { newExecutorOwnerId } from \"../mission-domain/execution-lease.js\";\nimport type { InferenceQueueStore, InferenceResourceLedger } from \"./inference-queue.js\";\nimport type {\n\tAcquireInferenceOutcome,\n\tAdmittedInference,\n\tEnqueueInferenceOutcome,\n\tInferenceAdmissionStatus,\n\tInferencePriority,\n\tInferenceRecoveryReport,\n\tInferenceRequestDependency,\n\tInferenceRequestLease,\n\tInferenceRequestRecord,\n\tInferenceRequestStatus,\n\tReleaseInferenceOutcome,\n\tRenewInferenceOutcome,\n\tSharedInferenceResource,\n\tSharedInferenceResourceStatus,\n\tSharedInferenceSchedulerStatus,\n} from \"./types.js\";\n\n// =============================================================================\n// Options\n// =============================================================================\n\nexport interface SharedInferenceSchedulerOptions {\n\tstore: InferenceQueueStore;\n\tnow?: () => number;\n\tsleep?: (ms: number) => Promise<void>;\n\t/** Lease lifetime for an admitted slot (default 10 minutes). */\n\tleaseDurationMs?: number;\n\tleaseIdFactory?: () => string;\n\trequestIdFactory?: () => string;\n\t/** Stable control-plane owner identity (default host+UUID, never PID). */\n\townerId?: string;\n\t/** Poll cadence while waiting for cross-process admission. */\n\twaitPollMs?: number;\n\t/** Optional queue-wait deadline (0 = unbounded). */\n\tqueueWaitTimeoutMs?: number;\n\t/** Deterministic priority policy knobs. */\n\tinteractiveBoost?: number;\n\tverificationBoost?: number;\n\tdependencyWeight?: number;\n\tdependencyCap?: number;\n\tagingIntervalMs?: number;\n\tagingWeight?: number;\n}\n\nconst DEFAULT_LEASE_DURATION_MS = 10 * 60_000;\nconst DEFAULT_WAIT_POLL_MS = 100;\nconst DEFAULT_INTERACTIVE_BOOST = 100;\nconst DEFAULT_VERIFICATION_BOOST = 200;\nconst DEFAULT_DEPENDENCY_WEIGHT = 50;\nconst DEFAULT_DEPENDENCY_CAP = 500;\nconst DEFAULT_AGING_INTERVAL_MS = 1000;\nconst DEFAULT_AGING_WEIGHT = 1;\n\n// =============================================================================\n// Priority (deterministic, explainable)\n// =============================================================================\n\ninterface PriorityBreakdown {\n\tbase: number;\n\tinteractive: number;\n\tverification: number;\n\tdependency: number;\n\taging: number;\n\teffective: number;\n}\n\nexport function isInferenceRequestLeaseActive(lease: InferenceRequestLease, now: number): boolean {\n\treturn lease.expiresAtMs > now;\n}\n\n// =============================================================================\n// Scheduler\n// =============================================================================\n\nexport class SharedInferenceScheduler {\n\tprivate readonly _store: InferenceQueueStore;\n\tprivate readonly _resources = new Map<string, SharedInferenceResource>();\n\tprivate readonly _now: () => number;\n\tprivate readonly _sleep: (ms: number) => Promise<void>;\n\tprivate readonly _leaseDurationMs: number;\n\tprivate readonly _leaseIdFactory: () => string;\n\tprivate readonly _requestIdFactory: () => string;\n\tprivate readonly _ownerId: string;\n\tprivate readonly _waitPollMs: number;\n\tprivate readonly _queueWaitTimeoutMs: number;\n\tprivate readonly _interactiveBoost: number;\n\tprivate readonly _verificationBoost: number;\n\tprivate readonly _dependencyWeight: number;\n\tprivate readonly _dependencyCap: number;\n\tprivate readonly _agingIntervalMs: number;\n\tprivate readonly _agingWeight: number;\n\tprivate _avoidableIdleSinceMs?: number;\n\tprivate _avoidableIdleTotalMs = 0;\n\n\tconstructor(options: SharedInferenceSchedulerOptions) {\n\t\tthis._store = options.store;\n\t\tthis._now = options.now ?? (() => Date.now());\n\t\tthis._sleep = options.sleep ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms)));\n\t\tthis._leaseDurationMs = options.leaseDurationMs ?? DEFAULT_LEASE_DURATION_MS;\n\t\tthis._leaseIdFactory = options.leaseIdFactory ?? (() => `inflease_${randomUUID()}`);\n\t\tthis._requestIdFactory = options.requestIdFactory ?? (() => `inference_${randomUUID()}`);\n\t\tthis._ownerId = options.ownerId ?? newExecutorOwnerId();\n\t\tthis._waitPollMs = options.waitPollMs ?? DEFAULT_WAIT_POLL_MS;\n\t\tthis._queueWaitTimeoutMs = options.queueWaitTimeoutMs ?? 0;\n\t\tthis._interactiveBoost = options.interactiveBoost ?? DEFAULT_INTERACTIVE_BOOST;\n\t\tthis._verificationBoost = options.verificationBoost ?? DEFAULT_VERIFICATION_BOOST;\n\t\tthis._dependencyWeight = options.dependencyWeight ?? DEFAULT_DEPENDENCY_WEIGHT;\n\t\tthis._dependencyCap = options.dependencyCap ?? DEFAULT_DEPENDENCY_CAP;\n\t\tthis._agingIntervalMs = options.agingIntervalMs ?? DEFAULT_AGING_INTERVAL_MS;\n\t\tthis._agingWeight = options.agingWeight ?? DEFAULT_AGING_WEIGHT;\n\t}\n\n\tget ownerId(): string {\n\t\treturn this._ownerId;\n\t}\n\n\tget store(): InferenceQueueStore {\n\t\treturn this._store;\n\t}\n\n\t// =========================================================================\n\t// Resource registry\n\t// =========================================================================\n\n\t/** Register (or refresh) a shared resource and ensure its ledger exists. */\n\tasync registerResource(resource: SharedInferenceResource): Promise<void> {\n\t\tif (resource.capacity < 1) throw new Error(`Resource ${resource.resourceId} capacity must be >= 1`);\n\t\tthis._resources.set(resource.resourceId, { ...resource });\n\t\tconst created = await this._store.create(resource.resourceId, {\n\t\t\tschemaVersion: 1,\n\t\t\tresourceId: resource.resourceId,\n\t\t\tcapacity: resource.capacity,\n\t\t\trunning: [],\n\t\t\tqueue: [],\n\t\t\thistory: [],\n\t\t\tnextSeq: 0,\n\t\t\tfencingToken: 0,\n\t\t\tcompletedCount: 0,\n\t\t\tcancelledCount: 0,\n\t\t\tfailedCount: 0,\n\t\t\tinterruptedCount: 0,\n\t\t\ttotalQueueWaitMs: 0,\n\t\t\ttotalInferenceMs: 0,\n\t\t\ttotalInputTokens: 0,\n\t\t\ttotalOutputTokens: 0,\n\t\t\tmaxQueueDepth: 0,\n\t\t\tupdatedAtMs: this._now(),\n\t\t\trevision: 1,\n\t\t});\n\t\tif (created.status === \"conflict\")\n\t\t\tthrow new Error(`Cannot register resource ${resource.resourceId}: ${created.error}`);\n\t}\n\n\t/** Resolve the shared resource a model targets, or undefined (non-shared path). */\n\tresourceFor(model: { provider: string; id: string }): SharedInferenceResource | undefined {\n\t\tfor (const resource of this._resources.values()) {\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\tlistResources(): SharedInferenceResource[] {\n\t\treturn [...this._resources.values()]\n\t\t\t.map((r) => ({ ...r }))\n\t\t\t.sort((a, b) => a.resourceId.localeCompare(b.resourceId));\n\t}\n\n\t// =========================================================================\n\t// Admission\n\t// =========================================================================\n\n\t/**\n\t * Non-blocking enqueue (or immediate admit) of an inference request. This is\n\t * the single authority for turning a request into either an admitted lease or\n\t * a durable queued entry; `acquire` composes it with a bounded wait loop.\n\t * The admission service uses `enqueue` + `admissionStatus` to expose a\n\t * request/wait protocol to remote clients without duplicating the algorithm.\n\t */\n\tasync enqueue(input: {\n\t\tlogicalAgentId: string;\n\t\tresource: SharedInferenceResource;\n\t\tmodel: { provider: string; id: string };\n\t\tinferenceRequestId?: string;\n\t\tmissionId?: string;\n\t\tassignmentId?: string;\n\t\texecutionId?: string;\n\t\tpriority?: InferencePriority;\n\t\tdependency?: InferenceRequestDependency;\n\t\testimatedInputTokens?: number;\n\t\tmaxOutputTokens?: number;\n\t}): Promise<EnqueueInferenceOutcome> {\n\t\tconst request = this._buildRequest(input);\n\t\tconst result = await this._enqueueOrAdmit(request.resourceId, request);\n\t\tif (result.status === \"admitted\") return { status: \"admitted\", admitted: result.admitted };\n\t\tif (result.status === \"queued\")\n\t\t\treturn { status: \"queued\", inferenceRequestId: result.inferenceRequestId, position: result.position };\n\t\t// The enqueue path can only produce admitted/queued; anything else is a\n\t\t// structural defect surfaced loudly rather than coerced.\n\t\tthrow new Error(`Unexpected enqueue outcome: ${(result as { status: string }).status}`);\n\t}\n\n\t/** Pollable admission state for a request owned by this scheduler instance. */\n\tasync admissionStatus(\n\t\tinferenceRequestId: string,\n\t\toptions: { ownerId?: string } = {},\n\t): Promise<InferenceAdmissionStatus> {\n\t\tconst ownerId = options.ownerId ?? this._ownerId;\n\t\tfor (const resourceId of await this._store.listResources()) {\n\t\t\tconst loaded = await this._store.load(resourceId);\n\t\t\tif (loaded.status !== \"ok\") continue;\n\t\t\tconst ledger = loaded.ledger;\n\t\t\tconst running = ledger.running.find(\n\t\t\t\t(r) => r.inferenceRequestId === inferenceRequestId && r.lease?.ownerId === ownerId,\n\t\t\t);\n\t\t\tif (running?.lease) return { status: \"admitted\", admitted: this._toAdmitted(running) };\n\t\t\tconst queuedIndex = ledger.queue.findIndex((r) => r.inferenceRequestId === inferenceRequestId);\n\t\t\tif (queuedIndex >= 0) {\n\t\t\t\treturn { status: \"queued\", inferenceRequestId, position: queuedIndex + 1 };\n\t\t\t}\n\t\t\tconst summary = ledger.history.find((r) => r.inferenceRequestId === inferenceRequestId);\n\t\t\tif (summary) return { status: \"terminal\", inferenceRequestId, state: summary.state };\n\t\t}\n\t\treturn { status: \"unknown\", inferenceRequestId };\n\t}\n\n\t/** Renew an active lease (keeps long-running generations from expiring mid-stream). */\n\tasync renew(admitted: AdmittedInference, options: { now?: number } = {}): Promise<RenewInferenceOutcome> {\n\t\tconst now = options.now ?? this._now();\n\t\tconst result = await this._store.mutate<RenewInferenceOutcome>(admitted.resourceId, (ledger) => {\n\t\t\tconst index = ledger.running.findIndex((r) => r.inferenceRequestId === admitted.inferenceRequestId);\n\t\t\tif (index < 0)\n\t\t\t\treturn { kind: \"noop\", value: { status: \"not_found\", inferenceRequestId: admitted.inferenceRequestId } };\n\t\t\tconst running = ledger.running[index];\n\t\t\tif (\n\t\t\t\t!running.lease ||\n\t\t\t\trunning.lease.leaseId !== admitted.lease.leaseId ||\n\t\t\t\trunning.lease.ownerId !== this._ownerId\n\t\t\t) {\n\t\t\t\treturn { kind: \"noop\", value: { status: \"not_found\", inferenceRequestId: admitted.inferenceRequestId } };\n\t\t\t}\n\t\t\tconst expiresAtMs = now + this._leaseDurationMs;\n\t\t\tconst next: InferenceResourceLedger = {\n\t\t\t\t...ledger,\n\t\t\t\trunning: ledger.running.map((r, i) => (i === index ? { ...r, lease: { ...r.lease!, expiresAtMs } } : r)),\n\t\t\t\tupdatedAtMs: now,\n\t\t\t\trevision: ledger.revision + 1,\n\t\t\t};\n\t\t\treturn {\n\t\t\t\tkind: \"write\",\n\t\t\t\tnext,\n\t\t\t\tvalue: {\n\t\t\t\t\tstatus: \"renewed\",\n\t\t\t\t\tadmitted: this._toAdmitted({ ...running, lease: { ...running.lease!, expiresAtMs } }),\n\t\t\t\t\texpiresAtMs,\n\t\t\t\t},\n\t\t\t};\n\t\t});\n\t\tif (result.status === \"missing\") return { status: \"not_found\", inferenceRequestId: admitted.inferenceRequestId };\n\t\tif (result.status === \"corrupt\")\n\t\t\tthrow new Error(`Inference queue for ${admitted.resourceId} is corrupt: ${result.diagnostic}`);\n\t\treturn result.value;\n\t}\n\n\t/**\n\t * Request an inference slot. If capacity is available the request is admitted\n\t * immediately; otherwise it is durably queued and the caller waits (bounded\n\t * poll, abortable, queue-timeout aware) until admission, cancellation, or\n\t * timeout.\n\t */\n\tasync acquire(input: {\n\t\tlogicalAgentId: string;\n\t\tresource: SharedInferenceResource;\n\t\tmodel: { provider: string; id: string };\n\t\tinferenceRequestId?: string;\n\t\tmissionId?: string;\n\t\tassignmentId?: string;\n\t\texecutionId?: string;\n\t\tpriority?: InferencePriority;\n\t\tdependency?: InferenceRequestDependency;\n\t\testimatedInputTokens?: number;\n\t\tmaxOutputTokens?: number;\n\t\tsignal?: AbortSignal;\n\t\tqueueWaitTimeoutMs?: number;\n\t}): Promise<AcquireInferenceOutcome> {\n\t\tconst inferenceRequestId = input.inferenceRequestId ?? this._requestIdFactory();\n\n\t\tconst first = await this.enqueue({\n\t\t\t...input,\n\t\t\tinferenceRequestId,\n\t\t});\n\t\tif (first.status === \"admitted\") return { status: \"admitted\", admitted: first.admitted };\n\n\t\tconst enqueuedAtMs = this._now();\n\t\tconst timeoutMs = input.queueWaitTimeoutMs ?? this._queueWaitTimeoutMs;\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 = this._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 state = await this.admissionStatus(inferenceRequestId);\n\t\t\tif (state.status === \"admitted\") return { status: \"admitted\", admitted: state.admitted };\n\t\t\tif (state.status === \"terminal\")\n\t\t\t\treturn { status: \"cancelled\", inferenceRequestId, reason: `terminal:${state.state}` };\n\t\t\tif (state.status === \"unknown\") {\n\t\t\t\treturn { status: \"cancelled\", inferenceRequestId, reason: \"removed_from_queue\" };\n\t\t\t}\n\n\t\t\tawait this._sleep(this._waitPollMs);\n\t\t}\n\t}\n\n\tprivate _buildRequest(input: {\n\t\tlogicalAgentId: string;\n\t\tresource: SharedInferenceResource;\n\t\tmodel: { provider: string; id: string };\n\t\tinferenceRequestId?: string;\n\t\tmissionId?: string;\n\t\tassignmentId?: string;\n\t\texecutionId?: string;\n\t\tpriority?: InferencePriority;\n\t\tdependency?: InferenceRequestDependency;\n\t\testimatedInputTokens?: number;\n\t\tmaxOutputTokens?: number;\n\t}): InferenceRequestRecord {\n\t\tconst now = this._now();\n\t\tconst requestId = input.inferenceRequestId ?? this._requestIdFactory();\n\t\treturn {\n\t\t\tschemaVersion: 1,\n\t\t\tinferenceRequestId: requestId,\n\t\t\tlogicalAgentId: input.logicalAgentId,\n\t\t\townerId: this._ownerId,\n\t\t\tmissionId: input.missionId,\n\t\t\tassignmentId: input.assignmentId,\n\t\t\texecutionId: input.executionId,\n\t\t\tresourceId: input.resource.resourceId,\n\t\t\tprovider: input.model.provider,\n\t\t\tmodel: input.model.id,\n\t\t\trequestedAtMs: now,\n\t\t\tenqueuedAtMs: now,\n\t\t\tpriority: input.priority ?? { base: 0 },\n\t\t\tdependency: input.dependency,\n\t\t\testimatedInputTokens: input.estimatedInputTokens,\n\t\t\tmaxOutputTokens: input.maxOutputTokens,\n\t\t\tstate: \"QUEUED\",\n\t\t};\n\t}\n\n\tprivate async _enqueueOrAdmit(\n\t\tresourceId: string,\n\t\trequest: InferenceRequestRecord,\n\t): Promise<AcquireInferenceOutcome> {\n\t\tconst result = await this._store.mutate<AcquireInferenceOutcome>(resourceId, (ledger) => {\n\t\t\tconst now = this._now();\n\t\t\tconst running = ledger.running.find((r) => r.inferenceRequestId === request.inferenceRequestId);\n\t\t\tif (running?.lease) {\n\t\t\t\treturn { kind: \"noop\", value: { status: \"admitted\" as const, admitted: this._toAdmitted(running) } };\n\t\t\t}\n\t\t\tconst queuedIndex = ledger.queue.findIndex((r) => r.inferenceRequestId === request.inferenceRequestId);\n\t\t\tif (queuedIndex >= 0) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"noop\",\n\t\t\t\t\tvalue: {\n\t\t\t\t\t\tstatus: \"queued\" as const,\n\t\t\t\t\t\tinferenceRequestId: request.inferenceRequestId,\n\t\t\t\t\t\tposition: queuedIndex + 1,\n\t\t\t\t\t},\n\t\t\t\t};\n\t\t\t}\n\n\t\t\tif (ledger.running.length < ledger.capacity) {\n\t\t\t\tconst next = this._admitInto(ledger, request, now, { incrementFence: true });\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"write\",\n\t\t\t\t\tnext,\n\t\t\t\t\tvalue: { status: \"admitted\" as const, admitted: this._admittedFrom(next, request.inferenceRequestId) },\n\t\t\t\t};\n\t\t\t}\n\n\t\t\tconst queued: InferenceRequestRecord = { ...request, state: \"QUEUED\" };\n\t\t\tconst queue = [...ledger.queue, queued];\n\t\t\tconst next: InferenceResourceLedger = {\n\t\t\t\t...ledger,\n\t\t\t\tqueue,\n\t\t\t\tmaxQueueDepth: Math.max(ledger.maxQueueDepth, queue.length),\n\t\t\t\tupdatedAtMs: now,\n\t\t\t\trevision: ledger.revision + 1,\n\t\t\t};\n\t\t\treturn {\n\t\t\t\tkind: \"write\",\n\t\t\t\tnext,\n\t\t\t\tvalue: {\n\t\t\t\t\tstatus: \"queued\" as const,\n\t\t\t\t\tinferenceRequestId: request.inferenceRequestId,\n\t\t\t\t\tposition: queue.length,\n\t\t\t\t},\n\t\t\t};\n\t\t});\n\n\t\tif (result.status === \"missing\") {\n\t\t\tthrow new Error(`Inference resource ${resourceId} is not registered`);\n\t\t}\n\t\tif (result.status === \"corrupt\") {\n\t\t\tthrow new Error(`Inference queue for ${resourceId} is corrupt: ${result.diagnostic}`);\n\t\t}\n\t\treturn result.value;\n\t}\n\n\t/** Release an admitted slot and promote queued work. Fenced by lease identity. */\n\tasync release(\n\t\tadmitted: AdmittedInference,\n\t\toutcome: {\n\t\t\tstate: \"COMPLETED\" | \"FAILED\" | \"CANCELLED\" | \"INTERRUPTED\";\n\t\t\tusage?: { input?: number; output?: number };\n\t\t\terrorMessage?: string;\n\t\t},\n\t): Promise<ReleaseInferenceOutcome> {\n\t\tconst result = await this._store.mutate<ReleaseInferenceOutcome>(admitted.resourceId, (ledger) => {\n\t\t\tconst now = this._now();\n\t\t\tconst index = ledger.running.findIndex((r) => r.inferenceRequestId === admitted.inferenceRequestId);\n\t\t\tif (index < 0) return { kind: \"noop\", value: { status: \"not_found\" as const } };\n\t\t\tconst running = ledger.running[index];\n\t\t\tif (\n\t\t\t\t!running.lease ||\n\t\t\t\trunning.lease.leaseId !== admitted.lease.leaseId ||\n\t\t\t\trunning.lease.ownerId !== this._ownerId\n\t\t\t) {\n\t\t\t\treturn { kind: \"noop\", value: { status: \"not_found\" as const } };\n\t\t\t}\n\n\t\t\tconst finishedAtMs = now;\n\t\t\tconst queueWaitMs = running.admittedAtMs\n\t\t\t\t? Math.max(0, running.admittedAtMs - running.enqueuedAtMs)\n\t\t\t\t: Math.max(0, now - running.enqueuedAtMs);\n\t\t\tconst inferenceWallMs = running.admittedAtMs ? Math.max(0, finishedAtMs - running.admittedAtMs) : 0;\n\t\t\tconst summary = {\n\t\t\t\tinferenceRequestId: running.inferenceRequestId,\n\t\t\t\tlogicalAgentId: running.logicalAgentId,\n\t\t\t\tresourceId: running.resourceId,\n\t\t\t\tprovider: running.provider,\n\t\t\t\tmodel: running.model,\n\t\t\t\trequestedAtMs: running.requestedAtMs,\n\t\t\t\tenqueuedAtMs: running.enqueuedAtMs,\n\t\t\t\tstate: outcome.state,\n\t\t\t\tadmittedAtMs: running.admittedAtMs,\n\t\t\t\tfinishedAtMs,\n\t\t\t\tqueueWaitMs,\n\t\t\t\tinferenceWallMs,\n\t\t\t\tinputTokens: outcome.usage?.input,\n\t\t\t\toutputTokens: outcome.usage?.output,\n\t\t\t\terrorMessage: outcome.errorMessage,\n\t\t\t};\n\n\t\t\tconst counters = {\n\t\t\t\tcompletedCount: outcome.state === \"COMPLETED\" ? ledger.completedCount + 1 : ledger.completedCount,\n\t\t\t\tcancelledCount: outcome.state === \"CANCELLED\" ? ledger.cancelledCount + 1 : ledger.cancelledCount,\n\t\t\t\tfailedCount: outcome.state === \"FAILED\" ? ledger.failedCount + 1 : ledger.failedCount,\n\t\t\t\tinterruptedCount: outcome.state === \"INTERRUPTED\" ? ledger.interruptedCount + 1 : ledger.interruptedCount,\n\t\t\t};\n\n\t\t\tconst next: InferenceResourceLedger = {\n\t\t\t\t...ledger,\n\t\t\t\trunning: ledger.running.filter((_, i) => i !== index),\n\t\t\t\thistory: [summary, ...ledger.history].slice(0, 1000),\n\t\t\t\t...counters,\n\t\t\t\ttotalQueueWaitMs: ledger.totalQueueWaitMs + queueWaitMs,\n\t\t\t\ttotalInferenceMs: ledger.totalInferenceMs + inferenceWallMs,\n\t\t\t\ttotalInputTokens: ledger.totalInputTokens + (outcome.usage?.input ?? 0),\n\t\t\t\ttotalOutputTokens: ledger.totalOutputTokens + (outcome.usage?.output ?? 0),\n\t\t\t\tupdatedAtMs: now,\n\t\t\t\trevision: ledger.revision + 1,\n\t\t\t};\n\n\t\t\t// Promote queued work up to free capacity.\n\t\t\tconst promoted = this._promoteQueued(next, now);\n\t\t\treturn {\n\t\t\t\tkind: \"write\",\n\t\t\t\tnext: promoted,\n\t\t\t\tvalue: { status: \"released\" as const, inferenceRequestId: admitted.inferenceRequestId },\n\t\t\t};\n\t\t});\n\n\t\tif (result.status === \"missing\") return { status: \"not_found\" };\n\t\tif (result.status === \"corrupt\")\n\t\t\tthrow new Error(`Inference queue for ${admitted.resourceId} is corrupt: ${result.diagnostic}`);\n\t\treturn result.value;\n\t}\n\n\t/** Cancel a queued (or locally-owned running) request. Queued requests never consume a slot afterward. */\n\tasync cancel(\n\t\tinferenceRequestId: string,\n\t): Promise<{ status: \"cancelled\" | \"not_found\" | \"running\"; inferenceRequestId: string }> {\n\t\tconst now = this._now();\n\t\tconst ids = await this._store.listResources();\n\t\tfor (const resourceId of ids) {\n\t\t\tconst result = await this._store.mutate<{ status: \"cancelled\" | \"not_found\" | \"running\" }>(\n\t\t\t\tresourceId,\n\t\t\t\t(ledger) => {\n\t\t\t\t\tconst qIndex = ledger.queue.findIndex((r) => r.inferenceRequestId === inferenceRequestId);\n\t\t\t\t\tif (qIndex >= 0) {\n\t\t\t\t\t\tconst request = ledger.queue[qIndex];\n\t\t\t\t\t\tconst summary = {\n\t\t\t\t\t\t\tinferenceRequestId: request.inferenceRequestId,\n\t\t\t\t\t\t\tlogicalAgentId: request.logicalAgentId,\n\t\t\t\t\t\t\tresourceId,\n\t\t\t\t\t\t\tprovider: request.provider,\n\t\t\t\t\t\t\tmodel: request.model,\n\t\t\t\t\t\t\trequestedAtMs: request.requestedAtMs,\n\t\t\t\t\t\t\tenqueuedAtMs: request.enqueuedAtMs,\n\t\t\t\t\t\t\tstate: \"CANCELLED\" as const,\n\t\t\t\t\t\t\tfinishedAtMs: now,\n\t\t\t\t\t\t\tqueueWaitMs: now - request.enqueuedAtMs,\n\t\t\t\t\t\t};\n\t\t\t\t\t\tconst next: InferenceResourceLedger = {\n\t\t\t\t\t\t\t...ledger,\n\t\t\t\t\t\t\tqueue: ledger.queue.filter((_, i) => i !== qIndex),\n\t\t\t\t\t\t\thistory: [summary, ...ledger.history].slice(0, 1000),\n\t\t\t\t\t\t\tcancelledCount: ledger.cancelledCount + 1,\n\t\t\t\t\t\t\ttotalQueueWaitMs: ledger.totalQueueWaitMs + (now - request.enqueuedAtMs),\n\t\t\t\t\t\t\tupdatedAtMs: now,\n\t\t\t\t\t\t\trevision: ledger.revision + 1,\n\t\t\t\t\t\t};\n\t\t\t\t\t\treturn { kind: \"write\", next, value: { status: \"cancelled\" as const } };\n\t\t\t\t\t}\n\t\t\t\t\tconst rIndex = ledger.running.findIndex((r) => r.inferenceRequestId === inferenceRequestId);\n\t\t\t\t\tif (rIndex >= 0 && ledger.running[rIndex].lease?.ownerId === this._ownerId) {\n\t\t\t\t\t\treturn { kind: \"noop\", value: { status: \"running\" as const } };\n\t\t\t\t\t}\n\t\t\t\t\treturn { kind: \"noop\", value: { status: \"not_found\" as const } };\n\t\t\t\t},\n\t\t\t);\n\t\t\tif (result.status === \"ok\" && result.value.status !== \"not_found\") {\n\t\t\t\treturn { status: result.value.status, inferenceRequestId };\n\t\t\t}\n\t\t}\n\t\treturn { status: \"not_found\", inferenceRequestId };\n\t}\n\n\t// =========================================================================\n\t// Recovery (restart reconciliation)\n\t// =========================================================================\n\n\t/** Reconcile expired RUNNING leases after a scheduler/process restart. Never fabricates completion. */\n\tasync recover(options: { now?: number } = {}): Promise<InferenceRecoveryReport> {\n\t\tconst now = options.now ?? this._now();\n\t\tconst ids = await this._store.listResources();\n\t\tconst report: InferenceRecoveryReport = {\n\t\t\tscannedResources: ids.length,\n\t\t\treconciledRequests: [],\n\t\t\tunchangedRunning: [],\n\t\t\tcorruptResources: [],\n\t\t\tactions: [],\n\t\t};\n\n\t\tfor (const resourceId of ids) {\n\t\t\tconst loaded = await this._store.load(resourceId);\n\t\t\tif (loaded.status === \"corrupt\") {\n\t\t\t\treport.corruptResources.push({ resourceId, diagnostic: loaded.diagnostic });\n\t\t\t\treport.actions.push(`resource '${resourceId}' is corrupt; surfaced, not recovered`);\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (loaded.status === \"missing\") continue;\n\t\t\tconst ledger = loaded.ledger;\n\n\t\t\tconst expired = ledger.running.filter((r) => !r.lease || !isInferenceRequestLeaseActive(r.lease, now));\n\t\t\tif (expired.length === 0) {\n\t\t\t\tfor (const r of ledger.running) report.unchangedRunning.push(r.inferenceRequestId);\n\t\t\t\tcontinue;\n\t\t\t}\n\n\t\t\tconst expiredIds = new Set(expired.map((r) => r.inferenceRequestId));\n\t\t\tconst result = await this._store.mutate<{ reconciled: string[] }>(resourceId, (current) => {\n\t\t\t\tconst mutationNow = this._now();\n\t\t\t\tconst stillExpired = current.running.filter(\n\t\t\t\t\t(r) => !r.lease || !isInferenceRequestLeaseActive(r.lease, now),\n\t\t\t\t);\n\t\t\t\tif (stillExpired.length === 0) return { kind: \"noop\", value: { reconciled: [] } };\n\t\t\t\tconst idsToExpire = new Set(stillExpired.map((r) => r.inferenceRequestId));\n\n\t\t\t\tconst history = stillExpired.map((r) => ({\n\t\t\t\t\tinferenceRequestId: r.inferenceRequestId,\n\t\t\t\t\tlogicalAgentId: r.logicalAgentId,\n\t\t\t\t\tresourceId,\n\t\t\t\t\tprovider: r.provider,\n\t\t\t\t\tmodel: r.model,\n\t\t\t\t\trequestedAtMs: r.requestedAtMs,\n\t\t\t\t\tenqueuedAtMs: r.enqueuedAtMs,\n\t\t\t\t\tstate: \"INTERRUPTED\" as const,\n\t\t\t\t\tadmittedAtMs: r.admittedAtMs,\n\t\t\t\t\tfinishedAtMs: mutationNow,\n\t\t\t\t\tqueueWaitMs: r.admittedAtMs\n\t\t\t\t\t\t? Math.max(0, r.admittedAtMs - r.enqueuedAtMs)\n\t\t\t\t\t\t: Math.max(0, mutationNow - r.enqueuedAtMs),\n\t\t\t\t\tinferenceWallMs: r.admittedAtMs ? Math.max(0, mutationNow - r.admittedAtMs) : 0,\n\t\t\t\t}));\n\n\t\t\t\tlet next: InferenceResourceLedger = {\n\t\t\t\t\t...current,\n\t\t\t\t\trunning: current.running.filter((r) => !idsToExpire.has(r.inferenceRequestId)),\n\t\t\t\t\thistory: [...history, ...current.history].slice(0, 1000),\n\t\t\t\t\tinterruptedCount: current.interruptedCount + stillExpired.length,\n\t\t\t\t\tupdatedAtMs: mutationNow,\n\t\t\t\t\trevision: current.revision + 1,\n\t\t\t\t};\n\t\t\t\tnext = this._promoteQueued(next, mutationNow);\n\t\t\t\treturn { kind: \"write\", next, value: { reconciled: stillExpired.map((r) => r.inferenceRequestId) } };\n\t\t\t});\n\n\t\t\tif (result.status === \"ok\") {\n\t\t\t\tfor (const id of result.value.reconciled) {\n\t\t\t\t\treport.reconciledRequests.push(id);\n\t\t\t\t\treport.actions.push(`inference request '${id}' reconciled RUNNING → INTERRUPTED (expired lease)`);\n\t\t\t\t}\n\t\t\t} else if (result.status === \"corrupt\") {\n\t\t\t\treport.corruptResources.push({ resourceId, diagnostic: result.diagnostic });\n\t\t\t\treport.actions.push(`resource '${resourceId}' became corrupt during recovery`);\n\t\t\t}\n\t\t\tfor (const r of ledger.running) {\n\t\t\t\tif (!expiredIds.has(r.inferenceRequestId)) report.unchangedRunning.push(r.inferenceRequestId);\n\t\t\t}\n\t\t}\n\n\t\treturn report;\n\t}\n\n\t// =========================================================================\n\t// Status / telemetry\n\t// =========================================================================\n\n\tasync status(): Promise<SharedInferenceSchedulerStatus> {\n\t\tconst now = this._now();\n\t\tconst resources: SharedInferenceResourceStatus[] = [];\n\t\tconst queue: InferenceRequestStatus[] = [];\n\t\tlet totalSlots = 0;\n\t\tlet busySlots = 0;\n\t\tlet idleSlots = 0;\n\t\tlet queueDepth = 0;\n\t\tlet completedCount = 0;\n\n\t\tfor (const resourceId of await this._store.listResources()) {\n\t\t\tconst loaded = await this._store.load(resourceId);\n\t\t\tif (loaded.status !== \"ok\") continue;\n\t\t\tconst ledger = loaded.ledger;\n\t\t\tconst registered = this._resources.get(resourceId);\n\t\t\tconst capacity = registered?.capacity ?? ledger.capacity;\n\n\t\t\tconst contextActive = ledger.running.reduce((sum, r) => sum + (r.estimatedInputTokens ?? 0), 0);\n\t\t\tconst contextQueued = ledger.queue.reduce((sum, r) => sum + (r.estimatedInputTokens ?? 0), 0);\n\n\t\t\tresources.push({\n\t\t\t\tresourceId,\n\t\t\t\tbackend: registered?.backend ?? ledger.running[0]?.provider ?? ledger.queue[0]?.provider ?? \"\",\n\t\t\t\tmodel: registered?.model ?? ledger.running[0]?.model ?? ledger.queue[0]?.model ?? \"\",\n\t\t\t\tlocation: registered?.location ?? \"\",\n\t\t\t\tcapacity,\n\t\t\t\tbusySlots: ledger.running.length,\n\t\t\t\tidleSlots: Math.max(0, capacity - ledger.running.length),\n\t\t\t\tqueueDepth: ledger.queue.length,\n\t\t\t\tcompletedCount: ledger.completedCount,\n\t\t\t\tcancelledCount: ledger.cancelledCount,\n\t\t\t\tfailedCount: ledger.failedCount,\n\t\t\t\tinterruptedCount: ledger.interruptedCount,\n\t\t\t\ttotalQueueWaitMs: ledger.totalQueueWaitMs,\n\t\t\t\ttotalInferenceMs: ledger.totalInferenceMs,\n\t\t\t\ttotalInputTokens: ledger.totalInputTokens,\n\t\t\t\ttotalOutputTokens: ledger.totalOutputTokens,\n\t\t\t\tmaxQueueDepth: ledger.maxQueueDepth,\n\t\t\t\tcontextTokensActive: contextActive,\n\t\t\t\tcontextTokensQueued: contextQueued,\n\t\t\t\tstate: registered?.state ?? \"available\",\n\t\t\t\tobservedAtMs: now,\n\t\t\t});\n\n\t\t\ttotalSlots += capacity;\n\t\t\tbusySlots += ledger.running.length;\n\t\t\tidleSlots += Math.max(0, capacity - ledger.running.length);\n\t\t\tqueueDepth += ledger.queue.length;\n\t\t\tcompletedCount += ledger.completedCount;\n\n\t\t\tfor (const r of ledger.running) {\n\t\t\t\tqueue.push(\n\t\t\t\t\tthis._requestStatus(\n\t\t\t\t\t\tr,\n\t\t\t\t\t\tnow,\n\t\t\t\t\t\tledger.running.findIndex((x) => x.inferenceRequestId === r.inferenceRequestId) + 1,\n\t\t\t\t\t),\n\t\t\t\t);\n\t\t\t}\n\t\t\tconst queued = this._sortedQueue(ledger.queue, now);\n\t\t\tfor (let i = 0; i < queued.length; i++) queue.push(this._requestStatus(queued[i]!, now, i + 1));\n\t\t}\n\n\t\tthis._observeAvoidableIdle(now, queueDepth, idleSlots);\n\n\t\tresources.sort((a, b) => a.resourceId.localeCompare(b.resourceId));\n\t\tqueue.sort((a, b) =>\n\t\t\ta.effectivePriority === b.effectivePriority\n\t\t\t\t? a.queuedAtMs - b.queuedAtMs\n\t\t\t\t: b.effectivePriority - a.effectivePriority,\n\t\t);\n\n\t\treturn {\n\t\t\tresources,\n\t\t\tqueue,\n\t\t\taggregate: {\n\t\t\t\ttotalSlots,\n\t\t\t\tbusySlots,\n\t\t\t\tidleSlots,\n\t\t\t\tqueueDepth,\n\t\t\t\tcompletedCount,\n\t\t\t\tavoidableIdleMs: this._avoidableIdleTotalMs,\n\t\t\t},\n\t\t\tagents: { runningInference: 0, waitingInference: 0, tooling: 0, parked: 0, runnable: 0, total: 0 },\n\t\t\tobservedAtMs: now,\n\t\t};\n\t}\n\n\t/** Deterministic avoidable-idle observation: queue non-empty + free slot + no admission. */\n\tprivate _observeAvoidableIdle(now: number, queueDepth: number, idleSlots: number): void {\n\t\tconst avoidable = queueDepth > 0 && idleSlots > 0;\n\t\tif (avoidable) {\n\t\t\tif (this._avoidableIdleSinceMs === undefined) this._avoidableIdleSinceMs = now;\n\t\t\tthis._avoidableIdleTotalMs += now - (this._avoidableIdleSinceMs ?? now);\n\t\t\tthis._avoidableIdleSinceMs = now;\n\t\t} else {\n\t\t\tif (this._avoidableIdleSinceMs !== undefined) {\n\t\t\t\tthis._avoidableIdleTotalMs += now - this._avoidableIdleSinceMs;\n\t\t\t}\n\t\t\tthis._avoidableIdleSinceMs = undefined;\n\t\t}\n\t}\n\n\t// =========================================================================\n\t// Internals\n\t// =========================================================================\n\n\tprivate _toAdmitted(request: InferenceRequestRecord): AdmittedInference {\n\t\treturn {\n\t\t\tinferenceRequestId: request.inferenceRequestId,\n\t\t\tresourceId: request.resourceId,\n\t\t\tlogicalAgentId: request.logicalAgentId,\n\t\t\tslot: request.lease!.slot,\n\t\t\tlease: request.lease!,\n\t\t};\n\t}\n\n\tprivate _admittedFrom(ledger: InferenceResourceLedger, inferenceRequestId: string): AdmittedInference {\n\t\tconst request = ledger.running.find((r) => r.inferenceRequestId === inferenceRequestId)!;\n\t\treturn this._toAdmitted(request);\n\t}\n\n\tprivate _admitInto(\n\t\tledger: InferenceResourceLedger,\n\t\trequest: InferenceRequestRecord,\n\t\tnow: number,\n\t\toptions: { incrementFence: boolean },\n\t): InferenceResourceLedger {\n\t\tconst slot = ledger.running.length;\n\t\tconst fencingToken = options.incrementFence ? ledger.fencingToken + 1 : ledger.fencingToken;\n\t\tconst lease: InferenceRequestLease = {\n\t\t\townerId: request.ownerId,\n\t\t\tleaseId: this._leaseIdFactory(),\n\t\t\tslot,\n\t\t\tacquiredAtMs: now,\n\t\t\texpiresAtMs: now + this._leaseDurationMs,\n\t\t};\n\t\tconst admitted: InferenceRequestRecord = {\n\t\t\t...request,\n\t\t\tstate: \"RUNNING\",\n\t\t\tlease,\n\t\t\tadmittedAtMs: now,\n\t\t};\n\t\treturn {\n\t\t\t...ledger,\n\t\t\trunning: [...ledger.running, admitted],\n\t\t\tfencingToken,\n\t\t\tupdatedAtMs: now,\n\t\t\trevision: ledger.revision + 1,\n\t\t};\n\t}\n\n\t/** Sort the queue deterministically and promote up to free capacity. */\n\tprivate _promoteQueued(ledger: InferenceResourceLedger, now: number): InferenceResourceLedger {\n\t\tlet next = ledger;\n\t\tconst freeSlots = ledger.capacity - ledger.running.length;\n\t\tif (freeSlots <= 0 || ledger.queue.length === 0) return next;\n\n\t\tconst sorted = this._sortedQueue(ledger.queue, now);\n\t\tconst promote = sorted.slice(0, freeSlots);\n\t\tconst remaining = sorted.slice(freeSlots);\n\n\t\tlet running = ledger.running;\n\t\tlet fencingToken = ledger.fencingToken;\n\t\tfor (const request of promote) {\n\t\t\tconst admitted: InferenceRequestRecord = {\n\t\t\t\t...request,\n\t\t\t\tstate: \"RUNNING\",\n\t\t\t\tlease: {\n\t\t\t\t\townerId: request.ownerId,\n\t\t\t\t\tleaseId: this._leaseIdFactory(),\n\t\t\t\t\tslot: running.length,\n\t\t\t\t\tacquiredAtMs: now,\n\t\t\t\t\texpiresAtMs: now + this._leaseDurationMs,\n\t\t\t\t},\n\t\t\t\tadmittedAtMs: now,\n\t\t\t};\n\t\t\trunning = [...running, admitted];\n\t\t\tfencingToken += 1;\n\t\t}\n\n\t\tnext = {\n\t\t\t...ledger,\n\t\t\trunning,\n\t\t\tqueue: remaining,\n\t\t\tfencingToken,\n\t\t\tupdatedAtMs: now,\n\t\t\trevision: ledger.revision + 1,\n\t\t};\n\t\treturn this._normalizeSlots(next);\n\t}\n\n\t/** Keep running[i].lease.slot === i. */\n\tprivate _normalizeSlots(ledger: InferenceResourceLedger): InferenceResourceLedger {\n\t\tconst running = ledger.running.map((r, i) => (r.lease ? { ...r, lease: { ...r.lease, slot: i } } : r));\n\t\treturn { ...ledger, running };\n\t}\n\n\tprivate _sortedQueue(queue: InferenceRequestRecord[], now: number): InferenceRequestRecord[] {\n\t\treturn [...queue].sort((a, b) => this._compare(a, b, now));\n\t}\n\n\tprivate _compare(a: InferenceRequestRecord, b: InferenceRequestRecord, now: number): number {\n\t\tconst pa = this._breakdown(a, now).effective;\n\t\tconst pb = this._breakdown(b, now).effective;\n\t\tif (pa !== pb) return pb - pa;\n\t\tif (a.enqueuedAtMs !== b.enqueuedAtMs) return a.enqueuedAtMs - b.enqueuedAtMs;\n\t\treturn a.inferenceRequestId < b.inferenceRequestId ? -1 : a.inferenceRequestId > b.inferenceRequestId ? 1 : 0;\n\t}\n\n\tprivate _breakdown(request: InferenceRequestRecord, now: number): PriorityBreakdown {\n\t\tconst interactive = request.priority.interactive ? this._interactiveBoost : 0;\n\t\tconst verification = request.priority.verification ? this._verificationBoost : 0;\n\t\tconst dependency = request.dependency\n\t\t\t? Math.min(request.dependency.unblocksCount * this._dependencyWeight, this._dependencyCap)\n\t\t\t: 0;\n\t\tconst ageMs = Math.max(0, now - request.enqueuedAtMs);\n\t\tconst aging = Math.floor(ageMs / this._agingIntervalMs) * this._agingWeight;\n\t\tconst base = request.priority.base;\n\t\treturn {\n\t\t\tbase,\n\t\t\tinteractive,\n\t\t\tverification,\n\t\t\tdependency,\n\t\t\taging,\n\t\t\teffective: base + interactive + verification + dependency + aging,\n\t\t};\n\t}\n\n\tprivate _requestStatus(request: InferenceRequestRecord, now: number, position: number): InferenceRequestStatus {\n\t\tconst breakdown = this._breakdown(request, now);\n\t\treturn {\n\t\t\tinferenceRequestId: request.inferenceRequestId,\n\t\t\tlogicalAgentId: request.logicalAgentId,\n\t\t\tresourceId: request.resourceId,\n\t\t\tstate: request.state,\n\t\t\tposition: request.state === \"QUEUED\" ? position : undefined,\n\t\t\teffectivePriority: breakdown.effective,\n\t\t\tbasePriority: breakdown.base,\n\t\t\tagingContribution: breakdown.aging,\n\t\t\tdependencyContribution: breakdown.dependency,\n\t\t\tqueuedAtMs: request.enqueuedAtMs,\n\t\t\trequestedAtMs: request.requestedAtMs,\n\t\t};\n\t}\n}\n"]}