{"version":3,"file":"worker-control-service.d.ts","sourceRoot":"","sources":["../../../src/core/worker-daemon/worker-control-service.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;GAoBG;AAEH,OAAO,KAAK,EACX,wBAAwB,EACxB,qBAAqB,EACrB,yBAAyB,EACzB,MAAM,6CAA6C,CAAC;AAIrD,OAAO,KAAK,EAAE,sBAAsB,EAAkB,oBAAoB,EAAE,MAAM,+BAA+B,CAAC;AAElH,OAAO,KAAK,EAAwB,mBAAmB,EAAE,MAAM,oCAAoC,CAAC;AAGpG,OAAO,KAAK,EAAE,sBAAsB,EAAE,MAAM,+CAA+C,CAAC;AAC5F,OAAO,EAIN,KAAK,iBAAiB,EAGtB,KAAK,oBAAoB,EACzB,KAAK,gBAAgB,EACrB,KAAK,kBAAkB,EACvB,KAAK,YAAY,EACjB,KAAK,aAAa,EAGlB,MAAM,mBAAmB,CAAC;AA4C3B,MAAM,WAAW,2BAA2B;IAC3C,qEAAqE;IACrE,UAAU,EAAE,MAAM,CAAC;IACnB,SAAS,EAAE,sBAAsB,CAAC;IAClC,WAAW,EAAE,wBAAwB,CAAC;IACtC,QAAQ,EAAE,mBAAmB,CAAC;IAC9B,2EAA2E;IAC3E,iBAAiB,EAAE,yBAAyB,CAAC;IAC7C;;;;OAIG;IACH,aAAa,CAAC,EAAE,qBAAqB,CAAC;IACtC;;;;;OAKG;IACH,QAAQ,CAAC,EAAE,sBAAsB,CAAC;IAClC,8DAA8D;IAC9D,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,8DAA8D;IAC9D,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,qEAAqE;IACrE,QAAQ,CAAC,EAAE,MAAM,CAAC;IAClB,gEAAgE;IAChE,eAAe,CAAC,EAAE,MAAM,CAAC;IACzB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,mEAAmE;IACnE,uBAAuB,CAAC,EAAE,MAAM,MAAM,CAAC;CACvC;AAMD,qBAAa,oBAAoB;IAChC,OAAO,CAAC,QAAQ,CAAC,WAAW,CAAS;IACrC,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAS;IACnC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAyB;IACpD,OAAO,CAAC,QAAQ,CAAC,YAAY,CAA2B;IACxD,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAsB;IAChD,OAAO,CAAC,QAAQ,CAAC,kBAAkB,CAA4B;IAC/D,OAAO,CAAC,QAAQ,CAAC,cAAc,CAAC,CAAwB;IACxD,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAC,CAAyB;IACpD,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAS;IACjC,OAAO,CAAC,QAAQ,CAAC,YAAY,CAAS;IACtC,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAC,CAAS;IACpC,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAC,CAAS;IAC3C,OAAO,CAAC,QAAQ,CAAC,IAAI,CAAe;IACpC,OAAO,CAAC,QAAQ,CAAC,wBAAwB,CAAe;IAExD,OAAO,CAAC,MAAM,CAAC,CAAuB;IACtC,OAAO,CAAC,SAAS,CAAC,CAAiB;IACnC,OAAO,CAAC,YAAY,CAAgC;IACpD,OAAO,CAAC,eAAe,CAAC,CAAiB;IACzC,OAAO,CAAC,UAAU,CAAC,CAAiB;IACpC,OAAO,CAAC,QAAQ,CAAS;IACzB,OAAO,CAAC,aAAa,CAAC,CAAkB;IACxC,OAAO,CAAC,SAAS,CAAC,CAA4B;IAC9C,OAAO,CAAC,UAAU,CAAC,CAAS;IAE5B,YAAY,OAAO,EAAE,2BAA2B,EAiB/C;IAED,IAAI,UAAU,IAAI,MAAM,CAEvB;IAED,IAAI,QAAQ,IAAI,MAAM,CAErB;IAED,IAAI,WAAW,IAAI,iBAAiB,CAEnC;IAED,IAAI,KAAK,IAAI,oBAAoB,GAAG,SAAS,CAE5C;IAMD;;;;;OAKG;IACG,KAAK,CAAC,OAAO,GAAE;QAAE,SAAS,CAAC,EAAE,OAAO,CAAC;QAAC,OAAO,CAAC,EAAE,OAAO,CAAA;KAAO,GAAG,OAAO,CAAC,kBAAkB,CAAC,CAyCjG;IAED;;;;OAIG;IACG,IAAI,CAAC,MAAM,CAAC,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAoCzC;IAMD;;;;;OAKG;IACG,OAAO,IAAI,OAAO,CAAC,gBAAgB,CAAC,CAezC;YAEa,YAAY;IAsD1B,gFAAgF;IAC1E,MAAM,IAAI,OAAO,CAAC,YAAY,CAAC,CAwCpC;IAMD;;;;;;OAMG;IACG,SAAS,IAAI,OAAO,CAAC,oBAAoB,CAAC,CAgE/C;YAMa,SAAS;YAmBT,gBAAgB;YAsBhB,UAAU;YASV,KAAK;YASL,oBAAoB;YAYpB,YAAY;YAMZ,kBAAkB;YAkBlB,iBAAiB;IAa/B,OAAO,CAAC,YAAY;CAapB;AAMD,wBAAsB,WAAW,CAAC,OAAO,EAAE;IAC1C,SAAS,EAAE,sBAAsB,CAAC;IAClC,WAAW,EAAE,wBAAwB,CAAC;IACtC,QAAQ,EAAE,mBAAmB,CAAC;CAC9B,GAAG,OAAO,CAAC;IAAE,OAAO,EAAE,aAAa,EAAE,CAAC;IAAC,OAAO,EAAE;QAAE,UAAU,EAAE,MAAM,CAAC;QAAC,UAAU,EAAE,MAAM,CAAA;KAAE,EAAE,CAAA;CAAE,CAAC,CAkC/F","sourcesContent":["/**\n * Worker Control Service (2.13.0).\n *\n * The durable daemon boundary that composes the existing primitives into a\n * long-lived worker loop. It NEVER re-implements executor identity, assignment\n * claim, execution fencing, mission lifecycle, evidence, or verification — it\n * consumes them:\n *\n *   - ExecutorControlService      → worker identity + incarnation + heartbeat.\n *   - AssignmentControlService    → discovery + claim + execution + completion.\n *   - DurableMissionCoordinator   → crash reconciliation + stale-owner interrupt.\n *   - ProcessMissionExecutor      → the existing local child-resume execution path.\n *   - buildAcceptanceCriteriaVerifier → deterministic completion promotion.\n *\n * Critical invariant (Jensen 3.0): logical agent concurrency != inference\n * concurrency. This worker owns an Assignment for its full logical lifetime, but\n * never reserves a Qwen inference slot for that lifetime. Inference is acquired\n * only while the child execution actually performs an inference request; the\n * future Shared Inference Scheduler will park/rehydrate the logical execution\n * (RUNNING → WAITING(INFERENCE) → RUNNING) without touching worker identity.\n */\n\nimport type {\n\tAssignmentControlService,\n\tBuildAssignedExecutor,\n\tBuildAssignedResumeLaunch,\n} from \"../assignment/assignment-control-service.js\";\nimport type { AssignmentRecord } from \"../assignment/assignment-types.js\";\nimport { AssignmentError } from \"../assignment/assignment-types.js\";\nimport { ExecutorRegistryError } from \"../executor-registry/executor-registry-types.js\";\nimport type { ExecutorControlService, ExecutorDetail, ExecutorRuntimeProof } from \"../executor-registry/index.js\";\nimport { DurableMissionCoordinator } from \"../mission-domain/durable-coordinator.js\";\nimport type { DurableMissionRecord, DurableMissionStore } from \"../mission-domain/durable-store.js\";\nimport type { MissionExecutor } from \"../mission-domain/mission-executor.js\";\nimport { isTerminalMissionState, type MissionState } from \"../mission-domain/mission-state.js\";\nimport type { ProcessMissionVerifier } from \"../mission-domain/process-mission-executor.js\";\nimport {\n\ttype WorkerActivity,\n\ttype WorkerCurrentAssignment,\n\ttype WorkerCurrentExecution,\n\ttype WorkerDaemonState,\n\ttype WorkerErrorCode,\n\ttype WorkerIdentity,\n\ttype WorkerRecoveryReport,\n\ttype WorkerRunOutcome,\n\ttype WorkerStartOutcome,\n\ttype WorkerStatus,\n\ttype WorkerSummary,\n\ttype WorkerWaitReason,\n\tworkerIdForExecutor,\n} from \"./worker-types.js\";\nimport { buildAcceptanceCriteriaVerifier } from \"./worker-verification.js\";\n\n// =============================================================================\n// Errors + helpers\n// =============================================================================\n\nfunction workerError(code: WorkerErrorCode, message: string, details: Record<string, unknown> = {}): never {\n\tconst error = new Error(message) as Error & { code: WorkerErrorCode; details: Readonly<Record<string, unknown>> };\n\terror.name = \"WorkerError\";\n\terror.code = code;\n\terror.details = Object.freeze({ ...details });\n\tthrow error;\n}\n\n/** Executor used only for read/recovery passes that never launch work. */\nconst NOOP_EXECUTOR: MissionExecutor = {\n\texecutorId: \"worker-noop\",\n\tasync launch(): Promise<never> {\n\t\tthrow new Error(\"noop executor cannot launch\");\n\t},\n\tasync awaitResult(): Promise<never> {\n\t\tthrow new Error(\"noop executor cannot await\");\n\t},\n\tasync cancel(): Promise<void> {},\n};\n\nfunction waitReasonFor(record: DurableMissionRecord): WorkerWaitReason | undefined {\n\tif (record.state !== \"WAITING\" && record.state !== \"BLOCKED\") return undefined;\n\tconst last = record.transitions[record.transitions.length - 1];\n\tconst reason = last?.reason ?? \"\";\n\tif (/inference/i.test(reason)) return \"INFERENCE\";\n\tif (/tool/i.test(reason)) return \"TOOL\";\n\tif (/external|process/i.test(reason)) return \"EXTERNAL_PROCESS\";\n\tif (/dependenc/i.test(reason)) return \"DEPENDENCY\";\n\tif (/operator/i.test(reason)) return \"OPERATOR\";\n\tif (/resource/i.test(reason)) return \"RESOURCE\";\n\treturn \"UNKNOWN\";\n}\n\n// =============================================================================\n// Options\n// =============================================================================\n\nexport interface WorkerControlServiceOptions {\n\t/** The logical executor this worker serves (workerId is derived). */\n\texecutorId: string;\n\texecutors: ExecutorControlService;\n\tassignments: AssignmentControlService;\n\tmissions: DurableMissionStore;\n\t/** Maps a durable child mission to the concrete local child CLI launch. */\n\tbuildResumeLaunch: BuildAssignedResumeLaunch;\n\t/**\n\t * Optional executor builder for a REMOTE executor. When set, the worker uses\n\t * it instead of the local `ProcessMissionExecutor` path. Built per-worker by\n\t * the CLI when the executor is bound to a remote target.\n\t */\n\tbuildExecutor?: BuildAssignedExecutor;\n\t/**\n\t * Optional verifier that promotes a clean exit-0 execution to SUCCEEDED.\n\t * When omitted, the worker builds one from the mission's declared acceptance\n\t * criteria (`buildAcceptanceCriteriaVerifier`). A mission with no verifiable\n\t * criteria therefore stays PARTIAL (unverified), never fabricated SUCCEEDED.\n\t */\n\tverifier?: ProcessMissionVerifier;\n\t/** Poll cadence for assignment discovery (default 1000ms). */\n\tpollMs?: number;\n\t/** Heartbeat cadence for worker liveness (default 3000ms). */\n\theartbeatMs?: number;\n\t/** Heartbeat expiry window (default from ExecutorControlService). */\n\texpiryMs?: number;\n\t/** Lease duration used for child execution (default 30 min). */\n\tleaseDurationMs?: number;\n\tnow?: () => number;\n\t/** Worker instance id factory (default host+UUID; never a PID). */\n\tworkerInstanceIdFactory?: () => string;\n}\n\n// =============================================================================\n// Service\n// =============================================================================\n\nexport class WorkerControlService {\n\tprivate readonly _executorId: string;\n\tprivate readonly _workerId: string;\n\tprivate readonly _executors: ExecutorControlService;\n\tprivate readonly _assignments: AssignmentControlService;\n\tprivate readonly _missions: DurableMissionStore;\n\tprivate readonly _buildResumeLaunch: BuildAssignedResumeLaunch;\n\tprivate readonly _buildExecutor?: BuildAssignedExecutor;\n\tprivate readonly _verifier?: ProcessMissionVerifier;\n\tprivate readonly _pollMs: number;\n\tprivate readonly _heartbeatMs: number;\n\tprivate readonly _expiryMs?: number;\n\tprivate readonly _leaseDurationMs?: number;\n\tprivate readonly _now: () => number;\n\tprivate readonly _workerInstanceIdFactory: () => string;\n\n\tprivate _proof?: ExecutorRuntimeProof;\n\tprivate _identity?: WorkerIdentity;\n\tprivate _daemonState: WorkerDaemonState = \"STOPPED\";\n\tprivate _heartbeatTimer?: NodeJS.Timeout;\n\tprivate _pollTimer?: NodeJS.Timeout;\n\tprivate _running = false;\n\tprivate _currentAbort?: AbortController;\n\tprivate _inFlight?: Promise<WorkerRunOutcome>;\n\tprivate _lastError?: string;\n\n\tconstructor(options: WorkerControlServiceOptions) {\n\t\tthis._executorId = options.executorId;\n\t\tthis._workerId = workerIdForExecutor(options.executorId);\n\t\tthis._executors = options.executors;\n\t\tthis._assignments = options.assignments;\n\t\tthis._missions = options.missions;\n\t\tthis._buildResumeLaunch = options.buildResumeLaunch;\n\t\tthis._buildExecutor = options.buildExecutor;\n\t\tthis._verifier = options.verifier;\n\t\tthis._pollMs = options.pollMs ?? 1000;\n\t\tthis._heartbeatMs = options.heartbeatMs ?? 3000;\n\t\tthis._expiryMs = options.expiryMs;\n\t\tthis._leaseDurationMs = options.leaseDurationMs;\n\t\tthis._now = options.now ?? (() => Date.now());\n\t\tthis._workerInstanceIdFactory =\n\t\t\toptions.workerInstanceIdFactory ??\n\t\t\t(() => `worker_${this._executorId}_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`);\n\t}\n\n\tget executorId(): string {\n\t\treturn this._executorId;\n\t}\n\n\tget workerId(): string {\n\t\treturn this._workerId;\n\t}\n\n\tget daemonState(): WorkerDaemonState {\n\t\treturn this._daemonState;\n\t}\n\n\tget proof(): ExecutorRuntimeProof | undefined {\n\t\treturn this._proof;\n\t}\n\n\t// =========================================================================\n\t// Lifecycle\n\t// =========================================================================\n\n\t/**\n\t * Register (idempotent) + activate a runtime incarnation + begin heartbeat\n\t * and polling. Exactly one worker runtime can be live per executor at a time;\n\t * a second concurrent `start()` for the same executor fails closed with\n\t * EXECUTOR_ALREADY_ACTIVE (the duplicate-daemon single-owner fence).\n\t */\n\tasync start(options: { reconcile?: boolean; polling?: boolean } = {}): Promise<WorkerStartOutcome> {\n\t\tif (this._daemonState !== \"STOPPED\") {\n\t\t\tworkerError(\"WORKER_ALREADY_STARTED\", `Worker ${this._workerId} is already ${this._daemonState}`);\n\t\t}\n\t\tthis._daemonState = \"STARTING\";\n\n\t\tlet activation: Awaited<ReturnType<ExecutorControlService[\"activateExecutor\"]>> | undefined;\n\t\ttry {\n\t\t\tactivation = await this._activate();\n\t\t} catch (error) {\n\t\t\tthis._daemonState = \"STOPPED\";\n\t\t\tif (error instanceof ExecutorRegistryError && error.code === \"EXECUTOR_ALREADY_ACTIVE\") {\n\t\t\t\tworkerError(\"EXECUTOR_ALREADY_ACTIVE\", `Executor ${this._executorId} already has a live worker runtime`, {\n\t\t\t\t\texecutorId: this._executorId,\n\t\t\t\t});\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\n\t\tif (!activation) {\n\t\t\tthis._daemonState = \"STOPPED\";\n\t\t\tworkerError(\"EXECUTOR_NOT_ACTIVE\", `Executor ${this._executorId} did not activate a runtime`);\n\t\t}\n\n\t\tif (options.reconcile !== false) {\n\t\t\ttry {\n\t\t\t\tawait this.reconcile();\n\t\t\t} catch (error) {\n\t\t\t\tthis._lastError = error instanceof Error ? error.message : String(error);\n\t\t\t}\n\t\t}\n\n\t\tthis._daemonState = \"RUNNING\";\n\t\tthis._heartbeatTimer = setInterval(() => void this._heartbeat(), this._heartbeatMs);\n\t\tif (options.polling !== false) this._pollTimer = setInterval(() => void this._poll(), this._pollMs);\n\t\treturn {\n\t\t\tidentity: this._identity!,\n\t\t\truntimeInstanceId: activation.runtimeInstanceId,\n\t\t\truntimeEpoch: activation.runtimeEpoch,\n\t\t\texpiresAtMs: activation.expiresAtMs,\n\t\t};\n\t}\n\n\t/**\n\t * Graceful shutdown: stop accepting work, stop heartbeats, abort any\n\t * in-flight execution honestly (the fenced child path persists CANCELLED,\n\t * never a fabricated success), and deactivate the runtime incarnation.\n\t */\n\tasync stop(reason?: string): Promise<void> {\n\t\tif (this._daemonState === \"STOPPED\") return;\n\t\tthis._daemonState = \"SHUTTING_DOWN\";\n\n\t\tif (this._pollTimer) clearInterval(this._pollTimer);\n\t\tif (this._heartbeatTimer) clearInterval(this._heartbeatTimer);\n\t\tthis._pollTimer = undefined;\n\t\tthis._heartbeatTimer = undefined;\n\n\t\tif (this._currentAbort) {\n\t\t\tthis._currentAbort.abort(reason ?? \"worker shutdown\");\n\t\t\tthis._currentAbort = undefined;\n\t\t}\n\n\t\t// Await the in-flight execution so its fenced terminal write (CANCELLED,\n\t\t// never a fabricated success) lands before the runtime is deactivated.\n\t\tif (this._inFlight) {\n\t\t\ttry {\n\t\t\t\tawait this._inFlight;\n\t\t\t} catch {\n\t\t\t\t// The terminal write is the authority; never mask shutdown.\n\t\t\t}\n\t\t}\n\n\t\tif (this._proof) {\n\t\t\ttry {\n\t\t\t\tawait this._executors.deactivateExecutor(this._proof);\n\t\t\t} catch (error) {\n\t\t\t\t// Deactivation is best-effort; liveness still expires via heartbeat.\n\t\t\t\tthis._lastError = error instanceof Error ? error.message : String(error);\n\t\t\t}\n\t\t}\n\n\t\tthis._proof = undefined;\n\t\tthis._identity = undefined;\n\t\tthis._daemonState = \"STOPPED\";\n\t}\n\n\t// =========================================================================\n\t// Work loop\n\t// =========================================================================\n\n\t/**\n\t * One discovery → claim → execute cycle. Deterministic ordering by\n\t * `createdAtMs` (then assignmentId). `--once` callers use this directly;\n\t * the daemon poll loop invokes it while idle. Never runs two executions\n\t * concurrently (initial concurrency policy = serial).\n\t */\n\tasync runOnce(): Promise<WorkerRunOutcome> {\n\t\tif (this._running) return { kind: \"idle\", observedAtMs: this._now() };\n\t\tif (!this._proof) {\n\t\t\tworkerError(\"WORKER_NOT_STARTED\", `Worker ${this._workerId} has not been started`);\n\t\t}\n\n\t\tthis._running = true;\n\t\tconst task = this._runOnceImpl();\n\t\tthis._inFlight = task;\n\t\ttry {\n\t\t\treturn await task;\n\t\t} finally {\n\t\t\tthis._running = false;\n\t\t\tthis._inFlight = undefined;\n\t\t}\n\t}\n\n\tprivate async _runOnceImpl(): Promise<WorkerRunOutcome> {\n\t\tconst eligible = await this._eligibleAssignments();\n\t\tif (eligible.length === 0) return { kind: \"idle\", observedAtMs: this._now() };\n\n\t\tconst assignment = eligible[0];\n\t\ttry {\n\t\t\t// Explicit claim step (ASSIGNED → ACCEPTED). An ACCEPTED assignment\n\t\t\t// (prior crash between accept and begin) is executed directly.\n\t\t\tif (assignment.state === \"ASSIGNED\") {\n\t\t\t\tawait this._assignments.acceptAssignment(assignment.assignmentId, this._proof!);\n\t\t\t}\n\n\t\t\tthis._currentAbort = new AbortController();\n\t\t\tconst verifier = this._verifier ?? (await this._verifierFor(assignment.missionId));\n\t\t\tconst started = await this._assignments.startAssignedMission(assignment.assignmentId, this._proof!, {\n\t\t\t\tbuildResumeLaunch: this._buildResumeLaunch,\n\t\t\t\tsignal: this._currentAbort.signal,\n\t\t\t\tverifier,\n\t\t\t\t...(this._buildExecutor ? { buildExecutor: this._buildExecutor } : {}),\n\t\t\t});\n\t\t\tthis._currentAbort = undefined;\n\t\t\tthis._lastError = undefined;\n\n\t\t\treturn {\n\t\t\t\tkind: \"executed\",\n\t\t\t\tassignmentId: assignment.assignmentId,\n\t\t\t\tmissionId: started.missionId,\n\t\t\t\tmissionState: started.missionState,\n\t\t\t\tattemptId: started.attemptId,\n\t\t\t\texecutionId: started.executionId,\n\t\t\t\tsuccess: started.success,\n\t\t\t\tobservedAtMs: this._now(),\n\t\t\t};\n\t\t} catch (error) {\n\t\t\tthis._currentAbort = undefined;\n\t\t\t// A claim race (another worker accepted first) is not a failure; it\n\t\t\t// is observed ownership, recorded and skipped safely.\n\t\t\tif (error instanceof AssignmentError && error.code === \"ASSIGNMENT_NOT_CURRENT\") {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"skipped\",\n\t\t\t\t\tassignmentId: assignment.assignmentId,\n\t\t\t\t\treason: error.message,\n\t\t\t\t\tobservedAtMs: this._now(),\n\t\t\t\t};\n\t\t\t}\n\t\t\tthis._lastError = error instanceof Error ? error.message : String(error);\n\t\t\tthrow error;\n\t\t}\n\t}\n\n\t// =========================================================================\n\t// Read model\n\t// =========================================================================\n\n\t/** Durable read model of this worker's identity, liveness, and current work. */\n\tasync status(): Promise<WorkerStatus> {\n\t\tlet detail: ExecutorDetail | undefined;\n\t\ttry {\n\t\t\tdetail = await this._executors.getExecutor(this._executorId);\n\t\t} catch {\n\t\t\tdetail = undefined;\n\t\t}\n\n\t\tconst identity: WorkerIdentity =\n\t\t\tthis._identity ??\n\t\t\t({\n\t\t\t\texecutorId: this._executorId,\n\t\t\t\tworkerId: this._workerId,\n\t\t\t\tworkerInstanceId: detail?.runtime?.runtimeInstanceId ?? \"\",\n\t\t\t\tworkerEpoch: detail?.runtimeEpoch ?? 0,\n\t\t\t\townerId: detail?.runtime?.ownerId ?? \"\",\n\t\t\t\thostname: detail?.runtime?.hostname,\n\t\t\t\tpid: detail?.runtime?.pid,\n\t\t\t\tstartedAtMs: detail?.runtime?.startedAtMs ?? 0,\n\t\t\t\tprocessStartedAtMs: detail?.runtime?.processStartedAtMs,\n\t\t\t\tjensenVersion: detail?.runtime?.jensenVersion,\n\t\t\t} as WorkerIdentity);\n\n\t\tconst liveness = detail?.status ?? \"OFFLINE\";\n\t\tconst currentAssignment = await this._currentAssignment();\n\t\tconst currentExecution = currentAssignment\n\t\t\t? await this._currentExecution(currentAssignment.missionId)\n\t\t\t: undefined;\n\t\tconst activity = this._activityFor(currentAssignment, currentExecution);\n\n\t\treturn {\n\t\t\tidentity,\n\t\t\tliveness,\n\t\t\tdaemonState: this._daemonState,\n\t\t\tactivity,\n\t\t\tcurrentAssignment,\n\t\t\tcurrentExecution,\n\t\t\tlastError: this._lastError,\n\t\t\tobservedAtMs: this._now(),\n\t\t};\n\t}\n\n\t// =========================================================================\n\t// Crash / restart reconciliation\n\t// =========================================================================\n\n\t/**\n\t * Conservative restart reconciliation. Runs the durable-store recovery\n\t * (expired leases → INTERRUPTED) and additionally revokes still-live leases\n\t * left behind by a prior worker incarnation of THIS executor (identified via\n\t * stale assignment `executionOwnerIdentity`). Never auto-runs side-effectful\n\t * work and never fabricates success.\n\t */\n\tasync reconcile(): Promise<WorkerRecoveryReport> {\n\t\tif (!this._proof) {\n\t\t\tworkerError(\"WORKER_NOT_STARTED\", `Worker ${this._workerId} has not been started`);\n\t\t}\n\t\tconst coordinator = new DurableMissionCoordinator(this._missions, NOOP_EXECUTOR, {\n\t\t\tnow: this._now,\n\t\t\tleaseDurationMs: this._leaseDurationMs,\n\t\t});\n\t\tconst durableRecovery = await coordinator.recover();\n\n\t\tconst staleInterrupted: WorkerRecoveryReport[\"staleInterrupted\"] = [];\n\t\tconst staleAssignmentsFailed: WorkerRecoveryReport[\"staleAssignmentsFailed\"] = [];\n\n\t\tconst { records } = await this._assignments.store.listRecords();\n\t\tconst staleExecuting = records\n\t\t\t.filter(\n\t\t\t\t(record) =>\n\t\t\t\t\trecord.executorId === this._executorId &&\n\t\t\t\t\trecord.state === \"EXECUTING\" &&\n\t\t\t\t\trecord.current &&\n\t\t\t\t\trecord.executionOwnerIdentity &&\n\t\t\t\t\trecord.executionOwnerIdentity.runtimeInstanceId !== this._proof!.runtimeInstanceId,\n\t\t\t)\n\t\t\t.sort((a, b) => a.createdAtMs - b.createdAtMs || (a.assignmentId < b.assignmentId ? -1 : 1));\n\n\t\tfor (const assignment of staleExecuting) {\n\t\t\tconst missionId = assignment.missionId;\n\t\t\tconst staleOwnerId = assignment.executionOwnerIdentity!.ownerId;\n\t\t\tconst missionLoaded = await this._missions.load(missionId);\n\n\t\t\tif (missionLoaded.status === \"ok\" && isTerminalMissionState(missionLoaded.record.state)) {\n\t\t\t\t// The child already reached a terminal result but the prior worker\n\t\t\t\t// died before completing the assignment: complete it honestly.\n\t\t\t\tconst record = missionLoaded.record;\n\t\t\t\tconst lastAttempt = record.attempts[record.attempts.length - 1];\n\t\t\t\tawait this._assignments.completeAssignment(assignment.assignmentId, {\n\t\t\t\t\tresultState: record.state,\n\t\t\t\t\tattemptId: record.resultExecutionId ? lastAttempt?.attemptId : assignment.consumedByAttemptId,\n\t\t\t\t\texecutionId: record.resultExecutionId ?? assignment.consumedByExecutionId,\n\t\t\t\t\treason: \"worker restart reconciliation: mission already terminal\",\n\t\t\t\t});\n\t\t\t\tcontinue;\n\t\t\t}\n\n\t\t\tconst interrupted = await coordinator.interruptStaleExecution(missionId, {\n\t\t\t\tstaleOwnerId,\n\t\t\t\treason: \"worker restart: prior execution owner is stale\",\n\t\t\t});\n\t\t\tif (interrupted.status === \"reconciled\") {\n\t\t\t\tstaleInterrupted.push({ missionId, previousState: interrupted.previousState ?? \"RUNNING\" });\n\t\t\t}\n\t\t\tconst failed = await this._assignments.interruptExecution(\n\t\t\t\tassignment.assignmentId,\n\t\t\t\t\"worker restart: prior execution interrupted\",\n\t\t\t);\n\t\t\tstaleAssignmentsFailed.push({ assignmentId: failed.assignmentId, missionId });\n\t\t}\n\n\t\treturn {\n\t\t\texecutorId: this._executorId,\n\t\t\tstaleInterrupted,\n\t\t\tstaleAssignmentsFailed,\n\t\t\tdurableRecovery,\n\t\t};\n\t}\n\n\t// =========================================================================\n\t// Internals\n\t// =========================================================================\n\n\tprivate async _activate(): Promise<Awaited<ReturnType<ExecutorControlService[\"activateExecutor\"]>>> {\n\t\ttry {\n\t\t\treturn await this._activateRuntime();\n\t\t} catch (error) {\n\t\t\tif (error instanceof ExecutorRegistryError && error.code === \"EXECUTOR_NOT_FOUND\") {\n\t\t\t\t// The executor is not registered yet: register a bare definition and\n\t\t\t\t// activate. If it already exists (e.g. operator-registered with\n\t\t\t\t// capabilities) it is never re-registered, so its definition is\n\t\t\t\t// preserved and the worker simply activates a runtime incarnation.\n\t\t\t\tawait this._executors.registerExecutor({\n\t\t\t\t\texecutorId: this._executorId,\n\t\t\t\t\tdisplayName: this._workerId,\n\t\t\t\t});\n\t\t\t\treturn this._activateRuntime();\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\t}\n\n\tprivate async _activateRuntime(): Promise<Awaited<ReturnType<ExecutorControlService[\"activateExecutor\"]>>> {\n\t\tconst activation = await this._executors.activateExecutor(this._executorId, {\n\t\t\truntimeInstanceId: this._workerInstanceIdFactory(),\n\t\t\texpiryMs: this._expiryMs,\n\t\t});\n\t\tthis._proof = activation.proof;\n\t\tconst runtime = activation.record.runtime;\n\t\tthis._identity = {\n\t\t\texecutorId: this._executorId,\n\t\t\tworkerId: this._workerId,\n\t\t\tworkerInstanceId: activation.runtimeInstanceId,\n\t\t\tworkerEpoch: activation.runtimeEpoch,\n\t\t\townerId: runtime?.ownerId ?? \"\",\n\t\t\thostname: runtime?.hostname,\n\t\t\tpid: runtime?.pid,\n\t\t\tstartedAtMs: runtime?.startedAtMs ?? this._now(),\n\t\t\tprocessStartedAtMs: runtime?.processStartedAtMs,\n\t\t\tjensenVersion: runtime?.jensenVersion,\n\t\t};\n\t\treturn activation;\n\t}\n\n\tprivate async _heartbeat(): Promise<void> {\n\t\tif (!this._proof || this._daemonState !== \"RUNNING\") return;\n\t\ttry {\n\t\t\tawait this._executors.heartbeatExecutor(this._proof, { expiryMs: this._expiryMs });\n\t\t} catch (error) {\n\t\t\tthis._lastError = error instanceof Error ? error.message : String(error);\n\t\t}\n\t}\n\n\tprivate async _poll(): Promise<void> {\n\t\tif (this._daemonState !== \"RUNNING\" || this._running) return;\n\t\ttry {\n\t\t\tawait this.runOnce();\n\t\t} catch (error) {\n\t\t\tthis._lastError = error instanceof Error ? error.message : String(error);\n\t\t}\n\t}\n\n\tprivate async _eligibleAssignments(): Promise<AssignmentRecord[]> {\n\t\tconst { records } = await this._assignments.store.listRecords();\n\t\treturn records\n\t\t\t.filter(\n\t\t\t\t(record) =>\n\t\t\t\t\trecord.executorId === this._executorId &&\n\t\t\t\t\trecord.current &&\n\t\t\t\t\t(record.state === \"ASSIGNED\" || record.state === \"ACCEPTED\"),\n\t\t\t)\n\t\t\t.sort((a, b) => a.createdAtMs - b.createdAtMs || (a.assignmentId < b.assignmentId ? -1 : 1));\n\t}\n\n\tprivate async _verifierFor(missionId: string): Promise<ProcessMissionVerifier | undefined> {\n\t\tconst loaded = await this._missions.load(missionId);\n\t\tif (loaded.status !== \"ok\") return undefined;\n\t\treturn buildAcceptanceCriteriaVerifier(loaded.record.request);\n\t}\n\n\tprivate async _currentAssignment(): Promise<WorkerCurrentAssignment | undefined> {\n\t\tconst { records } = await this._assignments.store.listRecords();\n\t\tconst current = records\n\t\t\t.filter((record) => record.executorId === this._executorId && record.current)\n\t\t\t.sort((a, b) => a.createdAtMs - b.createdAtMs || (a.assignmentId < b.assignmentId ? -1 : 1));\n\t\tif (current.length === 0) return undefined;\n\t\tconst record = current[0];\n\t\treturn {\n\t\t\tassignmentId: record.assignmentId,\n\t\t\tmissionId: record.missionId,\n\t\t\tstate: record.state,\n\t\t\tacceptedAtMs: record.acceptedAtMs,\n\t\t\texecutionStartedAtMs: record.executionStartedAtMs,\n\t\t\tattemptId: record.consumedByAttemptId,\n\t\t\texecutionId: record.consumedByExecutionId,\n\t\t};\n\t}\n\n\tprivate async _currentExecution(missionId: string): Promise<WorkerCurrentExecution | undefined> {\n\t\tconst loaded = await this._missions.load(missionId);\n\t\tif (loaded.status !== \"ok\") return undefined;\n\t\tconst record = loaded.record;\n\t\treturn {\n\t\t\tmissionId,\n\t\t\tmissionState: record.state,\n\t\t\tcurrentAttemptId: record.currentAttemptId,\n\t\t\tcurrentExecutionId: record.currentExecutionId,\n\t\t\twaitReason: waitReasonFor(record),\n\t\t};\n\t}\n\n\tprivate _activityFor(\n\t\tassignment: WorkerCurrentAssignment | undefined,\n\t\texecution: WorkerCurrentExecution | undefined,\n\t): WorkerActivity {\n\t\tif (!assignment) return \"IDLE\";\n\t\tif (assignment.state === \"ACCEPTED\") return \"CLAIMING\";\n\t\tif (assignment.state === \"EXECUTING\") {\n\t\t\tif (execution && execution.missionState === \"INTERRUPTED\") return \"INTERRUPTED\";\n\t\t\treturn \"EXECUTING\";\n\t\t}\n\t\tif (assignment.state === \"COMPLETED\" || assignment.state === \"FAILED\") return \"TERMINAL\";\n\t\treturn \"IDLE\";\n\t}\n}\n\n// =============================================================================\n// Worker list read model (cross-worker, durable-only)\n// =============================================================================\n\nexport async function listWorkers(options: {\n\texecutors: ExecutorControlService;\n\tassignments: AssignmentControlService;\n\tmissions: DurableMissionStore;\n}): Promise<{ entries: WorkerSummary[]; corrupt: { executorId: string; diagnostic: string }[] }> {\n\tconst result = await options.executors.listExecutors();\n\n\tconst { records: assignmentRecords } = await options.assignments.store.listRecords();\n\tconst currentByExecutor = new Map<string, AssignmentRecord>();\n\tfor (const record of assignmentRecords) {\n\t\tif (record.current) currentByExecutor.set(record.executorId, record);\n\t}\n\n\tconst entries: WorkerSummary[] = [];\n\tfor (const summary of result.entries) {\n\t\tconst current = currentByExecutor.get(summary.executorId);\n\t\tlet missionState: MissionState | undefined;\n\t\tif (current) {\n\t\t\tconst loaded = await options.missions.load(current.missionId);\n\t\t\tif (loaded.status === \"ok\") missionState = loaded.record.state;\n\t\t}\n\t\tentries.push({\n\t\t\texecutorId: summary.executorId,\n\t\t\tworkerId: workerIdForExecutor(summary.executorId),\n\t\t\tworkerInstanceId: summary.runtimeInstanceId,\n\t\t\tworkerEpoch: summary.runtimeEpoch,\n\t\t\thostname: summary.hostname,\n\t\t\tstatus: summary.status,\n\t\t\tlastHeartbeatAtMs: summary.lastHeartbeatAtMs,\n\t\t\texpiresAtMs: summary.expiresAtMs,\n\t\t\tcurrentAssignmentId: current?.assignmentId,\n\t\t\tcurrentAssignmentState: current?.state,\n\t\t\tcurrentMissionState: missionState,\n\t\t});\n\t}\n\tentries.sort((a, b) => (a.executorId < b.executorId ? -1 : a.executorId > b.executorId ? 1 : 0));\n\n\treturn { entries, corrupt: result.corrupt };\n}\n"]}