{"version":3,"file":"orchestration-mission-executor.d.ts","sourceRoot":"","sources":["../../../src/core/orchestration/orchestration-mission-executor.ts"],"names":[],"mappings":"AAAA;;;;;;;;GAQG;AAGH,OAAO,KAAK,EAAE,mBAAmB,EAAE,MAAM,oCAAoC,CAAC;AAC9E,OAAO,KAAK,EAAE,eAAe,EAAE,oBAAoB,EAAE,MAAM,uCAAuC,CAAC;AACnG,OAAO,EAAuB,KAAK,aAAa,EAAE,MAAM,qCAAqC,CAAC;AAC9F,OAAO,KAAK,EAAE,cAAc,EAAE,MAAM,sCAAsC,CAAC;AAC3E,OAAO,EAAuB,KAAK,aAAa,EAA0B,MAAM,qCAAqC,CAAC;AACtH,OAAO,KAAK,EAAE,8BAA8B,EAAE,MAAM,yBAAyB,CAAC;AAC9E,OAAO,KAAK,EAAE,mBAAmB,EAAE,MAAM,mBAAmB,CAAC;AAE7D,OAAO,KAAK,EAAqB,kBAAkB,EAAE,MAAM,YAAY,CAAC;AAExE,MAAM,WAAW,mCAAmC;IACnD,QAAQ,EAAE,mBAAmB,CAAC;IAC9B,KAAK,EAAE,kBAAkB,CAAC;IAC1B,4EAA4E;IAC5E,SAAS,EAAE,8BAA8B,CAAC;IAC1C,0EAA0E;IAC1E,YAAY,EAAE,mBAAmB,CAAC;IAClC,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,aAAa,CAAC,EAAE,MAAM,CAAC;IACvB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,kBAAkB,CAAC,EAAE,CAAC,OAAO,EAAE,cAAc,KAAK,MAAM,CAAC;CACzD;AA6BD,qBAAa,4BAA6B,YAAW,eAAe;IACnE,QAAQ,CAAC,UAAU,EAAE,MAAM,CAAC;IAC5B,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAqB;IAC5C,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAiC;IAC5D,OAAO,CAAC,QAAQ,CAAC,aAAa,CAAsB;IACpD,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAS;IACjC,OAAO,CAAC,QAAQ,CAAC,cAAc,CAAS;IACxC,OAAO,CAAC,QAAQ,CAAC,IAAI,CAAe;IACpC,OAAO,CAAC,QAAQ,CAAC,mBAAmB,CAAsC;IAC1E,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAsC;IAE9D,YAAY,OAAO,EAAE,mCAAmC,EASvD;IAEK,MAAM,CAAC,OAAO,EAAE,cAAc,EAAE,OAAO,GAAE,oBAAyB,GAAG,OAAO,CAAC,aAAa,CAAC,CAsDhG;IAEK,WAAW,CAAC,MAAM,EAAE,aAAa,EAAE,OAAO,GAAE;QAAE,MAAM,CAAC,EAAE,WAAW,CAAA;KAAO,GAAG,OAAO,CAAC,aAAa,CAAC,CAsFvG;IAEK,MAAM,CAAC,MAAM,EAAE,aAAa,EAAE,MAAM,SAA4B,GAAG,OAAO,CAAC,IAAI,CAAC,CAOrF;YAEa,QAAQ;IAOtB,OAAO,CAAC,MAAM;CA4Cd;AAED,wBAAgB,kCAAkC,CACjD,OAAO,EAAE,mCAAmC,GAC1C,4BAA4B,CAE9B","sourcesContent":["/**\n * Durable parent orchestration MissionExecutor.\n *\n * This adapter is deliberately subordinate to the Mission domain: the\n * DurableMissionCoordinator owns the parent lease and terminal commit, while\n * this executor drives the orchestration plan through the existing lifecycle\n * and child-execution authority ports. It never launches a process, opens SSH,\n * calls a provider, or writes parent mission state directly.\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport type { DurableMissionStore } from \"../mission-domain/durable-store.js\";\nimport type { MissionExecutor, MissionLaunchOptions } from \"../mission-domain/mission-executor.js\";\nimport { createMissionHandle, type MissionHandle } from \"../mission-domain/mission-handle.js\";\nimport type { MissionRequest } from \"../mission-domain/mission-request.js\";\nimport { createMissionResult, type MissionResult, type StructuredFailure } from \"../mission-domain/mission-result.js\";\nimport type { OrchestrationLifecycleExecutor } from \"./lifecycle-executor.js\";\nimport type { OrchestratorService } from \"./orchestrator.js\";\nimport { SchedulerWorkerDriverFailure } from \"./scheduler-worker-driver.js\";\nimport type { OrchestrationPlan, OrchestrationStore } from \"./types.js\";\n\nexport interface OrchestrationMissionExecutorOptions {\n\tmissions: DurableMissionStore;\n\tstore: OrchestrationStore;\n\t/** Existing lifecycle authority; no child execution is implemented here. */\n\tlifecycle: OrchestrationLifecycleExecutor;\n\t/** Reconciliation service used to unblock/materialize dependent nodes. */\n\torchestrator: OrchestratorService;\n\texecutorId?: string;\n\tpollMs?: number;\n\tmaxWallTimeMs?: number;\n\tnow?: () => number;\n\texecutionIdFactory?: (request: MissionRequest) => string;\n}\n\ninterface ActiveExecution {\n\trequest: MissionRequest;\n\texecutionId: string;\n\tstartedAtMs: number;\n\tcontroller: AbortController;\n\tlaunched: boolean;\n\tcancelled: boolean;\n}\n\nfunction sleep(ms: number, signal: AbortSignal): Promise<void> {\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\nexport class OrchestrationMissionExecutor implements MissionExecutor {\n\treadonly executorId: string;\n\tprivate readonly _store: OrchestrationStore;\n\tprivate readonly _lifecycle: OrchestrationLifecycleExecutor;\n\tprivate readonly _orchestrator: OrchestratorService;\n\tprivate readonly _pollMs: number;\n\tprivate readonly _maxWallTimeMs: number;\n\tprivate readonly _now: () => number;\n\tprivate readonly _executionIdFactory: (request: MissionRequest) => string;\n\tprivate readonly _active = new Map<string, ActiveExecution>();\n\n\tconstructor(options: OrchestrationMissionExecutorOptions) {\n\t\tthis.executorId = options.executorId ?? \"orchestration\";\n\t\tthis._store = options.store;\n\t\tthis._lifecycle = options.lifecycle;\n\t\tthis._orchestrator = options.orchestrator;\n\t\tthis._pollMs = options.pollMs ?? 250;\n\t\tthis._maxWallTimeMs = options.maxWallTimeMs ?? 30 * 60_000;\n\t\tthis._now = options.now ?? (() => Date.now());\n\t\tthis._executionIdFactory = options.executionIdFactory ?? (() => `exec_${randomUUID()}`);\n\t}\n\n\tasync launch(request: MissionRequest, options: MissionLaunchOptions = {}): Promise<MissionHandle> {\n\t\tconst contract = request.orchestrationExecution;\n\t\tif (!contract) throw new Error(`ORCHESTRATION_EXECUTION_REQUIRED: mission ${request.missionId}`);\n\t\tconst plan = await this.loadPlan(contract.orchestrationId);\n\t\tif (plan.parentMissionId !== request.missionId)\n\t\t\tthrow new Error(\n\t\t\t\t`ORCHESTRATION_PARENT_MISMATCH: ${contract.orchestrationId} belongs to ${plan.parentMissionId}`,\n\t\t\t);\n\n\t\tconst controller = new AbortController();\n\t\tif (options.signal) {\n\t\t\tif (options.signal.aborted) controller.abort(options.signal.reason);\n\t\t\telse options.signal.addEventListener(\"abort\", () => controller.abort(options.signal!.reason), { once: true });\n\t\t}\n\t\tconst active: ActiveExecution = {\n\t\t\trequest,\n\t\t\texecutionId: this._executionIdFactory(request),\n\t\t\tstartedAtMs: this._now(),\n\t\t\tcontroller,\n\t\t\tlaunched: false,\n\t\t\tcancelled: false,\n\t\t};\n\t\tthis._active.set(request.missionId, active);\n\t\ttry {\n\t\t\tawait this._orchestrator.reconcile(contract.orchestrationId);\n\t\t\tawait this._lifecycle.launchChildren(request.missionId);\n\t\t\tactive.launched = true;\n\t\t} catch (error) {\n\t\t\tthis._active.delete(request.missionId);\n\t\t\tthrow error;\n\t\t}\n\t\treturn createMissionHandle({\n\t\t\tmissionId: request.missionId,\n\t\t\tparentMissionId: request.parentMissionId,\n\t\t\tdepth: request.depth,\n\t\t\texecutionId: active.executionId,\n\t\t\tstate: \"RUNNING\",\n\t\t\tcreatedAtMs: request.createdAtMs,\n\t\t\tstartedAtMs: active.startedAtMs,\n\t\t\tcancel: (reason) =>\n\t\t\t\tthis.cancel(\n\t\t\t\t\t{\n\t\t\t\t\t\tmissionId: request.missionId,\n\t\t\t\t\t\tparentMissionId: request.parentMissionId,\n\t\t\t\t\t\tdepth: request.depth,\n\t\t\t\t\t\texecutionId: active.executionId,\n\t\t\t\t\t\tstate: \"RUNNING\",\n\t\t\t\t\t\tcreatedAtMs: request.createdAtMs,\n\t\t\t\t\t\tstartedAtMs: active.startedAtMs,\n\t\t\t\t\t\tcancel: async () => undefined,\n\t\t\t\t\t},\n\t\t\t\t\treason,\n\t\t\t\t),\n\t\t});\n\t}\n\n\tasync awaitResult(handle: MissionHandle, options: { signal?: AbortSignal } = {}): Promise<MissionResult> {\n\t\tconst active = this._active.get(handle.missionId);\n\t\tif (!active) throw new Error(`Unknown mission: ${handle.missionId}`);\n\t\tconst signal = active.controller.signal;\n\t\tif (options.signal) {\n\t\t\tif (options.signal.aborted) active.controller.abort(options.signal.reason);\n\t\t\telse\n\t\t\t\toptions.signal.addEventListener(\"abort\", () => active.controller.abort(options.signal!.reason), {\n\t\t\t\t\tonce: true,\n\t\t\t\t});\n\t\t}\n\t\tconst deadline = this._now() + this._maxWallTimeMs;\n\t\ttry {\n\t\t\tfor (;;) {\n\t\t\t\tif (signal.reason instanceof SchedulerWorkerDriverFailure)\n\t\t\t\t\treturn this.result(active, \"FAILED\", signal.reason.message);\n\t\t\t\tif (signal.aborted || active.cancelled) return this.result(active, \"CANCELLED\", \"orchestration cancelled\");\n\t\t\t\tif (this._now() >= deadline) {\n\t\t\t\t\tawait this._orchestrator\n\t\t\t\t\t\t.cancel(active.request.orchestrationExecution!.orchestrationId)\n\t\t\t\t\t\t.catch(() => undefined);\n\t\t\t\t\treturn this.result(active, \"TIMED_OUT\", \"orchestration deadline exceeded\");\n\t\t\t\t}\n\t\t\t\tconst contract = active.request.orchestrationExecution!;\n\t\t\t\tawait this._orchestrator.reconcile(contract.orchestrationId);\n\t\t\t\tawait this._lifecycle.launchChildren(active.request.missionId);\n\t\t\t\tconst plan = await this.loadPlan(contract.orchestrationId);\n\t\t\t\tconst statuses = [] as Array<{\n\t\t\t\t\trequirement: string;\n\t\t\t\t\tstatus: Awaited<ReturnType<OrchestrationLifecycleExecutor[\"childStatus\"]>>;\n\t\t\t\t}>;\n\t\t\t\tfor (const node of plan.nodes) {\n\t\t\t\t\tif (!node.childMissionId || !node.childSessionId) continue;\n\t\t\t\t\tstatuses.push({\n\t\t\t\t\t\trequirement: node.requirement,\n\t\t\t\t\t\tstatus: await this._lifecycle.childStatus(active.request.missionId, node.nodeId),\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t\tif (plan.decision === \"DIRECT\" && plan.nodes.length === 0)\n\t\t\t\t\treturn this.result(active, \"SUCCEEDED\", \"direct orchestration completed\", \"accepted\");\n\t\t\t\tconst required = statuses.filter((entry) => entry.requirement !== \"OPTIONAL\");\n\t\t\t\tconst requiredFailure = required.find(\n\t\t\t\t\t(entry) =>\n\t\t\t\t\t\tentry.status.terminal &&\n\t\t\t\t\t\tentry.status.missionState !== \"PARTIAL\" &&\n\t\t\t\t\t\t(!entry.status.success ||\n\t\t\t\t\t\t\tentry.status.verificationStatus === \"failed\" ||\n\t\t\t\t\t\t\tentry.status.completionDecision === \"rejected\"),\n\t\t\t\t);\n\t\t\t\tif (requiredFailure)\n\t\t\t\t\treturn this.result(active, \"FAILED\", `required child failed: ${requiredFailure.status.missionState}`);\n\t\t\t\tconst allRequiredTerminal = required.every((entry) => entry.status.terminal);\n\n\t\t\t\tconst allTerminal = statuses.length > 0 && statuses.every((entry) => entry.status.terminal);\n\t\t\t\tif (required.length === 0) {\n\t\t\t\t\tif (allTerminal)\n\t\t\t\t\t\treturn this.result(\n\t\t\t\t\t\t\tactive,\n\t\t\t\t\t\t\tstatuses.some((entry) => !entry.status.success) ? \"PARTIAL\" : \"SUCCEEDED\",\n\t\t\t\t\t\t\t\"all optional orchestration children reached terminal state\",\n\t\t\t\t\t\t);\n\t\t\t\t\tawait sleep(this._pollMs, signal);\n\t\t\t\t\tcontinue;\n\t\t\t\t}\n\t\t\t\tif (allRequiredTerminal && allTerminal) {\n\t\t\t\t\tconst optionalFailure = statuses.some(\n\t\t\t\t\t\t(entry) => entry.requirement === \"OPTIONAL\" && !entry.status.success,\n\t\t\t\t\t);\n\t\t\t\t\tconst unverifiedRequired = required.some(\n\t\t\t\t\t\t(entry) => entry.status.missionState === \"PARTIAL\" || entry.status.completionDecision !== \"accepted\",\n\t\t\t\t\t);\n\t\t\t\t\treturn this.result(\n\t\t\t\t\t\tactive,\n\t\t\t\t\t\toptionalFailure || unverifiedRequired ? \"PARTIAL\" : \"SUCCEEDED\",\n\t\t\t\t\t\toptionalFailure\n\t\t\t\t\t\t\t? \"required gates accepted; optional child failed\"\n\t\t\t\t\t\t\t: unverifiedRequired\n\t\t\t\t\t\t\t\t? \"orchestration completed but required child verification was not accepted\"\n\t\t\t\t\t\t\t\t: \"all required orchestration gates accepted\",\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t\tawait sleep(this._pollMs, signal);\n\t\t\t}\n\t\t} finally {\n\t\t\tthis._active.delete(handle.missionId);\n\t\t}\n\t}\n\n\tasync cancel(handle: MissionHandle, reason = \"orchestration cancelled\"): Promise<void> {\n\t\tconst active = this._active.get(handle.missionId);\n\t\tif (!active) return;\n\t\tactive.cancelled = true;\n\t\tactive.controller.abort(reason);\n\t\tconst orchestrationId = active.request.orchestrationExecution?.orchestrationId;\n\t\tif (orchestrationId) await this._orchestrator.cancel(orchestrationId).catch(() => undefined);\n\t}\n\n\tprivate async loadPlan(orchestrationId: string): Promise<OrchestrationPlan> {\n\t\tconst loaded = await this._store.load(orchestrationId);\n\t\tif (loaded.status === \"missing\") throw new Error(`ORCHESTRATION_NOT_FOUND: ${orchestrationId}`);\n\t\tif (loaded.status === \"corrupt\") throw new Error(`ORCHESTRATION_PLAN_CORRUPT: ${loaded.diagnostic}`);\n\t\treturn loaded.document.plan;\n\t}\n\n\tprivate result(\n\t\tactive: ActiveExecution,\n\t\tstate: MissionResult[\"state\"],\n\t\tsummary: string,\n\t\tcompletionDecision: MissionResult[\"completionDecision\"] = state === \"SUCCEEDED\"\n\t\t\t? \"accepted\"\n\t\t\t: state === \"PARTIAL\"\n\t\t\t\t? \"rejected\"\n\t\t\t\t: \"unavailable\",\n\t): MissionResult {\n\t\tconst finishedAtMs = this._now();\n\t\tconst executionOutcome =\n\t\t\tstate === \"SUCCEEDED\" || state === \"PARTIAL\"\n\t\t\t\t? \"COMPLETED\"\n\t\t\t\t: state === \"CANCELLED\"\n\t\t\t\t\t? \"CANCELLED\"\n\t\t\t\t\t: state === \"TIMED_OUT\"\n\t\t\t\t\t\t? \"TIMED_OUT\"\n\t\t\t\t\t\t: \"FAILED\";\n\t\tconst verification =\n\t\t\tstate === \"SUCCEEDED\" ? { status: \"verified\" as const, summary } : { status: \"unverified\" as const, summary };\n\t\tconst failures: StructuredFailure[] =\n\t\t\tstate === \"SUCCEEDED\" || state === \"PARTIAL\"\n\t\t\t\t? []\n\t\t\t\t: [\n\t\t\t\t\t\t{\n\t\t\t\t\t\t\tcategory: state === \"CANCELLED\" ? \"CANCELLED\" : state === \"TIMED_OUT\" ? \"TIMED_OUT\" : \"EXECUTION\",\n\t\t\t\t\t\t\tmessage: summary,\n\t\t\t\t\t\t},\n\t\t\t\t\t];\n\t\treturn createMissionResult({\n\t\t\tmissionId: active.request.missionId,\n\t\t\tparentMissionId: active.request.parentMissionId,\n\t\t\tdepth: active.request.depth,\n\t\t\tstate,\n\t\t\texecutionOutcome,\n\t\t\tverification,\n\t\t\tcompletionDecision,\n\t\t\tfailures,\n\t\t\texecutorDiagnostics: { executorId: this.executorId },\n\t\t\tstartedAtMs: active.startedAtMs,\n\t\t\tfinishedAtMs,\n\t\t});\n\t}\n}\n\nexport function createOrchestrationMissionExecutor(\n\toptions: OrchestrationMissionExecutorOptions,\n): OrchestrationMissionExecutor {\n\treturn new OrchestrationMissionExecutor(options);\n}\n"]}