import * as fs from "node:fs"; import * as fsp from "node:fs/promises"; import * as path from "node:path"; import { DEFAULT_CACHE, DEFAULT_PATHS } from "../../config/defaults.ts"; import { errors } from "../../errors.ts"; import type { TeamConfig } from "../../teams/team-config.ts"; import { createRunId, createTaskId } from "../../utils/ids.ts"; import { logInternalError } from "../../utils/internal-error.ts"; import { findRepoRoot, projectCrewRoot, userCrewRoot } from "../../utils/paths.ts"; import { assertSafePathId, resolveContainedRelativePath, resolveRealContainedPath } from "../../utils/safe-paths.ts"; import { toPiSessionId } from "../../utils/session-utils.ts"; import type { WorkflowConfig } from "../../workflows/workflow-config.ts"; import type { WriteDurability } from "../atomic-write.ts"; import { atomicWriteJson, atomicWriteJsonAsync, atomicWriteJsonCoalesced, flushPendingAtomicWrites, readJsonFile, } from "../atomic-write.ts"; import { canTransitionRunStatus } from "../contracts.ts"; import { withRunLock, withRunLockSync } from "../coordination/locks.ts"; import { appendEvent } from "../event-log/event-log.ts"; import type { RunModelContext, TeamRunManifest, TeamTaskState } from "../types.ts"; import { CURRENT_SCHEMA_VERSION } from "../types.ts"; import { unregisterActiveRun } from "./active-run-registry.ts"; import { extractTaskArray, loadTasksWithRecovery, loadTasksWithRecoveryAsync, quarantineCorruptFile } from "./manifest-io.ts"; export { loadManifestWithRecovery, loadTasksWithRecovery } from "./manifest-io.ts"; /** * stat() the manifest with a brief retry on Windows for the AV-scan window. * * On the GitHub Actions windows-latest runner, Windows Defender real-time * scanning can make a freshly-written manifest.json briefly invisible to * statSync (ENOENT) even though the write succeeded and the file is on disk. * loadRunManifestById is called right after createRunManifest in tests and in * production (e.g. refreshPersistedSubagentRecord), so without a retry the * caller sees a phantom "missing" run. * * On non-Windows, ENOENT means the file genuinely doesn't exist — passthrough * (throw immediately) with no retry. On Windows, ENOENT/EPERM/EBUSY/EAGAIN get * a handful of short retries (~30ms worst case) before giving up and throwing * so the caller's catch returns undefined as before. */ function statManifestWithWindowsRetry(manifestPath: string): fs.Stats { if (process.platform !== "win32") return fs.statSync(manifestPath); const retryable = new Set(["ENOENT", "EPERM", "EBUSY", "EAGAIN"]); for (let attempt = 0; attempt < 5; attempt++) { try { return fs.statSync(manifestPath); } catch (error) { const code = (error as NodeJS.ErrnoException).code; if (!retryable.has(code ?? "")) throw error; const end = Date.now() + Math.min(8, 1 * 2 ** attempt); while (Date.now() < end) { /* brief spin to ride out the AV scan window */ } } } return fs.statSync(manifestPath); // last attempt — let caller's catch handle ENOENT } export interface RunPaths { runId: string; stateRoot: string; artifactsRoot: string; manifestPath: string; tasksPath: string; eventsPath: string; } export interface ManifestCacheEntry { manifest: TeamRunManifest; tasks: TeamTaskState[]; manifestMtimeMs: number; manifestSize: number; tasksMtimeMs: number; tasksSize: number; cachedAt?: number; generation?: number; } // F1: per-stateRoot generation counter. The previous global counter was // incremented on every `invalidateRunCache` call regardless of which state's // cache was invalidated, so a write to run A caused a hot-path miss in run B's // cache even though B's on-disk files were untouched. With per-stateRoot // generations, run B only misses its own cache when its own writes happen. // NOTE: This generation counter is process-local. Cross-process consistency // (e.g., parent + child processes writing to the same run) relies on the // mtime/size checks in loadRunManifestById, not on this counter. const manifestCacheGeneration = new Map(); // R5-M1 (Round 5 MEDIUM-1): safety-net FIFO cap. stateRoots are few in // practice, but nothing else bounded this Map — it leaked ~1 entry per run // for the process lifetime (manifestCache itself has TTL + LRU eviction). const MANIFEST_CACHE_GENERATION_MAX = 64; function genOf(stateRoot: string): number { return manifestCacheGeneration.get(stateRoot) ?? 0; } const MANIFEST_CACHE_TTL_MS = 60 * 1000; // 60 seconds (FIX: increased from 15s for read-heavy workloads; render tick 160ms caused frequent misses) const LOAD_MANIFEST_RETRY_LIMIT = 5; // Configurable retry limit for mtime/size stability checks under contention const manifestCache = new Map(); /** @internal — exported for TTL-eviction unit testing (Round 19). */ export function __test__setManifestCache(stateRoot: string, entry: ManifestCacheEntry): void { setManifestCache(stateRoot, entry); } /** @internal — exported for TTL-eviction unit testing (Round 19). */ export function __test__getManifestCacheEntry(stateRoot: string): ManifestCacheEntry | undefined { return manifestCache.get(stateRoot); } /** @internal — the TTL in ms used for manifest cache eviction. */ export const MANIFEST_CACHE_TTL_MS_VALUE = MANIFEST_CACHE_TTL_MS; function setManifestCache(stateRoot: string, entry: ManifestCacheEntry): void { if (manifestCache.has(stateRoot)) manifestCache.delete(stateRoot); entry.cachedAt = Date.now(); // Stamp with the *current* per-root generation so a concurrent writer's // increment between read and write invalidates this cache hit. entry.generation = genOf(stateRoot); // FIX: Evict all stale entries by TTL before adding new entry. // This ensures entries that are never accessed still get evicted // based on TTL, not just entries that are hit. const now = Date.now(); for (const [key, val] of manifestCache.entries()) { if (val.cachedAt && now - val.cachedAt > MANIFEST_CACHE_TTL_MS) { manifestCache.delete(key); } } manifestCache.set(stateRoot, entry); while (manifestCache.size > DEFAULT_CACHE.manifestMaxEntries) { // FIX: Evict oldest entry by cachedAt (LRU), not insertion order. // cachedAt is set on both initial insertion and cache hits, so // frequently accessed entries bubble to the end and survive longer. let oldestKey: string | undefined; let oldestTime = Infinity; for (const [key, val] of manifestCache.entries()) { const t = val.cachedAt ?? 0; if (t < oldestTime) { oldestTime = t; oldestKey = key; } } if (!oldestKey) break; manifestCache.delete(oldestKey); } } function useProjectState(cwd: string): boolean { return findRepoRoot(cwd) !== undefined; } function invalidateRunCache(stateRoot: string): void { manifestCache.delete(stateRoot); // F1: only bump the generation of THIS stateRoot — sibling runs keep their // generation intact, so concurrent writes to run A no longer evict run B's // hot cache. // R5-M1: delete-then-set couples this generation entry to the cache entry it // guards (same lifecycle point) and refreshes its insertion order, so the // FIFO safety net below evicts the least-recently-invalidated stateRoot. const nextGen = genOf(stateRoot) + 1; manifestCacheGeneration.delete(stateRoot); manifestCacheGeneration.set(stateRoot, nextGen); while (manifestCacheGeneration.size > MANIFEST_CACHE_GENERATION_MAX) { const oldest = manifestCacheGeneration.keys().next().value; if (oldest === undefined) break; // Dropping a generation entry makes genOf read 0 — also drop that root's // manifest cache entry so a 0-stamped entry can never falsely match. manifestCache.delete(oldest); manifestCacheGeneration.delete(oldest); } } function scopeBaseRoot(cwd: string): string { return useProjectState(cwd) ? projectCrewRoot(cwd) : userCrewRoot(); } // P1-12: the containment verdict (resolveRealContainedPath ancestor walk, ~8-12 // syscalls) cannot change for a run dir that isn't replaced — memoize positive // results per (cwd, runId) with a short TTL. Negatives are NOT cached (so a // newly-created run is found promptly). Downstream manifest stat still catches // deletions, so a stale positive is harmless. const runStateRootCache = new Map(); const RUN_STATE_ROOT_TTL_MS = 10_000; const RUN_STATE_ROOT_CACHE_MAX = 256; function resolveRunStateRoot(cwd: string, runId: string): string | undefined { assertSafePathId("runId", runId); const key = `${cwd}\0${runId}`; const now = Date.now(); const cached = runStateRootCache.get(key); if (cached && cached.expiresAt > now) return cached.root; // F-L1 (2026-09-22): listing (run-index scopedRunRoots) UNIONS the user and // project run roots, but this resolution used scopeBaseRoot's single XOR // pick — a run living under the OTHER root (e.g. created by another cwd or // before the repo got a marker) listed fine yet FAILED every by-ID lookup: // team status / prune / forget / scheduler provenance all returned "not // found". Try the primary root first (hot path unchanged — one contained // resolve + cache hit), then fall back to the other root. The no-repo case // never falls back to a project path, mirroring scopedRunRoots' `if // (projectRoot)` guard — resolution sees exactly what listing sees. const candidates = useProjectState(cwd) ? [projectCrewRoot(cwd), userCrewRoot()] : [userCrewRoot()]; for (const root of candidates) { const runsRoot = path.join(root, DEFAULT_PATHS.state.runsSubdir); const scopedPath = resolveContainedRelativePath(runsRoot, runId, "runId"); // resolveRealContainedPath deliberately ACCEPTS a missing target (write-path // semantics: ENOENT is fine, only symlink/escape violations throw) — so // existence must be probed explicitly. Without this, the primary root // always "wins" with a phantom path and the fallback can never run. if (!fs.existsSync(path.join(scopedPath, DEFAULT_PATHS.state.manifestFile))) continue; try { resolveRealContainedPath(runsRoot, runId); } catch { continue; // symlink/escape violation under this root — try the other } if (runStateRootCache.size >= RUN_STATE_ROOT_CACHE_MAX) { const oldest = runStateRootCache.keys().next().value; if (oldest !== undefined) runStateRootCache.delete(oldest); } runStateRootCache.set(key, { root: scopedPath, expiresAt: now + RUN_STATE_ROOT_TTL_MS }); return scopedPath; } return undefined; } // PERF (2026-08-24): the artifacts containment verdict (existsSync + lstat + // resolveRealContainedPath ≈ 10-25 syscalls) cannot change for a run dir that // is not replaced — mirror the P1-12 runStateRootCache tradeoff: positive // verdicts cached 10s, negatives never cached (a newly created artifacts dir // must be found promptly). A stale positive is safe the same way P1-12 is: // downstream manifest stat catches deleted runs, and any write through // atomic-write re-runs isSymlinkSafeDirCached independently. const artifactsVerdictCache = new Map(); const ARTIFACTS_VERDICT_TTL_MS = 10_000; const ARTIFACTS_VERDICT_CACHE_MAX = 256; /** @internal — artifacts verdict cache introspection for unit tests. */ export function __test__artifactsVerdictCacheSize(): number { return artifactsVerdictCache.size; } /** @internal */ export function __test__clearArtifactsVerdictCache(): void { artifactsVerdictCache.clear(); } function validateRunManifestPaths(cwd: string, runId: string, manifest: TeamRunManifest, stateRoot: string, tasksPath: string): boolean { // Issue 2 fix: Reject manifests missing status field to prevent undefined // behavior in callers like canTransitionRunStatus(manifest.status, newStatus). if (!manifest.status || typeof manifest.status !== "string") return false; if ( manifest.runId !== runId || manifest.stateRoot !== stateRoot || manifest.tasksPath !== tasksPath || manifest.eventsPath !== path.join(stateRoot, "events.jsonl") ) return false; // F-L1 (2026-09-22): derive the artifacts parent from the RESOLVED stateRoot's // base root (…//state/runs/), not scopeBaseRoot — cross-root // runs (state under the fallback root, artifacts written beside it at // creation) must validate against their OWN root; scopeBaseRoot would // reject every fallback-resolved run and keep it invisible by ID. const baseRoot = path.resolve(stateRoot, "..", "..", ".."); const artifactsParent = path.join(baseRoot, DEFAULT_PATHS.state.artifactsSubdir); const expectedArtifactsRoot = resolveContainedRelativePath(artifactsParent, runId, "runId"); if (manifest.artifactsRoot !== expectedArtifactsRoot) return false; // PERF (2026-08-24): memoized verdict — see artifactsVerdictCache above. // Hit only after the cheap manifest-identity checks above, so a tampered // manifest (wrong paths/status) is still rejected without touching the memo. const verdictKey = `${cwd}\0${runId}`; const cachedVerdict = artifactsVerdictCache.get(verdictKey); if (cachedVerdict && cachedVerdict.expiresAt > Date.now()) return true; // Always validate artifactsRoot is not a symlink, even when manifest has // no artifacts entries. A symlinked artifactsRoot pointing outside the // artifacts parent is a security violation (could write to attacker- // controlled locations). This check catches the case where the run was // created legitimately but the artifactsRoot was later replaced with a // symlink between manifest creation and load. if (fs.existsSync(expectedArtifactsRoot)) { try { if (fs.lstatSync(expectedArtifactsRoot).isSymbolicLink()) return false; resolveRealContainedPath(artifactsParent, runId); } catch { return false; } } else if (manifest.artifacts && manifest.artifacts.length > 0) { // Has artifacts entries but directory doesn't exist - benign state for // runs still in progress. if (artifactsVerdictCache.size >= ARTIFACTS_VERDICT_CACHE_MAX) { const oldest = artifactsVerdictCache.keys().next().value; if (oldest !== undefined) artifactsVerdictCache.delete(oldest); } artifactsVerdictCache.set(verdictKey, { expiresAt: Date.now() + ARTIFACTS_VERDICT_TTL_MS }); return true; } if (artifactsVerdictCache.size >= ARTIFACTS_VERDICT_CACHE_MAX) { const oldest = artifactsVerdictCache.keys().next().value; if (oldest !== undefined) artifactsVerdictCache.delete(oldest); } artifactsVerdictCache.set(verdictKey, { expiresAt: Date.now() + ARTIFACTS_VERDICT_TTL_MS }); return true; } export function createRunPaths(cwd: string, runId = createRunId()): RunPaths { if (!cwd || typeof cwd !== "string") { throw new Error(`Invalid cwd: ${cwd}`); } assertSafePathId("runId", runId); const baseRoot = scopeBaseRoot(cwd); const stateRoot = resolveContainedRelativePath(path.join(baseRoot, DEFAULT_PATHS.state.runsSubdir), runId, "runId"); const artifactsRoot = resolveContainedRelativePath(path.join(baseRoot, DEFAULT_PATHS.state.artifactsSubdir), runId, "runId"); return { runId, stateRoot, artifactsRoot, manifestPath: path.join(stateRoot, DEFAULT_PATHS.state.manifestFile), tasksPath: path.join(stateRoot, DEFAULT_PATHS.state.tasksFile), eventsPath: path.join(stateRoot, DEFAULT_PATHS.state.eventsFile), }; } /** * Trailing phrases a step template leaves in front of `{goal}` — "Explore * the codebase for the goal: {goal}" should name the task ("Explore the * codebase"), not quote the whole run goal after every verb. The run goal is * run-level context (statusline shows it once); a task list that repeats it * per row reads like a roster of workers, not a list of work. */ const GOAL_TAILS = ["for the goal", "for goal", "the goal", "goal", "cho mục tiêu", "của mục tiêu", "mục tiêu", "for", "of", "cho", "của"]; function stripGoalTail(prefix: string): string { let name = prefix.replace(/[\s:—-]+$/, "").trim(); let changed = true; while (changed) { changed = false; const lowered = name.toLowerCase(); for (const tail of GOAL_TAILS) { if (lowered.endsWith(tail) && lowered.length > tail.length) { name = name .slice(0, name.length - tail.length) .replace(/[\s:—-]+$/, "") .trim(); changed = true; break; } } } return name; } export function createTasksFromWorkflow( runId: string, workflow: WorkflowConfig, team: TeamConfig, cwd: string, goal?: string, ): TeamTaskState[] { const stepToTaskId = new Map(workflow.steps.map((step, index) => [step.id, createTaskId(step.id, index)])); return workflow.steps.map((step, index) => { const role = team.roles.find((candidate) => candidate.name === step.role); const id = stepToTaskId.get(step.id) ?? createTaskId(step.id, index); const dependencies = step.dependsOn ?? []; const children = workflow.steps .filter((candidate) => candidate.dependsOn?.includes(step.id)) .map((candidate) => stepToTaskId.get(candidate.id)) .filter((childId): childId is string => childId !== undefined); // The plan's own words become the task's display identity: a heading // from the step body if present, else its first content line; the // remaining body is kept as `description` for detail surfaces (task // list detail line, view). Falls back to the step id when the body // carries no text of its own — `task` is optional on literal // WorkflowStep objects. // // `{goal}` in the title line is the RUN's goal, context for every // task, not this task's name. The template prefix before the // placeholder names the task; a placeholder-only body ("{goal}", // direct-run) adopts the goal itself as the title. `description` // lines substitute the placeholder (the prompt path substitutes // separately — unchanged). const bodyLines = typeof step.task === "string" && step.task.trim() ? step.task .split("\n") .map((line) => line.trim()) .filter(Boolean) : []; const heading = bodyLines.find((line) => line.startsWith("#")); const titleLine = heading ?? bodyLines[0]; const substitute = (text: string): string => (goal ? text.replaceAll("{goal}", goal) : text); let title: string; let restLines: string[]; const placeholderAt = titleLine?.indexOf("{goal}") ?? -1; if (titleLine && placeholderAt >= 0) { title = stripGoalTail(titleLine.slice(0, placeholderAt)) || goal || step.id; restLines = bodyLines.filter((line) => line !== titleLine); } else { title = (titleLine ?? step.id).replace(/^#+\s*/, ""); restLines = heading ? bodyLines.filter((line) => line !== titleLine) : bodyLines.slice(1); } const description = restLines.map(substitute).join("\n"); return { id, runId, stepId: step.id, role: step.role, agent: role?.agent ?? step.role, title: title.slice(0, 120), description: description && description !== title ? description : undefined, status: "queued", dependsOn: dependencies, cwd, model: step.model, graph: { taskId: id, parentId: dependencies[0] ? stepToTaskId.get(dependencies[0]) : undefined, children, dependencies: dependencies.map((dep) => stepToTaskId.get(dep) ?? dep), queue: dependencies.length ? "blocked" : "ready", }, }; }); } export function createRunManifest(params: { cwd: string; team: TeamConfig; workflow?: WorkflowConfig; goal: string; workspaceMode?: "single" | "worktree"; ownerSessionId?: string; runKind?: "team-run" | "goal-loop" | "dynamic-workflow"; /** round-14 P1-5: typed workflow arguments for .dwf.ts scripts (ctx.args()). */ args?: unknown; /** Deterministic-capture support (2026-09-22): pin the run ID and/or clock * so tools like docs/ui-samples/capture.ts produce byte-stable output. * Defaults preserve current behavior (generated id, wall clock). */ runId?: string; now?: () => Date; /** Model routing snapshot (Finding #1, battery 2026-09-29): detached runs * re-enter through background-runner with no ExtensionContext, so the * caller's model override / inherited session model / auth-filtered * catalogue must be persisted on the manifest at creation time. * Omitted when undefined so old manifests stay byte-identical. */ modelContext?: RunModelContext; }): { manifest: TeamRunManifest; tasks: TeamTaskState[]; paths: RunPaths } { const paths = createRunPaths(params.cwd, params.runId); const now = (params.now ? params.now() : new Date()).toISOString(); const tasks = params.workflow ? createTasksFromWorkflow(paths.runId, params.workflow, params.team, params.cwd, params.goal) : []; const manifest: TeamRunManifest = { schemaVersion: CURRENT_SCHEMA_VERSION, runId: paths.runId, sessionId: toPiSessionId(paths.runId), team: params.team.name, workflow: params.workflow?.name, goal: params.goal, status: "queued", workspaceMode: params.workspaceMode ?? params.team.workspaceMode ?? "single", createdAt: now, updatedAt: now, cwd: params.cwd, stateRoot: paths.stateRoot, artifactsRoot: paths.artifactsRoot, tasksPath: paths.tasksPath, eventsPath: paths.eventsPath, artifacts: [], ...(params.ownerSessionId ? { ownerSessionId: params.ownerSessionId } : {}), runKind: params.runKind ?? "team-run", ...(params.args !== undefined ? { args: params.args } : {}), ...(params.modelContext ? { modelContext: params.modelContext } : {}), }; fs.mkdirSync(paths.stateRoot, { recursive: true }); fs.mkdirSync(paths.artifactsRoot, { recursive: true }); // FIX: Use saveManifestAndTasksAtomicSync and check result to detect // partial failure. If tasksWritten=false when manifestWritten=true, // throw to ensure manifest and tasks are always consistent. const result = saveManifestAndTasksAtomicSync(manifest, tasks); if (!result.manifestWritten || !result.tasksWritten) { // Surface the underlying error message (result.error is String(err) from // saveManifestAndTasksAtomicSync). Passing it through errors.fileWrite as a // fake ErrnoException loses the message (reads .code → undefined → // "unknown"). Include it explicitly in the thrown message so CI logs and // production callers can see WHY the write failed instead of ": unknown". const cause = result.error ? `: ${result.error}` : ""; throw errors .fileWrite(paths.stateRoot, { code: "EWRITEFAIL", } as NodeJS.ErrnoException) .withContext( `saveManifestAndTasksAtomicSync: manifestWritten=${result.manifestWritten}, tasksWritten=${result.tasksWritten}${cause}`, ); } appendEvent(paths.eventsPath, { type: "run.created", runId: paths.runId, data: { team: params.team.name, workflow: params.workflow?.name }, metadata: { seq: 1, provenance: "team_runner", sessionIdentity: { title: params.team.name, workspace: params.cwd, purpose: params.goal, }, ownership: { owner: params.team.name, workflowScope: params.workflow?.name ?? "manual", watcherAction: "act", }, confidence: "high", }, }); invalidateRunCache(paths.stateRoot); return { manifest, tasks, paths }; } export function saveRunManifest(manifest: TeamRunManifest, options?: { allowTerminalExit?: boolean }): TeamRunManifest { // Finding 8 (2026-09-23): terminal-preserve at the WRITE layer. Mid-flight // savers (task-runner artifact/progress writes carrying a stale in-memory // "running" manifest) previously overwrote an externally-written terminal // status; merge/finalize had their own guards but every other save site did // not. Guard here so NO raw save can erase a terminal disk status. // FIX: Capture the cached tasks array + mtime/size BEFORE we invalidate the // cache. The previous implementation re-read tasks.json from disk after the // manifest write (a JSON.parse + fs.readFileSync per call), which defeated // the whole point of the manifest cache on the hot update path. Reusing the // already-cached entry is safe because: // 1. The cache entry's tasks array matches the tasks array reflected by // its tasksMtimeMs/tasksSize — the three values move together. // 2. After invalidate, saveRunTasks (or any other writer) would bump the // per-stateRoot generation, so a stale cache hit cannot serve. // 3. On cache miss, we fall back to tasks: [] with mtime/size 0 — the same // shape the original pre-FIX code produced — so behavior is unchanged. const cachedBeforeInvalidate = manifestCache.get(manifest.stateRoot); const cachedTasks = cachedBeforeInvalidate?.tasks ?? []; const cachedTasksMtimeMs = cachedBeforeInvalidate?.tasksMtimeMs ?? 0; const cachedTasksSize = cachedBeforeInvalidate?.tasksSize ?? 0; // FIX: Invalidate cache BEFORE atomic write. The order matters for crash // safety: if we invalidated after the write and crashed before invalidation, // the stale cache entry (up to MANIFEST_CACHE_TTL_MS old) could be served. // By invalidating first, the worst case is a cache miss forcing a disk read, // which is always safe. invalidateRunCache(manifest.stateRoot); const manifestPath = path.join(manifest.stateRoot, "manifest.json"); const effective = preserveDiskTerminalStatus(manifest, options?.allowTerminalExit); // REVIEW FIX (2026-09-10): reverted WI-2.2's coalesced conversion — // saveRunManifest is a SYNCHRONOUS persist by name/contract (tests assert // it, broker loadRunManifestById is a cross-process reader, and the // statSync-based cache repopulation below needs the real post-write // mtime/size). The 50ms coalesce window broke all three. atomicWriteJson(manifestPath, effective); // FIX: Re-populate cache with actual mtime/size so loadRunManifestById // doesn't miss the cache on next read. Without this, every load until // TTL expires would hit disk because cached 0 !== any real mtime. // NOTE: tasks is reused from the pre-invalidate cache snapshot above; if no // cache existed, tasks is [] (matching pre-FIX behavior). Callers that need // fresh tasks should call saveRunTasks or loadRunTasks separately. const manifestStat = fs.statSync(manifestPath); setManifestCache(manifest.stateRoot, { manifest: effective, tasks: cachedTasks, manifestMtimeMs: manifestStat.mtimeMs, manifestSize: manifestStat.size, tasksMtimeMs: cachedTasksMtimeMs, tasksSize: cachedTasksSize, }); return effective; } export async function saveRunManifestAsync(manifest: TeamRunManifest, options?: { allowTerminalExit?: boolean }): Promise { // FIX: Capture cached tasks array + mtime/size BEFORE invalidating, same // rationale as the sync saveRunManifest above. The async path previously // always set tasks: [] with mtime/size 0, so any cache hit was guaranteed // to look stale to loadRunManifestById until something else wrote tasks.json. const cachedBeforeInvalidate = manifestCache.get(manifest.stateRoot); const cachedTasks = cachedBeforeInvalidate?.tasks ?? []; const cachedTasksMtimeMs = cachedBeforeInvalidate?.tasksMtimeMs ?? 0; const cachedTasksSize = cachedBeforeInvalidate?.tasksSize ?? 0; // FIX: Invalidate cache BEFORE atomic write to prevent stale cache serving // after a crash. See saveRunManifest for full explanation. invalidateRunCache(manifest.stateRoot); const manifestPath = path.join(manifest.stateRoot, "manifest.json"); const effective = preserveDiskTerminalStatus(manifest, options?.allowTerminalExit); await atomicWriteJsonAsync(manifestPath, effective); // FIX: Re-populate cache with actual mtime/size. See saveRunManifest. // RACE GUARD: another concurrent async save (OPT-02) may unlink+rewrite // manifest.json between our atomicWriteJsonAsync and this stat. If stat // hits ENOENT, use fallback values — the cache will look stale on the // next load and re-read from disk, which is correct. let manifestStat: { mtimeMs: number; size: number }; try { manifestStat = await fs.promises.stat(manifestPath); } catch (statError) { const code = String((statError as NodeJS.ErrnoException).code ?? ""); if (code !== "ENOENT") throw statError; manifestStat = { mtimeMs: 0, size: 0 }; } setManifestCache(manifest.stateRoot, { manifest: effective, tasks: cachedTasks, manifestMtimeMs: manifestStat.mtimeMs, manifestSize: manifestStat.size, tasksMtimeMs: cachedTasksMtimeMs, tasksSize: cachedTasksSize, }); return effective; } /** * ST-4: Defense-in-depth guard — refuse to persist an empty tasks array over * a previously non-empty tasks file. Prevents the cascade where a corrupt * tasks.json (→ [] on load via readJsonFile catching SyntaxError) is then * persisted as [], permanently destroying the data. If the on-disk file has * tasks and the incoming array is empty, we log a warning and refuse. * * @returns true if the write should proceed, false if it was refused. */ function shouldPersistTasks(manifest: TeamRunManifest, tasks: TeamTaskState[]): boolean { if (tasks.length > 0) return true; const existing = extractTaskArray(readJsonFile(manifest.tasksPath)); if (existing.length > 0) { logInternalError( "state-store", new Error( `refusing to persist empty tasks over ${existing.length} existing task(s) — possible corrupt-load cascade (runId=${manifest.runId})`, ), undefined, "warn", ); return false; } return true; } export function saveRunTasks(manifest: TeamRunManifest, tasks: TeamTaskState[]): void { // ST-4: refuse to persist [] over a previously-non-empty tasks file. // Prevents the cascade where a corrupt tasks.json (→ [] on load) is then // persisted as [], permanently destroying the data. if (!shouldPersistTasks(manifest, tasks)) return; // FIX: Invalidate cache BEFORE atomic write to prevent stale cache serving. invalidateRunCache(manifest.stateRoot); // Guard: if the run state directory has been removed (prune/forget/cleanup), // silently return — there is nothing to persist and the run is gone. try { fs.statSync(manifest.stateRoot); } catch { return; } atomicWriteJson(manifest.tasksPath, tasks, { compact: true }); // FIX: Re-populate cache with actual mtime/size for manifest and tasks. // Note: We re-read manifest from disk to get its current mtime/size // since we only wrote tasks here. const manifestPath = path.join(manifest.stateRoot, "manifest.json"); let manifestStat: fs.Stats; let tasksStat: fs.Stats; try { manifestStat = fs.statSync(manifestPath); tasksStat = fs.statSync(manifest.tasksPath); } catch { // Run state disappeared between the guard above and now — give up gracefully. return; } // FIX: If cache was evicted, re-read manifest from disk rather than using // a minimal fallback. A stale minimal manifest with only runId populated // would cause manifest.status to be undefined, breaking status checks. // If the manifest cannot be read, throw an error — callers using withRunLock // should never hit this case, and the error indicates a serious problem. // FIX: Also check that the re-read manifest has a status field. If not, // fall back to the manifest parameter's status rather than serving a // degraded manifest that would break status transition checks. const cached = manifestCache.get(manifest.stateRoot); const manifestEntry = cached?.manifest ?? readJsonFile(manifestPath); if (!manifestEntry) { return; // Run deleted between guard and read } // Preserve current status from the manifest parameter if the on-disk // manifest is missing it (could be a partial write). if (!manifestEntry.status) { manifestEntry.status = manifest.status; } setManifestCache(manifest.stateRoot, { manifest: manifestEntry, tasks, manifestMtimeMs: manifestStat.mtimeMs, manifestSize: manifestStat.size, tasksMtimeMs: tasksStat.mtimeMs, tasksSize: tasksStat.size, }); } /** * 2.1 caller-migration helper: coalesced variant. Use only when the * caller does NOT immediately read tasks.json afterwards (the read would * see the previous on-disk content while the write is still buffered). * Bulk update paths that fan out into multiple writer call sites are the * intended use case. Single-update + read-update loops (e.g. * persistSingleTaskUpdate) should keep using saveRunTasks — OR call * `flushPendingAtomicWrites()` immediately before their read to force * any pending coalesced writes to land on disk first. * * (perf review 2026-07 F4) — now used by persistSingleTaskUpdate's * non-terminal checkpoint path; that caller calls flushPendingAtomicWrites() * before its read-modify-write load to defeat the stale-read window. * * ST-7: pass `skipCoalesce: true` for terminal task status transitions * (completed/failed/cancelled/needs_attention/skipped) to bypass the * 50ms coalesce window — a SIGKILL in that window would otherwise leave * tasks.json stale (showing "running") while events.jsonl already shows * the terminal event, causing false zombie detection / double-execution * on crash recovery. * * PERF round 2, Task 3: optional `durability` (default "full") overrides the * coalesced entry's durability. When the caller wants a non-terminal checkpoint * written without fsync (persistence.skipTasksFsync opt-in), pass * `durability: "best-effort"` — the coalesced entry stores it and the flush * forwards it to atomicWriteFile (see atomic-write.ts:997). Terminal * transitions (skipCoalesce=true) always fall to atomicWriteJson with full * durability regardless of this param. */ /** @internal */ export function saveRunTasksCoalesced( manifest: TeamRunManifest, tasks: TeamTaskState[], skipCoalesce: boolean = false, durability: WriteDurability = "full", ): void { // ST-4: refuse to persist [] over a previously-non-empty tasks file. if (!shouldPersistTasks(manifest, tasks)) return; // PERF (2026-08-24, Task 12): invalidating the WHOLE entry made every // loadRunManifestById after a coalesced save re-read + re-parse // manifest.json (24KB+) even though the manifest file did not change — // persistSingleTaskUpdate's next call (~500ms later) always paid it. Keep // the manifest half of the entry and zero only the tasks stamps, which is // the exact pre-existing signal for "tasks on disk may be stale" (coalesced // write not landed yet). The load path cooperates: the zeroed tasks stamps // force the slow path (tasks are always re-read from disk), but when the // retained manifest stamps still match a fresh stat of manifest.json — the // same mtime/size verification the fast path performs — the retry loop // reuses the cached manifest object and skips the re-read + re-parse. // Crash safety is unchanged: a zeroed tasks stamp can only cause a miss, // never a stale hit, and a concurrent manifest rewrite changes mtime/size // so the manifest reuse never serves stale content. Generation semantics // unchanged — setManifestCache stamps the CURRENT generation, so a // concurrent writer's bump still invalidates us. const cached = manifestCache.get(manifest.stateRoot); if (cached) { setManifestCache(manifest.stateRoot, { ...cached, tasks, tasksMtimeMs: 0, tasksSize: 0 }); } else { invalidateRunCache(manifest.stateRoot); } try { fs.statSync(manifest.stateRoot); } catch { return; } atomicWriteJsonCoalesced(manifest.tasksPath, tasks, undefined, { compact: true, durability }, skipCoalesce); } export async function saveRunTasksAsync(manifest: TeamRunManifest, tasks: TeamTaskState[]): Promise { // ST-4: refuse to persist [] over a previously-non-empty tasks file. if (!shouldPersistTasks(manifest, tasks)) return; // FIX: Invalidate cache BEFORE atomic write to prevent stale cache serving. invalidateRunCache(manifest.stateRoot); try { await fsp.access(manifest.stateRoot); } catch { return; } await atomicWriteJsonAsync(manifest.tasksPath, tasks, { compact: true }); } /** * Save manifest and tasks files with individual atomic writes. * FIX: Changed from Promise.all (parallel, non-jointly-atomic) to sequential * writes to ensure manifest is written before tasks. A crash between writes * leaves them in a known state (manifest is older, tasks is newer) which * loadRunManifestById's retry loop can detect via mtime comparison. * NOTE: There is no stale-reconciler component — the retry loop provides * best-effort detection of mid-write crashes by re-reading until mtime/size * are stable. For strict atomicity, callers should use withRunLock(). * FIX: Returns a result object so callers know which write step failed. * If manifest write succeeds but tasks write fails, the caller can recover. */ /** @internal */ interface SaveManifestAndTasksResult { manifestWritten: boolean; tasksWritten: boolean; error?: string; } async function saveManifestAndTasksAtomic(manifest: TeamRunManifest, tasks: TeamTaskState[]): Promise { let manifestWritten = false; let tasksWritten = false; try { await withRunLock(manifest, async () => { // FIX: Invalidate cache BEFORE writes to prevent stale cache serving. // Sequential writes instead of Promise.all to ensure manifest is // written before tasks. If a crash occurs between writes, manifest is // the older timestamp which stale-reconciler uses to detect inconsistency. invalidateRunCache(manifest.stateRoot); await atomicWriteJsonAsync(path.join(manifest.stateRoot, "manifest.json"), manifest); manifestWritten = true; await atomicWriteJsonAsync(manifest.tasksPath, tasks, { compact: true }); tasksWritten = true; }); } catch (err) { return { manifestWritten, tasksWritten, // FIX: Use String(err) to safely convert any thrown value (Error, // string, number, null, undefined) to a string. err.message would be // undefined for non-Error throwables, losing useful context. error: String(err), }; } return { manifestWritten: true, tasksWritten: true }; } /** @internal */ function saveManifestAndTasksAtomicSync(manifest: TeamRunManifest, tasks: TeamTaskState[]): SaveManifestAndTasksResult { let manifestWritten = false; let tasksWritten = false; try { withRunLockSync(manifest, () => { // FIX: Invalidate cache BEFORE writes to prevent stale cache serving. invalidateRunCache(manifest.stateRoot); atomicWriteJson(path.join(manifest.stateRoot, "manifest.json"), manifest); manifestWritten = true; atomicWriteJson(manifest.tasksPath, tasks, { compact: true }); tasksWritten = true; }); } catch (err) { return { manifestWritten, tasksWritten, // FIX: Use String(err) to safely convert any thrown value (Error, // string, number, null, undefined) to a string. err.message would be // undefined for non-Error throwables, losing useful context. error: String(err), }; } return { manifestWritten: true, tasksWritten: true }; } export interface UpdateRunStatusOptions { data?: Record; metadata?: Parameters[1]["metadata"]; /** Finding 8 (2026-09-23): allow leaving a TERMINAL disk status (resume is the * only legitimate terminal-exit flow). Defaults to false — a raw save with an * in-memory "running" manifest must never erase an externally-written * cancelled/failed/completed status (live race team_20260923175042: cancel * 17:51:08 → intermediate task-runner save re-wrote "running" → finalize * completed — the cancel was fully erased). */ allowTerminalExit?: boolean; } /** Finding 8: statuses a disk write must never silently leave. Mirrors * isRunTerminalPreserved (merge-loop.ts) — kept local to avoid a runtime dep * from the store layer to the runtime layer. */ const DISK_TERMINAL_STATUSES: ReadonlySet = new Set(["cancelled", "failed", "completed"]); /** Finding 8 write-layer guard: if the DISK manifest is terminal and the * incoming write carries a NON-terminal status (the erase class — a mid-flight * saver with a stale in-memory "running" manifest), preserve the disk * status/summary/updatedAt while keeping every other incoming field (artifacts, * usage, surface…). Terminal→terminal re-decisions (e.g. cancelling a run that * just completed — pinned by resume-cancel.test.ts) are LEGITIMATE and pass * through; they are governed by canTransitionRunStatus at the updateRunStatus * layer. Returns the effective manifest to persist. */ function preserveDiskTerminalStatus(manifest: TeamRunManifest, allowTerminalExit: boolean | undefined): TeamRunManifest { if (allowTerminalExit) return manifest; // Terminal→terminal re-decisions pass through (governed by // canTransitionRunStatus at the updateRunStatus layer); the guard applies // ONLY to the erase class: a NON-terminal incoming status over terminal disk. if (DISK_TERMINAL_STATUSES.has(manifest.status)) return manifest; try { const manifestPath = path.join(manifest.stateRoot, "manifest.json"); // Raw read (no cache): cross-process cancel writes must be seen NOW. const raw = fs.readFileSync(manifestPath, "utf-8"); const disk = JSON.parse(raw) as TeamRunManifest; if (!DISK_TERMINAL_STATUSES.has(disk.status) || disk.status === manifest.status) return manifest; return { ...manifest, status: disk.status, summary: disk.summary, updatedAt: disk.updatedAt }; } catch { return manifest; // no disk manifest yet (create) or unreadable — normal write } } export function updateRunStatus( manifest: TeamRunManifest, status: TeamRunManifest["status"], summary?: string, options: UpdateRunStatusOptions = {}, ): TeamRunManifest { if (!canTransitionRunStatus(manifest.status, status)) { throw errors.invalidStatusTransition(manifest.status, status); } const updated: TeamRunManifest = { ...manifest, status, updatedAt: new Date().toISOString(), summary: summary ?? manifest.summary, }; // Finding 8 (2026-09-23): the write-layer guard may PRESERVE a terminal disk // status (an external cancel/reconciler write the in-memory manifest never // saw). In that case do NOT emit run., do NOT flip the manifest — // record the refusal and return the preserved state. Resume passes // allowTerminalExit for its legitimate cancelled→running transition. const saved = saveRunManifest(updated, { allowTerminalExit: options.allowTerminalExit }); if (saved.status !== status) { appendEvent(saved.eventsPath, { type: "run.terminal_preserved", runId: saved.runId, message: `Preserved terminal status '${saved.status}'; refused in-memory transition to '${status}'.`, data: { preserved: saved.status, refused: status }, }); return saved; } // Unregister from active-run-index when run reaches a terminal status. // Without this, stale entries accumulate (e.g. integration tests in /tmp) and // Pi UI shows ghost "queued" runs that are actually completed/failed/cancelled. // Note: "blocked" is excluded because blocked runs can be unblocked later. if (status === "completed" || status === "failed" || status === "cancelled") { try { unregisterActiveRun(updated.runId); } catch { /* non-critical */ } } appendEvent(updated.eventsPath, { type: `run.${status}`, runId: updated.runId, message: summary, ...(options.data ? { data: options.data } : {}), metadata: { provenance: "team_runner", sessionIdentity: { title: updated.team, workspace: updated.cwd, purpose: updated.goal, }, ownership: { owner: updated.team, workflowScope: updated.workflow ?? "manual", watcherAction: "act", }, confidence: "high", ...options.metadata, }, }); return updated; } export function __test__manifestCacheSize(): number { return manifestCache.size; } export function __test__clearManifestCache(): void { manifestCache.clear(); manifestCacheGeneration.clear(); } /** * OPT-08 additive helper: drop the manifest cache entry for `stateRoot` after * flushing any pending coalesced atomic writes (process-wide). Forces the next * `loadRunManifestById` / `loadRunManifestByIdAsync` for this run to hit disk * instead of serving a possibly-stale cache hit. Idempotent: safe to call * multiple times, including for unknown stateRoots. * * This is an additive helper ONLY — it does NOT replace the load-bearing CAS * loop in `persistSingleTaskUpdate`, nor does it bypass `withRunLock` / * `fs.statSync`-based invalidation. It is a safe exit hatch for callers that * want a guaranteed cache drop (e.g. cleanup paths, end-of-run shutdown). * * Scope note (R10-2): the flush is scoped to the run's `tasks.json` — the only * file that (a) feeds this cache and (b) can have a pending coalesced write * (`manifest.json` is only ever written through immediate/durable paths, and * `agents.json`/`status.json` coalesced writes live outside this cache). * Unrelated coalesced writes for OTHER runs stay pending on their own timers; * unloading one run must not block on them. */ export async function unloadRun(stateRoot: string): Promise { // Flush first so any in-flight buffered write lands on disk before we drop // the cache entry. Otherwise a coalesced write could fire AFTER unloadRun // returns, re-populate the cache from disk, and re-stale it. flushPendingAtomicWrites(path.join(stateRoot, "tasks.json")); // invalidateRunCache bumps the per-stateRoot generation counter so even if // some in-process reader has a stale reference, the next cache lookup // misses (generation mismatch) and re-reads from disk. invalidateRunCache(stateRoot); } /** * OPT-08 additive helper: read-only inspection of the manifest cache. * Returns current size, configured limits, and per-stateRoot observability * (cache generation + age in ms). Pure read — never mutates cache state. * Useful for `tests`, dashboards, and pre-shutdown sanity checks. */ export interface ManifestCacheStats { size: number; maxEntries: number; ttlMs: number; perStateRoot: Record; } export function getManifestCacheStats(): ManifestCacheStats { const now = Date.now(); const perStateRoot: Record = {}; for (const [stateRoot, entry] of manifestCache.entries()) { perStateRoot[stateRoot] = { generation: entry.generation ?? genOf(stateRoot), // ageMs is 0 when cachedAt is unset (defensive — setManifestCache // always sets it today, but a future regression would otherwise // surface as `undefined` in the stats object). ageMs: entry.cachedAt ? now - entry.cachedAt : 0, }; } return { size: manifestCache.size, maxEntries: DEFAULT_CACHE.manifestMaxEntries, ttlMs: MANIFEST_CACHE_TTL_MS, perStateRoot, }; } async function readJsonFileAsync(filePath: string): Promise { try { return JSON.parse(await fs.promises.readFile(filePath, "utf-8")) as T; } catch (err) { const code = (err as NodeJS.ErrnoException).code; if (code !== "ENOENT" && code !== "ENOTDIR") { logInternalError("readJsonFileAsync", err, `filePath=${filePath}`); } return undefined; } } /** * Load a run manifest and its tasks by runId. * WARNING: This function provides best-effort consistency only. The sentinel-based * retry loop does NOT guarantee manifest/tasks consistency under contention — * a concurrent writer can complete a full write cycle between the final stat * and the read. For strict consistency, callers MUST wrap load+modify+save in * withRunLock(). Callers that need guaranteed consistency should use the lock. */ export function loadRunManifestById(cwd: string, runId: string): { manifest: TeamRunManifest; tasks: TeamTaskState[] } | undefined { const stateRoot = resolveRunStateRoot(cwd, runId); if (!stateRoot) return undefined; const manifestPath = path.join(stateRoot, "manifest.json"); const tasksPath = path.join(stateRoot, "tasks.json"); let manifestStat: fs.Stats; try { manifestStat = statManifestWithWindowsRetry(manifestPath); } catch { return undefined; } const cached = manifestCache.get(stateRoot); // Issue 1 fix: Note that the cache read-check-use sequence below is NOT atomic // under concurrent modification. While individual JS Map operations are atomic, // the sequence spanning multiple filesystem operations is not atomic. Callers MUST // NOT invoke loadRunManifestById concurrently with cache-modifying operations // (setManifestCache, invalidateRunCache) for the same stateRoot. The generation // counter provides some protection, but does not make the read-check-use atomic. let tasksStat: fs.Stats | undefined; try { tasksStat = fs.statSync(tasksPath); } catch { tasksStat = undefined; } const tasksMtimeMs = tasksStat?.mtimeMs ?? 0; if ( cached && cached.manifestMtimeMs === manifestStat.mtimeMs && cached.manifestSize === manifestStat.size && cached.tasksMtimeMs === tasksMtimeMs && cached.tasksSize === (tasksStat?.size ?? 0) && cached.generation === genOf(stateRoot) ) { // TTL eviction: expire stale entries even if mtime matches // FIX: Also evict entries where cachedAt is undefined — such entries are // effectively immortal otherwise (the `cached.cachedAt &&` check would skip // them every time). This can happen if a cache entry was created by // setManifestCache that didn't set cachedAt (shouldn't happen in current // code, but defensive against future regressions). if (!cached.cachedAt || Date.now() - cached.cachedAt > MANIFEST_CACHE_TTL_MS) { manifestCache.delete(stateRoot); } else if (!validateRunManifestPaths(cwd, runId, cached.manifest, stateRoot, tasksPath)) { manifestCache.delete(stateRoot); return undefined; } else if (!fs.existsSync(tasksPath)) { // Tasks file was deleted after cache was populated — this is a cache miss, // not a manifest inconsistency. The cache check passes because // tasksMtimeMs=0 and tasksSize=0 match a cache entry written when tasks.json // didn't exist. We fall through to the retry loop which will correctly // detect the missing file and return undefined. The alternative (treating // this as an inconsistency) would require failing the cache hit path, which // is unnecessary since the retry loop handles it correctly anyway. manifestCache.delete(stateRoot); return undefined; } else { return { manifest: cached.manifest, tasks: cached.tasks ?? [] }; } } // FIX: Sentinel-based retry loop for best-effort consistency. Re-stat and // re-read until mtime/size are stable. The retry limit is now configurable // via LOAD_MANIFEST_RETRY_LIMIT (was hardcoded 3). High contention can still // cause non-convergence — for strict consistency, callers MUST use withRunLock(). // Issue 3 fix: Made retry limit configurable instead of hardcoded 3. let attempts = 0; let manifest: TeamRunManifest | undefined; let tasks: TeamTaskState[] | undefined; while (attempts < LOAD_MANIFEST_RETRY_LIMIT) { const freshStat = fs.statSync(manifestPath); // PERF (2026-08-24, Task 12 realized): after saveRunTasksCoalesced the // cache keeps the manifest half of the entry while zeroing only the // tasks stamps, so this slow path runs solely to refresh tasks. When the // retained manifest stamps still match the fresh stat — the exact // mtime/size verification the fast path performs — reuse the cached // manifest object instead of re-reading + re-parsing manifest.json. // Stamp verification is NOT bypassed: any manifest.json rewrite changes // mtime/size and falls back to the disk read. The tasks re-read below // still always happens — zeroed tasks stamps must force it. manifest = cached && cached.manifestMtimeMs === freshStat.mtimeMs && cached.manifestSize === freshStat.size ? cached.manifest : readJsonFile(manifestPath); const freshTasksStat = fs.existsSync(tasksPath) ? fs.statSync(tasksPath) : undefined; tasks = loadTasksWithRecovery(tasksPath, manifest?.eventsPath ?? path.join(stateRoot, "events.jsonl"), manifest?.runId ?? runId); // If size/mtime didn't change between stat and read, we're consistent. if ( freshStat.mtimeMs === manifestStat.mtimeMs && freshStat.size === manifestStat.size && (!freshTasksStat || (freshTasksStat.mtimeMs === tasksStat?.mtimeMs && freshTasksStat.size === tasksStat?.size)) ) { break; } attempts += 1; manifestStat = freshStat; tasksStat = freshTasksStat; } // WARNING: Best-effort consistency only — retry loop detected mtime/size // instability. A concurrent writer can still complete a full write cycle // between the final stat and the read. Callers needing strict consistency // MUST use withRunLock() around load+modify+save. if (attempts > 0) { // Round 19: downgrade to debug — retry-loop instability is expected under // concurrent writes (live team runs constantly append to tasks.json). // This is best-effort by design; strict consistency requires withRunLock(). console.debug( `[state-store] loadRunManifestById: retry loop detected instability for run ${runId} after ${attempts} attempt(s) — best-effort only, use withRunLock() for strict consistency`, ); } // NOTE: manifest mtime may legitimately be >= tasks mtime because // saveManifestAndTasksAtomicSync writes manifest before tasks. However, // if a crash occurs AFTER tasks is written but BEFORE manifest is written, // tasks mtime would be > manifest mtime (the opposite). The retry loop // above detects this crash state by re-reading until mtime/size are stable // — the final stable state is what gets used. Because the retry loop // handles the crash case, we do NOT fail based on this comparison alone. // It does not indicate corruption on its own. // S-01: warn (do not throw) on schemaVersion mismatch — future version // bumps will add migration logic here. if (manifest && manifest.schemaVersion !== CURRENT_SCHEMA_VERSION) { logInternalError( "state-store", new Error( `Manifest schemaVersion mismatch: expected ${CURRENT_SCHEMA_VERSION}, got ${manifest.schemaVersion}. Run ${runId} may be incompatible.`, ), undefined, "warn", ); } // STATE-3: readJsonFile returns undefined for BOTH missing (ENOENT) and corrupt // (SyntaxError). A corrupt manifest currently makes the run silently invisible. // If the file EXISTS but readJsonFile returned undefined, it is corrupt → quarantine // it (preserve for diagnosis) + log + bail. Manifest reconstruction from events is // infeasible (run.created lacks manifest fields), so quarantine+log is the recovery. if (!manifest && fs.existsSync(manifestPath)) { quarantineCorruptFile(manifestPath); logInternalError( "state-store", new Error( `STATE-3: manifest.json for run ${runId} exists but is unparseable — quarantined. Run treated as missing; preserve the .corrupt-* file for diagnosis.`, ), undefined, "error", ); return undefined; } if (!manifest || !validateRunManifestPaths(cwd, runId, manifest, stateRoot, tasksPath)) return undefined; setManifestCache(stateRoot, { manifest, tasks: tasks ?? [], manifestMtimeMs: manifestStat.mtimeMs, manifestSize: manifestStat.size, tasksMtimeMs, tasksSize: tasksStat?.size ?? 0, }); return { manifest, tasks: tasks ?? [] }; } export async function loadRunManifestByIdAsync( cwd: string, runId: string, ): Promise<{ manifest: TeamRunManifest; tasks: TeamTaskState[] } | undefined> { const stateRoot = resolveRunStateRoot(cwd, runId); if (!stateRoot) return undefined; const manifestPath = path.join(stateRoot, "manifest.json"); const tasksPath = path.join(stateRoot, "tasks.json"); let manifestStat: fs.Stats; try { manifestStat = await fs.promises.stat(manifestPath); } catch { return undefined; } const cached = manifestCache.get(stateRoot); // Issue 1 fix: Note that the cache read-check-use sequence below is NOT atomic // under concurrent modification. Same as sync version — see loadRunManifestById. let tasksStat: fs.Stats | undefined; try { tasksStat = await fs.promises.stat(tasksPath); } catch { tasksStat = undefined; } const tasksMtimeMs = tasksStat?.mtimeMs ?? 0; if ( cached && cached.manifestMtimeMs === manifestStat.mtimeMs && cached.manifestSize === manifestStat.size && cached.tasksMtimeMs === tasksMtimeMs && cached.tasksSize === (tasksStat?.size ?? 0) && cached.generation === genOf(stateRoot) ) { // TTL eviction: expire stale entries even if mtime matches // FIX: Also evict entries where cachedAt is undefined — such entries are // effectively immortal otherwise (the `cached.cachedAt &&` check would skip // them every time). This can happen if a cache entry was created by // setManifestCache that didn't set cachedAt (shouldn't happen in current // code, but defensive against future regressions). if (!cached.cachedAt || Date.now() - cached.cachedAt > MANIFEST_CACHE_TTL_MS) { manifestCache.delete(stateRoot); } else if (!validateRunManifestPaths(cwd, runId, cached.manifest, stateRoot, tasksPath)) { manifestCache.delete(stateRoot); return undefined; } else if ( !(await fsp.access(tasksPath).then( () => true, () => false, )) ) { // Tasks file was deleted after cache was populated — do not serve stale cache. manifestCache.delete(stateRoot); return undefined; } else { return { manifest: cached.manifest, tasks: cached.tasks ?? [] }; } } // FIX: Sentinel-based retry loop to close TOCTOU window between stat and read. // Matches the pattern used in the sync loadRunManifestById. // Issue 3 fix: Made retry limit configurable instead of hardcoded 3. let manifest: TeamRunManifest | undefined; let tasks: TeamTaskState[] | undefined; let attempts = 0; while (attempts < LOAD_MANIFEST_RETRY_LIMIT) { const freshStat = await fs.promises.stat(manifestPath); // PERF (2026-08-24, Task 12 realized): async twin of the sync reuse — // after saveRunTasksCoalesced only the tasks stamps are zeroed, so when // the retained manifest stamps match the fresh stat (same verification // the fast path performs), reuse the cached manifest object instead of // re-reading + re-parsing manifest.json. A manifest.json rewrite changes // mtime/size and falls back to the disk read; the tasks re-read below // always happens. manifest = cached && cached.manifestMtimeMs === freshStat.mtimeMs && cached.manifestSize === freshStat.size ? cached.manifest : await readJsonFileAsync(manifestPath); const freshTasksStat = await fs.promises.stat(tasksPath).catch(() => undefined); tasks = await loadTasksWithRecoveryAsync( tasksPath, manifest?.eventsPath ?? path.join(stateRoot, "events.jsonl"), manifest?.runId ?? runId, ); // If size/mtime didn't change between stat and read, we're consistent. if ( freshStat.mtimeMs === manifestStat.mtimeMs && freshStat.size === manifestStat.size && (!freshTasksStat || (freshTasksStat.mtimeMs === tasksStat?.mtimeMs && freshTasksStat.size === tasksStat?.size)) ) { break; } attempts += 1; manifestStat = freshStat; tasksStat = freshTasksStat; } // WARNING: Best-effort consistency only — retry loop detected mtime/size // instability. A concurrent writer can still complete a full write cycle // between the final stat and the read. Callers needing strict consistency // MUST use withRunLock() around load+modify+save. if (attempts > 0) { // Round 19: downgrade to debug — retry-loop instability is expected under // concurrent writes (live team runs constantly append to tasks.json). // This is best-effort by design; strict consistency requires withRunLock(). console.debug( `[state-store] loadRunManifestByIdAsync: retry loop detected instability for run ${runId} after ${attempts} attempt(s) — best-effort only, use withRunLock() for strict consistency`, ); } // NOTE: manifest mtime may legitimately be >= tasks mtime because // saveManifestAndTasksAtomicSync writes manifest before tasks. However, // if a crash occurs AFTER tasks is written but BEFORE manifest is written, // tasks mtime would be > manifest mtime (the opposite). The retry loop // above detects this crash state by re-reading until mtime/size are stable // — the final stable state is what gets used. Because the retry loop // handles the crash case, we do NOT fail based on this comparison alone. // It does not indicate corruption on its own. // S-01: warn (do not throw) on schemaVersion mismatch — future version // bumps will add migration logic here. if (manifest && manifest.schemaVersion !== CURRENT_SCHEMA_VERSION) { logInternalError( "state-store", new Error( `Manifest schemaVersion mismatch: expected ${CURRENT_SCHEMA_VERSION}, got ${manifest.schemaVersion}. Run ${runId} may be incompatible.`, ), undefined, "warn", ); } // STATE-3 (async twin): readJsonFileAsync returns undefined for BOTH missing (ENOENT) // and corrupt (SyntaxError). Mirror the sync loadRunManifestById STATE-3 fix — if the // file EXISTS but the read returned undefined, it is corrupt → quarantine (preserve for // diagnosis) + log + bail. Manifest reconstruction from events is infeasible (run.created // lacks manifest fields), so quarantine+log is the recovery. if (!manifest && fs.existsSync(manifestPath)) { quarantineCorruptFile(manifestPath); logInternalError( "state-store", new Error( `STATE-3 async: manifest.json for run ${runId} exists but is unparseable — quarantined. Run treated as missing; preserve the .corrupt-* file for diagnosis.`, ), undefined, "error", ); return undefined; } if (!manifest || !validateRunManifestPaths(cwd, runId, manifest, stateRoot, tasksPath)) return undefined; setManifestCache(stateRoot, { manifest, tasks: tasks ?? [], manifestMtimeMs: manifestStat.mtimeMs, manifestSize: manifestStat.size, tasksMtimeMs, tasksSize: tasksStat?.size ?? 0, }); return { manifest, tasks: tasks ?? [] }; }