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:
parent
30c63af0e6
commit
f12b647ae8
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
) {
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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.",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Reference in New Issue