fix: reliably interrupt and resume legacy message queues (#13275)

## 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 <noreply@paperclip.ing>
This commit is contained in:
Dotta 2026-09-11 18:26:59 -05:00 committed by GitHub
parent 30c63af0e6
commit f12b647ae8
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
15 changed files with 392 additions and 155 deletions

View File

@ -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 queues 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

View File

@ -1546,6 +1546,10 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
if (
sessionId &&
!initial.proc.timedOut &&
!initial.proc.signal &&
// A started session can emit stale-rollout warnings for other threads.
// After Ctrl-C those warnings must not restart the cancelled turn.
!initial.parsed.sessionId &&
(initial.proc.exitCode ?? 0) !== 0 &&
isCodexUnknownSessionError(initial.proc.stdout, initial.rawStderr)
) {

View File

@ -698,6 +698,39 @@ describe("codex execute", () => {
}
});
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");

View File

@ -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<typeof import("../adapters/process/execute.js")>("../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<typeof actualProcess.execute>[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",

View File

@ -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<ReturnType<typeof seedQueue>>) {
const queueRunId = randomUUID();
const wake = await db

View File

@ -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<InvokableAgentSnapshot | null> {
@ -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);

View File

@ -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),

View File

@ -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",

View File

@ -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<string>();
// 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<typeof heartbeatRuns.$inferSelect, "id" | "companyId">,
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<void>((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

View File

@ -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();

View File

@ -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<IssueQueuedCommentQueue>(`/issues/${id}/queued-comments/interrupt`, data),
steerQueuedComment: (
id: string,
commentId: string,

View File

@ -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.",
);
});
});

View File

@ -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) {

View File

@ -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(

View File

@ -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<string>
@ -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<RunForIssue[]>(
queryKeys.issues.runs(ref),
),
liveRuns: queryClient.getQueryData<LiveRunForIssue[]>(
queryKeys.issues.liveRuns(ref),
),
activeRun: queryClient.getQueryData<ActiveRunForIssue | null>(
queryKeys.issues.activeRun(ref),
),
issue: queryClient.getQueryData<Issue>(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<RunForIssue[] | undefined>(
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: