diff --git a/AGENTS.md b/AGENTS.md index 24d1ae16..9284052f 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. 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 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. `.workit/recovery/` keeps at most three copies per record; `workit gc` prunes older leftovers. `state.recover` is not advertised to hosts (no shipped host supplies native recovery authority); the engine path remains for embedders that do. 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 dd88c2da..c9aff1c4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- `.workit/recovery/` no longer grows without bound: each task or workspace + record keeps its newest three recovery copies. `workit gc` (`--dry-run`, + `--json`) prunes copies left by older versions, removes stale temp files, and + collapses duplicate stored candidates in paused tasks (closed tasks are never + rewritten); `--dry-run` is read-only. +- `state.recover` is no longer advertised in host tool schemas or the CLI: it + requires native recovery authority that no shipped host supplies, so it could + only return `permission_denied`. - 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 diff --git a/README.md b/README.md index da0a81d5..fb45b95e 100644 --- a/README.md +++ b/README.md @@ -205,6 +205,7 @@ 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 gc [--dry-run] # prune .workit/recovery to the newest 3 copies per record workit [--payload ] [--task ] [--confirm] [--json] workit action --payload [--preview] [--confirm] [--json] workit handoff --task [--json] @@ -417,7 +418,12 @@ 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 +explicit escape hatch for a lock whose owner cannot be verified. Each snapshot +replacement keeps a copy of the previous bytes in `.workit/recovery/`, capped at +the newest three per task or workspace record; `workit gc` prunes copies left by +older versions, removes stale temp files, and collapses duplicate stored +candidates in paused tasks (closed tasks are never rewritten). It never deletes +the live task or workspace records, and `--dry-run` writes nothing. 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/task.ts b/packages/workit-cli/src/task.ts index f7ae0549..d295ffd8 100644 --- a/packages/workit-cli/src/task.ts +++ b/packages/workit-cli/src/task.ts @@ -58,7 +58,8 @@ export const TASK_ACTIONS = { decision: ["record", "revoke"], worker: ["assign", "report", "cancel"], writer: ["acquire", "release"], - state: ["export", "import", "recover"], + // state.recover is not exposed: no shipped host supplies native recovery authority. + state: ["export", "import"], } as const satisfies Record; type Stream = { write: (chunk: string) => void }; @@ -336,7 +337,6 @@ const CONSENT_ACTIONS = new Set([ "writer.acquire", "writer.release", "state.import", - "state.recover", ]); const needsConsent = (parsed: Parsed): boolean => diff --git a/packages/workit-cli/src/verbs/gc.ts b/packages/workit-cli/src/verbs/gc.ts new file mode 100644 index 00000000..d197de0c --- /dev/null +++ b/packages/workit-cli/src/verbs/gc.ts @@ -0,0 +1,37 @@ +// `workit gc`: bounded recovery state for the workspace root's .workit +// (WORKFLOW_WORKSPACE_ROOT, then --cwd/cwd). `--dry-run` is read-only. +import { TaskStore } from "@brainervirus/workit-core/src/core/task-store"; +import { emit, fail, ok, type Io } from "../output"; +import { workspaceRootFor } from "../task"; + +export async function run(argv: string[], io: Io): Promise { + const result = new TaskStore(workspaceRootFor({ cwd: io.cwd })).collectGarbage({ + dryRun: argv.includes("--dry-run"), + }); + if (!result.ok) + return emit( + io, + fail(result.code === "busy" ? "busy" : "failed", `${result.code}: ${result.error}`, { + data: { code: result.code, details: result.details }, + ...(result.code === "busy" ? { unblock: "retry `workit gc`" } : {}), + }), + ); + return emit(io, ok(result.data), (data) => { + const verb = data.dryRun ? "would remove" : "removed"; + const megabytes = (data.recovery.removedBytes / 1_048_576).toFixed(1); + const { candidates } = data; + return [ + `workit gc${data.dryRun ? " (dry run)" : ""}`, + `recovery: ${verb} ${data.recovery.removed} copies (${megabytes} MB), kept ${data.recovery.kept}`, + `temporary files: ${verb} ${data.temporary.removed}`, + `candidates: ${verb} ${candidates.removed} duplicates in ${candidates.tasks.length} tasks` + + (candidates.skippedActive.length + ? `; skipped active ${candidates.skippedActive.join(", ")}` + : "") + + (candidates.skippedClosed.length + ? `; skipped closed ${candidates.skippedClosed.join(", ")}` + : "") + + (candidates.failed.length ? `; failed ${candidates.failed.join(", ")}` : ""), + ]; + }); +} diff --git a/packages/workit-cli/src/verbs/registry.ts b/packages/workit-cli/src/verbs/registry.ts index 1a11daee..4d8cb147 100644 --- a/packages/workit-cli/src/verbs/registry.ts +++ b/packages/workit-cli/src/verbs/registry.ts @@ -71,6 +71,14 @@ export const VERBS: readonly VerbEntry[] = [ "Verify the offline installation health (--fix-lock clears a stale .workit metadata lock)", load: () => import("./doctor"), }, + { + name: "gc", + group: "setup", + usage: "workit gc [--dry-run] [--json]", + summary: + "Prune .workit/recovery copies beyond the cap and dedupe stored candidates in paused tasks", + load: () => import("./gc"), + }, { name: "uninstall", group: "setup", diff --git a/packages/workit-core/src/core.ts b/packages/workit-core/src/core.ts index 54e02505..aa7d5370 100644 --- a/packages/workit-core/src/core.ts +++ b/packages/workit-core/src/core.ts @@ -7,6 +7,7 @@ export { POLICY_VERSION, OPERATION_FAMILIES, operationSchemas, + advertisedOperationSchemas, operationJsonSchema, boundedOperationJsonSchema, OPERATION_SCHEMA_DEPTH, diff --git a/packages/workit-core/src/core/task-contract.ts b/packages/workit-core/src/core/task-contract.ts index 0d30c58f..bc59029f 100644 --- a/packages/workit-core/src/core/task-contract.ts +++ b/packages/workit-core/src/core/task-contract.ts @@ -1131,6 +1131,19 @@ export const operationSchemas = { state: z.discriminatedUnion("action", Object.values(stateOperations) as any), } as const; export type OperationRequest = z.infer<(typeof operationSchemas)[OperationFamily]>; + +/** + * Schemas advertised to hosts. `state.recover` needs host-supplied native + * recovery authority (`OperationContext.nativeRecovery`), which no shipped + * host provides, so advertising it only sends agents into a guaranteed + * permission_denied. parseOperation still accepts it for embedders that do + * supply that authority. + */ +const { recover: _unadvertisedRecover, ...advertisedStateOperations } = stateOperations; +export const advertisedOperationSchemas = { + ...operationSchemas, + state: z.discriminatedUnion("action", Object.values(advertisedStateOperations) as any), +} as const; export type TaskStartRequest = z.infer; const compiledOperationSchemas = Object.fromEntries( @@ -1234,7 +1247,7 @@ export function parseOperation(family: OperationFamily, input: unknown): Result< } export function operationJsonSchema(family: OperationFamily): z.core.JSONSchema.BaseSchema { - return z.toJSONSchema(operationSchemas[family], { target: "draft-2020-12" }); + return z.toJSONSchema(advertisedOperationSchemas[family], { target: "draft-2020-12" }); } /** diff --git a/packages/workit-core/src/core/task-store.ts b/packages/workit-core/src/core/task-store.ts index 08dbab8d..8f2853b7 100644 --- a/packages/workit-core/src/core/task-store.ts +++ b/packages/workit-core/src/core/task-store.ts @@ -46,6 +46,23 @@ import { } from "./store-lock"; export type { MetadataLock } from "./store-lock"; +/** Recovery copies kept per record (task or workspace); older copies are pruned. */ +export const RECOVERY_COPIES_PER_RECORD = 3; +/** A temp file older than this was left by a crashed writer. */ +const STALE_TEMPORARY_MS = 60 * 60_000; +const RECOVERY_NAME = /^(task|workspace)\.([^.]+)\.([0-9a-f]{64})\.json$/; +export type GarbageReport = { + dryRun: boolean; + recovery: { removed: number; removedBytes: number; kept: number }; + temporary: { removed: number }; + candidates: { + removed: number; + tasks: Id[]; + skippedActive: Id[]; + skippedClosed: Id[]; + failed: Id[]; + }; +}; export type TaskStoreOptions = { /** Total time a mutation retries a lock held by a live writer before `busy` * (default: `defaultLockTimeout()`, short for in-process hosts). */ @@ -258,15 +275,18 @@ 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 isTransientWindowsError = (error: unknown): boolean => { + if (process.platform !== "win32") return false; + const code = (error as { code?: unknown } | null)?.code; + return typeof code === "string" && TRANSIENT_WINDOWS_CODES.has(code); +}; 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; + if (attempt >= 20 || !isTransientWindowsError(error)) throw error; Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 5 * (attempt + 1)); } } @@ -281,6 +301,12 @@ const dropRoot = (root: string) => { else heldInProcess.delete(root); }; +/** Drop earlier copies of a repeated candidate ID; content is identical by ID. */ +const dedupeCandidates = (candidates: TaskRecord["candidates"]): TaskRecord["candidates"] => { + const last = new Map(candidates.map((candidate, index) => [candidate.id, index])); + return candidates.filter((candidate, index) => last.get(candidate.id) === index); +}; + export class TaskStore { readonly root: string; private readonly lockTimeoutMs: number; @@ -805,7 +831,7 @@ export class TaskStore { try { const candidates: RecoveryCandidate[] = []; for (const name of fs.readdirSync(this.recoveryDir)) { - const match = /^(task|workspace)\.([^.]+)\.([0-9a-f]{64})\.json$/.exec(name); + const match = RECOVERY_NAME.exec(name); if (match) candidates.push({ target: match[1] as "task" | "workspace", @@ -1056,9 +1082,23 @@ export class TaskStore { timeoutMs: Math.max(0, deadline - Date.now()), }); } catch (error) { - // Losing a reclaim race to another process is contention, not damage. + // Losing a reclaim race to another process, or Windows refusing the + // lock file while another process opens or deletes it, is contention. const code = (error as { code?: unknown })?.code; - if (code !== "file_lock_stale" || Date.now() >= deadline) throw error; + if ( + (code !== "file_lock_stale" && !isTransientWindowsError(error)) || + Date.now() >= deadline + ) + throw error; + // Windows also denies the create while a just-released lock is still + // pending delete, and that entry is already invisible to lstat. So a + // denial with no visible holder is told apart by probing whether the + // directory accepts a new file: if it does not, this is a permission + // problem (read-only attribute, ACL) and fails fast. + if (code !== "file_lock_stale" && !this.lockHolderPresent() && !this.lockDirWritable()) + throw error; + if (code !== "file_lock_stale") + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 5); } } } @@ -1131,7 +1171,7 @@ export class TaskStore { } dropRoot(this.root); try { - handle.release(); + retryTransient(() => handle.release()); } catch (error) { const released = this.lockFailure(error); const releaseError = released.ok ? "metadata lock release failed" : released.error; @@ -1228,7 +1268,7 @@ export class TaskStore { if (handle) { dropRoot(this.root); try { - handle.release(); + retryTransient(() => handle.release()); } catch (error) { const releaseFailure = this.lockFailure(error); const releaseError = releaseFailure.ok @@ -1242,13 +1282,52 @@ export class TaskStore { return result; } + /** Whether the lock's directory accepts a new file right now. */ + private lockDirWritable(): boolean { + const probe = `${this.lockPath}.${process.pid}.${randomUUID()}.probe`; + try { + fs.closeSync(fs.openSync(probe, "wx", 0o600)); + } catch { + return false; + } + try { + fs.rmSync(probe, { force: true }); + } catch {} + return true; + } + + /** A lock file (or an entry Windows is still tearing down) exists. */ + private lockHolderPresent(): boolean { + try { + fs.lstatSync(this.lockPath); + return true; + } catch (error) { + return (error as { code?: unknown } | null)?.code !== "ENOENT"; + } + } + 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" || code === "file_lock_stale") { + // A sharing-class denial is contention only when someone holds the lock. + if (isTransientWindowsError(error) && !this.lockHolderPresent() && !this.lockDirWritable()) + return failure( + "storage_error", + `cannot create the workspace metadata lock (${code}); no other Workit call holds it`, + { + path: this.lockPath, + guidance: `Check permissions and the read-only attribute on ${path.dirname(this.lockPath)}.`, + }, + ); + if ( + code === "EEXIST" || + code === "file_lock_timeout" || + code === "file_lock_stale" || + isTransientWindowsError(error) + ) { let lock: MetadataLock | null = null; try { lock = parseMetadataLockOrNull(fs.readFileSync(this.lockPath, "utf8")); @@ -1285,9 +1364,16 @@ export class TaskStore { } private initializeMutationStorage() { - fs.mkdirSync(this.tasksDir, { recursive: true }); - fs.mkdirSync(this.recoveryDir, { recursive: true }); - fs.writeFileSync(this.gitignorePath, "*\n"); + retryTransient(() => fs.mkdirSync(this.tasksDir, { recursive: true })); + retryTransient(() => fs.mkdirSync(this.recoveryDir, { recursive: true })); + // Rewriting an unchanged .gitignore on every call makes concurrent + // writers collide on it (Windows sharing violations); write it only when + // it is missing or different. + let current: string | null = null; + try { + current = retryTransient(() => fs.readFileSync(this.gitignorePath, "utf8")); + } catch {} + if (current !== "*\n") retryTransient(() => fs.writeFileSync(this.gitignorePath, "*\n")); } private replaceSnapshot(file: string, value: unknown, previous: unknown): Result { @@ -1389,8 +1475,13 @@ export class TaskStore { const id = target === "task" ? path.basename(file, ".json") : "workspace"; const destination = path.join(this.recoveryDir, `${target}.${id}.${digestBytes(bytes)}.json`); if (fs.existsSync(destination)) { - if (digestBytes(fs.readFileSync(destination)) === digestBytes(bytes)) return; - throw new Error("recovery copy already exists with different bytes"); + if (digestBytes(retryTransient(() => fs.readFileSync(destination))) !== digestBytes(bytes)) + throw new Error("recovery copy already exists with different bytes"); + // Re-saved bytes are the newest copy again for pruning purposes. + const stamp = new Date(); + retryTransient(() => fs.utimesSync(destination, stamp, stamp)); + this.pruneRecovery(`${target}.${id}.`, destination); + return; } let temporary: string | undefined; try { @@ -1405,6 +1496,7 @@ export class TaskStore { retryTransient(() => fs.renameSync(temporary!, destination)); temporary = undefined; this.fsyncDirectory(this.recoveryDir); + this.pruneRecovery(`${target}.${id}.`, destination); } finally { if (temporary) try { @@ -1413,6 +1505,144 @@ export class TaskStore { } } + /** + * Keep the newest RECOVERY_COPIES_PER_RECORD copies of one record (always + * including `keep`, the copy just written). Pruning is best-effort: a + * failure never fails the mutation that triggered it. + */ + private pruneRecovery(prefix: string, keep: string) { + try { + const copies = this.recoveryCopies(prefix); + const surplus = this.surplusCopies(copies, path.basename(keep)); + for (const copy of surplus) fs.rmSync(copy.path, { force: true }); + } catch {} + } + + private recoveryCopies( + prefix = "", + ): { name: string; path: string; group: string; mtimeNs: bigint }[] { + if (!fs.existsSync(this.recoveryDir)) return []; + const copies = []; + for (const name of fs.readdirSync(this.recoveryDir)) { + if (!name.startsWith(prefix)) continue; + const match = RECOVERY_NAME.exec(name); + if (!match) continue; + const file = path.join(this.recoveryDir, name); + try { + const stat = fs.lstatSync(file, { bigint: true }); + if (!stat.isFile()) continue; + copies.push({ name, path: file, group: `${match[1]}.${match[2]}`, mtimeNs: stat.mtimeNs }); + } catch {} + } + return copies; + } + + /** Copies of one record beyond the cap, oldest first out; `keep` is never surplus. */ + private surplusCopies(copies: T[], keep?: string) { + const ordered = copies.toSorted((left, right) => + left.name === keep + ? -1 + : right.name === keep + ? 1 + : left.mtimeNs === right.mtimeNs + ? left.name.localeCompare(right.name) + : left.mtimeNs > right.mtimeNs + ? -1 + : 1, + ); + return ordered.slice(RECOVERY_COPIES_PER_RECORD); + } + + /** + * `workit gc`: prune recovery copies beyond the per-record cap, remove temp + * files left by crashed writers, and collapse duplicate stored candidates in + * paused tasks. Live records (tasks/*.json, workspace.json) are never + * deleted; candidate dedupe keeps every candidate ID and the latest position. + * Closed tasks are history and are never rewritten. A dry run is read-only: + * no lock, no directory creation, no .gitignore rewrite. + */ + collectGarbage(options: { dryRun?: boolean } = {}): Result { + const dryRun = options.dryRun === true; + const report: GarbageReport = { + dryRun, + recovery: { removed: 0, removedBytes: 0, kept: 0 }, + temporary: { removed: 0 }, + candidates: { removed: 0, tasks: [], skippedActive: [], skippedClosed: [], failed: [] }, + }; + if (!fs.existsSync(this.workitDir)) return success(null, null, report); + const sweep = (): Result => { + try { + const groups = new Map>(); + for (const copy of this.recoveryCopies()) + groups.set(copy.group, [...(groups.get(copy.group) ?? []), copy]); + for (const copies of groups.values()) { + const surplus = this.surplusCopies(copies); + report.recovery.kept += copies.length - surplus.length; + for (const copy of surplus) { + report.recovery.removed += 1; + try { + report.recovery.removedBytes += fs.statSync(copy.path).size; + } catch {} + if (!dryRun) fs.rmSync(copy.path, { force: true }); + } + } + const nowMs = Date.now(); + for (const directory of [this.recoveryDir, this.tasksDir, this.workitDir]) { + if (!fs.existsSync(directory)) continue; + for (const name of fs.readdirSync(directory)) { + if (!name.endsWith(".tmp") && !name.endsWith(".probe")) continue; + const file = path.join(directory, name); + const stat = fs.lstatSync(file); + if (!stat.isFile() || nowMs - stat.mtimeMs <= STALE_TEMPORARY_MS) continue; + report.temporary.removed += 1; + if (!dryRun) fs.rmSync(file, { force: true }); + } + } + return success(null, null, null); + } catch (error) { + return failure("storage_error", `garbage collection failed: ${String(error)}`, { + path: this.recoveryDir, + }); + } + }; + // withLock initializes storage (mkdir, .gitignore); a dry run must not. + const pruned = dryRun ? sweep() : this.withLock(sweep); + if (!pruned.ok) return pruned; + const tasks = this.listTasks(); + if (!tasks.ok) return tasks; + for (const task of tasks.data) { + const deduped = dedupeCandidates(task.candidates); + const removed = task.candidates.length - deduped.length; + if (removed === 0) continue; + // An active task belongs to a live session; rewriting it would bump the + // revision under that session's feet. + if (task.status === "active") { + report.candidates.skippedActive.push(task.id); + continue; + } + // Closed records are immutable history (docs/workit-v1/contracts.md). + if (task.status === "closed") { + report.candidates.skippedClosed.push(task.id); + continue; + } + if (!dryRun) { + const written = this.mutateTask(task.id, task.revision, (current) => + success(current.revision, null, { + ...current, + candidates: dedupeCandidates(current.candidates), + }), + ); + if (!written.ok) { + report.candidates.failed.push(task.id); + continue; + } + } + report.candidates.removed += removed; + report.candidates.tasks.push(task.id); + } + return success(null, null, report); + } + private findRecovery( target: "task" | "workspace", taskId: Id | null, @@ -1479,7 +1709,7 @@ export class TaskStore { private snapshotBytes(file: string): Buffer | null { try { - return fs.readFileSync(file); + return retryTransient(() => fs.readFileSync(file)); } catch { return null; } diff --git a/packages/workit-opencode/src/tools/workit.ts b/packages/workit-opencode/src/tools/workit.ts index 9f752bb6..b8e8eecf 100644 --- a/packages/workit-opencode/src/tools/workit.ts +++ b/packages/workit-opencode/src/tools/workit.ts @@ -4,7 +4,7 @@ import { TaskStore, canonicalJson, failure, - operationSchemas, + advertisedOperationSchemas, OPERATION_SCHEMA_DEPTH, canonicalFieldsDescription, parseOperation, @@ -413,7 +413,7 @@ const boundedSchema = (schema: any, depth: number, field: string): any => { const operationShapeFor = (family: OperationFamily): Record => { const options = ( - operationSchemas[family] as unknown as { + advertisedOperationSchemas[family] as unknown as { options: Array<{ shape: Record }>; } ).options; diff --git a/test/workit-cli/router.test.ts b/test/workit-cli/router.test.ts index 2bf443a9..d4d33e27 100644 --- a/test/workit-cli/router.test.ts +++ b/test/workit-cli/router.test.ts @@ -315,6 +315,7 @@ const JSON_ERROR_PATHS: Record = { upgrade: [["--bogus"]], launch: [[], ["nohost"]], doctor: [["--fix-lock", "--force"]], + gc: [["--dry-run"]], uninstall: [[]], cutover: [[], ["bogus"]], task: [[], ["bogus"]], diff --git a/test/workit-cli/task-commands.test.ts b/test/workit-cli/task-commands.test.ts index 887231c7..4d3b6a48 100644 --- a/test/workit-cli/task-commands.test.ts +++ b/test/workit-cli/task-commands.test.ts @@ -50,7 +50,7 @@ test("the CLI exposes exactly the eight families and 24 actions", () => { "writer", "state", ]); - expect(Object.values(TASK_ACTIONS).flat()).toHaveLength(24); + expect(Object.values(TASK_ACTIONS).flat()).toHaveLength(23); }); test("CLI external action previews its exact descriptor and refuses headless mutation", async () => { diff --git a/test/workit-core/recovery-gc.test.ts b/test/workit-core/recovery-gc.test.ts new file mode 100644 index 00000000..8976f327 --- /dev/null +++ b/test/workit-core/recovery-gc.test.ts @@ -0,0 +1,271 @@ +import { expect, setDefaultTimeout, test } from "bun:test"; +import { spawnSync } from "node:child_process"; +import { + existsSync, + mkdtempSync, + rmSync, + readdirSync, + readFileSync, + statSync, + utimesSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { join, resolve } from "node:path"; +import { RECOVERY_COPIES_PER_RECORD, TaskStore } from "@/packages/workit-core/src/core/task-store"; +import { + boundedOperationJsonSchema, + candidateDigest, + operationJsonSchema, + parseOperation, + sha256, + success, + type Candidate, + type TaskRecord, +} from "@/packages/workit-core/src/core/task-contract"; +import { ref, scope } from "./task-fixtures"; + +// Each test makes several fsync'd store writes; a cold windows-latest runner +// took 8 s for one of them, past bun's 5 s default. +setDefaultTimeout(30_000); + +// Bounded recovery (spec "Bounded state"): recovery copies are capped per +// record, `workit gc` prunes what older versions left behind, and nothing a +// current read needs is ever deleted. + +const provenance = { + kind: "host_observed" as const, + host: "workit_cli" as const, + session: null, + workerId: null, + receipts: [], +}; + +const startedStore = () => { + const store = new TaskStore(mkdtempSync(join(tmpdir(), "workit-gc-"))); + const created = store.create({ + expectedWorkspaceRevision: null, + provenance, + intent: { objective: "gc test", scope: scope(), authorityRefs: [ref()] }, + }); + if (!created.ok) throw new Error(created.error); + return { store, task: created.data }; +}; + +const recoveryDir = (store: TaskStore) => join(store.root, ".workit", "recovery"); +const copiesFor = (store: TaskStore, prefix: string) => + readdirSync(recoveryDir(store)).filter((name) => name.startsWith(prefix)); + +const touch = (task: TaskRecord, summary: string) => + success(task.revision, null, { ...task, progress: { ...task.progress, summary } }); + +const writesLeaveBoundedCopies = (writes: number) => { + const { store, task } = startedStore(); + const file = join(store.root, ".workit", "tasks", `${task.id}.json`); + let revision = task.revision; + let previous = ""; + for (let write = 0; write < writes; write += 1) { + previous = readFileSync(file, "utf8"); + const result = store.mutateTask(task.id, revision, (current) => touch(current, `w${write}`)); + if (!result.ok) throw new Error(result.error); + revision = result.data.revision; + } + const copies = copiesFor(store, `task.${task.id}.`); + expect(copies.length).toBeLessThanOrEqual(RECOVERY_COPIES_PER_RECORD); + expect(copies).toContain(`task.${task.id}.${sha256(previous)}.json`); +}; + +// 20 writes already exceed the cap several times over; the 1,000-write +// version is an opt-in soak (WORKIT_SOAK=1) because it is fsync-bound. +test( + `Given 20 writes to one task, Then at most ${RECOVERY_COPIES_PER_RECORD} recovery copies remain and the latest prior bytes are kept`, + () => writesLeaveBoundedCopies(20), + 60_000, +); + +test.skipIf(process.env.WORKIT_SOAK !== "1")( + `Given 1,000 writes to one task (soak), Then at most ${RECOVERY_COPIES_PER_RECORD} recovery copies remain`, + () => writesLeaveBoundedCopies(1000), + 300_000, +); + +const seedStaleCopies = (store: TaskStore, prefix: string, count: number) => { + const names: string[] = []; + for (let index = 0; index < count; index += 1) { + const bytes = `{"stale":${index}}\n`; + const name = `${prefix}${sha256(bytes)}.json`; + writeFileSync(join(recoveryDir(store), name), bytes); + // Oldest first: index 0 is the oldest copy. + const when = new Date(Date.UTC(2026, 0, 1, 0, 0, index)); + utimesSync(join(recoveryDir(store), name), when, when); + names.push(name); + } + return names; +}; + +test("Given a recovery dir with N stale copies per record, When gc runs, Then each record keeps only the newest copies up to the cap", () => { + const { store, task } = startedStore(); + const taskCopies = seedStaleCopies(store, `task.${task.id}.`, 40); + const workspaceCopies = seedStaleCopies(store, "workspace.workspace.", 25); + const before = readdirSync(recoveryDir(store)).length; + const result = store.collectGarbage(); + expect(result.ok).toBe(true); + if (!result.ok) throw new Error(result.error); + expect(copiesFor(store, `task.${task.id}.`).toSorted()).toEqual(taskCopies.slice(-3).toSorted()); + expect(copiesFor(store, "workspace.workspace.").toSorted()).toEqual( + workspaceCopies.slice(-3).toSorted(), + ); + expect(result.data.recovery.removed).toBe(before - readdirSync(recoveryDir(store)).length); +}); + +test("Given gc --dry-run, Then it reports what it would remove and writes nothing at all", () => { + const { store, task } = startedStore(); + seedStaleCopies(store, `task.${task.id}.`, 10); + const workit = join(store.root, ".workit"); + // No lock, no storage initialization: the .gitignore stays deleted. + rmSync(join(workit, ".gitignore")); + const listing = () => + readdirSync(workit, { recursive: true }) + .map(String) + .toSorted() + .map((name) => `${name}:${statSync(join(workit, name)).mtimeMs}`); + const before = listing(); + const result = store.collectGarbage({ dryRun: true }); + expect(result).toMatchObject({ ok: true, data: { dryRun: true } }); + if (!result.ok) throw new Error(result.error); + expect(result.data.recovery.removed).toBeGreaterThan(0); + expect(listing()).toEqual(before); + expect(existsSync(join(workit, ".gitignore"))).toBe(false); + expect(existsSync(join(workit, "metadata.lock"))).toBe(false); +}); + +test("Given gc runs, Then every current read returns the same records", () => { + const { store, task } = startedStore(); + seedStaleCopies(store, `task.${task.id}.`, 12); + const tasksDir = join(store.root, ".workit", "tasks"); + const bytesBefore = readdirSync(tasksDir).map((name) => readFileSync(join(tasksDir, name))); + const workspaceBefore = readFileSync(join(store.root, ".workit", "workspace.json")); + const listBefore = store.listTasks(); + expect(store.collectGarbage().ok).toBe(true); + expect(store.listTasks()).toEqual(listBefore); + expect(store.readTask(task.id)).toMatchObject({ ok: true, data: { revision: task.revision } }); + expect(readdirSync(tasksDir).map((name) => readFileSync(join(tasksDir, name)))).toEqual( + bytesBefore, + ); + expect(readFileSync(join(store.root, ".workit", "workspace.json"))).toEqual(workspaceBefore); +}); + +const candidate = (head: string): Candidate => { + const value = { + id: "0".repeat(64), + scope: scope(), + completeness: "known" as const, + files: [], + environment: [], + head, + }; + return { ...value, id: candidateDigest(value) }; +}; + +test("Given paused, active, and closed tasks with duplicate stored candidates, When gc runs, Then only the paused task collapses to its latest positions", () => { + const { store, task } = startedStore(); + const [a, b] = [candidate("a"), candidate("b")]; + const paused = store.mutateTask(task.id, task.revision, (current) => + success(current.revision, null, { + ...current, + status: "paused", + pauseReason: "fixture", + candidates: [a, b, a, b, a], + }), + ); + if (!paused.ok) throw new Error(paused.error); + const active = store.create({ + expectedWorkspaceRevision: store.readWorkspace().ok + ? ((store.readWorkspace() as any).data.revision as string) + : null, + provenance, + intent: { objective: "active", scope: scope(), authorityRefs: [ref()] }, + }); + if (!active.ok) throw new Error(active.error); + const activeWithDuplicates = store.mutateTask(active.data.id, active.data.revision, (current) => + success(current.revision, null, { ...current, candidates: [a, a] }), + ); + if (!activeWithDuplicates.ok) throw new Error(activeWithDuplicates.error); + const workspace = store.readWorkspace(); + if (!workspace.ok || !workspace.data) throw new Error("workspace missing"); + const closedTask = store.create({ + expectedWorkspaceRevision: workspace.data.revision, + provenance, + intent: { objective: "closed", scope: scope(), authorityRefs: [ref()] }, + }); + if (!closedTask.ok) throw new Error(closedTask.error); + const closed = store.mutateTask(closedTask.data.id, closedTask.data.revision, (current) => + success(current.revision, null, { ...current, status: "closed", candidates: [b, b] }), + ); + if (!closed.ok) throw new Error(closed.error); + const closedFile = join(store.root, ".workit", "tasks", `${closedTask.data.id}.json`); + const closedBytes = readFileSync(closedFile); + + const result = store.collectGarbage(); + if (!result.ok) throw new Error(result.error); + expect(result.data.candidates).toMatchObject({ + removed: 3, + skippedActive: [active.data.id], + skippedClosed: [closedTask.data.id], + }); + expect(readFileSync(closedFile)).toEqual(closedBytes); + const after = store.readTask(task.id); + if (!after.ok) throw new Error(after.error); + expect(after.data.candidates.map((item) => item.head)).toEqual(["b", "a"]); + expect(after.data.candidates.at(-1)).toEqual(paused.data.candidates.at(-1)); + const untouched = store.readTask(active.data.id); + if (!untouched.ok) throw new Error(untouched.error); + expect(untouched.data.candidates).toHaveLength(2); +}); + +test("Given a stale leftover temp file, When gc runs, Then it is removed while fresh temp files stay", () => { + const { store, task } = startedStore(); + const stale = join(recoveryDir(store), `task.${task.id}.${"a".repeat(64)}.json.1.x.tmp`); + const fresh = join(recoveryDir(store), `task.${task.id}.${"b".repeat(64)}.json.1.y.tmp`); + writeFileSync(stale, "partial"); + writeFileSync(fresh, "partial"); + utimesSync(stale, new Date(0), new Date(0)); + const result = store.collectGarbage(); + expect(result).toMatchObject({ ok: true, data: { temporary: { removed: 1 } } }); + expect(() => statSync(stale)).toThrow(); + expect(statSync(fresh).isFile()).toBe(true); +}); + +const cliEntry = resolve(import.meta.dir, "../../packages/workit-cli/src/main.ts"); + +test("Given stale recovery copies, When `workit gc --json` runs, Then it prunes to the cap and reports the count", () => { + const { store, task } = startedStore(); + seedStaleCopies(store, `task.${task.id}.`, 9); + const run = spawnSync(process.execPath, [cliEntry, "gc", "--json"], { + cwd: store.root, + encoding: "utf8", + env: { ...process.env, HOME: store.root }, + }); + expect(run.status, run.stderr).toBe(0); + const report = JSON.parse(run.stdout); + expect(report).toMatchObject({ ok: true, data: { recovery: { removed: 6 } } }); + expect(copiesFor(store, `task.${task.id}.`)).toHaveLength(3); +}); + +test("Given no host supplies native recovery authority, Then state.recover is not advertised but stays parseable for the engine", () => { + const advertised = JSON.stringify(operationJsonSchema("state")); + expect(advertised).not.toContain('"recover"'); + expect(advertised).toContain('"export"'); + expect(JSON.stringify(boundedOperationJsonSchema("state"))).not.toContain('"recover"'); + expect( + parseOperation("state", { + schemaVersion: 1, + action: "recover", + target: "workspace", + expectedBytes: "a".repeat(64), + snapshotDigest: "b".repeat(64), + reason: "engine path", + authorityRefs: [], + }).ok, + ).toBe(true); +}); diff --git a/test/workit-core/store-lock.test.ts b/test/workit-core/store-lock.test.ts index 23334708..7e9a930e 100644 --- a/test/workit-core/store-lock.test.ts +++ b/test/workit-core/store-lock.test.ts @@ -1,6 +1,7 @@ import { expect, test } from "bun:test"; import { spawn, spawnSync } from "node:child_process"; import { + chmodSync, existsSync, mkdirSync, mkdtempSync, @@ -167,19 +168,26 @@ const workerScript = (root: string, taskId: string, calls: number) => ` import { TaskStore } from ${JSON.stringify(storeModule)}; const store = new TaskStore(${JSON.stringify(root)}); const codes = {}; +const errors = []; +const expected = new Set(["ok", "busy", "revision_conflict"]); 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, - })); + let result = current; + if (current.ok) + 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; + // Keep the detail of anything unexpected so a CI failure names the fs op. + if (!expected.has(code) && errors.length < 10) + errors.push({ step: current.ok ? "mutate" : "read", code, error: result.error, details: result.details }); + if (!current.ok) break; } codes[code] = (codes[code] ?? 0) + 1; } -process.stdout.write(JSON.stringify(codes)); +process.stdout.write(JSON.stringify({ codes, errors })); `; test("Given three processes each making 40 writes to one task, When they contend, Then none returns recovery_required", async () => { @@ -187,7 +195,7 @@ test("Given three processes each making 40 writes to one task, When they contend const runs = await Promise.all( [0, 1, 2].map( () => - new Promise>((done, fail) => { + new Promise<{ codes: Record; errors: unknown[] }>((done, fail) => { const child = spawn(process.execPath, ["-e", workerScript(store.root, task.id, 40)], { stdio: ["ignore", "pipe", "pipe"], }); @@ -207,13 +215,22 @@ test("Given three processes each making 40 writes to one task, When they contend ); 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); + for (const [code, count] of Object.entries(run.codes)) + totals[code] = (totals[code] ?? 0) + count; + const errors = JSON.stringify(runs.flatMap((run) => run.errors)); + // The invariant: contention never reports recovery_required, and every + // non-ok result is a retryable code. How many calls exhaust the short + // in-process wait budget (busy) depends on runner speed (a windows-latest + // run measured 84 ok / 36 busy), so the progress floor is deliberately + // weak: at least one process's worth of writes must land. + expect(totals.recovery_required ?? 0, errors).toBe(0); + const retryable = new Set(["ok", "busy", "revision_conflict"]); + expect( + Object.keys(totals).filter((code) => !retryable.has(code)), + errors, + ).toEqual([]); + expect(Object.values(totals).reduce((sum, count) => sum + count, 0)).toBe(120); + expect(totals.ok ?? 0).toBeGreaterThanOrEqual(40); }, 60_000); const storeModule2 = resolve(import.meta.dir, "../../packages/workit-core/src/core/store-lock.ts"); @@ -357,3 +374,38 @@ test.skipIf(!localLockHost().includes("#"))( expect(existsSync(lockPath)).toBe(false); }, ); + +// Permission denials are not contention: with nothing holding the lock, a +// read-only .workit must fail fast with storage_error and a permissions hint, +// never spend the budget and report busy. Windows reports these as +// EPERM/EACCES, so the Windows path is exercised by simulating the platform. +const readOnlyWorkit = !( + process.platform === "win32" || + (typeof process.getuid === "function" && process.getuid() === 0) +); +for (const simulateWindows of [false, true]) + test.skipIf(!readOnlyWorkit)( + `Given a read-only .workit and no lock holder${simulateWindows ? " (simulated Windows)" : ""}, When a write runs, Then it fails fast with storage_error and a permissions hint`, + () => { + const { store, task } = startedStore({ lockTimeoutMs: 2_000 }); + const workit = join(store.root, ".workit"); + const platform = Object.getOwnPropertyDescriptor(process, "platform")!; + chmodSync(workit, 0o555); + try { + if (simulateWindows) Object.defineProperty(process, "platform", { value: "win32" }); + const started = performance.now(); + const result = store.mutateTask(task.id, task.revision, identity); + const elapsed = performance.now() - started; + Object.defineProperty(process, "platform", platform); + expect(result).toMatchObject({ ok: false, code: "storage_error" }); + if (simulateWindows) + expect(result).toMatchObject({ + details: { guidance: expect.stringContaining("read-only attribute") }, + }); + expect(elapsed).toBeLessThan(1_000); + } finally { + Object.defineProperty(process, "platform", platform); + chmodSync(workit, 0o755); + } + }, + ); diff --git a/test/workit-mcp/server.test.ts b/test/workit-mcp/server.test.ts index 3a8eca6e..dc0e5eac 100644 --- a/test/workit-mcp/server.test.ts +++ b/test/workit-mcp/server.test.ts @@ -331,8 +331,12 @@ test("MCP resolves the trusted provider root and leaves read-only task listing b test("MCP publishes every core action in every family without a second action table", async () => { const families = new Map>(); for (const fixture of operationCorpus()) { + const action = String((fixture.input as { action: string }).action); + // state.recover parses but is not advertised: no shipped host supplies + // native recovery authority. + if (fixture.family === "state" && action === "recover") continue; const actions = families.get(fixture.family) ?? new Set(); - actions.add(String((fixture.input as { action: string }).action)); + actions.add(action); families.set(fixture.family, actions); } const { client, server } = await connect("cursor", { diff --git a/test/workit-opencode/task11-repair.test.ts b/test/workit-opencode/task11-repair.test.ts index 0a560ff0..9f9148bd 100644 --- a/test/workit-opencode/task11-repair.test.ts +++ b/test/workit-opencode/task11-repair.test.ts @@ -4,6 +4,7 @@ import { join } from "node:path"; import { tmpdir } from "node:os"; import { TaskStore, WorkitCore, success } from "@/packages/workit-core/src/core"; import { runDoctor } from "@/packages/workit-core/src/core/doctor"; +import { binDirWithRuntimes } from "@/test/shared/helpers/doctor-fixture"; import { scope, taskStartRequest } from "@/test/workit-core/task-fixtures"; import { server as plugin } from "@/packages/workit-opencode/src/index"; import { NativeReceiptStore } from "@/packages/workit-opencode/src/tools/workit"; @@ -1598,11 +1599,15 @@ test("doctor checks the OpenCode SDK pin in devDependencies", () => { devDependencies: { "@opencode-ai/plugin": "1.0.0" }, }), ); + // Offline and hermetic: only node and bun on PATH, so the doctor's + // registry (npm view) and provider identity (gh/glab) probes cannot reach + // the network and stall the test on a slow runner. const report = runDoctor({ host: "opencode", dev: root, home: root, configDir: join(root, "config"), + env: { ...process.env, HOME: root, PATH: binDirWithRuntimes(root) }, }); const versions = report.checks.find((check) => check.id === "versions"); expect(versions?.status).toBe("fail");