From 30c63af0e6197e4d903e3f62572cc34472ee9aba Mon Sep 17 00:00:00 2001 From: Dotta <34892728+cryppadotta@users.noreply.github.com> Date: Fri, 11 Sep 2026 18:25:42 -0500 Subject: [PATCH 1/2] fix(cli): recover abandoned workspace build locks (#13288) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - The test-drive command starts an isolated instance for local testing. > - Source startup builds shared packages before it starts the server. > - An interrupted build can leave an empty lock directory. > - Later starts wait without output and fail after 60 seconds. > - This pull request recovers abandoned locks and shows build progress. > - Local testing can start again without manual lock removal. ## Linked Issues or Issue Description **What happened?** `pnpm paperclipai test-drive` stopped at “Starting Paperclip server…” in a source checkout. A leftover plugin build lock caused a silent 60-second wait and then a timeout. **Expected behavior** Startup should recover an abandoned build lock. It should show when it waits for a live build. An interrupted or failed build should not leave partial output that the next start accepts as complete. **Steps to reproduce** 1. Leave an empty `node_modules/.cache/paperclip-plugin-build-deps.lock` directory after an interrupted build. 2. Make the shared or plugin SDK build output out of date. 3. Run `pnpm paperclipai test-drive --api-key placeholder --no-browser` with a fresh data directory. 4. Observe the silent wait at server startup. **Paperclip version or commit** Reproduced at `2083bf6f9`. **Deployment mode** Local source checkout with an isolated embedded PostgreSQL instance. Related work: #12894 added test-drive. #12898 restored its credential inputs. Neither change handles abandoned workspace build locks. No duplicate fix was found. ## What Changed - Publish a lock directory with an owner record in one rename. - Recover locks after their owner and compiler exit. Recover legacy empty locks after two minutes. - Keep the lock until the compiler stops on SIGINT or SIGTERM. - Print build and lock-wait progress. - Record source, dependency, compiler-config, and output content fingerprints only after a successful compile. Recover partial output even when modification times are unchanged. - Add 12 process-level regression tests and update the development guide. ## Verification - `node --test scripts/__tests__/ensure-plugin-build-deps.test.mjs`: 12 tests pass. - `pnpm exec vitest run --config cli/vitest.config.ts cli/src/__tests__/test-drive.test.ts`: 32 tests pass. - `pnpm --filter paperclipai typecheck`: passed. - `pnpm --filter paperclipai build`: passed. - Live smoke tests: fresh startup and startup with an abandoned lock both reach ready state. The API and UI respond. The command creates the company and CEO and enables worktree execution. Test instances stop cleanly. - Full repository `pnpm -r typecheck` and `pnpm build`: passed. - Full Vitest suite coverage completed using the repository-supported server, chat, workspace, and serialized shards. The initial local run needed the fresh-worktree fake native-provider binary built and focused reruns for port/socket races and load-related timeouts; all affected tests passed on rerun. Suites skipped by fail-fast exits were run separately and passed. The initial serial `pnpm test:run` was stopped in favor of these shards. - Greptile: 5/5 on commit `8b5a790c2af06a52b5dc76e5f52331966df990b8`, with all review threads resolved. - CI: 31 checks passed and two Storybook checks intentionally skipped. The initial workspace and browser jobs were interrupted by runner shutdowns; both passed on the second attempt. Build, typecheck, canary dry run, all general and serialized tests, all browser shards, security checks, and final verification summaries are green. [CI run](https://github.com/paperclipai/paperclip/actions/runs/34654730783) ## Risks - This changes shared source-build locking for the CLI and plugin SDK commands. - Legacy locks have no owner identity. Recovery uses a two-minute age threshold for empty legacy directories. - Startup reads and hashes source and output files to verify the build cache. Identical direct builds reuse the cache. Changed or partial output requires a rebuild. - A reused process ID can delay recovery. Live owner or compiler processes keep their lock. - No database, API, or UI contract changes. ## Model Used OpenAI GPT-6 in Codex, with reasoning, tool use, code execution, and process-level testing. A more specific API model identifier and context-window size are not exposed in this session. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip --- doc/DEVELOPING.md | 13 + .../ensure-plugin-build-deps.test.mjs | 238 +++++++++++++++++ scripts/ensure-plugin-build-deps.mjs | 246 +++++++++++++----- 3 files changed, 432 insertions(+), 65 deletions(-) create mode 100644 scripts/__tests__/ensure-plugin-build-deps.test.mjs diff --git a/doc/DEVELOPING.md b/doc/DEVELOPING.md index 8927e66d41..3983005523 100644 --- a/doc/DEVELOPING.md +++ b/doc/DEVELOPING.md @@ -417,6 +417,19 @@ configs with `database.mode: postgres`, suppresses the invocation directory's guard. The selected instance's own environment file still loads. The command selects the first available loopback port at or above `3100`. +Source-checkout startup builds the shared and plugin SDK packages when needed. +It prints build progress and any wait for another build. Interrupted builds +release their lock after the compiler stops; later startups recover locks whose +owner and compiler have exited. Empty locks from older versions are recovered +once they are at least two minutes old. The command remains in the foreground +after printing its ready URL to serve the instance; use Ctrl-C to stop it. +Each package gets a completion marker only after a successful build. A hard +kill leaves that marker absent, so the next startup rebuilds partial output. +The marker records source and output content fingerprints, so recovery does +not depend on filesystem timestamp precision. Direct `tsc` builds that produce +identical output reuse the marker. Changed or partial output is rebuilt once +before later startups reuse the completed build. + Claude uses `ANTHROPIC_API_KEY`; Codex uses `OPENAI_API_KEY`; OpenCode uses `OPENROUTER_API_KEY` and requires an `openrouter/...` model. `--api-key-env` can name a different source variable while the agent still receives the diff --git a/scripts/__tests__/ensure-plugin-build-deps.test.mjs b/scripts/__tests__/ensure-plugin-build-deps.test.mjs new file mode 100644 index 0000000000..99778f75f3 --- /dev/null +++ b/scripts/__tests__/ensure-plugin-build-deps.test.mjs @@ -0,0 +1,238 @@ +import assert from "node:assert/strict"; +import { spawn, spawnSync } from "node:child_process"; +import { once } from "node:events"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { setTimeout as sleep } from "node:timers/promises"; +import test from "node:test"; + +function fixture(t) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), "paperclip-build-lock-")); + t.after(() => fs.rmSync(root, { recursive: true, force: true })); + fs.mkdirSync(path.join(root, "scripts")); + fs.copyFileSync(new URL("../ensure-plugin-build-deps.mjs", import.meta.url), path.join(root, "scripts/ensure-plugin-build-deps.mjs")); + const compiler = path.join(root, "node_modules/typescript/bin/tsc"); + fs.mkdirSync(path.dirname(compiler), { recursive: true }); + fs.writeFileSync(compiler, ` +const fs = require("node:fs"); +const path = require("node:path"); +const target = path.dirname(process.argv[3]); +const active = path.resolve("compiler-active"); +try { fs.mkdirSync(active); } catch { process.exit(42); } +process.on("exit", () => fs.rmSync(active, { recursive: true, force: true })); +process.on("SIGTERM", () => process.exit(143)); +process.on("SIGINT", () => process.exit(130)); +fs.appendFileSync("builds", target + "\\n"); +fs.mkdirSync(path.join(target, "dist"), { recursive: true }); +// Deliberately write index.js before the compiler finishes emitting the rest. +fs.writeFileSync(path.join(target, "dist/index.js"), "export {};\\n"); +setTimeout(() => { + if (fs.existsSync("fail")) process.exit(2); + fs.writeFileSync(path.join(target, "dist/complete"), "done"); +}, Number(process.env.BUILD_DELAY ?? 20)); +`); + for (const target of ["packages/shared", "packages/plugins/sdk"]) { + fs.mkdirSync(path.join(root, target, "src"), { recursive: true }); + fs.writeFileSync(path.join(root, target, "src/index.ts"), "export {};\n"); + fs.writeFileSync(path.join(root, target, "tsconfig.json"), "{}"); + } + const lock = path.join(root, "node_modules/.cache/paperclip-plugin-build-deps.lock"); + const launch = (env = {}) => { + const child = spawn(process.execPath, ["scripts/ensure-plugin-build-deps.mjs"], { + cwd: root, env: { ...process.env, ...env }, stdio: ["ignore", "pipe", "pipe"], + }); + let output = ""; + child.stdout.on("data", (data) => { output += data; }); + child.stderr.on("data", (data) => { output += data; }); + const done = once(child, "close").then(([code]) => ({ code, output })); + t.after(() => { if (child.exitCode === null) child.kill("SIGTERM"); }); + return { child, done }; + }; + return { root, lock, launch }; +} + +async function until(predicate) { + const deadline = Date.now() + 5000; + while (!predicate()) { + assert.ok(Date.now() < deadline, "condition timed out"); + await sleep(10); + } +} + +test("recovers the old empty lock left by interrupted startup", async (t) => { + const f = fixture(t); + fs.mkdirSync(f.lock, { recursive: true }); + const old = new Date(Date.now() - 180_000); + fs.utimesSync(f.lock, old, old); + const result = await f.launch().done; + assert.equal(result.code, 0, result.output); + assert.match(result.output, /Recovered abandoned/); + assert.match(result.output, /Building @paperclipai\/shared/); + assert.equal(fs.existsSync(f.lock), false); +}); + +test("concurrent startups recover a dead owner and build only once", async (t) => { + const f = fixture(t); + const dead = spawnSync(process.execPath, ["-e", "" ]).pid; + fs.mkdirSync(f.lock, { recursive: true }); + fs.writeFileSync(path.join(f.lock, `owner-${dead}-old.json`), JSON.stringify({ pid: dead })); + const results = await Promise.all(Array.from({ length: 4 }, () => f.launch({ BUILD_DELAY: "150" }).done)); + for (const result of results) assert.equal(result.code, 0, result.output); + assert.equal(fs.readFileSync(path.join(f.root, "builds"), "utf8").trim().split("\n").length, 2); +}); + +test("does not accept partially emitted output while another compiler holds the lock", async (t) => { + const f = fixture(t); + const first = f.launch({ BUILD_DELAY: "200" }); + await until(() => fs.existsSync(path.join(f.root, "packages/plugins/sdk/dist/index.js"))); + const second = await f.launch().done; + assert.equal(second.code, 0, second.output); + assert.match(second.output, /Waiting for another workspace build/); + assert.equal(fs.existsSync(path.join(f.root, "packages/plugins/sdk/dist/complete")), true); + assert.equal((await first.done).code, 0); +}); + +test("preserves a live compiler's lock even when its parent has exited", async (t) => { + const f = fixture(t); + const dead = spawnSync(process.execPath, ["-e", ""]).pid; + fs.mkdirSync(f.lock, { recursive: true }); + fs.writeFileSync(path.join(f.lock, `owner-${dead}-old.json`), JSON.stringify({ pid: dead, childPid: process.pid })); + const run = f.launch(); + await sleep(200); + assert.equal(fs.existsSync(path.join(f.root, "builds")), false); + run.child.kill("SIGTERM"); + assert.equal((await run.done).code, 143); + assert.equal(fs.existsSync(f.lock), true); +}); + +test("termination stops the compiler and releases the lock for the next startup", async (t) => { + const f = fixture(t); + const run = f.launch({ BUILD_DELAY: "10000" }); + await until(() => fs.existsSync(path.join(f.root, "compiler-active"))); + run.child.kill("SIGTERM"); + assert.equal((await run.done).code, 143); + assert.equal(fs.existsSync(f.lock), false); + assert.equal(fs.existsSync(path.join(f.root, "compiler-active")), false); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.equal(fs.existsSync(path.join(f.root, "packages/shared/dist/complete")), true); +}); + +test("failed compilation releases the lock and is rebuilt on retry", async (t) => { + const f = fixture(t); + fs.writeFileSync(path.join(f.root, "fail"), ""); + assert.equal((await f.launch().done).code, 2); + assert.equal(fs.existsSync(f.lock), false); + fs.unlinkSync(path.join(f.root, "fail")); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.equal(fs.existsSync(path.join(f.root, "packages/shared/dist/complete")), true); +}); + +test("hard-killed compilation cannot leave partial output accepted on recovery", async (t) => { + const f = fixture(t); + // First establish valid completion markers, then start a rebuild. + assert.equal((await f.launch().done).code, 0); + const output = path.join(f.root, "packages/shared/dist/index.js"); + const complete = path.join(f.root, "packages/shared/dist/complete"); + fs.unlinkSync(output); + fs.unlinkSync(complete); + const run = f.launch({ BUILD_DELAY: "10000" }); + await until(() => fs.existsSync(output)); + const owner = JSON.parse(fs.readFileSync(path.join(f.lock, fs.readdirSync(f.lock)[0]), "utf8")); + run.child.kill("SIGKILL"); + process.kill(owner.childPid, "SIGKILL"); + await run.done; + await until(() => { + try { process.kill(owner.childPid, 0); return false; } + catch (error) { return error.code === "ESRCH"; } + }); + // SIGKILL cannot run the fixture compiler's exit hook either. + fs.rmSync(path.join(f.root, "compiler-active"), { recursive: true, force: true }); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.match(retry.output, /Recovered abandoned/); + assert.equal(fs.existsSync(complete), true); +}); + +test("rebuilds changed direct output once, then reuses the completed build", async (t) => { + const f = fixture(t); + assert.equal((await f.launch().done).code, 0); + const builds = fs.readFileSync(path.join(f.root, "builds"), "utf8"); + // A successful or interrupted direct tsc invocation updates index.js without + // changing our marker. Neither can certify that all output was emitted. + for (const target of ["packages/shared", "packages/plugins/sdk"]) { + fs.appendFileSync(path.join(f.root, target, "dist/index.js"), "// direct build changed output\n"); + } + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + const rebuilt = fs.readFileSync(path.join(f.root, "builds"), "utf8"); + assert.equal(rebuilt.trim().split("\n").length, builds.trim().split("\n").length + 2); + const next = await f.launch().done; + assert.equal(next.code, 0, next.output); + assert.equal(fs.readFileSync(path.join(f.root, "builds"), "utf8"), rebuilt); + assert.doesNotMatch(next.output, /Building/); +}); + +test("rejects partial output from an interrupted direct compiler", async (t) => { + const f = fixture(t); + assert.equal((await f.launch().done).code, 0); + const target = path.join(f.root, "packages/shared"); + fs.unlinkSync(path.join(target, "dist/complete")); + const completionTime = fs.statSync(path.join(target, "dist/.paperclip-build-complete")).mtimeMs; + await sleep(20); + const compiler = spawn(process.execPath, ["node_modules/typescript/bin/tsc", "-p", path.join(target, "tsconfig.json")], { + cwd: f.root, env: { ...process.env, BUILD_DELAY: "10000" }, stdio: "ignore", + }); + const closed = once(compiler, "close"); + t.after(() => { if (compiler.exitCode === null) compiler.kill("SIGKILL"); }); + await until(() => fs.statSync(path.join(target, "dist/index.js")).mtimeMs > completionTime); + compiler.kill("SIGKILL"); + await closed; + fs.rmSync(path.join(f.root, "compiler-active"), { recursive: true, force: true }); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.equal(fs.existsSync(path.join(target, "dist/complete")), true); +}); + +test("detects partial output even when all modification times are unchanged", async (t) => { + const f = fixture(t); + assert.equal((await f.launch().done).code, 0); + const target = path.join(f.root, "packages/shared"); + const output = path.join(target, "dist/index.js"); + const oldTime = fs.statSync(output).mtime; + fs.writeFileSync(output, "// incomplete direct build\n"); + fs.utimesSync(output, oldTime, oldTime); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.equal(fs.readFileSync(output, "utf8"), "export {};\n"); +}); + +test("reuses identical direct output regardless of its timestamps", async (t) => { + const f = fixture(t); + assert.equal((await f.launch().done).code, 0); + const builds = fs.readFileSync(path.join(f.root, "builds"), "utf8"); + const output = path.join(f.root, "packages/shared/dist/index.js"); + fs.writeFileSync(output, fs.readFileSync(output)); + const newer = new Date(Date.now() + 1000); + fs.utimesSync(output, newer, newer); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.equal(fs.readFileSync(path.join(f.root, "builds"), "utf8"), builds); + assert.doesNotMatch(retry.output, /Building/); +}); + +test("shared source changes invalidate both shared and dependent SDK output", async (t) => { + const f = fixture(t); + assert.equal((await f.launch().done).code, 0); + const source = path.join(f.root, "packages/shared/src/index.ts"); + const oldTime = fs.statSync(source).mtime; + fs.appendFileSync(source, "export const changed = true;\n"); + fs.utimesSync(source, oldTime, oldTime); + const retry = await f.launch().done; + assert.equal(retry.code, 0, retry.output); + assert.match(retry.output, /Building @paperclipai\/shared/); + assert.match(retry.output, /Building @paperclipai\/plugin-sdk/); + assert.equal(fs.readFileSync(path.join(f.root, "builds"), "utf8").trim().split("\n").length, 4); +}); diff --git a/scripts/ensure-plugin-build-deps.mjs b/scripts/ensure-plugin-build-deps.mjs index 5b19024dd9..356e0e0f80 100644 --- a/scripts/ensure-plugin-build-deps.mjs +++ b/scripts/ensure-plugin-build-deps.mjs @@ -1,9 +1,11 @@ #!/usr/bin/env node -import { spawnSync } from "node:child_process"; +import { spawn } from "node:child_process"; import fs from "node:fs"; import path from "node:path"; import { fileURLToPath } from "node:url"; +import { createHash, randomUUID } from "node:crypto"; +import { setTimeout as sleep } from "node:timers/promises"; const scriptDir = path.dirname(fileURLToPath(import.meta.url)); const rootDir = path.resolve(scriptDir, ".."); @@ -16,14 +18,18 @@ const buildTargets = [ { name: "@paperclipai/shared", output: path.join(rootDir, "packages/shared/dist/index.js"), + completion: path.join(rootDir, "packages/shared/dist/.paperclip-build-complete"), sourceDir: path.join(rootDir, "packages/shared/src"), tsconfig: path.join(rootDir, "packages/shared/tsconfig.json"), + dependencies: [], }, { name: "@paperclipai/plugin-sdk", output: path.join(rootDir, "packages/plugins/sdk/dist/index.js"), + completion: path.join(rootDir, "packages/plugins/sdk/dist/.paperclip-build-complete"), sourceDir: path.join(rootDir, "packages/plugins/sdk/src"), tsconfig: path.join(rootDir, "packages/plugins/sdk/tsconfig.json"), + dependencies: [0], }, ]; @@ -31,102 +37,212 @@ if (!fs.existsSync(tscCliPath)) { throw new Error(`TypeScript CLI not found at ${tscCliPath}`); } -function newestSourceMtimeMs(sourceDir) { - let newest = 0; - +function directoryFingerprint(directory, exclude) { + const hash = createHash("sha256"); function visit(dir) { - for (const entry of fs.readdirSync(dir, { withFileTypes: true })) { + for (const entry of fs.readdirSync(dir, { withFileTypes: true }).sort((a, b) => a.name.localeCompare(b.name))) { const entryPath = path.join(dir, entry.name); + if (entryPath === exclude) continue; if (entry.isDirectory()) { visit(entryPath); - continue; + } else if (entry.isFile()) { + const content = fs.readFileSync(entryPath); + hash.update(JSON.stringify([path.relative(directory, entryPath), content.length])); + hash.update(content); } - if (!/\.(tsx?|json)$/.test(entry.name)) continue; - newest = Math.max(newest, fs.statSync(entryPath).mtimeMs); } } + visit(directory); + return hash.digest("hex"); +} - visit(sourceDir); - return newest; +function sourceFingerprint(target) { + const hash = createHash("sha256"); + hash.update(directoryFingerprint(target.sourceDir)); + for (const config of [ + target.tsconfig, + path.join(path.dirname(target.tsconfig), "package.json"), + path.join(rootDir, "tsconfig.json"), + path.join(rootDir, "tsconfig.base.json"), + path.join(rootDir, "node_modules/typescript/package.json"), + ]) { + if (fs.existsSync(config)) hash.update(fs.readFileSync(config)); + } + for (const dependency of target.dependencies) hash.update(sourceFingerprint(buildTargets[dependency])); + return hash.digest("hex"); +} + +function outputFingerprint(target) { + return directoryFingerprint(path.dirname(target.output), target.completion); } function needsBuild(target) { if (!fs.existsSync(target.output)) return true; - const outputMtime = fs.statSync(target.output).mtimeMs; - return newestSourceMtimeMs(target.sourceDir) > outputMtime; + try { + const completed = JSON.parse(fs.readFileSync(target.completion, "utf8")); + // Content fingerprints detect partial direct builds even on filesystems + // with coarse timestamps, while identical successful direct builds reuse + // the certified output without another compile. + return completed.sources !== sourceFingerprint(target) + || completed.outputs !== outputFingerprint(target); + } catch (error) { + if (error.code === "ENOENT" || error instanceof SyntaxError) return true; + throw error; + } } function allOutputsCurrent() { return buildTargets.every((target) => !needsBuild(target)); } -function sleep(ms) { - Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, ms); -} - -function waitForLockRelease() { - const startedAt = Date.now(); - while (Date.now() - startedAt < lockTimeoutMs) { - if (!fs.existsSync(lockDir)) { - return; - } - if (allOutputsCurrent()) { - return; - } - sleep(lockPollMs); - } - - throw new Error(`Timed out waiting for plugin build dependency lock at ${lockDir}`); -} - -if (allOutputsCurrent()) { - process.exit(0); -} - -fs.mkdirSync(path.dirname(lockDir), { recursive: true }); - +// Publish an already-populated directory so another contender never mistakes a +// newly acquired lock for an abandoned, ownerless lock. Never recursively remove +// the shared path: another process may have acquired it since we last read it. +const ownerFile = `owner-${process.pid}-${randomUUID()}.json`; +let child = null; +let stoppingSignal = null; let holdsLock = false; -let exitCode = 0; -try { + +function processAlive(pid) { + if (!Number.isInteger(pid) || pid <= 0) return true; try { - fs.mkdirSync(lockDir); - holdsLock = true; + process.kill(pid, 0); + return true; } catch (error) { - if (error && typeof error === "object" && "code" in error && error.code === "EEXIST") { - waitForLockRelease(); - if (!allOutputsCurrent()) { - throw new Error("Plugin build dependency lock released before all outputs were created"); - } - process.exit(0); - } + return error.code !== "ESRCH"; + } +} + +function removeOwner(file) { + try { + fs.unlinkSync(path.join(lockDir, file)); + } catch (error) { + if (error.code === "ENOENT") return; throw error; } + try { + fs.rmdirSync(lockDir); + } catch (error) { + if (!["ENOENT", "ENOTEMPTY", "EEXIST"].includes(error.code)) throw error; + } +} - for (const target of buildTargets) { - if (!needsBuild(target)) { - continue; +function releaseLock() { + if (!holdsLock) return; + removeOwner(ownerFile); + holdsLock = false; +} + +function recoverAbandonedLock() { + try { + const entries = fs.readdirSync(lockDir); + if (entries.length === 0) { + // Older versions wrote no owner. Allow their bounded CLI build to finish + // before reclaiming an empty directory left by interruption or timeout. + if (Date.now() - fs.statSync(lockDir).mtimeMs < 120_000) return; + fs.rmdirSync(lockDir); + } else if (entries.length === 1 && /^owner-.*\.json$/.test(entries[0])) { + const owner = JSON.parse(fs.readFileSync(path.join(lockDir, entries[0]), "utf8")); + if (processAlive(owner.pid) || (owner.childPid && processAlive(owner.childPid))) return; + removeOwner(entries[0]); + } else { + return; } + console.log("[paperclip] Recovered abandoned workspace build lock."); + } catch (error) { + if (["ENOENT", "ENOTEMPTY", "EEXIST"].includes(error.code) || error instanceof SyntaxError) return; + throw error; + } +} - const result = spawnSync(process.execPath, [tscCliPath, "-p", target.tsconfig], { +async function acquireLock() { + fs.mkdirSync(path.dirname(lockDir), { recursive: true }); + const candidate = fs.mkdtempSync(`${lockDir}.candidate-`); + fs.writeFileSync(path.join(candidate, ownerFile), JSON.stringify({ pid: process.pid })); + const startedAt = Date.now(); + let reportedWait = false; + try { + while (!stoppingSignal) { + // Do not replace a fresh empty lock held by an older script. + recoverAbandonedLock(); + if (!fs.existsSync(lockDir)) { + try { + fs.renameSync(candidate, lockDir); + holdsLock = true; + return; + } catch (error) { + if (!["ENOTEMPTY", "EEXIST", "EPERM"].includes(error.code)) throw error; + } + } + if (!reportedWait) { + console.log(`[paperclip] Waiting for another workspace build (${lockDir})...`); + reportedWait = true; + } + if (Date.now() - startedAt >= lockTimeoutMs) { + throw new Error(`Timed out waiting for workspace build lock at ${lockDir}. Another build may still be running.`); + } + await sleep(lockPollMs); + } + } finally { + fs.rmSync(candidate, { recursive: true, force: true }); + } +} + +async function build(target) { + console.log(`[paperclip] Building ${target.name}...`); + // A hard kill bypasses cleanup. Only a completed compile may restore this + // marker, so recovery never trusts index.js emitted partway through a build. + fs.rmSync(target.completion, { force: true }); + const sources = sourceFingerprint(target); + const code = await new Promise((resolve, reject) => { + child = spawn(process.execPath, [tscCliPath, "-p", target.tsconfig], { cwd: rootDir, stdio: "inherit", }); + // A hard-killed parent must not let a successor race its surviving compiler. + fs.writeFileSync(path.join(lockDir, ownerFile), JSON.stringify({ pid: process.pid, childPid: child.pid })); + child.once("error", (error) => { + fs.rmSync(target.output, { force: true }); + reject(error); + }); + child.once("close", (code) => { + child = null; + resolve(code ?? 1); + }); + }); + // tsc emits index.js before it finishes the package. A failed or interrupted + // compile must not make the next startup accept that partial build as current. + if (code !== 0) fs.rmSync(target.output, { force: true }); + else fs.writeFileSync(target.completion, JSON.stringify({ sources, outputs: outputFingerprint(target) }) + "\n"); + return code; +} - if (result.error) { - throw result.error; - } +if (allOutputsCurrent() && !fs.existsSync(lockDir)) { + process.exit(0); +} - if (result.status !== 0) { - exitCode = result.status ?? 1; - break; +// Keep the lock until the compiler has stopped, including when the foreground +// CLI's build timeout terminates this helper. +for (const signal of ["SIGINT", "SIGTERM"]) { + process.on(signal, () => { + stoppingSignal = signal; + child?.kill(signal); + }); +} +process.once("exit", releaseLock); + +let exitCode = 0; +try { + await acquireLock(); + if (holdsLock) { + for (const target of buildTargets) { + if (stoppingSignal) break; + if (!needsBuild(target)) continue; + exitCode = await build(target); + if (exitCode !== 0) break; } } } finally { - if (holdsLock) { - fs.rmSync(lockDir, { recursive: true, force: true }); - } -} - -if (exitCode !== 0) { - process.exit(exitCode); + releaseLock(); } +process.exitCode = stoppingSignal === "SIGINT" ? 130 : stoppingSignal ? 143 : exitCode; From f12b647ae865c5dd93e54e1be64648574fb46609 Mon Sep 17 00:00:00 2001 From: Dotta <34892728+cryppadotta@users.noreply.github.com> Date: Fri, 11 Sep 2026 18:26:59 -0500 Subject: [PATCH 2/2] fix: reliably interrupt and resume legacy message queues (#13275) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - A task can collect more messages while its agent works. > - Legacy runners must stop the active process before they can receive those messages. > - The old Interrupt action cancelled the run but could leave the queue idle and hidden. > - Codex could also classify a cancelled run as successful or start a fresh process after cancellation. > - This pull request joins cancellation, preserves the provider session, and dispatches the current queue after cleanup. > - The benefit is reliable interruption with the saved message order, edits, and deletions. ## Linked Issues or Issue Description **What happened?** Interrupt could strand a legacy message queue. The UI could hide pending messages after the run stopped. A Codex signal exit could race the cancellation write. A stale session warning could also trigger a fresh process after an interrupted resume. **Expected behavior** Interrupt stops the active turn and sends the remaining messages once, in their saved order. Deleted messages stay deleted. An interrupted Codex turn keeps its session and does not restart itself. **Steps to reproduce** 1. Assign a task to a legacy Codex agent that runs a long command. 2. Queue three messages. Edit one, discard another, and move the last message first. 3. Click Interrupt in the queue. 4. Repeat the interruption while the resumed session runs another command. Related work: Refs #13160, which moves native queue steering into the wake-queue module. This change fixes legacy interruption and keeps native steering unchanged. ## What Changed - Add a revision-checked, company-scoped endpoint for legacy queue interruption. - Promote only the requested queue after the provider stops and releases its lease. Retry its persisted interrupt intent from the scheduler after a promotion error or server restart. - Keep pending legacy queues visible after a run stops. Use server state for the interrupt result. - Serialize owned process cancellation before classifying the adapter result. Preserve late session and log metadata. Acknowledge cancellation only when an actual process or process group was owned; scheduler placeholders retain their normal release policy. - Send Ctrl-C to legacy Codex. Prevent missing-session fallback once the session has started. - Add cancellation race, multi-actor queue order, durable retry, resume fallback, and stale request regression tests. Document the behavior. ## Verification - Real browser tests passed with legacy Codex CLI and ACP engines, using Codex 0.153.4 and gpt-5.6-sol. - All three automated ACP browser scenarios passed locally: immediate Interrupt delivery, no replay of an unfinished write, and pause requiring Resume. Updated the old test expectation that required a separate “go” after Interrupt. - Browser tests covered queued edits, deletion, reordering, deleting the final message, and repeated interruption. - Two consecutive CLI interrupts kept one provider session. Both stopped processes exited. The final message arrived once. - `pnpm -r typecheck` passed. - `pnpm check:token-gates` passed. - All 346 post-review scheduling, recovery, queue-route, archived-company, worktree-suppression, and stale-queue regression tests passed. - All 318 process-recovery and durable-chat tests passed after the final cancellation guard. - Codex adapter, queue UI, issue-page, and OpenAPI contract tests passed. - `pnpm build` passed. - Full local suite coverage completed with `PAPERCLIP_IN_WORKTREE=false`, using the stable runner and its CI shards: 618 general server suites, all 145 serialized server suites, and all workspace groups. Every failing suite passed a targeted rerun after the fixes, rebuilding the native test fixture, correcting macOS temporary-path setup, or retrying setup/timing failures. Existing skips remain. - The original monolithic run reported failures before the final fixes; its failed suites were rerun rather than rerunning all 618 suites again. The final process-recovery/durable-chat regression run passed all 318 tests. - All CI checks passed for `e30eaf787f23a5511a3cb3cdb5abbccab9ed001d`: [run 34654820774, attempt 2](https://github.com/paperclipai/paperclip/actions/runs/34654820774/attempts/2), including typecheck, build, all test shards, E2E, and canary. The signoff and Cursor sandbox tests each hit a timeout in the initial attempt; both suites passed locally, and both failed shards passed their single CI rerun. All three corrected ACP browser scenarios passed in CI. - Greptile reviewed `e30eaf787f23a5511a3cb3cdb5abbccab9ed001d`: 5/5, no open review threads. ## Risks Cancellation order affects local adapters. The tests cover signal exits, graceful exits, adapter exceptions, termination errors, and cancellation write errors. Embedded adapters keep their cancellation controls. Ordinary run cancellation and task pause keep their distinct queue policies. No database migration is required. ## Model Used OpenAI Codex, GPT-6, with reasoning, tool use, browser testing, and code execution. The exact serving model ID and context-window size are not exposed in this session. The live test runner used OpenAI gpt-5.6-sol. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip --- doc/execution-semantics.md | 2 + .../codex-local/src/server/execute.ts | 4 + .../src/__tests__/codex-local-execute.test.ts | 33 ++++ .../heartbeat-process-recovery.test.ts | 152 ++++++++++++++++-- .../issue-queued-comments-routes.test.ts | 21 +++ .../modules/wake-queue/adapters/postgres.ts | 23 ++- server/src/routes/issues.ts | 41 +++++ server/src/routes/openapi.ts | 25 +++ server/src/services/heartbeat.ts | 84 ++++++++-- tests/e2e/acp-stop-continuation.spec.ts | 10 +- ui/src/api/issues.ts | 4 + .../task-chat/TaskChatQueuedMessages.test.tsx | 2 +- .../task-chat/TaskChatQueuedMessages.tsx | 4 +- ui/src/pages/IssueDetail.test.tsx | 24 ++- ui/src/pages/IssueDetail.tsx | 118 ++------------ 15 files changed, 392 insertions(+), 155 deletions(-) diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index 2ecf5fe725..34794de52d 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -357,6 +357,8 @@ A board comment can be an interrupt, an ownership change, both, or neither. Pape An interrupt stops the current live execution path for the issue. It does not, by itself, select the next owner. If an active run is interrupted by the board, the run may still terminate with the underlying `cancelled` status, but the issue activity and wake context should make the operator intent visible as an interruption rather than an unexplained runtime failure. +For legacy runners, **Interrupt** on a queued message stops the active run and explicitly continues the pending queue after execution cleanup. It validates the queue revision and target run, then dispatches the requested queue’s current message bodies in their saved order. Other actors’ queues cannot consume that interrupt. The persisted interrupt intent is retried by the scheduler after a promotion error or server restart until that queue is dispatched or discarded. Edits and discards remain authoritative until dispatch; deleting the final message must not create an empty continuation. Pending messages remain visible after a run stops. Cancelling only the run preserves the queue for a later explicit wake; pausing the task retains its separate queue-cancellation behavior. Native same-turn steering keeps its separate acknowledgement protocol. Legacy Codex uses Ctrl-C to stop its tool sessions and cannot retry a missing-session fallback after the provider has confirmed that the session started. + An ownership change selects who owns the issue after the comment is committed: - setting `assigneeAgentId` makes the named agent the owner diff --git a/packages/adapters/codex-local/src/server/execute.ts b/packages/adapters/codex-local/src/server/execute.ts index 70152abe3f..39c0d05560 100644 --- a/packages/adapters/codex-local/src/server/execute.ts +++ b/packages/adapters/codex-local/src/server/execute.ts @@ -1546,6 +1546,10 @@ export async function execute(ctx: AdapterExecutionContext): Promise { } }); + it.each([true, false])("retries missing resume only before a session starts (started=%s)", async (started) => { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "paperclip-codex-resume-stop-")); + const commandPath = path.join(root, "codex"); + const attemptsPath = path.join(root, "attempts"); + await seedSharedCodexAuth(root); + await fs.writeFile(commandPath, `#!/usr/bin/env node +const fs = require("node:fs"); +fs.appendFileSync(${JSON.stringify(attemptsPath)}, "attempt\\n"); +if (process.argv.includes("resume")) { + console.error("state db missing rollout path for thread unrelated-old-thread"); + ${started ? 'console.log(JSON.stringify({ type: "thread.started", thread_id: "existing-session" }));' : ''} + process.exitCode = 1; +} else { + console.log(JSON.stringify({ type: "thread.started", thread_id: "fresh-session" })); + console.log(JSON.stringify({ type: "turn.completed", usage: { input_tokens: 1, output_tokens: 1 } })); +} +`, "utf8"); + await fs.chmod(commandPath, 0o755); + try { + const result = await execute({ + runId: `resume-stop-${started}`, + agent: { id: "agent-1", companyId: "company-1", name: "Codex", adapterType: "codex_local", adapterConfig: { engine: "cli" } }, + runtime: { sessionId: "existing-session", sessionParams: null, sessionDisplayId: "existing-session", taskKey: null }, + config: { engine: "cli", command: commandPath, cwd: root, promptTemplate: "Test resume." }, + context: {}, onLog: async () => {}, + }); + expect((await fs.readFile(attemptsPath, "utf8")).trim().split("\n")).toHaveLength(started ? 1 : 2); + expect(result.sessionId).toBe(started ? "existing-session" : "fresh-session"); + } finally { + await fs.rm(root, { recursive: true, force: true }); + } + }); + it("classifies mid-turn harness crashes as retryable transient upstream errors", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "paperclip-codex-execute-harness-crash-")); const workspace = path.join(root, "workspace"); diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index ff517559ba..f77f756285 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -6699,6 +6699,126 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(repairWakeups).toHaveLength(0); }); + it("dispatches interrupted CLI input after the executor releases its lease", async () => { + const actualProcess = await vi.importActual("../adapters/process/execute.js"); + const { companyId, agentId, issueId, runId } = await seedRunFixture({ + runtimeMode: "legacy", adapterType: "codex_local", agentStatus: "idle", runStatus: "queued", + }); + await db.update(agents).set({ adapterConfig: { + command: process.execPath, args: ["-e", "console.log('ready');setInterval(() => {}, 1000)"], graceSec: 1, + } }).where(eq(agents.id, agentId)); + mockAdapterExecute.mockImplementationOnce((async (input: unknown) => + actualProcess.execute(input as Parameters[0])) as typeof mockAdapterExecute); + const heartbeat = heartbeatService(db); + await heartbeat.resumeQueuedRuns(); + expect(await waitForValue(async () => runningProcesses.get(runId))).toBeTruthy(); + const [comment] = await db.insert(issueComments).values({ companyId, issueId, authorUserId: "responsible-user", body: "continue" }).returning(); + const [wake] = await db.insert(agentWakeupRequests).values({ + companyId, agentId, source: "automation", reason: "issue_commented", status: "deferred_issue_execution", + requestedByActorType: "user", requestedByActorId: "responsible-user", + payload: { issueId, commentId: comment!.id, _paperclipWakeContext: { issueId, wakeReason: "issue_commented", wakeCommentIds: [comment!.id] } }, + }).returning(); + await heartbeat.cancelRun(runId, "Interrupt queued input", { + errorCode: "operator_interrupted", suppressImmediateRecovery: true, + resultJson: { operatorInterrupted: true, queuedCommentInterruptQueueId: wake!.id }, + }); + await heartbeat.drainActiveRunExecutions(); + const [updated] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wake!.id)); + expect(updated!.runId).toBeTruthy(); + expect(updated!.runId).not.toBe(runId); + expect((await heartbeat.getRun(updated!.runId!))!.contextSnapshot?.wakeCommentIds).toEqual([comment!.id]); + }); + + it("retries durable queue interruption after a promotion failure on a fresh service", async () => { + const { companyId, agentId, issueId, runId } = await seedRunFixture({ + runtimeMode: "legacy", adapterType: "codex_local", agentStatus: "idle", runStatus: "cancelled", + }); + const [comment] = await db.insert(issueComments).values({ companyId, issueId, authorUserId: "responsible-user", body: "retry this input" }).returning(); + const [wake] = await db.insert(agentWakeupRequests).values({ + companyId, agentId, source: "automation", reason: "issue_commented", status: "deferred_issue_execution", + requestedByActorType: "user", requestedByActorId: "responsible-user", + payload: { issueId, commentId: comment!.id, _paperclipWakeContext: { issueId, wakeReason: "issue_commented", wakeCommentIds: [comment!.id] } }, + }).returning(); + await db.update(heartbeatRuns).set({ resultJson: { + queuedCommentInterruptQueueId: wake!.id, + executionCancellation: { state: "acknowledged" }, + conversationContinuation: "continue_conversation_v1", + } }).where(eq(heartbeatRuns.id, runId)); + const failedPromotion = vi.spyOn(db, "transaction").mockRejectedValueOnce(new Error("temporary queue promotion outage")); + try { + await heartbeatService(db).resumeQueuedRuns(); + expect(failedPromotion).toHaveBeenCalled(); + expect((await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wake!.id)))[0]!.status).toBe("deferred_issue_execution"); + } finally { + failedPromotion.mockRestore(); + } + const restarted = heartbeatService(db); + await restarted.resumeQueuedRuns(); + await restarted.drainActiveRunExecutions(); + const [updated] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wake!.id)); + expect(updated!.runId).toBeTruthy(); + expect((await restarted.getRun(updated!.runId!))!.contextSnapshot?.wakeCommentIds).toEqual([comment!.id]); + await restarted.resumeQueuedRuns(); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId))).toHaveLength(2); + }); + + it.each(["pending", "discarded", "wrong queue"] as const)( + "resumes only the authorized %s queue after an acknowledged legacy interrupt", + async (state) => { + const { companyId, agentId, issueId, runId } = await seedRunFixture({ + runtimeMode: "legacy", adapterType: "codex_local", agentStatus: "running", + }); + const heartbeat = heartbeatService(db); + const comments = await db.insert(issueComments).values([ + { companyId, issueId, authorUserId: "responsible-user", body: "First, edited" }, + { companyId, issueId, authorUserId: "responsible-user", body: "Deleted" }, + { companyId, issueId, authorUserId: "responsible-user", body: "Third, moved first" }, + ]).returning(); + const commentIds = [comments[2]!.id, comments[0]!.id]; + const [deferred] = await db.insert(agentWakeupRequests).values({ + companyId, agentId, source: "automation", reason: "issue_commented", + status: state === "discarded" ? "cancelled" : "deferred_issue_execution", + requestedByActorType: "user", requestedByActorId: "responsible-user", + payload: { issueId, commentId: commentIds[0], _paperclipWakeContext: { + issueId, wakeReason: "issue_commented", wakeCommentIds: commentIds, + } }, + }).returning(); + // A different actor's older queue must not consume this interrupt. + const [otherComment] = await db.insert(issueComments).values({ + companyId, issueId, authorUserId: "other-user", body: "Other actor's input", + }).returning(); + await db.insert(agentWakeupRequests).values({ + companyId, agentId, source: "automation", reason: "issue_commented", + status: "deferred_issue_execution", requestedAt: new Date(0), + requestedByActorType: "user", requestedByActorId: "other-user", + payload: { issueId, commentId: otherComment!.id, _paperclipWakeContext: { + issueId, wakeReason: "issue_commented", wakeCommentIds: [otherComment!.id], + } }, + }); + await heartbeat.cancelRun(runId, "Interrupt queued messages", { + suppressImmediateRecovery: true, errorCode: "operator_interrupted", + resultJson: { + operatorInterrupted: true, + queuedCommentInterruptQueueId: state === "wrong queue" ? randomUUID() : deferred!.id, + executionCancellation: { state: "acknowledged" }, + executionRecovery: { kind: "interrupted", providerStopped: true, sessionPreserved: true, actionOutcomes: "settled" }, + }, + }); + await heartbeat.drainActiveRunExecutions(); + const runs = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId)); + const successors = runs.filter((run) => run.id !== runId) + .sort((left, right) => left.createdAt.getTime() - right.createdAt.getTime()); + expect(successors).toHaveLength(state === "pending" ? 2 : 0); + if (state === "pending") { + expect(successors[0]!.contextSnapshot?.wakeCommentIds).toEqual(commentIds); + // Only the requested turn's normal completion can drain the other queue. + expect(successors[1]!.contextSnapshot?.wakeCommentIds).toEqual([otherComment!.id]); + await heartbeat.cancelRun(runId, "Duplicate interrupt"); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId))).toHaveLength(3); + } + }, + ); + it("preserves deferred input on a clean Stop and adopts it once on the next explicit comment", async () => { const { companyId, agentId, issueId, runId } = await seedRunFixture({ runtimeMode: "legacy", agentStatus: "running" }); const heartbeat = heartbeatService(db); @@ -7090,7 +7210,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { ); }); - it.each([ + it.each(([ { mode: "signal", graceful: false, failure: null }, { mode: "graceful exit", graceful: true, failure: null }, { mode: "adapter exception", graceful: false, failure: null }, @@ -7111,9 +7231,12 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { graceful: true, failure: "write", }, - ] as const)( - "settles an owned process Stop before classifying its $mode", - async ({ mode, graceful, failure }) => { + ] as const).flatMap((scenario) => + (["process", "codex_local"] as const).map((adapterType) => ({ ...scenario, adapterType })), + ))( + "settles an owned $adapterType Stop before classifying its $mode", + async ({ mode, graceful, failure, adapterType }) => { + const stopSignal = adapterType === "codex_local" ? "SIGINT" : "SIGTERM"; const actualProcess = await vi.importActual< typeof import("../adapters/process/execute.js") >("../adapters/process/execute.js"); @@ -7163,7 +7286,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { throw new Error("owned termination unconfirmed"); }); const { runId, agentId } = await seedRunFixture({ - adapterType: "process", + adapterType, agentStatus: "idle", runStatus: "queued", includeIssue: false, @@ -7175,7 +7298,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { command: process.execPath, args: [ "-e", - `${graceful ? "process.on('SIGTERM', () => process.exit(0));" : ""} console.log('stop ready'); setInterval(() => {}, 1000)`, + `${graceful ? `process.on('${stopSignal}', () => process.exit(0));` : ""} console.log('stop ready'); setInterval(() => {}, 1000)`, ], graceSec: 1, }, @@ -7204,7 +7327,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(await waitForValue(async () => observedResult)).toMatchObject( graceful ? { exitCode: 0, signal: null } - : { exitCode: null, signal: "SIGTERM" }, + : { exitCode: null, signal: stopSignal }, ); // The process utility already removed its child record on close. A new // service instance must still join the original cancellation owner. @@ -7227,11 +7350,10 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect((await heartbeat.getRun(runId))?.status).toBe("running"); expect(duplicateSettled).toBe(false); if (failure === "write") { - writeSpy = vi - .spyOn(db, "transaction") - .mockRejectedValueOnce( - new Error("owned cancellation write unavailable"), - ); + const error = new Error("owned cancellation write unavailable"); + writeSpy = adapterType === "codex_local" + ? vi.spyOn(db, "update").mockImplementationOnce(() => { throw error; }) + : vi.spyOn(db, "transaction").mockRejectedValueOnce(error); } } finally { releaseTermination(); @@ -7446,7 +7568,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { ); expect(mockTerminateLocalService).toHaveBeenCalledWith( expect.objectContaining({ pid: 12345, processGroupId: null }), - { forceAfterMs: 1000 }, + { forceAfterMs: 1000, signal: "SIGINT" }, ); expect(runningProcesses.has(runId)).toBe(false); } finally { @@ -7479,7 +7601,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(mockTerminateLocalService).toHaveBeenCalledWith( expect.objectContaining({ pid: 12_346, processGroupId: null }), - { forceAfterMs: 2_000 }, + { forceAfterMs: 2_000, signal: "SIGINT" }, ); expect(runningProcesses.has(runId)).toBe(false); }); @@ -7515,7 +7637,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(outcome).toMatchObject({ status: "succeeded", errorCode: null }); expect(mockTerminateLocalService).toHaveBeenCalledWith( expect.objectContaining({ pid: 12_347, processGroupId: null }), - { forceAfterMs: 2_000 }, + { forceAfterMs: 2_000, signal: "SIGINT" }, ); await expect(heartbeat.getRun(runId)).resolves.toMatchObject({ status: "succeeded", diff --git a/server/src/__tests__/issue-queued-comments-routes.test.ts b/server/src/__tests__/issue-queued-comments-routes.test.ts index b2f3b45a88..7e4a70ef73 100644 --- a/server/src/__tests__/issue-queued-comments-routes.test.ts +++ b/server/src/__tests__/issue-queued-comments-routes.test.ts @@ -175,6 +175,27 @@ describeEmbeddedPostgres("issue queued-comment routes", () => { return { companyId, agentId, issueId, runId, wakeId, commentIds }; } + it.each(["stale revision", "native run", "different issue"] as const)( + "rejects queued interruption for a %s without stopping the run", + async (scenario) => { + const seeded = await seedQueue(); + if (scenario !== "native run") { + await db.update(agents).set({ adapterType: "codex_local" }).where(eq(agents.id, seeded.agentId)); + await db.update(heartbeatRuns).set({ runtimeMode: "legacy" }).where(eq(heartbeatRuns.id, seeded.runId)); + } + const client = app(seeded.companyId); + const queue = await request(client).get(`/api/issues/${seeded.issueId}/queued-comments`).expect(200); + if (scenario === "different issue") { + await db.update(heartbeatRuns).set({ contextSnapshot: { issueId: randomUUID() } }).where(eq(heartbeatRuns.id, seeded.runId)); + } + await request(client).post(`/api/issues/${seeded.issueId}/queued-comments/interrupt`).send({ + queueId: seeded.wakeId, targetRunId: seeded.runId, + revision: scenario === "stale revision" ? "stale" : queue.body.revision, + }).expect(409); + expect((await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, seeded.runId)))[0]!.status).toBe("running"); + }, + ); + async function promoteQueue(seeded: Awaited>) { const queueRunId = randomUUID(); const wake = await db diff --git a/server/src/modules/wake-queue/adapters/postgres.ts b/server/src/modules/wake-queue/adapters/postgres.ts index abeccf1b12..adc27cc9f5 100644 --- a/server/src/modules/wake-queue/adapters/postgres.ts +++ b/server/src/modules/wake-queue/adapters/postgres.ts @@ -176,6 +176,9 @@ function buildHost(_tx: Db, deps: WakeQueuePostgresAdapterDeps): WakeQueueHost { function buildTransaction(tx: Db, deps: WakeQueuePostgresAdapterDeps, db: Db, run: HeartbeatRunRow): WakeQueueTransaction { const treeControlSvc = issueTreeControlService(tx); const issuesSvc = issueService(tx); + const interruptQueueId = run.runtimeMode !== "native" && run.status === "cancelled" + ? readNonEmptyString(run.resultJson?.queuedCommentInterruptQueueId) + : null; return { async findInvokableAgent({ companyId, agentId }): Promise { @@ -199,6 +202,8 @@ function buildTransaction(tx: Db, deps: WakeQueuePostgresAdapterDeps, db: Db, ru eq(agentWakeupRequests.companyId, companyId), eq(agentWakeupRequests.status, DEFERRED_WAKE_STATUS), sql`${agentWakeupRequests.payload} ->> 'issueId' = ${issueId}`, + interruptQueueId ? eq(agentWakeupRequests.id, interruptQueueId) : undefined, + interruptQueueId ? eq(agentWakeupRequests.agentId, run.agentId) : undefined, ), ) .orderBy(asc(agentWakeupRequests.requestedAt)) @@ -1000,6 +1005,20 @@ export function createPostgresWakeQueueAdapter(db: Db, deps: WakeQueuePostgresAd const issueRow = (contextIssueId ? candidateIssues.find((candidate) => candidate.id === contextIssueId) : candidateIssues[0]) ?? null; + // A queue interrupt authorizes only its original pending queue. Replays + // after dispatch or deleting the final message cannot launch other work. + const interruptQueueId = run.runtimeMode !== "native" + ? readNonEmptyString(run.resultJson?.queuedCommentInterruptQueueId) + : null; + const [interruptedQueue] = interruptQueueId && issueRow + ? await tx.select({ id: agentWakeupRequests.id }).from(agentWakeupRequests).where(and( + eq(agentWakeupRequests.id, interruptQueueId), + eq(agentWakeupRequests.companyId, run.companyId), + eq(agentWakeupRequests.agentId, run.agentId), + eq(agentWakeupRequests.status, "deferred_issue_execution"), + sql`${agentWakeupRequests.payload}->>'issueId' = ${issueRow.id}`, + )).limit(1) + : []; const preDrainFacts: PreDrainFacts = { issueRowPresent: issueRow !== null, executionRunIdMatchesRun: !issueRow || !issueRow.executionRunId || issueRow.executionRunId === run.id, @@ -1013,7 +1032,9 @@ export function createPostgresWakeQueueAdapter(db: Db, deps: WakeQueuePostgresAd // next explicit wake adopts those messages atomically when it // queues a run. executionCancellationAcknowledged: - run.status === "cancelled" && parseObject(run.resultJson?.executionCancellation).state === "acknowledged", + run.status === "cancelled" && + parseObject(run.resultJson?.executionCancellation).state === "acknowledged" && + !interruptedQueue, }; const preDrain = decidePreDrain(preDrainFacts); diff --git a/server/src/routes/issues.ts b/server/src/routes/issues.ts index 124e06a8a8..566678cf0c 100644 --- a/server/src/routes/issues.ts +++ b/server/src/routes/issues.ts @@ -15201,6 +15201,47 @@ export function issueRoutes( }, ); + router.post( + "/issues/:id/queued-comments/interrupt", + validate(queuedCommentSteeringTargetSchema), + async (req, res) => { + assertBoard(req); + if (!req.actor.userId) throw forbidden("Board user context required"); + const issue = await getAccessibleResource(req, res, svc.getById(req.params.id as string), "Issue not found"); + if (!issue) return; + const actor = getActorInfo(req); + await db.transaction(async (tx) => { + const locked = await lockQueuedCommentState({ + tx, issue, actor, queueId: req.body.queueId, targetRunId: req.body.targetRunId, + }); + assertQueueMutationTarget({ queue: locked.queue, queueId: req.body.queueId, revision: req.body.revision }); + if (locked.queue.protocol !== "legacy" || locked.activeRun?.agentId !== issue.assigneeAgentId) { + throw conflict("This queue does not support legacy interruption"); + } + }); + // Never hold the issue lock while joining the adapter. Queue edits and + // discards stay authoritative until the dispatcher claims the successor. + const options = operatorInterruptCancelOptions({ issueId: issue.id, actor }); + await heartbeat.cancelRun(req.body.targetRunId, "Interrupted to send queued messages", { + ...options, + suppressImmediateRecovery: true, + resultJson: { ...options.resultJson, queuedCommentInterruptQueueId: req.body.queueId }, + }); + await logActivity(db, { + companyId: issue.companyId, actorType: actor.actorType, actorId: actor.actorId, + agentId: actor.agentId, runId: actor.runId, agentApiKeyId: actor.agentApiKeyId, + action: "issue.queued_comments_interrupted", entityType: "issue", entityId: issue.id, + details: { queueId: req.body.queueId, targetRunId: req.body.targetRunId }, + }); + const currentIssue = await svc.getById(issue.id); + const queue = await buildQueuedCommentQueue({ + executor: db, issue: currentIssue ?? issue, + activeRun: await resolveActiveIssueRun(currentIssue ?? issue), actor, + }); + res.json(await runRedactions.redactForIssue(issue.companyId, issue.id, queue)); + }, + ); + router.post( "/issues/:id/queued-comments/:commentId/steer", validate(queuedCommentSteeringTargetSchema), diff --git a/server/src/routes/openapi.ts b/server/src/routes/openapi.ts index 3e72e26186..1028baf59a 100644 --- a/server/src/routes/openapi.ts +++ b/server/src/routes/openapi.ts @@ -6527,6 +6527,31 @@ registry.registerPath({ }, }); +registry.registerPath({ + method: "post", + path: "/api/issues/{id}/queued-comments/interrupt", + tags: ["issues"], + summary: "Interrupt the active legacy run and continue its queued comments", + request: { + params: z.object({ id: z.string() }), + body: jsonBody( + z.object({ + queueId: z.string().min(1), + revision: z.string().min(1), + targetRunId: z.string().min(1), + }), + ), + }, + responses: { + 200: r.ok(), + 400: r.badRequest, + 401: r.unauthorized, + 403: r.forbidden, + 404: r.notFound, + 409: r.conflict, + }, +}); + registry.registerPath({ method: "post", path: "/api/issues/{id}/queued-comments/{commentId}/steer", diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 2825c56d71..84289b5753 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -1210,11 +1210,11 @@ const SESSIONED_LOCAL_ADAPTERS = new Set([ // Routes and the scheduler construct separate heartbeatService instances, but // they must agree on in-process adapter executions when reaping stale runs. const activeRunExecutions = new Set(); -// A process adapter's signal exit can race the operator cancellation CAS while +// A legacy process adapter's signal exit can race the operator cancellation CAS while // its owned process group is still being joined. Keep that exit from becoming // a successful result (or a competing failure) before Stop settles. This is an // in-process ordering barrier, not durable cancellation or provider authority. -// Other adapters can have independently proven terminal results after a signal. +// Embedded adapters use their own cancellation control and acknowledgement. const processRunCancellationSettlements = new Map< string, { @@ -8631,6 +8631,7 @@ async function terminateHeartbeatRunProcess(input: { pid: number | null | undefined; processGroupId: number | null | undefined; graceMs?: number; + signal?: NodeJS.Signals; }) { const pid = input.pid ?? null; const processGroupId = input.processGroupId ?? null; @@ -8649,7 +8650,7 @@ async function terminateHeartbeatRunProcess(input: { ? processGroupId : null, }, - input.graceMs ? { forceAfterMs: input.graceMs } : undefined, + { forceAfterMs: input.graceMs, signal: input.signal }, ); } @@ -18577,6 +18578,31 @@ export function heartbeatService( await resumeExecutionWaitComments(); const cutoff = await getWorktreeExecutionCutoff(); + // The cancellation marker is durable intent. Retry while its exact queue + // is still deferred, including after a failed cleanup promotion or restart. + // Normal admission still checks process ownership, leases, pauses, and scope. + const interruptedQueues = await db + .select({ id: heartbeatRuns.id, companyId: heartbeatRuns.companyId }) + .from(agentWakeupRequests) + .innerJoin(heartbeatRuns, and( + sql`${heartbeatRuns.resultJson}->>'queuedCommentInterruptQueueId' = ${agentWakeupRequests.id}::text`, + eq(heartbeatRuns.companyId, agentWakeupRequests.companyId), + eq(heartbeatRuns.agentId, agentWakeupRequests.agentId), + )) + .innerJoin(companies, eq(companies.id, heartbeatRuns.companyId)) + .where(and( + eq(agentWakeupRequests.status, "deferred_issue_execution"), + eq(heartbeatRuns.status, "cancelled"), + eq(heartbeatRuns.runtimeMode, "legacy"), + eq(companies.status, "active"), + cutoff ? gte(heartbeatRuns.createdAt, cutoff) : undefined, + )); + for (const run of interruptedQueues) { + await releaseIssueExecutionAndPromote(run, { suppressImmediateRecovery: true }).catch((err) => { + logger.error({ err, runId: run.id }, "failed to retry interrupted comment queue"); + }); + } + const queuedRuns = await db .select({ agentId: heartbeatRuns.agentId }) .from(heartbeatRuns) @@ -23562,10 +23588,8 @@ export function heartbeatService( } } const processCancellation = - agent.adapterType === "process" - ? (processRunCancellationSettlements.get(run.id) ?? - failedProcessRunCancellations.get(run.id)) - : undefined; + processRunCancellationSettlements.get(run.id) ?? + failedProcessRunCancellations.get(run.id); await processCancellation?.settled; let outcome: RunSessionOutcome; const latestRun = await getRun(run.id); @@ -23587,7 +23611,7 @@ export function heartbeatService( } else if ( (adapterResult.exitCode ?? 0) === 0 && !adapterResult.errorMessage && - !(agent.adapterType === "process" && adapterResult.signal) && + !adapterResult.signal && !processCancellation?.failed ) { outcome = "succeeded"; @@ -23784,9 +23808,11 @@ export function heartbeatService( // adapter's semantic result, usage, logs, or presentation decision. // Only complete the late metadata write when the reconciler chose the // same terminal status; a conflicting terminal outcome remains owned - // by the path that won the compare-and-set. + // by the path that won the compare-and-set. Owned legacy cancellation + // likewise keeps the provider session, logs, and usage after Stop wins. if ( - adapterResult.nativeFinalization && + (adapterResult.nativeFinalization || + (processCancellation && !processCancellation.failed && status === "cancelled")) && persistedRunWrite.run?.status === status ) { persistedRun = await db @@ -24279,9 +24305,7 @@ export function heartbeatService( } // A process adapter may throw while its owned Stop is joining the // child. Let the cancellation write settle before attempting failure. - if (agent.adapterType === "process") { - await processRunCancellationSettlements.get(run.id)?.settled; - } + await processRunCancellationSettlements.get(run.id)?.settled; const message = redactCurrentUserText( err instanceof Error ? err.message : "Unknown adapter failure", await getCurrentUserRedactionOptions(), @@ -24886,6 +24910,18 @@ export function heartbeatService( }); } } + // Interrupting a queued message explicitly authorizes the pending queue. + // Retry its normal promotion after leases and adapter cleanup have settled; + // the earlier terminal write can still have an execution blocker here. + if ( + latestRun?.status === "cancelled" && + latestRun.runtimeMode !== "native" && + readNonEmptyString(latestRun.resultJson?.queuedCommentInterruptQueueId) + ) { + await releaseIssueExecutionAndPromote(latestRun, { suppressImmediateRecovery: true }).catch((err) => { + logger.error({ err, runId: run.id }, "failed to promote interrupted comment queue after cleanup"); + }); + } activeRunExecutions.delete(run.id); // A failed owned Stop remains visible until this exact executor settles, // including a graceful exit result arriving after the cancellation error. @@ -24909,7 +24945,7 @@ export function heartbeatService( } async function releaseIssueExecutionAndPromote( - run: typeof heartbeatRuns.$inferSelect, + run: Pick, options: { suppressImmediateRecovery?: boolean } = {}, ) { try { @@ -27468,8 +27504,8 @@ export function heartbeatService( try { let releaseProcessCancellation: (() => void) | undefined; const processCancellationSettlement = - agent?.adapterType === "process" && run.runtimeMode !== "native" && + !control && running ? { settled: new Promise((resolve) => { @@ -27528,6 +27564,9 @@ export function heartbeatService( await terminateHeartbeatRunProcess({ pid: running.child.pid, processGroupId: running.processGroupId, + // Codex handles Ctrl-C by cancelling its tool sessions. SIGTERM + // can leave commands in their separate process groups alive. + signal: !control && agent?.adapterType === "codex_local" ? "SIGINT" : undefined, graceMs: cancellationTerminationGraceMs( running.graceSec, options.terminationGraceMs, @@ -27585,6 +27624,21 @@ export function heartbeatService( resultJson: { ...persistedCancellationResult, ...(resultJson ?? {}), + // A scheduler placeholder has no process to acknowledge. + // Preserve its normal release policy instead of treating + // it as an operator stop of provider work. + ...(processCancellationSettlement && agent && running && ( + (Number.isInteger(running.child.pid) && (running.child.pid ?? 0) > 0) || + (Number.isInteger(running.processGroupId) && (running.processGroupId ?? 0) > 0) + ) + ? mergeRunStopMetadataForAgent(agent, "cancelled", { + resultJson: { + ...resultJson, + executionCancellation: { state: "acknowledged", acknowledgedAt: finishedAt.toISOString() }, + }, + errorCode, errorMessage: reason, + }) + : {}), // The native cancellation helper may have advanced a durable // pending intent to its acknowledged state after `run` was // first read. Never let that stale snapshot overwrite the diff --git a/tests/e2e/acp-stop-continuation.spec.ts b/tests/e2e/acp-stop-continuation.spec.ts index c23ee8e3c1..4384ef7f02 100644 --- a/tests/e2e/acp-stop-continuation.spec.ts +++ b/tests/e2e/acp-stop-continuation.spec.ts @@ -10,7 +10,7 @@ async function json(response: APIResponse) { } for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false }, { unfinishedWrite: true, pause: false }, { unfinishedWrite: false, pause: true }]) { - test(`embedded ACP Stop: ${unfinishedWrite ? "unknown action continues without replaying the write" : pause ? "composer pause requires Resume before continuation" : "go continues the same session with queued input"}`, async ({ page, request }) => { + test(`embedded ACP Stop: ${unfinishedWrite ? "Interrupt continues without replaying the write" : pause ? "composer pause requires Resume before continuation" : "Interrupt delivers queued input in the same session"}`, async ({ page, request }) => { test.setTimeout(120_000); const root = await mkdtemp(path.join(os.tmpdir(), "paperclip-stop-browser-")); const company = await json(await request.post("/api/companies", { data: { name: `ACP Stop ${Date.now()}` } })); @@ -38,7 +38,7 @@ for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false await expect.poll(async () => JSON.stringify(await json(await request.get(`/api/issues/${issue.id}/queued-comments`)))) .toContain("List my recent Drive files."); - // Run-level Stop leaves the task unpaused; composer Stop additionally pauses the task. + // Interrupt sends the queue immediately; composer Stop pauses the task. let stopped; if (pause) { await page.getByRole("button", { name: "Stop", exact: true }).click(); @@ -65,19 +65,15 @@ for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false const dialog = page.getByRole("dialog"); await dialog.getByRole("checkbox").check(); await dialog.getByRole("button", { name: "Resume work", exact: true }).click(); - } else { - await editor.fill("go"); - await page.getByRole("button", { name: "Send", exact: true }).click(); } await expect(page.getByText("Answered the pending follow-up once.", { exact: false })).toBeVisible({ timeout: 30_000 }); await expect.poll(async () => (await json(await request.get(`/api/issues/${issue.id}/live-runs`))).length).toBe(0); const prompts = (await readFile(path.join(root, "prompts"), "utf8")).trim().split("\n").map(line => JSON.parse(line)); expect(prompts).toHaveLength(2); expect(new Set(prompts.map(prompt => prompt.sessionId)).size).toBe(1); - // Resume delivers the queued follow-up in the same provider session. + // Interrupt or Resume delivers the queued follow-up without another message. const continuationPrompts = pause ? prompts.slice(1) : [prompts.at(-1)]; expect(JSON.stringify(continuationPrompts)).toContain("List my recent Drive files."); - if (!pause) expect(JSON.stringify(continuationPrompts)).toContain("go"); expect(await readFile(path.join(root, "completed"), "utf8")).toBe("follow-up\n"); const completedIssue = await json(await request.get(`/api/issues/${issue.id}`)); expect(completedIssue.executionBlocker).toBeNull(); diff --git a/ui/src/api/issues.ts b/ui/src/api/issues.ts index 26110545d8..d3c2dca602 100644 --- a/ui/src/api/issues.ts +++ b/ui/src/api/issues.ts @@ -384,6 +384,10 @@ export const issuesApi = { `/issues/${id}/queued-comments/order`, data, ), + interruptQueuedComments: ( + id: string, + data: { queueId: string; targetRunId: string; revision: string }, + ) => api.post(`/issues/${id}/queued-comments/interrupt`, data), steerQueuedComment: ( id: string, commentId: string, diff --git a/ui/src/components/task-chat/TaskChatQueuedMessages.test.tsx b/ui/src/components/task-chat/TaskChatQueuedMessages.test.tsx index 755a6754b7..16e03e8457 100644 --- a/ui/src/components/task-chat/TaskChatQueuedMessages.test.tsx +++ b/ui/src/components/task-chat/TaskChatQueuedMessages.test.tsx @@ -324,7 +324,7 @@ describe("TaskChatQueuedMessages", () => { ), ).not.toBeNull(); expect(container.textContent).toContain( - "Active turn interrupted. Message remains queued.", + "Interruption requested. Queued messages will continue after the active turn stops.", ); }); }); diff --git a/ui/src/components/task-chat/TaskChatQueuedMessages.tsx b/ui/src/components/task-chat/TaskChatQueuedMessages.tsx index 725ef7ee4c..286d64b976 100644 --- a/ui/src/components/task-chat/TaskChatQueuedMessages.tsx +++ b/ui/src/components/task-chat/TaskChatQueuedMessages.tsx @@ -147,7 +147,7 @@ function SortableQueuedMessage({ type="button" onClick={onInterrupt} disabled={busy || !queue.targetRunId || !onInterrupt} - title="Interrupt the active turn; this message stays queued" + title="Interrupt the active turn and send queued messages" className="flex h-7 shrink-0 items-center gap-1.5 rounded-md px-2 text-xs font-medium text-muted-foreground transition-colors hover:bg-accent hover:text-foreground disabled:opacity-40" data-testid={`task-chat-queued-interrupt-${entry.comment.id}`} > @@ -334,7 +334,7 @@ export function TaskChatQueuedMessages({ action === "steer" ? "Message steered into the active turn." : action === "interrupt" - ? "Active turn interrupted. Message remains queued." + ? "Interruption requested. Queued messages will continue after the active turn stops." : "Queued message discarded.", ); } catch (error) { diff --git a/ui/src/pages/IssueDetail.test.tsx b/ui/src/pages/IssueDetail.test.tsx index 439732e8bb..12e889590f 100644 --- a/ui/src/pages/IssueDetail.test.tsx +++ b/ui/src/pages/IssueDetail.test.tsx @@ -52,6 +52,7 @@ const mockIssuesApi = vi.hoisted(() => ({ listFeedbackVotes: vi.fn(), listInteractions: vi.fn(), getQueuedComments: vi.fn(), + interruptQueuedComments: vi.fn(), editQueuedComment: vi.fn(), reorderQueuedComments: vi.fn(), steerQueuedComment: vi.fn(), @@ -1320,6 +1321,7 @@ describe("IssueDetail", () => { entries: [], }), ); + mockIssuesApi.interruptQueuedComments.mockReset().mockResolvedValue(createQueuedCommentQueue()); mockIssuesApi.editQueuedComment.mockResolvedValue( createQueuedCommentQueue(), ); @@ -3674,14 +3676,19 @@ describe("IssueDetail", () => { body: "Queued run message", }); + mockIssuesApi.getQueuedComments.mockResolvedValue(createQueuedCommentQueue({ + targetRunId: "run-queued", protocol: "legacy", steeringDisposition: "unsupported", + })); await act(async () => { await persistedProps.onInterruptQueued( persistedComment!.queueTargetRunId!, ); }); - expect(mockHeartbeatsApi.cancel).toHaveBeenCalledWith("run-queued"); - mockHeartbeatsApi.cancel.mockClear(); + expect(mockIssuesApi.interruptQueuedComments).toHaveBeenCalledWith("PAP-1", { + queueId: "wake-queue-1", revision: "queue-revision-1", targetRunId: "run-queued", + }); + expect(mockHeartbeatsApi.cancel).not.toHaveBeenCalled(); }); it("projects a native follow-up into the steering well before the post resolves", async () => { @@ -3876,15 +3883,16 @@ describe("IssueDetail", () => { queueTargetRunId: "run-original", }); + mockIssuesApi.getQueuedComments.mockResolvedValue(createQueuedCommentQueue({ + targetRunId: "run-replacement", protocol: "legacy", steeringDisposition: "unsupported", + })); await act(async () => { - await replacementProps.onInterruptQueued( + await expect(replacementProps.onInterruptQueued( optimisticComment!.queueTargetRunId!, - ); + )).rejects.toThrow("The queued messages changed"); }); - expect(mockHeartbeatsApi.cancel).toHaveBeenCalledWith("run-original"); - expect(mockHeartbeatsApi.cancel).not.toHaveBeenCalledWith( - "run-replacement", - ); + expect(mockIssuesApi.interruptQueuedComments).not.toHaveBeenCalled(); + expect(mockHeartbeatsApi.cancel).not.toHaveBeenCalled(); await act(async () => { postedComment.resolve( diff --git a/ui/src/pages/IssueDetail.tsx b/ui/src/pages/IssueDetail.tsx index a1cd68e001..76d6d9933d 100644 --- a/ui/src/pages/IssueDetail.tsx +++ b/ui/src/pages/IssueDetail.tsx @@ -1485,7 +1485,7 @@ const IssueDetailChatTab = memo(function IssueDetailChatTab({ const queuedCommentQueueEnabled = !classicTaskInterfaceEnabled && runtimeSelectionKnown && - Boolean(liveRuntimeRun || assigneeUsesPaperclipRunner); + Boolean(liveRuntimeRun || issueAssigneeAgentId); const { data: authoritativeQueuedCommentQueue } = useQuery({ queryKey: queryKeys.issues.queuedComments(issueId), queryFn: async () => @@ -1494,7 +1494,8 @@ const IssueDetailChatTab = memo(function IssueDetailChatTab({ issueId, ), enabled: queuedCommentQueueEnabled, - refetchInterval: queuedCommentQueueEnabled ? 1000 : false, + refetchInterval: (query) => queuedCommentQueueEnabled && + (liveRuntimeRun || query.state.data?.entries.length) ? 1000 : false, }); const [consumedQueuedCommentIds, setConsumedQueuedCommentIds] = useState< ReadonlySet @@ -4925,93 +4926,14 @@ export function IssueDetail({ tasksTab }: { tasksTab?: TaskSidePanelProps["tasks }); const interruptQueuedComment = useMutation({ - mutationFn: (runId: string) => heartbeatsApi.cancel(runId), - onMutate: async (runId) => { - await Promise.all( - issueCacheRefs.flatMap((ref) => [ - queryClient.cancelQueries({ queryKey: queryKeys.issues.runs(ref) }), - queryClient.cancelQueries({ - queryKey: queryKeys.issues.liveRuns(ref), - }), - queryClient.cancelQueries({ - queryKey: queryKeys.issues.activeRun(ref), - }), - queryClient.cancelQueries({ queryKey: queryKeys.issues.detail(ref) }), - ]), - ); - - const previousRunState = issueCacheRefs.map((ref) => ({ - ref, - runs: queryClient.getQueryData( - queryKeys.issues.runs(ref), - ), - liveRuns: queryClient.getQueryData( - queryKeys.issues.liveRuns(ref), - ), - activeRun: queryClient.getQueryData( - queryKeys.issues.activeRun(ref), - ), - issue: queryClient.getQueryData(queryKeys.issues.detail(ref)), - })); - const previousLocalQueuedCommentRunIds = locallyQueuedCommentRunIds; - const cachedActiveRun = - previousRunState.find((state) => state.activeRun?.id === runId) - ?.activeRun ?? - previousRunState.find((state) => state.activeRun)?.activeRun ?? - null; - const liveRunList = dedupeLiveRunsById( - previousRunState.flatMap((state) => state.liveRuns ?? []), - ); - const interruptibleIssueRun = resolveInterruptibleIssueRun( - cachedActiveRun, - liveRunList, - ); - const targetRun = - cachedActiveRun?.id === runId - ? cachedActiveRun - : (liveRunList?.find((run) => run.id === runId) ?? - interruptibleIssueRun ?? - null); - - if (targetRun) { - const interruptedAt = new Date().toISOString(); - for (const ref of issueCacheRefs) { - queryClient.setQueryData( - queryKeys.issues.runs(ref), - (current) => - upsertInterruptedRun(current, targetRun, interruptedAt), - ); - } + mutationFn: async (runId: string) => { + const queue = await issuesApi.getQueuedComments(issueId!); + if (!queue.queueId || queue.targetRunId !== runId) { + throw new Error("The queued messages changed. Refresh and try again."); } - - for (const ref of issueCacheRefs) { - queryClient.setQueryData( - queryKeys.issues.liveRuns(ref), - (current: LiveRunForIssue[] | undefined) => - removeLiveRunById(current, runId), - ); - queryClient.setQueryData( - queryKeys.issues.activeRun(ref), - (current: ActiveRunForIssue | null | undefined) => - current?.id === runId ? null : current, - ); - queryClient.setQueryData( - queryKeys.issues.detail(ref), - (current: Issue | undefined) => - clearIssueExecutionRun(current, runId), - ); - } - setLocallyQueuedCommentRunIds((current) => { - const next = new Map( - [...current].filter(([, targetRunId]) => targetRunId !== runId), - ); - return next.size === current.size ? current : next; + return issuesApi.interruptQueuedComments(issueId!, { + queueId: queue.queueId, revision: queue.revision, targetRunId: runId, }); - - return { - previousRunState, - previousLocalQueuedCommentRunIds, - }; }, onSuccess: () => { invalidateIssueDetail(); @@ -5022,25 +4944,9 @@ export function IssueDetail({ tasksTab }: { tasksTab?: TaskSidePanelProps["tasks tone: "success", }); }, - onError: (err, _runId, context) => { - for (const state of context?.previousRunState ?? []) { - queryClient.setQueryData(queryKeys.issues.runs(state.ref), state.runs); - queryClient.setQueryData( - queryKeys.issues.liveRuns(state.ref), - state.liveRuns, - ); - queryClient.setQueryData( - queryKeys.issues.activeRun(state.ref), - state.activeRun, - ); - queryClient.setQueryData( - queryKeys.issues.detail(state.ref), - state.issue, - ); - } - if (context?.previousLocalQueuedCommentRunIds) { - setLocallyQueuedCommentRunIds(context.previousLocalQueuedCommentRunIds); - } + onError: (err) => { + invalidateIssueDetail(); + invalidateIssueRunState(); pushToast({ title: "Interrupt failed", body: