{"version":3,"file":"registry.d.ts","sourceRoot":"","sources":["../../../src/core/jobs/registry.ts"],"names":[],"mappings":"AAMA,OAAO,KAAK,EACX,gBAAgB,EAChB,kBAAkB,EAClB,sBAAsB,EACtB,mBAAmB,EAEnB,cAAc,EACd,aAAa,EACb,uBAAuB,EACvB,MAAM,YAAY,CAAC;AAgBpB,MAAM,WAAW,eAAe;IAC/B,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,UAAU,EAAE,MAAM,CAAC;IACnB,IAAI,CAAC,EAAE,MAAM,EAAE,CAAC;IAChB,GAAG,CAAC,EAAE,MAAM,CAAC;IACb,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,SAAS,CAAC,EAAE,sBAAsB,CAAC;IACnC,GAAG,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,MAAM,CAAC,CAAC;IAC7B,aAAa,CAAC,EAAE,MAAM,CAAC;IACvB,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,WAAW,CAAC,EAAE,MAAM,OAAO,CAAC,OAAO,CAAC,GAAG,OAAO,CAAC;IAC/C,WAAW,CAAC,EAAE,MAAM,CAAC;CACrB;AAED,MAAM,WAAW,4BAA4B;IAC5C,UAAU,EAAE,MAAM,CAAC;IACnB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,SAAS,CAAC,EAAE,OAAO,CAAC;CACpB;AAED;;;;GAIG;AACH,qBAAa,qBAAqB;IACjC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAS;IACpC,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAS;IACjC,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAS;IAChC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAS;IACpC,OAAO,CAAC,QAAQ,CAAC,WAAW,CAAC,CAAS;IACtC,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAC,CAAS;IACrC,OAAO,CAAC,QAAQ,CAAC,GAAG,CAAe;IACnC,OAAO,CAAC,QAAQ,CAAC,SAAS,CAAU;IACpC,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAkE;IAE1F,YAAY,IAAI,EAAE,4BAA4B,EAS7C;IAED,OAAO,CAAC,UAAU;IAIlB,OAAO,CAAC,UAAU;IAIlB,OAAO,CAAC,UAAU;IAIZ,IAAI,IAAI,OAAO,CAAC,IAAI,CAAC,CAG1B;YAEa,WAAW;YAKX,OAAO;IAKf,IAAI,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC,mBAAmB,GAAG,IAAI,CAAC,CAO7D;IAEK,IAAI,IAAI,OAAO,CAAC,mBAAmB,EAAE,CAAC,CAgB3C;YAGa,QAAQ;IAyBtB;;;;OAIG;IACG,KAAK,CAAC,OAAO,EAAE,eAAe,GAAG,OAAO,CAAC,mBAAmB,CAAC,CAmElE;YAEa,MAAM;YAqBN,OAAO;IAgBrB,OAAO,CAAC,WAAW;IA2BnB,8EAA8E;IACxE,MAAM,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC,uBAAuB,GAAG,IAAI,CAAC,CA2BnE;IAED,yDAAyD;IACnD,IAAI,CAAC,GAAG,EAAE,cAAc,GAAG,OAAO,CAAC,aAAa,CAAC,CAyBtD;IAED;;;;OAIG;IACG,IAAI,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC,mBAAmB,GAAG,IAAI,CAAC,CAkC7D;IAED,8DAA8D;IACxD,OAAO,CAAC,KAAK,EAAE,MAAM,EAAE,KAAK,SAAmB,GAAG,OAAO,CAAC,mBAAmB,GAAG,IAAI,CAAC,CAmC1F;IAED;;;;OAIG;IACG,KAAK,CAAC,KAAK,EAAE,MAAM,EAAE,QAAQ,EAAE,gBAAgB,GAAG,OAAO,CAAC,mBAAmB,GAAG,IAAI,CAAC,CAiC1F;IAED,wEAAwE;IAClE,cAAc,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC,mBAAmB,GAAG,IAAI,CAAC,CAMvE;IAED;;;OAGG;IACG,kBAAkB,CACvB,YAAY,CAAC,EAAE,MAAM,EACrB,YAAY,CAAC,EAAE,SAAS,GAAG,SAAS,GAAG,QAAQ,GAAG,SAAS,GACzD,OAAO,CAAC;QACV,WAAW,EAAE,OAAO,CAAC;QACrB,cAAc,CAAC,EAAE,MAAM,CAAC;KACxB,CAAC,CAmBD;IAED,+CAA+C;IACzC,QAAQ,IAAI,OAAO,CAAC,IAAI,CAAC,CAO9B;IAEK,UAAU,IAAI,OAAO,CAAC,kBAAkB,EAAE,CAAC,CAUhD;IAEK,MAAM,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAKzC;IAED,oEAAoE;IAC9D,QAAQ,CAAC,WAAW,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAEjD;CACD","sourcesContent":["import { spawn } from \"node:child_process\";\nimport { randomUUID } from \"node:crypto\";\nimport { createWriteStream } from \"node:fs\";\nimport { mkdir, readFile, rm, writeFile } from \"node:fs/promises\";\nimport nodePath from \"node:path\";\nimport { identityMatches, readProcessIdentity } from \"./process-identity.js\";\nimport type {\n\tAdoptionEvidence,\n\tBackgroundJobEvent,\n\tBackgroundJobOwnership,\n\tBackgroundJobRecord,\n\tBackgroundJobState,\n\tJobLogsRequest,\n\tJobLogsResult,\n\tJobStatusClassification,\n} from \"./types.js\";\n\nconst DEFAULT_MAX_LOG_BYTES = 2 * 1024 * 1024; // 2 MiB per stream\nconst SECRET_ARG_RE = /(--[a-z0-9_-]*=?|-[a-z]=?)(.*)/i;\nconst SECRET_TOKENS = [\"token\", \"secret\", \"password\", \"passwd\", \"api_key\", \"apikey\", \"authorization\", \"auth\", \"key\"];\n\nfunction sanitizeArgs(args: string[]): string[] {\n\treturn args.map((a) => {\n\t\tconst m = a.match(SECRET_ARG_RE);\n\t\tif (m && SECRET_TOKENS.some((t) => m[1].toLowerCase().includes(t)) && m[2]) {\n\t\t\treturn `${m[1]}${m[1].endsWith(\"=\") ? \"\" : \" \"}<redacted>`;\n\t\t}\n\t\treturn a;\n\t});\n}\n\nexport interface StartJobOptions {\n\tjobId?: string;\n\texecutable: string;\n\targs?: string[];\n\tcwd?: string;\n\tworkspaceId?: string;\n\townerRunId?: string;\n\townership?: BackgroundJobOwnership;\n\tenv?: Record<string, string>;\n\trestartPolicy?: string;\n\tstartupTimeoutMs?: number;\n\thealthCheck?: () => Promise<boolean> | boolean;\n\tmaxLogBytes?: number;\n}\n\nexport interface BackgroundJobRegistryOptions {\n\tstorageDir: string;\n\tworkspaceId?: string;\n\townerRunId?: string;\n\tnow?: () => number;\n\tisWindows?: boolean;\n}\n\n/**\n * Durable, authoritative background-job registry. Jobs are recorded before or\n * atomically with launch, owned through process-tree/group primitives, stopped\n * only after identity verification, and never kill unrelated processes.\n */\nexport class BackgroundJobRegistry {\n\tprivate readonly storageDir: string;\n\tprivate readonly jobsDir: string;\n\tprivate readonly logDir: string;\n\tprivate readonly eventsPath: string;\n\tprivate readonly workspaceId?: string;\n\tprivate readonly ownerRunId?: string;\n\tprivate readonly now: () => number;\n\tprivate readonly isWindows: boolean;\n\tprivate readonly running = new Map<string, { kill: () => void; exited: Promise<void> }>();\n\n\tconstructor(opts: BackgroundJobRegistryOptions) {\n\t\tthis.storageDir = opts.storageDir;\n\t\tthis.workspaceId = opts.workspaceId;\n\t\tthis.ownerRunId = opts.ownerRunId;\n\t\tthis.now = opts.now ?? (() => Date.now());\n\t\tthis.isWindows = opts.isWindows ?? process.platform === \"win32\";\n\t\tthis.jobsDir = nodePath.join(this.storageDir, \"jobs\");\n\t\tthis.logDir = nodePath.join(this.storageDir, \"job-logs\");\n\t\tthis.eventsPath = nodePath.join(this.storageDir, \"job-events.jsonl\");\n\t}\n\n\tprivate recordPath(jobId: string): string {\n\t\treturn nodePath.join(this.jobsDir, `${jobId}.json`);\n\t}\n\n\tprivate stdoutPath(jobId: string): string {\n\t\treturn nodePath.join(this.logDir, `${jobId}.stdout.log`);\n\t}\n\n\tprivate stderrPath(jobId: string): string {\n\t\treturn nodePath.join(this.logDir, `${jobId}.stderr.log`);\n\t}\n\n\tasync init(): Promise<void> {\n\t\tawait mkdir(this.jobsDir, { recursive: true });\n\t\tawait mkdir(this.logDir, { recursive: true });\n\t}\n\n\tprivate async appendEvent(ev: BackgroundJobEvent): Promise<void> {\n\t\tawait mkdir(nodePath.dirname(this.eventsPath), { recursive: true });\n\t\tawait writeFile(this.eventsPath, `${JSON.stringify(ev)}\\n`, { flag: \"a\" }).catch(() => {});\n\t}\n\n\tprivate async persist(record: BackgroundJobRecord): Promise<void> {\n\t\tawait mkdir(this.jobsDir, { recursive: true });\n\t\tawait writeFile(this.recordPath(record.jobId), JSON.stringify(record, null, 2), { mode: 0o600 });\n\t}\n\n\tasync read(jobId: string): Promise<BackgroundJobRecord | null> {\n\t\ttry {\n\t\t\tconst raw = await readFile(this.recordPath(jobId), \"utf-8\");\n\t\t\treturn JSON.parse(raw) as BackgroundJobRecord;\n\t\t} catch {\n\t\t\treturn null;\n\t\t}\n\t}\n\n\tasync list(): Promise<BackgroundJobRecord[]> {\n\t\tconst { readdir } = await import(\"node:fs/promises\");\n\t\tlet files: string[] = [];\n\t\ttry {\n\t\t\tfiles = await readdir(this.jobsDir);\n\t\t} catch {\n\t\t\treturn [];\n\t\t}\n\t\tconst out: BackgroundJobRecord[] = [];\n\t\tfor (const f of files) {\n\t\t\tif (!f.endsWith(\".json\")) continue;\n\t\t\tconst rec = await this.read(f.slice(0, -5));\n\t\t\tif (rec) out.push(rec);\n\t\t}\n\t\tout.sort((a, b) => a.startedAt.localeCompare(b.startedAt));\n\t\treturn out;\n\t}\n\n\t/** Durable registration before launch. Returns the new record. */\n\tprivate async register(options: StartJobOptions): Promise<BackgroundJobRecord> {\n\t\tconst jobId = options.jobId ?? `job-${randomUUID().slice(0, 12)}`;\n\t\tconst nowIso = new Date(this.now()).toISOString();\n\t\tconst record: BackgroundJobRecord = {\n\t\t\tjobId,\n\t\t\townerRunId: options.ownerRunId ?? this.ownerRunId,\n\t\t\tworkspaceId: options.workspaceId ?? this.workspaceId,\n\t\t\townership: options.ownership,\n\t\t\tcommandIdentity: [options.executable, ...(options.args ?? [])].join(\" \"),\n\t\t\texecutable: options.executable,\n\t\t\tsanitizedArguments: sanitizeArgs(options.args ?? []),\n\t\t\tcwd: options.cwd ?? process.cwd(),\n\t\t\tprocessIdentity: \"starting\",\n\t\t\tstartedAt: nowIso,\n\t\t\tstate: \"starting\",\n\t\t\thealth: \"unknown\",\n\t\t\trestartPolicy: options.restartPolicy,\n\t\t\tlogArtifactId: jobId,\n\t\t\trestartCount: 0,\n\t\t};\n\t\tawait this.persist(record);\n\t\tawait this.appendEvent({ event: \"BACKGROUND_JOB_REGISTERED\", jobId, at: this.now() });\n\t\treturn record;\n\t}\n\n\t/**\n\t * Start a new durable background job. The record is written before launch,\n\t * then the process tree is spawned and its identity captured. Returns only\n\t * after authoritative startup status is known.\n\t */\n\tasync start(options: StartJobOptions): Promise<BackgroundJobRecord> {\n\t\tconst record = await this.register(options);\n\t\tconst stdoutPath = this.stdoutPath(record.jobId);\n\t\tconst stderrPath = this.stderrPath(record.jobId);\n\t\tconst maxLogBytes = options.maxLogBytes ?? DEFAULT_MAX_LOG_BYTES;\n\n\t\tconst proc = spawn(options.executable, options.args ?? [], {\n\t\t\tcwd: options.cwd ?? process.cwd(),\n\t\t\tenv: options.env ?? process.env,\n\t\t\tdetached: !this.isWindows,\n\t\t\twindowsHide: this.isWindows,\n\t\t\tstdio: [\"ignore\", \"pipe\", \"pipe\"],\n\t\t\tshell: false,\n\t\t});\n\n\t\tconst stdout = createWriteStream(stdoutPath, { flags: \"a\" });\n\t\tconst stderr = createWriteStream(stderrPath, { flags: \"a\" });\n\t\tlet stdoutBytes = 0;\n\t\tlet stderrBytes = 0;\n\t\tproc.stdout?.on(\"data\", (d: Buffer) => {\n\t\t\tif (stdoutBytes < maxLogBytes) {\n\t\t\t\tstdoutBytes += Math.min(d.length, maxLogBytes - stdoutBytes);\n\t\t\t\tstdout.write(d.subarray(0, maxLogBytes - (stdoutBytes - d.length)));\n\t\t\t}\n\t\t});\n\t\tproc.stderr?.on(\"data\", (d: Buffer) => {\n\t\t\tif (stderrBytes < maxLogBytes) {\n\t\t\t\tstderrBytes += Math.min(d.length, maxLogBytes - stderrBytes);\n\t\t\t\tstderr.write(d.subarray(0, maxLogBytes - (stderrBytes - d.length)));\n\t\t\t}\n\t\t});\n\n\t\tlet _settledExit: { code: number | null } | null = null;\n\t\tconst exited = new Promise<void>((resolve) => {\n\t\t\tproc.once(\"exit\", (code) => {\n\t\t\t\t_settledExit = { code };\n\t\t\t\tvoid this.onExit(record.jobId, code, stdout, stderr).then(resolve);\n\t\t\t});\n\t\t\tproc.once(\"error\", (err) => {\n\t\t\t\tvoid this.onError(record.jobId, err, stdout, stderr).then(() => resolve());\n\t\t\t});\n\t\t});\n\n\t\t// Capture authoritative identity after spawn.\n\t\tconst identity = await readProcessIdentity(proc.pid ?? 0).catch(() => null);\n\t\tconst updated: BackgroundJobRecord = {\n\t\t\t...record,\n\t\t\tprocessIdentity: String(proc.pid),\n\t\t\tprocessTreeIdentity: this.isWindows ? undefined : `pgid:${proc.pid}`,\n\t\t\tprocessStartIdentity: identity?.startIdentity,\n\t\t\tstate: \"running\",\n\t\t\thealth: \"healthy\",\n\t\t};\n\t\tawait this.persist(updated);\n\t\tawait this.appendEvent({\n\t\t\tevent: \"BACKGROUND_JOB_STARTED\",\n\t\t\tjobId: record.jobId,\n\t\t\tprocessIdentity: String(proc.pid),\n\t\t\tat: this.now(),\n\t\t});\n\n\t\tthis.running.set(record.jobId, {\n\t\t\tkill: () => this.killProcess(proc.pid ?? 0, record.jobId),\n\t\t\texited,\n\t\t});\n\n\t\treturn updated;\n\t}\n\n\tprivate async onExit(\n\t\tjobId: string,\n\t\tcode: number | null,\n\t\tstdout: NodeJS.WritableStream,\n\t\tstderr: NodeJS.WritableStream,\n\t): Promise<void> {\n\t\tstdout.end();\n\t\tstderr.end();\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return;\n\t\tconst state: BackgroundJobState = code === 0 ? \"exited\" : \"failed\";\n\t\tconst updated: BackgroundJobRecord = { ...record, state, exitCode: code ?? undefined, health: \"unknown\" };\n\t\tawait this.persist(updated);\n\t\tawait this.appendEvent(\n\t\t\tcode === 0\n\t\t\t\t? { event: \"BACKGROUND_JOB_EXITED\", jobId, exitCode: code ?? 0, at: this.now() }\n\t\t\t\t: { event: \"BACKGROUND_JOB_FAILED\", jobId, reason: `exit_code_${code}`, at: this.now() },\n\t\t);\n\t\tthis.running.delete(jobId);\n\t}\n\n\tprivate async onError(\n\t\tjobId: string,\n\t\terr: Error,\n\t\tstdout: NodeJS.WritableStream,\n\t\tstderr: NodeJS.WritableStream,\n\t): Promise<void> {\n\t\tstdout.end();\n\t\tstderr.end();\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return;\n\t\tconst updated: BackgroundJobRecord = { ...record, state: \"failed\", health: \"unknown\" };\n\t\tawait this.persist(updated);\n\t\tawait this.appendEvent({ event: \"BACKGROUND_JOB_FAILED\", jobId, reason: err.message, at: this.now() });\n\t\tthis.running.delete(jobId);\n\t}\n\n\tprivate killProcess(pid: number, _jobId: string): void {\n\t\tif (this.isWindows) {\n\t\t\tconst { execFile } = require(\"node:child_process\") as typeof import(\"node:child_process\");\n\t\t\texecFile(\"taskkill\", [\"/T\", \"/F\", \"/PID\", String(pid)], { windowsHide: true }, () => {});\n\t\t\treturn;\n\t\t}\n\t\ttry {\n\t\t\t// Kill the process group (negative pid) so owned descendants are\n\t\t\t// terminated, never unrelated processes outside the group.\n\t\t\tprocess.kill(-pid, \"SIGTERM\");\n\t\t} catch {\n\t\t\ttry {\n\t\t\t\tprocess.kill(pid, \"SIGTERM\");\n\t\t\t} catch {\n\t\t\t\t// already gone\n\t\t\t}\n\t\t}\n\t\t// Force kill fallback.\n\t\tsetTimeout(() => {\n\t\t\ttry {\n\t\t\t\tprocess.kill(-pid, \"SIGKILL\");\n\t\t\t} catch {\n\t\t\t\t/* noop */\n\t\t\t}\n\t\t}, 2000);\n\t}\n\n\t/** Classify verified live status of a job, checking real process identity. */\n\tasync status(jobId: string): Promise<JobStatusClassification | null> {\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return null;\n\t\tif (record.state === \"stopped\" || record.state === \"exited\" || record.state === \"failed\") {\n\t\t\treturn { kind: \"exited_with_code\", record, exitCode: record.exitCode ?? 0 };\n\t\t}\n\t\tconst pid = Number(record.processIdentity);\n\t\tif (!Number.isFinite(pid) || pid <= 0) return { kind: \"adoption_required\", record };\n\t\tconst live = await readProcessIdentity(pid);\n\t\tif (!live.alive) {\n\t\t\treturn { kind: \"recorded_running_but_missing\", record };\n\t\t}\n\t\tconst check = identityMatches(live, {\n\t\t\tcommandIdentity: record.commandIdentity,\n\t\t\tprocessStartIdentity: record.processStartIdentity,\n\t\t});\n\t\tif (!check.match) {\n\t\t\tconst updated: BackgroundJobRecord = { ...record, state: \"adoption_required\", health: \"degraded\" };\n\t\t\tawait this.persist(updated);\n\t\t\treturn {\n\t\t\t\tkind: \"process_identity_mismatch\",\n\t\t\t\trecord: updated,\n\t\t\t\treason: check.reason ?? \"identity_mismatch\",\n\t\t\t\tadoptionRequired: true,\n\t\t\t};\n\t\t}\n\t\treturn { kind: \"recorded_running_and_alive\", record };\n\t}\n\n\t/** Bounded log retrieval. Logs are untrusted content. */\n\tasync logs(req: JobLogsRequest): Promise<JobLogsResult> {\n\t\tconst tailLines = req.tailLines ?? 200;\n\t\tconst maxBytes = req.maxBytes ?? 64 * 1024;\n\t\tconst readTail = async (path: string): Promise<string> => {\n\t\t\ttry {\n\t\t\t\tconst content = await readFile(path, \"utf-8\");\n\t\t\t\tconst sinceSlice = req.since\n\t\t\t\t\t? content\n\t\t\t\t\t\t\t.split(\"\\n\")\n\t\t\t\t\t\t\t.filter((l) => l.length > 0)\n\t\t\t\t\t\t\t.slice(0)\n\t\t\t\t\t\t\t.join(\"\\n\")\n\t\t\t\t\t: content;\n\t\t\t\t// No timestamp filtering available; cap by bytes then lines.\n\t\t\t\tconst byBytes = sinceSlice.length > maxBytes ? sinceSlice.slice(-maxBytes) : sinceSlice;\n\t\t\t\tconst lines = byBytes.split(\"\\n\");\n\t\t\t\treturn lines.slice(-tailLines).join(\"\\n\");\n\t\t\t} catch {\n\t\t\t\treturn \"\";\n\t\t\t}\n\t\t};\n\t\tconst stdoutStream = req.stream === \"stderr\" ? \"\" : await readTail(this.stdoutPath(req.jobId));\n\t\tconst stderrStream = req.stream === \"stdout\" ? \"\" : await readTail(this.stderrPath(req.jobId));\n\t\tconst truncated = stdoutStream.length >= maxBytes || stderrStream.length >= maxBytes;\n\t\treturn { jobId: req.jobId, stdout: stdoutStream, stderr: stderrStream, truncated };\n\t}\n\n\t/**\n\t * Stop an owned job. Verifies ownership identity before terminating the\n\t * owned process tree; never kills an unrelated process that merely shares a\n\t * name.\n\t */\n\tasync stop(jobId: string): Promise<BackgroundJobRecord | null> {\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return null;\n\t\tif ([\"stopped\", \"exited\", \"failed\"].includes(record.state)) return record;\n\t\tawait this.appendEvent({ event: \"BACKGROUND_JOB_STOP_REQUESTED\", jobId, at: this.now() });\n\n\t\tconst pid = Number(record.processIdentity);\n\t\tif (Number.isFinite(pid) && pid > 0) {\n\t\t\tconst live = await readProcessIdentity(pid);\n\t\t\tconst check = identityMatches(live, {\n\t\t\t\tcommandIdentity: record.commandIdentity,\n\t\t\t\tprocessStartIdentity: record.processStartIdentity,\n\t\t\t});\n\t\t\tif (live.alive && !check.match) {\n\t\t\t\t// PID belongs to a different process now: never kill it.\n\t\t\t\tconst updated: BackgroundJobRecord = { ...record, state: \"adoption_required\", health: \"degraded\" };\n\t\t\t\tawait this.persist(updated);\n\t\t\t\treturn updated;\n\t\t\t}\n\t\t\tif (live.alive) {\n\t\t\t\tthis.killProcess(pid, jobId);\n\t\t\t\t// Wait for terminal state (bounded).\n\t\t\t\tconst running = this.running.get(jobId);\n\t\t\t\tif (running) {\n\t\t\t\t\tawait Promise.race([running.exited, sleep(5000)]);\n\t\t\t\t} else {\n\t\t\t\t\tawait sleep(1500);\n\t\t\t\t}\n\t\t\t}\n\t\t}\n\t\tconst updated: BackgroundJobRecord = { ...record, state: \"stopped\", health: \"unknown\" };\n\t\tawait this.persist(updated);\n\t\tawait this.appendEvent({ event: \"BACKGROUND_JOB_STOPPED\", jobId, at: this.now() });\n\t\treturn updated;\n\t}\n\n\t/** Restart a job: new process identity, lineage preserved. */\n\tasync restart(jobId: string, cause = \"manual_restart\"): Promise<BackgroundJobRecord | null> {\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return null;\n\t\tawait this.stop(jobId);\n\t\tconst restarts = record.restarts ?? [];\n\t\tconst newRecord = await this.start({\n\t\t\tjobId,\n\t\t\texecutable: record.executable,\n\t\t\targs: record.sanitizedArguments,\n\t\t\tcwd: record.cwd,\n\t\t\tworkspaceId: record.workspaceId,\n\t\t\townerRunId: record.ownerRunId,\n\t\t\townership: record.ownership,\n\t\t\trestartPolicy: record.restartPolicy,\n\t\t});\n\t\tnewRecord.restarts = [\n\t\t\t...restarts,\n\t\t\t{\n\t\t\t\tpreviousProcessIdentity: record.processIdentity,\n\t\t\t\tcause,\n\t\t\t\tat: new Date(this.now()).toISOString(),\n\t\t\t\tnewProcessIdentity: newRecord.processIdentity,\n\t\t\t},\n\t\t];\n\t\tnewRecord.restartCount = (record.restartCount ?? 0) + 1;\n\t\tnewRecord.state = \"running\";\n\t\tawait this.persist(newRecord);\n\t\tawait this.appendEvent({\n\t\t\tevent: \"BACKGROUND_JOB_RESTARTED\",\n\t\t\tjobId,\n\t\t\tpreviousProcessIdentity: record.processIdentity,\n\t\t\tnewProcessIdentity: newRecord.processIdentity,\n\t\t\tat: this.now(),\n\t\t});\n\t\treturn newRecord;\n\t}\n\n\t/**\n\t * Conservative adoption. Adopts a process only with strong identity\n\t * evidence (executable, command line, cwd, start time). Never adopts by PID\n\t * or name alone. Refuses without matching evidence.\n\t */\n\tasync adopt(jobId: string, evidence: AdoptionEvidence): Promise<BackgroundJobRecord | null> {\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return null;\n\t\tconst pid = Number(record.processIdentity);\n\t\tif (!Number.isFinite(pid) || pid <= 0) return null;\n\t\tconst live = await readProcessIdentity(pid);\n\t\tif (!live.alive) return null;\n\n\t\tconst cmdlineMatches = evidence.commandLine\n\t\t\t? live.commandLine.trim() === evidence.commandLine.trim() ||\n\t\t\t\t(evidence.commandLine &&\n\t\t\t\t\tlive.commandLine.includes(evidence.executable.split(/[\\\\/]/).pop() ?? evidence.executable))\n\t\t\t: evidence.executable &&\n\t\t\t\tlive.commandLine &&\n\t\t\t\tlive.commandLine.split(/\\s+/)[0]?.includes(evidence.executable.split(/[\\\\/]/).pop() ?? evidence.executable);\n\t\tif (!cmdlineMatches) return null;\n\n\t\tconst startMatch =\n\t\t\tevidence.startTimeMs === undefined ||\n\t\t\tlive.startIdentity === undefined ||\n\t\t\tMath.abs(live.startIdentity - evidence.startTimeMs / 1000) < 30;\n\t\tif (!startMatch) return null;\n\n\t\tconst updated: BackgroundJobRecord = {\n\t\t\t...record,\n\t\t\tstate: \"running\",\n\t\t\thealth: \"healthy\",\n\t\t\tprocessStartIdentity: live.startIdentity,\n\t\t\tcommandIdentity: live.commandLine || record.commandIdentity,\n\t\t};\n\t\tawait this.persist(updated);\n\t\tawait this.appendEvent({ event: \"BACKGROUND_JOB_ADOPTED\", jobId, processIdentity: String(pid), at: this.now() });\n\t\treturn updated;\n\t}\n\n\t/** Refuse adoption for a process lacking matching identity evidence. */\n\tasync refuseAdoption(jobId: string): Promise<BackgroundJobRecord | null> {\n\t\tconst record = await this.read(jobId);\n\t\tif (!record) return null;\n\t\tconst updated: BackgroundJobRecord = { ...record, state: \"adoption_required\" };\n\t\tawait this.persist(updated);\n\t\treturn updated;\n\t}\n\n\t/**\n\t * Long-horizon completion gate. A step requiring a running/healthy job\n\t * cannot be considered complete while the required job state is unresolved.\n\t */\n\tasync gateStepCompletion(\n\t\trequireJobId?: string,\n\t\trequireState?: \"running\" | \"healthy\" | \"exited\" | \"stopped\",\n\t): Promise<{\n\t\tcanComplete: boolean;\n\t\tblockingReason?: string;\n\t}> {\n\t\tif (!requireJobId) return { canComplete: true };\n\t\tconst record = await this.read(requireJobId);\n\t\tif (!record) return { canComplete: false, blockingReason: `job ${requireJobId} not found` };\n\t\tif (requireState === \"running\" || requireState === \"healthy\") {\n\t\t\tconst status = await this.status(requireJobId);\n\t\t\tif (status?.kind === \"recorded_running_and_alive\") return { canComplete: true };\n\t\t\treturn {\n\t\t\t\tcanComplete: false,\n\t\t\t\tblockingReason: `job ${requireJobId} is ${record.state} (${status?.kind ?? record.state})`,\n\t\t\t};\n\t\t}\n\t\tif (record.state !== requireState) {\n\t\t\treturn {\n\t\t\t\tcanComplete: false,\n\t\t\t\tblockingReason: `job ${requireJobId} is ${record.state}, expected ${requireState}`,\n\t\t\t};\n\t\t}\n\t\treturn { canComplete: true };\n\t}\n\n\t/** Shutdown: stop all owned jobs; no leaks. */\n\tasync shutdown(): Promise<void> {\n\t\tconst records = await this.list();\n\t\tfor (const r of records) {\n\t\t\tif ([\"stopped\", \"exited\", \"failed\"].includes(r.state)) continue;\n\t\t\tawait this.stop(r.jobId).catch(() => {});\n\t\t}\n\t\tthis.running.clear();\n\t}\n\n\tasync readEvents(): Promise<BackgroundJobEvent[]> {\n\t\ttry {\n\t\t\tconst raw = await readFile(this.eventsPath, \"utf-8\");\n\t\t\treturn raw\n\t\t\t\t.split(\"\\n\")\n\t\t\t\t.filter(Boolean)\n\t\t\t\t.map((l) => JSON.parse(l) as BackgroundJobEvent);\n\t\t} catch {\n\t\t\treturn [];\n\t\t}\n\t}\n\n\tasync remove(jobId: string): Promise<void> {\n\t\tawait this.stop(jobId).catch(() => {});\n\t\tawait rm(this.recordPath(jobId), { force: true }).catch(() => {});\n\t\tawait rm(this.stdoutPath(jobId), { force: true }).catch(() => {});\n\t\tawait rm(this.stderrPath(jobId), { force: true }).catch(() => {});\n\t}\n\n\t/** Move durable state files (survives restarts across sessions). */\n\tasync relocate(_storageDir: string): Promise<void> {\n\t\tthrow new Error(\"relocate is not implemented; durable state stays in the original storageDir\");\n\t}\n}\n\nfunction sleep(ms: number): Promise<void> {\n\treturn new Promise((r) => setTimeout(r, ms));\n}\n"]}