Files
PaperClipAI/server/src/__tests__/native-session-resumption.test.ts
DottaandPaperclip 96bf004a79 fix: use persisted state for lifecycle continuation and retry budgets (#13888)
## Thinking Path

> - Paperclip is the open source app people use to manage AI agents for
work.
> - Its control plane decides when a task can continue, wait, stop, or
complete.
> - Legacy continuation could change when an agent changed its wording
without changing task state.
> - Shared attempt counts also let repair and infrastructure retries
affect each other's limits.
> - This pull request uses persisted state and separate, bounded
allowances for these decisions.
> - If automatic repair stops, the task explains what happened and
offers a guarded retry.
> - Paired tests and real-provider evaluations verify that Stop,
approvals, ownership, and spending limits remain authoritative.

## Linked Issues or Issue Description

Related work: Refs #13761, Refs #11126, Refs #13610. These cover
obsolete continuation dispatch and retry storms. Open and closed issues
and PRs were searched for related lifecycle, continuation, and retry
work.

**What happened?**
Legacy continuation depended on English wording and progress heuristics.
Repair, failure retry, and productive continuation could consume shared
counts. When bounded repair stopped, the task showed a technical
recovery message without a clear next action.

**Expected behavior**
Persisted disposition and owned execution paths determine the next
action. Missing disposition prompts bounded agent repair. Explicit work
mode determines planning mode. Narrative changes and raw activity counts
cannot replenish allowances. An exhausted repair shows a readable
notice. An explicit retry checks current controls and preserves the
assigned agent.

**Steps to reproduce**
Run `pnpm test:lifecycle-baseline`. The paired probes keep structured
state constant while varying completion, planning, blocker, and progress
prose. Run the explicit `lifecycle-baseline` and
`continuation-accounting` Product E2E suites for real-provider coverage.
In Storybook, open **Design previews / Recovery notice** to inspect the
production component's normal, pending, acknowledged, unavailable,
failure, and mobile states.

## What Changed

- Hide the image attachment button, icon, and drop/paste hint in answer
composers. Image paste and drop support remains available.
- Merge current master and retain both browser regression sets. Use a
production-stamped service worker in the offline recovery browser
fixture.
- Share one state-based legacy continuation decision across immediate,
delayed, and recovered dispatch. Bind bounded repairs to their source
run and episode.
- Remove title and description wording from work-mode authority. Agents
can still write requested plans in execution mode.
- Persist separate failure-retry and productive-continuation counters.
Disposition repair and resource waits cannot consume or reset those
allowances.
- Validate delayed repair identity, then recheck current gates before
provider dispatch. Fence native startup cancellation.
- Show **Agent needs attention**, a plain-language explanation, **Retry
agent**, and expandable details in both task interfaces. Report request
progress, acknowledgement, and errors inline.
- Store typed recovery notice metadata. Recognize older active notices
only through exact stored action and run IDs. Notice text never grants
retry authority.
- Use the existing recovery-action endpoint for retry. Recheck current
action, status, owner, agent availability, dependencies, active runs,
pending questions and confirmations, approvals, pause controls, and
budget. Duplicate requests do not wake twice.
- Add component, page, route, database, contract, and Storybook
coverage. Keep the scenario inventory and executable evals here.
Historical reports and snapshots live in the [commit-pinned
paperclip-evals
archive](https://github.com/paperclipai/paperclip-evals/blob/ce3e5afcd4a1184650f586a2b5b8be5874c66c8b/experiments/2026-09-lifecycle-authority/README.md).
- Preserve unsaved project fields while the same project URL changes to
its canonical alias. Do not reuse data across projects or companies.
This separate fix addresses the repeated repository-editor browser
failure without changing the browser test.
- Keep the development service worker from intercepting Vite module
reloads. Update the connection-intent browser fixture to record progress
and completion through the agent API.

## Verification

Merge preparation on September 25, commit
`c1e8e4b7ddd9fbc4913ed55ce21b8e12906c2f97`:

- Merged master `bd2030932` and resolved the browser test-list conflict
by keeping both sets of regressions.
- Deterministic lifecycle baseline: 1,090/1,090 assertions passed; no
failures, skips, or missing selected evidence. Unit 423, runner 184,
database integration 397, grading 86.
- Browser support: 17/17 passed. The offline recovery test first failed
with an unstamped development worker, then passed with the production
stamp. Its assertions are unchanged.
- Focused interaction UI and offline fallback tests: 19/19 passed.
Verified the custom-answer composer in Storybook: no attachment controls
or hint; entering an answer enables Next.
- Recursive typecheck, production build, token gates, and diff checks
passed. The worktree is clean. No new real-provider campaign was run.
- Current CI and review: [Current PR CI
passed](https://github.com/paperclipai/paperclip/actions/runs/36166011243):
55 successful checks and two optional Storybook skips. Greptile scored
this exact commit 5/5. Hiding the question attachment controls is an
intentional UI change; paste/drop remains available.

Earlier recovery UI verification, commit
`21be0fec0e90e86b6d662b8ee4831847cd041cdb`:

- Recursive typecheck, production build, token gates, and diff checks
passed.
- Focused UI coverage: 338 tests passed across six suites (336 before
the interaction guard, with the two affected suites rerun at 149 passed
after it). Covers both task interfaces, the real page mutation,
pending/error acknowledgement, stale state, and unavailable controls.
- Recovery database integration: 352 tests passed before the interaction
guard. The complete recovery-action and mutation-route suites passed 181
tests after it. The two new pending question/confirmation regressions
failed before the fix and passed afterward, including
resolved-interaction controls. Shared validator suite: 31 passed. E2E
catalog suites: 34 passed.
- Browser inspection passed for light/dark themes, mobile layout,
expandable details, pending retry, acknowledgement, failure, and
disabled retry. Storybook renders the production component; its request
is simulated.
- The broad local run hit two chat callback-order wait failures and was
stopped after all CI unit/database/runner shards passed. Both local
failures passed when rerun without the competing full-suite process.
- CI exposed a repeated project-repository draft-loss race during
canonical redirects. A new unit regression failed before the fix; all
nine project-page tests now pass, including controls for other projects
and companies. Both unchanged repository browser tests passed against a
fresh local server. UI typecheck, production UI build, and token gates
passed after this fix.
- [Earlier PR CI
passed](https://github.com/paperclipai/paperclip/actions/runs/36072486798)
on `21be0fec0e90e86b6d662b8ee4831847cd041cdb`: 55 successful checks, two
optional Storybook skips, and no failed or pending checks. The
repository browser shard passed with the production fix. Greptile is 5/5
on this exact commit with no unresolved review threads. The PR is
mergeable.

Historical, source-qualified lifecycle evidence:

- Lifecycle baseline: 1,074 assertions. Native session coverage: 447
tests. Product E2E support: 515 tests. Browser support: 11 tests. Full
earlier verification is retained in the archive.
- [Real-provider campaign: 8/8 passed, zero
retries](https://d1p6rlowie26tp.cloudfront.net/runner-e2e/campaigns/gha-35881382080-1/index.html),
source `e88d210417280140b44a36449027290adcb1aeaa`. Evidence and cleanup
checks passed. This includes deliberately exhausted repair cases that
correctly remain blocked; it does not mean every task finished Done.
This campaign predates the recovery UI change.
- Archive migration verified all 16 original JSON files byte-for-byte
and all 24 checksum entries. App tests do not need private archive
access. [Archive PR
#27](https://github.com/paperclipai/paperclip-evals/pull/27) is merged.

## Risks

- Agents that omit durable disposition receive at most two repair
attempts by default. Prose-only completion exposes missing state rather
than silently changing scheduling.
- A retry is an explicit board action. The server rechecks current
controls. A successful response confirms the task returned to To do; it
does not claim that the provider has already started.
- Existing notice metadata remains valid. Only older active notices with
matching structured evidence receive the new UI. Historical notices
without that evidence keep their existing rendering. No schema migration
is required.
- Old run records require conservative retry accounting. Tests cover old
counters, alternating retry lanes, restarts, and exhausted repairs.
- Historical snapshots require private `paperclip-evals` access. The app
index retains public campaign links. Live campaigns qualify specific
sources and scenarios; no new real-provider campaign has run for the
recovery UI commit.

> This fixes existing lifecycle and recovery behavior and does not
duplicate planned core work.

## Model Used

OpenAI GPT-6 through Codex assisted implementation, reasoning, code
execution, and review. The exact serving model ID and context window are
not exposed in this task. Historical real-provider evaluations used
Codex model `gpt-5.6-sol`, separately from the implementation assistant.

## Checklist

- [x] I have included a thinking path that traces from project context
to this change
- [x] I have specified the model used (with version and capability
details)
- [x] I have checked ROADMAP.md and confirmed this PR does not duplicate
planned core work
- [x] I have searched GitHub for duplicate or related PRs and linked
them above
- [x] I have either (a) linked existing issues with `Fixes: #` / `Closes
#` / `Refs #` OR (b) described the issue in-PR following the relevant
issue template
- [x] I have not referenced internal/instance-local Paperclip issues or
links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip`
URLs)
- [x] My branch name describes the change and contains no internal
Paperclip ticket id or instance-derived details
- [x] I have run tests locally and they pass
- [x] I have added or updated tests where applicable
- [x] I have updated relevant documentation to reflect my changes
- [x] I have considered and documented any risks above
- [x] All Paperclip CI gates are green
- [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge

---------

Co-authored-by: Paperclip <noreply@paperclip.ing>
2026-09-25 15:28:11 -07:00

1455 lines
49 KiB
TypeScript

import { spawn } from "node:child_process";
import { randomUUID } from "node:crypto";
import { once } from "node:events";
import { fileURLToPath } from "node:url";
import { afterAll, beforeAll, describe, expect, it, vi } from "vitest";
import { eq, inArray } from "drizzle-orm";
import {
activityLog,
agentWakeupRequests,
agents,
companies,
completionContracts,
createDb,
executionWorkspaces,
heartbeatRunEvents,
heartbeatRuns,
issueRecoveryActions,
issueComments,
issueWorkProducts,
issues,
nativeRunFinalizations,
nativeRunResults,
projectWorkspaces,
projects,
statusDecisionEffects,
statusDecisions,
workAssessments,
} from "@paperclipai/db";
import {
type NativeExecutionInputV1,
type NativeExecutionInput,
type NativeSession,
type NativeSessionBackend,
type PersistedNativeSession,
type PrpEvent,
} from "@paperclipai/paperclip-runner";
import {
CONTROL_PLANE_CONFORMANCE_RESULT,
CONTROL_PLANE_CONFORMANCE_TERMINAL,
} from "../vendor/paperclip-runner/testing.js";
import { startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js";
import { drainHeartbeatRunsToQuiescence } from "./helpers/drain-heartbeat-runs.js";
import { waitForPendingRunFailureReports } from "../services/run-failure-report.js";
import {
claimNativeSessionResumptions,
dispatchNativeSessionResumptions,
} from "../services/native-runtime/native-finalization-reconciler.js";
const legacyAdapterExecute = vi.hoisted(() => vi.fn(async () => ({
exitCode: 0,
signal: null,
timedOut: false,
summary: "Fresh flag-off run completed through legacy.",
resultJson: { summary: "fresh legacy after persisted native recovery" },
provider: "test",
model: "legacy-test",
})));
vi.mock("../adapters/index.js", async () => {
const actual = await vi.importActual<typeof import("../adapters/index.js")>("../adapters/index.js");
return {
...actual,
getServerAdapter: vi.fn(() => ({
type: "codex_local",
execute: legacyAdapterExecute,
supportsLocalAgentJwt: false,
})),
};
});
const mockCaptureRunFailure = vi.hoisted(() => vi.fn());
vi.mock("../sentry.js", async () => {
const actual = await vi.importActual<typeof import("../sentry.js")>("../sentry.js");
return { ...actual, captureRunFailure: mockCaptureRunFailure };
});
import { heartbeatService } from "../services/heartbeat.js";
import { instanceSettingsService } from "../services/instance-settings.js";
describe("P6-25 pre-result native session recovery", () => {
let temporary: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>> | null = null;
let db: ReturnType<typeof createDb>;
const companyId = "79000000-0000-4000-8000-000000000001";
const agentId = "79000000-0000-4000-8000-000000000002";
const issueId = "79000000-0000-4000-8000-000000000003";
const runId = "79000000-0000-4000-8000-000000000004";
const cancelledRunId = "79000000-0000-4000-8000-000000000005";
const exhaustedRunId = "79000000-0000-4000-8000-000000000006";
const missingCheckpointRunId = "79000000-0000-4000-8000-000000000007";
const initialRunId = "79000000-0000-4000-8000-000000000008";
const bootstrapRetryRunId = "79000000-0000-4000-8000-000000000009";
const observedExpiredRunId = "79000000-0000-4000-8000-000000000010";
const observedLivePidRunId = "79000000-0000-4000-8000-000000000011";
const persistedProfile = {
mode: "native",
nativeExecutionInput: { schema: "paperclip.native-execution-input.v1", binding: { runId } },
sessionCheckpoint: {
backendKind: "codex_app_server",
sessionId: "persisted-session",
identity: { runId },
providerSessionId: "provider-session",
activeTurnId: "active-turn",
},
};
beforeAll(async () => {
temporary = await startEmbeddedPostgresTestDatabase("paperclip-native-resume-");
db = createDb(temporary.connectionString);
await db.insert(companies).values({ id: companyId, name: "Native resume", issuePrefix: "NRR" });
await db.insert(agents).values({
id: agentId,
companyId,
name: "Native resume agent",
adapterType: "codex_local",
status: "running",
});
await db.insert(issues).values({
id: issueId,
companyId,
title: "Resume the same native run",
status: "in_progress",
assigneeAgentId: agentId,
workMode: "standard",
});
await db.insert(heartbeatRuns).values([
{
id: runId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: persistedProfile,
contextSnapshot: { issueId },
},
{
id: cancelledRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "cancelled",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: persistedProfile,
contextSnapshot: { issueId },
},
{
id: exhaustedRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "failed",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: persistedProfile,
contextSnapshot: { issueId },
},
{
id: missingCheckpointRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: missingCheckpointRunId },
},
},
contextSnapshot: { issueId },
},
{
id: initialRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: initialRunId },
},
},
contextSnapshot: { issueId },
},
{
id: bootstrapRetryRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "failed",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: bootstrapRetryRunId },
},
},
errorCode: "provider_initialize_timeout",
contextSnapshot: { issueId },
},
...[observedExpiredRunId, observedLivePidRunId].map((observedRunId) => ({
id: observedRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running" as const,
runtimeMode: "native" as const,
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
...persistedProfile,
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: observedRunId },
},
sessionCheckpoint: {
...persistedProfile.sessionCheckpoint,
identity: { runId: observedRunId },
},
},
contextSnapshot: { issueId },
})),
]);
await db.insert(nativeRunFinalizations).values([
{ runId, companyId, issueId, phase: "retryable_failure", attempt: 1, nextAttemptAt: new Date(0) },
{ runId: cancelledRunId, companyId, issueId, phase: "retryable_failure", attempt: 1, nextAttemptAt: new Date(0) },
{ runId: exhaustedRunId, companyId, issueId, phase: "terminal_failure", attempt: 3 },
{ runId: missingCheckpointRunId, companyId, issueId, phase: "retryable_failure", attempt: 1 },
{ runId: initialRunId, companyId, issueId, phase: "observed", attempt: 0 },
{
runId: bootstrapRetryRunId,
companyId,
issueId,
phase: "retryable_failure",
attempt: 1,
nextAttemptAt: new Date(0),
failureCode: "native_session_interrupted",
failureDetail: {
message: "provider_initialize_timeout: provider=opencode stage=health",
originalFailureCode: "provider_initialize_timeout",
recoveryMode: "bootstrap_retry",
providerSessionEstablished: false,
providerEventsExist: false,
checkpointExists: false,
},
},
...[observedExpiredRunId, observedLivePidRunId].map((observedRunId) => ({
runId: observedRunId,
companyId,
issueId,
phase: "observed" as const,
attempt: 2,
leaseOwner: "prior-native-owner",
leaseExpiresAt: new Date(0),
})),
]);
}, 30_000);
afterAll(async () => temporary?.cleanup());
it("filters durable ownership holds before the bounded resume candidate limit", async () => {
const heldRunIds = Array.from(
{ length: 26 },
(_, index) =>
`79100000-0000-4000-8000-${String(index).padStart(12, "0")}`,
);
const eligibleRunId = "79100000-0000-4000-8000-000000000099";
const candidateRunIds = [...heldRunIds, eligibleRunId];
await db.insert(heartbeatRuns).values(
candidateRunIds.map((candidateRunId) => ({
id: candidateRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
nativePhase:
candidateRunId === eligibleRunId
? "retryable_failure"
: "terminal_failure",
errorCode:
candidateRunId === eligibleRunId
? null
: "native_execution_ownership_unverified",
runnerProfileJson: {
...persistedProfile,
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: candidateRunId },
},
sessionCheckpoint: {
...persistedProfile.sessionCheckpoint,
identity: { runId: candidateRunId },
},
},
})),
);
await db.insert(nativeRunFinalizations).values(
candidateRunIds.map((candidateRunId) => ({
runId: candidateRunId,
companyId,
issueId,
phase: "retryable_failure",
attempt: 1,
nextAttemptAt: new Date(0),
leaseExpiresAt: new Date(0),
})),
);
try {
expect(
await claimNativeSessionResumptions({
db,
runnerInstanceId: "bounded-reaper",
runIds: candidateRunIds,
limit: 1,
}),
).toEqual([
{
runId: eligibleRunId,
leaseOwner: expect.stringContaining("bounded-reaper:resume:"),
},
]);
const held = await db
.select()
.from(nativeRunFinalizations)
.where(inArray(nativeRunFinalizations.runId, heldRunIds));
expect(held).toHaveLength(26);
expect(
held.every(
(row) => row.phase === "retryable_failure" && row.leaseOwner === null,
),
).toBe(true);
} finally {
await db
.delete(nativeRunFinalizations)
.where(inArray(nativeRunFinalizations.runId, candidateRunIds));
await db
.delete(heartbeatRuns)
.where(inArray(heartbeatRuns.id, candidateRunIds));
}
});
it("wins one database lease for the original result-less run without consulting the flag", async () => {
const results = await Promise.all([
claimNativeSessionResumptions({ db, runnerInstanceId: "reaper-a", runIds: [runId] }),
claimNativeSessionResumptions({ db, runnerInstanceId: "reaper-b", runIds: [runId] }),
]);
expect(results.flat()).toHaveLength(1);
expect(results.flat()[0]).toMatchObject({ runId });
await expect(db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId))).resolves.toEqual([
expect.objectContaining({ id: runId, status: "running", runtimeMode: "native" }),
]);
await expect(db.select().from(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, runId))).resolves.toEqual([
expect.objectContaining({ runId, phase: "observed", resultId: null, attempt: 1 }),
]);
});
it("dispatches the persisted run id and lease to the live same-run resume consumer", async () => {
await db.update(nativeRunFinalizations).set({
phase: "retryable_failure",
leaseOwner: null,
leaseExpiresAt: null,
nextAttemptAt: new Date(0),
}).where(eq(nativeRunFinalizations.runId, runId));
const dispatched: Array<{ runId: string; leaseOwner: string }> = [];
await expect(dispatchNativeSessionResumptions({
db,
runnerInstanceId: "heartbeat-reaper",
runIds: [runId],
dispatch: (claim) => dispatched.push(claim),
})).resolves.toHaveLength(1);
expect(dispatched).toEqual([{ runId, leaseOwner: expect.stringContaining("heartbeat-reaper:resume:") }]);
await expect(db.select().from(heartbeatRuns)).resolves.toHaveLength(8);
await expect(db.select().from(nativeRunFinalizations)).resolves.toHaveLength(8);
});
it("uses provider-neutral checkpoint-free bootstrap retry only when durable evidence proves no session existed", async () => {
await expect(claimNativeSessionResumptions({
db,
runnerInstanceId: "reaper",
runIds: [bootstrapRetryRunId],
})).resolves.toEqual([
{ runId: bootstrapRetryRunId, leaseOwner: expect.stringContaining("reaper:resume:") },
]);
});
it("never claims an expired observed coordinator without explicit retryable failure", async () => {
const dispatched: Array<{ runId: string; leaseOwner: string }> = [];
await expect(dispatchNativeSessionResumptions({
db,
runnerInstanceId: "replacement-reaper",
runIds: [observedExpiredRunId],
dispatch: (claim) => dispatched.push(claim),
})).resolves.toEqual([]);
expect(dispatched).toEqual([]);
await expect(db.select().from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, observedExpiredRunId))).resolves.toEqual([
expect.objectContaining({
phase: "observed",
attempt: 2,
leaseOwner: "prior-native-owner",
leaseExpiresAt: new Date(0),
}),
]);
});
it("blocks ambiguous observed ownership and a live unrelated persisted PID without replacement effects", async () => {
const unrelatedProcess = spawn(
process.execPath,
["-e", "setInterval(() => {}, 1_000)"],
{ stdio: "ignore" },
);
await once(unrelatedProcess, "spawn");
try {
await db.update(heartbeatRuns).set({
processPid: unrelatedProcess.pid!,
processStartedAt: new Date("2026-08-09T04:00:00.000Z"),
}).where(eq(heartbeatRuns.id, observedLivePidRunId));
const backendFactory = vi.fn((): NativeSessionBackend => ({
async descriptor() {
return {
kind: "mock",
name: "unexpected-observed-recovery",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
async openSession() {
throw new Error("observed ownership must not open a provider session");
},
async recoverSession() {
throw new Error("observed ownership must not recover a provider session");
},
}));
const heartbeat = heartbeatService(db, {
runtimeEnv: { PAPERCLIP_INSTANCE_ID: "observed-owner-test" },
nativeSessionBackendFactory: backendFactory,
});
const reaped = await heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 });
expect(reaped.runIds).not.toContain(observedExpiredRunId);
expect(reaped.runIds).not.toContain(observedLivePidRunId);
await heartbeat.drainActiveRunExecutions();
expect(backendFactory).not.toHaveBeenCalled();
expect(() => process.kill(unrelatedProcess.pid!, 0)).not.toThrow();
await expect(db.select().from(heartbeatRuns).where(eq(
heartbeatRuns.id,
observedExpiredRunId,
))).resolves.toEqual([
expect.objectContaining({
status: "running",
errorCode: "native_execution_ownership_unverified",
}),
]);
await expect(db.select().from(heartbeatRuns).where(eq(
heartbeatRuns.id,
observedLivePidRunId,
))).resolves.toEqual([
expect.objectContaining({
status: "running",
processPid: unrelatedProcess.pid,
errorCode: "native_execution_ownership_unverified",
}),
]);
for (const observedRunId of [observedExpiredRunId, observedLivePidRunId]) {
await expect(db.select().from(nativeRunFinalizations).where(eq(
nativeRunFinalizations.runId,
observedRunId,
))).resolves.toEqual([
expect.objectContaining({
phase: "observed",
attempt: 2,
leaseOwner: "prior-native-owner",
leaseExpiresAt: new Date(0),
}),
]);
await expect(db.select().from(nativeRunResults).where(eq(
nativeRunResults.runId,
observedRunId,
))).resolves.toHaveLength(0);
await expect(db.select().from(workAssessments).where(eq(
workAssessments.runId,
observedRunId,
))).resolves.toHaveLength(0);
}
} finally {
if (
unrelatedProcess.exitCode === null &&
unrelatedProcess.signalCode === null
) {
const exited = once(unrelatedProcess, "exit");
unrelatedProcess.kill("SIGKILL");
await exited;
}
}
});
it("does not resume cancelled or exhausted runs and fails closed without a checkpoint", async () => {
await expect(claimNativeSessionResumptions({
db,
runnerInstanceId: "reaper",
runIds: [cancelledRunId, exhaustedRunId, missingCheckpointRunId],
})).resolves.toEqual([]);
await expect(db.select().from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, missingCheckpointRunId))).resolves.toEqual([
expect.objectContaining({ phase: "terminal_failure", failureCode: "native_session_interrupted" }),
]);
await expect(db.select().from(issueRecoveryActions)
.where(eq(issueRecoveryActions.sourceIssueId, issueId))).resolves.toEqual([
expect.objectContaining({ cause: "native_session_interrupted", wakePolicy: null }),
]);
});
it("does not mistake the pre-first-attempt observed coordinator for an orphan", async () => {
await expect(claimNativeSessionResumptions({
db,
runnerInstanceId: "reaper",
runIds: [initialRunId],
})).resolves.toEqual([]);
await expect(db.select().from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, initialRunId))).resolves.toEqual([
expect.objectContaining({ phase: "observed", attempt: 0, failureCode: null }),
]);
});
it("sends exactly one Sentry event when a sweep resolves a run from running to failed", async () => {
const freshRunId = "79000000-0000-4000-8000-000000000101";
await db.insert(heartbeatRuns).values({
id: freshRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: freshRunId },
},
},
contextSnapshot: { issueId },
});
await db.insert(nativeRunFinalizations).values({
runId: freshRunId,
companyId,
issueId,
phase: "retryable_failure",
attempt: 1,
});
const captureCallsBefore = mockCaptureRunFailure.mock.calls.length;
await claimNativeSessionResumptions({ db, runnerInstanceId: "reaper", runIds: [freshRunId] });
// The reconciler reports asynchronously; an unrelated database round trip
// does not guarantee that callback has completed.
await waitForPendingRunFailureReports();
expect(mockCaptureRunFailure.mock.calls.slice(captureCallsBefore)).toHaveLength(1);
const newCaptures = mockCaptureRunFailure.mock.calls.slice(captureCallsBefore);
expect(newCaptures[0]?.[0]).toMatchObject({ runId: freshRunId, runStatus: "failed" });
});
it("sends no Sentry event when a sweep finds a run that is already failed", async () => {
const freshRunId = "79000000-0000-4000-8000-000000000102";
await db.insert(heartbeatRuns).values({
id: freshRunId,
companyId,
agentId,
nativeIssueId: issueId,
status: "failed",
runtimeMode: "native",
runtimeModeResolvedAt: new Date(),
runnerProfileJson: {
nativeExecutionInput: {
...persistedProfile.nativeExecutionInput,
binding: { runId: freshRunId },
},
},
contextSnapshot: { issueId },
});
await db.insert(nativeRunFinalizations).values({
runId: freshRunId,
companyId,
issueId,
phase: "retryable_failure",
attempt: 1,
});
const captureCallsBefore = mockCaptureRunFailure.mock.calls.length;
await claimNativeSessionResumptions({ db, runnerInstanceId: "reaper", runIds: [freshRunId] });
await waitForPendingRunFailureReports();
expect(mockCaptureRunFailure.mock.calls.slice(captureCallsBefore)).toHaveLength(0);
});
});
describe.each(["unchanged", "newer_active", "stale_idle"] as const)(
"P6-25 persisted reaper-to-finalization recovery (%s)",
(variant) => {
const newerRequest = variant !== "unchanged";
const staleIdleCheckpoint = variant === "stale_idle";
let temporary: Awaited<
ReturnType<typeof startEmbeddedPostgresTestDatabase>
> | null = null;
let db: ReturnType<typeof createDb>;
let staleProviderProcess: ReturnType<typeof spawn> | null = null;
const companyId = randomUUID();
const agentId = randomUUID();
const projectId = randomUUID();
const projectWorkspaceId = randomUUID();
const executionWorkspaceId = randomUUID();
const newerExecutionWorkspaceId = randomUUID();
const issueId = randomUUID();
const freshIssueId = randomUUID();
const runId = randomUUID();
const contractId = randomUUID();
const workProductId = randomUUID();
const sessionId = randomUUID();
const runnerInstanceId = randomUUID();
const newerCommentId = randomUUID();
const newerWakeKey = `newer-comment:${newerCommentId}`;
const turnId = "provider-active-turn";
const providerSessionId = "provider-existing-session";
const repoRoot = fileURLToPath(new URL("../../../", import.meta.url));
const contract = {
revision: "phase6-recovery-v1",
objective: "Recover the persisted provider turn",
criteria: [
{ id: "objective", requirement: "Complete through same-run recovery" },
],
};
const contractSha = "phase6-recovery-contract";
const evidenceRef = `work_product:${workProductId}`;
const result = structuredClone(CONTROL_PLANE_CONFORMANCE_RESULT);
result.completionClaim.contractRevision = contract.revision;
result.completionClaim.criteria[0]!.evidenceRefs = [evidenceRef];
result.evidence = [{ kind: "work_product", ref: evidenceRef }];
result.verification[0]!.artifactRef = evidenceRef;
result.summary = "Recovered the already-active provider turn.";
const terminal = {
...CONTROL_PLANE_CONFORMANCE_TERMINAL,
reportedWorkDisposition: result.reportedWorkDisposition,
};
const execution: NativeExecutionInputV1 = {
schema: "paperclip.native-execution-input.v1",
binding: { companyId, runId, issueId, agentId, executionWorkspaceId },
task: {
identifier: "NRR-1",
title: "Recover one native heartbeat",
description: null,
workMode: "standard",
},
workspace: {
cwd: repoRoot,
repoUrl: null,
repoRef: null,
branchName: null,
},
session: {
normalizedSessionId: sessionId,
driverKind: "codex_app_server",
protocolVersion: 1,
},
provider: { kind: "codex", model: null },
completionContract: {
id: contractId,
sha256: contractSha,
schemaVersion: "paperclip.completion-contract.v1",
contract,
},
interactionResponses: [],
credentialBindings: [],
};
const checkpoint: PersistedNativeSession = {
backendKind: "mock",
sessionId: "driver-existing-session",
identity: { companyId, runId, issueId, agentId, sessionId },
providerSessionId,
cursor: "1",
activeTurnId: staleIdleCheckpoint ? null : turnId,
pendingRuntimeRequests: [],
lineage: [],
};
const providerTerminalEvent: PrpEvent = {
schema: "paperclip.prp.event.v1",
sourceEventId: `${runnerInstanceId}:provider-terminal`,
sourceSeq: 1,
sourceInstanceId: runnerInstanceId,
sourceKind: "runner",
runId,
normalizedSessionId: sessionId,
turnId,
eventType: "turn.completed",
schemaVersion: 1,
priority: 0,
emittedAt: "2026-08-09T04:30:00.000Z",
payload: {},
};
const openSession = vi.fn(async () => {
throw new Error(
"same-run recovery must not open a second provider session",
);
});
const startTurn = vi.fn<NativeSession["startTurn"]>(async () => {
throw new Error(
"recovering old work must not start a replacement provider turn",
);
});
const close = vi.fn(async () => undefined);
const recoverSession = vi.fn(async (persisted: PersistedNativeSession) => {
expect(persisted).toMatchObject({
providerSessionId,
activeTurnId: staleIdleCheckpoint ? null : turnId,
});
expect(staleProviderProcess).not.toBeNull();
expect(
staleProviderProcess!.exitCode !== null ||
staleProviderProcess!.signalCode !== null,
).toBe(true);
// Model the real launch/checkpoint gap: the durable DB snapshot may still
// look idle, while the recovered provider has already started the old turn.
const recoveredSnapshot: PersistedNativeSession = {
...structuredClone(checkpoint),
cursor: "2",
activeTurnId: turnId,
};
const session: NativeSession = {
identity: () => structuredClone(checkpoint.identity),
async capabilities() {
return {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
};
},
async *events() {
yield providerTerminalEvent;
},
startTurn,
async result() {
return { result, terminal, turnId };
},
async snapshot() {
return structuredClone(recoveredSnapshot);
},
close,
};
return { recovered: true, session };
});
const backend: NativeSessionBackend = {
async descriptor() {
return {
kind: "mock",
name: "persisted-recovery-backend",
version: "1",
capabilities: {
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
},
};
},
openSession,
recoverSession,
};
beforeAll(async () => {
temporary = await startEmbeddedPostgresTestDatabase(
"paperclip-native-reaper-e2e-",
);
db = createDb(temporary.connectionString);
await instanceSettingsService(db).updateExperimental({
enableNativeRunner: newerRequest,
});
await db.insert(companies).values({
id: companyId,
name: "Native same-run recovery",
issuePrefix: "NRR",
status: "active",
defaultResponsibleUserId: "responsible-user",
});
await db
.insert(projects)
.values({
id: projectId,
companyId,
name: "Recovery project",
status: "active",
});
await db.insert(projectWorkspaces).values({
id: projectWorkspaceId,
companyId,
projectId,
name: "Recovery workspace",
cwd: repoRoot,
isPrimary: true,
});
await db.insert(agents).values({
id: agentId,
companyId,
name: "Native recovery agent",
adapterType: "paperclip_runner",
// Match the persisted fixture workspace's strategy explicitly. A
// metadata-free workspace is not proof of the current config; an actual
// strategy change correctly refuses immutable native-input rebinding.
adapterConfig: { workspaceStrategy: { type: "project_primary" } },
status: "active",
runtimeConfig: {
heartbeat: { wakeOnDemand: true, maxConcurrentRuns: 1 },
nativeRunner: {
mode: "native",
backend: "codex_app_server",
protocolVersion: 1,
},
},
});
await db.insert(issues).values({
id: issueId,
companyId,
projectId,
projectWorkspaceId,
issueNumber: 1,
identifier: "NRR-1",
title: "Recover one native heartbeat",
status: "in_progress",
assigneeAgentId: agentId,
workMode: "standard",
});
await db.insert(executionWorkspaces).values({
id: executionWorkspaceId,
companyId,
projectId,
projectWorkspaceId,
sourceIssueId: issueId,
mode: "shared_workspace",
strategyType: "project_primary",
name: "Persisted recovery workspace",
status: "active",
cwd: repoRoot,
providerType: "local_fs",
});
await db.insert(executionWorkspaces).values({
id: newerExecutionWorkspaceId,
companyId,
projectId,
projectWorkspaceId,
sourceIssueId: issueId,
mode: "shared_workspace",
strategyType: "project_primary",
name: "Newer issue workspace",
status: "active",
cwd: repoRoot,
providerType: "local_fs",
});
await db
.update(issues)
.set({
// Simulate a newer run moving the issue-level pointer before the older native run is
// recovered. The older run must still restore its own immutable workspace binding.
executionWorkspaceId: newerExecutionWorkspaceId,
executionWorkspacePreference: "reuse_existing",
executionWorkspaceSettings: { mode: "shared_workspace" },
})
.where(eq(issues.id, issueId));
await db.insert(completionContracts).values({
id: contractId,
companyId,
issueId,
revision: 1,
schemaVersion: "paperclip.completion-contract.v1",
policyVersion: "phase6-v1",
risk: "standard",
completionAuthority: "server_arbiter",
incompleteCriteriaPolicy: "preserve_non_terminal",
contractJson: contract,
canonicalSha256: contractSha,
createdByActorType: "system",
createdByActorId: "test",
});
await db.insert(issueWorkProducts).values({
id: workProductId,
companyId,
issueId,
type: "artifact",
provider: "paperclip",
title: "Recovered result evidence",
status: "ready_for_review",
reviewState: "approved",
});
await db.insert(heartbeatRuns).values({
id: runId,
companyId,
agentId,
nativeIssueId: issueId,
status: "running",
runtimeMode: "native",
runtimeModeResolverVersion: "phase6-v1",
runtimeModeReason: "eligible_opt_in",
runtimeModeResolvedAt: new Date("2026-08-09T04:00:00.000Z"),
runnerProfileJson: {
mode: "native",
backend: "codex_app_server",
protocolVersion: 1,
nativeExecutionInput: execution,
sessionCheckpoint: checkpoint,
},
runnerInstanceId,
nativeSessionId: sessionId,
driverKind: "codex_app_server",
driverVersion: "phase6-v1",
completionContractId: contractId,
completionContractSha256: contractSha,
nativePhase: "retryable_failure",
nativePhaseUpdatedAt: new Date("2026-08-09T04:00:00.000Z"),
startedAt: new Date("2026-08-09T04:00:00.000Z"),
contextSnapshot: { issueId, taskId: issueId, skipIssueComment: true },
});
await db
.update(issues)
.set({ executionRunId: runId })
.where(eq(issues.id, issueId));
await db.insert(nativeRunFinalizations).values({
runId,
companyId,
issueId,
phase: "retryable_failure",
attempt: 1,
failureCode: "native_session_interrupted",
nextAttemptAt: new Date(0),
});
if (newerRequest)
await db.insert(issueComments).values({
id: newerCommentId,
companyId,
issueId,
authorUserId: "responsible-user",
body: "Verify the newest user instruction before completing this task.",
});
}, 30_000);
afterAll(async () => {
if (
staleProviderProcess &&
staleProviderProcess.exitCode === null &&
staleProviderProcess.signalCode === null
) {
staleProviderProcess.kill("SIGKILL");
}
if (temporary) {
await drainHeartbeatRunsToQuiescence(
db,
heartbeatService(db, {
runtimeEnv: { PAPERCLIP_INSTANCE_ID: "phase6-recovery-test" },
nativeSessionBackendFactory: () => backend,
}),
);
await temporary.cleanup();
}
});
it("preserves the old provider contract and admits newer direction only in a separate turn", async () => {
legacyAdapterExecute.mockClear();
staleProviderProcess = spawn(
process.execPath,
["-e", "setInterval(() => {}, 1_000)"],
{ stdio: "ignore" },
);
await once(staleProviderProcess, "spawn");
expect(staleProviderProcess.pid).toEqual(expect.any(Number));
await db
.update(heartbeatRuns)
.set({
processPid: staleProviderProcess.pid!,
processStartedAt: new Date("2026-08-09T04:00:00.000Z"),
})
.where(eq(heartbeatRuns.id, runId));
const factoryInputs: NativeExecutionInput[] = [];
const freshStartTurn = vi.fn<NativeSession["startTurn"]>();
const freshBackend = (
executionInput: NativeExecutionInput,
): NativeSessionBackend => {
const freshRunId = executionInput.binding.runId;
const freshTurnId = `new-direction:${freshRunId}`;
const freshIdentity = {
...checkpoint.identity,
runId: freshRunId,
sessionId: executionInput.session.normalizedSessionId,
};
let releaseStart!: () => void;
const started = new Promise<void>((resolve) => {
releaseStart = resolve;
});
let freshResult = structuredClone(result);
const session: NativeSession = {
identity: () => freshIdentity,
capabilities: async () => ({
resume: true,
typedEvents: true,
steering: false,
interruption: true,
structuredResult: true,
}),
async *events() {
await started;
const [freshRun] = await db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, freshRunId));
yield {
...providerTerminalEvent,
sourceEventId: `${freshRun.runnerInstanceId}:fresh-terminal`,
sourceInstanceId: freshRun.runnerInstanceId!,
runId: freshRunId,
normalizedSessionId: freshIdentity.sessionId,
turnId: freshTurnId,
};
},
async startTurn(input) {
freshStartTurn(input);
const delivered = JSON.parse(input.message.text);
expect(delivered.task.prompt).toContain(
"Verify the newest user instruction",
);
const currentContract =
delivered.completionContract as typeof contract;
freshResult.completionClaim.contractRevision =
currentContract.revision;
freshResult.completionClaim.criteria = currentContract.criteria.map(
(criterion) => ({
...freshResult.completionClaim.criteria[0]!,
criterionId: criterion.id,
}),
);
freshResult.summary =
"Completed the separately delivered newer instruction.";
releaseStart();
return { turnId: freshTurnId };
},
async result() {
return { result: freshResult, terminal, turnId: freshTurnId };
},
async snapshot() {
return {
...checkpoint,
identity: freshIdentity,
providerSessionId: `fresh:${freshRunId}`,
activeTurnId: null,
cursor: "0",
};
},
async close() {},
};
return {
async descriptor() {
return {
...(await backend.descriptor()),
runtimeContextCapabilities: {
instructions: "native",
skills: "native",
mcp: "native",
},
};
},
async openSession(input) {
expect(input.identity).toEqual(freshIdentity);
return session;
},
async recoverSession() {
throw new Error(
"newer user cause must not reuse the prior active turn",
);
},
};
};
const backendFactory = vi.fn((input: NativeExecutionInput) => {
factoryInputs.push(structuredClone(input));
return input.binding.runId === runId ? backend : freshBackend(input);
});
const heartbeat = heartbeatService(db, {
runtimeEnv: { PAPERCLIP_INSTANCE_ID: "phase6-recovery-test" },
nativeSessionBackendFactory: backendFactory,
});
if (newerRequest) {
// Restore the independently admitted receipt that was deferred while the
// original provider was running, before this simulated process loss.
await db.insert(agentWakeupRequests).values({
companyId,
agentId,
status: "deferred_issue_execution",
source: "automation",
triggerDetail: "system",
reason: "issue_commented",
idempotencyKey: newerWakeKey,
requestedByActorType: "user",
requestedByActorId: "responsible-user",
payload: {
issueId,
commentId: newerCommentId,
_paperclipWakeContext: {
issueId,
taskId: issueId,
commentId: newerCommentId,
wakeCommentIds: [newerCommentId],
skipIssueComment: true,
},
},
});
const [receipt] = await db
.select()
.from(agentWakeupRequests)
.where(eq(agentWakeupRequests.idempotencyKey, newerWakeKey));
expect(receipt.status).toBe("deferred_issue_execution");
}
await expect(
heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 }),
).resolves.not.toContain(runId);
await heartbeat.drainActiveRunExecutions();
expect(backendFactory).not.toHaveBeenCalled();
expect(recoverSession).not.toHaveBeenCalled();
expect(() => process.kill(staleProviderProcess!.pid!, 0)).not.toThrow();
await expect(
db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)),
).resolves.toEqual([
expect.objectContaining({
id: runId,
status: "running",
processPid: staleProviderProcess.pid,
errorCode: "native_execution_ownership_unverified",
}),
]);
await expect(
db
.select()
.from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, runId)),
).resolves.toEqual([
expect.objectContaining({
phase: "retryable_failure",
attempt: 1,
leaseOwner: null,
}),
]);
const unrelatedProcessExit = once(staleProviderProcess, "exit");
staleProviderProcess.kill("SIGKILL");
await unrelatedProcessExit;
await db
.update(nativeRunFinalizations)
.set({
leaseOwner: null,
leaseExpiresAt: null,
})
.where(eq(nativeRunFinalizations.runId, runId));
await expect(
heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 }),
).resolves.not.toContain(runId);
await heartbeat.drainActiveRunExecutions();
if (newerRequest) {
await drainHeartbeatRunsToQuiescence(db, heartbeat);
const original = factoryInputs.find(
(input) => input.binding.runId === runId,
);
expect(original?.completionContract).toEqual(
execution.completionContract,
);
expect(original?.task.prompt).toBe(execution.task.title);
expect(startTurn).not.toHaveBeenCalled();
const [accepted] = await db
.select()
.from(nativeRunResults)
.where(eq(nativeRunResults.runId, runId));
expect(accepted?.completionContractId).toBe(contractId);
const [receipt] = await db
.select()
.from(agentWakeupRequests)
.where(eq(agentWakeupRequests.idempotencyKey, newerWakeKey));
expect(receipt.payload).toMatchObject({ commentId: newerCommentId });
expect(receipt.status).toBe("completed");
expect(receipt.runId).not.toBe(runId);
expect(freshStartTurn).toHaveBeenCalledOnce();
const [freshRun] = await db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, receipt.runId!));
expect(freshRun).toMatchObject({
status: "succeeded",
nativePhase: "committed",
});
expect(
await db
.select()
.from(nativeRunResults)
.where(eq(nativeRunResults.runId, receipt.runId!)),
).toHaveLength(1);
expect(
factoryInputs.filter((input) => input.binding.runId === runId),
).toHaveLength(1);
await heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 });
await drainHeartbeatRunsToQuiescence(db, heartbeat);
expect(freshStartTurn).toHaveBeenCalledOnce();
expect(
await db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.companyId, companyId)),
).toHaveLength(2);
return;
}
const recoveryState = {
run: await db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId)),
coordinator: await db
.select()
.from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, runId)),
};
expect(
backendFactory.mock.calls.length,
JSON.stringify(recoveryState),
).toBe(1);
expect(recoverSession).toHaveBeenCalledOnce();
expect(openSession).not.toHaveBeenCalled();
expect(startTurn).not.toHaveBeenCalled();
expect(close).toHaveBeenCalledOnce();
expect(legacyAdapterExecute).not.toHaveBeenCalled();
await expect(
db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.companyId, companyId)),
).resolves.toEqual([
expect.objectContaining({
id: runId,
runtimeMode: "native",
status: "succeeded",
nativePhase: "committed",
processPid: null,
processGroupId: null,
processStartedAt: null,
}),
]);
await expect(
db
.select()
.from(nativeRunResults)
.where(eq(nativeRunResults.runId, runId)),
).resolves.toHaveLength(1);
await expect(
db
.select()
.from(workAssessments)
.where(eq(workAssessments.runId, runId)),
).resolves.toHaveLength(1);
const decisions = await db
.select()
.from(statusDecisions)
.where(eq(statusDecisions.issueId, issueId));
expect(decisions).toEqual([
expect.objectContaining({
reasonCode: "completion_contract_satisfied",
toStatus: "done",
applicationState: "applied",
}),
]);
const effects = await db
.select()
.from(statusDecisionEffects)
.where(eq(statusDecisionEffects.issueId, issueId));
expect(new Set(effects.map((effect) => effect.decisionId))).toEqual(
new Set([decisions[0]!.id]),
);
expect(effects.map((effect) => effect.effectKind).sort()).toEqual([
"issue_status_projection",
"release_checkout",
]);
await expect(
db.select().from(issues).where(eq(issues.id, issueId)),
).resolves.toEqual([
expect.objectContaining({
status: "done",
statusVersion: 1,
lastStatusDecisionId: decisions[0]!.id,
executionWorkspaceId: newerExecutionWorkspaceId,
}),
]);
await expect(
db
.select()
.from(executionWorkspaces)
.where(eq(executionWorkspaces.companyId, companyId)),
).resolves.toHaveLength(2);
await expect(
db
.select()
.from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, runId)),
).resolves.toEqual([
expect.objectContaining({
phase: "committed",
resultId: expect.any(String),
assessmentId: expect.any(String),
decisionId: decisions[0]!.id,
}),
]);
await expect(
db
.select()
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runId)),
).resolves.toEqual(
expect.arrayContaining([
expect.objectContaining({ eventType: "turn.completed" }),
expect.objectContaining({ eventType: "run.result.accepted" }),
expect.objectContaining({ eventType: "run.terminal" }),
]),
);
await expect(
db.select().from(activityLog).where(eq(activityLog.entityId, issueId)),
).resolves.toEqual(
expect.arrayContaining([
expect.objectContaining({ action: "issue.updated" }),
]),
);
await heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 });
await heartbeat.drainActiveRunExecutions();
expect(backendFactory).toHaveBeenCalledOnce();
await expect(
db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.companyId, companyId)),
).resolves.toHaveLength(1);
await expect(
db
.select()
.from(nativeRunResults)
.where(eq(nativeRunResults.runId, runId)),
).resolves.toHaveLength(1);
await expect(
db
.select()
.from(workAssessments)
.where(eq(workAssessments.runId, runId)),
).resolves.toHaveLength(1);
await expect(
db
.select()
.from(statusDecisions)
.where(eq(statusDecisions.issueId, issueId)),
).resolves.toHaveLength(1);
// The persisted Paperclip Runner run above remains recoverable while the
// flag is off. Switching the agent back to a direct adapter now proves a
// fresh run ignores the stale native profile and stays on the legacy path.
await db
.update(agents)
.set({ adapterType: "codex_local" })
.where(eq(agents.id, agentId));
await db.insert(issues).values({
id: freshIssueId,
companyId,
projectId,
projectWorkspaceId,
issueNumber: 2,
identifier: "NRR-2",
title: "Start only after the native kill switch is off",
status: "in_progress",
assigneeAgentId: agentId,
workMode: "standard",
});
const directExecute = legacyAdapterExecute.getMockImplementation()!;
legacyAdapterExecute.mockImplementationOnce(async () => {
await db.update(issues).set({ status: "done" }).where(eq(issues.id, freshIssueId));
return directExecute();
});
const fresh = await heartbeat.wakeup(agentId, {
source: "automation",
triggerDetail: "system",
reason: "issue_commented",
payload: { issueId: freshIssueId },
contextSnapshot: {
issueId: freshIssueId,
taskId: freshIssueId,
skipIssueComment: true,
},
});
expect(fresh).not.toBeNull();
await drainHeartbeatRunsToQuiescence(db, heartbeat);
expect(legacyAdapterExecute).toHaveBeenCalledOnce();
await expect(
db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, fresh!.id)),
).resolves.toEqual([
expect.objectContaining({
agentId,
runtimeMode: "legacy",
runtimeModeReason: "direct_adapter",
status: "succeeded",
}),
]);
await expect(
db
.select()
.from(nativeRunFinalizations)
.where(eq(nativeRunFinalizations.runId, fresh!.id)),
).resolves.toHaveLength(0);
await expect(
db
.select()
.from(nativeRunResults)
.where(eq(nativeRunResults.runId, fresh!.id)),
).resolves.toHaveLength(0);
await expect(
db
.select({ runtimeConfig: agents.runtimeConfig })
.from(agents)
.where(eq(agents.id, agentId)),
).resolves.toEqual([
expect.objectContaining({
runtimeConfig: expect.objectContaining({
nativeRunner: {
mode: "native",
backend: "codex_app_server",
protocolVersion: 1,
},
}),
}),
]);
}, 30_000);
},
);