diff --git a/AGENTS.md b/AGENTS.md index b8d6a3d8..24d1ae16 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -115,7 +115,7 @@ process around it. `gh`/`glab` CLI identity, not separate GitHub/GitLab token files. Doctor also fails when OpenCode's frozen `@latest` package cache lags the published `workit-opencode` (delete the cache dir and restart OpenCode). -- Task state lives under `.workit/` in the session directory, even when that directory is not a Git repository; never edit it directly. Git/hosting actions accept `cwd` for an action-time target checkout without prior attachment; the coordinator task keeps the writer, while a conflicting writer in the target checkout blocks mutation. A managed action holds the target metadata lock through settlement so another writer cannot acquire during the effect. Use the eight shared operation families (`workit_task`, `workit_policy`, `workit_evidence`, `workit_finding`, `workit_decision`, `workit_worker`, `workit_writer`, `workit_state`) with closed `action` enums. CLI surface is `workit ` (hyphenated actions). There are no `workit flow` aliases. +- Task state lives under `.workit/` in the session directory, even when that directory is not a Git repository; never edit it directly. A stale `metadata.lock` (dead/reused pid in the same host, pid namespace and boot; anything else only past its TTL) is reclaimed by the next write or cleared with `workit doctor --fix-lock` (`--force --yes` for an unverifiable lock); contention with a live holder returns retryable `busy`, and `recovery_required` is reserved for genuine state damage. Git/hosting actions accept `cwd` for an action-time target checkout without prior attachment; the coordinator task keeps the writer, while a conflicting writer in the target checkout blocks mutation. A managed action holds the target metadata lock through settlement so another writer cannot acquire during the effect. Use the eight shared operation families (`workit_task`, `workit_policy`, `workit_evidence`, `workit_finding`, `workit_decision`, `workit_worker`, `workit_writer`, `workit_state`) with closed `action` enums. CLI surface is `workit ` (hyphenated actions). There are no `workit flow` aliases. - `task.list` defaults to a compact, 20-item active/paused projection; bounded closed/all history is opt-in with `status` and `limit`, and `task.inspect` defaults to `summary`. Closed views evaluate their captured closure candidate and carry no current writer lease. - Workit tracking is optional: direct investigation, questions, non-Git work, and routine reversible edits need no task, policy assessment, or writer calls. Start one compact task only when handoff, dependencies, concurrent actors, or meaningful decisions make continuity useful; assess or reassess when the relevant rules/evidence require it. `policy.preview` stays read-only. - Routine authorized branch/commit work chooses native host Git/shell tools from the outset when managed coordination or reconciliation is unnecessary. Resolve target conventions; never switch paths after a denial or uncertain managed effect to evade safeguards. A local-commit endpoint creates no PR-readiness or task-closure ceremony. The distributed bootstrap carries this routing guidance. diff --git a/CHANGELOG.md b/CHANGELOG.md index 6cee7bfc..dd88c2da 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- A `.workit/metadata.lock` left by a dead process (or a reused pid, or a + foreign-host/namespace lock past its TTL) no longer bricks the store: the next + write reclaims it. Locks carry the pid namespace and boot id so a container + sharing the hostname is never judged by this host's process table, and a + lock from before a reboot is reclaimed at once. + Contention with a live writer retries briefly (250 ms in-process, 2 s in the + CLI) and returns the retryable `busy` code instead of `recovery_required`. + `workit doctor` reports a stale lock (`workspace_lock`); `workit doctor + --fix-lock` clears it under the reclaim guard, and `--force --yes` clears an + unverifiable lock explicitly. - Compact task context retains the newest decisions and surfaces bounded, redacted choice summaries instead of selecting an arbitrary UUID-ordered set. - Distributed cross-repository guidance binds unfinished work to its checkout, diff --git a/README.md b/README.md index 85e245c9..da0a81d5 100644 --- a/README.md +++ b/README.md @@ -203,6 +203,8 @@ workit upgrade --apply --confirm # apply a reviewed preview workit upgrade --cli --apply --confirm # also update an existing global CLI workit launch pi --auto-upgrade -- # update before starting Pi workit doctor # offline installation health report (--json for machines) +workit doctor --fix-lock # clear a stale .workit metadata lock (WORKFLOW_WORKSPACE_ROOT or cwd) +workit doctor --fix-lock --force [--yes] # clear a lock whose owner cannot be verified workit [--payload ] [--task ] [--confirm] [--json] workit action --payload [--preview] [--confirm] [--json] workit handoff --task [--json] @@ -403,7 +405,19 @@ The canonical target, relevant Git/remote state, and effective `gh`/`glab` account are checked again before a remote effect. The coordinator owns the Workit writer; an independently held writer in the target checkout remains a real conflict, and managed actions hold that checkout's Workit metadata lock -through effect settlement so a writer cannot acquire mid-action. New branch +through effect settlement so a writer cannot acquire mid-action. A metadata +lock whose owner is gone (dead or reused pid) is reclaimed by the next write. +A lock records its host plus, on Linux, its pid namespace and boot id; a lock +from another host, container namespace, boot, or an older Workit version cannot +be checked against this process table and is reclaimed only after a 10-minute +TTL (a lock from the same host and pid namespace but an earlier boot is +reclaimed at once). `workit doctor` warns when such an unverifiable lock has +blocked writes for over 30 s and prints `workit doctor --fix-lock --force --yes`. A write that meets a live holder retries briefly (250 ms inside host +plugins and the MCP server, 2 s in the CLI) and then returns the retryable +`busy` code, never `recovery_required`. `workit doctor` warns about a stale +lock and `workit doctor --fix-lock` clears it under the same reclaim guard +writers use; `--force` (with `--yes` or an interactive confirmation) is the +explicit escape hatch for a lock whose owner cannot be verified. New branch setup shows both the existing local base SHA and remote base SHA in its approval, rechecks them, and creates only from an approved commit. Workit does not reject Git-valid branch names or user commit diff --git a/packages/workit-cli/src/index.tsx b/packages/workit-cli/src/index.tsx index 46d745b8..f39f8c31 100755 --- a/packages/workit-cli/src/index.tsx +++ b/packages/workit-cli/src/index.tsx @@ -8,6 +8,12 @@ import { createLogger } from "@brainervirus/workit-core/src/core/logger"; import { EVENT, errorDetail } from "@brainervirus/workit-core/src/core/boundary"; import { setDiagnosticLogger } from "@brainervirus/workit-core/src/core/config"; import { runDoctor } from "@brainervirus/workit-core/src/core/doctor"; +import { + clearStaleMetadataLock, + inspectMetadataLock, + setDefaultLockTimeout, +} from "@brainervirus/workit-core/src/core/store-lock"; +import { createInterface } from "node:readline/promises"; import { applySetupPreview, buildSetupPreview, @@ -24,7 +30,7 @@ import { } from "@brainervirus/workit-core/src/core/uninstall"; import { applyWizardBranchPolicy } from "./logic"; import { runCutoverCommand } from "./cutover-cli"; -import { runActionCommand, runTaskCommand, TASK_FAMILIES } from "./task"; +import { runActionCommand, runTaskCommand, TASK_FAMILIES, workspaceRootFor } from "./task"; import { externalActionHelp } from "@brainervirus/workit-core/src/core"; import { runLaunchCommand, runUpgradeCommand } from "./upgrade"; @@ -58,6 +64,8 @@ Usage: workit upgrade Preview upgrades (--apply --confirm; --hosts=a,b; --cli for the CLI) workit launch [--auto-upgrade] [-- args] Upgrade before host startup workit doctor Verify the offline installation health (add --json for a machine-readable report) + --fix-lock clears a stale .workit metadata lock in the workspace root + --fix-lock --force [--yes] clears it even when its owner cannot be verified workit uninstall Remove workit host registrations interactively (~/.config/workit is kept) workit cutover Preview or apply an explicit v1 cutover (apply requires --confirm) ${COMMAND_DESCRIPTIONS.map(([cmd, desc]) => ` ${cmd.padEnd(helpColumn)}${desc}`).join("\n")} @@ -310,11 +318,54 @@ export async function runUninstall() { // `workit doctor` (DG-07): offline engine, human or --json report, exit code // reflects the health. Never writes the report to stderr (the logger owns that). -function runDoctorCommand(args: string[]) { - const report = runDoctor({ host: "cli", cwd: process.cwd() }); +// Explicit escape hatch for a lock whose owner cannot be verified (no process +// start time, a foreign pid namespace): show the holder, then require --yes or +// an interactive confirmation. +async function confirmForcedLockClear(root: string, args: string[]): Promise { + const lock = inspectMetadataLock(root); + if (!lock.present) return true; + const owner = lock.owner + ? `pid ${lock.owner.pid} on ${lock.owner.host} (start ${lock.owner.processStart ?? "unknown"})` + : "unreadable lock"; + console.log(`fix-lock --force: ${lock.path} is held by ${owner}: ${lock.reason}`); + if (args.includes("--yes")) return true; + if (process.stdin.isTTY !== true) { + console.log("fix-lock --force: refusing without --yes outside an interactive terminal"); + return false; + } + const rl = createInterface({ input: process.stdin, output: process.stdout }); + try { + const answer = await rl.question("Remove this lock even if its holder may be alive? [y/N] "); + return /^y(es)?$/i.test(answer.trim()); + } finally { + rl.close(); + } +} + +async function runDoctorCommand(args: string[]) { + // --fix-lock runs first so the report reflects the cleaned state. Without + // --force it removes only a lock whose owner is provably gone. + const root = workspaceRootFor(); + const force = args.includes("--force"); + let fixLock: ReturnType | null = null; + if (args.includes("--fix-lock")) { + if (force && !(await confirmForcedLockClear(root, args))) process.exit(1); + fixLock = clearStaleMetadataLock(root, { force }); + } + const report = runDoctor({ host: "cli", cwd: process.cwd(), workspaceRoot: root }); if (args.includes("--json")) { - console.log(JSON.stringify(report, null, 2)); + console.log(JSON.stringify(fixLock ? { ...report, fixLock } : report, null, 2)); } else { + if (fixLock) { + const what = fixLock.cleared + ? `cleared ${force ? "" : "stale "}lock ${fixLock.path} (${fixLock.reason})` + : fixLock.state === "absent" + ? "no metadata lock to clear" + : `kept lock ${fixLock.path}: ${fixLock.skipped ?? fixLock.reason}`; + console.log( + `fix-lock: ${what}${fixLock.guardCleared ? "; removed abandoned reclaim guard" : ""}`, + ); + } console.log( `workit doctor — ${report.ok ? "healthy" : "problems found"} (${report.offline ? "offline" : "online"})`, ); @@ -334,6 +385,9 @@ if (import.meta.main) { const args = process.argv.slice(2); const [subcommand] = args; setDiagnosticLogger(logger); + // The CLI owns its process, so a contended write may wait longer than an + // in-process host could afford before reporting busy. + setDefaultLockTimeout(2_000); logger.info(EVENT.initialization, { host: "cli", command: subcommand }); // The CLI owns its process: uncaught failures are logged and surfaced with a // nonzero exit instead of a silent crash (DG-04). @@ -351,7 +405,7 @@ if (import.meta.main) { } else if (subcommand === "init") { await runInit(); } else if (subcommand === "doctor") { - runDoctorCommand(args); + await runDoctorCommand(args); } else if ((TASK_FAMILIES as readonly string[]).includes(subcommand)) { process.exit(await runTaskCommand(args)); } else if (subcommand === "action") { diff --git a/packages/workit-cli/src/task.ts b/packages/workit-cli/src/task.ts index 5ae75cc6..f7ae0549 100644 --- a/packages/workit-cli/src/task.ts +++ b/packages/workit-cli/src/task.ts @@ -46,6 +46,10 @@ import type { Provenance } from "@brainervirus/workit-core/src/core/task-contrac import { canonicalJson, type Result } from "@brainervirus/workit-core/src/core/task-contract"; export const TASK_FAMILIES = OPERATION_FAMILIES; + +/** The Workit store root for a CLI command: explicit root, WORKFLOW_WORKSPACE_ROOT, then cwd. */ +export const workspaceRootFor = (deps: { root?: string; cwd?: string } = {}): string => + deps.root ?? process.env.WORKFLOW_WORKSPACE_ROOT ?? deps.cwd ?? process.cwd(); export const TASK_ACTIONS = { task: ["start", "list", "inspect", "revise", "progress", "pause", "resume", "close"], policy: ["assess", "preview", "explain"], @@ -390,7 +394,7 @@ export async function runTaskCommand(argv: string[], deps: TaskCliDeps = {}): Pr ); return parsed.usage ? 2 : 1; } - const root = deps.root ?? process.env.WORKFLOW_WORKSPACE_ROOT ?? deps.cwd ?? process.cwd(); + const root = workspaceRootFor(deps); let observedConfirmation = parsed.parsed.observedConfirmation; if (needsConsent(parsed.parsed) && !parsed.parsed.confirmed) { const consent = await observeConsent(deps); @@ -650,7 +654,7 @@ export async function runActionCommand(argv: string[], deps: TaskCliDeps = {}): else printHuman(parsed, deps); return 2; } - const root = deps.root ?? process.env.WORKFLOW_WORKSPACE_ROOT ?? deps.cwd ?? process.cwd(); + const root = workspaceRootFor(deps); let resolved = resolveExternalActionRequest(root, parsed.data); if (!resolved.ok) { if (json) jsonResult(outOf(deps), resolved); diff --git a/packages/workit-core/src/core/doctor.ts b/packages/workit-core/src/core/doctor.ts index 324748e1..9413eb39 100644 --- a/packages/workit-core/src/core/doctor.ts +++ b/packages/workit-core/src/core/doctor.ts @@ -19,6 +19,7 @@ import { import os from "node:os"; import path from "node:path"; import { SUPPORT_MATRIX } from "./support-matrix"; +import { inspectMetadataLock } from "./store-lock"; import { bundleHashOfFile, isEphemeralCachePath } from "./runtime-identity"; import { EVENT } from "./boundary"; import { getDiagnosticLogger, isConfigObject } from "./config"; @@ -61,6 +62,7 @@ export type DoctorCheckId = | "duplicate_registration" | "malformed_config" | "workspace_mismatch" + | "workspace_lock" | "credential_metadata" | "github_identity" | "gitlab_identity" @@ -110,6 +112,8 @@ export type DoctorOptions = { /** Checkout containing packages/ (monorepo or share clone). */ dev?: string; cwd?: string; + /** Workit store root for the lock check (default: WORKFLOW_WORKSPACE_ROOT, then cwd). */ + workspaceRoot?: string; opencodeConfig?: string; /** OpenCode npm `@latest` package cache root (test seam). */ opencodePackageCacheDir?: string; @@ -129,6 +133,7 @@ type Resolved = { configDir: string; stateDir: string; cwd: string; + workspaceRoot: string; dev: string | null; opencodeConfig: string; opencodePackageCacheDir: string; @@ -179,6 +184,7 @@ const resolve = (options: DoctorOptions): Resolved => { configDir, stateDir, cwd, + workspaceRoot: options.workspaceRoot ?? env.WORKFLOW_WORKSPACE_ROOT ?? cwd, dev, opencodeConfig: options.opencodeConfig ?? path.join(home, ".config", "opencode", "opencode.json"), @@ -1693,6 +1699,44 @@ const checkManagedContentConflict = (res: Resolved): DoctorCheck => { }; }; +// The checkout's `.workit/metadata.lock`. Writes reclaim a stale lock by +// themselves, so a stale lock is a warning with an explicit cleanup command. +const BLOCKING_LOCK_WARN_MS = 30_000; +const checkWorkspaceLock = (res: Resolved): DoctorCheck => { + const lock = inspectMetadataLock(res.workspaceRoot); + const fix = "workit doctor --fix-lock"; + if (lock.guard === "abandoned") + return { + id: "workspace_lock", + status: "warn", + detail: `abandoned lock reclaim guard at ${lock.path}.reclaim`, + fix, + }; + if (lock.state === "absent") + return { id: "workspace_lock", status: "pass", detail: "no metadata lock held" }; + if (lock.state === "stale") + return { + id: "workspace_lock", + status: "warn", + detail: `stale metadata lock at ${lock.path}: ${lock.reason}`, + fix, + }; + // An unverifiable owner (other host, pid namespace, or an older Workit's + // lock) that has blocked writes this long needs an explicit decision. + if (lock.state === "unknown" && (lock.ageMs ?? 0) > BLOCKING_LOCK_WARN_MS) + return { + id: "workspace_lock", + status: "warn", + detail: `metadata lock at ${lock.path} has blocked writes for ${Math.round((lock.ageMs ?? 0) / 1000)}s and its owner cannot be verified: ${lock.reason}`, + fix: "workit doctor --fix-lock --force --yes", + }; + return { + id: "workspace_lock", + status: "pass", + detail: `metadata lock ${lock.reason} (writes retry, then report busy)`, + }; +}; + const RUN_CHECKS: Array<(res: Resolved) => DoctorCheck> = [ checkRuntime, checkVersions, @@ -1705,6 +1749,7 @@ const RUN_CHECKS: Array<(res: Resolved) => DoctorCheck> = [ checkDuplicateRegistration, checkMalformedConfig, checkWorkspaceMismatch, + checkWorkspaceLock, checkCredentialMetadata, checkGithubIdentity, checkGitLabIdentity, diff --git a/packages/workit-core/src/core/methods.ts b/packages/workit-core/src/core/methods.ts index 5cde8c9f..3ae48171 100644 --- a/packages/workit-core/src/core/methods.ts +++ b/packages/workit-core/src/core/methods.ts @@ -110,6 +110,11 @@ Start a record once for an explicit tracked objective; assess or reassess only when policy selection or changed evidence/constraints requires it. Omitted expectedRevision and expectedWorkspaceRevision use current values; explicit values are still concurrency-checked, so never copy revisions between calls. +A busy result means another live Workit call holds the checkout lock: retry the +same call; it is not a recovery condition. A lock left by a dead process is +reclaimed on the next write, and \`workit doctor --fix-lock\` clears it on demand. +A revision_conflict on a call that omitted expectedRevision is contention too: +re-read the record and retry the call. A solo edit does not need writer acquisition; use it when concurrent checkout writers need coordination. Record only observed facts and checks. Evidence can become stale when its bound candidate changes; reconcile findings against the diff --git a/packages/workit-core/src/core/store-lock.ts b/packages/workit-core/src/core/store-lock.ts new file mode 100644 index 00000000..9fa110ae --- /dev/null +++ b/packages/workit-core/src/core/store-lock.ts @@ -0,0 +1,301 @@ +import { spawnSync } from "node:child_process"; +import fs from "node:fs"; +import { hostname } from "node:os"; +import path from "node:path"; +import * as z from "zod"; +import { canonicalJson } from "./task-contract"; + +/** + * Ownership rules for a checkout's `.workit/metadata.lock`. + * + * The lock is a short mutex around one store mutation (or one managed effect). + * A lock whose owner is gone is reclaimed automatically; a lock whose owner is + * alive is contention, which callers report as retryable `busy`. + */ + +export type MetadataLock = { + pid: number; + processStart: string | null; + host: string; + nonce: string; + externalAction?: true; +}; + +const metadataLockSchema = z + .object({ + pid: z.number().int().nonnegative().safe(), + processStart: z.string().nullable(), + host: z.string().min(1), + nonce: z.string().min(1), + externalAction: z.literal(true).optional(), + }) + .strict(); + +/** A lock from another host cannot be checked for liveness; trust it this long. */ +export const FOREIGN_LOCK_TTL_MS = 10 * 60_000; +/** An empty or unparseable lock is a writer mid-create; after this it is debris. */ +export const UNREADABLE_LOCK_TTL_MS = 30_000; +/** A reclaim guard lives for microseconds; one older than this was abandoned. */ +export const RECLAIM_GUARD_TTL_MS = 30_000; +/** + * Total time a mutation waits for a live holder before returning `busy`. The + * wait blocks the calling thread, so in-process hosts (OpenCode, MCP, Pi) keep + * the short default; the CLI raises it for its own process. + */ +let defaultLockTimeoutMs = 250; +export const defaultLockTimeout = (): number => defaultLockTimeoutMs; +export const setDefaultLockTimeout = (ms: number): void => { + defaultLockTimeoutMs = ms; +}; + +export const parseMetadataLock = (raw: string): MetadataLock => { + let value: unknown; + try { + value = JSON.parse(raw); + } catch { + throw Object.assign(new Error("metadata lock is invalid"), { code: "metadata_lock_invalid" }); + } + const parsed = metadataLockSchema.safeParse(value); + if (!parsed.success) + throw Object.assign(new Error("metadata lock is invalid"), { code: "metadata_lock_invalid" }); + return parsed.data; +}; + +/** Lenient variant for the acquire loop: an unreadable lock is classified by age. */ +export const parseMetadataLockOrNull = (raw: string): MetadataLock | null => { + try { + return parseMetadataLock(raw); + } catch { + return null; + } +}; + +export const sameMetadataLock = (left: unknown, right: MetadataLock): boolean => { + const parsed = metadataLockSchema.safeParse(left); + return parsed.success && canonicalJson(parsed.data) === canonicalJson(right); +}; + +/** Field 22 (starttime) of /proc//stat; parsed after the last ")" so a comm with spaces cannot shift it. */ +export const parseProcStatStart = (stat: string): string | null => { + const close = stat.lastIndexOf(")"); + if (close < 0) return null; + // After ")": field 3 (state) is index 0, so field 22 is index 19. + return ( + stat + .slice(close + 1) + .trim() + .split(/\s+/)[19] ?? null + ); +}; + +export const processStartOf = (pid: number): string | null => { + if (process.platform === "linux") { + try { + return parseProcStatStart(fs.readFileSync(`/proc/${pid}/stat`, "utf8")); + } catch { + return null; + } + } + if (process.platform === "darwin" || process.platform === "freebsd") { + try { + const run = spawnSync("ps", ["-o", "lstart=", "-p", String(pid)], { + encoding: "utf8", + timeout: 1_000, + }); + const value = run.status === 0 ? run.stdout.trim() : ""; + return value || null; + } catch { + return null; + } + } + return null; +}; + +const readTrimmed = (file: string): string | null => { + try { + return fs.readFileSync(file, "utf8").trim() || null; + } catch { + return null; + } +}; + +/** + * Identity of the pid space this process lives in: hostname plus, on Linux, + * the pid-namespace inode and boot id. Containers that share the hostname + * (`--network host`) but not the pid namespace get a different identity, so + * their pids are never checked against this process table. Folded into the + * existing `host` string so older Workit versions still parse the lock. + */ +let cachedLockHost: string | null = null; +export const localLockHost = (): string => { + if (cachedLockHost !== null) return cachedLockHost; + let pidns: string | null = null; + try { + pidns = /\[(\d+)\]/.exec(fs.readlinkSync("/proc/self/ns/pid"))?.[1] ?? null; + } catch {} + const boot = readTrimmed("/proc/sys/kernel/random/boot_id"); + cachedLockHost = pidns || boot ? `${hostname()}#${pidns ?? "?"}:${boot ?? "?"}` : hostname(); + return cachedLockHost; +}; + +const pidAlive = (pid: number): boolean => { + if (!Number.isSafeInteger(pid) || pid <= 0) return false; + try { + process.kill(pid, 0); + return true; + } catch (error) { + // EPERM: the process exists but belongs to another user. + return (error as { code?: unknown }).code === "EPERM"; + } +}; + +export type LockOwnerState = { + /** live: wait. stale: reclaim. unknown: wait (cannot prove the owner is gone). */ + state: "live" | "stale" | "unknown"; + reason: string; +}; + +export const classifyLockOwner = ( + payload: unknown, + ageMs: number | null, + localHost: string = localLockHost(), +): LockOwnerState => { + const lock = metadataLockSchema.safeParse(payload); + if (!lock.success) + return ageMs !== null && ageMs > UNREADABLE_LOCK_TTL_MS + ? { state: "stale", reason: "unreadable lock left behind" } + : { state: "unknown", reason: "lock is being written" }; + const { pid, processStart, host } = lock.data; + // Same host and pid namespace but another boot: the machine rebooted since + // the lock was taken, so its owner cannot still be running. + const [lockName, lockSpace] = host.split("#"); + const [localName, localSpace] = localHost.split("#"); + if (lockSpace && localSpace && lockName === localName) { + const [lockNs, lockBoot] = lockSpace.split(":"); + const [localNs, localBoot] = localSpace.split(":"); + if (lockNs === localNs && lockNs !== "?" && lockBoot !== "?" && lockBoot !== localBoot) + return { state: "stale", reason: "lock was taken before this machine rebooted" }; + } + // Another host, another pid namespace, or a lock written by an older Workit + // without namespace identity: its pid cannot be checked here. + if (host !== localHost) + return ageMs !== null && ageMs > FOREIGN_LOCK_TTL_MS + ? { state: "stale", reason: `lock from ${host} is older than its TTL` } + : { state: "unknown", reason: `lock is held from ${host}; its pid cannot be checked here` }; + if (!pidAlive(pid)) return { state: "stale", reason: `pid ${pid} is not running` }; + const currentStart = processStartOf(pid); + if (processStart !== null && currentStart !== null && processStart !== currentStart) + return { state: "stale", reason: `pid ${pid} now belongs to a different process` }; + return { state: "live", reason: `held by running pid ${pid}` }; +}; + +const ageOf = (file: string, nowMs: number): number | null => { + try { + return nowMs - fs.lstatSync(file).mtimeMs; + } catch { + return null; + } +}; + +export const lockPathFor = (root: string) => path.join(root, ".workit", "metadata.lock"); + +/** Remove a reclaim guard abandoned by a crashed reclaimer. Returns true when removed. */ +export const clearAbandonedReclaimGuard = (lockPath: string, nowMs = Date.now()): boolean => { + const guard = `${lockPath}.reclaim`; + const age = ageOf(guard, nowMs); + if (age === null || age <= RECLAIM_GUARD_TTL_MS) return false; + try { + fs.rmdirSync(guard); + return true; + } catch { + return false; + } +}; + +export type MetadataLockStatus = { + path: string; + present: boolean; + owner: MetadataLock | null; + state: LockOwnerState["state"] | "absent"; + reason: string; + guard: "absent" | "fresh" | "abandoned"; + /** Exact lock bytes that were classified (for compare-before-remove). */ + raw?: string; + /** Age of the lock file in milliseconds, when present. */ + ageMs?: number | null; +}; + +/** Read-only inspection of a checkout's metadata lock (doctor surface). */ +export const inspectMetadataLock = (root: string, nowMs = Date.now()): MetadataLockStatus => { + const lockPath = lockPathFor(root); + const guardAge = ageOf(`${lockPath}.reclaim`, nowMs); + const guard = + guardAge === null ? "absent" : guardAge > RECLAIM_GUARD_TTL_MS ? "abandoned" : "fresh"; + let raw: string; + try { + raw = fs.readFileSync(lockPath, "utf8"); + } catch { + return { + path: lockPath, + present: false, + owner: null, + state: "absent", + reason: "no lock", + guard, + }; + } + const owner = parseMetadataLockOrNull(raw); + const ageMs = ageOf(lockPath, nowMs); + const verdict = classifyLockOwner(owner, ageMs); + return { path: lockPath, present: true, owner, ...verdict, guard, raw, ageMs }; +}; + +export type ClearLockOutcome = MetadataLockStatus & { + cleared: boolean; + guardCleared: boolean; + /** Why the lock was kept, when it was present and not cleared. */ + skipped?: string; +}; + +/** + * Clear a stale metadata lock (or, with `force`, any lock) and an abandoned + * reclaim guard. Removal holds the same `.reclaim` guard that writers take + * before reclaiming, so no writer can replace the lock between the final + * byte check and the unlink; a fresh guard means a reclaim is already in + * progress and the lock is left alone. + */ +export const clearStaleMetadataLock = ( + root: string, + options: { force?: boolean; nowMs?: number } = {}, +): ClearLockOutcome => { + const nowMs = options.nowMs ?? Date.now(); + const status = inspectMetadataLock(root, nowMs); + const guardCleared = status.guard === "abandoned" && clearAbandonedReclaimGuard(status.path); + const outcome = { ...status, cleared: false, guardCleared }; + if (!status.present) return outcome; + if (status.state !== "stale" && !options.force) return { ...outcome, skipped: status.reason }; + const guard = `${status.path}.reclaim`; + try { + fs.mkdirSync(guard); + } catch { + return { ...outcome, skipped: "a reclaim is in progress" }; + } + try { + const before = fs.readFileSync(status.path, "utf8"); + if (before !== status.raw) return { ...outcome, skipped: "the lock changed" }; + if (!options.force) { + const verdict = classifyLockOwner(parseMetadataLockOrNull(before), ageOf(status.path, nowMs)); + if (verdict.state !== "stale") return { ...outcome, skipped: verdict.reason }; + } + if (fs.readFileSync(status.path, "utf8") !== before) + return { ...outcome, skipped: "the lock changed" }; + fs.rmSync(status.path); + return { ...outcome, cleared: true }; + } catch (error) { + return { ...outcome, skipped: `could not clear: ${String(error)}` }; + } finally { + try { + fs.rmdirSync(guard); + } catch {} + } +}; diff --git a/packages/workit-core/src/core/task-contract.ts b/packages/workit-core/src/core/task-contract.ts index e8fa973c..dcae8504 100644 --- a/packages/workit-core/src/core/task-contract.ts +++ b/packages/workit-core/src/core/task-contract.ts @@ -1048,6 +1048,8 @@ export type ErrorCode = | "capability_unavailable" | "requirements_unsatisfied" | "writer_conflict" + /** Retryable: another live Workit call holds the checkout's metadata lock. */ + | "busy" | "recovery_required" | "storage_error" | "external_outcome_unknown"; diff --git a/packages/workit-core/src/core/task-store.ts b/packages/workit-core/src/core/task-store.ts index 870823e3..f81b4b24 100644 --- a/packages/workit-core/src/core/task-store.ts +++ b/packages/workit-core/src/core/task-store.ts @@ -1,7 +1,6 @@ import { createHash, randomUUID } from "node:crypto"; import { AsyncLocalStorage } from "node:async_hooks"; import * as fs from "node:fs"; -import { hostname } from "node:os"; import path from "node:path"; import * as z from "zod"; import { packageRoot } from "./package-root"; @@ -33,6 +32,24 @@ import { type Utc, type WorkspaceRecord, } from "./task-contract"; +import { + classifyLockOwner, + defaultLockTimeout, + localLockHost, + clearAbandonedReclaimGuard, + parseMetadataLock, + parseMetadataLockOrNull, + processStartOf, + sameMetadataLock, + type MetadataLock, +} from "./store-lock"; + +export type { MetadataLock } from "./store-lock"; +export type TaskStoreOptions = { + /** Total time a mutation retries a lock held by a live writer before `busy` + * (default: `defaultLockTimeout()`, short for in-process hosts). */ + lockTimeoutMs?: number; +}; export type MutationContext = { now: Utc; revision: Revision }; export type TaskMutation = (task: TaskRecord, context: MutationContext) => Result; @@ -72,13 +89,6 @@ export type RecoveryInput = { writer: WorkspaceRecord["writer"], ) => Result; }; -export type MetadataLock = { - pid: number; - processStart: string | null; - host: string; - nonce: string; - externalAction?: true; -}; export type ProcessEvidence = { state: "stopped" | "accounted_for"; pid: number; @@ -96,15 +106,6 @@ const processEvidenceSchema = z .nullable(), }) .strict(); -const metadataLockSchema = z - .object({ - pid: z.number().int().nonnegative().safe(), - processStart: z.string().nullable(), - host: z.string().min(1), - nonce: z.string().min(1), - externalAction: z.literal(true).optional(), - }) - .strict(); /** One task's listing facts, kept in `.workit/index.json` so per-turn host * hooks can find a session's task without parsing every full record. */ export type TaskIndexEntry = { @@ -224,22 +225,6 @@ const isObject = (value: unknown): value is Record => const validId = (value: string): boolean => /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/.test(value); const validDigest = (value: string): boolean => /^[0-9a-f]{64}$/.test(value); -const parseMetadataLock = (raw: string): MetadataLock => { - let value: unknown; - try { - value = JSON.parse(raw); - } catch { - throw Object.assign(new Error("metadata lock is invalid"), { code: "metadata_lock_invalid" }); - } - const parsed = metadataLockSchema.safeParse(value); - if (!parsed.success) - throw Object.assign(new Error("metadata lock is invalid"), { code: "metadata_lock_invalid" }); - return parsed.data; -}; -const sameMetadataLock = (left: unknown, right: MetadataLock): boolean => { - const parsed = metadataLockSchema.safeParse(left); - return parsed.success && canonicalJson(parsed.data) === canonicalJson(right); -}; export const sameDirectoryIdentity = (left: string, right: string): boolean => { if (!path.isAbsolute(left) || !path.isAbsolute(right)) return false; const normalizedLeft = path.resolve(left); @@ -268,13 +253,40 @@ export const sameDirectoryIdentity = (left: string, right: string): boolean => { } }; type LockSnapshot = { raw: string; data: MetadataLock }; +const TRANSIENT_WINDOWS_CODES = new Set(["EPERM", "EACCES", "EBUSY"]); +/** Windows briefly refuses to replace or open a file that another process is + * reading or renaming at that instant. That is contention, not damage: retry + * for about a second before surfacing the error. */ +const retryTransient = (run: () => T): T => { + if (process.platform !== "win32") return run(); + for (let attempt = 0; ; attempt += 1) { + try { + return run(); + } catch (error) { + const code = (error as { code?: unknown } | null)?.code; + if (attempt >= 20 || typeof code !== "string" || !TRANSIENT_WINDOWS_CODES.has(code)) + throw error; + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 5 * (attempt + 1)); + } + } +}; const externalActionLockRoots = new AsyncLocalStorage>(); +/** Roots whose metadata lock this process holds; waiting on them can only time out. */ +const heldInProcess = new Map(); +const holdRoot = (root: string) => heldInProcess.set(root, (heldInProcess.get(root) ?? 0) + 1); +const dropRoot = (root: string) => { + const count = (heldInProcess.get(root) ?? 1) - 1; + if (count > 0) heldInProcess.set(root, count); + else heldInProcess.delete(root); +}; export class TaskStore { readonly root: string; + private readonly lockTimeoutMs: number; - constructor(root: string) { + constructor(root: string, options: TaskStoreOptions = {}) { this.root = fs.existsSync(root) ? fs.realpathSync(root) : path.resolve(root); + this.lockTimeoutMs = options.lockTimeoutMs ?? defaultLockTimeout(); } readTask(taskId: Id): Result { @@ -932,6 +944,7 @@ export class TaskStore { return replaced.ok ? success(value.revision, null, value) : replaced; }, this.recoveryLockOptions(lock.data, evidence.data, recoveryGate), + "recovery_required", ); return result; } catch (error) { @@ -989,24 +1002,62 @@ export class TaskStore { } } + /** + * Mutation lock: a holder that is gone (dead pid, reused pid, or a foreign or + * unreadable lock past its TTL) is reclaimed; a live holder is waited on + * briefly and then reported as retryable `busy`. + */ private metadataLockOptions(): FileLockSyncAcquireOptions { + const inProcess = heldInProcess.has(this.root); return { lockPath: this.lockPath, staleMs: Number.MAX_SAFE_INTEGER, - timeoutMs: 0, - retry: { retries: 0 }, - staleRecovery: "fail-closed", - shouldReclaim: () => false, - parsePayload: parseMetadataLock, + timeoutMs: inProcess ? 0 : this.lockTimeoutMs, + retry: inProcess + ? { retries: 0 } + : { minTimeout: 5, maxTimeout: 100, factor: 1.5, randomize: true }, + staleRecovery: "remove-if-unchanged", + shouldReclaim: ({ payload, nowMs }) => { + let ageMs: number | null = null; + try { + ageMs = nowMs - fs.lstatSync(this.lockPath).mtimeMs; + } catch {} + return classifyLockOwner(payload, ageMs).state === "stale"; + }, + // The library re-checks the bytes before removal, so a lock replaced + // after classification is never deleted. + shouldRemoveStaleLock: () => true, + parsePayload: parseMetadataLockOrNull, payload: () => ({ pid: process.pid, - processStart: this.processStart(process.pid), - host: hostname(), + processStart: processStartOf(process.pid), + host: localLockHost(), nonce: randomUUID(), }), }; } + private acquireMetadataLock( + options: FileLockSyncAcquireOptions, + ): FileLockSyncHandle { + clearAbandonedReclaimGuard(this.lockPath); + // One budget for the whole acquisition: a lost reclaim race retries with + // the remaining time, never a fresh timeout. + const deadline = Date.now() + (options.timeoutMs ?? 0); + while (true) { + try { + return acquireFileLockSync(this.workspacePath, { + ...options, + timeoutMs: Math.max(0, deadline - Date.now()), + }); + } catch (error) { + // Losing a reclaim race to another process is contention, not damage. + const code = (error as { code?: unknown })?.code; + if (code !== "file_lock_stale" || Date.now() >= deadline) throw error; + } + } + } + private externalActionLockOptions(): FileLockSyncAcquireOptions { const options = this.metadataLockOptions(); return { ...options, payload: () => ({ ...options.payload(), externalAction: true }) }; @@ -1032,10 +1083,11 @@ export class TaskStore { } let handle: FileLockSyncHandle; try { - handle = acquireFileLockSync(this.workspacePath, this.externalActionLockOptions()); + handle = this.acquireMetadataLock(this.externalActionLockOptions()); } catch (error) { return this.lockFailure(error); } + holdRoot(this.root); let result: Result; try { if (!handle.verifyStillHeld()) @@ -1072,6 +1124,7 @@ export class TaskStore { path: this.lockPath, }); } + dropRoot(this.root); try { handle.release(); } catch (error) { @@ -1097,6 +1150,9 @@ export class TaskStore { gate: { reclaimed: boolean }, ): FileLockSyncAcquireOptions { const options = this.metadataLockOptions(); + options.timeoutMs = 0; + options.retry = { retries: 0 }; + options.parsePayload = parseMetadataLock; options.staleRecovery = "remove-if-unchanged"; options.shouldReclaim = ({ payload }) => Boolean( @@ -1127,6 +1183,7 @@ export class TaskStore { private withLock( operation: (handle: FileLockSyncHandle) => Result, options: FileLockSyncAcquireOptions = this.metadataLockOptions(), + contention: "busy" | "recovery_required" = "busy", ): Result { try { this.initializeMutationStorage(); @@ -1141,10 +1198,11 @@ export class TaskStore { "metadata lock operation did not produce a result", ); try { - handle = acquireFileLockSync(this.workspacePath, options); + handle = this.acquireMetadataLock(options); } catch (error) { - result = this.lockFailure(error); + result = this.lockFailure(error, contention); } + if (handle) holdRoot(this.root); if (handle) { try { if (!handle.verifyStillHeld()) @@ -1163,6 +1221,7 @@ export class TaskStore { } } if (handle) { + dropRoot(this.root); try { handle.release(); } catch (error) { @@ -1178,16 +1237,33 @@ export class TaskStore { return result; } - private lockFailure(error: unknown): Result { + private lockFailure( + error: unknown, + contention: "busy" | "recovery_required" = "busy", + ): Result { const value = error as { code?: unknown; message?: unknown }; const code = typeof value?.code === "string" ? value.code : ""; - if (code === "EEXIST" || code === "file_lock_timeout") { - const lock = this.readLockSnapshot(); - if (lock.ok && lock.data?.data.externalAction) + if (code === "EEXIST" || code === "file_lock_timeout" || code === "file_lock_stale") { + let lock: MetadataLock | null = null; + try { + lock = parseMetadataLockOrNull(fs.readFileSync(this.lockPath, "utf8")); + } catch {} + if (lock?.externalAction) return failure("writer_conflict", "workspace is reserved by a managed external action", { outcome: "not_started", path: this.lockPath, }); + if (contention === "busy") + return failure( + "busy", + `workspace metadata lock is held by another Workit call${lock ? ` (pid ${lock.pid} on ${lock.host})` : ""}; retry shortly`, + { + outcome: "not_started", + path: this.lockPath, + guidance: + "Retry the same call. If it stays busy, run `workit doctor --fix-lock` to clear a lock left by a dead process.", + }, + ); } const recovery = code === "EEXIST" || @@ -1224,7 +1300,7 @@ export class TaskStore { } finally { fs.closeSync(fd); } - fs.renameSync(temporary, file); + retryTransient(() => fs.renameSync(temporary!, file)); temporary = undefined; this.fsyncDirectory(path.dirname(file)); if (path.dirname(file) === this.tasksDir) this.indexTaskWrite(file, value as TaskRecord); @@ -1283,7 +1359,7 @@ export class TaskStore { }), { mode: 0o600, flag: "wx" }, ); - fs.renameSync(temporary, this.indexPath); + retryTransient(() => fs.renameSync(temporary, this.indexPath)); } catch { try { fs.unlinkSync(temporary); @@ -1321,7 +1397,7 @@ export class TaskStore { } finally { fs.closeSync(fd); } - fs.renameSync(temporary, destination); + retryTransient(() => fs.renameSync(temporary!, destination)); temporary = undefined; this.fsyncDirectory(this.recoveryDir); } finally { @@ -1406,7 +1482,7 @@ export class TaskStore { private readRecord(file: string, schema: { safeParse(value: unknown): any }) { try { - const bytes = fs.readFileSync(file, "utf8"); + const bytes = retryTransient(() => fs.readFileSync(file, "utf8")); return { exists: true, result: this.parseBytes(bytes, schema) }; } catch (error: any) { if (error?.code === "ENOENT") @@ -1446,14 +1522,6 @@ export class TaskStore { return failure("recovery_required", "snapshot does not satisfy its schema"); } - private processStart(pid: number): string | null { - try { - return fs.readFileSync(`/proc/${pid}/stat`, "utf8").split(" ")[21] ?? null; - } catch { - return null; - } - } - private conflict(expected: Revision, actual: Revision): Result { return failure( "revision_conflict", diff --git a/test/artifacts/reliability-report.test.ts b/test/artifacts/reliability-report.test.ts index 82172b75..a70d4f28 100644 --- a/test/artifacts/reliability-report.test.ts +++ b/test/artifacts/reliability-report.test.ts @@ -41,13 +41,14 @@ test("default report aggregates the deterministic candidate and an isolated doct expect(report.candidate.map((c) => c.sha256)).toEqual(packs.map((p) => p.sha256)); // The default env-isolated doctor (node+bun on PATH, no git) is deterministic: // exactly the utility check fails (D11/D13); codex_pin passes (absent). - // Counts include both provider identity checks (pass with no Git remote). + // Counts include both provider identity checks (pass with no Git remote) + // and the workspace_lock check (pass with no metadata lock). expect(report.doctor).toEqual({ ok: false, - passed: 19, + passed: 20, warned: 0, failed: 1, - total: 20, + total: 21, fixes: 1, }); expect(report.logs).toEqual({ files: 0, events: 0 }); @@ -73,13 +74,14 @@ test("report doctor counts are exact against a controlled isolated fixture", () }, }); // node+bun on PATH but no git: exactly the utility check fails; codex_pin passes (absent). - // Counts include both provider identity checks (pass with no Git remote). + // Counts include both provider identity checks (pass with no Git remote) + // and the workspace_lock check (pass with no metadata lock). expect(report.doctor).toEqual({ ok: false, - passed: 19, + passed: 20, warned: 0, failed: 1, - total: 20, + total: 21, fixes: 1, }); } finally { diff --git a/test/workit-cli/doctor.test.ts b/test/workit-cli/doctor.test.ts index 2815391b..9b4a45b0 100644 --- a/test/workit-cli/doctor.test.ts +++ b/test/workit-cli/doctor.test.ts @@ -1,6 +1,7 @@ import { afterAll, expect, test } from "bun:test"; import { spawnSync } from "node:child_process"; -import { mkdirSync, rmSync, writeFileSync } from "node:fs"; +import { existsSync, mkdirSync, readFileSync, rmSync, utimesSync, writeFileSync } from "node:fs"; +import { localLockHost } from "@/packages/workit-core/src/core/store-lock"; import path from "node:path"; import { fileURLToPath } from "node:url"; import type { DoctorReport } from "@/packages/workit-core/src/core/doctor"; @@ -16,11 +17,13 @@ const cliEntry = path.join(repoRoot, "packages/workit-cli/src/index.tsx"); const fixture = makeDoctorFixture(); afterAll(() => fixture.cleanup()); -const runCli = (args: string[], cwd: string) => +const runCli = (args: string[], cwd: string, extraEnv: Record = {}) => spawnSync("bun", [cliEntry, ...args], { cwd, + stdio: ["ignore", "pipe", "pipe"], env: { ...process.env, + ...extraEnv, HOME: fixture.home, WORKFLOW_TOOLKIT_CONFIG: fixture.configDir, WORKFLOW_TOOLKIT_STATE: fixture.stateDir, @@ -64,3 +67,129 @@ test("workit doctor (text) prints per-check lines and no JSON to stdout", () => expect(() => JSON.parse(text.stdout)).toThrow(); expect(text.stdout).toMatch(/stale_pin/); }); + +const deadPid = (): number => + Number( + spawnSync(process.execPath, ["-e", "process.stdout.write(String(process.pid))"], { + encoding: "utf8", + }).stdout, + ); + +test("Given a stale lock left by a dead pid, When workit doctor runs, Then it warns and names --fix-lock", () => { + const lockPath = path.join(fixture.cwd, ".workit", "metadata.lock"); + mkdirSync(path.dirname(lockPath), { recursive: true }); + writeFileSync( + lockPath, + JSON.stringify({ pid: deadPid(), processStart: "1", host: localLockHost(), nonce: "n" }), + ); + try { + const result = runCli(["doctor", "--json"], fixture.cwd); + const report = JSON.parse(result.stdout) as DoctorReport; + const check = report.checks.find((c) => c.id === "workspace_lock"); + expect(check).toMatchObject({ status: "warn", fix: "workit doctor --fix-lock" }); + expect(check?.detail).toContain("is not running"); + expect(existsSync(lockPath)).toBe(true); + } finally { + rmSync(path.join(fixture.cwd, ".workit"), { recursive: true, force: true }); + } +}); + +test("Given a stale lock and an abandoned reclaim guard, When workit doctor --fix-lock runs, Then both are cleared", () => { + const lockPath = path.join(fixture.cwd, ".workit", "metadata.lock"); + mkdirSync(`${lockPath}.reclaim`, { recursive: true }); + utimesSync(`${lockPath}.reclaim`, new Date(0), new Date(0)); + writeFileSync( + lockPath, + JSON.stringify({ pid: deadPid(), processStart: "1", host: localLockHost(), nonce: "n" }), + ); + try { + const result = runCli(["doctor", "--fix-lock"], fixture.cwd); + expect(result.stdout).toContain("fix-lock: cleared stale lock"); + expect(result.stdout).toContain("removed abandoned reclaim guard"); + expect(existsSync(lockPath)).toBe(false); + expect(existsSync(`${lockPath}.reclaim`)).toBe(false); + expect(result.stdout).toMatch(/ok {3}workspace_lock — no metadata lock held/); + } finally { + rmSync(path.join(fixture.cwd, ".workit"), { recursive: true, force: true }); + } +}); + +test("Given a lock held by a live process, When workit doctor --fix-lock runs, Then the lock is kept", () => { + const lockPath = path.join(fixture.cwd, ".workit", "metadata.lock"); + mkdirSync(path.dirname(lockPath), { recursive: true }); + const bytes = JSON.stringify({ + pid: process.pid, + processStart: null, + host: localLockHost(), + nonce: "live", + }); + writeFileSync(lockPath, bytes); + try { + const result = runCli(["doctor", "--json", "--fix-lock"], fixture.cwd); + const report = JSON.parse(result.stdout) as DoctorReport & { fixLock: { cleared: boolean } }; + expect(report.fixLock.cleared).toBe(false); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); + } finally { + rmSync(path.join(fixture.cwd, ".workit"), { recursive: true, force: true }); + } +}); + +test("Given WORKFLOW_WORKSPACE_ROOT points at another checkout, When workit doctor --fix-lock runs, Then it clears that checkout's stale lock", () => { + const other = path.join(fixture.root, "other-workspace"); + const lockPath = path.join(other, ".workit", "metadata.lock"); + mkdirSync(path.dirname(lockPath), { recursive: true }); + writeFileSync( + lockPath, + JSON.stringify({ pid: deadPid(), processStart: "1", host: localLockHost(), nonce: "n" }), + ); + try { + const result = runCli(["doctor", "--fix-lock"], fixture.cwd, { + WORKFLOW_WORKSPACE_ROOT: other, + }); + expect(result.stdout).toContain(`cleared stale lock ${lockPath}`); + expect(existsSync(lockPath)).toBe(false); + } finally { + rmSync(other, { recursive: true, force: true }); + } +}); + +test("Given a lock whose owner cannot be verified, When workit doctor --fix-lock --force runs without --yes or a TTY, Then it refuses and keeps the lock; with --yes it clears it", () => { + const lockPath = path.join(fixture.cwd, ".workit", "metadata.lock"); + mkdirSync(path.dirname(lockPath), { recursive: true }); + const bytes = JSON.stringify({ pid: 1, processStart: null, host: "elsewhere", nonce: "n" }); + writeFileSync(lockPath, bytes); + try { + const refused = runCli(["doctor", "--fix-lock", "--force"], fixture.cwd); + expect(refused.status).toBe(1); + expect(refused.stdout).toContain("held by pid 1 on elsewhere"); + expect(refused.stdout).toContain("refusing without --yes"); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); + const forced = runCli(["doctor", "--fix-lock", "--force", "--yes"], fixture.cwd); + expect(forced.stdout).toContain("fix-lock: cleared lock"); + expect(existsSync(lockPath)).toBe(false); + } finally { + rmSync(path.join(fixture.cwd, ".workit"), { recursive: true, force: true }); + } +}); + +test("Given an unverifiable lock that has blocked writes for over 30s, When workit doctor runs, Then it warns with the exact force command; a fresh one passes", () => { + const lockPath = path.join(fixture.cwd, ".workit", "metadata.lock"); + mkdirSync(path.dirname(lockPath), { recursive: true }); + writeFileSync( + lockPath, + JSON.stringify({ pid: 1, processStart: null, host: "elsewhere", nonce: "n" }), + ); + try { + const fresh = JSON.parse(runCli(["doctor", "--json"], fixture.cwd).stdout) as DoctorReport; + expect(fresh.checks.find((c) => c.id === "workspace_lock")?.status).toBe("pass"); + const minuteAgo = new Date(Date.now() - 60_000); + utimesSync(lockPath, minuteAgo, minuteAgo); + const blocked = JSON.parse(runCli(["doctor", "--json"], fixture.cwd).stdout) as DoctorReport; + expect(blocked.checks.find((c) => c.id === "workspace_lock")).toMatchObject({ + status: "warn", + fix: "workit doctor --fix-lock --force --yes", + }); + } finally { + rmSync(path.join(fixture.cwd, ".workit"), { recursive: true, force: true }); + } +}); diff --git a/test/workit-core/install-scripts.test.ts b/test/workit-core/install-scripts.test.ts index a6d51309..a80d3be2 100644 --- a/test/workit-core/install-scripts.test.ts +++ b/test/workit-core/install-scripts.test.ts @@ -244,6 +244,7 @@ function copyCoreSources(stub: string) { "legacy-ownership.ts", "safe-write.ts", "task-contract.ts", + "store-lock.ts", "runtime-identity.ts", ]) { const src = diff --git a/test/workit-core/store-lock.test.ts b/test/workit-core/store-lock.test.ts new file mode 100644 index 00000000..23334708 --- /dev/null +++ b/test/workit-core/store-lock.test.ts @@ -0,0 +1,359 @@ +import { expect, test } from "bun:test"; +import { spawn, spawnSync } from "node:child_process"; +import { + existsSync, + mkdirSync, + mkdtempSync, + readFileSync, + utimesSync, + writeFileSync, +} from "node:fs"; +import { hostname, tmpdir } from "node:os"; +import { + localLockHost, + parseProcStatStart, + processStartOf, +} from "@/packages/workit-core/src/core/store-lock"; +import { join, resolve } from "node:path"; +import { TaskStore } from "@/packages/workit-core/src/core/task-store"; +import { success, type TaskRecord } from "@/packages/workit-core/src/core/task-contract"; +import { ref, scope } from "./task-fixtures"; + +// Lock reclaim and contention (spec "Stale lock" / "Contention"): a lock held +// by a dead or replaced process is reclaimed automatically, while contention +// between live writers is retryable `busy`, never `recovery_required`. + +const provenance = { + kind: "host_observed" as const, + host: "workit_cli" as const, + session: null, + workerId: null, + receipts: [], +}; +const identity = (task: TaskRecord) => success(task.revision, null, task); + +const startedStore = (options?: { lockTimeoutMs?: number }) => { + const store = new TaskStore(mkdtempSync(join(tmpdir(), "workit-lock-")), options); + const created = store.create({ + expectedWorkspaceRevision: null, + provenance, + intent: { objective: "lock test", scope: scope(), authorityRefs: [ref()] }, + }); + if (!created.ok) throw new Error(created.error); + return { store, task: created.data, lockPath: join(store.root, ".workit", "metadata.lock") }; +}; + +const processStart = processStartOf; + +const deadPid = (): number => { + const child = spawnSync(process.execPath, ["-e", "process.stdout.write(String(process.pid))"], { + encoding: "utf8", + }); + return Number(child.stdout); +}; + +const writeLock = (lockPath: string, payload: Record) => + writeFileSync(lockPath, `${JSON.stringify(payload, null, 2)}\n`); + +test("Given a lock held by a dead pid, When a write runs, Then the lock is reclaimed and the write succeeds", () => { + const { store, task, lockPath } = startedStore(); + writeLock(lockPath, { pid: deadPid(), processStart: "1", host: localLockHost(), nonce: "dead" }); + const result = store.mutateTask(task.id, task.revision, identity); + expect(result.ok).toBe(true); + expect(existsSync(lockPath)).toBe(false); +}); + +test.skipIf(process.platform !== "linux")( + "Given a lock whose pid now belongs to a different process, When a write runs, Then the lock is reclaimed and the write succeeds", + () => { + const { store, task, lockPath } = startedStore(); + writeLock(lockPath, { + pid: process.pid, + processStart: `${processStart(process.pid)}0`, + host: localLockHost(), + nonce: "reused", + }); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); + expect(existsSync(lockPath)).toBe(false); + }, +); + +test("Given a lock from another host older than the TTL, When a write runs, Then the lock is reclaimed and the write succeeds", () => { + const { store, task, lockPath } = startedStore(); + writeLock(lockPath, { pid: 999999, processStart: "old", host: "another-host", nonce: "x" }); + utimesSync(lockPath, new Date(0), new Date(0)); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); +}); + +test("Given a fresh lock from another host, When a write runs, Then it returns busy and keeps the lock", () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 150 }); + const bytes = `${JSON.stringify({ pid: 1, processStart: "x", host: "another-host", nonce: "y" })}\n`; + writeFileSync(lockPath, bytes); + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ + ok: false, + code: "busy", + }); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); +}); + +test("Given a lock held by a live writer, When another write runs, Then it returns retryable busy, not recovery_required, and the lock stays", async () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 200 }); + const holder = spawn("sleep", ["30"], { stdio: "ignore" }); + try { + await new Promise((done) => setTimeout(done, 50)); + const pid = holder.pid!; + writeLock(lockPath, { + pid, + processStart: processStart(pid), + host: localLockHost(), + nonce: "live", + }); + const bytes = readFileSync(lockPath, "utf8"); + const result = store.mutateTask(task.id, task.revision, identity); + expect(result).toMatchObject({ ok: false, code: "busy" }); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); + } finally { + holder.kill("SIGKILL"); + } +}); + +test("Given a live writer that releases its lock within the retry window, When another write runs, Then the write succeeds", async () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 3000 }); + const holder = spawn("sleep", ["30"], { stdio: "ignore" }); + await new Promise((done) => setTimeout(done, 50)); + const pid = holder.pid!; + writeLock(lockPath, { + pid, + processStart: processStart(pid), + host: localLockHost(), + nonce: "brief", + }); + const releaser = spawn( + process.execPath, + ["-e", `setTimeout(() => require("node:fs").rmSync(${JSON.stringify(lockPath)}), 300)`], + { stdio: "ignore" }, + ); + try { + // The synchronous retry loop blocks this event loop, so the release comes + // from another process. + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); + } finally { + holder.kill("SIGKILL"); + releaser.kill("SIGKILL"); + } +}); + +test("Given an abandoned reclaim guard older than the TTL, When a write runs, Then the guard is cleared and the write succeeds", () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 200 }); + mkdirSync(`${lockPath}.reclaim`); + utimesSync(`${lockPath}.reclaim`, new Date(0), new Date(0)); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); + expect(existsSync(`${lockPath}.reclaim`)).toBe(false); +}); + +test("Given a corrupt task record, When a write runs, Then recovery_required still surfaces", () => { + const { store, task } = startedStore(); + const file = join(store.root, ".workit", "tasks", `${task.id}.json`); + writeFileSync(file, "{broken"); + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ + ok: false, + code: "recovery_required", + }); + expect(readFileSync(file, "utf8")).toBe("{broken"); +}); + +const storeModule = resolve(import.meta.dir, "../../packages/workit-core/src/core/task-store.ts"); +const workerScript = (root: string, taskId: string, calls: number) => ` +import { TaskStore } from ${JSON.stringify(storeModule)}; +const store = new TaskStore(${JSON.stringify(root)}); +const codes = {}; +for (let call = 0; call < ${calls}; call += 1) { + let code = "revision_conflict"; + for (let attempt = 0; attempt < 50 && code === "revision_conflict"; attempt += 1) { + const current = store.readTask(${JSON.stringify(taskId)}); + if (!current.ok) { code = current.code; break; } + const result = store.mutateTask(current.data.id, current.data.revision, (task) => ({ + ok: true, revision: task.revision, workspaceRevision: null, data: task, + })); + code = result.ok ? "ok" : result.code; + } + codes[code] = (codes[code] ?? 0) + 1; +} +process.stdout.write(JSON.stringify(codes)); +`; + +test("Given three processes each making 40 writes to one task, When they contend, Then none returns recovery_required", async () => { + const { store, task } = startedStore(); + const runs = await Promise.all( + [0, 1, 2].map( + () => + new Promise>((done, fail) => { + const child = spawn(process.execPath, ["-e", workerScript(store.root, task.id, 40)], { + stdio: ["ignore", "pipe", "pipe"], + }); + let out = ""; + let err = ""; + child.stdout.on("data", (chunk) => (out += chunk)); + child.stderr.on("data", (chunk) => (err += chunk)); + child.on("close", () => { + try { + done(JSON.parse(out)); + } catch { + fail(new Error(`worker output: ${out} ${err}`)); + } + }); + }), + ), + ); + const totals: Record = {}; + for (const run of runs) + for (const [code, count] of Object.entries(run)) totals[code] = (totals[code] ?? 0) + count; + expect(totals.recovery_required ?? 0).toBe(0); + // Under heavy machine load a writer may exhaust its retry window: that is + // the retryable `busy`, never anything else. + expect(Object.keys(totals).filter((code) => code !== "ok" && code !== "busy")).toEqual([]); + expect((totals.ok ?? 0) + (totals.busy ?? 0)).toBe(120); + expect(totals.ok ?? 0).toBeGreaterThan(100); +}, 60_000); + +const storeModule2 = resolve(import.meta.dir, "../../packages/workit-core/src/core/store-lock.ts"); + +test("Given a lock from a container that shares the hostname but not the pid namespace, When a write runs, Then its pid is not checked here and the write is busy", () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 150 }); + // pid 1 is alive on this host too; the namespace suffix marks it as foreign. + const bytes = `${JSON.stringify({ + pid: 1, + processStart: "1", + host: `${hostname()}#4026599999:other-boot`, + nonce: "container", + })}\n`; + writeFileSync(lockPath, bytes); + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ + ok: false, + code: "busy", + }); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); + // Past the foreign TTL the same lock is reclaimable. + utimesSync(lockPath, new Date(0), new Date(0)); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); +}); + +test.skipIf(process.platform !== "linux")( + "Given a lock written by an older Workit without namespace identity, When a write runs, Then it is treated as foreign until its TTL", + () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 150 }); + writeLock(lockPath, { pid: deadPid(), processStart: null, host: hostname(), nonce: "legacy" }); + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ code: "busy" }); + utimesSync(lockPath, new Date(0), new Date(0)); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); + }, +); + +test("Given a process name with spaces and parentheses, When /proc stat is parsed, Then the start time is read after the last parenthesis", () => { + const fields = Array.from({ length: 30 }, (_, index) => String(index + 3)); + expect(parseProcStatStart(`4242 (we ird) (name) ${fields.join(" ")}`)).toBe("22"); +}); + +test("Given an in-process host with the default budget, When a live holder keeps the lock, Then busy returns within the short budget", async () => { + const { store, task, lockPath } = startedStore(); + const fresh = new TaskStore(store.root); + const holder = spawn("sleep", ["30"], { stdio: "ignore" }); + try { + await new Promise((done) => setTimeout(done, 50)); + const pid = holder.pid!; + writeLock(lockPath, { + pid, + processStart: processStart(pid), + host: localLockHost(), + nonce: "x", + }); + const started = performance.now(); + expect(fresh.mutateTask(task.id, task.revision, identity)).toMatchObject({ code: "busy" }); + expect(performance.now() - started).toBeLessThan(1_000); + } finally { + holder.kill("SIGKILL"); + } +}); + +test("Given doctor --fix-lock is preempted while a writer reclaims the same stale lock, Then the two never hold the lock at once", async () => { + const { store, task, lockPath } = startedStore(); + writeLock(lockPath, { pid: deadPid(), processStart: null, host: localLockHost(), nonce: "x" }); + const run = (script: string) => + new Promise((done) => { + const child = spawn(process.execPath, ["-e", script], { stdio: ["ignore", "pipe", "pipe"] }); + let out = ""; + child.stdout.on("data", (chunk) => (out += chunk)); + child.on("close", () => done(out)); + }); + // The doctor pauses 400 ms right before removing the lock (models preemption). + const doctor = run(` + import fs from "node:fs"; + const rm = fs.rmSync; + fs.rmSync = (p, ...rest) => { + if (String(p).endsWith("metadata.lock")) Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 400); + return rm(p, ...rest); + }; + const { clearStaleMetadataLock } = await import(${JSON.stringify(storeModule2)}); + console.log(JSON.stringify(clearStaleMetadataLock(${JSON.stringify(store.root)}))); + `); + const writer = (label: string, holdMs: number) => ` + const { TaskStore } = await import(${JSON.stringify(storeModule)}); + const s = new TaskStore(${JSON.stringify(store.root)}, { lockTimeoutMs: 5000 }); + const t = s.readTask(${JSON.stringify(task.id)}); + let span = [0, 0]; + const r = s.mutateTask(t.data.id, t.data.revision, (x) => { + span[0] = Date.now(); + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, ${holdMs}); + span[1] = Date.now(); + return { ok: true, revision: x.revision, workspaceRevision: null, data: x }; + }); + console.log(JSON.stringify({ label: ${JSON.stringify(label)}, code: r.ok ? "ok" : r.code, span })); + `; + await new Promise((done) => setTimeout(done, 150)); + const first = run(writer("W", 800)); + await new Promise((done) => setTimeout(done, 500)); + const second = run(writer("X", 100)); + const [doctorOut, w, x] = await Promise.all([doctor, first, second]); + const spans = [w, x].map((out) => JSON.parse(out.trim().split("\n").at(-1)!)); + expect(JSON.parse(doctorOut.trim())).toMatchObject({ cleared: true }); + // Revision conflicts are fine; overlapping critical sections are not. + const held = spans.filter((item) => item.span[0] > 0).toSorted((a, b) => a.span[0] - b.span[0]); + for (let index = 1; index < held.length; index += 1) + expect(held[index].span[0]).toBeGreaterThanOrEqual(held[index - 1].span[1]); + expect(spans.every((item) => item.code === "ok" || item.code === "revision_conflict")).toBe(true); +}, 30_000); + +// Linux only: where there is no pid-namespace identity (macOS, Windows) a +// plain-hostname lock is this host's own format and pid checks apply. +test.skipIf(!localLockHost().includes("#"))( + "Given a plain-hostname lock naming live pid 1 with a mismatched start time (a container's lock), When a write runs, Then it stays busy and is not reclaimed", + () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 150 }); + // Pre-namespace locks carried only hostname(); judging pid 1 against this + // host's process table would steal a live container's lock. + const bytes = `${JSON.stringify({ pid: 1, processStart: "1", host: hostname(), nonce: "c" })}\n`; + writeFileSync(lockPath, bytes); + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ + ok: false, + code: "busy", + }); + expect(readFileSync(lockPath, "utf8")).toBe(bytes); + }, +); + +test.skipIf(!localLockHost().includes("#"))( + "Given a lock from the same host and pid namespace but an earlier boot, When a write runs, Then it is reclaimed immediately", + () => { + const { store, task, lockPath } = startedStore({ lockTimeoutMs: 150 }); + const [name, space] = localLockHost().split("#"); + const [namespace] = space.split(":"); + writeLock(lockPath, { + pid: process.pid, + processStart: processStart(process.pid), + host: `${name}#${namespace}:00000000-0000-0000-0000-000000000000`, + nonce: "before-reboot", + }); + expect(store.mutateTask(task.id, task.revision, identity).ok).toBe(true); + expect(existsSync(lockPath)).toBe(false); + }, +); diff --git a/test/workit-core/task-store.test.ts b/test/workit-core/task-store.test.ts index e32119ea..7fc98c33 100644 --- a/test/workit-core/task-store.test.ts +++ b/test/workit-core/task-store.test.ts @@ -231,7 +231,7 @@ test("corrupt current bytes are reported without replacement", () => { expect(readFileSync(file, "utf8")).toBe("{broken"); }); -test("unsupported snapshots and leftover locks stay inspectable", () => { +test("unsupported snapshots stay inspectable and a leftover stale lock no longer blocks writes", () => { const { store, task } = startedStore(); const workit = join(store.root, ".workit"); const workspaceFile = join(workit, "workspace.json"); @@ -249,11 +249,11 @@ test("unsupported snapshots and leftover locks stay inspectable", () => { }); writeFileSync(lockPath, lockBytes); utimesSync(lockPath, new Date(0), new Date(0)); - expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ - ok: false, - code: "recovery_required", - }); + expect(store.readTask(task.id)).toMatchObject({ ok: true }); expect(readFileSync(lockPath, "utf8")).toBe(lockBytes); + // A foreign-host lock past its TTL has no provable owner: the write reclaims it. + expect(store.mutateTask(task.id, task.revision, identity)).toMatchObject({ ok: true }); + expect(existsSync(lockPath)).toBe(false); }); test("recovery restores validated bytes with a fresh revision", () => {