From bd6caf51bb0171b2a699565225d4e93698c25a66 Mon Sep 17 00:00:00 2001 From: Devin Foley Date: Fri, 25 Sep 2026 11:34:31 -0700 Subject: [PATCH] fix: preserve restore failure results and stop unsafe retries (#14035) Preserve agent output and earlier execution errors when workspace restore fails. Report the restore phase and confirmed saved-plan links. Require verified repair before retrying unsafe archives, while preserving approval states and the retry budget. Verified with full CI, 506 focused regression tests, and Greptile 5/5 with all review threads resolved. Co-Authored-By: Paperclip --- doc/run-log-events.md | 26 +++++++ .../src/workspace-restore-merge.test.ts | 18 +++++ .../src/workspace-restore-merge.ts | 18 ++++- .../src/workspace-restore-result.test.ts | 56 ++++++++++++++ .../src/workspace-restore-result.ts | 68 +++++++++++++++++ .../grok-local/src/server/execute.test.ts | 27 +++++-- .../adapters/grok-local/src/server/execute.ts | 37 +++------ packages/shared/src/index.ts | 2 + .../shared/src/types/execution-projection.ts | 2 + packages/shared/src/validators/issue.ts | 1 + packages/shared/src/workspace-restore.ts | 21 +++++ server/src/__tests__/heartbeat-list.test.ts | 6 ++ .../heartbeat-retry-scheduling.test.ts | 24 ++++++ server/src/services/activity.ts | 29 ++++++- .../src/services/conversation-continuation.ts | 4 +- server/src/services/execution-blocker.ts | 12 ++- .../services/execution-recovery-resolution.ts | 9 ++- .../explicit-native-continuation.test.ts | 19 +++++ .../services/explicit-native-continuation.ts | 3 + server/src/services/heartbeat.ts | 44 +++++++++-- .../legacy-execution-recovery.test.ts | 15 ++++ .../src/services/legacy-execution-recovery.ts | 23 ++++-- .../native-safe-replacement.test.ts | 76 +++++++++++++++++++ ui/src/components/TaskChatThread.test.tsx | 38 +++++++++- ui/src/components/TaskChatThread.tsx | 35 ++++++--- .../components/task-chat/TaskChatMarker.tsx | 7 +- .../components/task-chat/task-chat-model.ts | 1 + ui/src/lib/workspace-restore-marker.test.ts | 16 ++++ ui/src/lib/workspace-restore-marker.ts | 18 +++++ 29 files changed, 595 insertions(+), 60 deletions(-) create mode 100644 packages/adapter-utils/src/workspace-restore-result.test.ts create mode 100644 packages/adapter-utils/src/workspace-restore-result.ts create mode 100644 packages/shared/src/workspace-restore.ts create mode 100644 ui/src/lib/workspace-restore-marker.test.ts create mode 100644 ui/src/lib/workspace-restore-marker.ts diff --git a/doc/run-log-events.md b/doc/run-log-events.md index 6a99289429..7a7833dbd8 100644 --- a/doc/run-log-events.md +++ b/doc/run-log-events.md @@ -167,6 +167,32 @@ matching receipt in PostgreSQL and project only the run's issue/task identifiers from its context, so historical duplicate receipts cannot multiply run contexts in server memory. Existing duplicate events do not require deletion or migration. +### Workspace restore failures + +Legacy adapter results can carry `workspaceRestoreFailure` with the code +`restore_permission_denied`, `restore_lock_timeout`, `restore_unsafe_archive`, +or `restore_failed`. The heartbeat records `workspace_restore_failed` and keeps +the run failed and the workspace-finalization barrier closed. Available output, +usage, session metadata, and the previous execution outcome survive settlement. +`executionBeforeRestore` retains the earlier error code, exit code, signal, and +timeout flag. The ordinary redacted error field retains an earlier error message. + +The chat reports the restore phase separately from a missing final response. +A saved-plan link requires a stored document and its run-bound revision or a +matching run-bound review record. Older unclassified failures use neutral wording. Diagnostics show only +a validated relative member path, never an archive link target or host path. + +An unsafe archive or an outbound confinement refusal keeps the existing execution recovery hold, including across +conversation resets. It cannot start another model turn until an operator uses +the existing recovery action to record `executionReconciliation` with +`workspaceRepairEvidence` (20–12000 characters). This evidence must describe +verified safe staging or repair for the referenced failed run. It does not grant +plan approval. Saved comments, document revisions, and confirmation IDs and +states stay unchanged. Recovery uses the existing delivery identity and links +the successor to the original failed run. Repair does not reset the automatic +retry budget. Transient failures retain the existing +bounded retry policy. Archive confinement remains required. + ## Codex resume usage snapshot The native runner retains a bounded local `harness.diagnostic` event with code diff --git a/packages/adapter-utils/src/workspace-restore-merge.test.ts b/packages/adapter-utils/src/workspace-restore-merge.test.ts index 516c6dcc81..58d7fe3bbc 100644 --- a/packages/adapter-utils/src/workspace-restore-merge.test.ts +++ b/packages/adapter-utils/src/workspace-restore-merge.test.ts @@ -126,6 +126,24 @@ describe("workspace restore merge", () => { }); describe("classifyWorkspaceRestoreFailure", () => { + it.each([ + "Daytona syncOut refusing tarball with an unparseable entry listing: private listing", + "Daytona syncOut refusing unparseable or ambiguous symlink entry: private listing", + "Daytona syncOut refusing unparseable or ambiguous hardlink entry: private listing", + "Daytona syncOut refusing tarball member that escapes the extraction dir: ../private", + "Daytona syncOut refusing tarball link whose target escapes the extraction dir: link -> /private", + "Daytona sync source path is not a confined absolute path: ../private", + "Daytona sync source path escapes the workspace remote dir: /private", + ...[40, 41, 42, 44, 45].map((code) => `Daytona outbound symlink-escape guard command failed (exit ${code}): private detail`), + ])("holds the deterministic confinement refusal: %s", (message) => { + expect(classifyWorkspaceRestoreFailure(new Error(message))).toBe("restore_unsafe_archive"); + expect(describeWorkspaceRestoreFailure(classifyWorkspaceRestoreFailure(new Error(message)))).not.toContain("private"); + }); + + it("preserves the generic policy for other outbound command failures", () => { + expect(classifyWorkspaceRestoreFailure(new Error("Daytona outbound symlink-escape guard command failed (exit 1): transport failed"))).toBe("restore_failed"); + }); + it("maps an EACCES error to restore_permission_denied", () => { const error: NodeJS.ErrnoException = new Error("permission denied"); error.code = "EACCES"; diff --git a/packages/adapter-utils/src/workspace-restore-merge.ts b/packages/adapter-utils/src/workspace-restore-merge.ts index 92992e684c..ad2b7f424c 100644 --- a/packages/adapter-utils/src/workspace-restore-merge.ts +++ b/packages/adapter-utils/src/workspace-restore-merge.ts @@ -206,6 +206,7 @@ export const WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE = "ERR_WORKSPACE_RESTORE_LOCK_T export type WorkspaceRestoreFailureCode = | "restore_permission_denied" | "restore_lock_timeout" + | "restore_unsafe_archive" | "restore_failed"; /** @@ -222,13 +223,24 @@ export type WorkspaceRestoreOutcome = * `EACCES` and `EPERM` to a permission failure, the merge-lock timeout * (matched by {@link WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE}, never by the error * message text) to a lock-timeout failure, and every other error to a generic - * failure. Never reads or returns `Error.message`, a filesystem path, or a - * process id. + * failure. The known Daytona confinement diagnostic also identifies unsafe + * archives across plugin transports that retain only a message. Never returns + * raw messages, paths or process IDs. */ export function classifyWorkspaceRestoreFailure(error: unknown): WorkspaceRestoreFailureCode { const code = error && typeof error === "object" ? (error as NodeJS.ErrnoException).code : undefined; if (code === "EACCES" || code === "EPERM") return "restore_permission_denied"; if (code === WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE) return "restore_lock_timeout"; + const message = error instanceof Error ? error.message : ""; + const archiveRefused = /Daytona syncOut refusing (?:tarball (?:with an unparseable entry listing|(?:link whose target|member that) escapes the extraction dir)|unparseable or ambiguous (?:sym|hard)link entry)/.test(message); + const outboundPathRefused = /Daytona sync source path (?:is not a confined absolute path|escapes the workspace remote dir):/.test(message); + // These are the fail-closed guard's own exit codes. Transport/command failures + // with other exit codes retain the existing transient failure policy. + const outboundGuardRefused = /Daytona outbound symlink-escape guard command failed \(exit (?:40|41|42|44|45)\)/.test(message); + if (code === "WORKSPACE_RESTORE_UNSAFE_ARCHIVE" || + archiveRefused || outboundPathRefused || outboundGuardRefused) { + return "restore_unsafe_archive"; + } return "restore_failed"; } @@ -246,6 +258,8 @@ export function describeWorkspaceRestoreFailure(code: WorkspaceRestoreFailureCod return "the restore could not write to the workspace (permission denied)"; case "restore_lock_timeout": return "the restore timed out waiting for the workspace merge lock"; + case "restore_unsafe_archive": + return "the archive contains an unsafe link or path; workspace repair is required"; case "restore_failed": return "the restore failed"; } diff --git a/packages/adapter-utils/src/workspace-restore-result.test.ts b/packages/adapter-utils/src/workspace-restore-result.test.ts new file mode 100644 index 0000000000..98b611789c --- /dev/null +++ b/packages/adapter-utils/src/workspace-restore-result.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from "vitest"; +import { applyWorkspaceRestoreFailure, withWorkspaceRestore } from "./workspace-restore-result.js"; +import type { AdapterExecutionResult } from "./types.js"; + +const completed: AdapterExecutionResult = { + exitCode: 0, signal: null, timedOut: false, + summary: "The plan is ready.", sessionId: "session", usage: { inputTokens: 4, outputTokens: 9 }, + resultJson: { requestId: "request" }, +}; +const unsafe = new Error("Daytona syncOut refusing tarball link whose target escapes the extraction dir: .claude/skills/paperclip -> /tmp/private-clone/secret-key"); + +describe("workspace restore settlement", () => { + it("retains the completed result and safe member while excluding the unsafe target", async () => { + const result = await withWorkspaceRestore(async () => completed, async () => { throw unsafe; }); + expect(result).toMatchObject({ + summary: completed.summary, sessionId: "session", usage: completed.usage, + errorCode: "workspace_restore_failed", timedOut: false, + resultJson: { requestId: "request", workspaceRestoreFailure: "restore_unsafe_archive", + workspaceRestorePath: ".claude/skills/paperclip", finalResponseRecorded: true }, + }); + expect(JSON.stringify(result)).not.toMatch(/private-clone|secret-key|\/tmp/); + expect(applyWorkspaceRestoreFailure(result)).toBe(result); + }); + + it("retains the previous execution failure and restore classification", async () => { + const result = await withWorkspaceRestore(async () => ({ ...completed, exitCode: 2, errorCode: "model_error", errorMessage: "Model request failed." }), async () => { throw unsafe; }); + expect(result.errorMessage).toContain("Model request failed."); + expect(result.errorMessage).toContain("Workspace restore failed."); + expect(result.resultJson?.executionBeforeRestore).toMatchObject({ errorCode: "model_error", exitCode: 2 }); + }); + + it("retains a thrown execution error when restore also fails", async () => { + const result = await withWorkspaceRestore(async () => { throw new Error("Process failed."); }, async () => { throw unsafe; }); + expect(result.errorMessage).toContain("Process failed."); + expect(result.resultJson).toMatchObject({ finalResponseRecorded: false, executionBeforeRestore: { errorCode: "adapter_failed" } }); + }); + + it("keeps timeout evidence while making restore the terminal failure phase", async () => { + const result = await withWorkspaceRestore(async () => ({ ...completed, timedOut: true }), async () => { throw unsafe; }); + expect(result.timedOut).toBe(false); + expect(result.resultJson?.executionBeforeRestore).toMatchObject({ timedOut: true }); + }); + + it("preserves the original result or thrown error after a clean restore", async () => { + expect(await withWorkspaceRestore(async () => completed, async () => {})).toBe(completed); + const error = new Error("execution failed"); + await expect(withWorkspaceRestore(async () => { throw error; }, async () => {})).rejects.toBe(error); + }); + + it.each(["/tmp/private/file", "../escape", "C:\\private\\file"])("omits unsafe member %s", async (member) => { + const result = await withWorkspaceRestore(async () => completed, async () => { + throw new Error(`Daytona syncOut refusing tarball link whose target escapes the extraction dir: ${member} -> /secret`); + }); + expect(result.resultJson?.workspaceRestorePath).toBeUndefined(); + }); +}); diff --git a/packages/adapter-utils/src/workspace-restore-result.ts b/packages/adapter-utils/src/workspace-restore-result.ts new file mode 100644 index 0000000000..6db30e1879 --- /dev/null +++ b/packages/adapter-utils/src/workspace-restore-result.ts @@ -0,0 +1,68 @@ +import { hasWorkspaceRestoreFailure, safeWorkspaceRestorePath } from "@paperclipai/shared"; +import type { AdapterExecutionResult } from "./types.js"; +import { classifyWorkspaceRestoreFailure } from "./workspace-restore-merge.js"; + +/** A completed model turn does not imply that required workspace files arrived. */ +export function applyWorkspaceRestoreFailure(result: AdapterExecutionResult): AdapterExecutionResult { + if (!hasWorkspaceRestoreFailure(result.resultJson) || result.errorCode === "workspace_restore_failed") return result; + return { + ...result, + timedOut: false, + errorCode: "workspace_restore_failed", + errorMessage: [result.errorMessage, "Workspace restore failed. Workspace files need recovery."].filter(Boolean).join(" "), + resultJson: { + ...result.resultJson, + executionBeforeRestore: { + errorCode: result.errorCode ?? null, + exitCode: result.exitCode, + signal: result.signal, + timedOut: result.timedOut, + }, + finalResponseRecorded: typeof result.resultJson?.finalResponseRecorded === "boolean" + ? result.resultJson.finalResponseRecorded : Boolean(result.summary?.trim()), + }, + }; +} + +/** Preserve a pending result (or earlier error) when copy-back fails. */ +export async function withWorkspaceRestore( + execute: () => Promise, + restore: () => Promise, +): Promise { + let result: AdapterExecutionResult | undefined; + let executionError: unknown; + let executionThrew = false; + try { + result = await execute(); + } catch (error) { + executionError = error; + executionThrew = true; + } + try { + await restore(); + } catch (error) { + const code = classifyWorkspaceRestoreFailure(error); + // Parse the archive member only. Never expose the unsafe link target. + const member = error instanceof Error + ? /Daytona syncOut refusing tarball link whose target escapes the extraction dir: (.+?) -> /.exec(error.message)?.[1] + : null; + const relativePath = safeWorkspaceRestorePath(member); + return applyWorkspaceRestoreFailure({ + ...(result ?? { + exitCode: null, + signal: null, + timedOut: false, + errorCode: "adapter_failed", + // The server applies its ordinary execution-error redaction to this field. + errorMessage: executionError instanceof Error ? executionError.message : "Adapter execution failed.", + }), + resultJson: { + ...result?.resultJson, + workspaceRestoreFailure: code, + ...(relativePath ? { workspaceRestorePath: relativePath } : {}), + }, + }); + } + if (executionThrew) throw executionError; + return result!; +} diff --git a/packages/adapters/grok-local/src/server/execute.test.ts b/packages/adapters/grok-local/src/server/execute.test.ts index 7cd6e07fb9..b919195514 100644 --- a/packages/adapters/grok-local/src/server/execute.test.ts +++ b/packages/adapters/grok-local/src/server/execute.test.ts @@ -757,12 +757,21 @@ describe("grok_local execute", () => { expect(await pathExists(stagedDir)).toBe(false); }); - it("removes the staged home when the workspace restore rejects during teardown", async () => { + it.each(["completed", "failed", "timed_out"])("preserves %s output and removes the staged home when restore fails", async (state) => { delete process.env.XAI_API_KEY; mocks.state.isRemote = true; await seedHostGrokAuth("{}"); let stagedDir = ""; - runProcessMock.mockImplementation(async () => makeSuccessfulRunResult()); + runProcessMock.mockImplementation(async () => ({ + ...makeSuccessfulRunResult(), + exitCode: state === "failed" ? 2 : 0, + timedOut: state === "timed_out", + stderr: state === "failed" ? "Model request failed." : "", + stdout: [JSON.stringify({ type: "text", data: "Saved output." }), JSON.stringify({ + type: "end", sessionId: "sess-1", requestId: "req-1", stopReason: state === "completed" ? "EndTurn" : null, + usage: { input_tokens: 4, output_tokens: 9 }, + })].join("\n"), + })); prepareRuntimeMock.mockImplementationOnce(async (input: { assets?: Array<{ localDir: string }> }) => { stagedDir = input.assets?.[0]?.localDir ?? ""; return { @@ -774,9 +783,17 @@ describe("grok_local execute", () => { }; }); - await expect(execute(await makeCtx("run-remote-teardown-restore-reject", await makeTempRoot()))).rejects.toThrow( - "restore failed", - ); + const result = await execute(await makeCtx("run-remote-teardown-restore-reject", await makeTempRoot())); + expect(result).toMatchObject({ + errorCode: "workspace_restore_failed", + sessionId: "sess-1", + summary: "Saved output.", + usage: { inputTokens: 4, outputTokens: 9 }, + resultJson: { workspaceRestoreFailure: "restore_failed", requestId: "req-1", finalResponseRecorded: state === "completed", + executionBeforeRestore: { exitCode: state === "failed" ? 2 : 0, timedOut: state === "timed_out" } }, + }); + if (state === "failed") expect(result.errorMessage).toContain("Model request failed."); + if (state === "timed_out") expect(result.errorMessage).toContain("Timed out after"); expect(stagedDir).not.toBe(""); expect(await pathExists(stagedDir)).toBe(false); diff --git a/packages/adapters/grok-local/src/server/execute.ts b/packages/adapters/grok-local/src/server/execute.ts index 4cd9934750..76c85afac3 100644 --- a/packages/adapters/grok-local/src/server/execute.ts +++ b/packages/adapters/grok-local/src/server/execute.ts @@ -1,3 +1,4 @@ +import { withWorkspaceRestore } from "@paperclipai/adapter-utils/workspace-restore-result"; import fs from "node:fs/promises"; import path from "node:path"; import { fileURLToPath } from "node:url"; @@ -255,7 +256,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise => { const envConfig = parseObject(config.env); const env: Record = { ...buildPaperclipEnv(agent), @@ -592,17 +593,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise { - if (attempt.proc.timedOut) { - return { - exitCode: attempt.proc.exitCode, - signal: attempt.proc.signal, - timedOut: true, - errorMessage: `Timed out after ${timeoutSec}s`, - clearSession: clearSessionOnMissingSession, - }; - } - - const failed = (attempt.proc.exitCode ?? 0) !== 0; + const failed = attempt.proc.timedOut || (attempt.proc.exitCode ?? 0) !== 0; const parsedError = typeof attempt.parsed.errorMessage === "string" ? attempt.parsed.errorMessage.trim() : ""; const stderrLine = firstNonEmptyLine(attempt.proc.stderr); const fallbackErrorMessage = @@ -631,8 +622,8 @@ export async function execute(ctx: AdapterExecutionContext): Promise { await restoreRemoteWorkspace?.(); }); } finally { - // Remove the staged GROK_HOME allowlist temp dir first, before the - // `Promise.all` below. A rejecting member of that `Promise.all` (for - // example a failed workspace restore) throws out of this `finally` and - // skips every statement after it, so the removal must run before that - // await to hold on every exit path (teardown AND error), never only the - // happy path. Cleanup failure is logged, not fatal — a leaked temp dir - // must not crash the run. + // Cleanup runs after settlement on both success and failure. if (stagedGrokHomeDir) { await fs.rm(stagedGrokHomeDir, { recursive: true, force: true }).catch(async (error) => { await onLog( @@ -696,9 +686,6 @@ export async function execute(ctx: AdapterExecutionContext): Promise | null | undefined): boolean { + return WORKSPACE_RESTORE_FAILURE_CODES.some((code) => result?.workspaceRestoreFailure === code); +} + +/** Display only ordinary repository paths, never targets, host paths or temp IDs. */ +export function safeWorkspaceRestorePath(value: unknown): string | null { + if (typeof value !== "string" || value.length > 180) return null; + const relative = value.replace(/^\.\//, ""); + if (!relative || !/^[a-zA-Z0-9_.\/-]+$/.test(relative) || relative.startsWith("/")) return null; + if (relative.split("/").some((part) => !part || part === "." || part === "..")) return null; + if (/(?:[a-f0-9]{8}-[a-f0-9-]{27,}|[a-zA-Z0-9_-]{32,}|(?:^|\/)(?:tmp|temp|paperclip-clone)[^/]*)(?:\/|$)/i.test(relative)) return null; + return relative; +} diff --git a/server/src/__tests__/heartbeat-list.test.ts b/server/src/__tests__/heartbeat-list.test.ts index e45e4e958a..465514840b 100644 --- a/server/src/__tests__/heartbeat-list.test.ts +++ b/server/src/__tests__/heartbeat-list.test.ts @@ -264,6 +264,9 @@ describeEmbeddedPostgres("heartbeat list", () => { summary: "completed", stdout: oversizedStdout, nestedHuge: { payload: oversizedNestedPayload }, + workspaceRestoreFailure: "restore_unsafe_archive", + finalResponseRecorded: true, + executionBeforeRestore: { errorCode: "model_error", exitCode: 2, timedOut: false }, }, }); @@ -275,6 +278,9 @@ describeEmbeddedPostgres("heartbeat list", () => { truncated: true, truncationReason: "oversized_result_json", stdoutTruncated: true, + workspaceRestoreFailure: "restore_unsafe_archive", + finalResponseRecorded: true, + executionBeforeRestore: { errorCode: "model_error", exitCode: 2, timedOut: false }, }); expect(typeof result?.stdout).toBe("string"); expect((result?.stdout as string).length).toBeLessThan(oversizedStdout.length); diff --git a/server/src/__tests__/heartbeat-retry-scheduling.test.ts b/server/src/__tests__/heartbeat-retry-scheduling.test.ts index 3289bf0774..faac738638 100644 --- a/server/src/__tests__/heartbeat-retry-scheduling.test.ts +++ b/server/src/__tests__/heartbeat-retry-scheduling.test.ts @@ -250,6 +250,30 @@ describeEmbeddedPostgres("heartbeat bounded retry scheduling", () => { }); } + + 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"); + await seedRetryFixture({ runId, companyId, agentId, now, errorCode: "workspace_restore_failed", resultJson: { + workspaceRestoreFailure: classification, conversationContinuation: "continue_conversation_v1", errorFamily: "transient_upstream", + } }); + const issueId = randomUUID(); + await db.insert(issues).values({ id: issueId, companyId, title: "Restore fixture", status: "in_progress", assigneeAgentId: agentId }); + await db.update(heartbeatRuns).set({ contextSnapshot: { issueId, wakeReason: "issue_assigned" } }).where(eq(heartbeatRuns.id, runId)); + const result = await heartbeat.scheduleBoundedRetry(runId, { now, random: () => 0 }); + if (classification === "restore_unsafe_archive") { + expect(result).toMatchObject({ outcome: "not_scheduled", errorCode: "legacy_execution_requires_reconciliation" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + } else { + expect(result.outcome).toBe("scheduled"); + if (result.outcome !== "scheduled" || !result.run) throw new Error("Expected a bounded retry"); + await db.update(heartbeatRuns).set({ status: "failed", errorCode: "workspace_restore_failed", scheduledRetryAttempt: 2, + resultJson: { workspaceRestoreFailure: classification, conversationContinuation: "continue_conversation_v1", errorFamily: "transient_upstream" }, + }).where(eq(heartbeatRuns.id, result.run.id)); + expect(await heartbeat.scheduleBoundedRetry(result.run.id, { now, random: () => 0 })).toMatchObject({ outcome: "retry_exhausted" }); + } + }); + it("reuses one failure successor across concurrent and repeated scheduling", async () => { const runId = randomUUID(), companyId = randomUUID(), agentId = randomUUID(); const now = new Date("2026-04-20T12:00:00.000Z"); diff --git a/server/src/services/activity.ts b/server/src/services/activity.ts index fef52abaea..99afc0a57d 100644 --- a/server/src/services/activity.ts +++ b/server/src/services/activity.ts @@ -5,6 +5,7 @@ import { activityLog, agents, documentRevisions, + documents, environmentLeases, environments, heartbeatRunEvents, @@ -15,7 +16,7 @@ import { issueWorkProducts, workspaceOperations, } from "@paperclipai/db"; -import { ISSUE_CONTINUATION_SUMMARY_DOCUMENT_KEY } from "@paperclipai/shared"; +import { hasWorkspaceRestoreFailure, safeWorkspaceRestorePath, ISSUE_CONTINUATION_SUMMARY_DOCUMENT_KEY } from "@paperclipai/shared"; import { logger } from "../middleware/logger.js"; import { visibleIssueCondition } from "./issue-visibility.js"; import { classifyRunLiveness } from "./run-liveness.js"; @@ -87,6 +88,13 @@ export function activityService(db: Db) { when ${heartbeatRuns.resultJson} is null then null else jsonb_strip_nulls(jsonb_build_object( 'conversationReset', ${heartbeatRuns.resultJson} -> 'conversationReset', + 'workspaceRestoreFailure', case when ${heartbeatRuns.resultJson} ->> 'workspaceRestoreFailure' + in ('restore_permission_denied', 'restore_lock_timeout', 'restore_unsafe_archive', 'restore_failed') + then ${heartbeatRuns.resultJson} -> 'workspaceRestoreFailure' end, + 'workspaceRestorePath', case when length(${heartbeatRuns.resultJson} ->> 'workspaceRestorePath') <= 180 + then ${heartbeatRuns.resultJson} -> 'workspaceRestorePath' end, + 'finalResponseRecorded', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'finalResponseRecorded') = 'boolean' + then ${heartbeatRuns.resultJson} -> 'finalResponseRecorded' end, 'billingType', coalesce(${heartbeatRuns.resultJson} -> 'billingType', ${heartbeatRuns.resultJson} -> 'billing_type'), 'billing_type', coalesce(${heartbeatRuns.resultJson} -> 'billing_type', ${heartbeatRuns.resultJson} -> 'billingType'), 'costUsd', coalesce( @@ -489,6 +497,16 @@ export function activityService(db: Db) { } const executionByRunId = await executionProjectionsForRuns(db, companyId, runIds); + // Only stored, current plan revisions can support a saved-plan link. + // Do not trust an adapter's claim that it wrote a document. + const [savedPlan] = runs.some((run) => hasWorkspaceRestoreFailure(run.resultJson)) + ? await db.select({ revisionId: documentRevisions.id, runId: documentRevisions.createdByRunId }) + .from(issueDocuments) + .innerJoin(documents, and(eq(documents.id, issueDocuments.documentId), eq(documents.companyId, companyId))) + .innerJoin(documentRevisions, and(eq(documentRevisions.id, documents.latestRevisionId), eq(documentRevisions.documentId, documents.id), eq(documentRevisions.companyId, companyId))) + .where(and(eq(issueDocuments.companyId, companyId), eq(issueDocuments.issueId, issueId), eq(issueDocuments.key, "plan"))) + .limit(1) + : []; return runs.map((run) => { const leaseRow = leaseByRunId.get(run.runId); const leaseMetadata = leaseRow?.lease.metadata ?? null; @@ -500,6 +518,15 @@ export function activityService(db: Db) { : null; return { ...run, + resultJson: run.resultJson ? { + ...run.resultJson, + ...(Object.hasOwn(run.resultJson, "workspaceRestorePath") ? { + workspaceRestorePath: safeWorkspaceRestorePath(run.resultJson.workspaceRestorePath), + } : {}), + ...(hasWorkspaceRestoreFailure(run.resultJson) ? { + ...(savedPlan?.runId === run.runId ? { savedPlanRevisionId: savedPlan.revisionId } : {}), + } : {}), + } : null, execution: executionByRunId.get(run.runId) ?? null, environment: leaseRow ? { diff --git a/server/src/services/conversation-continuation.ts b/server/src/services/conversation-continuation.ts index 7ff78ff6e9..3dd737357f 100644 --- a/server/src/services/conversation-continuation.ts +++ b/server/src/services/conversation-continuation.ts @@ -16,7 +16,7 @@ export function isConversationAdapter(adapterType: string): boolean { export const CONVERSATION_CONTINUATION_POLICY = "continue_conversation_v1"; export function hasConversationContinuationPolicy(result: Record | null | undefined): boolean { - return result?.conversationContinuation === CONVERSATION_CONTINUATION_POLICY; + return result?.workspaceRestoreFailure !== "restore_unsafe_archive" && result?.conversationContinuation === CONVERSATION_CONTINUATION_POLICY; } /** Persisted by the server when it claims the run, before remote provisioning. */ @@ -52,6 +52,7 @@ export async function historicalAdapterType(db: Db, run: typeof heartbeatRuns.$i } export async function runUsedConversationAdapter(db: Db, run: typeof heartbeatRuns.$inferSelect): Promise { + if (run.resultJson?.workspaceRestoreFailure === "restore_unsafe_archive") return false; if (hasConversationContinuationPolicy(run.resultJson)) return true; const adapterType = await historicalAdapterType(db, run); return adapterType !== null && isConversationAdapter(adapterType); @@ -70,6 +71,7 @@ export function conversationRecoveryActionPredicate() { and ${heartbeatRuns.id}::text = ${issueRecoveryActions.evidence}->>'runId' and coalesce(${heartbeatRuns.nativeIssueId}::text, ${heartbeatRuns.contextSnapshot}->>'issueId') = ${issueRecoveryActions.sourceIssueId}::text and ${heartbeatRuns.runtimeMode} = 'legacy' + and coalesce(${heartbeatRuns.resultJson}->>'workspaceRestoreFailure', '') <> 'restore_unsafe_archive' and ${inArray(heartbeatRuns.status, ['failed', 'timed_out', 'interrupted', 'cancelled'])} and ${conversationRunPredicate()} and ${or( diff --git a/server/src/services/execution-blocker.ts b/server/src/services/execution-blocker.ts index 14d1e69e13..f757466799 100644 --- a/server/src/services/execution-blocker.ts +++ b/server/src/services/execution-blocker.ts @@ -19,9 +19,17 @@ export async function getExecutionBlocker(db: Db, companyId: string, issueId: st boundaryId: issues.conversationBoundaryCommentId }).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 + // survives conversation boundaries until the existing repair/reconciliation path clears it. + const [restoreHold] = await db.select().from(issueRecoveryActions).where(and( + eq(issueRecoveryActions.companyId, companyId), + eq(issueRecoveryActions.sourceIssueId, issueId), + executionBlockerPredicate(), + sql`${issueRecoveryActions.evidence}->>'workspaceRestoreFailure' = 'restore_unsafe_archive'`, + )).orderBy(desc(issueRecoveryActions.updatedAt)).limit(1); // A persisted user /new is an ordered context command, not a retry of uncertain work. // The normal issue execution lock still serializes it behind any active turn. - if (conversation?.agentId && options?.conversationResetCommentId) { + if (!restoreHold && conversation?.agentId && options?.conversationResetCommentId) { const [command] = await db.select().from(issueComments).where(and( eq(issueComments.companyId, companyId), eq(issueComments.issueId, issueId), eq(issueComments.id, options.conversationResetCommentId), @@ -36,7 +44,7 @@ export async function getExecutionBlocker(db: Db, companyId: string, issueId: st const ownership = await getConversationOwnershipBlocker(db, companyId, issueId); if (ownership) return { ...ownership, recoveryActionId: null }; - const [action] = await db.select().from(issueRecoveryActions).where(and( + const [action] = restoreHold ? [restoreHold] : await db.select().from(issueRecoveryActions).where(and( eq(issueRecoveryActions.companyId, companyId), eq(issueRecoveryActions.sourceIssueId, issueId), executionBlockerPredicate(), diff --git a/server/src/services/execution-recovery-resolution.ts b/server/src/services/execution-recovery-resolution.ts index e37aa0a66c..703d241291 100644 --- a/server/src/services/execution-recovery-resolution.ts +++ b/server/src/services/execution-recovery-resolution.ts @@ -1,3 +1,4 @@ +import { hasWorkspaceRestoreFailure } from "@paperclipai/shared"; import { randomUUID } from "node:crypto"; import { conversationRecoveryActionPredicate, getConversationOwnershipBlocker } from "./conversation-continuation.js"; import { persistActivity } from "./activity-log.js"; @@ -70,6 +71,10 @@ export async function validateExecutionReconciliation(input: { "The recovery source or task owner changed. Inspect the current execution before continuing.", ); } + if (hasWorkspaceRestoreFailure(run.resultJson) && + (!decision.workspaceRepairEvidence || decision.workspaceRepairEvidence.trim().length < 20)) { + throw conflict("Verify safe workspace staging or repair and record workspaceRepairEvidence before continuing this run."); + } for (const pid of [ run.processPid, run.processGroupId ? -run.processGroupId : null, @@ -471,7 +476,9 @@ export async function settleUnrecoverableExecutions( (!task.executionRunId || task.executionRunId === run.id) && (!task.checkoutRunId || task.checkoutRunId === run.id); const note = current - ? "Automatic recovery stopped. Recorded work is preserved; actions with unverified outcomes will not be repeated." + ? hasWorkspaceRestoreFailure(run.resultJson) + ? "Workspace repair required. Verify safe staging or repair before continuing. Saved work and approval decisions remain in force." + : "Automatic recovery stopped. Recorded work is preserved; actions with unverified outcomes will not be repeated." : "Recovery closed because the task's owner, execution, or status changed. No work was replayed."; let nativeFailureBlock = action.evidence.nativeFailureBlock; if (current) { diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index 936a81c1a8..3c8c869137 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -643,6 +643,25 @@ const support = await getEmbeddedPostgresTestSupport(); expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toBeNull(); }); + + it.each(["issue_commented", "retry_failed_run"])("does not bypass an unsafe restore hold through %s", async reason => { + const f = await seed(); + await db.update(agents).set({ adapterType: "grok_local" }).where(eq(agents.id, f.agentId)); + await db.update(heartbeatRuns).set({ runtimeMode: "legacy", processPid: null, + errorCode: "workspace_restore_failed", resultJson: { workspaceRestoreFailure: "restore_unsafe_archive", conversationContinuation: "continue_conversation_v1" }, + }).where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); + await db.update(issueRecoveryActions).set({ cause: "legacy_execution_requires_reconciliation" }).where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + const blocked = vi.fn(); + expect(await db.transaction(tx => admitExplicitNativeContinuation({ ...f, reason, + commentId: reason === "issue_commented" ? f.commentId : null, + failedRunId: reason === "retry_failed_run" ? f.sourceRunId : null, + onBlocked: blocked, db: tx as unknown as typeof db, + }))).toBeNull(); + expect(blocked).toHaveBeenCalledWith("workspace_repair_required", expect.any(String)); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ runId: f.sourceRunId }); + }); + it.each(["issue_commented", "retry_failed_run"])("continues a legacy Daytona run lost before adapter.invoke: %s", async reason => { const f = await seed(); await db.update(agents).set({ adapterType: "claude_local" }).where(eq(agents.id, f.agentId)); diff --git a/server/src/services/explicit-native-continuation.ts b/server/src/services/explicit-native-continuation.ts index 2063bcf68a..5c6c0349d1 100644 --- a/server/src/services/explicit-native-continuation.ts +++ b/server/src/services/explicit-native-continuation.ts @@ -194,6 +194,9 @@ export async function admitExplicitNativeContinuation(input: { if (!lockedRun || lockedRun.status !== run.status || lockedRun.agentId !== run.agentId || lockedRun.finishedAt?.getTime() !== run.finishedAt.getTime()) return null; run = lockedRun; + if (run.resultJson?.workspaceRestoreFailure === "restore_unsafe_archive") { + return blocked("workspace_repair_required", "Verify safe workspace staging or repair before continuing. Your message is saved."); + } const cancelledStartup = await isCancelledNativeStartup(db, run, coordinator); if (cancelledStartup) cancelledStartupIds.add(run.id); if (run.runtimeMode !== "native" && !unusedAdmission && !legacyUserTurn && !cancelledStartup) return null; diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 48a252089a..3acfcd71d4 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -1,3 +1,5 @@ +import { applyWorkspaceRestoreFailure } from "@paperclipai/adapter-utils/workspace-restore-result"; +import { hasWorkspaceRestoreFailure } from "@paperclipai/shared"; import { externalConversationStateSql, nonIdleSlackIssueCondition } from "./slack-conversation-state.js"; import { settleSlackConversation } from "./slack-conversation-lifecycle.js"; import { publicChatTaskUrl } from "./chat-task-url.js"; @@ -3423,6 +3425,21 @@ const heartbeatRunSafeResultJsonColumn = sql | null>` 'error', left(${heartbeatRuns.resultJson} ->> 'error', ${HEARTBEAT_RUN_RESULT_SUMMARY_MAX_CHARS}), 'stdout', left(${heartbeatRuns.resultJson} ->> 'stdout', ${HEARTBEAT_RUN_RESULT_OUTPUT_MAX_CHARS}), 'stderr', left(${heartbeatRuns.resultJson} ->> 'stderr', ${HEARTBEAT_RUN_RESULT_OUTPUT_MAX_CHARS}), + 'workspaceRestoreFailure', case when ${heartbeatRuns.resultJson} ->> 'workspaceRestoreFailure' + in ('restore_permission_denied', 'restore_lock_timeout', 'restore_unsafe_archive', 'restore_failed') + then ${heartbeatRuns.resultJson} -> 'workspaceRestoreFailure' end, + 'finalResponseRecorded', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'finalResponseRecorded') = 'boolean' + then ${heartbeatRuns.resultJson} -> 'finalResponseRecorded' end, + 'executionBeforeRestore', case when jsonb_typeof(${heartbeatRuns.resultJson} -> 'executionBeforeRestore') = 'object' + then jsonb_strip_nulls(jsonb_build_object( + 'errorCode', left(${heartbeatRuns.resultJson} #>> '{executionBeforeRestore,errorCode}', 128), + 'exitCode', case when jsonb_typeof(${heartbeatRuns.resultJson} #> '{executionBeforeRestore,exitCode}') = 'number' + and length(${heartbeatRuns.resultJson} #>> '{executionBeforeRestore,exitCode}') < 16 + then ${heartbeatRuns.resultJson} #> '{executionBeforeRestore,exitCode}' end, + 'signal', left(${heartbeatRuns.resultJson} #>> '{executionBeforeRestore,signal}', 50), + 'timedOut', case when jsonb_typeof(${heartbeatRuns.resultJson} #> '{executionBeforeRestore,timedOut}') = 'boolean' + then ${heartbeatRuns.resultJson} #> '{executionBeforeRestore,timedOut}' end + )) end, 'stdoutTruncated', case when length(${heartbeatRuns.resultJson} ->> 'stdout') > ${HEARTBEAT_RUN_RESULT_OUTPUT_MAX_CHARS} then to_jsonb(true) @@ -24382,9 +24399,9 @@ export function heartbeatService( if (!guardedDispatch.dispatched) return; adapterResult = await guardedDispatch.resultPromise; } - // Adapter returned cleanly, which means its workspace-restore finally - // block also ran without throwing. Record the workspace_finalize - // barrier so dependents that share this executionWorkspace can wake. + adapterResult = applyWorkspaceRestoreFailure(adapterResult); + // A returned result can include a failed restore. Keep the workspace + // barrier closed until required files have been restored. // If recording the barrier itself fails, propagate as a run failure // rather than silently leaving dependents stranded behind a missing // finalize row. @@ -24400,15 +24417,16 @@ export function heartbeatService( eq(heartbeatRuns.status, "running"), ), ); - await recordWorkspaceFinalize("succeeded"); + const workspaceFinalizeStatus = hasWorkspaceRestoreFailure(adapterResult.resultJson) ? "failed" : "succeeded"; + await recordWorkspaceFinalize(workspaceFinalizeStatus); if (adapterResult.nativeFinalization) { adapterResult.nativeFinalization.workspaceFinalizeStatus = - "succeeded"; + workspaceFinalizeStatus; try { const finalized = await finalizeNativeRun({ db, runId: run.id, - workspaceFinalizeStatus: "succeeded", + workspaceFinalizeStatus, preserveProviderAttempt: Boolean(nativeWorkspaceSync), }); await dispatchPendingNativeStatusWakeups({ @@ -26933,6 +26951,7 @@ export function heartbeatService( } let reconciledSourceRunId: string | null = null; + let reconciledRestoreRetryCount: number | null = null; if (executionReconciliationWake) { const actionId = readNonEmptyString( enrichedContextSnapshot.recoveryActionId, @@ -27036,6 +27055,15 @@ export function heartbeatService( } if (action.evidence.continuationDelivery !== "pending") return { kind: "skipped" as const }; + const [reconciledRun] = await tx.select().from(heartbeatRuns).where(and( + eq(heartbeatRuns.companyId, issue.companyId), eq(heartbeatRuns.id, sourceRunId), + )); + if (hasWorkspaceRestoreFailure(reconciledRun?.resultJson)) { + if ((readNonEmptyString(decision.workspaceRepairEvidence)?.length ?? 0) < 20) + return { kind: "skipped" as const }; + // Repair does not reset the remaining automatic retry budget. + reconciledRestoreRetryCount = executionFailureRetryCount(reconciledRun!); + } reconciledSourceRunId = sourceRunId; } @@ -27982,6 +28010,10 @@ export function heartbeatService( ...(reconciledSourceRunId ? { retryOfRunId: reconciledSourceRunId } : {}), + ...(reconciledRestoreRetryCount !== null ? { + scheduledRetryAttempt: reconciledRestoreRetryCount, + scheduledRetryReason: "transient_failure", + } : {}), }) .returning() .then((rows) => rows[0]); diff --git a/server/src/services/legacy-execution-recovery.test.ts b/server/src/services/legacy-execution-recovery.test.ts index 6df61bf3d0..52397bf506 100644 --- a/server/src/services/legacy-execution-recovery.test.ts +++ b/server/src/services/legacy-execution-recovery.test.ts @@ -71,3 +71,18 @@ it("retries a busy AI subscription only when no provider work started", () => { expect(legacyExecutionNeedsReconciliation({ ...waiting, resultJson: {} })).toBe(true); expect(legacyExecutionNeedsReconciliation({ ...waiting, resultJson: { executionRecovery: { kind: "ai_connection_wait", providerWorkStarted: true } } })).toBe(true); }); + + +it.each(["failed", "timed_out", "cancelled"])("holds an unsafe archive after %s even with conversation or bootstrap evidence", (status) => { + expect(legacyExecutionNeedsReconciliation({ runtimeMode: "legacy", status, errorCode: "workspace_restore_failed", resultJson: { + workspaceRestoreFailure: "restore_unsafe_archive", conversationContinuation: "continue_conversation_v1", + executionRecovery: { kind: "bootstrap", providerWorkStarted: false }, stopReason: "max_turns", + } })).toBe(true); +}); + + +it("retains conversation retry eligibility for a transient restore lock timeout", () => { + expect(legacyExecutionNeedsReconciliation({ runtimeMode: "legacy", status: "failed", errorCode: "workspace_restore_failed", resultJson: { + workspaceRestoreFailure: "restore_lock_timeout", conversationContinuation: "continue_conversation_v1", + } })).toBe(false); +}); diff --git a/server/src/services/legacy-execution-recovery.ts b/server/src/services/legacy-execution-recovery.ts index 1cdbeea720..d79c2d2865 100644 --- a/server/src/services/legacy-execution-recovery.ts +++ b/server/src/services/legacy-execution-recovery.ts @@ -1,7 +1,8 @@ +import { hasWorkspaceRestoreFailure } from "@paperclipai/shared"; import { normalizeMaxTurnStopReason } from "./heartbeat-stop-metadata.js"; import { hasConversationContinuationPolicy } from "./conversation-continuation.js"; import { randomUUID } from "node:crypto"; -import { and, eq, inArray, sql } from "drizzle-orm"; +import { and, eq, inArray, or, sql } from "drizzle-orm"; import { heartbeatRuns, issueRecoveryActions, issues, type Db } from "@paperclipai/db"; import { issueRecoveryActionService } from "./issue-recovery-actions.js"; import { parseIssueExecutionState } from "./issue-execution-policy.js"; @@ -20,6 +21,8 @@ export function legacyExecutionNeedsReconciliation( !["failed", "timed_out", "interrupted", "cancelled"].includes(run.status) ) return false; + // A fresh model turn cannot repair or verify unrestored files. + if (run.resultJson?.workspaceRestoreFailure === "restore_unsafe_archive") return true; // A fresh conversation turn lets the agent decide what remains. The retry // scheduler, not an action-outcome hold, owns the automatic attempt limit. if (hasConversationContinuationPolicy(run.resultJson)) return false; @@ -116,13 +119,21 @@ export async function terminalizeLegacyExecution(input: { !["done", "cancelled"].includes(task.status) ) { // Periodic stranded-work checks may revisit this terminal run before its - // reconciled continuation is dispatched. Preserve the recorded decision. + // reconciled continuation is dispatched. Preserve the recorded decision + // and an existing unsafe-workspace hold instead of creating another one. const [reconciled] = await tx.select({ id: issueRecoveryActions.id }) .from(issueRecoveryActions).where(and( eq(issueRecoveryActions.companyId, run.companyId), eq(issueRecoveryActions.sourceIssueId, task.id), eq(issueRecoveryActions.status, "resolved"), - sql`${issueRecoveryActions.evidence}->'executionReconciliation'->>'runId' = ${run.id}`, + or( + sql`${issueRecoveryActions.evidence}->'executionReconciliation'->>'runId' = ${run.id}`, + and( + sql`${issueRecoveryActions.evidence}->>'runId' = ${run.id}`, + sql`${issueRecoveryActions.evidence}->>'workspaceRestoreFailure' = 'restore_unsafe_archive'`, + sql`${issueRecoveryActions.evidence}->'automaticRecovery'->>'replay' = 'blocked'`, + ), + ), )).limit(1); if (reconciled) return updated; await issueRecoveryActionService(tx as unknown as Db).upsertSourceScoped({ @@ -137,11 +148,13 @@ export async function terminalizeLegacyExecution(input: { runId: run.id, ...(isCurrentReviewer ? { reviewParticipantAgentId: run.agentId } : {}), originalFailureCode: updated.errorCode, + ...(hasWorkspaceRestoreFailure(updated.resultJson) ? { workspaceRestoreFailure: updated.resultJson!.workspaceRestoreFailure } : {}), adapterRecovery: "unsupported_or_unknown", attempt: executionFailureRetryCount(run) + 1, }, - nextAction: - "Inspect the stopped provider and recorded actions, then reconcile their outcomes before continuing. This adapter has not established a safe resume checkpoint.", + nextAction: hasWorkspaceRestoreFailure(updated.resultJson) + ? "Verify safe workspace staging or repair, then reconcile the stopped run before continuing. Saved work and approval decisions remain in force." + : "Inspect the stopped provider and recorded actions, then reconcile their outcomes before continuing. This adapter has not established a safe resume checkpoint.", maxAttempts: 3, wakePolicy: null, supersedeOnIdentityChange: true, diff --git a/server/src/services/native-runtime/native-safe-replacement.test.ts b/server/src/services/native-runtime/native-safe-replacement.test.ts index 3d6484a1b3..732958d503 100644 --- a/server/src/services/native-runtime/native-safe-replacement.test.ts +++ b/server/src/services/native-runtime/native-safe-replacement.test.ts @@ -1,7 +1,9 @@ +import { getExecutionBlocker } from "../execution-blocker.js"; import { createRunDispatch, deriveCommentId } from "../../modules/run-dispatch/index.js"; import { buildExecutionContinuation } from "../execution-continuation.js"; import { issueService } from "../issues.js"; import { activityService } from "../activity.js"; +import { instanceSettingsService } from "../instance-settings.js"; import { buildPaperclipWakePayload, heartbeatService } from "../heartbeat.js"; import { legacyExecutionNeedsReconciliation, terminalizeLegacyExecution } from "../legacy-execution-recovery.js"; import { deliverExecutionStatuses } from "../execution-status-delivery.js"; @@ -27,6 +29,10 @@ import { heartbeatRuns, issueRecoveryActions, issueComments, + issueThreadInteractions, + issueDocuments, + documents, + documentRevisions, issues, nativeRunFinalizations, } from "@paperclipai/db"; @@ -361,6 +367,76 @@ const support = externalDatabaseUrl expect(legacyExecutionNeedsReconciliation({ ...run, status: "failed", resultJson: { executionRecovery: { kind: "bootstrap", providerWorkStarted: false } } })).toBe(false); expect(legacyExecutionNeedsReconciliation({ ...run, status: "failed", scheduledRetryAttempt: 2, resultJson: { executionRecovery: { kind: "bootstrap", providerWorkStarted: false } } })).toBe(true); }); + + it.each(["pending", "accepted", "rejected", "expired"] as const)("preserves %s approval and saved work while unsafe restore recovery is repeated", async (status) => { + const source = await seed(); + const documentId = randomUUID(), revisionId = randomUUID(), interactionId = randomUUID(); + await db.insert(documents).values({ id: documentId, companyId: source.companyId, title: "Plan", latestBody: "Only the two approved changes.", latestRevisionId: revisionId }); + await db.insert(documentRevisions).values({ id: revisionId, companyId: source.companyId, documentId, revisionNumber: 1, body: "Only the two approved changes.", createdByRunId: source.runId }); + await db.insert(issueDocuments).values({ companyId: source.companyId, issueId: source.issueId, documentId, key: "plan" }); + await db.insert(issueComments).values({ companyId: source.companyId, issueId: source.issueId, authorAgentId: source.agentId, createdByRunId: source.runId, body: "The plan is saved." }); + await db.insert(issueThreadInteractions).values({ + id: interactionId, companyId: source.companyId, issueId: source.issueId, kind: "request_confirmation", status, + sourceRunId: source.runId, createdByAgentId: source.agentId, + resolvedByUserId: status === "accepted" || status === "rejected" ? "board-user" : null, + resolvedAt: status === "pending" ? null : new Date(), + payload: { version: 1, prompt: "Review the scoped plan.", target: { type: "issue_document", key: "plan", revisionId } }, + result: status === "pending" ? null : { version: 1, outcome: status === "expired" ? "superseded_by_comment" : status }, + }); + const [run] = await db.update(heartbeatRuns).set({ + runtimeMode: "legacy", errorCode: "workspace_restore_failed", scheduledRetryAttempt: 1, scheduledRetryReason: "transient_failure", + runnerProfileJson: { adapterDispatch: { adapterType: "grok_local" } }, + resultJson: { workspaceRestoreFailure: "restore_unsafe_archive", conversationContinuation: "continue_conversation_v1", summary: "The plan is saved.", workspaceRestorePath: "/tmp/private-clone/file", savedPlanRevisionId: randomUUID() }, + }).where(eq(heartbeatRuns.id, source.runId)).returning(); + const [resetComment] = await db.insert(issueComments).values({ companyId: source.companyId, issueId: source.issueId, authorUserId: "board-user", body: "/new" }).returning(); + await db.update(issues).set({ conversationAgentId: source.agentId, conversationUserId: "board-user", conversationState: "active", conversationBoundaryCommentId: resetComment.id }).where(eq(issues.id, source.issueId)); + const activityRuns = await activityService(db).runsForIssue(source.companyId, source.issueId); + expect(activityRuns.find((row) => row.runId === source.runId)?.resultJson).toMatchObject({ workspaceRestoreFailure: "restore_unsafe_archive", savedPlanRevisionId: revisionId, workspaceRestorePath: null }); + const before = await db.select().from(issueThreadInteractions).where(eq(issueThreadInteractions.id, interactionId)); + const savedComments = await db.select().from(issueComments).where(eq(issueComments.issueId, source.issueId)); + for (let attempt = 0; attempt < 2; attempt++) { + await terminalizeLegacyExecution({ db, run, status: "failed" }); + await settleUnrecoverableExecutions(db); + } + const [task] = await db.select().from(issues).where(eq(issues.id, source.issueId)); + expect(task.status).toBe("blocked"); + const actions = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, source.issueId)); + expect(actions).toHaveLength(1); + const [action] = actions; + expect(action).toMatchObject({ outcome: "blocked", evidence: { automaticRecovery: { replay: "blocked" }, workspaceRestoreFailure: "restore_unsafe_archive" } }); + expect(action.nextAction).toContain("Workspace repair required"); + expect(await getExecutionBlocker(db, source.companyId, source.issueId, { conversationResetCommentId: resetComment.id })).toMatchObject({ runId: source.runId }); + const decision = { runId: source.runId, providerStopped: true as const, actionOutcome: "mixed" as const, outcomeEvidence: "Verified saved plan and comment. Workspace copy-back failed." }; + const input = { db, companyId: source.companyId, issueId: source.issueId, agentId: source.agentId, sourceRunId: source.runId, decision }; + await expect(validateExecutionReconciliation(input)).rejects.toThrow("workspaceRepairEvidence"); + const verified = { ...decision, workspaceRepairEvidence: "Verified fresh confined staging after repairing the unsafe fixture link." }; + await expect(validateExecutionReconciliation({ ...input, decision: verified })).resolves.toMatchObject({ id: source.runId }); + await markExecutionReconciliation(db, action, verified, "operator"); + // Occupy the only agent slot. Exercise real admission without a provider. + await instanceSettingsService(db).updateExperimental({ enableAgentChat: true }); + await db.update(agents).set({ status: "active", adapterType: "grok_local", runtimeConfig: { heartbeat: { wakeOnDemand: true, maxConcurrentRuns: 1 } } }).where(eq(agents.id, source.agentId)); + await db.update(issues).set({ responsibleUserId: "board-user" }).where(eq(issues.id, source.issueId)); + await db.insert(heartbeatRuns).values({ companyId: source.companyId, agentId: source.agentId, status: "running" }); + const heartbeat = heartbeatService(db); + const wake = vi.fn(heartbeat.wakeup); + await deliverReconciledExecutions(db, wake as unknown as Parameters[1]); + await deliverReconciledExecutions(db, wake as unknown as Parameters[1]); + expect(wake).toHaveBeenCalledTimes(1); + const [successor] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, source.runId)); + expect(successor).toMatchObject({ status: "queued", scheduledRetryAttempt: 1, scheduledRetryReason: "transient_failure" }); + await db.update(heartbeatRuns).set({ status: "succeeded" }).where(eq(heartbeatRuns.id, successor.id)); + const [original] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, source.runId)); + expect(original).toMatchObject({ status: "failed", errorCode: "workspace_restore_failed", resultJson: run.resultJson }); + const continuation = await buildExecutionContinuation({ db, companyId: source.companyId, issueId: source.issueId, agentId: source.agentId, context: { previousRunId: source.runId }, summary: null, exposeLowTrustRaw: false }); + if (status === "accepted" || status === "rejected") expect(continuation.humanResponses).toContainEqual(expect.objectContaining({ id: interactionId, status })); + else expect(continuation.humanResponses).not.toContainEqual(expect.objectContaining({ id: interactionId })); + expect(await db.select().from(issueThreadInteractions).where(eq(issueThreadInteractions.issueId, source.issueId))).toEqual(before); + expect(await db.select().from(issueComments).where(eq(issueComments.issueId, source.issueId))).toEqual(savedComments); + expect(await db.select().from(documentRevisions).where(eq(documentRevisions.documentId, documentId))).toHaveLength(1); + const [savedDocument] = await db.select().from(documents).where(eq(documents.id, documentId)); + expect(savedDocument).toMatchObject({ latestBody: "Only the two approved changes.", latestRevisionId: revisionId }); + }); + it("does not reopen a reconciled legacy run while continuation is pending", async () => { const source = await seed(); const [run] = await db.update(heartbeatRuns).set({ runtimeMode: "legacy", status: "cancelled" }).where(eq(heartbeatRuns.id, source.runId)).returning(); diff --git a/ui/src/components/TaskChatThread.test.tsx b/ui/src/components/TaskChatThread.test.tsx index 1df30e6039..2767b68cc3 100644 --- a/ui/src/components/TaskChatThread.test.tsx +++ b/ui/src/components/TaskChatThread.test.tsx @@ -666,6 +666,42 @@ describe("TaskChatThread runtime transcript selection", () => { ).toBeNull(); }); + + it.each(["matching", "document_only", "absent", "other_revision"] as const)("reports a restore failure with %s saved-plan evidence", (evidence) => { + if (evidence !== "absent") planState.data = planDocument(); + render( {}} issueStatus="blocked" + interactions={evidence === "absent" || evidence === "document_only" ? [] : [planReviewInteraction("accepted", evidence === "matching" ? "revision-3" : "revision-2", "restore-run")]} + onRetryFailedRun={vi.fn()} linkedRuns={[{ + runId: "restore-run", runtimeMode: "legacy", status: "failed", errorCode: "workspace_restore_failed", + agentId: "agent-1", agentName: "Runner", adapterType: "grok_local", + createdAt: "2026-08-25T18:00:00.000Z", startedAt: "2026-08-25T18:00:00.000Z", finishedAt: "2026-08-25T18:00:02.000Z", + resultJson: { workspaceRestoreFailure: "restore_unsafe_archive", workspaceRestorePath: ".claude/skills/paperclip", finalResponseRecorded: false, + ...(evidence === "document_only" ? { savedPlanRevisionId: "revision-3" } : {}), + }, + }]} />); + const marker = container.querySelector('[data-testid="task-chat-collapsible-marker"]'); + expect(marker?.textContent).toContain("Workspace restore failed"); + flushSync(() => marker!.querySelector('button[aria-expanded]')!.click()); + expect(marker?.textContent).toContain("Workspace files need recovery"); + expect(marker?.textContent).toContain("No final response was recorded"); + expect(marker?.textContent).toContain(".claude/skills/paperclip"); + expect(marker?.querySelector('a[href*="/runs/restore-run"]')).not.toBeNull(); + expect(marker?.textContent.includes("after the plan was saved")).toBe(evidence === "matching" || evidence === "document_only"); + expect(Boolean(marker?.querySelector('a[href*="document-plan"]'))).toBe(evidence === "matching" || evidence === "document_only"); + expect(marker?.querySelector('[data-testid="task-chat-run-failed-try-again"]')).toBeNull(); + }); + + it("uses neutral wording for an unclassified historical legacy failure", () => { + render( {}} linkedRuns={[{ + runId: "old-run", runtimeMode: "legacy", status: "failed", errorCode: "adapter_failed", + agentId: "agent-1", agentName: "Runner", adapterType: "grok_local", + createdAt: "2026-08-25T18:00:00.000Z", startedAt: null, finishedAt: "2026-08-25T18:00:02.000Z", + }]} />); + expect(container.textContent).toContain("The run failed"); + expect(container.textContent).not.toContain("before returning an answer"); + expect(container.textContent).not.toContain("Workspace restore failed"); + }); + it("projects the saved Plan inline at its native write_document boundary", () => { planState.data = planDocument({ updatedAt: new Date("2026-08-25T18:00:02.000Z"), @@ -1543,7 +1579,7 @@ describe("TaskChatThread runtime transcript selection", () => { container.querySelector( '[data-testid="task-chat-collapsible-marker-details"]', )?.textContent, - ).toContain("before returning an answer"); + ).toContain("The runner stopped (provider_transport_failed)."); expect( container.querySelector('[data-testid="task-chat-final-response"]'), ).toBeNull(); diff --git a/ui/src/components/TaskChatThread.tsx b/ui/src/components/TaskChatThread.tsx index 6c468c2562..6ec28f29e8 100644 --- a/ui/src/components/TaskChatThread.tsx +++ b/ui/src/components/TaskChatThread.tsx @@ -1,3 +1,5 @@ +import { hasWorkspaceRestoreFailure } from "@paperclipai/shared"; +import { workspaceRestoreMarkerDetail } from "@/lib/workspace-restore-marker"; import type { ActivityEvent } from "@paperclipai/shared"; import { useProjectCreatedItems } from "@/hooks/useProjectCreatedItems"; import { skillCreatedItems } from "@/components/task-chat/skill-created-items"; @@ -1578,13 +1580,13 @@ export function TaskChatThread(props: TaskChatThreadProps) { ? "Run timed out" : "Run failed"; const responseBoundary = sourceHasNativeResponse - ? "after returning a final response" - : "before returning an answer"; + ? " after returning a final response" + : ""; const detail = source.status === "cancelled" - ? `The run was cancelled ${responseBoundary}.` + ? `The run was cancelled${responseBoundary}.` : source.status === "interrupted" - ? `The run was interrupted ${responseBoundary}.` + ? `The run was interrupted${responseBoundary}.` : code === "native_provider_approval_required" ? "This operation requires approval, but this runner has no interactive approval handler. Review the operation and update the agent's permission setting before retrying." : code === "native_provider_model_rejected" @@ -1595,8 +1597,8 @@ export function TaskChatThread(props: TaskChatThreadProps) { : code === "provider_frame_too_large" ? "Provider output exceeded the safe limit." : source.status === "timed_out" - ? `The runner timed out ${responseBoundary} (${code}).` - : `The runner stopped ${responseBoundary} (${code}).`; + ? `The runner timed out${responseBoundary} (${code}).` + : `The runner stopped${responseBoundary} (${code}).`; const id = `${source.id}:failure`; const runAgent = meta?.agentId ? agentMap?.get(meta.agentId) @@ -1640,20 +1642,26 @@ export function TaskChatThread(props: TaskChatThreadProps) { : canRetryFailedRun ? "You can retry this message now." : "Your message is preserved."; + const restoreFailed = hasWorkspaceRestoreFailure(meta?.resultJson); + const savedPlan = Boolean(planDocument && (meta?.resultJson?.savedPlanRevisionId === planDocument.latestRevisionId || interactions?.some((interaction) => + interaction.sourceRunId === source.id && interactionTargetsPlanRevision(interaction, planDocument), + ))); const aiRequest = interactions?.find((interaction) => interaction.kind === "connection_intent" && interaction.payload.purpose === "ai" && interaction.sourceRunId === source.id); - const detail = aiRequest + const detail = restoreFailed + ? workspaceRestoreMarkerDetail({ result: meta?.resultJson, savedPlan, hasResponse: sourceHasPresentationComment || Boolean(acceptedSummary) }) + : aiRequest ? aiRequest.status === "pending" ? "The selected AI account is unavailable. Fix it in the connection card." : "This run stopped because its AI account was unavailable." : source.status === "cancelled" ? code === "execution_reconciliation_required" ? "The previous execution must be checked before this task can continue. Your message is preserved. View the stopped run for details." - : "Execution was stopped before returning an answer." + : "Execution was stopped." : code === "provider_frame_too_large" ? `Provider output exceeded the safe limit. ${retryDetail}` : code.startsWith("workspace_git_scan_") ? `Workspace setup failed before the agent started. ${retryDetail}` - : `The runner stopped before returning an answer (${code}). ${retryDetail}`; + : `The run failed (${code}). ${retryDetail}`; const id = `${source.id}:failure`; entriesWithFailures.push({ ms: toMs(meta?.finishedAt ?? meta?.startedAt ?? meta?.createdAt), @@ -1663,8 +1671,14 @@ export function TaskChatThread(props: TaskChatThreadProps) { id, kind: "marker", variant: "interrupted", - label: source.status === "cancelled" ? (meta?.startedAt ? "Stopped" : "Couldn't start") : "Run failed", + label: restoreFailed ? "Workspace restore failed" : source.status === "cancelled" ? (meta?.startedAt ? "Stopped" : "Couldn't start") : "Run failed", runId: source.status === "cancelled" ? undefined : source.id, + ...(restoreFailed ? { + retryable: meta?.resultJson?.workspaceRestoreFailure !== "restore_unsafe_archive", + collapsible: true, + runHref: meta?.agentId ? `/agents/${encodeURIComponent(agentMap?.get(meta.agentId)?.urlKey ?? meta.agentId)}/runs/${encodeURIComponent(source.id)}` : undefined, + planHref: savedPlan ? "#document-plan" : undefined, + } : {}), tone: source.status === "cancelled" ? "neutral" : "error", detail, }, @@ -2049,6 +2063,7 @@ export function TaskChatThread(props: TaskChatThreadProps) { steeringAnchorsByRun, legacyTimelineAnchorsByRun, hasBrief, + planDocument, planDocumentSourceRunId, planTurnItem, agentMap, diff --git a/ui/src/components/task-chat/TaskChatMarker.tsx b/ui/src/components/task-chat/TaskChatMarker.tsx index 2163df9887..29ccd9f149 100644 --- a/ui/src/components/task-chat/TaskChatMarker.tsx +++ b/ui/src/components/task-chat/TaskChatMarker.tsx @@ -97,8 +97,13 @@ export function TaskChatMarker({ {item.detail} ) : null} - {item.runHref || onTryAgain ? ( + {item.runHref || item.planHref || onTryAgain ? (
+ {item.planHref ? ( + + ) : null} {item.runHref ? (