{"version":3,"file":"file-inference-queue-store.d.ts","sourceRoot":"","sources":["../../../src/core/shared-inference/file-inference-queue-store.ts"],"names":[],"mappings":"AAAA;;;;;;;;;GASG;AAOH,OAAO,EACN,KAAK,2BAA2B,EAChC,KAAK,yBAAyB,EAC9B,KAAK,2BAA2B,EAChC,KAAK,uBAAuB,EAC5B,KAAK,mBAAmB,EACxB,KAAK,uBAAuB,EAG5B,MAAM,sBAAsB,CAAC;AAU9B,MAAM,WAAW,8BAA8B;IAC9C,IAAI,EAAE,MAAM,CAAC;IACb,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,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,gBAAgB;YAIV,WAAW;YAaX,OAAO;YAQP,UAAU;YAoBV,YAAY;IA4CpB,MAAM,CAAC,UAAU,EAAE,MAAM,EAAE,MAAM,EAAE,uBAAuB,GAAG,OAAO,CAAC,2BAA2B,CAAC,CAWtG;IAEK,IAAI,CAAC,UAAU,EAAE,MAAM,GAAG,OAAO,CAAC,yBAAyB,CAAC,CAajE;IAEK,MAAM,CAAC,CAAC,EACb,UAAU,EAAE,MAAM,EAClB,QAAQ,EAAE,CAAC,OAAO,EAAE,uBAAuB,KAAK,uBAAuB,CAAC,CAAC,CAAC,GACxE,OAAO,CAAC,2BAA2B,CAAC,CAAC,CAAC,CAAC,CAoBzC;IAEK,aAAa,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC,CAevC;CACD;AAED,wBAAgB,6BAA6B,CAAC,IAAI,GAAE,MAAoC,GAAG,uBAAuB,CAEjH","sourcesContent":["/**\n * File inference queue store (3.0.0 foundation).\n *\n * Cross-process durable ledger store. Concurrency and durability follow the\n * exact pattern of FileDurableMissionStore: per-write unique temp files,\n * fsync + atomic rename, and a per-(root, resourceId) proper-lockfile critical\n * section so two OS processes can never interleave a read-validate-write\n * mutation. Corruption is surfaced structurally; a corrupt ledger never\n * fabricates a free slot or a completed inference.\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 InferenceLedgerCreateResult,\n\ttype InferenceLedgerLoadResult,\n\ttype InferenceLedgerMutateResult,\n\ttype InferenceLedgerMutation,\n\ttype InferenceQueueStore,\n\ttype InferenceResourceLedger,\n\tisSafeResourceId,\n\tparseInferenceResourceLedger,\n} from \"./inference-queue.js\";\n\nconst RECORD_SUFFIX = \".inference.json\";\nconst ATOMIC_SUFFIX = \".tmp\";\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 FileInferenceQueueStoreOptions {\n\troot: string;\n\tstoreId?: string;\n\tlockStaleMs?: number;\n\tlockRetries?: number;\n\tlockMinTimeoutMs?: number;\n\tlockMaxTimeoutMs?: number;\n}\n\nexport function defaultInferenceQueueRoot(): string {\n\tconst env = process.env.JENSEN_INFERENCE_QUEUE_DIR;\n\tif (env?.trim()) return env.trim();\n\treturn path.join(os.homedir(), \".jensen\", \"shared-inference\");\n}\n\nexport class FileInferenceQueueStore implements InferenceQueueStore {\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: FileInferenceQueueStoreOptions) {\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(resourceId: string): string {\n\t\treturn path.join(this.root, `${resourceId}${RECORD_SUFFIX}`);\n\t}\n\n\tprivate assertResourceId(resourceId: string): void {\n\t\tif (!isSafeResourceId(resourceId)) throw new Error(`Unsafe resource id for inference queue store: ${resourceId}`);\n\t}\n\n\tprivate async writeAtomic(target: string, content: string): Promise<void> {\n\t\tawait fsp.mkdir(this.root, { recursive: true });\n\t\tconst tmp = `${target}.${randomUUID()}${ATOMIC_SUFFIX}`;\n\t\tconst fh = await fsp.open(tmp, \"w\");\n\t\ttry {\n\t\t\tawait fh.writeFile(content, \"utf8\");\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(resourceId: string): Promise<string | undefined> {\n\t\ttry {\n\t\t\treturn await fsp.readFile(this.resolve(resourceId), \"utf8\");\n\t\t} catch {\n\t\t\treturn undefined;\n\t\t}\n\t}\n\n\tprivate async readParsed(\n\t\tresourceId: string,\n\t): Promise<\n\t\t| { status: \"ok\"; ledger: InferenceResourceLedger }\n\t\t| { status: \"missing\" }\n\t\t| { status: \"corrupt\"; diagnostic: string }\n\t> {\n\t\tconst raw = await this.readRaw(resourceId);\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\", diagnostic: \"ledger is not valid JSON\" };\n\t\t}\n\t\tconst result = parseInferenceResourceLedger(parsed);\n\t\tif (!result.ok) return { status: \"corrupt\", diagnostic: result.diagnostic };\n\t\treturn { status: \"ok\", ledger: result.ledger };\n\t}\n\n\tprivate async withFileLock<T>(resourceId: string, fn: () => Promise<T>): Promise<T> {\n\t\tthis.assertResourceId(resourceId);\n\t\tawait fsp.mkdir(this.root, { recursive: true });\n\t\tconst target = this.resolve(resourceId);\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});\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 Error(`Timed out acquiring inference queue lock for resource ${resourceId}`);\n\t\t\t}\n\t\t\tif (code === \"ENOTEMPTY\" || code === \"ENOTDIR\") {\n\t\t\t\tthrow new Error(`Corrupt inference queue lock metadata for resource ${resourceId}`);\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\n\t\ttry {\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// Best-effort release.\n\t\t\t\t}\n\t\t\t}\n\t\t}\n\t}\n\n\tasync create(resourceId: string, ledger: InferenceResourceLedger): Promise<InferenceLedgerCreateResult> {\n\t\tthis.assertResourceId(resourceId);\n\t\treturn this.withFileLock(resourceId, async () => {\n\t\t\tconst existing = await this.load(resourceId);\n\t\t\tif (existing.status === \"ok\") return { status: \"idempotent\", ledger: existing.ledger };\n\t\t\tif (existing.status === \"corrupt\") {\n\t\t\t\treturn { status: \"conflict\", error: `existing ledger is corrupt: ${existing.diagnostic}` };\n\t\t\t}\n\t\t\tawait this.writeAtomic(this.resolve(resourceId), JSON.stringify(ledger, null, 2));\n\t\t\treturn { status: \"created\" };\n\t\t});\n\t}\n\n\tasync load(resourceId: string): Promise<InferenceLedgerLoadResult> {\n\t\tthis.assertResourceId(resourceId);\n\t\tconst raw = await this.readRaw(resourceId);\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\", resourceId, diagnostic: \"ledger is not valid JSON\" };\n\t\t}\n\t\tconst result = parseInferenceResourceLedger(parsed);\n\t\tif (!result.ok) return { status: \"corrupt\", resourceId, diagnostic: result.diagnostic };\n\t\treturn { status: \"ok\", ledger: result.ledger };\n\t}\n\n\tasync mutate<T>(\n\t\tresourceId: string,\n\t\tmutation: (current: InferenceResourceLedger) => InferenceLedgerMutation<T>,\n\t): Promise<InferenceLedgerMutateResult<T>> {\n\t\tthis.assertResourceId(resourceId);\n\t\treturn this.withFileLock(resourceId, async () => {\n\t\t\tconst current = await this.readParsed(resourceId);\n\t\t\tif (current.status === \"missing\") return { status: \"missing\" };\n\t\t\tif (current.status === \"corrupt\") {\n\t\t\t\treturn { status: \"corrupt\", resourceId, diagnostic: current.diagnostic };\n\t\t\t}\n\t\t\tconst output = mutation(current.ledger);\n\t\t\tif (output.kind === \"noop\") return { status: \"ok\", value: output.value };\n\n\t\t\tconst nextValidation = parseInferenceResourceLedger(output.next);\n\t\t\tif (!nextValidation.ok) {\n\t\t\t\tthrow new Error(\n\t\t\t\t\t`Inference mutation produced an invalid ledger for ${resourceId}: ${nextValidation.diagnostic}`,\n\t\t\t\t);\n\t\t\t}\n\t\t\tawait this.writeAtomic(this.resolve(resourceId), JSON.stringify(output.next, null, 2));\n\t\t\treturn { status: \"ok\", value: output.value };\n\t\t});\n\t}\n\n\tasync listResources(): 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\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 (isSafeResourceId(id)) ids.push(id);\n\t\t}\n\t\treturn ids.sort();\n\t}\n}\n\nexport function createFileInferenceQueueStore(root: string = defaultInferenceQueueRoot()): FileInferenceQueueStore {\n\treturn new FileInferenceQueueStore({ root });\n}\n"]}