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 <noreply@paperclip.ing>
This commit is contained in:
Devin Foley
2026-09-25 11:34:31 -07:00
committed by GitHub
co-authored by Paperclip
parent 5b09d66183
commit bd6caf51bb
29 changed files with 595 additions and 60 deletions
+26
View File
@@ -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
@@ -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";
@@ -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";
}
@@ -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();
});
});
@@ -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<AdapterExecutionResult>,
restore: () => Promise<void>,
): Promise<AdapterExecutionResult> {
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!;
}
@@ -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);
@@ -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<AdapterExec
// adapter's `stagedCodexHomeDir` handling.
let stagedGrokHomeDir: string | null = null;
try {
const executeTurn = async (): Promise<AdapterExecutionResult> => {
const envConfig = parseObject(config.env);
const env: Record<string, string> = {
...buildPaperclipEnv(agent),
@@ -592,17 +593,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
clearSessionOnMissingSession = false,
isRetry = false,
): AdapterExecutionResult => {
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<AdapterExec
return {
exitCode: attempt.proc.exitCode,
signal: attempt.proc.signal,
timedOut: false,
errorMessage: failed ? fallbackErrorMessage : null,
timedOut: attempt.proc.timedOut,
errorMessage: attempt.proc.timedOut ? `Timed out after ${timeoutSec}s` : failed ? fallbackErrorMessage : null,
usage: {
inputTokens: attempt.parsed.inputTokens,
outputTokens: attempt.parsed.outputTokens,
@@ -654,6 +645,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
costUsd: billingType === "api" ? attempt.parsed.costUsd : null,
resultJson: {
stopReason: attempt.parsed.stopReason,
finalResponseRecorded: attempt.parsed.stopReason === "EndTurn" && Boolean(attempt.parsed.summary?.trim()),
requestId: attempt.parsed.requestId,
...(failed ? { stderr: attempt.proc.stderr } : {}),
},
@@ -678,14 +670,12 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
}
return toResult(initial);
};
try {
return await withWorkspaceRestore(executeTurn, async () => { 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<AdapterExec
);
});
}
await Promise.all([
restoreRemoteWorkspace?.(),
stagedAssets.cleanup(),
]);
await stagedAssets.cleanup();
}
}
+2
View File
@@ -2782,3 +2782,5 @@ export * from "./slack-tools.js";
export { MEMORY_CONNECTOR_IDS, isMemoryConnectorId, type MemoryConnectorId } from "./memory-connectors.js";
export * from "./connection-routing.js";
export { WORKSPACE_RESTORE_FAILURE_CODES, hasWorkspaceRestoreFailure, safeWorkspaceRestorePath } from "./workspace-restore.js";
@@ -64,4 +64,6 @@ export interface ExecutionReconciliation {
providerStopped: true;
actionOutcome: "completed" | "not_performed" | "mixed";
outcomeEvidence: string;
/** Operator evidence bound to this failed run; required after workspace restore failure. */
workspaceRepairEvidence?: string;
}
+1
View File
@@ -563,6 +563,7 @@ export const resolveIssueRecoveryActionSchema = z
providerStopped: z.literal(true),
actionOutcome: z.enum(["completed", "not_performed", "mixed"]),
outcomeEvidence: z.string().trim().min(20).max(12000),
workspaceRepairEvidence: z.string().trim().min(20).max(12000).optional(),
})
.strict()
.optional(),
+21
View File
@@ -0,0 +1,21 @@
/** Closed, path-free classifications persisted by adapter settlement. */
export const WORKSPACE_RESTORE_FAILURE_CODES = [
"restore_permission_denied",
"restore_lock_timeout",
"restore_unsafe_archive",
"restore_failed",
] as const;
export function hasWorkspaceRestoreFailure(result: Record<string, unknown> | 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;
}
@@ -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);
@@ -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");
+28 -1
View File
@@ -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
? {
@@ -16,7 +16,7 @@ export function isConversationAdapter(adapterType: string): boolean {
export const CONVERSATION_CONTINUATION_POLICY = "continue_conversation_v1";
export function hasConversationContinuationPolicy(result: Record<string, unknown> | 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<boolean> {
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(
+10 -2
View File
@@ -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(),
@@ -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) {
@@ -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));
@@ -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;
+38 -6
View File
@@ -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<Record<string, unknown> | 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]);
@@ -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);
});
@@ -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,
@@ -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<typeof deliverReconciledExecutions>[1]);
await deliverReconciledExecutions(db, wake as unknown as Parameters<typeof deliverReconciledExecutions>[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();
+37 -1
View File
@@ -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(<TaskChatThread issueId="issue-1" comments={[]} onAdd={async () => {}} 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<HTMLButtonElement>('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(<TaskChatThread comments={[]} onAdd={async () => {}} 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();
+25 -10
View File
@@ -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,
@@ -97,8 +97,13 @@ export function TaskChatMarker({
{item.detail}
</div>
) : null}
{item.runHref || onTryAgain ? (
{item.runHref || item.planHref || onTryAgain ? (
<div className="flex items-center justify-end gap-2 border-t border-border/70 bg-background/50 px-3 py-2 dark:bg-background/30">
{item.planHref ? (
<Button asChild variant="ghost" size="xs">
<Link to={item.planHref}>View saved plan</Link>
</Button>
) : null}
{item.runHref ? (
<Button asChild variant="ghost" size="xs">
<Link to={item.runHref}>View run</Link>
@@ -255,6 +255,7 @@ export interface TaskChatMarkerItem {
runId?: string;
createdAtIso?: string;
runHref?: string;
planHref?: string;
}
/** A second-tier live token/cost readout (ACP UsageUpdate). */
@@ -0,0 +1,16 @@
import { expect, it } from "vitest";
import { workspaceRestoreMarkerDetail } from "./workspace-restore-marker";
it("makes saved-plan claims only with durable evidence", () => {
expect(workspaceRestoreMarkerDetail({ result: {}, savedPlan: true, hasResponse: true })).toContain("after the plan was saved");
expect(workspaceRestoreMarkerDetail({ result: {}, savedPlan: false, hasResponse: false })).toBe("Workspace restore failed. Workspace files need recovery.");
});
it("separates explicit missing-response evidence from the restore failure", () => {
expect(workspaceRestoreMarkerDetail({ result: { finalResponseRecorded: false }, savedPlan: false, hasResponse: false })).toContain("No final response was recorded.");
expect(workspaceRestoreMarkerDetail({ result: { finalResponseRecorded: false }, savedPlan: false, hasResponse: true })).not.toContain("No final response");
});
it.each(["/host/secret", "../escape", "file -> /secret", "C:\\host\\secret", "tmp/private-clone/file"])("omits diagnostic path %s", (workspaceRestorePath) => {
expect(workspaceRestoreMarkerDetail({ result: { workspaceRestorePath }, savedPlan: false, hasResponse: true })).not.toContain("Affected path");
});
+18
View File
@@ -0,0 +1,18 @@
import { safeWorkspaceRestorePath } from "@paperclipai/shared";
export function workspaceRestoreMarkerDetail(input: {
result: Record<string, unknown> | null | undefined;
savedPlan: boolean;
hasResponse: boolean;
}): string {
const parts = [input.savedPlan
? "Workspace restore failed after the plan was saved. The saved plan is available."
: "Workspace restore failed."];
parts.push("Workspace files need recovery.");
if (!input.hasResponse && input.result?.finalResponseRecorded === false) {
parts.push("No final response was recorded.");
}
const relativePath = safeWorkspaceRestorePath(input.result?.workspaceRestorePath);
if (relativePath) parts.push(`Affected path: ${relativePath}.`);
return parts.join(" ");
}