From 7a52dcdc74ef63a25e4b3e13415dfee2b0209fbe Mon Sep 17 00:00:00 2001 From: Dotta <34892728+cryppadotta@users.noreply.github.com> Date: Fri, 2 Oct 2026 13:47:59 -0500 Subject: [PATCH] fix: repair MCP validation and cancelled execution recovery (#14951) ## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - The tool gateway gives agents access to connected services. Recovery controls what happens when a run stops. > - Generated tool names can exceed the provider limit after the MCP client adds its prefix. > - The same invalid definition can fail each automatic retry. A cancelled run can also hold saved messages without showing its cause. > - This pull request bounds tool names, stops configuration retries, and retains cancellation evidence. > - It shows the stopped run and admits saved input only after the existing safety checks pass. > - The benefit is a clear recovery path that preserves operator Stop and prevents duplicate message delivery. ## Linked Issues or Issue Description **What happened?** A long connected MCP tool name makes the provider reject the entire request. Automatic recovery repeats the invalid request. Separately, unexpected legacy cancellations can leave saved input behind a recovery hold. The notice does not identify the stopped run or its cause. **Expected behavior** Complete MCP names fit the provider limit. Tool-definition errors require configuration repair. Cancelled runs retain their source and reason. The recovery notice shows the cause and saved-message count. Verified unexpected cancellations can start a fresh turn through the existing admission checks. **Steps to reproduce** 1. Assign an App gallery connection with a long application key and tool name to a Claude agent. 2. Start a run. The provider rejects a name over 128 characters, including its MCP prefix. 3. For cancellation recovery, stop a legacy provider turn without an operator Stop request and send a user message while the recovery hold is active. 4. Inspect the recovery notice and the deferred message queue. **Paperclip version or commit** Rebased onto master at `cf8ad63c806685bfd7c48e3ed4a919d61a7c55f1`. **Deployment mode** Hosted or self-hosted server with legacy Claude or Codex execution. Related public work: - Refs #14017. That PR caps name segments. This PR preserves existing short names and uses stable hash aliases for long complete names. It also covers classification and recovery. - Refs #4510. That PR adds a cancellation-source column. This PR records bounded evidence in the existing run result, without a migration. - Refs #12552 and #4506. Those PRs suppress recovery after operator cancellation. This PR preserves operator intent and uses the existing continuation gates. ## What Changed - Bound gateway names with the full provider prefix in the 128-character budget. Retain the original upstream tool name for dispatch and permissions. - Classify invalid tool definitions as configuration failures before diagnostic redaction. Stop automatic retries and continuation attempts for that error code. - Persist cancellation source, expectedness, initiator, reason, and time. Preserve recorded Stop intent when adapter results arrive. Report unexpected started cancellations with closed diagnostic labels. - Show the run cause, saved-message count, and Inspect run link. Offer Continue for eligible unexpected cancellations. Require verified provider stop, empty tool inventory, ownership, and the existing pause, budget, approval, and dependency gates. Use the existing queue for single delivery. - Add regression coverage and update the execution, MCP gateway, and run-log documentation. ## Verification - `pnpm -r typecheck` and `pnpm build` passed. - `pnpm check:token-gates` passed. - Ran `pnpm test:run` and completed its workspace and serialized groups. Initial resource and timing failures passed on isolated reruns. All 149 serialized route suites passed. - Reran the changed server, adapter, and UI suites after the rebase. Coverage includes long-name upstream dispatch, configuration retry suppression, cancellation evidence retention, privacy labels, oversized run projection, and concurrent saved-message delivery. - `pnpm test:e2e tests/e2e/legacy-failure-continuation.spec.ts` passed all six browser scenarios. The recovery notice shows the run cause and inspection link, and each recovery entry point reaches one new response. - Added database-backed checks for active, removed, paused, unavailable, and disabled chat connections. The final continuation and recovery-notice suites passed 167 tests. Externally bound chats hide board Continue and show a usable next action. - All 55 GitHub checks passed on `42afbf1371dcaeb72646e3d8f65c19ff7cddf8de`. Two unrelated Storybook jobs were skipped by their normal conditions. Greptile reviewed that commit at 5/5 with no findings and no open review threads. ## Risks - Long tool names change to aliases. Existing short names stay compatible. The original connection and upstream name remain the dispatch authority. - Invalid tool definitions no longer get automatic retries. An operator must repair the configuration before a new attempt. - Continuation changes apply only to positively identified unexpected legacy cancellations with complete empty tool inventory. Operator Stop, unknown historical cancellations, outstanding tools, and unverified provider termination keep their holds. - No database migration. The added projection fields are optional. Cancellation reason and initiator IDs remain local run evidence; Sentry receives only closed source and initiator-type labels and expectedness. ## Model Used - OpenAI GPT-6 through Codex, with reasoning, repository editing, shell execution, and GitHub tool use. The runtime does not expose the exact model variant or context-window size. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip --- doc/MCP-ACCESS-GOVERNANCE.md | 7 ++ doc/execution-semantics.md | 16 ++++ doc/run-log-events.md | 15 ++++ .../src/acpx-engine/execute.test.ts | 16 ++++ .../adapter-utils/src/acpx-engine/execute.ts | 4 +- .../terminal-session-failure.test.ts | 22 ++++- .../acpx-engine/terminal-session-failure.ts | 18 ++++ packages/adapter-utils/src/types.ts | 1 + .../shared/src/types/execution-projection.ts | 5 ++ ...at-direct-adapter-native-isolation.test.ts | 22 +++++ server/src/__tests__/heartbeat-list.test.ts | 12 +++ .../heartbeat-retry-scheduling.test.ts | 10 +++ server/src/__tests__/tool-gateway.test.ts | 25 ++++++ server/src/routes/agents.ts | 12 ++- server/src/sentry.ts | 2 +- .../__tests__/run-failure-report.test.ts | 19 +++++ server/src/services/execution-blocker.ts | 62 +++++++++++++- .../explicit-native-continuation.test.ts | 80 +++++++++++++++++- .../services/explicit-native-continuation.ts | 4 +- server/src/services/heartbeat.ts | 82 ++++++++++++++---- server/src/services/mcp-tool-names.test.ts | 17 ++++ server/src/services/mcp-tool-names.ts | 17 ++++ .../provider-failure-classification.test.ts | 6 ++ server/src/services/recovery/service.ts | 3 +- server/src/services/run-cancellation.test.ts | 34 ++++++++ server/src/services/run-cancellation.ts | 84 +++++++++++++++++++ .../src/services/run-failure-diagnostics.ts | 10 +++ server/src/services/run-failure-report.ts | 6 +- server/src/services/tool-gateway.ts | 7 +- tests/e2e/legacy-failure-continuation.spec.ts | 7 +- .../ExecutionBlockerNotice.test.tsx | 41 ++++++++- ui/src/components/ExecutionBlockerNotice.tsx | 22 +++-- 32 files changed, 640 insertions(+), 48 deletions(-) create mode 100644 server/src/services/mcp-tool-names.test.ts create mode 100644 server/src/services/mcp-tool-names.ts create mode 100644 server/src/services/run-cancellation.test.ts create mode 100644 server/src/services/run-cancellation.ts 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}

}
); }