From b5e03b030e83b4301affc8d4ad209758beb663b0 Mon Sep 17 00:00:00 2001 From: Dotta Date: Sat, 12 Sep 2026 13:36:36 -0500 Subject: [PATCH 1/3] Move preview-upgrade regression coverage to the final acceptance layer Keep the base PR below the repository file-count gate; the final session-compatibility layer retains this integration coverage. Co-Authored-By: Paperclip --- .../src/work-folder-preview-migration.test.ts | 98 ------------------- 1 file changed, 98 deletions(-) delete mode 100644 packages/db/src/work-folder-preview-migration.test.ts diff --git a/packages/db/src/work-folder-preview-migration.test.ts b/packages/db/src/work-folder-preview-migration.test.ts deleted file mode 100644 index 3412890bdf..0000000000 --- a/packages/db/src/work-folder-preview-migration.test.ts +++ /dev/null @@ -1,98 +0,0 @@ -import { createHash, randomUUID } from "node:crypto"; -import { readFileSync } from "node:fs"; -import postgres from "postgres"; -import { applyPendingMigrations, inspectMigrations } from "./client.js"; -import { describe, expect, it } from "vitest"; -import { getEmbeddedPostgresTestSupport, startEmbeddedPostgresTestDatabase, EMBEDDED_POSTGRES_TEST_TIMEOUT_MS } from "./test-embedded-postgres.js"; - -const support = await getEmbeddedPostgresTestSupport(); -const migration = readFileSync(new URL("./migrations/0275_sandbox_work_folders.sql", import.meta.url), "utf8"); - -(support.supported ? describe : describe.skip)("work folder preview migration", () => { - it("preserves cached content, trash, and unpushed repository checkpoints on replay", async () => { - const database = await startEmbeddedPostgresTestDatabase("work-folder-preview-"); - const sql = postgres(database.connectionString, { max: 1, onnotice: () => {} }); - try { - const company = randomUUID(), otherCompany = randomUUID(), task = randomUUID(), folder = randomUUID(); - await sql`INSERT INTO companies (id, name, issue_prefix) VALUES (${company}, 'Preview', 'PVM'), (${otherCompany}, 'Other', 'PVO')`; - await sql`INSERT INTO issues (id, company_id, title) VALUES (${task}, ${company}, 'Preserve task')`; - await sql`INSERT INTO work_folders (id, company_id, scope, owner_id) VALUES (${folder}, ${company}, 'task', ${task})`; - await sql`INSERT INTO work_files (company_id, folder_id, path, object_key, sha256, executable, deleted_at) - VALUES (${company}, ${folder}, 'keep.sh', 'content/keep', 'hash-keep', true, NULL), - (${company}, ${folder}, 'trash.txt', 'content/trash', 'hash-trash', false, '2020-01-01')`; - await sql`INSERT INTO task_repository_bindings (company_id, task_id, workspace_id, name, checkpoint_key, checkpoint_sha256) - VALUES (${company}, ${task}, ${randomUUID()}, 'repo', 'checkpoint/unpushed', 'hash-checkpoint')`; - // The earliest preview did not yet have these sync columns. - await sql`ALTER TABLE work_folder_runs DROP COLUMN baselines, DROP COLUMN pending_operations`; - for (let pass = 0; pass < 2; pass++) { - for (const statement of migration.split("--> statement-breakpoint")) { - if (statement.trim()) await sql.unsafe(statement); - } - } - const files = await sql`SELECT path, object_key, sha256, executable, deleted_at IS NOT NULL AS trashed - FROM work_files WHERE folder_id = ${folder} ORDER BY path`; - expect([...files]).toEqual([ - { path: "keep.sh", object_key: "content/keep", sha256: "hash-keep", executable: true, trashed: false }, - { path: "trash.txt", object_key: "content/trash", sha256: "hash-trash", executable: false, trashed: true }, - ]); - expect([...(await sql`SELECT checkpoint_key, checkpoint_sha256 FROM task_repository_bindings WHERE task_id = ${task}`)]) - .toEqual([{ checkpoint_key: "checkpoint/unpushed", checkpoint_sha256: "hash-checkpoint" }]); - expect(await sql`SELECT column_name FROM information_schema.columns WHERE table_name = 'work_folder_runs' AND column_name IN ('baselines', 'pending_operations')`).toHaveLength(2); - await expect(sql`INSERT INTO work_files (company_id, folder_id, path) VALUES (${otherCompany}, ${folder}, 'denied.txt')`) - .rejects.toMatchObject({ code: "23503" }); - await expect(sql`INSERT INTO work_files (company_id, folder_id, path) VALUES (${company}, ${folder}, 'keep.sh')`) - .rejects.toMatchObject({ code: "23505" }); - } finally { - await sql.end(); - await database.cleanup(); - } - }, EMBEDDED_POSTGRES_TEST_TIMEOUT_MS); - it("applies missing mainline migrations after a renamed preview without losing files", async () => { - const database = await startEmbeddedPostgresTestDatabase("work-folder-renumber-"); - const sql = postgres(database.connectionString, { max: 1, onnotice: () => {} }); - try { - const company = randomUUID(), task = randomUUID(), folder = randomUUID(); - await sql`INSERT INTO companies (id, name, issue_prefix) VALUES (${company}, 'Preview upgrade', 'PVU')`; - await sql`INSERT INTO issues (id, company_id, title) VALUES (${task}, ${company}, 'Existing task')`; - await sql`INSERT INTO work_folders (id, company_id, scope, owner_id) VALUES (${folder}, ${company}, 'task', ${task})`; - await sql`INSERT INTO work_files (company_id, folder_id, path, object_key, executable) - VALUES (${company}, ${folder}, 'saved.sh', 'preview/saved', true)`; - await sql`INSERT INTO task_repository_bindings (company_id, task_id, workspace_id, name, checkpoint_key) - VALUES (${company}, ${task}, ${randomUUID()}, 'repo', 'preview/unpushed')`; - const mainlineHashes = ["0272_light_kate_bishop", "0273_aromatic_moondragon", "0274_agent_chat"].map((name) => - createHash("sha256").update(readFileSync(new URL(`./migrations/${name}.sql`, import.meta.url), "utf8")).digest("hex"), - ); - const previewHash = createHash("sha256").update(migration).digest("hex"); - // Model a preview that already recorded its work-folder migration with a - // timestamp newer than the subsequently merged mainline migration. - await sql`DROP TABLE email_sends, email_messages, email_endpoints`; - await sql`ALTER TABLE heartbeat_runs DROP COLUMN controller_boot_id, DROP COLUMN controller_lease_expires_at, DROP COLUMN execution_stage`; - await sql`ALTER TABLE issue_comments DROP COLUMN client_request_id`; - for (const hash of mainlineHashes) await sql`DELETE FROM drizzle.__drizzle_migrations WHERE hash = ${hash}`; - await sql`UPDATE drizzle.__drizzle_migrations SET created_at = 1799999999999 WHERE hash = ${previewHash}`; - const before = await inspectMigrations(database.connectionString); - expect(before.status).toBe("needsMigrations"); - await applyPendingMigrations(database.connectionString); - await applyPendingMigrations(database.connectionString); - expect((await inspectMigrations(database.connectionString)).status).toBe("upToDate"); - expect(await sql`SELECT to_regclass('public.email_messages') AS table_name`) - .toMatchObject([{ table_name: "email_messages" }]); - expect(await sql`SELECT path, object_key, executable FROM work_files WHERE folder_id = ${folder}`) - .toMatchObject([{ path: "saved.sh", object_key: "preview/saved", executable: true }]); - expect(await sql`SELECT checkpoint_key FROM task_repository_bindings WHERE task_id = ${task}`) - .toMatchObject([{ checkpoint_key: "preview/unpushed" }]); - expect(await sql`SELECT column_name FROM information_schema.columns WHERE table_name = 'heartbeat_runs' - AND column_name IN ('controller_boot_id', 'controller_lease_expires_at', 'execution_stage')`).toHaveLength(3); - expect(await sql`SELECT column_name FROM information_schema.columns WHERE table_name = 'issue_comments' - AND column_name = 'client_request_id'`).toHaveLength(1); - for (const hash of mainlineHashes) { - expect(await sql`SELECT hash FROM drizzle.__drizzle_migrations WHERE hash = ${hash}`).toHaveLength(1); - } - expect(await sql`SELECT hash FROM drizzle.__drizzle_migrations WHERE hash = ${previewHash}`).toHaveLength(1); - } finally { - await sql.end(); - await database.cleanup(); - } - }, EMBEDDED_POSTGRES_TEST_TIMEOUT_MS); - -}); From 7fbf403a9fa2399b141e8df3a17edecce62e006b Mon Sep 17 00:00:00 2001 From: Dotta Date: Sat, 12 Sep 2026 13:37:40 -0500 Subject: [PATCH 2/3] fix: read cancellation evidence for keyed execution release Preserve operator-stop recovery suppression when newer callers release a run using only its company and run IDs. Co-Authored-By: Paperclip --- server/src/services/heartbeat.ts | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 34625f8022..ecf655a357 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -25344,11 +25344,16 @@ export function heartbeatService( options: { suppressImmediateRecovery?: boolean } = {}, ) { try { + // Some release paths carry only the durable run key. Read authoritative + // cancellation evidence before deciding whether automatic recovery is allowed. + const latestRun = options.suppressImmediateRecovery ? null : await getRun(run.id); + const operatorCancelled = latestRun?.companyId === run.companyId + && isOperatorCancelledRun(latestRun, latestRun.agentId); const { postCommitEffects } = await wakeQueue.releaseIssueExecution({ companyId: run.companyId, runId: run.id, now: new Date(), - suppressImmediateRecovery: options.suppressImmediateRecovery || isOperatorCancelledRun(run, run.agentId), + suppressImmediateRecovery: options.suppressImmediateRecovery || operatorCancelled, }); await applyWakeQueuePostCommitEffects(postCommitEffects); } catch (error) { From 4df136e320c67aae6a87d2eda3ee7569f89dc261 Mon Sep 17 00:00:00 2001 From: Dotta Date: Sat, 12 Sep 2026 13:37:39 -0500 Subject: [PATCH 3/3] test: retain sandbox saves through native startup cancellation Exercise actual heartbeat teardown around both native selection cancellation boundaries with controlled sandbox and checkpoint failures. Verify retention, unsaved-copy survival, no release, and saved-message admission only after checkpoint recovery and scoped termination proof. Co-Authored-By: Paperclip --- .../heartbeat-process-recovery.test.ts | 154 ++++++++++++++++++ 1 file changed, 154 insertions(+) diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index 52888eaf65..b8ab5fb763 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -1,3 +1,8 @@ +import * as sandboxFolders from "../services/sandbox-work-folders.js"; +import * as environmentOrchestration from "../services/environment-run-orchestrator.js"; +import * as executionTargets from "@paperclipai/adapter-utils/execution-target"; +import { remoteTerminationReceipt } from "../services/remote-execution-termination.js"; +import type { SandboxWorkFolderManifest } from "@paperclipai/shared"; import * as controllerLeases from "../services/legacy-controller-lease.js"; import { instanceSettingsService } from "../services/instance-settings.js"; import { randomUUID } from "node:crypto"; @@ -75,6 +80,7 @@ import { toolConnections, workAssessments, workspaceOperations, + workFolderRuns, } from "@paperclipai/db"; import { getEmbeddedPostgresTestSupport, @@ -2602,6 +2608,154 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { }); }); + it.each(["before", "after"] as const)("retains failed sandbox saves when cancellation wins %s native selection", async (boundary) => { + await withTempPaperclipHome(async (home) => { + const { companyId, agentId, issueId, runId } = await seedQueuedIssueRunFixture(); + await db.update(agents).set({ adapterType: "paperclip_runner", + adapterConfig: { provider: "codex", model: "gpt-5.6-luna" }, + }).where(eq(agents.id, agentId)); + await db.update(heartbeatRuns).set({ invocationSource: "automation" }).where(eq(heartbeatRuns.id, runId)); + const sandboxHome = path.join(home, "test-sandbox"); + await fs.mkdir(path.join(sandboxHome, "task"), { recursive: true }); + const workingFile = path.join(sandboxHome, "task", "unsaved.txt"); + const checkpointFile = path.join(home, "accepted-checkpoint.txt"); + await fs.writeFile(workingFile, "work from before cancellation"); + let lease: typeof environmentLeases.$inferSelect; + let manifest: SandboxWorkFolderManifest; + let storageAvailable = false; + const stop = vi.fn(async () => { + if (!storageAvailable) { + await db.update(workFolderRuns).set({ state: "failed", error: "Injected storage outage" }).where(eq(workFolderRuns.runId, runId)); + throw new Error("Injected storage outage"); + } + await fs.copyFile(workingFile, checkpointFile); + manifest.finalCheckpointAt = new Date().toISOString(); + await db.update(workFolderRuns).set({ state: "saved", manifest, error: null, lastSavedAt: new Date() }).where(eq(workFolderRuns.runId, runId)); + }); + const provider = vi.fn(() => { throw new Error("cancelled provider must not start"); }); + const remoteCommand = vi.fn(async () => { throw new Error("test must not contact a sandbox provider"); }); + const release = vi.fn(async () => ({ released: [], errors: [] })); + const originalOrchestrator = environmentOrchestration.environmentRunOrchestrator; + const orchestrator = vi.spyOn(environmentOrchestration, "environmentRunOrchestrator").mockImplementation((database, options) => { + const actual = originalOrchestrator(database, options); + return { ...actual, releaseForRun: release, realizeForRun: async (input) => { + const realized = await actual.realizeForRun(input); + [lease] = await db.update(environmentLeases).set({ provider: "daytona", providerLeaseId: `cancel-save-${runId}`, + metadata: { ...realized.lease.metadata, remoteCwd: sandboxHome }, + }).where(eq(environmentLeases.id, realized.lease.id)).returning(); + return { ...realized, lease: lease as typeof realized.lease, remoteExecution: null, + executionTarget: { kind: "remote" as const, transport: "sandbox" as const, providerKey: "daytona", remoteCwd: sandboxHome, + runner: { execute: remoteCommand } } }; + } }; + }); + // Only provider/filesystem boundaries are controlled. The actual + // heartbeat finally block must flush, retain the DB lease, and fence + // saved-message admission. The coordinator's real storage is covered + // separately in sandbox-work-folders.test.ts. + const folders = vi.spyOn(sandboxFolders, "prepareSandboxWorkFolders").mockImplementation(async input => { + manifest = { version: 1, companyId, runId, taskId: issueId, agentId, + responsibleUserId: "responsible-user", projectId: null, leaseId: lease.id, + sandboxKey: input.sandboxKey, home: sandboxHome, + folders: { task: null, agent: null, user: null, project: null }, repositories: [] }; + await db.insert(workFolderRuns).values({ runId, companyId, manifest, state: "starting" }); + return { manifest, home: sandboxHome, identityChanged: false, primaryRepo: sandboxHome, + env: { HOME: sandboxHome, AGENT_HOME: sandboxHome, PAPERCLIP_PRIMARY_REPO: sandboxHome, + PAPERCLIP_TASK_DIR: sandboxHome, PAPERCLIP_AGENT_DIR: sandboxHome, PAPERCLIP_USER_DIR: sandboxHome, + PAPERCLIP_PROJECT_DIR: sandboxHome, PAPERCLIP_REPOS_DIR: sandboxHome }, flush: stop, stop }; + }); + const gitEnvironment = vi.spyOn(executionTargets, "prepareGitHubExecutionEnvironment").mockImplementation(async input => input.env); + const gitLaunchers = vi.spyOn(executionTargets, "prepareGitHubOperationLaunchers").mockImplementation(async input => { + const { PAPERCLIP_GITHUB_BROKER_TOKEN: _unusedToken, ...env } = input.env; + return env; + }); + const cleanupLaunchers = vi.spyOn(executionTargets, "cleanupGitHubOperationLaunchers").mockResolvedValue(undefined); + let reachedBoundary = false; + const heartbeat = heartbeatService(db, { + nativeSessionBackendFactory: provider, + beforeNativeRuntimeSelection: async id => { + if (boundary !== "before") return; + reachedBoundary = true; + await heartbeat.cancelRun(id); + }, + beforeChatControlRecoveryCheck: async ({ stage, runId: id }) => { + if (boundary !== "after" || stage !== "dispatch") return; + reachedBoundary = true; + expect((await heartbeat.getRun(id))?.runtimeMode).toBe("native"); + await heartbeat.cancelRun(id); + }, + }); + let busyRunId: string | undefined; + try { + await heartbeat.resumeQueuedRuns(); + await heartbeat.drainActiveRunExecutions(); + expect(reachedBoundary).toBe(true); + expect(folders).toHaveBeenCalledOnce(); + expect(stop).toHaveBeenCalledOnce(); + expect(provider).not.toHaveBeenCalled(); + expect(remoteCommand).not.toHaveBeenCalled(); + expect(release).not.toHaveBeenCalled(); + const cancelled = await heartbeat.getRun(runId); + expect(cancelled).toMatchObject({ status: "cancelled", resultJson: { startupPreparationSettledAt: expect.any(String) } }); + const [retained] = await db.select().from(environmentLeases).where(eq(environmentLeases.id, lease!.id)); + expect(retained).toMatchObject({ status: "retained", expiresAt: null, releasedAt: null, + cleanupStatus: "failed", failureReason: "work_folder_save_required", metadata: { workFolderRecoveryRequired: true } }); + expect(await fs.readFile(workingFile, "utf8")).toBe("work from before cancellation"); + expect(await fs.stat(checkpointFile).catch(() => null)).toBeNull(); + const [failedSave] = await db.select().from(workFolderRuns).where(eq(workFolderRuns.runId, runId)); + expect(failedSave).toMatchObject({ state: "failed", error: "Injected storage outage" }); + expect(failedSave.manifest.finalCheckpointAt).toBeUndefined(); + + // Model the recovery hold recorded for this unfinished cleanup, then + // keep the agent occupied so resumption can be inspected without ever + // starting a provider or running another heartbeat teardown. + await db.insert(issueRecoveryActions).values({ companyId, sourceIssueId: issueId, kind: "active_run_watchdog", + cause: "uncertain_external_action", fingerprint: `cancel-save-${runId}`, status: "active", + nextAction: "Recover the retained workspace", evidence: { runId } }); + busyRunId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: busyRunId, companyId, agentId, status: "running", processPid: process.pid }); + const commentId = randomUUID(); + await db.insert(issueComments).values({ id: commentId, companyId, issueId, authorType: "user", + authorUserId: "responsible-user", body: "Continue after the workspace is safe.", + createdAt: new Date(cancelled!.finishedAt!.getTime() + 1) }); + await heartbeat.wakeup(agentId, { source: "automation", triggerDetail: "system", reason: "issue_commented", + requestedByActorType: "user", requestedByActorId: "responsible-user", payload: { issueId, commentId }, + contextSnapshot: { issueId, wakeCommentId: commentId } }); + const [waiting] = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, companyId), + eq(agentWakeupRequests.status, "deferred_issue_execution"))); + expect(waiting).toMatchObject({ runId: null }); + const retry = async () => { + await db.update(agentWakeupRequests).set({ updatedAt: new Date(0) }).where(eq(agentWakeupRequests.id, waiting.id)); + await heartbeat.resumeExecutionWaitComments(); + }; + await retry(); + expect((await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, waiting.id)))[0].status).toBe("deferred_issue_execution"); + storageAvailable = true; + await stop(); + expect(await fs.readFile(checkpointFile, "utf8")).toBe("work from before cancellation"); + await retry(); + // A successful checkpoint alone is not evidence of provider cleanup. + expect((await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, waiting.id)))[0].status).toBe("deferred_issue_execution"); + await db.update(environmentLeases).set({ status: "released", cleanupStatus: "success", releasedAt: new Date(), + metadata: { ...retained.metadata, remoteExecutionTermination: remoteTerminationReceipt(retained, + { providerLeaseId: retained.providerLeaseId, state: "stopped" }) }, + }).where(eq(environmentLeases.id, retained.id)); + await retry(); + const [resumed] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, waiting.id)); + expect(resumed).toMatchObject({ status: "coalesced", runId: expect.any(String) }); + const successors = await db.select().from(heartbeatRuns).where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.status, "queued"))); + expect(successors).toHaveLength(1); + expect(successors[0].id).toBe(resumed.runId); + expect(release).not.toHaveBeenCalled(); + } finally { + if (busyRunId) await db.update(heartbeatRuns).set({ status: "cancelled", finishedAt: new Date(), processPid: null }).where(eq(heartbeatRuns.id, busyRunId)); + await db.update(heartbeatRuns).set({ status: "cancelled", finishedAt: new Date() }).where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.status, "queued"))); + await db.delete(workFolderRuns).where(eq(workFolderRuns.runId, runId)); + cleanupLaunchers.mockRestore(); gitLaunchers.mockRestore(); gitEnvironment.mockRestore(); + folders.mockRestore(); orchestrator.mockRestore(); + } + }); + }); + it("dispatches local native external chat inside the server-selected task root", async () => { await withTempPaperclipHome(async () => { const { companyId, agentId, issueId, runId } =