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: