diff --git a/doc/MCP-ACCESS-GOVERNANCE.md b/doc/MCP-ACCESS-GOVERNANCE.md index 561d50f754..be1fcaebb4 100644 --- a/doc/MCP-ACCESS-GOVERNANCE.md +++ b/doc/MCP-ACCESS-GOVERNANCE.md @@ -125,6 +125,13 @@ Paperclip plays two roles in the MCP graph, and confusing them is the most commo Operators usually mean *gateway* when they say "MCP access governance". For Paperclip-managed local adapter runs, Paperclip writes adapter MCP config that points at named gateway endpoints with short-lived scoped bearer tokens. Policies, approvals, and the audit log only exist for calls that enter gateway mode. +Connected tool names reserve the `mcp__paperclip-assigned__` provider prefix +within the 128-character limit. Short existing names stay compatible. Longer +names use a readable prefix and a stable hash of the connection, application, +and upstream tool identity; duplicate catalog entries also include their entry +identity. Gateway metadata retains the original upstream name for dispatch, +permissions, and audit. An alias never changes which connection executes a call. + V1 does not claim host-wide MCP enforcement. If an unmanaged external client, hand-edited adapter config, or process outside the Paperclip-controlled workspace calls an upstream MCP server directly, Paperclip can warn about known overlapping config entries but cannot prevent or audit that bypass. Treat managed MCP config as a control-plane containment feature for Paperclip-launched agents, not as an endpoint firewall for the operator's whole machine. ## Managed connections diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index df10f22d84..e6bd5fcf92 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -1475,6 +1475,22 @@ Resumed sessions keep the existing message-delta path; fresh sessions receive th full covered history. Stable wording and bounded references avoid adding another full brief on each comment, but provider cache hits must be measured separately. +### Legacy cancellation recovery + +Legacy recovery notices expose the original stopped run, its failure reason, +saved-message count, and the recorded next action. All run-bound holds offer +Inspect run. A positively identified unexpected provider cancellation can request +a fresh continuation only with a complete empty tool inventory. Admission still +proves provider termination and applies pause, budget, approval, dependency, +and ownership gates. Saved input uses the existing ordered, single-delivery queue; +prior completed actions and uncertain outcomes remain history. An operator Stop +or ambiguous historical cancellation does not automatically release that queue. +Externally bound chat conversations use a new chat message to continue. Their +recovery notices show that guidance and do not offer the board Continue action. +Removed chat connections direct the operator to inspect the run and create a new +task. Unavailable connections direct the operator to restore access or create a +new task. A retained conversation record does not prove the chat is available. + ### Native finalization recovery display A native finalization retry uses the existing recovery record but does not imply diff --git a/doc/run-log-events.md b/doc/run-log-events.md index adcc57ed69..dd962b0aa2 100644 --- a/doc/run-log-events.md +++ b/doc/run-log-events.md @@ -239,6 +239,21 @@ section in the Observability contract. ## Execution recovery +Cancelled runs retain `resultJson.cancellation`: a closed `source` label +(`operator`, `queued_message`, `shutdown`, `provider`, `transport`, +`control_plane`, or `unknown`), whether the stop was expected, the initiator, +reason, and recording time. Recorded stop intent survives adapter completion. +The local lifecycle event includes this evidence. Started cancellations without +an expected stop are also reported to Sentry; its cancellation diagnostics contain +only source, expectedness, and initiator type, never initiator IDs or reason text. +Historical ambiguous cancellations stay `unknown` and do not authorize replay. + +Provider tool-definition validation failures use +`provider_tool_definition_invalid` / `configuration`. Automatic retry and +continuation recovery stop until the configuration is repaired. Classification +uses raw provider diagnostics in memory before redaction; stored diagnostics +remain redacted and bounded. + Provider identity diagnostics remain in the local run log. They record the notification method, expected and received thread/turn identifiers, and the classification (root, verified descendant, stale, unrelated informational, or invalid authoritative). They omit the original provider payload and credentials. Repeated informational notices are bounded. Ignored unrelated Codex notifications use `harness.diagnostic` with code diff --git a/packages/adapter-utils/src/acpx-engine/execute.test.ts b/packages/adapter-utils/src/acpx-engine/execute.test.ts index 07d6070d40..e5762316d2 100644 --- a/packages/adapter-utils/src/acpx-engine/execute.test.ts +++ b/packages/adapter-utils/src/acpx-engine/execute.test.ts @@ -3036,6 +3036,22 @@ describe("gemini ACP flag selection", () => { }); describe("ACP activity diagnostics", () => { + it("classifies rejected definitions before redaction and overrides a transient adapter label", async () => { + const root = await makeTempRoot(), cwd = path.join(root, "worktree"); + await fs.mkdir(cwd, { recursive: true }); + const execute = createAcpxEngineExecutor({ + classifyTerminalSessionFailure: () => ({ errorCode: "claude_transient_upstream", errorFamily: "transient_upstream" }), + createRuntime: () => ({ ...buildRuntime(), startTurn: (options: { onTerminalSessionFailure: (failure: unknown) => void }) => { + options.onTerminalSessionFailure({ category: "service", details: "API Error: 400 tools.17.custom.name: String should have at most 128 characters" }); + return { events: (async function* () {})(), result: Promise.resolve({ status: "failed", error: new Error("turn failed") }), cancel: async () => {} }; + } }) as never, + }); + const result = await execute({ runId: "invalid-definition", agent: { id: "agent-1", companyId: "company-1" }, runtime: {}, + config: { agent: "custom", agentCommand: "node ./fake-acp.js", stateDir: path.join(root, "state"), cwd, env: { SECRET: "128" } }, + context: {}, onMeta: async () => {}, onLog: async () => {} } as never); + expect(result).toMatchObject({ errorCode: "provider_tool_definition_invalid", errorFamily: "configuration" }); + expect(JSON.stringify(result.resultJson?.terminalSessionFailure)).not.toContain("128"); + }); it.each(["terminal", "relay_error", "no_events"])( "snapshots %s activity before usage reads, failure logging, and cleanup", async (outcome) => { const root = await makeTempRoot(); diff --git a/packages/adapter-utils/src/acpx-engine/execute.ts b/packages/adapter-utils/src/acpx-engine/execute.ts index b47667cd63..b65db8efb3 100644 --- a/packages/adapter-utils/src/acpx-engine/execute.ts +++ b/packages/adapter-utils/src/acpx-engine/execute.ts @@ -40,6 +40,7 @@ import { captureLocalProcess, capturedProcessExited, killCapturedLocalProcess } import type { DuplexLossReason } from "../duplex-observability.js"; import { DUPLEX_CHANNEL_LOST_ERROR_CODE } from "../bridge-transport-contract.js"; import { + classifyToolDefinitionFailure, formatTerminalSessionFailure, sanitizeTerminalSessionFailure, type AcpxTerminalSessionFailure, @@ -4786,10 +4787,11 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { // Diagnostics are retained even when an adapter has no recovery // classifier. Redact before bounding so partial secrets cannot leak. onTerminalSessionFailure: (failure: AcpxTerminalSessionFailure) => { + terminalFailureClassification = classifyToolDefinitionFailure(failure) + ?? deps.classifyTerminalSessionFailure?.(failure, new Date(now())) ?? null; terminalSessionFailure = sanitizeTerminalSessionFailure( failure, prepared.env, ctx.authToken, parseObject(ctx.config.env), ); - terminalFailureClassification = deps.classifyTerminalSessionFailure?.(failure, new Date(now())) ?? null; }, }); activeTurn = turn; diff --git a/packages/adapter-utils/src/acpx-engine/terminal-session-failure.test.ts b/packages/adapter-utils/src/acpx-engine/terminal-session-failure.test.ts index dbe8cea1ed..e1fdad45b7 100644 --- a/packages/adapter-utils/src/acpx-engine/terminal-session-failure.test.ts +++ b/packages/adapter-utils/src/acpx-engine/terminal-session-failure.test.ts @@ -1,7 +1,27 @@ import { describe, expect, it } from "vitest"; -import { formatTerminalSessionFailure, sanitizeTerminalSessionFailure } from "./terminal-session-failure.js"; +import { classifyToolDefinitionFailure, formatTerminalSessionFailure, sanitizeTerminalSessionFailure } from "./terminal-session-failure.js"; describe("terminal session failure diagnostics", () => { + it.each([ + "API Error: 400 tools.17.custom.name: String should have at most 128 characters", + "tools[0].input_schema: invalid schema", + "Invalid schema for function 'read-file': unsupported type", + ])("rejects unchanged tool definitions: %s", (details) => { + const failure = { category: "service", details }; + expect(classifyToolDefinitionFailure(failure)).toEqual({ + errorCode: "provider_tool_definition_invalid", errorFamily: "configuration", + }); + // Classification precedes arbitrary short-secret redaction. + expect(sanitizeTerminalSessionFailure(failure, { SECRET: "128" }).details).not.toContain("128"); + }); + it.each([ + "HTTP 503 service unavailable", + "tool call arguments are invalid: name must be a string", + "prompt is too long", + "tools.0.arguments.name: String should have at most 128 characters", + ])("keeps unrelated errors on their existing recovery path: %s", (details) => { + expect(classifyToolDefinitionFailure({ category: "service", details })).toBeNull(); + }); it("keeps the provider category, message, request id and stack", () => { const failure = { category: "service", diff --git a/packages/adapter-utils/src/acpx-engine/terminal-session-failure.ts b/packages/adapter-utils/src/acpx-engine/terminal-session-failure.ts index 8491aa56f8..c83e23e31c 100644 --- a/packages/adapter-utils/src/acpx-engine/terminal-session-failure.ts +++ b/packages/adapter-utils/src/acpx-engine/terminal-session-failure.ts @@ -1,4 +1,5 @@ import { redactDiagnosticText, REDACTED_COMMAND_TEXT_VALUE } from "../command-redaction.js"; +import type { AdapterExecutionResult } from "../types.js"; export interface AcpxTerminalSessionFailure { category: string; @@ -11,6 +12,23 @@ export interface AcpxTerminalSessionFailureDiagnostic extends AcpxTerminalSessio } const CATEGORIES = new Set(["connection", "access", "limit", "service", "request", "unknown"]); + +/** Inspect raw provider text before redaction can remove status codes or limits. + * Only definition paths qualify: invalid call arguments and overloads remain + * subject to their existing recovery contracts. + */ +export function classifyToolDefinitionFailure( + failure: AcpxTerminalSessionFailure, +): Pick | null { + const text = `${failure.title ?? ""}\n${failure.details ?? ""}`; + const definitionPath = /\btools(?:\.\d+|\[\d+\])(?:\.(?:custom|function))?\.(?:name|input_schema|parameters|description|type)\b/i; + const validation = /at most|too long|max(?:imum)?[_ ]?length|invalid|not valid|must|should|schema|validation|unsupported/i; + const namedDefinition = /(?:invalid|unsupported)\s+(?:tool|function)\s+(?:definition|schema|name)|invalid schema for (?:function|tool)/i; + if ((definitionPath.test(text) && validation.test(text)) || namedDefinition.test(text)) { + return { errorCode: "provider_tool_definition_invalid", errorFamily: "configuration" }; + } + return null; +} // Leave room under the server's 64 KiB run-log chunk limit even when every // retained character needs JSON escaping. The transcript stores the text once. const FIELD_LIMITS = { title: 4096, details: 24576 } as const; diff --git a/packages/adapter-utils/src/types.ts b/packages/adapter-utils/src/types.ts index 50b3653c24..e8cf9307cd 100644 --- a/packages/adapter-utils/src/types.ts +++ b/packages/adapter-utils/src/types.ts @@ -67,6 +67,7 @@ export interface AdapterRuntimeServiceReport { } export type AdapterExecutionErrorFamily = + | "configuration" | "transient_upstream" | "provider_quota" | "model_refusal" diff --git a/packages/shared/src/types/execution-projection.ts b/packages/shared/src/types/execution-projection.ts index 1736db6589..60090799e9 100644 --- a/packages/shared/src/types/execution-projection.ts +++ b/packages/shared/src/types/execution-projection.ts @@ -5,6 +5,11 @@ export interface ExecutionBlocker { agentId: string | null; cause: string; nextAction: string; + runStatus?: string | null; + runError?: string | null; + /** Candidate only: admission must still verify termination and all gates. */ + canContinue?: boolean; + savedMessageCount?: number; } /** Presentation of existing execution records, not a second task status machine. */ diff --git a/server/src/__tests__/heartbeat-direct-adapter-native-isolation.test.ts b/server/src/__tests__/heartbeat-direct-adapter-native-isolation.test.ts index db5dcb7ea0..1987e463cb 100644 --- a/server/src/__tests__/heartbeat-direct-adapter-native-isolation.test.ts +++ b/server/src/__tests__/heartbeat-direct-adapter-native-isolation.test.ts @@ -196,4 +196,26 @@ describeEmbeddedPostgres("direct adapter native-runner isolation", () => { expect(nativeRows.every((rows) => rows.length === 0)).toBe(true); }, ); + + it("retains an unsolicited provider cancellation through finalization", async () => { + const companyId = randomUUID(), agentId = randomUUID(); + execute.mockImplementation(async (context) => { + expect(context.signal?.aborted).toBe(false); + return { exitCode: 1, signal: null, timedOut: false, errorMessage: "ACP turn cancelled", + resultJson: { status: "cancelled", acpToolInventoryComplete: true, acpPendingToolCount: 0 } }; + }); + await db.insert(companies).values({ id: companyId, name: "Provider stop", + issuePrefix: `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`, + requireBoardApprovalForNewAgents: false, defaultResponsibleUserId: "responsible-user" }); + await db.insert(agents).values({ id: agentId, companyId, name: "Claude", role: "engineer", status: "idle", + adapterType: "claude_local", adapterConfig: {}, runtimeConfig: {}, permissions: {} }); + const queued = await heartbeat.invoke(agentId, "on_demand", {}, "manual"); + const finished = await waitForRunToFinish(heartbeat, queued!.id); + expect(finished).toMatchObject({ status: "cancelled", runtimeMode: "legacy", resultJson: { + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: expect.stringContaining("Provider cancelled execution") }, + acpToolInventoryComplete: true, acpPendingToolCount: 0, + } }); + expect(execute).toHaveBeenCalledOnce(); + }); }); diff --git a/server/src/__tests__/heartbeat-list.test.ts b/server/src/__tests__/heartbeat-list.test.ts index 6b8407e284..c73d2b7b66 100644 --- a/server/src/__tests__/heartbeat-list.test.ts +++ b/server/src/__tests__/heartbeat-list.test.ts @@ -288,6 +288,12 @@ describeEmbeddedPostgres("heartbeat list", () => { privateSyncMetadata: oversizedNestedPayload, }, workspaceRestoreFailure: "restore_unsafe_archive", + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: "Provider cancelled execution ".repeat(50), recordedAt: "2026-10-02T15:00:00.000Z", + privateMetadata: oversizedNestedPayload }, + acpToolInventoryComplete: true, + acpPendingToolCount: 0, + errorFamily: "configuration", finalResponseRecorded: true, executionBeforeRestore: { errorCode: "model_error", exitCode: 2, timedOut: false }, }, @@ -312,6 +318,11 @@ describeEmbeddedPostgres("heartbeat list", () => { storageWarning: "Agent storage is full. Runs can continue.".repeat(50).slice(0, 1024), }, workspaceRestoreFailure: "restore_unsafe_archive", + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: "Provider cancelled execution ".repeat(50).slice(0, 512), recordedAt: "2026-10-02T15:00:00.000Z" }, + acpToolInventoryComplete: true, + acpPendingToolCount: 0, + errorFamily: "configuration", finalResponseRecorded: true, executionBeforeRestore: { errorCode: "model_error", exitCode: 2, timedOut: false }, }); @@ -319,6 +330,7 @@ describeEmbeddedPostgres("heartbeat list", () => { expect((result?.stdout as string).length).toBeLessThan(oversizedStdout.length); expect(result).not.toHaveProperty("nestedHuge"); expect(result?.instructionSave).not.toHaveProperty("privateSyncMetadata"); + expect(result?.cancellation).not.toHaveProperty("privateMetadata"); expect(result?.terminalSessionFailure).not.toHaveProperty("privateMetadata"); const diagnostic = result?.terminalSessionFailure as { details: string }; expect(diagnostic.details).toContain("[truncated for run retrieval; full text in run error/transcript]"); diff --git a/server/src/__tests__/heartbeat-retry-scheduling.test.ts b/server/src/__tests__/heartbeat-retry-scheduling.test.ts index 115902a079..697b6e7fb5 100644 --- a/server/src/__tests__/heartbeat-retry-scheduling.test.ts +++ b/server/src/__tests__/heartbeat-retry-scheduling.test.ts @@ -251,6 +251,16 @@ describeEmbeddedPostgres("heartbeat bounded retry scheduling", () => { } + it("never schedules invalid provider definitions even when a caller supplies a retry policy", async () => { + const runId = randomUUID(), companyId = randomUUID(), agentId = randomUUID(); + const now = new Date("2026-04-20T12:00:00.000Z"); + await seedRetryFixture({ runId, companyId, agentId, now, errorCode: "provider_tool_definition_invalid", + errorFamily: "transient_upstream" }); + expect(await heartbeat.scheduleBoundedRetry(runId, { now, retryReason: "transient_failure", maxAttempts: 9 })) + .toMatchObject({ outcome: "not_scheduled", reason: expect.stringContaining("Repair") }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + }); + it.each(["restore_unsafe_archive", "restore_lock_timeout"])("keeps the existing retry budget for %s", async (classification) => { const runId = randomUUID(), companyId = randomUUID(), agentId = randomUUID(); const now = new Date("2026-04-20T12:00:00.000Z"); diff --git a/server/src/__tests__/tool-gateway.test.ts b/server/src/__tests__/tool-gateway.test.ts index 51ee450db0..ee91153c37 100644 --- a/server/src/__tests__/tool-gateway.test.ts +++ b/server/src/__tests__/tool-gateway.test.ts @@ -2077,6 +2077,31 @@ rl.on("line", (line) => { ])); }); + it("dispatches a bounded alias to the original long upstream tool name", async () => { + const company = await createCompany(db); + const agent = await createAgent(db, company.id); + const { run } = await createIssueAndRun(db, company.id, agent.id); + const upstreamName = "file_download_batch_exports_count_rows_create"; + const fake = await startFakeRemoteMcpServer((request) => { + expect(request.body).toMatchObject({ method: "tools/call", params: { name: upstreamName, arguments: {} } }); + return { body: { jsonrpc: "2.0", id: request.body?.id, result: { content: [{ type: "text", text: "ok" }] } } }; + }); + try { + const { catalogEntry } = await createRemoteMcpTool(db, company.id, { + applicationKey: `app-gallery:posthog:${randomUUID()}`, toolName: upstreamName, url: fake.url, + }); + await allowAllToolsForAgent(db, company.id, agent.id); + const gateway = createTestToolGatewayService(db); + const session = await gateway.createSession({ companyId: company.id, agentId: agent.id, runId: run.id }); + const tool = (await gateway.listToolsForSession(session.token)).find((entry) => entry.catalogEntryId === catalogEntry.id)!; + expect(`mcp__paperclip-assigned__${tool.name}`.length).toBeLessThanOrEqual(128); + expect((await gateway.listToolsForSession(session.token)).find((entry) => entry.catalogEntryId === catalogEntry.id)?.name).toBe(tool.name); + await gateway.executeTool({ sessionToken: session.token, tool: tool.name, parameters: {} }); + } finally { + await fake.close(); + } + }); + it.each([ ["local_trusted", { deploymentMode: "local_trusted" as const, deploymentExposure: "private" as const }], ["authenticated/private", { deploymentMode: "authenticated" as const, deploymentExposure: "private" as const }], diff --git a/server/src/routes/agents.ts b/server/src/routes/agents.ts index 686ada2e48..c3f1b8a422 100644 --- a/server/src/routes/agents.ts +++ b/server/src/routes/agents.ts @@ -13,6 +13,7 @@ import { assertAiConnectionCreateAccess, canInstallSharedAiConnectionForNewAgent import { isAiConnectionCompatible } from "@paperclipai/shared"; import { applyConnectorSkills, resolveConnectorAssignments, annotateConnectorSkills, isConnectorSkill } from "../services/connector-runtime.js"; import { getExecutionBlocker } from "../services/execution-blocker.js"; +import { canContinueCancelledRun } from "../services/run-cancellation.js"; import { paperclipRunnerTransitionConfig, normalizeLegacyRunnerProvider, isPaperclipRunnerProvider } from "@paperclipai/adapter-utils"; import { executionProjectionForRun, executionProjectionsForRuns } from "../services/execution-projection.js"; import { selectDashboardRunIds } from "../services/dashboard-run-selection.js"; @@ -6005,8 +6006,8 @@ export function agentRoutes( ) { throw notFound("Failed run not found"); } - if (!["failed", "timed_out"].includes(failedRun.status)) { - throw conflict("Only a failed run can be retried."); + if (!["failed", "timed_out"].includes(failedRun.status) && !canContinueCancelledRun(failedRun)) { + throw conflict("Only a failed run or verified unexpected cancellation can start a new attempt."); } if (failedRun.runtimeMode === "native" && failedRun.errorCode === "native_session_cleanup_quarantined") { throw conflict("The stopped native session requires cleanup and reconciliation before a new attempt.", { @@ -6018,6 +6019,10 @@ export function agentRoutes( typeof failedContext.issueId === "string" ? failedContext.issueId : null; + if (failedRun.status === "cancelled" && (!issueId || + (await getExecutionBlocker(db, agent.companyId, issueId))?.runId !== failedRun.id)) { + throw conflict("The stopped run no longer owns this task's recovery hold."); + } if (issueId) { const issue = await issueService(db).getById(issueId); if (!issue || issue.companyId !== agent.companyId) throw notFound("Task not found"); @@ -6073,6 +6078,9 @@ export function agentRoutes( .then((rows) => rows[0]) : null; if (chatBinding) { + if (failedRun.status === "cancelled") { + throw conflict("Send a new chat message to continue this stopped conversation."); + } if (!options.chatRunRetries || !req.actor.userId) { throw conflict("Chat retry authorization is unavailable.", { code: "chat_failed_run_retry_requires_authorized_context", diff --git a/server/src/sentry.ts b/server/src/sentry.ts index 0a4ad5d678..665cab811f 100644 --- a/server/src/sentry.ts +++ b/server/src/sentry.ts @@ -126,7 +126,7 @@ export function captureException(error: unknown): void { } /** The run status values that mark a run as a genuine terminal failure. */ -export type RunFailureStatus = "failed" | "timed_out"; +export type RunFailureStatus = "failed" | "timed_out" | "cancelled"; /** * The diagnostic values `captureRunFailure` sends with a terminal-failure diff --git a/server/src/services/__tests__/run-failure-report.test.ts b/server/src/services/__tests__/run-failure-report.test.ts index 2b7e184e48..eddef8827b 100644 --- a/server/src/services/__tests__/run-failure-report.test.ts +++ b/server/src/services/__tests__/run-failure-report.test.ts @@ -261,6 +261,25 @@ describeEmbeddedPostgres("reportRunFailure", () => { expect(mockCaptureRunFailure).not.toHaveBeenCalled(); }); + it.each(["provider", "unknown"])("reports an unexpected started cancellation from %s", async source => { + await seedCompanyAndAgent(); + const run = buildRun({ status: "cancelled", startedAt: new Date(0), finishedAt: new Date(1000), + resultJson: { cancellation: { source, expected: false, initiator: { type: "provider", id: "private-actor" }, + reason: "private-reason", recordedAt: new Date(1000).toISOString() } } }); + await reportRunFailure(db, run); + expect(mockCaptureRunFailure).toHaveBeenCalledOnce(); + expect(mockCaptureRunFailure.mock.calls[0][0]).toMatchObject({ runStatus: "cancelled", + diagnostics: { execution: { cancellationSource: source, cancellationExpected: false } } }); + expect(JSON.stringify(mockCaptureRunFailure.mock.calls)).not.toMatch(/private-actor|private-reason/); + }); + it("does not report an operator's Stop as a failure", async () => { + await seedCompanyAndAgent(); + await reportRunFailure(db, buildRun({ status: "cancelled", startedAt: new Date(0), resultJson: { + cancellation: { source: "operator", expected: true, initiator: { type: "user", id: "board" }, + reason: "Stop", recordedAt: new Date().toISOString() }, + } })); + expect(mockCaptureRunFailure).not.toHaveBeenCalled(); + }); it("sends agents.adapterType from the loaded agent row", async () => { await seedCompanyAndAgent({ adapterType: "claude_managed" }); diff --git a/server/src/services/execution-blocker.ts b/server/src/services/execution-blocker.ts index f757466799..325b9b9d4f 100644 --- a/server/src/services/execution-blocker.ts +++ b/server/src/services/execution-blocker.ts @@ -1,7 +1,9 @@ import { and, desc, eq, gt, inArray, not, or, sql } from "drizzle-orm"; import { conversationRecoveryActionPredicate, getConversationOwnershipBlocker } from "./conversation-continuation.js"; import { z } from "zod"; -import { heartbeatRuns, issueComments, issues, issueRecoveryActions, type Db } from "@paperclipai/db"; +import { agentWakeupRequests, chatConversations, chatEndpoints, heartbeatRuns, issueComments, issues, issueRecoveryActions, toolConnections, type Db } from "@paperclipai/db"; +import { canContinueCancelledRun, readRunCancellation } from "./run-cancellation.js"; +import { queuedCommentIdsFromWakePayload } from "./issue-queued-comment-queue.js"; import { EXECUTION_RECONCILIATION_CAUSES, type ExecutionBlocker } from "@paperclipai/shared"; /** Resolved recovery bookkeeping can still carry an effective no-replay hold. */ @@ -16,7 +18,8 @@ export function executionBlockerPredicate() { export async function getExecutionBlocker(db: Db, companyId: string, issueId: string, options?: { conversationResetCommentId?: string | null }): Promise { const [conversation] = await db.select({ agentId: issues.conversationAgentId, - boundaryId: issues.conversationBoundaryCommentId }).from(issues).where(and( + boundaryId: issues.conversationBoundaryCommentId, status: issues.status, + assigneeAgentId: issues.assigneeAgentId }).from(issues).where(and( eq(issues.companyId, companyId), eq(issues.id, issueId), )).limit(1); // Resetting model context cannot make an unsafe workspace safe. This hold @@ -53,9 +56,56 @@ export async function getExecutionBlocker(db: Db, companyId: string, issueId: st if (!action) return null; const parsedRunId = z.string().guid().safeParse(action.evidence.runId ?? action.evidence.sourceRunId); const runId = parsedRunId.success ? parsedRunId.data : null; - const [run] = runId ? await db.select({ agentId: heartbeatRuns.agentId }).from(heartbeatRuns).where(and( + const [run] = runId ? await db.select().from(heartbeatRuns).where(and( eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.id, runId), )).limit(1) : []; + const queue = await db.select({ payload: agentWakeupRequests.payload }).from(agentWakeupRequests).where(and( + eq(agentWakeupRequests.companyId, companyId), eq(agentWakeupRequests.status, "deferred_issue_execution"), + sql`${agentWakeupRequests.payload}->>'issueId' = ${issueId}`, + )); + const savedIds = [...new Set(queue.flatMap(entry => { + const context = entry.payload?._paperclipWakeContext as Record | undefined; + return [...queuedCommentIdsFromWakePayload(entry.payload), entry.payload?.commentId, context?.wakeCommentId] + .filter((id): id is string => z.string().guid().safeParse(id).success); + }))]; + const saved = savedIds.length ? await db.select({ id: issueComments.id }).from(issueComments).where(and( + eq(issueComments.companyId, companyId), eq(issueComments.issueId, issueId), + inArray(issueComments.id, savedIds), eq(issueComments.authorType, "user"), sql`${issueComments.deletedAt} is null`, + sql`not exists (select 1 from ${heartbeatRuns} consumed where consumed.company_id = ${companyId} + and consumed.context_snapshot->>'issueId' = ${issueId} + and not (consumed.status = 'cancelled' and consumed.started_at is null) + and (consumed.context_snapshot->>'wakeCommentId' = ${issueComments.id}::text + or consumed.context_snapshot->'wakeCommentIds' @> jsonb_build_array(${issueComments.id}::text)))`, + )) : []; + const cancellation = run ? readRunCancellation(run.resultJson) : null; + const runError = cancellation?.reason ?? (run?.status === "cancelled" + ? "Execution was cancelled; its source was not recorded." + : run?.error); + const eligibleContinuation = Boolean(run && conversation?.assigneeAgentId === run.agentId && + !["done", "cancelled"].includes(conversation.status) && canContinueCancelledRun(run)); + const [chatBinding] = eligibleContinuation ? await db.select({ + state: chatConversations.state, endpointStatus: chatEndpoints.status, + connectionStatus: toolConnections.status, connectionEnabled: toolConnections.enabled, + }).from(chatConversations).leftJoin(chatEndpoints, and( + eq(chatEndpoints.companyId, companyId), eq(chatEndpoints.id, chatConversations.endpointId), + )).leftJoin(toolConnections, and( + eq(toolConnections.companyId, companyId), eq(toolConnections.id, chatEndpoints.connectionId), + )).where(and( + eq(chatConversations.companyId, companyId), eq(chatConversations.issueId, issueId), + )).limit(1) : []; + let nextAction = action.nextAction; + if (chatBinding) { + if (chatBinding.state === "endpoint_removed" || !chatBinding.endpointStatus || chatBinding.endpointStatus === "archived") { + nextAction = "This chat connection was removed. Inspect the stopped run and create a new task to continue the work."; + } else if (chatBinding.state === "unavailable" || !["active", "verifying"].includes(chatBinding.endpointStatus) || + chatBinding.connectionStatus !== "active" || !chatBinding.connectionEnabled) { + nextAction = "This chat connection is unavailable. Restore access in Apps or create a new task to continue the work."; + } else { + nextAction = "Send a new chat message to continue this conversation."; + } + } else if (run?.status === "cancelled" && !eligibleContinuation && action.cause === "legacy_execution_requires_reconciliation") { + nextAction += " Inspect the run before sending a new message to request continuation."; + } return { recoveryActionId: action.id, @@ -63,6 +113,10 @@ export async function getExecutionBlocker(db: Db, companyId: string, issueId: st // A stopped reviewer can differ from the task owner who receives the work back. agentId: run?.agentId ?? null, cause: action.cause, - nextAction: action.nextAction, + nextAction, + runStatus: run?.status ?? null, + runError: runError?.slice(0, 1024) ?? null, + canContinue: eligibleContinuation && !chatBinding, + savedMessageCount: saved.length, }; } diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index d7aae395d6..830c3ff264 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -12,7 +12,7 @@ import { join } from "node:path"; import { and, eq } from "drizzle-orm"; import { beforeAll, afterAll, describe, it, expect, vi } from "vitest"; import { - approvals, issueApprovals, issueThreadInteractions, + approvals, issueApprovals, issueThreadInteractions, chatConversations, chatEndpoints, toolApplications, toolConnections, agentWakeupRequests, agents, companies, createDb, heartbeatRunEvents, heartbeatRuns, issueComments, issueRecoveryActions, issues, nativeRunFinalizations, nativeRunResults, completionContracts, environmentLeases, environments, issueRelations, issueTreeHolds, issueTreeHoldMembers, } from "@paperclipai/db"; @@ -326,12 +326,18 @@ const support = await getEmbeddedPostgresTestSupport(); else await expect(result).rejects.toThrow("continuation_user_authorization_missing"); }); - it.each(["valid", "wrong_actor", "consumed", "discarded", "operator_stop", "already_delivered", "earlier_delivered", "unstarted_cancelled", "unstarted_cancelled_metadata", "foreign_queue"])("validates automatic saved-message delivery: %s", async kind => { + it.each(["valid", "provider_cancel", "pending_tool", "unknown_cancel", "wrong_actor", "consumed", "discarded", "operator_stop", "already_delivered", "earlier_delivered", "unstarted_cancelled", "unstarted_cancelled_metadata", "foreign_queue"])("validates automatic saved-message delivery: %s", async kind => { const f = await seed(); await db.update(agents).set({ adapterType: "claude_local" }).where(eq(agents.id, f.agentId)); await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); await db.update(heartbeatRuns).set({ runtimeMode: "legacy", nativeIssueId: null, - status: kind === "operator_stop" ? "cancelled" : "failed", + status: ["operator_stop", "provider_cancel", "pending_tool", "unknown_cancel"].includes(kind) ? "cancelled" : "failed", + startedAt: new Date("2026-09-11T08:00:00Z"), + ...(["provider_cancel", "pending_tool"].includes(kind) ? { resultJson: { + acpToolInventoryComplete: true, acpPendingToolCount: kind === "pending_tool" ? 1 : 0, + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: "Provider cancelled execution", recordedAt: new Date().toISOString() }, + } } : {}), contextSnapshot: { issueId: f.issueId, ...(kind === "already_delivered" ? { wakeCommentIds: [f.commentId] } : {}) }, }).where(eq(heartbeatRuns.id, f.sourceRunId)); await db.update(issueRecoveryActions).set({ cause: "legacy_execution_requires_reconciliation" }) @@ -367,7 +373,7 @@ const support = await getEmbeddedPostgresTestSupport(); return admitExplicitNativeContinuation({ ...f, actorId: kind === "wrong_actor" ? "someone-else" : f.actorId, db: tx as unknown as typeof db, queuedCommentRequestId: queueId, dryRun: true }); }); - if (kind === "valid" || kind.startsWith("unstarted_cancelled")) expect(result).toMatchObject({ previousRunId: f.sourceRunId, commentId: f.commentId }); + if (kind === "valid" || kind === "provider_cancel" || kind.startsWith("unstarted_cancelled")) expect(result).toMatchObject({ previousRunId: f.sourceRunId, commentId: f.commentId }); else expect(result).toBeNull(); }); @@ -379,6 +385,72 @@ const support = await getEmbeddedPostgresTestSupport(); return result; }); + it.each(["active", "removed", "unavailable", "paused", "disabled"])("shows usable recovery guidance for external chat cancellations: %s", async kind => { + const f = await seed(); + await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); + await db.update(heartbeatRuns).set({ runtimeMode: "legacy", status: "cancelled", nativeIssueId: null, + startedAt: new Date("2026-09-11T08:00:00Z"), resultJson: { + acpToolInventoryComplete: true, acpPendingToolCount: 0, + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: "Provider cancelled execution", recordedAt: new Date().toISOString() }, + }, + }).where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.update(issueRecoveryActions).set({ cause: "legacy_execution_requires_reconciliation" }) + .where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ canContinue: true }); + const applicationId = randomUUID(), connectionId = randomUUID(), endpointId = randomUUID(); + await db.insert(toolApplications).values({ id: applicationId, companyId: f.companyId, name: "Slack", type: "chat" }); + await db.insert(toolConnections).values({ id: connectionId, companyId: f.companyId, applicationId, + name: "Slack", uid: connectionId, transport: "chat_sdk", connectionPurpose: "channel", status: "active", enabled: kind !== "disabled" }); + await db.insert(chatEndpoints).values({ id: endpointId, companyId: f.companyId, connectionId, + provider: "slack", publicId: endpointId, assignedAgentId: f.agentId, + status: kind === "removed" ? "archived" : kind === "paused" ? "paused" : "active" }); + await db.insert(chatConversations).values({ companyId: f.companyId, endpointId, issueId: f.issueId, + externalConversationId: "channel", externalLabel: "#recovery", + state: kind === "removed" ? "endpoint_removed" : kind === "unavailable" ? "unavailable" : "active" }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ + runId: f.sourceRunId, canContinue: false, + nextAction: kind === "active" ? "Send a new chat message to continue this conversation." + : kind === "removed" ? "This chat connection was removed. Inspect the stopped run and create a new task to continue the work." + : "This chat connection is unavailable. Restore access in Apps or create a new task to continue the work.", + }); + }); + + it("delivers saved messages once after provider cancellation, preserving source history", async () => { + const f = await seed(); + await db.update(agents).set({ adapterType: "claude_local" }).where(eq(agents.id, f.agentId)); + // Hold agent capacity so the test observes durable admission without launching a provider. + await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); + await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); + const resultJson = { status: "cancelled", acpToolInventoryComplete: true, acpPendingToolCount: 0, + cancellation: { source: "provider", expected: false, initiator: { type: "provider" }, + reason: "Provider cancelled execution", recordedAt: new Date().toISOString() }, + completedAction: "preserve-this-record" }; + await db.update(heartbeatRuns).set({ runtimeMode: "legacy", status: "cancelled", nativeIssueId: null, + startedAt: new Date("2026-09-11T08:00:00Z"), resultJson, + }).where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.update(issueRecoveryActions).set({ cause: "legacy_execution_requires_reconciliation" }) + .where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + const queueId = randomUUID(); + await db.insert(agentWakeupRequests).values({ id: queueId, companyId: f.companyId, agentId: f.agentId, + source: "automation", reason: "issue_commented", status: "deferred_issue_execution", + requestedByActorType: "user", requestedByActorId: f.actorId, updatedAt: new Date(0), + payload: { issueId: f.issueId, commentId: f.commentId, executionWait: { reason: "execution_recovery" }, + _paperclipWakeContext: { issueId: f.issueId, wakeReason: "issue_commented", wakeCommentIds: [f.commentId], wakeCommentId: f.commentId } }, + }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ savedMessageCount: 1, canContinue: true }); + await Promise.all([heartbeatService(db).resumeExecutionWaitComments(), heartbeatService(db).resumeExecutionWaitComments()]); + await heartbeatService(db).resumeExecutionWaitComments(); + const runs = await db.select().from(heartbeatRuns).where(and(eq(heartbeatRuns.companyId, f.companyId), eq(heartbeatRuns.status, "queued"))); + expect(runs).toHaveLength(1); + expect(runs[0].contextSnapshot).toMatchObject({ forceFreshSession: true, previousRunId: f.sourceRunId, wakeCommentIds: [f.commentId] }); + const [queue] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, queueId)); + expect(queue).toMatchObject({ status: "coalesced", runId: runs[0].id }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toBeNull(); + const [source] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.sourceRunId)); + expect(source.resultJson).toEqual(resultJson); + }); + it.each(["verified", "unproven", "changed", "dry_run", "retry", "duplicate"])( "requires exact local cleanup for a new turn after worker loss (%s)", async mode => { const f = await seed(); diff --git a/server/src/services/explicit-native-continuation.ts b/server/src/services/explicit-native-continuation.ts index 2eafaf4d3e..eb47212a0d 100644 --- a/server/src/services/explicit-native-continuation.ts +++ b/server/src/services/explicit-native-continuation.ts @@ -1,4 +1,5 @@ import { appendHeartbeatRunEvent } from "./heartbeat-run-events.js"; +import { canContinueCancelledRun } from "./run-cancellation.js"; import { readQueuedInteractionResponse } from "./queued-interaction-response.js"; import { isCancelledNativeStartup } from "./cancelled-native-startup.js"; import { hasNativeLocalProcessStop, hasHistoricalSuspendedNativeSession } from "./native-local-process-stop.js"; @@ -171,7 +172,8 @@ export async function admitExplicitNativeContinuation(input: { !(queuedInterrupt && response?.source.requiresFreshSession && run.runtimeMode === "native")) return null; // Saved input is a request for a new turn, never permission to undo an // operator Stop or redeliver a message already consumed by this run. - if (queuedRequest && !queuedInterrupt && ((run.status === "cancelled" && !unusedAdmission) || + if (retry && run.status === "cancelled" && !canContinueCancelledRun(run)) return blocked("cancelled_by_operator", "Inspect the stopped run and send a new message to continue."); + if (queuedRequest && !queuedInterrupt && ((run.status === "cancelled" && !unusedAdmission && !canContinueCancelledRun(run)) || run.contextSnapshot?.wakeCommentId === commentId || (Array.isArray(run.contextSnapshot?.wakeCommentIds) && run.contextSnapshot.wakeCommentIds.includes(commentId)))) return null; if (legacyUserTurn) { diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 6deec22e7a..ca08e34d3c 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -5,6 +5,8 @@ import { AgentDirectoryReuseInvalidatedError, isAgentDirectoryCopy } from "./age import type { PaperclipTurnContext } from "@paperclipai/adapter-utils/server-utils"; import { restoreNativeWorkspaceBestEffort } from "./native-runtime/native-workspace-best-effort.js"; +import { ASSIGNED_MCP_SERVER_NAME } from "./mcp-tool-names.js"; +import { cancellationResultJson, canContinueCancelledRun, readRunCancellation, requestedRunCancellation } from "./run-cancellation.js"; import { withNativeWorkspaceFinalizationOwnership, NativeWorkspaceFinalizationBusyError, @@ -3495,6 +3497,27 @@ const heartbeatRunSafeResultJsonColumn = sql | null>` 'workspaceRestoreFailure', case when ${heartbeatRuns.resultJson} ->> 'workspaceRestoreFailure' in ('restore_permission_denied', 'restore_lock_timeout', 'restore_unsafe_archive', 'restore_failed') then ${heartbeatRuns.resultJson} -> 'workspaceRestoreFailure' end, + 'cancellation', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'cancellation') = 'object' + then jsonb_strip_nulls(jsonb_build_object( + 'source', case when ${heartbeatRuns.resultJson} #>> '{cancellation,source}' + in ('operator', 'queued_message', 'shutdown', 'provider', 'transport', 'control_plane', 'unknown') + then ${heartbeatRuns.resultJson} #> '{cancellation,source}' end, + 'expected', case when jsonb_typeof(${heartbeatRuns.resultJson} #> '{cancellation,expected}') = 'boolean' + then ${heartbeatRuns.resultJson} #> '{cancellation,expected}' end, + 'initiator', jsonb_strip_nulls(jsonb_build_object( + 'type', case when ${heartbeatRuns.resultJson} #>> '{cancellation,initiator,type}' in ('user', 'agent', 'system', 'provider') + then ${heartbeatRuns.resultJson} #> '{cancellation,initiator,type}' end, + 'id', left(${heartbeatRuns.resultJson} #>> '{cancellation,initiator,id}', 128) + )), + 'reason', left(${heartbeatRuns.resultJson} #>> '{cancellation,reason}', 512), + 'recordedAt', left(${heartbeatRuns.resultJson} #>> '{cancellation,recordedAt}', 64) + )) end, + 'acpToolInventoryComplete', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'acpToolInventoryComplete') = 'boolean' + then ${heartbeatRuns.resultJson} -> 'acpToolInventoryComplete' end, + 'acpPendingToolCount', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'acpPendingToolCount') = 'number' + and length(${heartbeatRuns.resultJson} ->> 'acpPendingToolCount') < 16 + then ${heartbeatRuns.resultJson} -> 'acpPendingToolCount' end, + 'errorFamily', left(${heartbeatRuns.resultJson} ->> 'errorFamily', 32), 'finalResponseRecorded', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'finalResponseRecorded') = 'boolean' then ${heartbeatRuns.resultJson} -> 'finalResponseRecorded' end, 'executionBeforeRestore', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'executionBeforeRestore') = 'object' @@ -4872,7 +4895,7 @@ export async function buildPaperclipRuntimeMcpServers(input: { return [ { - name: "paperclip-assigned", + name: ASSIGNED_MCP_SERVER_NAME, url: `${paperclipApiBaseUrl()}/mcp/gateways/${gateway!.gatewayPublicId}`, token: token.token, connectionId: `assignment:${assignmentDigest}`, @@ -10604,6 +10627,12 @@ export function heartbeatService( if (!sourceId || !isUuidLike(sourceId)) continue; const run = await getRun(sourceId); if (!run || run.companyId !== wake.companyId || run.agentId !== wake.agentId) continue; + if (canContinueCancelledRun(run)) { + await resumeSavedLegacyComments(wake.companyId, wake.id).catch(err => { + logger.warn({ err, runId: run.id }, "failed to resume saved input after provider cancellation"); + }); + continue; + } await resumeRemoteStopComments(run, wake.id).catch(err => { logger.warn({ err, runId: run.id }, "failed to resume saved execution-wait message"); }); @@ -12857,6 +12886,9 @@ export function heartbeatService( // Cancelling a queued run that never acquired provider execution is // positive bootstrap evidence. It must not hold unrelated queued messages. + if (previousStatus && (status === "cancelled" || status === "interrupted")) { + patch = { ...patch, resultJson: cancellationResultJson(previousStatus, status, patch?.resultJson, patch?.errorCode, patch?.error) }; + } if ( status === "cancelled" && previousStatus?.status === "queued" && @@ -12941,6 +12973,10 @@ export function heartbeatService( .where(eq(heartbeatRuns.id, runId)) .then((rows) => rows[0] ?? null); + if (previousStatus && (status === "cancelled" || status === "interrupted")) { + patch = { ...patch, resultJson: cancellationResultJson(previousStatus, status, patch?.resultJson, patch?.errorCode, patch?.error) }; + } + // Cancelling a queued run that never acquired provider execution is // positive bootstrap evidence. It must not hold unrelated queued messages. if ( @@ -15357,6 +15393,11 @@ export function heartbeatService( delayMs?: number; }, ) { + if (run.errorCode === "provider_tool_definition_invalid") { + return { outcome: "not_scheduled" as const, + reason: "Repair the invalid tool definitions before starting a new attempt.", + issueId: readNonEmptyString(run.contextSnapshot?.issueId) }; + } if (Array.isArray(run.contextSnapshot?.chatCompletionDeliveryIds) && run.contextSnapshot.chatCompletionDeliveryIds.some(id => typeof id === "string")) { return { outcome: "not_scheduled" as const, reason: "The completion outbox owns this reply's retry budget and publication identity.", @@ -25255,6 +25296,8 @@ export function heartbeatService( : "failed"; } else if (adapterResult.timedOut) { outcome = "timed_out"; + } else if (adapterResult.resultJson?.status === "cancelled") { + outcome = "cancelled"; } else if ( (adapterResult.exitCode ?? 0) === 0 && !adapterResult.errorMessage && @@ -25287,7 +25330,7 @@ export function heartbeatService( const normalizedUsage = sessionUsageResolution.normalizedUsage; const runErrorMessage = outcome === "cancelled" - ? (latestRun?.error ?? adapterResult.errorMessage ?? "Cancelled") + ? redactCurrentUserText(latestRun?.error ?? adapterResult.errorMessage ?? "Cancelled", currentUserRedactionOptions) : outcome === "succeeded" ? null : redactCurrentUserText( @@ -25402,7 +25445,7 @@ export function heartbeatService( } as Record) : null; - const persistedResultJson = mergeHeartbeatRunResultJson( + const persistedResultJson = cancellationResultJson(latestRun ?? run, outcome, mergeHeartbeatRunResultJson( mergeRunStopMetadataForAgent(agent, outcome, { resultJson: mergeAdapterRecoveryMetadata({ resultJson: { @@ -25422,7 +25465,7 @@ export function heartbeatService( errorMessage: runErrorMessage, }), adapterResult.summary ?? null, - ); + ), runErrorCode, runErrorMessage); const finalRunPatch: Partial = { finishedAt: new Date(), @@ -25467,6 +25510,7 @@ export function heartbeatService( .update(heartbeatRuns) .set({ ...finalRunPatch, + resultJson: cancellationResultJson(persistedRunWrite.run, status, finalRunPatch.resultJson, runErrorCode, runErrorMessage), finishedAt: persistedRunWrite.run.finishedAt ?? finalRunPatch.finishedAt, updatedAt: new Date(), @@ -25519,6 +25563,7 @@ export function heartbeatService( payload: { status, exitCode: adapterResult.exitCode, + ...(readRunCancellation(finalizedRun.resultJson) ? { cancellation: readRunCancellation(finalizedRun.resultJson) } : {}), }, }); try { @@ -26811,7 +26856,7 @@ export function heartbeatService( if (opts.requestedByActorType !== "user" || !opts.requestedByActorId || reason !== "retry_failed_run" || source !== "on_demand" || triggerDetail !== "manual" || !failed || failed.companyId !== agent.companyId || failed.agentId !== agentId || - !["failed", "timed_out"].includes(failed.status) || + (!["failed", "timed_out"].includes(failed.status) && !canContinueCancelledRun(failed)) || (failed.nativeIssueId ?? readNonEmptyString(failed.contextSnapshot?.issueId)) !== issueId) { throw conflict("The selected failed run cannot be retried for this task."); } @@ -29386,6 +29431,7 @@ export function heartbeatService( return run; const agent = await getAgent(run.agentId); const errorCode = options.errorCode ?? "cancelled"; + const cancellation = requestedRunCancellation(options.resultJson ?? {}, reason); const pendingProcessCancellation = processRunCancellationSettlements.get( run.id, @@ -29415,6 +29461,7 @@ export function heartbeatService( resultJson: sql`coalesce(${heartbeatRuns.resultJson}, '{}'::jsonb) || ${JSON.stringify(options.errorCode === "issue_reassigned" && options.resultJson?.reassignmentStopConfirmed === true ? { reassignmentStopRequested: true } : {})}::jsonb || + ${JSON.stringify({ cancellation })}::jsonb || jsonb_build_object('startupCancellation', jsonb_build_object( 'requestedAt', ${new Date().toISOString()}::text, 'beforeNativeSelection', ${heartbeatRuns.runtimeMode} = 'legacy' @@ -29428,7 +29475,7 @@ export function heartbeatService( if (!fenced) return getRun(runId); run = fenced; } - const resultJson = agent + const resultJson = { ...(agent ? { ...mergeRunStopMetadataForAgent(agent, "cancelled", { resultJson: parseObject(run.resultJson), @@ -29437,7 +29484,7 @@ export function heartbeatService( }), ...(options.resultJson ?? {}), } - : options.resultJson; + : options.resultJson), cancellation }; try { let releaseProcessCancellation: (() => void) | undefined; @@ -29639,7 +29686,7 @@ export function heartbeatService( stream: "system", level: "warn", message: options.eventMessage ?? "run cancelled", - ...(options.eventPayload ? { payload: options.eventPayload } : {}), + payload: { ...options.eventPayload, cancellation: readRunCancellation(cancelled.resultJson) }, }); await releaseIssueExecutionAndPromote(cancelled, { suppressImmediateRecovery: options.suppressImmediateRecovery, @@ -29682,6 +29729,9 @@ export function heartbeatService( continue; } if (run.runtimeMode === "native") { + await db.update(heartbeatRuns).set({ resultJson: + sql`coalesce(${heartbeatRuns.resultJson}, '{}'::jsonb) || ${JSON.stringify({ cancellation: requestedRunCancellation({}, reason) })}::jsonb`, + }).where(and(eq(heartbeatRuns.id, run.id), eq(heartbeatRuns.status, run.status))); await cancelHeartbeatNativeRun({ db, runId: run.id, @@ -29699,15 +29749,13 @@ export function heartbeatService( finishedAt: new Date(), error: reason, errorCode, - ...(agent - ? { - resultJson: mergeRunStopMetadataForAgent(agent, "cancelled", { - resultJson: persistedCancellationResult, - errorCode, - errorMessage: reason, - }), - } - : {}), + resultJson: { + ...persistedCancellationResult, + ...(agent ? mergeRunStopMetadataForAgent(agent, "cancelled", { + resultJson: persistedCancellationResult, errorCode, errorMessage: reason, + }) : {}), + cancellation: readRunCancellation(persistedCancellationResult) ?? requestedRunCancellation({}, reason), + }, }); await setWakeupStatus(run.wakeupRequestId, "cancelled", { diff --git a/server/src/services/mcp-tool-names.test.ts b/server/src/services/mcp-tool-names.test.ts new file mode 100644 index 0000000000..293d2d6d4b --- /dev/null +++ b/server/src/services/mcp-tool-names.test.ts @@ -0,0 +1,17 @@ +import { describe, expect, it } from "vitest"; +import { ASSIGNED_MCP_SERVER_NAME, boundedMcpToolName } from "./mcp-tool-names.js"; + +describe("provider MCP name budget", () => { + it("keeps short names compatible", () => { + expect(boundedMcpToolName("mcp.drive-12345678:read-file", ["connection", "read_file"])) + .toBe("mcp.drive-12345678:read-file"); + }); + it("bounds the complete provider name and retains unique, stable aliases", () => { + const name = `mcp.app-gallery-posthog-${"a".repeat(40)}-cb60822c:file-download-batch-exports-count-rows-create`; + const first = boundedMcpToolName(name, ["connection-one", "file_download_batch_exports_count_rows_create"]); + expect(`mcp__${ASSIGNED_MCP_SERVER_NAME}__${first}`.length).toBeLessThanOrEqual(128); + expect(first).toContain("file-download"); + expect(boundedMcpToolName(name, ["connection-one", "file_download_batch_exports_count_rows_create"])).toBe(first); + expect(boundedMcpToolName(name, ["connection-two", "file_download_batch_exports_count_rows_create"])).not.toBe(first); + }); +}); diff --git a/server/src/services/mcp-tool-names.ts b/server/src/services/mcp-tool-names.ts new file mode 100644 index 0000000000..e31c26f3b7 --- /dev/null +++ b/server/src/services/mcp-tool-names.ts @@ -0,0 +1,17 @@ +import { createHash } from "node:crypto"; + +// Claude includes this namespace in its 128-character tool-name limit. +// Gateway punctuation is normalized one-for-one by the MCP client. +export const ASSIGNED_MCP_SERVER_NAME = "paperclip-assigned"; +export const MAX_GATEWAY_TOOL_NAME_LENGTH = 128 - `mcp__${ASSIGNED_MCP_SERVER_NAME}__`.length; + +/** Preserve existing short names; aliases still resolve through gateway metadata. */ +export function boundedMcpToolName(name: string, identity: readonly string[]): string { + if (name.length <= MAX_GATEWAY_TOOL_NAME_LENGTH) return name; + const hash = createHash("sha256").update(JSON.stringify(identity)).digest("hex").slice(0, 16); + const separator = name.indexOf(":"); + const namespace = name.slice(0, separator < 0 ? name.length : separator).slice(0, 32); + const tool = separator < 0 ? "tool" : name.slice(separator + 1); + const budget = MAX_GATEWAY_TOOL_NAME_LENGTH - namespace.length - hash.length - 2; + return `${namespace}:${tool.slice(0, budget)}-${hash}`; +} diff --git a/server/src/services/recovery/provider-failure-classification.test.ts b/server/src/services/recovery/provider-failure-classification.test.ts index 17226f5e94..d58bcf42a9 100644 --- a/server/src/services/recovery/provider-failure-classification.test.ts +++ b/server/src/services/recovery/provider-failure-classification.test.ts @@ -172,4 +172,10 @@ describe("classifyAdapterFailureForRecovery", () => { resultJson: null, })).toBeNull(); }); + it("holds invalid provider definitions for repair even with stale transient metadata", () => { + const run = { errorCode: "provider_tool_definition_invalid", error: "Tool name is too long.", + resultJson: { errorFamily: "transient_upstream" } }; + expect(classifyAdapterFailureForRecovery(run)).toEqual({ kind: "configuration_incomplete" }); + expect(classifyContinuationFailure(run as never)).toMatchObject({ kind: "non_retryable", maxAttempts: 0 }); + }); }); diff --git a/server/src/services/recovery/service.ts b/server/src/services/recovery/service.ts index a238785b07..6f7ac9d9ca 100644 --- a/server/src/services/recovery/service.ts +++ b/server/src/services/recovery/service.ts @@ -490,6 +490,7 @@ const TRANSIENT_INFRA_CONTINUATION_ERROR_CODES = new Set([ ]); const NON_RETRYABLE_CONTINUATION_ERROR_CODES = new Set([ + "provider_tool_definition_invalid", "adapter_engine_unavailable", "agent_not_invokable", "agent_not_found", @@ -624,7 +625,7 @@ export function classifyAdapterFailureForRecovery( ): AdapterFailureRecoveryClassification { // An engine prerequisite cannot be repaired by asking the same unavailable // engine to retry. Use the existing configuration-blocker path. - if (latestRun.errorCode === "adapter_engine_unavailable") { + if (latestRun.errorCode === "adapter_engine_unavailable" || latestRun.errorCode === "provider_tool_definition_invalid") { return { kind: "configuration_incomplete" }; } if ( diff --git a/server/src/services/run-cancellation.test.ts b/server/src/services/run-cancellation.test.ts new file mode 100644 index 0000000000..18b66bcadc --- /dev/null +++ b/server/src/services/run-cancellation.test.ts @@ -0,0 +1,34 @@ +import { describe, expect, it } from "vitest"; +import { cancellationResultJson, canContinueCancelledRun, isUnexpectedRunCancellation, requestedRunCancellation } from "./run-cancellation.js"; + +describe("cancellation provenance", () => { + const run = { status: "cancelled", runtimeMode: "legacy" as const, startedAt: new Date(0), finishedAt: new Date(1), + error: "Cancelled", errorCode: "cancelled", resultJson: null }; + it("preserves operator intent over adapter completion metadata", () => { + const cancellation = requestedRunCancellation({ cancelledByUserId: "board", cancelledByActorType: "user" }, "Stop"); + const resultJson = cancellationResultJson({ ...run, resultJson: { cancellation } }, "cancelled", { status: "cancelled" }); + expect(resultJson?.cancellation).toEqual(cancellation); + expect(isUnexpectedRunCancellation({ ...run, resultJson })).toBe(false); + expect(canContinueCancelledRun({ ...run, resultJson })).toBe(false); + }); + it("records queue interruption and shutdown as expected system actions", () => { + expect(requestedRunCancellation({ queuedCommentInterruptQueueId: "queue", cancelledByUserId: "board" }, "Send messages")) + .toMatchObject({ source: "queued_message", expected: true, initiator: { type: "user", id: "board" } }); + expect(cancellationResultJson(run, "interrupted", {}, "server_shutdown_interrupted")?.cancellation) + .toMatchObject({ source: "shutdown", expected: true }); + expect(requestedRunCancellation({ queuedCommentInterruptQueueId: "queue", interruptedByActorType: "user", interruptedByActorId: "board" }, "Send messages")) + .toMatchObject({ source: "queued_message", initiator: { type: "user", id: "board" } }); + }); + it("requires positive provider evidence and a complete empty tool inventory for saved input", () => { + const resultJson = cancellationResultJson(run, "cancelled", { status: "cancelled", acpToolInventoryComplete: true, acpPendingToolCount: 0 }); + expect(resultJson?.cancellation).toMatchObject({ source: "provider", expected: false }); + expect(canContinueCancelledRun({ ...run, resultJson })).toBe(true); + expect(canContinueCancelledRun({ ...run, resultJson: { ...resultJson, acpPendingToolCount: 1 } })).toBe(false); + expect(canContinueCancelledRun({ ...run, resultJson: { ...resultJson, acpToolInventoryComplete: false } })).toBe(false); + expect(canContinueCancelledRun({ ...run, resultJson: {} })).toBe(false); + expect(isUnexpectedRunCancellation(run)).toBe(true); + expect(isUnexpectedRunCancellation({ ...run, resultJson: { + executionRecovery: { kind: "bootstrap", providerWorkStarted: false }, + } })).toBe(false); + }); +}); diff --git a/server/src/services/run-cancellation.ts b/server/src/services/run-cancellation.ts new file mode 100644 index 0000000000..b8cb8681a4 --- /dev/null +++ b/server/src/services/run-cancellation.ts @@ -0,0 +1,84 @@ +import type { heartbeatRuns } from "@paperclipai/db"; + +type Run = Pick; +const SOURCES = ["operator", "queued_message", "shutdown", "provider", "transport", "control_plane", "unknown"] as const; +export type CancellationSource = typeof SOURCES[number]; +export interface RunCancellation { + source: CancellationSource; + expected: boolean; + initiator: { type: "user" | "agent" | "system" | "provider"; id?: string }; + reason: string; + recordedAt: string; +} + +/** Server-owned local evidence. Sentry receives only the closed labels below. */ +export function readRunCancellation(result: Run["resultJson"] | undefined): RunCancellation | null { + const value = result?.cancellation; + if (!value || typeof value !== "object" || Array.isArray(value)) return null; + const entry = value as Record; + const initiator = entry.initiator as Record | undefined; + if (!SOURCES.includes(entry.source as CancellationSource) || typeof entry.expected !== "boolean" || + !initiator || !["user", "agent", "system", "provider"].includes(initiator.type as string) || + typeof entry.reason !== "string" || typeof entry.recordedAt !== "string") return null; + return { source: entry.source as CancellationSource, expected: entry.expected, + initiator: { type: initiator.type as RunCancellation["initiator"]["type"], + ...(typeof initiator.id === "string" ? { id: initiator.id } : {}) }, + reason: entry.reason.slice(0, 512), recordedAt: entry.recordedAt }; +} + +export function requestedRunCancellation(result: Run["resultJson"], reason: string): RunCancellation { + const queued = typeof result?.queuedCommentInterruptQueueId === "string"; + const actorType = result?.cancelledByActorType ?? result?.interruptedByActorType; + const user = actorType === "user" || actorType === "board" || typeof result?.cancelledByUserId === "string"; + const agent = actorType === "agent"; + const actorId = result?.cancelledByUserId ?? result?.cancelledByAgentId ?? result?.interruptedByActorId; + const source: CancellationSource = queued ? "queued_message" : user ? "operator" : "control_plane"; + return { source, expected: true, initiator: { + type: user ? "user" : agent ? "agent" : "system", + ...(typeof actorId === "string" ? { id: actorId } : {}), + }, reason: reason.slice(0, 512), recordedAt: new Date().toISOString() }; +} + +/** Preserve recorded stop intent when adapter completion replaces resultJson. */ +export function cancellationResultJson(run: Run, status: string, patch?: Run["resultJson"], errorCode?: string | null, error?: string | null): Run["resultJson"] { + if (status !== "cancelled" && status !== "interrupted") return patch ?? run.resultJson; + const result = { ...run.resultJson, ...patch }; + const previous = readRunCancellation(run.resultJson); + let cancellation = previous?.expected ? previous : readRunCancellation(patch) ?? previous; + const code = errorCode ?? run.errorCode; + if (!cancellation) { + if (code === "server_shutdown_interrupted") { + cancellation = { source: "shutdown", expected: true, initiator: { type: "system" }, + reason: "Server shutdown interrupted execution", recordedAt: new Date().toISOString() }; + } else if (result.cancelledByActorType || result.cancelledByUserId || result.interruptedByActorType || result.queuedCommentInterruptQueueId) { + cancellation = requestedRunCancellation(result, error ?? run.error ?? "Execution stopped"); + } else { + const source = code === "duplex_channel_lost" ? "transport" : result.status === "cancelled" ? "provider" : "unknown"; + const description = source === "provider" ? "Provider cancelled execution" + : source === "transport" ? "Connection to the execution provider was lost" + : "Execution stopped without a recorded cancellation request"; + const detail = error ?? run.error; + cancellation = { source, expected: false, initiator: { type: source === "provider" ? "provider" : "system" }, + reason: (!detail || /^cancelled\.?$/i.test(detail) ? description : `${description}: ${detail}`).slice(0, 512), + recordedAt: new Date().toISOString() }; + } + } + return { ...result, cancellation }; +} + +/** Eligibility only; admission still verifies process/lease shutdown and gates. + * Historical ambiguous cancellations and outstanding tools stay held. + */ +export function canContinueCancelledRun(run: Run): boolean { + const cancellation = readRunCancellation(run.resultJson); + return run.status === "cancelled" && run.runtimeMode === "legacy" && Boolean(run.startedAt && run.finishedAt) && + cancellation?.expected === false && ["provider", "transport"].includes(cancellation.source) && + run.resultJson?.acpToolInventoryComplete === true && run.resultJson?.acpPendingToolCount === 0; +} + +export function isUnexpectedRunCancellation(run: Run): boolean { + const bootstrap = run.resultJson?.executionRecovery as Record | undefined; + return run.status === "cancelled" && Boolean(run.startedAt) && + !(bootstrap?.kind === "bootstrap" && bootstrap.providerWorkStarted === false) && + readRunCancellation(run.resultJson)?.expected !== true; +} diff --git a/server/src/services/run-failure-diagnostics.ts b/server/src/services/run-failure-diagnostics.ts index cf5e7ad390..f6711e4766 100644 --- a/server/src/services/run-failure-diagnostics.ts +++ b/server/src/services/run-failure-diagnostics.ts @@ -1,4 +1,5 @@ import type { heartbeatRuns } from "@paperclipai/db"; +import { readRunCancellation } from "./run-cancellation.js"; import { WORKSPACE_RESTORE_FAILURE_CODES } from "@paperclipai/shared"; import { redactDiagnosticText } from "@paperclipai/adapter-utils/command-redaction"; import { redactCurrentUserText } from "../log-redaction.js"; @@ -143,6 +144,15 @@ export function collectRunFailureDiagnostics(run: Run, options: RunFailureReport if (Number.isFinite(durationMs) && durationMs >= 0) execution.durationMs = durationMs; } const result = run.resultJson; + const cancellation = readRunCancellation(result); + if (cancellation) { + execution.cancellationSource = cancellation.source; + execution.cancellationExpected = cancellation.expected; + execution.cancellationInitiatorType = cancellation.initiator.type; + } else if (run.status === "cancelled") { + execution.cancellationSource = "unknown"; + execution.cancellationExpected = false; + } Object.assign(execution, scalars(result, [ "mode", "stopReason", "timeoutFired", "timeoutSource", "timeoutConfigured", "effectiveTimeoutSec", "errorFamily", diff --git a/server/src/services/run-failure-report.ts b/server/src/services/run-failure-report.ts index 0c044b6615..2ed6528611 100644 --- a/server/src/services/run-failure-report.ts +++ b/server/src/services/run-failure-report.ts @@ -10,6 +10,7 @@ import { type RunFailureReportOptions, } from "./run-failure-diagnostics.js"; import { logger } from "../middleware/logger.js"; +import { isUnexpectedRunCancellation } from "./run-cancellation.js"; type HeartbeatRun = typeof heartbeatRuns.$inferSelect; @@ -30,7 +31,7 @@ const pendingRunFailureReports = new Set>(); const PENDING_REPORT_DRAIN_TIMEOUT_MS = 5_000; function isRunFailureStatus(status: string): status is RunFailureStatus { - return status === "failed" || status === "timed_out"; + return status === "failed" || status === "timed_out" || status === "cancelled"; } function readTaskId(run: HeartbeatRun): string | null { @@ -41,7 +42,7 @@ function readTaskId(run: HeartbeatRun): string | null { /** * Report a terminal run failure to Sentry. Returns at once for any status - * other than `failed` and `timed_out`. Never throws — a Sentry failure or a + * other than failures and unexpected started cancellations. Never throws — a Sentry failure or a * database read failure must not change the caller's control flow. * * Call this beside the caller's own terminal-status write, with @@ -52,6 +53,7 @@ function readTaskId(run: HeartbeatRun): string | null { */ export function reportRunFailure(db: Db, run: HeartbeatRun, options: RunFailureReportOptions = {}): Promise { if (!isRunFailureStatus(run.status)) return Promise.resolve(); + if (run.status === "cancelled" && !isUnexpectedRunCancellation(run)) return Promise.resolve(); const runStatus = run.status; const report = captureTerminalRunFailure(db, run, runStatus, options); pendingRunFailureReports.add(report); diff --git a/server/src/services/tool-gateway.ts b/server/src/services/tool-gateway.ts index 3dd5d20c24..c236fed1d6 100644 --- a/server/src/services/tool-gateway.ts +++ b/server/src/services/tool-gateway.ts @@ -1,4 +1,5 @@ import { isInsufficientConnectionScope, INSUFFICIENT_CONNECTION_SCOPE_MESSAGE } from "./connection-permission-errors.js"; +import { boundedMcpToolName } from "./mcp-tool-names.js"; import { browserUseService } from "./browser-use.js"; import { isBrowserUseConnection } from "./browser-use-client.js"; import { COGNEE_STDIO_TEMPLATE, cogneeCloudUrl, callCogneeCloud } from "./cognee-connection.js"; @@ -1384,10 +1385,14 @@ export function createToolGatewayService( ); } const baseName = baseNames[index]!; - const gatewayToolName = + const unboundedGatewayToolName = baseNameCounts.get(baseName)! > 1 ? `${baseName}-${shortStableId(catalogEntry.id)}` : baseName; + const gatewayToolName = boundedMcpToolName(unboundedGatewayToolName, [ + connection.id, application.id, catalogEntry.toolName, + ...(baseNameCounts.get(baseName)! > 1 ? [catalogEntry.id] : []), + ]); const applicationKey = application.applicationKey ?? null; const inputSchema = projectedConnectionToolInputSchema( connection, diff --git a/tests/e2e/legacy-failure-continuation.spec.ts b/tests/e2e/legacy-failure-continuation.spec.ts index c15bc556b9..e57369bf24 100644 --- a/tests/e2e/legacy-failure-continuation.spec.ts +++ b/tests/e2e/legacy-failure-continuation.spec.ts @@ -73,8 +73,11 @@ for (const action of ["task_retry", "thread_retry", "inbox_retry", "message", "q await page.goto(action === "inbox_retry" ? `/${company.issuePrefix}/inbox/all` : taskUrl); if (action === "task_retry") { const notice = page.getByRole("status", { name: "Task recovery" }); - await expect(notice).toHaveText("Automatic recovery of this task stopped.Retry"); - await expect(notice.getByRole("link")).toHaveCount(0); + await expect(notice).toContainText("Server restarted during startup"); + await expect(notice).toContainText("Automatic recovery stopped."); + await expect(notice.getByRole("link", { name: "Inspect run" })) + .toHaveAttribute("href", new RegExp(`/agents/${agent.id}/runs/${sourceRunId}$`)); + await expect(notice.getByRole("button", { name: "Retry", exact: true })).toBeEnabled(); const presentation = await notice.evaluate(element => { const style = getComputedStyle(element); return { border: style.borderTopWidth, background: style.backgroundColor }; diff --git a/ui/src/components/ExecutionBlockerNotice.test.tsx b/ui/src/components/ExecutionBlockerNotice.test.tsx index 3e8d42b420..9627098a37 100644 --- a/ui/src/components/ExecutionBlockerNotice.test.tsx +++ b/ui/src/components/ExecutionBlockerNotice.test.tsx @@ -32,12 +32,47 @@ describe("stopped task recovery notice", () => { await act(async () => { await new Promise(resolve => setTimeout(resolve, 10)); }); }); afterEach(async () => { await act(async () => root.unmount()); client.clear(); container.remove(); }); - it("shows only the requested sentence and Retry, inside a distinct recovery container", () => { + it("keeps the recovery guidance and an inspection link alongside Retry", () => { const notice = container.querySelector('[role="status"][aria-label="Task recovery"]')!; - expect(notice.textContent).toBe("Automatic recovery of this task stopped.Retry"); + expect(notice.textContent).toContain("Recorded work is preserved"); + expect(notice.textContent).toContain("Retry"); expect(notice.classList.contains("border")).toBe(true); expect(notice.classList.contains("bg-muted")).toBe(true); - expect(notice.querySelector("a")).toBeNull(); + expect(notice.querySelector("a")?.getAttribute("href")).toBe("/agents/agent/runs/failed-run"); + }); + it.each([false, true])("explains cancelled runs and saved input; continue eligibility %s", async canContinue => { + await act(async () => root.render( + + )); + expect(container.textContent).toContain("Provider cancelled execution"); + expect(container.textContent).toContain("2 saved messages are waiting"); + expect(container.querySelector("a")?.getAttribute("href")).toContain("cancelled-run"); + expect(container.querySelector("button")?.textContent ?? null).toBe(canContinue ? "Continue" : null); + if (!canContinue) expect(container.textContent).toContain("sending a new message to request continuation"); + if (canContinue) { + vi.mocked(agentsApi.retryFailedRun).mockResolvedValue({} as never); + await act(async () => container.querySelector("button")!.click()); + expect(agentsApi.retryFailedRun).toHaveBeenCalledWith("agent", "cancelled-run", "company"); + } + }); + it("keeps removed-chat guidance without suggesting a message to an unavailable destination", async () => { + await act(async () => root.render( + + )); + expect(container.textContent).toContain("create a new task"); + expect(container.textContent).not.toContain("sending a new message"); + expect(container.querySelector("button")).toBeNull(); + expect(container.querySelector("a")?.textContent).toBe("Inspect run"); }); it("keeps the required next action for other reconciliation causes", async () => { await act(async () => root.render( diff --git a/ui/src/components/ExecutionBlockerNotice.tsx b/ui/src/components/ExecutionBlockerNotice.tsx index 527ef3c6c8..d80d11b6a3 100644 --- a/ui/src/components/ExecutionBlockerNotice.tsx +++ b/ui/src/components/ExecutionBlockerNotice.tsx @@ -13,7 +13,7 @@ export function ExecutionBlockerNotice({ companyId, issueId, blocker, onRetried onRetried: () => void; }) { const queryClient = useQueryClient(); - const { data: runs } = useQuery({ + const { data: runs, error: runsError } = useQuery({ queryKey: queryKeys.issues.runs(issueId), queryFn: () => activityApi.runsForIssue(issueId), }); @@ -22,7 +22,7 @@ export function ExecutionBlockerNotice({ companyId, issueId, blocker, onRetried const requiresInspection = blocker.cause === "native_continuation_requires_reconciliation" || blocker.cause === "native_session_cleanup_quarantined"; const retry = useMutation({ - mutationFn: () => agentsApi.retryFailedRun(failedRun!.agentId, failedRun!.runId, companyId), + mutationFn: () => agentsApi.retryFailedRun(blocker.agentId!, blocker.runId!, companyId), onSuccess: () => { onRetried(); for (const queryKey of [queryKeys.issues.detail(issueId), queryKeys.issues.runs(issueId), @@ -33,22 +33,28 @@ export function ExecutionBlockerNotice({ companyId, issueId, blocker, onRetried }); return (
- {blocker.cause === "legacy_execution_requires_reconciliation" - ? "Automatic recovery of this task stopped." - : `${requiresInspection ? "Recovery needed. " : ""}${blocker.nextAction}`} - {requiresInspection && blocker.agentId && blocker.runId && ( +
+

Recovery needed.{blocker.runError ? ` ${blocker.runError}` : ""}

+

{blocker.nextAction}

+ {Boolean(blocker.savedMessageCount) && ( +

{blocker.savedMessageCount} saved {blocker.savedMessageCount === 1 ? "message is" : "messages are"} waiting for recovery.

+ )} +
+ {blocker.agentId && blocker.runId && ( )} - {!requiresInspection && failedRun && ( + {!requiresInspection && blocker.agentId && blocker.runId && + ((blocker.cause === "legacy_execution_requires_reconciliation" && failedRun) || blocker.canContinue) && ( )} {retry.isError && (

{retry.error.message}

)} + {runsError &&

{runsError.message}

}
); }