{"version":3,"file":"file-durable-mission-store.d.ts","sourceRoot":"","sources":["../../../src/core/mission-durable/file-durable-mission-store.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;GAmCG;AAOH,OAAO,EACN,KAAK,0BAA0B,EAC/B,KAAK,wBAAwB,EAC7B,KAAK,0BAA0B,EAC/B,KAAK,sBAAsB,EAC3B,KAAK,oBAAoB,EACzB,KAAK,yBAAyB,EAC9B,KAAK,wBAAwB,EAC7B,KAAK,mBAAmB,EAIxB,MAAM,oCAAoC,CAAC;AAkB5C,MAAM,WAAW,8BAA8B;IAC9C,IAAI,EAAE,MAAM,CAAC;IACb,2CAA2C;IAC3C,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,oDAAoD;IACpD,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,gBAAgB,CAAC,EAAE,MAAM,CAAC;CAC1B;AAED,wBAAgB,yBAAyB,IAAI,MAAM,CAIlD;AAED,qBAAa,uBAAwB,YAAW,mBAAmB;IAClE,QAAQ,CAAC,OAAO,EAAE,MAAM,CAAC;IACzB,OAAO,CAAC,QAAQ,CAAC,IAAI,CAAS;IAC9B,OAAO,CAAC,QAAQ,CAAC,WAAW,CAAS;IACrC,OAAO,CAAC,QAAQ,CAAC,WAAW,CAAS;IACrC,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAS;IAC1C,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAS;IAE1C,YAAY,OAAO,EAAE,8BAA8B,EAOlD;IAED,OAAO,CAAC,OAAO;IAIf,OAAO,CAAC,eAAe;YAMT,WAAW;YAiBX,OAAO;YAgBP,YAAY;YAqEZ,UAAU;IAoBlB,MAAM,CAAC,MAAM,EAAE,oBAAoB,GAAG,OAAO,CAAC,0BAA0B,CAAC,CAwB9E;IAEK,IAAI,CAAC,SAAS,EAAE,MAAM,GAAG,OAAO,CAAC,wBAAwB,CAAC,CAiB/D;IAEK,IAAI,CACT,MAAM,EAAE,oBAAoB,EAC5B,OAAO,GAAE,yBAA8B,GACrC,OAAO,CAAC,wBAAwB,CAAC,CAqCnC;IAEK,MAAM,CAAC,CAAC,EACb,SAAS,EAAE,MAAM,EACjB,QAAQ,EAAE,CAAC,OAAO,EAAE,oBAAoB,KAAK,sBAAsB,CAAC,CAAC,CAAC,GACpE,OAAO,CAAC,0BAA0B,CAAC,CAAC,CAAC,CAAC,CAuBxC;IAEK,YAAY,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC,CAkBtC;IAEK,uBAAuB,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC,CASjD;IAEK,YAAY,CAAC,eAAe,EAAE,MAAM,GAAG,OAAO,CAAC,MAAM,EAAE,CAAC,CAW7D;CACD;AAED;;GAEG;AACH,wBAAgB,6BAA6B,CAAC,IAAI,GAAE,MAAoC,GAAG,uBAAuB,CAEjH","sourcesContent":["/**\n * File Durable Mission Store (2.4.0).\n *\n * Concrete local implementation of the DurableMissionStore port. Records are\n * written atomically (unique temp + fsync + rename) so a process exit never\n * leaves a partially-written authoritative record. On load, records are\n * schema-validated and corruption is surfaced structurally — never silently\n * dropped and never fabricated into a default success.\n *\n * File naming is derived only from the stable `missionId` (a safe path\n * component), never from untrusted objective text.\n *\n * Concurrency model (cross-process):\n *   - Each write uses a per-write unique temp path (never shared), so two\n *     writers cannot truncate/rename one another's temp file.\n *   - Every mutation (`create`, `save`, `mutate`) is serialized by a\n *     per-(root, missionId) cross-process file lock (proper-lockfile, atomic\n *     exclusive lock-directory creation). Two independent OS processes cannot\n *     interleave a read-check-write critical section.\n *   - The optimistic `revision` compare-and-save therefore is atomic across\n *     processes: exactly one writer wins, the stale writer receives a\n *     structural `{ status: \"stale\" }`.\n *   - Execution-authoritative saves may additionally carry a `leaseProof`;\n *     the store verifies the current durable lease under the same critical\n *     section and rejects stale owners (see `ExecutionLease`).\n *\n * Durability policy:\n *   - The temp file is fsynced, then atomically renamed over the target.\n *   - The containing directory is NOT fsynced after rename, matching the\n *     strongest existing Jensen convention (MissionFileStore and MissionStore\n *     also omit directory fsync). On POSIX this means the file CONTENT is\n *     durable and the rename is atomic, but a sudden power loss could in\n *     theory lose the directory entry. This is the same guarantee the rest of\n *     the Jensen durability layer provides; it is not weakened or strengthened\n *     here.\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport * as fsp from \"node:fs/promises\";\nimport * as os from \"node:os\";\nimport * as path from \"node:path\";\nimport lockfile from \"proper-lockfile\";\nimport {\n\ttype DurableMissionCreateResult,\n\ttype DurableMissionLoadResult,\n\ttype DurableMissionMutateResult,\n\ttype DurableMissionMutation,\n\ttype DurableMissionRecord,\n\ttype DurableMissionSaveOptions,\n\ttype DurableMissionSaveResult,\n\ttype DurableMissionStore,\n\tisSafeMissionId,\n\tmissionRequestsEqual,\n\tparseDurableMissionRecord,\n} from \"../mission-domain/durable-store.js\";\nimport { ExecutionOwnershipError } from \"../mission-domain/execution-lease.js\";\nimport { isTerminalMissionState } from \"../mission-domain/mission-state.js\";\n\nconst RECORD_SUFFIX = \".mission.json\";\nconst ATOMIC_SUFFIX = \".tmp\";\n\n/**\n * Default cross-process mutation lock tuning. The lock is held only for a short\n * read-validate-write critical section (milliseconds), never across model\n * inference or execution. A crashed lock holder is recovered via proper-lockfile\n * staleness after `lockStaleMs`.\n */\nconst DEFAULT_LOCK_STALE_MS = 30_000;\nconst DEFAULT_LOCK_RETRIES = 8;\nconst DEFAULT_LOCK_MIN_TIMEOUT_MS = 20;\nconst DEFAULT_LOCK_MAX_TIMEOUT_MS = 250;\n\nexport interface FileDurableMissionStoreOptions {\n\troot: string;\n\t/** Store identity (defaults to \"file\"). */\n\tstoreId?: string;\n\t/** Cross-process mutation lock staleness window. */\n\tlockStaleMs?: number;\n\tlockRetries?: number;\n\tlockMinTimeoutMs?: number;\n\tlockMaxTimeoutMs?: number;\n}\n\nexport function defaultDurableMissionRoot(): string {\n\tconst env = process.env.JENSEN_DURABLE_MISSION_STORE;\n\tif (env?.trim()) return env.trim();\n\treturn path.join(os.homedir(), \".jensen\", \"durable-missions\");\n}\n\nexport class FileDurableMissionStore implements DurableMissionStore {\n\treadonly storeId: string;\n\tprivate readonly root: string;\n\tprivate readonly lockStaleMs: number;\n\tprivate readonly lockRetries: number;\n\tprivate readonly lockMinTimeoutMs: number;\n\tprivate readonly lockMaxTimeoutMs: number;\n\n\tconstructor(options: FileDurableMissionStoreOptions) {\n\t\tthis.root = path.resolve(options.root);\n\t\tthis.storeId = options.storeId ?? \"file\";\n\t\tthis.lockStaleMs = options.lockStaleMs ?? DEFAULT_LOCK_STALE_MS;\n\t\tthis.lockRetries = options.lockRetries ?? DEFAULT_LOCK_RETRIES;\n\t\tthis.lockMinTimeoutMs = options.lockMinTimeoutMs ?? DEFAULT_LOCK_MIN_TIMEOUT_MS;\n\t\tthis.lockMaxTimeoutMs = options.lockMaxTimeoutMs ?? DEFAULT_LOCK_MAX_TIMEOUT_MS;\n\t}\n\n\tprivate resolve(missionId: string): string {\n\t\treturn path.join(this.root, `${missionId}${RECORD_SUFFIX}`);\n\t}\n\n\tprivate assertMissionId(missionId: string): void {\n\t\tif (!isSafeMissionId(missionId)) {\n\t\t\tthrow new Error(`Unsafe mission id for durable file store: ${missionId}`);\n\t\t}\n\t}\n\n\tprivate async writeAtomic(target: string, content: string): Promise<void> {\n\t\tawait fsp.mkdir(this.root, { recursive: true });\n\t\t// Unique per-write temp path so concurrent writers can never truncate,\n\t\t// delete, or rename one another's temporary file. The suffix still ends\n\t\t// in `.tmp` so `listMissions` skips incomplete writes after a crash.\n\t\tconst tmp = `${target}.${randomUUID()}${ATOMIC_SUFFIX}`;\n\t\tawait fsp.writeFile(tmp, content, \"utf8\");\n\t\t// fsync the temp file before the atomic rename over the target.\n\t\tconst fh = await fsp.open(tmp, \"r\");\n\t\ttry {\n\t\t\tawait fh.sync();\n\t\t} finally {\n\t\t\tawait fh.close();\n\t\t}\n\t\tawait fsp.rename(tmp, target);\n\t}\n\n\tprivate async readRaw(missionId: string): Promise<string | undefined> {\n\t\ttry {\n\t\t\treturn await fsp.readFile(this.resolve(missionId), \"utf8\");\n\t\t} catch {\n\t\t\treturn undefined;\n\t\t}\n\t}\n\n\t/**\n\t * Acquire the per-mission cross-process mutation lock, run `fn`, release.\n\t *\n\t * The lock protects only the short read-validate-write critical section. It\n\t * is never held across executor/model work. Lock acquisition has bounded\n\t * retries/backoff; exhaustion is a structured LOCK_TIMEOUT, and a\n\t * compromised lock is CORRUPT_LOCK_METADATA.\n\t */\n\tprivate async withFileLock<T>(missionId: string, fn: () => Promise<T>): Promise<T> {\n\t\tthis.assertMissionId(missionId);\n\t\tawait fsp.mkdir(this.root, { recursive: true });\n\t\tconst target = this.resolve(missionId);\n\n\t\tlet compromised = false;\n\t\tlet release: (() => Promise<void>) | undefined;\n\t\ttry {\n\t\t\trelease = await lockfile.lock(target, {\n\t\t\t\trealpath: false,\n\t\t\t\tstale: this.lockStaleMs,\n\t\t\t\tretries: {\n\t\t\t\t\tretries: this.lockRetries,\n\t\t\t\t\tfactor: 2,\n\t\t\t\t\tminTimeout: this.lockMinTimeoutMs,\n\t\t\t\t\tmaxTimeout: this.lockMaxTimeoutMs,\n\t\t\t\t\trandomize: true,\n\t\t\t\t},\n\t\t\t\tonCompromised: () => {\n\t\t\t\t\tcompromised = true;\n\t\t\t\t},\n\t\t\t});\n\t\t} catch (error) {\n\t\t\tconst code =\n\t\t\t\ttypeof error === \"object\" && error !== null && \"code\" in error\n\t\t\t\t\t? (error as { code?: unknown }).code\n\t\t\t\t\t: undefined;\n\t\t\tif (code === \"ELOCKED\") {\n\t\t\t\tthrow new ExecutionOwnershipError(\n\t\t\t\t\t\"LOCK_TIMEOUT\",\n\t\t\t\t\t`Timed out acquiring mutation lock for mission ${missionId}`,\n\t\t\t\t\t{ missionId },\n\t\t\t\t);\n\t\t\t}\n\t\t\t// A non-empty lock directory is corrupt lock metadata (proper-lockfile\n\t\t\t// can only remove an empty lock dir during stale recovery). Surface it\n\t\t\t// structurally rather than fabricating ownership; the operator/test may\n\t\t\t// remove the corrupt `.lock` directory and retry.\n\t\t\tif (code === \"ENOTEMPTY\" || code === \"ENOTDIR\") {\n\t\t\t\tthrow new ExecutionOwnershipError(\n\t\t\t\t\t\"CORRUPT_LOCK_METADATA\",\n\t\t\t\t\t`Corrupt mutation lock metadata for mission ${missionId}`,\n\t\t\t\t\t{ missionId },\n\t\t\t\t);\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\n\t\ttry {\n\t\t\tif (compromised) {\n\t\t\t\tthrow new ExecutionOwnershipError(\n\t\t\t\t\t\"CORRUPT_LOCK_METADATA\",\n\t\t\t\t\t`Mutation lock for mission ${missionId} was compromised`,\n\t\t\t\t\t{ missionId },\n\t\t\t\t);\n\t\t\t}\n\t\t\treturn await fn();\n\t\t} finally {\n\t\t\tif (release) {\n\t\t\t\ttry {\n\t\t\t\t\tawait release();\n\t\t\t\t} catch {\n\t\t\t\t\t// A compromised/recovered lock may already be invalid; releasing\n\t\t\t\t\t// is best-effort and never turns a successful mutation into an error.\n\t\t\t\t}\n\t\t\t}\n\t\t}\n\t}\n\n\tprivate async readParsed(\n\t\tmissionId: string,\n\t): Promise<{ status: \"ok\"; record: DurableMissionRecord } | { status: \"missing\" } | { status: \"corrupt\" }> {\n\t\tconst raw = await this.readRaw(missionId);\n\t\tif (raw === undefined) return { status: \"missing\" };\n\t\tlet parsed: unknown;\n\t\ttry {\n\t\t\tparsed = JSON.parse(raw);\n\t\t} catch {\n\t\t\treturn { status: \"corrupt\" };\n\t\t}\n\t\tconst result = parseDurableMissionRecord(parsed);\n\t\tif (!result.ok) return { status: \"corrupt\" };\n\t\treturn { status: \"ok\", record: result.record };\n\t}\n\n\t// =========================================================================\n\t// DurableMissionStore\n\t// =========================================================================\n\n\tasync create(record: DurableMissionRecord): Promise<DurableMissionCreateResult> {\n\t\tthis.assertMissionId(record.missionId);\n\n\t\treturn this.withFileLock(record.missionId, async () => {\n\t\t\tconst existing = await this.load(record.missionId);\n\t\t\tif (existing.status === \"ok\") {\n\t\t\t\tif (missionRequestsEqual(existing.record.request, record.request)) {\n\t\t\t\t\treturn { status: \"idempotent\", record: existing.record };\n\t\t\t\t}\n\t\t\t\treturn {\n\t\t\t\t\tstatus: \"conflict\",\n\t\t\t\t\terror: \"missionId already exists with a different immutable MissionRequest\",\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (existing.status === \"corrupt\") {\n\t\t\t\treturn {\n\t\t\t\t\tstatus: \"conflict\",\n\t\t\t\t\terror: `existing durable record is corrupt: ${existing.diagnostic}`,\n\t\t\t\t};\n\t\t\t}\n\n\t\t\tawait this.writeAtomic(this.resolve(record.missionId), JSON.stringify(record, null, 2));\n\t\t\treturn { status: \"created\" };\n\t\t});\n\t}\n\n\tasync load(missionId: string): Promise<DurableMissionLoadResult> {\n\t\tthis.assertMissionId(missionId);\n\t\tconst raw = await this.readRaw(missionId);\n\t\tif (raw === undefined) return { status: \"missing\" };\n\n\t\tlet parsed: unknown;\n\t\ttry {\n\t\t\tparsed = JSON.parse(raw);\n\t\t} catch {\n\t\t\treturn { status: \"corrupt\", missionId, diagnostic: \"record is not valid JSON\" };\n\t\t}\n\n\t\tconst result = parseDurableMissionRecord(parsed);\n\t\tif (!result.ok) {\n\t\t\treturn { status: \"corrupt\", missionId, diagnostic: result.diagnostic };\n\t\t}\n\t\treturn { status: \"ok\", record: result.record };\n\t}\n\n\tasync save(\n\t\trecord: DurableMissionRecord,\n\t\toptions: DurableMissionSaveOptions = {},\n\t): Promise<DurableMissionSaveResult> {\n\t\tthis.assertMissionId(record.missionId);\n\n\t\treturn this.withFileLock(record.missionId, async () => {\n\t\t\tif (options.expectedRevision !== undefined || options.leaseProof !== undefined) {\n\t\t\t\tconst current = await this.readParsed(record.missionId);\n\t\t\t\tif (current.status !== \"ok\") {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tstatus: \"stale\",\n\t\t\t\t\t\texpectedRevision: options.expectedRevision ?? 0,\n\t\t\t\t\t\tactualRevision: undefined,\n\t\t\t\t\t};\n\t\t\t\t}\n\n\t\t\t\tif (options.expectedRevision !== undefined && current.record.revision !== options.expectedRevision) {\n\t\t\t\t\treturn {\n\t\t\t\t\t\tstatus: \"stale\",\n\t\t\t\t\t\texpectedRevision: options.expectedRevision,\n\t\t\t\t\t\tactualRevision: current.record.revision,\n\t\t\t\t\t};\n\t\t\t\t}\n\n\t\t\t\tif (options.leaseProof !== undefined) {\n\t\t\t\t\tconst lease = current.record.lease;\n\t\t\t\t\tif (!lease) return { status: \"lease_not_found\" };\n\t\t\t\t\tif (\n\t\t\t\t\t\tlease.leaseId !== options.leaseProof.leaseId ||\n\t\t\t\t\t\tlease.fencingToken !== options.leaseProof.fencingToken\n\t\t\t\t\t) {\n\t\t\t\t\t\treturn { status: \"stale_owner\", leaseId: lease.leaseId, fencingToken: lease.fencingToken };\n\t\t\t\t\t}\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tawait this.writeAtomic(this.resolve(record.missionId), JSON.stringify(record, null, 2));\n\t\t\treturn { status: \"saved\" };\n\t\t});\n\t}\n\n\tasync mutate<T>(\n\t\tmissionId: string,\n\t\tmutation: (current: DurableMissionRecord) => DurableMissionMutation<T>,\n\t): Promise<DurableMissionMutateResult<T>> {\n\t\tthis.assertMissionId(missionId);\n\n\t\treturn this.withFileLock(missionId, async () => {\n\t\t\tconst current = await this.readParsed(missionId);\n\t\t\tif (current.status === \"missing\") return { status: \"missing\" };\n\t\t\tif (current.status === \"corrupt\") {\n\t\t\t\treturn { status: \"corrupt\", missionId, diagnostic: \"record is not valid or schema-invalid\" };\n\t\t\t}\n\n\t\t\tconst output = mutation(current.record);\n\t\t\tif (output.kind === \"noop\") return { status: \"ok\", value: output.value };\n\n\t\t\t// Defense in depth: never persist a mutation result that does not\n\t\t\t// round-trip through the canonical schema validator.\n\t\t\tconst nextValidation = parseDurableMissionRecord(output.next);\n\t\t\tif (!nextValidation.ok) {\n\t\t\t\tthrow new Error(`Mutation produced an invalid record for ${missionId}: ${nextValidation.diagnostic}`);\n\t\t\t}\n\n\t\t\tawait this.writeAtomic(this.resolve(missionId), JSON.stringify(output.next, null, 2));\n\t\t\treturn { status: \"ok\", value: output.value };\n\t\t});\n\t}\n\n\tasync listMissions(): Promise<string[]> {\n\t\tlet entries: string[];\n\t\ttry {\n\t\t\tentries = await fsp.readdir(this.root);\n\t\t} catch {\n\t\t\treturn [];\n\t\t}\n\t\tconst ids: string[] = [];\n\t\tfor (const entry of entries) {\n\t\t\t// Skip incomplete temp files from a crash mid-write; the authoritative\n\t\t\t// target (if any) is what matters. Lock directories also never end in\n\t\t\t// the record suffix, so they are naturally excluded.\n\t\t\tif (entry.endsWith(ATOMIC_SUFFIX)) continue;\n\t\t\tif (!entry.endsWith(RECORD_SUFFIX)) continue;\n\t\t\tconst id = entry.slice(0, -RECORD_SUFFIX.length);\n\t\t\tif (isSafeMissionId(id)) ids.push(id);\n\t\t}\n\t\treturn ids.sort();\n\t}\n\n\tasync listNonterminalMissions(): Promise<string[]> {\n\t\tconst ids = await this.listMissions();\n\t\tconst nonterminal: string[] = [];\n\t\tfor (const id of ids) {\n\t\t\tconst loaded = await this.load(id);\n\t\t\tif (loaded.status !== \"ok\") continue; // corrupt/missing surfaced via load, not listing\n\t\t\tif (!isTerminalMissionState(loaded.record.state)) nonterminal.push(id);\n\t\t}\n\t\treturn nonterminal;\n\t}\n\n\tasync listChildren(parentMissionId: string): Promise<string[]> {\n\t\tthis.assertMissionId(parentMissionId);\n\t\tconst ids = await this.listMissions();\n\t\tconst children: string[] = [];\n\t\tfor (const id of ids) {\n\t\t\tconst loaded = await this.load(id);\n\t\t\tif (loaded.status === \"ok\" && loaded.record.parentMissionId === parentMissionId) {\n\t\t\t\tchildren.push(id);\n\t\t\t}\n\t\t}\n\t\treturn children;\n\t}\n}\n\n/**\n * Convenience factory using the default Jensen state directory.\n */\nexport function createFileDurableMissionStore(root: string = defaultDurableMissionRoot()): FileDurableMissionStore {\n\treturn new FileDurableMissionStore({ root });\n}\n"]}