import { randomUUID } from "node:crypto"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { ownerFileStem, parseOwnerRecord, readOwnerRecord } from "../../shared/owner-record.ts"; import { processLiveness, processStartKey, startKeyScheme, type ProcessLiveness } from "../../shared/process-identity.ts"; import { DIRS, type SubagentState } from "../../shared/types.ts"; const MAX_CLAIMED_PREDECESSOR_SESSIONS = 8; const DEFAULT_LIVENESS_TTL_MS = 2_000; /** Initial claim plus re-claims after a dead claimer. Exhausting this fails safe (no takeover). */ const MAX_CLAIM_GENERATIONS = 4; const MAX_TAKEN_RESULTS = 256; const MAX_VERDICT_CACHE_ENTRIES = 64; const OWNER_RECORD_GC_AGE_MS = 7 * 24 * 60 * 60 * 1000; const CLAIM_GC_HARD_AGE_MS = 30 * 24 * 60 * 60 * 1000; const STALE_TEMP_AGE_MS = 10 * 60 * 1000; const GC_INTERVAL_MS = 10 * 60 * 1000; const GC_MAX_ENTRIES_EXAMINED = 256; const GC_MAX_DELETIONS = 16; const CLAIMS_DIR_NAME = "claims"; type ResultDeliveryState = Pick; /** `dead` is the only verdict that permits takeover; `unknown` covers every kind of missing or ambiguous proof. */ export type OwnerVerdict = "dead" | "live" | "unknown"; export interface ResultTakeoverDeps { /** Directory holding `.json` owner records and `claims/`. Defaults to `DIRS.owners`, read per call. */ ownersDir?: string | (() => string); hostname?: () => string; processLiveness?: (pid: number) => ProcessLiveness; processStartKey?: (pid: number) => string | undefined; now?: () => number; /** How long a per-owner verdict is reused before liveness is probed again (probing may spawn `ps`). */ livenessTtlMs?: number; } export interface ResultDeliveryOwnership { /** * True when this process owns the result. Pass `runKey` (the result's run id) so a result left by a provably dead * previous owner can be taken over; that call atomically claims it (exactly once across processes). */ owns(sessionId: string, completionOwnerId: unknown, runKey?: string): boolean; /** Side-effect-free twin of the takeover half of `owns`: would `owns(..., runKey)` succeed? */ canTakeOver(sessionId: string, completionOwnerId: unknown, runKey?: string): boolean; claimedSessionIds(): readonly string[]; claimPredecessor(previousSessionFile: string | undefined, previousRuntimeSessionId: string | null): boolean; clear(): void; } interface ClaimMarker { claimedBy: string; previousOwner: string; claimedAt: number; } type MarkerRead = { state: "missing" } | { state: "invalid" } | { state: "ok"; claimedBy: string }; function errorCode(error: unknown): string | undefined { return typeof error === "object" && error !== null && "code" in error ? (error as NodeJS.ErrnoException).code : undefined; } export function createResultDeliveryOwnership(state: ResultDeliveryState, deps: ResultTakeoverDeps = {}): ResultDeliveryOwnership { const claimed = new Map(); /** runKey -> proof that this process already holds the claim; keeps `owns` a cheap read after the first call. */ const takenOver = new Map(); const verdicts = new Map(); const now = deps.now ?? Date.now; const hostname = deps.hostname ?? os.hostname; const livenessOf = deps.processLiveness ?? ((pid: number) => processLiveness(pid)); const startKeyOf = deps.processStartKey ?? processStartKey; const livenessTtlMs = deps.livenessTtlMs ?? DEFAULT_LIVENESS_TTL_MS; let lastGcAt = Number.NEGATIVE_INFINITY; const currentOwner = (): string | undefined => typeof state.completionOwnerId === "string" && state.completionOwnerId ? state.completionOwnerId : undefined; const ownersDir = (): string => typeof deps.ownersDir === "function" ? deps.ownersDir() : deps.ownersDir ?? DIRS.owners; const claimsDir = (): string => path.join(ownersDir(), CLAIMS_DIR_NAME); /** Same session condition the in-process rule has always used; takeover never widens it. */ const sessionIsOurs = (sessionId: string, owner: string): boolean => sessionId === state.currentSessionId || claimed.get(sessionId) === owner; const ownerVerdictFromDisk = (ownerId: string): OwnerVerdict => { const record = readOwnerRecord(ownersDir(), ownerId); if (!record) return "unknown"; return recordVerdict(record.pid, record.startKey, record.hostname); }; const recordVerdict = (pid: number, recordedStartKey: string | undefined, recordedHostname: string): OwnerVerdict => { if (recordedHostname !== hostname()) return "unknown"; const liveness = livenessOf(pid); if (liveness === "dead") return "dead"; if (liveness !== "alive") return "unknown"; if (recordedStartKey) { const current = startKeyOf(pid); // A pid reused by another process has a different start time. Keys from different schemes (for example // `/proc` vs `ps`) are not comparable, so a scheme change proves nothing. if (current && startKeyScheme(current) === startKeyScheme(recordedStartKey) && current !== recordedStartKey) return "dead"; } return "live"; }; const verdictOf = (ownerId: string): OwnerVerdict => { const at = now(); const cached = verdicts.get(ownerId); if (cached && at >= cached.at && at - cached.at < livenessTtlMs) return cached.verdict; let verdict: OwnerVerdict; try { verdict = ownerVerdictFromDisk(ownerId); } catch { verdict = "unknown"; } verdicts.delete(ownerId); verdicts.set(ownerId, { verdict, at }); while (verdicts.size > MAX_VERDICT_CACHE_ENTRIES) verdicts.delete(verdicts.keys().next().value!); return verdict; }; const claimPath = (runKey: string, generation: number): string | undefined => { const stem = ownerFileStem(runKey); // `%g` cannot occur in an encodeURIComponent() stem, so generation files never collide with another key. return stem ? path.join(claimsDir(), generation <= 1 ? `${stem}.json` : `${stem}%g${generation}.json`) : undefined; }; const readMarker = (file: string): MarkerRead => { try { const parsed = JSON.parse(fs.readFileSync(file, "utf-8")) as Partial | null; return parsed && typeof parsed.claimedBy === "string" && parsed.claimedBy ? { state: "ok", claimedBy: parsed.claimedBy } : { state: "invalid" }; } catch (error) { return errorCode(error) === "ENOENT" ? { state: "missing" } : { state: "invalid" }; } }; /** * Exclusive create with complete content: write a private temp file, then hard-link it into place (link fails * with EEXIST if any claimant got there first, and a crash can never leave a half-written marker). Filesystems * without hard links fall back to `wx`. */ const createMarkerExclusively = (file: string, marker: ClaimMarker): boolean => { fs.mkdirSync(path.dirname(file), { recursive: true, mode: 0o700 }); const body = `${JSON.stringify(marker)}\n`; const temp = `${file}.${process.pid}.${randomUUID()}.tmp`; fs.writeFileSync(temp, body, { encoding: "utf-8", mode: 0o600, flag: "wx" }); try { fs.linkSync(temp, file); return true; } catch (error) { const code = errorCode(error); if (code === "EEXIST") return false; if (code !== "EPERM" && code !== "ENOSYS" && code !== "ENOTSUP" && code !== "EOPNOTSUPP" && code !== "EXDEV") throw error; try { fs.writeFileSync(file, body, { encoding: "utf-8", mode: 0o600, flag: "wx" }); return true; } catch (fallbackError) { if (errorCode(fallbackError) === "EEXIST") return false; throw fallbackError; } } finally { fs.rmSync(temp, { force: true }); } }; /** * Claim chain: generation 1 is created by the first taker; if its claimer is later proven dead, the next taker * creates generation 2 (exclusively), and so on. Nothing is ever renamed or overwritten, so two processes that * both saw a dead claimer cannot both win: exactly one create succeeds per generation. */ const claimResult = (runKey: string, previousOwner: string, owner: string, mode: "claim" | "peek"): boolean => { for (let generation = 1; generation <= MAX_CLAIM_GENERATIONS; generation += 1) { const file = claimPath(runKey, generation); if (!file) return false; let marker = readMarker(file); if (marker.state === "missing") { if (mode === "peek") return true; if (createMarkerExclusively(file, { claimedBy: owner, previousOwner, claimedAt: now() })) return true; marker = readMarker(file); } if (marker.state !== "ok") return false; if (marker.claimedBy === owner) return true; if (verdictOf(marker.claimedBy) !== "dead") return false; } return false; }; const collectGarbage = (): void => { const at = now(); if (at - lastGcAt < GC_INTERVAL_MS) return; lastGcAt = at; const owner = currentOwner(); const root = ownersDir(); let examined = 0; let deleted = 0; const old = (file: string, ageMs: number): boolean => { try { return at - fs.statSync(file).mtimeMs > ageMs; } catch { return false; } }; try { for (const name of fs.readdirSync(root)) { if (examined >= GC_MAX_ENTRIES_EXAMINED || deleted >= GC_MAX_DELETIONS) break; if (!name.endsWith(".json")) continue; examined += 1; const file = path.join(root, name); if (!old(file, OWNER_RECORD_GC_AGE_MS)) continue; const record = (() => { try { return parseOwnerRecord(JSON.parse(fs.readFileSync(file, "utf-8"))); } catch { return undefined; } })(); if (record?.completionOwnerId === owner) continue; if (record && recordVerdict(record.pid, record.startKey, record.hostname) !== "dead") continue; fs.rmSync(file, { force: true }); deleted += 1; } } catch { // Garbage collection is best effort and must never affect delivery. } try { const claims = claimsDir(); for (const name of fs.readdirSync(claims)) { if (examined >= GC_MAX_ENTRIES_EXAMINED || deleted >= GC_MAX_DELETIONS) break; examined += 1; const file = path.join(claims, name); if (name.endsWith(".tmp")) { if (old(file, STALE_TEMP_AGE_MS)) { fs.rmSync(file, { force: true }); deleted += 1; } continue; } if (!name.endsWith(".json") || !old(file, OWNER_RECORD_GC_AGE_MS)) continue; // Only the newest generation of a key may go, or an older dead claim could be re-created beneath a live one. const match = /^(.*?)(?:%g(\d+))?\.json$/.exec(name); if (!match) continue; const generation = match[2] ? Number(match[2]) : 1; if (fs.existsSync(path.join(claims, `${match[1]}%g${generation + 1}.json`))) continue; if (!old(file, CLAIM_GC_HARD_AGE_MS)) { const marker = readMarker(file); if (marker.state === "ok" && (marker.claimedBy === owner || verdictOf(marker.claimedBy) !== "dead")) continue; } fs.rmSync(file, { force: true }); deleted += 1; } } catch { // Best effort, as above. } }; const takeOver = (sessionId: string, completionOwnerId: unknown, runKey: string | undefined, mode: "claim" | "peek"): boolean => { const owner = currentOwner(); if (!owner || typeof completionOwnerId !== "string" || !completionOwnerId || completionOwnerId === owner) return false; if (typeof runKey !== "string" || !runKey || typeof sessionId !== "string" || !sessionId) return false; if (!sessionIsOurs(sessionId, owner)) return false; const taken = takenOver.get(runKey); if (taken && taken.by === owner && taken.sessionId === sessionId && taken.previousOwner === completionOwnerId) return true; try { collectGarbage(); if (verdictOf(completionOwnerId) !== "dead") return false; if (!claimResult(runKey, completionOwnerId, owner, mode)) return false; } catch { // Any filesystem surprise means we cannot prove exclusive ownership: fail safe. return false; } if (mode === "claim") { takenOver.delete(runKey); takenOver.set(runKey, { sessionId, previousOwner: completionOwnerId, by: owner }); while (takenOver.size > MAX_TAKEN_RESULTS) takenOver.delete(takenOver.keys().next().value!); } return true; }; return { owns(sessionId, completionOwnerId, runKey) { const owner = currentOwner(); if (!owner) return false; if (completionOwnerId === owner) return sessionId === state.currentSessionId || claimed.get(sessionId) === owner; return takeOver(sessionId, completionOwnerId, runKey, "claim"); }, canTakeOver(sessionId, completionOwnerId, runKey) { return takeOver(sessionId, completionOwnerId, runKey, "peek"); }, claimedSessionIds() { const owner = currentOwner(); if (!owner) return []; return [...claimed].flatMap(([sessionId, claimedOwner]) => claimedOwner === owner ? [sessionId] : []); }, claimPredecessor(previousSessionFile, previousRuntimeSessionId) { const owner = currentOwner(); if (!owner || !previousSessionFile || !previousRuntimeSessionId) return false; if (previousSessionFile !== previousRuntimeSessionId || state.currentSessionId !== previousRuntimeSessionId) return false; claimed.delete(previousRuntimeSessionId); claimed.set(previousRuntimeSessionId, owner); while (claimed.size > MAX_CLAIMED_PREDECESSOR_SESSIONS) claimed.delete(claimed.keys().next().value!); return true; }, clear() { claimed.clear(); takenOver.clear(); verdicts.clear(); }, }; }