diff --git a/docs/supervisor.md b/docs/supervisor.md index fa80e888..5cb0147c 100644 --- a/docs/supervisor.md +++ b/docs/supervisor.md @@ -1,6 +1,6 @@ # Supervisor (jobs + shells) -> **Status: unsupported experimental.** These surfaces are source-available but outside the supported v3 product contract. They are hidden from default `hack --help` (see `hack help --all`) and print a warning when invoked. +> **Status: unsupported experimental.** These surfaces are source-available but outside the supported product contract. They are hidden from default `hack --help` (see `hack help --all`) and print a warning when invoked. The supervisor is the execution engine behind remote workflows. It can run commands as jobs, stream logs/events, and host PTY-backed shells. The CLI exposes it locally via `hack x supervisor` @@ -10,6 +10,16 @@ This page documents a beta-adjacent execution surface. Use [Beta workflows](beta.md) for the guided remote path and [Extensions & reference](reference.md) for the rest of the command and API material. +Cancellation is coordinated with the job runner. A successful cancel response waits +for the runner's terminal metadata and event, not process exit or log drain; it does not +race a second writer against process exit. Once terminal persistence has started, +a later cancellation returns `not_running` and preserves the completed/failed outcome. +Repeated requests accepted before finalization share one cancellation outcome and event. +If terminal persistence fails, cancellation and the run promise report that failure; +the service logs it without replacing the claimed outcome with `failed`. Metadata or +event history can be incomplete, so a persistence error is not a successful acknowledgement. +No automatic retry appends another terminal event after a partial write. + ## Local usage ```bash diff --git a/src/control-plane/extensions/supervisor/runner.ts b/src/control-plane/extensions/supervisor/runner.ts index 3fbb72a7..409bd0f2 100644 --- a/src/control-plane/extensions/supervisor/runner.ts +++ b/src/control-plane/extensions/supervisor/runner.ts @@ -10,6 +10,8 @@ export type JobRunResult = { export type JobSpawnListener = (opts: { readonly proc: SpawnedProcess; + /** Request cancellation before terminal status persistence begins. */ + readonly cancel: () => Promise; }) => void; type SpawnedProcess = ReturnType; @@ -23,6 +25,7 @@ type SpawnedProcess = ReturnType; * @param opts.cwd - Optional working directory for the process. * @param opts.env - Optional environment overrides. * @param opts.onSpawn - Optional hook with the spawned process handle. + * @param opts.onTerminalClaim - Observe the chosen outcome before persistence yields. * @returns Final job status and exit code. */ export async function runJob(opts: { @@ -32,6 +35,7 @@ export async function runJob(opts: { readonly cwd?: string; readonly env?: Record; readonly onSpawn?: JobSpawnListener; + readonly onTerminalClaim?: (opts: { readonly status: JobStatus }) => void; }): Promise { const meta = await opts.jobStore.readJobMeta({ jobId: opts.jobId }); if (!meta) { @@ -79,7 +83,48 @@ export async function runJob(opts: { type: "job.started", payload: { pid: proc.pid }, }); - opts.onSpawn?.({ proc }); + + let terminalStatus: JobStatus | undefined; + let terminalWrite: Promise | undefined; + const finish = (input: { + readonly status: JobStatus; + readonly exitCode?: number; + }): Promise => { + if (terminalWrite) { + return terminalWrite; + } + // Claim the outcome before storage yields; both paths share one writer. + terminalStatus = input.status; + opts.onTerminalClaim?.({ status: input.status }); + terminalWrite = (async () => { + await opts.jobStore.updateJobStatus({ + jobId: opts.jobId, + status: input.status, + }); + await opts.jobStore.appendEvent({ + jobId: opts.jobId, + type: `job.${input.status}`, + ...(input.exitCode === undefined + ? {} + : { payload: { exitCode: input.exitCode } }), + }); + return input.status; + })(); + return terminalWrite; + }; + opts.onSpawn?.({ + proc, + cancel: async () => { + if (terminalStatus && terminalStatus !== "cancelled") { + return false; + } + if (!terminalStatus) { + proc.kill(); + } + await finish({ status: "cancelled" }); + return true; + }, + }); const paths = opts.jobStore.getJobPaths({ jobId: opts.jobId }); const stdoutTask = pipeStreamToFiles({ @@ -94,19 +139,10 @@ export async function runJob(opts: { const exitCode = await proc.exited; await Promise.all([stdoutTask, stderrTask]); - const metaAfter = await opts.jobStore.readJobMeta({ jobId: opts.jobId }); - if (metaAfter?.status === "cancelled") { - return { jobId: opts.jobId, status: "cancelled", exitCode }; - } - - const status: JobStatus = exitCode === 0 ? "completed" : "failed"; - await opts.jobStore.updateJobStatus({ jobId: opts.jobId, status }); - await opts.jobStore.appendEvent({ - jobId: opts.jobId, - type: status === "completed" ? "job.completed" : "job.failed", - payload: { exitCode }, + const status = await finish({ + status: exitCode === 0 ? "completed" : "failed", + exitCode, }); - return { jobId: opts.jobId, status, exitCode }; } diff --git a/src/control-plane/extensions/supervisor/service.ts b/src/control-plane/extensions/supervisor/service.ts index 93b843b1..65284557 100644 --- a/src/control-plane/extensions/supervisor/service.ts +++ b/src/control-plane/extensions/supervisor/service.ts @@ -40,7 +40,7 @@ export type SupervisorService = { readonly env?: Record; }) => Promise; /** - * Attempt to cancel a running job. + * Request cancellation and wait for its terminal metadata and event. * * @param opts.projectDir - Project .hack directory. * @param opts.jobId - Job id to cancel. @@ -76,15 +76,18 @@ export type SupervisorService = { * Create a supervisor service for managing jobs and their metadata. * * @param opts.logger - Optional logger override. + * @param opts.createStore - Optional job-store factory. * @returns Supervisor service helpers. */ export function createSupervisorService(opts?: { readonly logger?: Logger; + readonly createStore?: typeof createJobStore; }): SupervisorService { const logger = opts?.logger ?? baseLogger; + const createStore = opts?.createStore ?? createJobStore; const runningJobs = new Map< string, - { readonly proc: ReturnType } + { readonly cancel: () => Promise } >(); const createJob = async (input: { @@ -96,7 +99,7 @@ export function createSupervisorService(opts?: { readonly cwd?: string; readonly env?: Record; }): Promise => { - const store = await createJobStore({ projectDir: input.projectDir }); + const store = await createStore({ projectDir: input.projectDir }); const jobId = randomUUID(); const meta = await store.createJob({ jobId, @@ -106,17 +109,24 @@ export function createSupervisorService(opts?: { projectName: input.projectName, }); + let terminalClaimed = false; const run = runJob({ jobStore: store, jobId, command: input.command, cwd: input.cwd, env: input.env, - onSpawn: ({ proc }) => { - runningJobs.set(jobId, { proc }); + onTerminalClaim: () => { + terminalClaimed = true; + }, + onSpawn: ({ cancel }) => { + runningJobs.set(jobId, { cancel }); }, }) .catch(async (error) => { + if (terminalClaimed) { + throw error; + } logger.error({ message: `Job failed: ${formatError(error)}` }); await store.updateJobStatus({ jobId, status: "failed" }); await store.appendEvent({ @@ -131,6 +141,13 @@ export function createSupervisorService(opts?: { runningJobs.delete(jobId); }); + // Background callers need not await run; observing errors here keeps the + // returned promise rejected without creating an unhandled daemon rejection. + void run.catch((error) => { + logger.error({ + message: `Job persistence failed: ${formatError(error)}`, + }); + }); return { jobId, meta, run }; }; @@ -141,7 +158,7 @@ export function createSupervisorService(opts?: { readonly projectDir: string; readonly jobId: string; }): Promise => { - const store = await createJobStore({ projectDir }); + const store = await createStore({ projectDir }); const meta = await store.readJobMeta({ jobId }); if (!meta) { return { ok: false, status: "not_found" }; @@ -152,10 +169,9 @@ export function createSupervisorService(opts?: { return { ok: false, status: "not_running" }; } - running.proc.kill(); - await store.updateJobStatus({ jobId, status: "cancelled" }); - await store.appendEvent({ jobId, type: "job.cancelled" }); - + if (!(await running.cancel())) { + return { ok: false, status: "not_running" }; + } return { ok: true, status: "cancelled" }; }; @@ -166,7 +182,7 @@ export function createSupervisorService(opts?: { readonly projectDir: string; readonly jobId: string; }): Promise => { - const store = await createJobStore({ projectDir }); + const store = await createStore({ projectDir }); return await store.readJobMeta({ jobId }); }; @@ -175,7 +191,7 @@ export function createSupervisorService(opts?: { }: { readonly projectDir: string; }): Promise => { - const store = await createJobStore({ projectDir }); + const store = await createStore({ projectDir }); const entries = await safeReadDir(store.jobsRoot); const metas = await Promise.all( entries.map((jobId) => store.readJobMeta({ jobId })) diff --git a/tests/supervisor-runner.test.ts b/tests/supervisor-runner.test.ts index 8199df87..7e6cd48c 100644 --- a/tests/supervisor-runner.test.ts +++ b/tests/supervisor-runner.test.ts @@ -83,3 +83,165 @@ test("runJob records failed status for non-zero exit", async () => { const types = events.map((event) => event.type); expect(types).toContain("job.failed"); }); + +test("cancellation survives process exit before terminal persistence and is recorded once", async () => { + tempDir = await mkdtemp(join(tmpdir(), "hack-supervisor-runner-")); + const store = await createJobStore({ projectDir: join(tempDir, ".hack") }); + await store.createJob({ + jobId: "cancel-before-write", + runner: "generic", + command: [process.execPath, "-e", "setTimeout(() => {}, 5000)"], + }); + const terminalWrite = Promise.withResolvers(); + const allowWrite = Promise.withResolvers(); + let cancel: () => Promise = async () => false; + let cancellations: Promise[] = []; + let processExited: Promise = Promise.resolve(0); + const run = runJob({ + jobId: "cancel-before-write", + jobStore: { + ...store, + updateJobStatus: async (opts) => { + if (["cancelled", "completed", "failed"].includes(opts.status)) { + terminalWrite.resolve(); + await allowWrite.promise; + } + return await store.updateJobStatus(opts); + }, + }, + onSpawn: (control) => { + cancel = control.cancel; + processExited = control.proc.exited; + cancellations = [cancel(), cancel()]; + }, + }); + try { + await terminalWrite.promise; + await processExited; + // The killed process has exited, but no terminal metadata has been saved. + expect( + (await store.readJobMeta({ jobId: "cancel-before-write" }))?.status + ).toBe("running"); + } finally { + allowWrite.resolve(); + await run; + } + expect(await Promise.all(cancellations)).toEqual([true, true]); + expect((await run).status).toBe("cancelled"); + expect( + (await store.readJobMeta({ jobId: "cancel-before-write" }))?.status + ).toBe("cancelled"); + const events = await store.readEvents({ jobId: "cancel-before-write" }); + expect(events.map((event) => event.type)).toEqual([ + "job.created", + "job.starting", + "job.started", + "job.cancelled", + ]); + expect(events.map((event) => event.seq)).toEqual([1, 2, 3, 4]); +}); + +test("late cancellation cannot replace a completed outcome while its write is pending", async () => { + tempDir = await mkdtemp(join(tmpdir(), "hack-supervisor-runner-")); + const store = await createJobStore({ projectDir: join(tempDir, ".hack") }); + await store.createJob({ + jobId: "completion-first", + runner: "generic", + command: [process.execPath, "-e", "process.exit(0)"], + }); + const terminalWrite = Promise.withResolvers(); + const allowWrite = Promise.withResolvers(); + let cancel: () => Promise = async () => false; + const run = runJob({ + jobId: "completion-first", + jobStore: { + ...store, + updateJobStatus: async (opts) => { + if (opts.status === "completed") { + terminalWrite.resolve(); + await allowWrite.promise; + } + return await store.updateJobStatus(opts); + }, + }, + onSpawn: (control) => { + cancel = control.cancel; + }, + }); + try { + await terminalWrite.promise; + expect(await cancel()).toBe(false); + } finally { + allowWrite.resolve(); + await run; + } + expect((await run).status).toBe("completed"); + expect((await store.readJobMeta({ jobId: "completion-first" }))?.status).toBe( + "completed" + ); + const events = await store.readEvents({ jobId: "completion-first" }); + expect(events.map((event) => event.type)).toEqual([ + "job.created", + "job.starting", + "job.started", + "job.completed", + ]); +}); + +test("cancellation acknowledgement does not wait for a process that ignores SIGTERM", async () => { + tempDir = await mkdtemp(join(tmpdir(), "hack-supervisor-runner-")); + const readyPath = join(tempDir, "ready"); + const store = await createJobStore({ projectDir: join(tempDir, ".hack") }); + await store.createJob({ + jobId: "ignores-signal", + runner: "generic", + command: [ + process.execPath, + "-e", + `process.on("SIGTERM", () => {}); await Bun.write(${JSON.stringify(readyPath)}, "ready"); setInterval(() => {}, 1000);`, + ], + }); + let control: + | Parameters[0]["onSpawn"]>>[0] + | undefined; + const run = runJob({ + jobStore: store, + jobId: "ignores-signal", + onSpawn: (value) => { + control = value; + }, + }); + let exited = false; + const observedRun = run.finally(() => { + exited = true; + }); + let deadlineTimer: ReturnType | undefined; + try { + const deadline = Date.now() + 2000; + while (!(await Bun.file(readyPath).exists())) { + if (Date.now() > deadline) { + throw new Error("Child did not become ready"); + } + await Bun.sleep(10); + } + if (!control) { + throw new Error("Missing process control"); + } + const deadlinePromise = new Promise((_resolve, reject) => { + deadlineTimer = setTimeout( + () => reject(new Error("Cancellation waited for process exit")), + 2000 + ); + }); + expect(await Promise.race([control.cancel(), deadlinePromise])).toBe(true); + expect(exited).toBe(false); + expect((await store.readJobMeta({ jobId: "ignores-signal" }))?.status).toBe( + "cancelled" + ); + } finally { + clearTimeout(deadlineTimer); + control?.proc.kill("SIGKILL"); + await observedRun; + } + expect((await run).status).toBe("cancelled"); +}); diff --git a/tests/supervisor-service.test.ts b/tests/supervisor-service.test.ts index 5bf27f39..f950fee3 100644 --- a/tests/supervisor-service.test.ts +++ b/tests/supervisor-service.test.ts @@ -1,8 +1,9 @@ import { afterEach, expect, test } from "bun:test"; -import { mkdir, mkdtemp, rm } from "node:fs/promises"; +import { appendFile, mkdir, mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { createJobStore } from "../src/control-plane/extensions/supervisor/job-store.ts"; import { createSupervisorService } from "../src/control-plane/extensions/supervisor/service.ts"; import { readTextFile } from "../src/lib/fs.ts"; @@ -70,9 +71,21 @@ test("Supervisor service cancels running jobs", async () => { const cancel = await service.cancelJob({ projectDir, jobId: created.jobId }); expect(cancel.ok).toBe(true); - + // A successful cancellation response includes durable terminal state. + expect( + (await service.getJob({ projectDir, jobId: created.jobId }))?.status + ).toBe("cancelled"); + const store = await createJobStore({ projectDir }); + const events = await store.readEvents({ jobId: created.jobId }); + expect(events.filter((event) => event.type === "job.cancelled")).toHaveLength( + 1 + ); + expect(events.some((event) => event.type === "job.failed")).toBe(false); const result = await created.run; expect(result.status).toBe("cancelled"); + expect(await service.cancelJob({ projectDir, jobId: created.jobId })).toEqual( + { ok: false, status: "not_running" } + ); const job = await service.getJob({ projectDir, jobId: created.jobId }); expect(job?.status).toBe("cancelled"); @@ -107,3 +120,116 @@ async function waitForJobStatus(opts: { } throw new Error(`Timed out waiting for job status: ${opts.status}`); } + +for (const failurePoint of ["status", "event-before", "event-after"] as const) { + test(`cancellation persistence failure at ${failurePoint} preserves the claimed outcome`, async () => { + tempDir = await mkdtemp(join(tmpdir(), "hack-supervisor-service-")); + const projectDir = join(tempDir, ".hack"); + const failure = new Error(`Injected ${failurePoint} failure`); + const service = createSupervisorService({ + createStore: async (opts) => { + const store = await createJobStore(opts); + return { + ...store, + updateJobStatus: async (input) => { + if (input.status === "cancelled" && failurePoint === "status") { + throw failure; + } + return await store.updateJobStatus(input); + }, + appendEvent: async (input) => { + if (input.type !== "job.cancelled") { + return await store.appendEvent(input); + } + if (failurePoint === "event-before") { + throw failure; + } + if (failurePoint === "event-after") { + const meta = await store.readJobMeta({ jobId: input.jobId }); + if (!meta) { + throw new Error("Missing test job"); + } + // Fail after the event append, before its sequence reaches metadata. + await appendFile( + store.getJobPaths(input).eventsPath, + `${JSON.stringify({ + seq: meta.lastEventSeq + 1, + ts: new Date().toISOString(), + type: input.type, + })}\n` + ); + throw failure; + } + return await store.appendEvent(input); + }, + }; + }, + }); + const created = await service.createJob({ + projectDir, + runner: "generic", + command: [process.execPath, "-e", "setTimeout(() => {}, 5000)"], + }); + const outcome = created.run.catch((error: unknown) => error); + await waitForJobStatus({ + service, + projectDir, + jobId: created.jobId, + status: "running", + timeoutMs: 10_000, + }); + await expect( + service.cancelJob({ projectDir, jobId: created.jobId }) + ).rejects.toThrow(failure.message); + expect(await outcome).toBe(failure); + const stored = await service.getJob({ projectDir, jobId: created.jobId }); + expect(stored?.status).toBe( + failurePoint === "status" ? "running" : "cancelled" + ); + const store = await createJobStore({ projectDir }); + const events = await store.readEvents({ jobId: created.jobId }); + expect(stored?.lastEventSeq).toBe(3); + expect(events.filter((event) => event.type === "job.failed")).toHaveLength( + 0 + ); + expect( + events.filter((event) => event.type === "job.cancelled") + ).toHaveLength(failurePoint === "event-after" ? 1 : 0); + }); +} + +test("completion event persistence failure does not replace completed with failed", async () => { + tempDir = await mkdtemp(join(tmpdir(), "hack-supervisor-service-")); + const projectDir = join(tempDir, ".hack"); + const failure = new Error("Injected completion event failure"); + const service = createSupervisorService({ + createStore: async (opts) => { + const store = await createJobStore(opts); + return { + ...store, + appendEvent: async (input) => { + const event = await store.appendEvent(input); + if (input.type === "job.completed") { + throw failure; + } + return event; + }, + }; + }, + }); + const created = await service.createJob({ + projectDir, + runner: "generic", + command: [process.execPath, "-e", "process.exit(0)"], + }); + await expect(created.run).rejects.toThrow(failure.message); + expect( + (await service.getJob({ projectDir, jobId: created.jobId }))?.status + ).toBe("completed"); + const store = await createJobStore({ projectDir }); + expect( + (await store.readEvents({ jobId: created.jobId })).map( + (event) => event.type + ) + ).toEqual(["job.created", "job.starting", "job.started", "job.completed"]); +});