{"version":3,"file":"scheduler-worker-driver.d.ts","sourceRoot":"","sources":["../../../src/core/orchestration/scheduler-worker-driver.ts"],"names":[],"mappings":"AAAA;;;;;;;;;GASG;AAEH,OAAO,KAAK,EAAE,wBAAwB,EAAE,MAAM,6CAA6C,CAAC;AAC5F,OAAO,KAAK,EAAE,mBAAmB,EAAE,MAAM,oCAAoC,CAAC;AAE9E,OAAO,KAAK,EAAE,uBAAuB,EAAE,MAAM,2CAA2C,CAAC;AAEzF,OAAO,KAAK,EAAE,oBAAoB,EAAE,MAAM,4CAA4C,CAAC;AAEvF,MAAM,WAAW,qCAAqC;IACrD,QAAQ,EAAE,mBAAmB,CAAC;IAC9B,SAAS,EAAE,uBAAuB,CAAC;IACnC,WAAW,EAAE,wBAAwB,CAAC;CACtC;AAED,MAAM,WAAW,oCAAoC;IACpD,QAAQ,EAAE,MAAM,EAAE,CAAC;IACnB,mBAAmB,EAAE,MAAM,EAAE,CAAC;IAC9B,gBAAgB,EAAE,MAAM,EAAE,CAAC;IAC3B,6BAA6B,EAAE,MAAM,EAAE,CAAC;IACxC,SAAS,EAAE,MAAM,EAAE,CAAC;CACpB;AAED;;;;;GAKG;AACH,qBAAa,8BAA8B;IAC1C,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAsB;IAChD,OAAO,CAAC,QAAQ,CAAC,UAAU,CAA0B;IACrD,OAAO,CAAC,QAAQ,CAAC,YAAY,CAA2B;IAExD,YAAY,OAAO,EAAE,qCAAqC,EAIzD;IAEK,OAAO,CAAC,eAAe,EAAE,MAAM,GAAG,OAAO,CAAC,oCAAoC,CAAC,CA6DpF;CACD;AAED,MAAM,WAAW,4BAA4B;IAC5C,SAAS,EAAE,uBAAuB,CAAC;IACnC,MAAM,EAAE,oBAAoB,CAAC;IAC7B,4FAA4F;IAC5F,eAAe,EAAE,8BAA8B,CAAC;IAChD,0EAA0E;IAC1E,QAAQ,CAAC,EAAE,MAAM,CAAC;IAClB,4EAA4E;IAC5E,cAAc,CAAC,EAAE,MAAM,CAAC;CACxB;AAED,MAAM,WAAW,+BAA+B;IAC/C,MAAM,CAAC,EAAE,WAAW,CAAC;IACrB,eAAe,CAAC,EAAE,MAAM,CAAC;CACzB;AAED,qFAAqF;AACrF,qBAAa,4BAA6B,SAAQ,KAAK;IACtD,QAAQ,CAAC,KAAK,EAAE,OAAO,CAAC;IAExB,YAAY,OAAO,EAAE,MAAM,EAAE,KAAK,CAAC,EAAE,OAAO,EAI3C;CACD;AAoCD;;;GAGG;AACH,qBAAa,qBAAqB;IACjC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAA0B;IACrD,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAuB;IAC/C,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAiC;IAClE,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAS;IACnC,OAAO,CAAC,QAAQ,CAAC,eAAe,CAAS;IACzC,OAAO,CAAC,OAAO,CAAS;IAExB,YAAY,OAAO,EAAE,4BAA4B,EAUhD;IAED,6EAA6E;IACvE,OAAO,CAAC,CAAC,EACd,SAAS,EAAE,CAAC,MAAM,EAAE,WAAW,KAAK,OAAO,CAAC,CAAC,CAAC,EAC9C,OAAO,GAAE,+BAAoC,GAC3C,OAAO,CAAC,CAAC,CAAC,CA+EZ;YAEa,KAAK;CASnB;AAED,wBAAgB,2BAA2B,CAAC,OAAO,EAAE,4BAA4B,GAAG,qBAAqB,CAExG","sourcesContent":["/**\n * Bounded in-process Scheduler -> Worker application driver.\n *\n * The driver supplies the optional application seam between a parent\n * orchestration execution and the existing control-plane services. It starts\n * one WorkerControlService, serially runs scheduler ticks followed by worker\n * passes, and stops the worker when the parent operation completes. It never\n * writes mission state or launches children itself; Scheduler, Assignment,\n * Worker, and the parent DurableMissionCoordinator retain their authorities.\n */\n\nimport type { AssignmentControlService } from \"../assignment/assignment-control-service.js\";\nimport type { DurableMissionStore } from \"../mission-domain/durable-store.js\";\nimport { isTerminalMissionState } from \"../mission-domain/mission-state.js\";\nimport type { SchedulerControlService } from \"../scheduler/scheduler-control-service.js\";\nimport { SchedulingError } from \"../scheduler/scheduler-types.js\";\nimport type { WorkerControlService } from \"../worker-daemon/worker-control-service.js\";\n\nexport interface SchedulerWorkerTerminalCleanupOptions {\n\tmissions: DurableMissionStore;\n\tscheduler: SchedulerControlService;\n\tassignments: AssignmentControlService;\n}\n\nexport interface SchedulerWorkerTerminalCleanupReport {\n\tchildren: string[];\n\treleasedAssignments: string[];\n\tcancelledIntents: string[];\n\tcompletedExecutingAssignments: string[];\n\tresiduals: string[];\n}\n\n/**\n * Bounded parent-terminal cleanup over the existing control-plane authorities.\n * It never writes mission state: active execution is stopped by Worker, pending\n * intent is cancelled by Scheduler, and assignment ownership is released or\n * completed by Assignment using the terminal mission result as evidence.\n */\nexport class SchedulerWorkerTerminalCleanup {\n\tprivate readonly _missions: DurableMissionStore;\n\tprivate readonly _scheduler: SchedulerControlService;\n\tprivate readonly _assignments: AssignmentControlService;\n\n\tconstructor(options: SchedulerWorkerTerminalCleanupOptions) {\n\t\tthis._missions = options.missions;\n\t\tthis._scheduler = options.scheduler;\n\t\tthis._assignments = options.assignments;\n\t}\n\n\tasync cleanup(parentMissionId: string): Promise<SchedulerWorkerTerminalCleanupReport> {\n\t\tconst children = await this._missions.listChildren(parentMissionId);\n\t\tconst report: SchedulerWorkerTerminalCleanupReport = {\n\t\t\tchildren,\n\t\t\treleasedAssignments: [],\n\t\t\tcancelledIntents: [],\n\t\t\tcompletedExecutingAssignments: [],\n\t\t\tresiduals: [],\n\t\t};\n\n\t\tfor (const missionId of children) {\n\t\t\tconst current = await this._assignments.getCurrentForMission(missionId);\n\t\t\tif (current?.state === \"ASSIGNED\" || current?.state === \"ACCEPTED\") {\n\t\t\t\tawait this._assignments.releaseAssignment(current.assignmentId);\n\t\t\t\treport.releasedAssignments.push(current.assignmentId);\n\t\t\t} else if (current?.state === \"EXECUTING\") {\n\t\t\t\tconst mission = await this._missions.load(missionId);\n\t\t\t\tif (mission.status === \"ok\" && mission.record.result && isTerminalMissionState(mission.record.state)) {\n\t\t\t\t\tconst lastAttempt = mission.record.attempts[mission.record.attempts.length - 1];\n\t\t\t\t\tawait this._assignments.completeAssignment(current.assignmentId, {\n\t\t\t\t\t\tresultState: mission.record.state,\n\t\t\t\t\t\tattemptId: lastAttempt?.attemptId,\n\t\t\t\t\t\texecutionId: mission.record.resultExecutionId ?? lastAttempt?.executionId,\n\t\t\t\t\t\treason: \"parent terminal cleanup: mission already terminal\",\n\t\t\t\t\t});\n\t\t\t\t\treport.completedExecutingAssignments.push(current.assignmentId);\n\t\t\t\t} else {\n\t\t\t\t\treport.residuals.push(\n\t\t\t\t\t\t`EXECUTING assignment ${current.assignmentId} for nonterminal child ${missionId} cannot be interrupted safely`,\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t}\n\n\t\t\ttry {\n\t\t\t\tconst intent = await this._scheduler.getIntentForMission(missionId);\n\t\t\t\tif (intent.state === \"PENDING\" || intent.state === \"UNSCHEDULABLE\") {\n\t\t\t\t\tawait this._scheduler.cancelIntent(missionId);\n\t\t\t\t\treport.cancelledIntents.push(intent.intentId);\n\t\t\t\t}\n\t\t\t} catch (error) {\n\t\t\t\tif (!(error instanceof SchedulingError) || error.code !== \"INTENT_NOT_FOUND\") throw error;\n\t\t\t}\n\t\t}\n\n\t\tfor (const missionId of children) {\n\t\t\tconst current = await this._assignments.getCurrentForMission(missionId);\n\t\t\tif (current?.state === \"ASSIGNED\" || current?.state === \"ACCEPTED\" || current?.state === \"EXECUTING\")\n\t\t\t\treport.residuals.push(\n\t\t\t\t\t`current ${current.state} assignment ${current.assignmentId} remains for child ${missionId}`,\n\t\t\t\t);\n\t\t\ttry {\n\t\t\t\tconst intent = await this._scheduler.getIntentForMission(missionId);\n\t\t\t\tif (intent.state === \"PENDING\")\n\t\t\t\t\treport.residuals.push(`executable PENDING intent ${intent.intentId} remains for child ${missionId}`);\n\t\t\t} catch (error) {\n\t\t\t\tif (!(error instanceof SchedulingError) || error.code !== \"INTENT_NOT_FOUND\") throw error;\n\t\t\t}\n\t\t}\n\t\tif (report.residuals.length > 0)\n\t\t\tthrow new Error(`PARENT_TERMINAL_CLEANUP_INCOMPLETE: ${report.residuals.join(\"; \")}`);\n\t\treturn report;\n\t}\n}\n\nexport interface SchedulerWorkerDriverOptions {\n\tscheduler: SchedulerControlService;\n\tworker: WorkerControlService;\n\t/** Mandatory parent-terminal cleanup over the existing Scheduler/Assignment authorities. */\n\tterminalCleanup: SchedulerWorkerTerminalCleanup;\n\t/** Maximum number of scheduler/worker passes for one parent operation. */\n\tmaxTicks?: number;\n\t/** Delay between passes. Defaults to 0 for application/test composition. */\n\ttickIntervalMs?: number;\n}\n\nexport interface SchedulerWorkerDriverRunOptions {\n\tsignal?: AbortSignal;\n\tparentMissionId?: string;\n}\n\n/** Structured scheduler/worker driver failure, distinct from caller cancellation. */\nexport class SchedulerWorkerDriverFailure extends Error {\n\treadonly cause: unknown;\n\n\tconstructor(message: string, cause?: unknown) {\n\t\tsuper(message);\n\t\tthis.name = \"SchedulerWorkerDriverFailure\";\n\t\tthis.cause = cause;\n\t}\n}\n\nfunction isParentFailedOutcome(result: unknown): boolean {\n\treturn (\n\t\ttypeof result === \"object\" &&\n\t\tresult !== null &&\n\t\t\"state\" in result &&\n\t\t(result as { state?: unknown }).state === \"FAILED\"\n\t);\n}\n\nfunction isParentTerminalCleanupRequired(result: unknown): boolean {\n\tif (typeof result !== \"object\" || result === null || !(\"state\" in result)) return true;\n\tconst state = (result as { state?: unknown }).state;\n\treturn state === \"CANCELLED\" || state === \"FAILED\" || state === \"TIMED_OUT\" || state === \"CRASHED\";\n}\n\nfunction sleep(ms: number, signal: AbortSignal): Promise<void> {\n\tif (ms <= 0) return Promise.resolve();\n\treturn new Promise((resolve) => {\n\t\tif (signal.aborted) {\n\t\t\tresolve();\n\t\t\treturn;\n\t\t}\n\t\tconst timer = setTimeout(resolve, ms);\n\t\tsignal.addEventListener(\n\t\t\t\"abort\",\n\t\t\t() => {\n\t\t\t\tclearTimeout(timer);\n\t\t\t\tresolve();\n\t\t\t},\n\t\t\t{ once: true },\n\t\t);\n\t});\n}\n\n/**\n * Drives a parent operation while an optional local worker consumes its child\n * assignments. One driver instance owns one serialized worker run at a time.\n */\nexport class SchedulerWorkerDriver {\n\tprivate readonly _scheduler: SchedulerControlService;\n\tprivate readonly _worker: WorkerControlService;\n\tprivate readonly _terminalCleanup: SchedulerWorkerTerminalCleanup;\n\tprivate readonly _maxTicks: number;\n\tprivate readonly _tickIntervalMs: number;\n\tprivate _active = false;\n\n\tconstructor(options: SchedulerWorkerDriverOptions) {\n\t\tthis._scheduler = options.scheduler;\n\t\tthis._worker = options.worker;\n\t\tthis._terminalCleanup = options.terminalCleanup;\n\t\tthis._maxTicks = options.maxTicks ?? 10_000;\n\t\tthis._tickIntervalMs = options.tickIntervalMs ?? 0;\n\t\tif (!Number.isSafeInteger(this._maxTicks) || this._maxTicks < 1)\n\t\t\tthrow new Error(\"SCHEDULER_WORKER_DRIVER_INVALID: maxTicks must be a positive safe integer\");\n\t\tif (!Number.isFinite(this._tickIntervalMs) || this._tickIntervalMs < 0)\n\t\t\tthrow new Error(\"SCHEDULER_WORKER_DRIVER_INVALID: tickIntervalMs must be non-negative\");\n\t}\n\n\t/** Execute an operation with the bounded Scheduler -> Worker pump active. */\n\tasync execute<T>(\n\t\toperation: (signal: AbortSignal) => Promise<T>,\n\t\toptions: SchedulerWorkerDriverRunOptions = {},\n\t): Promise<T> {\n\t\tif (this._active) throw new Error(\"SCHEDULER_WORKER_DRIVER_ACTIVE: driver is already running\");\n\t\tthis._active = true;\n\t\tconst controller = new AbortController();\n\t\tconst callerSignal = options.signal;\n\t\tconst abortFromCaller = () => controller.abort(callerSignal?.reason);\n\t\tif (callerSignal) {\n\t\t\tif (callerSignal.aborted) controller.abort(callerSignal.reason);\n\t\t\telse callerSignal.addEventListener(\"abort\", abortFromCaller, { once: true });\n\t\t}\n\n\t\tlet pumpFailure: SchedulerWorkerDriverFailure | undefined;\n\t\tlet pump: Promise<void> | undefined;\n\t\tlet result: T | undefined;\n\t\tlet operationCompleted = false;\n\t\tlet cleanupRequired = false;\n\t\tlet primaryError: unknown;\n\t\tconst finalizationErrors: unknown[] = [];\n\t\ttry {\n\t\t\ttry {\n\t\t\t\tawait this._worker.start({ reconcile: false, polling: false });\n\t\t\t\tpump = this._pump(controller.signal).catch((error: unknown) => {\n\t\t\t\t\tconst failure = new SchedulerWorkerDriverFailure(\"scheduler/worker driver failed\", error);\n\t\t\t\t\tpumpFailure = failure;\n\t\t\t\t\tcontroller.abort(failure);\n\t\t\t\t});\n\t\t\t\tresult = await operation(controller.signal);\n\t\t\t\toperationCompleted = true;\n\t\t\t\tcleanupRequired = isParentTerminalCleanupRequired(result);\n\t\t\t} catch (error) {\n\t\t\t\tprimaryError = error;\n\t\t\t\tcontroller.abort(error);\n\t\t\t}\n\n\t\t\tif (pumpFailure !== undefined || primaryError !== undefined) cleanupRequired = true;\n\t\t} finally {\n\t\t\tcontroller.abort(\"scheduler/worker driver cleanup\");\n\t\t\t// stop() aborts the worker's in-flight child before awaiting it. This\n\t\t\t// ordering prevents cleanup from waiting on a runOnce that cleanup\n\t\t\t// itself must cancel. Finalization errors are collected so they cannot\n\t\t\t// replace the operation or pump failure that caused the shutdown.\n\t\t\ttry {\n\t\t\t\tawait this._worker.stop(\"scheduler/worker driver cleanup\");\n\t\t\t} catch (error) {\n\t\t\t\tfinalizationErrors.push(error);\n\t\t\t}\n\t\t\tif (pump) {\n\t\t\t\ttry {\n\t\t\t\t\tawait pump;\n\t\t\t\t} catch (error) {\n\t\t\t\t\tfinalizationErrors.push(error);\n\t\t\t\t}\n\t\t\t}\n\t\t\tif (pumpFailure !== undefined) cleanupRequired = true;\n\t\t\tif ((cleanupRequired || pumpFailure !== undefined) && options.parentMissionId) {\n\t\t\t\ttry {\n\t\t\t\t\tawait this._terminalCleanup.cleanup(options.parentMissionId);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tfinalizationErrors.push(error);\n\t\t\t\t}\n\t\t\t}\n\t\t\tif (callerSignal) callerSignal.removeEventListener(\"abort\", abortFromCaller);\n\t\t\tthis._active = false;\n\t\t}\n\n\t\t// A coordinator result with FAILED state is the durable authority for a\n\t\t// pump failure. Do not replace it with a caller-only exception. For\n\t\t// non-durable callers, retain the pump failure as the primary error.\n\t\tif (pumpFailure !== undefined && !isParentFailedOutcome(result) && !callerSignal?.aborted)\n\t\t\tprimaryError = pumpFailure;\n\t\tif (operationCompleted && result !== undefined && isParentTerminalCleanupRequired(result)) cleanupRequired = true;\n\n\t\tif (primaryError !== undefined && finalizationErrors.length > 0)\n\t\t\tthrow new AggregateError([primaryError, ...finalizationErrors], \"scheduler/worker driver finalization failed\");\n\t\tif (primaryError !== undefined) throw primaryError;\n\t\tif (finalizationErrors.length === 1) throw finalizationErrors[0];\n\t\tif (finalizationErrors.length > 1)\n\t\t\tthrow new AggregateError(finalizationErrors, \"scheduler/worker driver finalization failed\");\n\t\treturn result as T;\n\t}\n\n\tprivate async _pump(signal: AbortSignal): Promise<void> {\n\t\tfor (let tick = 0; tick < this._maxTicks && !signal.aborted; tick++) {\n\t\t\tawait this._scheduler.runTick();\n\t\t\tif (signal.aborted) break;\n\t\t\tawait this._worker.runOnce();\n\t\t\tawait sleep(this._tickIntervalMs, signal);\n\t\t}\n\t\tif (!signal.aborted) throw new Error(\"SCHEDULER_WORKER_DRIVER_LIMIT: bounded tick limit exhausted\");\n\t}\n}\n\nexport function createSchedulerWorkerDriver(options: SchedulerWorkerDriverOptions): SchedulerWorkerDriver {\n\treturn new SchedulerWorkerDriver(options);\n}\n"]}