fix: preserve warm Codex turns with incremental managed file checkpoints (#14735)

## Thinking Path

> - Paperclip manages AI agents and keeps their instructions and files
durable.
> - Native Codex runners can keep a process alive between compatible
turns.
> - Managed file collection stopped that process after each turn, which
defeated warm reuse.
> - Agent folders can contain large images and other files, so full
copies on every turn are expensive.
> - This change keeps one managed directory for the live session and
saves only file changes after each turn.
> - Ownership, authorization, instruction changes, and process
retirement still control when reuse is safe.

## Linked Issues or Issue Description

Related: #13710 introduced native warm session reuse. This fixes managed
file collection that still forced those sessions to stop. No duplicate
open PR or issue was found.

**What happened?**

With managed instructions and warm native Codex enabled, consecutive
turns reused a Daytona sandbox but started a new runner process each
time. The managed directory collector required process termination
before saving files.

**Expected behavior**

Compatible turns keep the same process and managed `AGENT_HOME`. Each
completed turn saves added, changed, and deleted files before the next
turn starts. Unchanged large files do not transfer again.

**Steps to reproduce**

1. Use a native Codex agent with managed instructions and a reusable
Daytona environment.
2. Enable warm session reuse and run three turns on the same task.
3. Write a large binary on the first turn, edit a small note on each
turn, and delete a file on the second turn.
4. Compare process identity across turns and read the canonical files
through the public agent-files API.

**Paperclip version or commit**

Reproduced on `d30b03bd8c17604cdab1533eeeeb087aba30e8b1`.

**Deployment mode**

Local server with remote Daytona execution; cloud native runner uses the
same path.

## What Changed

- Retain the managed directory only for the verified owner of a live
native Codex session.
- Checkpoint each completed turn before releasing the session for reuse.
Retry unstable captures, then stop and collect when a warm checkpoint
cannot be validated.
- Compare metadata and cached hashes, stream only changed file payloads,
record deletions, and validate path, content, quota, and authorization
before saving.
- Rotate sessions when canonical files, loaded instructions,
credentials, or launch policy change. Fence stale collection and cleanup
callbacks from later owners.
- Keep cleanup and recovery aware of the current session owner. Recheck
canonical files under the writer lock at handoff, attach the successor
collector before fallible bookkeeping, and emit one final save receipt
on checkpoint fallback. Preserve storage warnings across unchanged
checkpoints.
- Add regression coverage and a three-turn Daytona test with independent
public API file checks, an unchanged 8 MiB binary, deletion checks, and
strict process identity checks.
- Document checkpoint consistency, lifecycle behavior, and local run-log
counters.
- Replace a timing assumption in the Daytona teardown test with explicit
transfer-arrival gates after CI exposed an unset release callback.

## Verification

- Full local `pnpm -r typecheck` and `pnpm build` passed. Server checks
were repeated after the final storage-warning fix.
- Runner E2E typecheck and 749 runner E2E unit tests passed.
- Focused file checkpoint, directory ownership, instruction collection,
native session, and merge tests passed. After review fixes, the
managed-directory and native-session suites passed 550 tests, including
intervening canonical edits, same-run fresh restore, failed handoff
collection, and one-call fallback collection. Server typecheck passed
again. The Daytona plugin suite passed 218 tests. The quota-warning
regression failed before the fix and passed afterward.
- Three real Daytona campaigns passed before the final handoff review
fixes. The latest kept PID 547 across all three turns. The first
checkpoint copied 8,388,635 bytes; the next two copied 36 and 54 bytes.
Public API reads verified the binary, note contents, and deletion after
every turn. Test cleanup deleted the sandbox.
- The final head was also deployed to an isolated cloud staging instance
and passed three UI-triggered native Codex turns with managed
instructions. All three retained the same process ID/start time, native
session, provider session, runner instance, and Daytona sandbox.
Checkpoints copied 8,388,643 bytes on turn 1, then only 52 and 78 bytes
on turns 2 and 3; those warm captures also hashed only 52 and 78 bytes.
Independent canonical API reads verified every byte of the unchanged 8
MiB binary and the exact note contents after every turn; the deleted
file returned 404 after turns 2 and 3. After restoring the original
lifecycle and agent-auth configuration, removing the temporary secret,
pausing the test agent, and deleting both test sandboxes, independent
canonical API reads still verified the entire binary, the final 78-byte
three-line note, and the deletion. The native runner flag remained
enabled and the final serving revision remained the PR head.
- Two earlier staging attempts are preserved as failures and are
excluded from the acceptance result: a saved ChatGPT login failed with a
provider routing 401, and its subsequent stopped-sandbox retry failed
before provider startup with a closed-lease admission error. The
successful campaign used a fresh sandbox and a temporary encrypted
API-key binding. The stopped-lease retry remains unexplained; this
campaign does not establish recovery of that failed sandbox.
- All [Paperclip CI
gates](https://github.com/paperclipai/paperclip/actions/runs/36750397355)
pass on `26ef2ef56a389259246809805c0b34a4747eb86b`, including full test
partitions, build, typecheck, runner verification, E2E shards, and the
Canary clean public-npm install. Greptile reviewed that exact head at
5/5 with no unresolved review threads or outstanding findings.
- Full local repository coverage used the existing CI partitions, but
the 40,000-file Git streaming stress test timed out and its local retry
was interrupted by macOS thermal emergency sleep; this is not a green
full local suite claim. The exact stress test passed on the final head
in [CI server shard
2/12](https://github.com/paperclipai/paperclip/actions/runs/36750397355/job/110008294290),
in 111.9 seconds.
- Repeat the live test with configured credentials and a Linux runner
artifact: `pnpm test:e2e:runner -- --id
daytona-warm-continuity.runner-codex.daytona.warm-three-turn`.

## Risks

- This is a file-level checkpoint, not an atomic snapshot of the whole
folder. Background writes after a capture are saved by the next
checkpoint or final stopped collection.
- Metadata scans still visit all paths. Modified files transfer in full;
unchanged files do not rehash or transfer.
- Incorrect ownership or reuse could collect the wrong directory. Run
ownership fences, current authorization, stable capture validation, and
stopped collection fallbacks are covered by tests.
- Warm reuse remains opt-in. No database migration or fleet default
changes.

## Model Used

OpenAI GPT-6 through Codex, with reasoning, code editing, tool use, and
test execution. The exact serving model ID and context-window size are
not exposed by this session.

## Checklist

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

---------

Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
Dotta
2026-09-30 13:33:37 -05:00
committed by GitHub
co-authored by Paperclip
parent d432dc7fa3
commit b54b2dc35c
23 changed files with 1195 additions and 92 deletions
+44 -22
View File
@@ -23,7 +23,10 @@ The canonical host directory keeps its existing physical location:
any-supported-file
file-sync/ controller-only operational state
adopted.json
runs/<run>/live/ isolated writable copy while a run executes
canonical-manifest.json metadata/hash cache, no file contents
runs/<initial-run>/
owner.json current run owning the writable copy
live/ writable copy for the provider lifetime
```
The process starts in its existing task workspace. `AGENT_HOME` points to the
@@ -75,13 +78,16 @@ the run detail also shows warnings from its save receipt. The editor only shows
preserved instruction-only candidates that may need review, alongside errors
from the current browser edit. Later successful saves do not erase run history.
Larger folders take longer to hash, copy, and transfer on each run. There is one
canonical folder plus temporary working copies for currently active runs (and
remote staging when the transport needs it). No additional captured tree is
created. Terminal runs remove their private trees and baseline metadata, keeping
only a small receipt. Restart recovery retries interrupted cleanup without
removing a running provider's files. These are not aggregate disk quotas; the
operator still provisions storage for agents and the configured run concurrency.
Initial restoration copies the canonical folder once. Warm native Codex turns
reuse that working directory. Checkpoints enumerate file metadata, hash files
whose identity/size/mode/mtime/ctime changed, and temporarily copy and transfer
only changed file contents. Deletions and empty directories travel as manifest
entries. An unchanged image is neither rehashed nor recopied after its first
checkpoint. A modified file is transferred in full; this is a file-level delta,
not block-level deduplication. Canonical hash caches and manifests contain no
file contents. Temporary checkpoint payloads are removed after application.
The operator still provisions storage for the canonical folders, active working
copies and changed-file payloads; these are not aggregate disk quotas.
## Run lifecycle
@@ -90,23 +96,39 @@ operator still provisions storage for agents and the configured run concurrency.
revision history.
2. Stage the copy through the existing workspace transport. Point `AGENT_HOME`
and instruction guidance at that registered root.
3. At the provider's verified checkpoint-and-stop boundary, retrieve the entire
directory into the existing working copy before releasing its environment.
3. For warm native Codex, capture a manifest and changed-file payload at each
terminal turn boundary before admitting another turn. Check file metadata
before and after streaming and hash the captured payload independently on
the host. Retry an unstable checkpoint up to three times; if it cannot be
validated, close the owned provider and perform the stopped collector. Other
execution paths retain their stopped-provider collection boundary.
4. Recheck the responsible user's current authorization. Under the same agent
lock used by editor writes, apply only files changed or deleted relative to
the starting baseline. For a competing edit or deletion of the same file,
the last synchronization to acquire the lock wins. Unchanged files do not
overwrite another run's changes; newly added unrelated files survive.
5. Record the outcome and remove temporary copies for successful and failed
runs. No per-run file versions, conflict copies, or review queue accumulate.
The next run starts with the current directory.
the last acknowledged baseline. For a competing edit or deletion of the same
file, the last synchronization to acquire the lock wins. Unchanged files do
not overwrite another run's changes; newly added unrelated files survive.
5. Record the save receipt and advance the baseline only after application. A
warm session keeps its directory and hands ownership to the next run using
a controller-owned marker. Old callbacks cannot collect or remove the next
owner's files. On session retirement, collect any later writes and remove
the private directory. No per-run file versions or conflict copies accumulate.
The whole-directory contract closes the provider process to establish a safe
collection boundary, including child processes. It preserves the provider's
resumable conversation. Only the loaded instruction entry participates in the
new runtime instruction digest; adding or editing another file does not change
that digest. Relative supporting files are read from `AGENT_HOME`, not from the
read-only prompt snapshot.
These are validated **per-file checkpoints**, not an atomic snapshot of arbitrary
background writers across an entire directory. Writes after a checkpoint remain
pending until the next checkpoint or verified session retirement. A save receipt
acknowledges only the captured bytes. Lost remote bytes or missing stop proof
cannot become a successful save.
Warm reuse requires the actual live session, the same remote environment and
provider lease, and unchanged canonical files since its last checkpoint. An
editor or another task changing canonical files retires that session before a
fresh copy is restored. A directory-path mismatch also forces retirement. Only
the loaded instruction entry participates in the runtime instruction digest;
ordinary memory/image edits do not change it. Changes to loaded instructions,
policy, credentials or provider configuration may still replace the process.
Relative supporting files are read from `AGENT_HOME`, not the read-only prompt
snapshot. External instruction bundles remain read-only and use their existing
lifecycle.
The editor supplies the hash of the file it read. A stale browser save returns
409 and retains the user's unsaved draft. Run synchronization itself uses
+17
View File
@@ -318,3 +318,20 @@ The message distinguishes an automatic retry from work that is no longer eligibl
This pre-provider wait records `ai_connection_busy` on the cancelled run and does
not consume the provider-failure retry allowance. The event contains no credentials
and creates no Telemetry or OpenTelemetry export.
## Managed Agent File Save Receipts
The server writes `instruction_save` after managed file collection or a warm
turn checkpoint. The payload includes the save state, instruction entry path,
storage warning, and error code/message. Agent-directory receipts identify
`contract: "agent_files"` and the applied candidate hash. Legacy instruction
receipts instead identify the saved revision.
A validated warm checkpoint reports `saved` or `unchanged`, even though its
working directory remains owned by the live session. An unstable checkpoint
reports `pending_collection` until stopped collection produces a final receipt.
Successful checkpoints can include `checkpointStats`: `scannedEntries`,
`hashedBytes`, `copiedFiles`, and `copiedBytes`. These counts describe that
capture, not cumulative traffic or an atomic snapshot of background writers.
They contain no file contents. The receipt remains in the instance run log;
it adds no Paperclip Telemetry or OpenTelemetry export.
@@ -571,15 +571,18 @@ export async function mergeDirectoryWithBaseline(input: {
conflictPolicy?: "reject";
beforeApply?: () => Promise<void>;
afterApply?: () => Promise<void>;
/** Caller holds the target's writer lock and validated an immutable sparse
* source. Unchanged entries need no payload and are never copied. */
snapshots?: { source: DirectorySnapshot; current: DirectorySnapshot };
}): Promise<void> {
const options = { exclude: input.baseline.exclude, ignoredPaths: input.baseline.ignoredPaths, diskBacked: true };
const source = await captureDirectorySnapshot(input.sourceDir, options);
const source = input.snapshots?.source ?? await captureDirectorySnapshot(input.sourceDir, options);
try {
await withDirectoryMergeLock(input.targetDir, async (canonicalTargetDir) => {
await input.beforeApply?.();
// Strict preflight must see excluded children before a directory is
// replaced. The merge still applies only the filtered source/baseline.
const current = await captureDirectorySnapshot(canonicalTargetDir,
const current = input.snapshots?.current ?? await captureDirectorySnapshot(canonicalTargetDir,
input.conflictPolicy === "reject" ? { exclude: [], diskBacked: true } : options);
try {
if (input.conflictPolicy === "reject") {
@@ -5524,16 +5524,22 @@ describe("daytona native file-sync hooks", () => {
const sandbox = createMockSandbox({ id: "sandbox-123" });
// Hold the inbound upload and the outbound download open at the same time, so
// the shared lease has two active sync calls when teardown starts.
let uploadArrived!: () => void;
const uploadStarted = new Promise<void>((resolve) => { uploadArrived = resolve; });
let releaseUpload!: () => void;
sandbox.fs.uploadFiles.mockImplementation(async () => {
await new Promise<void>((resolve) => {
releaseUpload = resolve;
uploadArrived();
});
});
let downloadArrived!: () => void;
const downloadStarted = new Promise<void>((resolve) => { downloadArrived = resolve; });
let releaseDownload!: () => void;
sandbox.fs.downloadFiles.mockImplementation(async (requests: Array<{ source: string; destination: string }>) => {
await new Promise<void>((resolve) => {
releaseDownload = resolve;
downloadArrived();
});
return Promise.all(
requests.map(async (request) => {
@@ -5550,9 +5556,9 @@ describe("daytona native file-sync hooks", () => {
const outboundCall = plugin.definition.onEnvironmentSyncOut?.(
syncOutParams({ operationId: "out-active", sourcePath: `${REMOTE_DIR}/out.txt`, targetPath: outboundTarget }),
);
// Let both sync calls register on the activity gate and reach their hung
// transfer, so teardown sees a refCount of two.
await new Promise((resolve) => setTimeout(resolve, 0));
// Wait for the actual transfers. One event-loop tick does not guarantee
// that the inbound filesystem reads have finished on a busy runner.
await Promise.all([uploadStarted, downloadStarted]);
const destroyCall = plugin.definition.onEnvironmentDestroyLease?.({
driverKey: "daytona",
@@ -18,6 +18,7 @@ import { resolveManagedInstructionsRoot } from "../services/agent-instructions.j
import { buildNativeRuntimeContext } from "../services/native-runtime/runtime-context.js";
import type { EnvironmentRuntimeService } from "../services/environment-runtime.js";
import { remoteTerminationReceipt } from "../services/remote-execution-termination.js";
import { AgentDirectoryReuseInvalidatedError, agentDirectoryWorkingCopyService } from "../services/agent-directory-working-copies.js";
describe("persistent agent directories", () => {
let database: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>>;
@@ -31,10 +32,10 @@ describe("persistent agent directories", () => {
const initial = "\uFEFF# Original\r\n☃\n";
const target = () => ({ companyId, agentId });
const board = () => ({ type: "board" as const, userId, source: "session" as const });
async function run() {
async function run(options: { warm?: boolean; reuseRunId?: string } = {}) {
const runId = randomUUID();
await db.insert(heartbeatRuns).values({ id: runId, companyId, agentId, invocationSource: "on_demand", responsibleUserId: userId });
return (await copies.prepare({ ...target(), runId, cwd: home }))!;
return (await copies.prepare({ ...target(), runId, cwd: home, ...options }))!;
}
beforeAll(async () => {
home = await fs.realpath(await fs.mkdtemp(path.join(os.tmpdir(), "instruction-working-copies-")));
@@ -165,16 +166,18 @@ describe("persistent agent directories", () => {
expect(await fs.readFile(path.join(root, "next-task.txt"), "utf8")).toBe("still working");
});
it("warns on each run at the file limit and clears the warning after ordinary agent cleanup", async () => {
it.each([false, true])("warns on each run at the file limit and clears the warning after ordinary agent cleanup (warm=%s)", async (warm) => {
await sparseFile(path.join(root, "full.bin"), MAX_AGENT_FILE_BYTES);
const first = await run();
const first = await run({ warm });
expect(first.receipt?.storageWarning).toContain("256 MiB");
expect(instructionWorkingCopyGuidance(first)).toContain("Runs can continue");
const unchanged = await copies.collectStopped({ companyId, runId: first.runId });
expect(unchanged).toMatchObject({ state: "unchanged", errorCode: null });
const unchanged = await copies[warm ? "checkpointWarm" : "collectStopped"]({ companyId, runId: first.runId });
expect(unchanged).toMatchObject({ state: warm ? "warm_saved" : "unchanged", errorCode: null });
if (warm) expect(unchanged?.receipt?.checkpointState).toBe("unchanged");
expect(unchanged?.receipt?.storageWarning).toContain("Agent storage is full");
copies = agentInstructionWorkingCopyService(db);
const second = await run();
const second = await run({ warm, ...(warm ? { reuseRunId: first.runId } : {}) });
if (warm) expect(second.localRoot).toBe(first.localRoot);
expect(second.receipt?.storageWarning).toContain("Agent storage is full");
await fs.unlink(path.join(second.localRoot, "full.bin"));
await fs.writeFile(path.join(second.localRoot, "task-output.txt"), "work continues");
@@ -443,6 +446,150 @@ describe("persistent agent directories", () => {
expect(after.context.aggregateDigest).toBe(before.context.aggregateDigest);
expect(after.context.instructions.bundle.fileCount).toBe(1);
});
it("checkpoints three managed turns into one live directory, fences old cleanup, and persists after retirement", async () => {
let copy = await run({ warm: true });
const original = copy;
for (let turn = 1; turn <= 3; turn++) {
await fs.appendFile(path.join(copy.localRoot, "memory.txt"), `turn ${turn}\n`);
const saved = (await copies.checkpointWarm({ companyId, runId: copy.runId }))!;
expect(saved.state).toBe("warm_saved");
expect(saved.processStoppedAt).toBeNull();
expect(await agentFileStore(db).read(companyId, agentId, "memory.txt", board())).toEqual(Buffer.from(Array.from({ length: turn }, (_, i) => `turn ${i + 1}\n`).join("")));
await copies.release(companyId, copy.runId);
expect(await fs.stat(copy.localRoot)).toBeTruthy();
if (turn < 3) {
expect(await copies.canReuseWarm(companyId, agentId, copy.runId)).toBe(true);
copy = await run({ warm: true, reuseRunId: copy.runId });
expect(copy.executionRoot).toBe(original.executionRoot);
await copies.collectStopped({ companyId, runId: original.runId });
await copies.release(companyId, original.runId);
expect(await fs.stat(copy.localRoot)).toBeTruthy();
}
}
// A child may write after the last turn's checkpoint; session retirement
// must collect this delta using the current owner, not the first run.
await fs.writeFile(path.join(copy.localRoot, "late.txt"), "after turn");
expect((await copies.collectStopped({ companyId, runId: copy.runId }))?.state).toBe("saved");
await expect(fs.stat(copy.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
const fresh = await run();
expect(await fs.readFile(path.join(fresh.localRoot, "memory.txt"), "utf8")).toBe("turn 1\nturn 2\nturn 3\n");
expect(await fs.readFile(path.join(fresh.localRoot, "late.txt"), "utf8")).toBe("after turn");
});
it("rechecks canonical edits at handoff and releases the writer lock before retirement and fresh restore", async () => {
const first = await run({ warm: true });
await copies.checkpointWarm({ companyId, runId: first.runId });
expect(await copies.canReuseWarm(companyId, agentId, first.runId)).toBe(true);
await agentFileStore(db).write({ ...target(), path: "editor.txt", bytes: Buffer.from("new canonical content"), baseHash: null }, board());
const runId = randomUUID();
await db.insert(heartbeatRuns).values({ id: runId, companyId, agentId, invocationSource: "on_demand", responsibleUserId: userId });
await expect(copies.prepare({ ...target(), runId, cwd: home, warm: true, reuseRunId: first.runId })).rejects.toBeInstanceOf(AgentDirectoryReuseInvalidatedError);
expect(JSON.parse(await fs.readFile(path.join(path.dirname(first.localRoot), "owner.json"), "utf8")).runId).toBe(first.runId);
// Both paths need the canonical writer lock. A handoff failure must release
// it before orchestration stops the old session and restores a fresh copy.
expect((await copies.collectStopped({ companyId, runId: first.runId }))?.state).toBe("unchanged");
const next = (await copies.prepare({ ...target(), runId, cwd: home, warm: true }))!;
expect(next.localRoot).not.toBe(first.localRoot);
expect(await fs.readFile(path.join(next.localRoot, "editor.txt"), "utf8")).toBe("new canonical content");
}, 10_000);
it("attaches the successor collector before fallible post-handoff bookkeeping", async () => {
const first = await run({ warm: true });
await copies.checkpointWarm({ companyId, runId: first.runId });
const runId = randomUUID();
await db.insert(heartbeatRuns).values({ id: runId, companyId, agentId, invocationSource: "on_demand", responsibleUserId: userId });
let collectorRunId = first.runId;
const handoff = vi.fn((copy: typeof first) => { collectorRunId = copy.runId; });
const directories = agentDirectoryWorkingCopyService(db, copies.get, async () => { throw new Error("injected post-handoff database failure"); });
await expect(directories.prepare({ ...target(), runId, cwd: home, warm: true, reuseRunId: first.runId, onWarmHandoff: handoff }))
.rejects.toThrow("injected post-handoff database failure");
expect(handoff).toHaveBeenCalledOnce();
expect(collectorRunId).toBe(runId);
await fs.writeFile(path.join(first.localRoot, "last-write.txt"), "saved after failed handoff");
// The stale collector cannot touch the successor; the attached collector
// still saves and cleans the actual owner after process retirement.
await copies.collectStopped({ companyId, runId: first.runId });
expect((await fs.stat(first.localRoot)).isDirectory()).toBe(true);
expect((await copies.collectStopped({ companyId, runId: collectorRunId }))?.state).toBe("saved");
expect(await fs.readFile(path.join(root, "last-write.txt"), "utf8")).toBe("saved after failed handoff");
await expect(fs.stat(first.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
}, 10_000);
it("retires loaded instruction edits while allowing ordinary personal-file edits to stay warm", async () => {
const copy = await run({ warm: true });
await fs.writeFile(path.join(copy.localRoot, entryFile), "# New loaded policy\n");
expect((await copies.checkpointWarm({ companyId, runId: copy.runId }))?.state).toBe("warm_saved");
expect(await copies.canReuseWarm(companyId, agentId, copy.runId)).toBe(false);
await copies.collectStopped({ companyId, runId: copy.runId });
const fresh = await run({ warm: true });
expect(fresh.executionRoot).not.toBe(copy.executionRoot);
expect(await fs.readFile(path.join(fresh.localRoot, entryFile), "utf8")).toBe("# New loaded policy\n");
});
it("refreshes after editor changes and preserves concurrent unrelated files", async () => {
const copy = await run({ warm: true });
await fs.writeFile(path.join(copy.localRoot, "mine.txt"), "from run");
await agentFileStore(db).write({ ...target(), path: "board.txt", bytes: Buffer.from("from board"), baseHash: null }, board());
expect((await copies.checkpointWarm({ companyId, runId: copy.runId }))?.state).toBe("warm_saved");
expect(await copies.canReuseWarm(companyId, agentId, copy.runId)).toBe(false);
expect(await fs.readFile(path.join(root, "board.txt"), "utf8")).toBe("from board");
expect(await fs.readFile(path.join(root, "mine.txt"), "utf8")).toBe("from run");
await copies.collectStopped({ companyId, runId: copy.runId });
});
it("keeps an invalid warm checkpoint unsaved and requests stopped collection", async () => {
const copy = await run({ warm: true });
await fs.symlink(root, path.join(copy.localRoot, "escape"));
const failed = await copies.checkpointWarm({ companyId, runId: copy.runId });
expect(failed?.state).toBe("prepared");
expect(failed?.errorCode).toBe("AGENT_FILES_CHECKPOINT_UNSTABLE");
expect(failed?.processStoppedAt).toBeNull();
expect(await copies.canReuseWarm(companyId, agentId, copy.runId)).toBe(false);
await fs.rm(path.join(copy.localRoot, "escape"));
expect((await copies.collectStopped({ companyId, runId: copy.runId }))?.state).toBe("unchanged");
});
it("rejects over-quota warm changes without replacing saved bytes, then permits a clean future run", async () => {
const copy = await run({ warm: true });
const handle = await fs.open(path.join(copy.localRoot, "too-large.bin"), "w");
await handle.truncate(MAX_AGENT_FILE_BYTES + 1); await handle.close();
expect((await copies.checkpointWarm({ companyId, runId: copy.runId }))?.errorCode).toBe("AGENT_FILES_CHECKPOINT_UNSTABLE");
const stopped = (await copies.collectStopped({ companyId, runId: copy.runId }))!;
expect(stopped).toMatchObject({ state: "unavailable", errorCode: "AGENT_FILES_LIMIT_EXCEEDED" });
expect(stopped.receipt?.storageWarning).toContain("Agent storage is full");
expect(await fs.readFile(path.join(root, entryFile), "utf8")).toBe(initial);
await expect(fs.stat(copy.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
expect(await run({ warm: true })).toBeTruthy();
});
it("rechecks authorization even for an unchanged warm checkpoint", async () => {
const copy = await run({ warm: true });
await copies.checkpointWarm({ companyId, runId: copy.runId });
await db.delete(principalPermissionGrants).where(eq(principalPermissionGrants.companyId, companyId));
await db.update(companyMemberships).set({ membershipRole: "member" }).where(eq(companyMemberships.companyId, companyId));
expect((await copies.checkpointWarm({ companyId, runId: copy.runId }))?.errorCode).toBe("AGENT_FILES_CHECKPOINT_UNSTABLE");
expect(await copies.canReuseWarm(companyId, agentId, copy.runId)).toBe(false);
expect((await copies.collectStopped({ companyId, runId: copy.runId }))?.state).toBe("unavailable");
});
it("reclaims a warm remote owner after verified sandbox destruction without claiming unsaved tail bytes", async () => {
const copy = await run({ warm: true });
await fs.writeFile(path.join(copy.localRoot, "saved.txt"), "checkpoint");
await copies.checkpointWarm({ companyId, runId: copy.runId });
const environmentId = randomUUID(), leaseId = randomUUID(), remoteCwd = "/fixture/task";
const lease = { id: leaseId, companyId, environmentId, heartbeatRunId: copy.runId, provider: "daytona", providerLeaseId: "destroyed-warm" };
await db.insert(environments).values({ id: environmentId, name: environmentId, driver: "sandbox" });
await db.insert(environmentLeases).values({ ...lease, status: "released", releasedAt: new Date(), cleanupStatus: "success",
metadata: { remoteExecutionTermination: remoteTerminationReceipt(lease, { providerLeaseId: lease.providerLeaseId, state: "destroyed" }) } });
await db.update(heartbeatRuns).set({ status: "succeeded", runtimeMode: "native" }).where(eq(heartbeatRuns.id, copy.runId));
const saved = (await copies.get(companyId, copy.runId))!;
await db.update(agentInstructionWorkingCopies).set({ location: `remote:${environmentId}`,
executionRoot: path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", agentId, copy.runId),
receipt: { ...saved.receipt, cleanup: { leaseId, remoteCwd } } }).where(eq(agentInstructionWorkingCopies.runId, copy.runId));
copies = agentInstructionWorkingCopyService(db);
await copies.recoverStopped();
const recovered = (await copies.get(companyId, copy.runId))!;
expect(recovered).toMatchObject({ state: "unavailable", errorCode: "AGENT_FILES_FINAL_COLLECTION_UNAVAILABLE" });
expect(recovered.receipt?.cleanupPending).toBe(false);
expect(await fs.readFile(path.join(root, "saved.txt"), "utf8")).toBe("checkpoint");
await expect(fs.stat(copy.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
});
it("uses the workspace transport to restore after destruction of the remote filesystem", async () => {
const remoteCwd = path.join(home, "remote-task");
await fs.mkdir(remoteCwd, { recursive: true });
@@ -496,6 +643,52 @@ describe("persistent agent directories", () => {
await copies.release(companyId, second.runId);
});
it("checkpoints and reuses the remote directory through the real transport without copying unchanged bytes", async () => {
const remoteCwd = path.join(home, "warm-remote-task");
await fs.mkdir(remoteCwd, { recursive: true });
const runner: import("@paperclipai/adapter-utils/command-managed-runtime").CommandManagedRuntimeRunner = {
execute: async input => {
const startedAt = new Date().toISOString();
const env = { ...process.env, ...input.env };
const args = [...(input.args ?? [])];
if (input.stdin != null && (args[0] === "-c" || args[0] === "-lc")) {
env.PAPERCLIP_TEST_STDIN = input.stdin;
args[1] = `printf '%s' "$PAPERCLIP_TEST_STDIN" | (${args[1]})`;
}
try {
const result = await execFile(input.command, args, { cwd: input.cwd, env, timeout: input.timeoutMs, maxBuffer: 32 * 1024 * 1024 });
return { exitCode: 0, signal: null, timedOut: false, stdout: result.stdout, stderr: result.stderr, pid: null, startedAt };
} catch (error) {
const e = error as { code?: number; signal?: NodeJS.Signals; stdout?: string; stderr?: string };
return { exitCode: typeof e.code === "number" ? e.code : 1, signal: e.signal ?? null, timedOut: false, stdout: e.stdout ?? "", stderr: e.stderr ?? "", pid: null, startedAt };
}
},
};
const executionTarget = { kind: "remote" as const, transport: "sandbox" as const, environmentId: randomUUID(), remoteCwd, runner };
const prepare = async (reuseRunId?: string) => {
const runId = randomUUID();
await db.insert(heartbeatRuns).values({ id: runId, companyId, agentId, invocationSource: "on_demand", responsibleUserId: userId });
return (await copies.prepare({ ...target(), runId, cwd: home, target: executionTarget, warm: true, reuseRunId }))!;
};
let copy = await prepare();
const original = copy;
await fs.writeFile(path.join(copy.executionRoot, "image.bin"), Buffer.alloc(32768, 93));
for (let turn = 1; turn <= 3; turn++) {
await fs.appendFile(path.join(copy.executionRoot, "memory.txt"), `turn ${turn}\n`);
const saved = (await copies.checkpointWarm({ companyId, runId: copy.runId, target: executionTarget }))!;
expect(saved, JSON.stringify({ code: saved.errorCode, message: saved.errorMessage })).toMatchObject({ state: "warm_saved", errorCode: null });
if (turn > 1) expect(saved.receipt?.checkpointStats).toMatchObject({ copiedFiles: 1, copiedBytes: turn * 7, hashedBytes: turn * 7 });
expect(await fs.readFile(path.join(root, "memory.txt"), "utf8")).toBe(Array.from({ length: turn }, (_, i) => `turn ${i + 1}\n`).join(""));
expect((await fs.readdir(path.dirname(copy.localRoot))).filter(name => name.startsWith("checkpoint-"))).toEqual([]);
expect((await fs.readdir(path.join(copy.executionRoot, ".paperclip-runtime"))).filter(name => name.startsWith("checkpoint-"))).toEqual([]);
if (turn < 3) { copy = await prepare(copy.runId); expect(copy.executionRoot).toBe(original.executionRoot); }
}
await copies.collectStopped({ companyId, runId: original.runId, target: executionTarget });
expect(await fs.stat(copy.executionRoot)).toBeTruthy();
await copies.collectStopped({ companyId, runId: copy.runId, target: executionTarget });
await expect(fs.stat(copy.executionRoot)).rejects.toMatchObject({ code: "ENOENT" });
});
it("stages SSH agent files at the registered root without a nested task workspace", async () => {
const remoteCwd = path.join(home, "ssh-task");
const exclude = vi.spyOn(executionTargetTools, "runAdapterExecutionTargetShellCommand").mockResolvedValue({ exitCode: 0, signal: null, timedOut: false, stdout: "", stderr: "" });
@@ -0,0 +1,132 @@
import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { captureAgentFiles } from "../services/scripts/agent-file-checkpoint.mjs";
import { captureAgentFileCheckpoint, validateAgentFileCheckpoint } from "../services/agent-file-checkpoints.js";
describe("incremental managed file checkpoints", () => {
let root: string;
beforeEach(async () => { root = await fs.realpath(await fs.mkdtemp(path.join(os.tmpdir(), "agent-checkpoint-"))); await fs.mkdir(path.join(root, "live")); });
afterEach(async () => { vi.restoreAllMocks(); await fs.rm(root, { recursive: true, force: true }); });
const live = () => path.join(root, "live");
it("does not read or copy a large unchanged binary when only a note changes", async () => {
const image = Buffer.alloc(64 * 1024 * 1024, 0x93);
await fs.writeFile(path.join(live(), "image.bin"), image);
await fs.writeFile(path.join(live(), "note.txt"), "one");
const first = await captureAgentFiles(live());
await fs.writeFile(path.join(live(), "note.txt"), "two");
const delta = await captureAgentFileCheckpoint({ runId: "fixture", localRoot: live(), executionRoot: live(), previous: first.manifest });
try {
expect(delta.stats.hashedBytes).toBe(3);
expect(delta.stats.copiedBytes).toBe(3);
expect(delta.stats.copiedFiles).toBe(1);
expect(await fs.readFile(path.join(delta.directory, "files", "note.txt"), "utf8")).toBe("two");
await expect(fs.stat(path.join(delta.directory, "files", "image.bin"))).rejects.toMatchObject({ code: "ENOENT" });
expect(delta.manifest.entries.find(([p]) => p === "image.bin")?.[1].hash).toBe(first.manifest.entries.find(([p]) => p === "image.bin")?.[1].hash);
} finally { await delta.cleanup(); }
});
it("records deletions and empty directories without resending unchanged data", async () => {
await fs.writeFile(path.join(live(), "deleted"), "gone");
await fs.writeFile(path.join(live(), "kept"), "keep");
const first = await captureAgentFiles(live());
await fs.rm(path.join(live(), "deleted"));
await fs.mkdir(path.join(live(), "empty"));
const next = await captureAgentFiles(live(), first.manifest, path.join(root, "delta"));
expect(next.stats.hashedBytes).toBe(0);
expect(next.stats.copiedBytes).toBe(0);
expect(next.manifest.entries.map(([p]) => p)).toEqual(["empty", "kept"]);
});
it("detects same-size writes with a restored mtime", async () => {
const file = path.join(live(), "note");
await fs.writeFile(file, "aaaa");
const before = await fs.stat(file);
const first = await captureAgentFiles(live());
await fs.writeFile(file, "bbbb");
await fs.utimes(file, before.atime, before.mtime);
const next = await captureAgentFiles(live(), first.manifest, path.join(root, "delta"));
expect(next.stats.copiedBytes).toBe(4);
expect(next.manifest.entries[0]![1].hash).not.toBe(first.manifest.entries[0]![1].hash);
});
it.each(["symlink", "hardlink"])("refuses a %s without publishing a manifest", async kind => {
const outside = path.join(root, "outside"); await fs.writeFile(outside, "private");
if (kind === "symlink") await fs.symlink(outside, path.join(live(), "link"));
else await fs.link(outside, path.join(live(), "link"));
await expect(captureAgentFiles(live(), undefined, path.join(root, "delta"))).rejects.toMatchObject({ code: "AGENT_FILES_UNSAFE_PATH" });
await expect(fs.stat(path.join(root, "delta", "checkpoint.json"))).rejects.toMatchObject({ code: "ENOENT" });
});
it("rejects a sparse oversized file before reading its bytes", async () => {
const f = await fs.open(path.join(live(), "large"), "w"); await f.truncate(256 * 1024 * 1024 + 1); await f.close();
await expect(captureAgentFiles(live())).rejects.toMatchObject({ code: "AGENT_FILES_LIMIT_EXCEEDED" });
});
it("validates received contents independently of the remote manifest", async () => {
await fs.writeFile(path.join(live(), "note"), "expected");
const output = path.join(root, "delta");
await captureAgentFiles(live(), undefined, output);
await fs.writeFile(path.join(output, "files", "note"), "tampered");
await expect(validateAgentFileCheckpoint(output, { entries: [] })).rejects.toThrow("payload mismatch");
});
it.each(["../escape", "/absolute", "nested/.paperclip-runtime/secret"])("rejects a forged checkpoint path %s", async name => {
const output = path.join(root, "delta"); await fs.mkdir(output);
await fs.writeFile(path.join(output, "checkpoint.json"), JSON.stringify({ version: 1, entries: [[name, { kind: "dir" }]] }));
await expect(validateAgentFileCheckpoint(output, { entries: [] })).rejects.toThrow("AGENT_FILES_UNSAFE_PATH");
});
it.skipIf(process.platform !== "darwin")("accepts root-owned macOS temp aliases while rejecting user symlink roots", async () => {
await fs.writeFile(path.join(live(), "note"), "one");
const alias = live().replace(/^\/private(?=\/(var|tmp)\/)/, "");
expect((await captureAgentFiles(alias)).manifest.entries).toHaveLength(1);
await fs.symlink(live(), path.join(root, "user-link"));
await expect(captureAgentFiles(path.join(root, "user-link"))).rejects.toThrow("AGENT_FILES_UNSAFE_PATH");
});
it("rejects a file growing while its bytes are streamed", async () => {
const filename = path.join(live(), "writer");
await fs.writeFile(filename, Buffer.alloc(1024 * 1024, 1));
const original = fs.open.bind(fs);
let inject = true;
vi.spyOn(fs, "open").mockImplementation(async (...args) => {
const handle = await original(...args);
if (String(args[0]) === filename && inject) {
inject = false;
const stream = handle.createReadStream.bind(handle);
handle.createReadStream = ((options: Parameters<typeof stream>[0]) => (async function* () {
let first = true;
for await (const chunk of stream(options)) {
if (first) { first = false; await fs.appendFile(filename, Buffer.alloc(1024 * 1024, 2)); }
yield chunk;
}
})()) as typeof handle.createReadStream;
}
return handle;
});
await expect(captureAgentFiles(live(), undefined, path.join(root, "delta"))).rejects.toMatchObject({ code: "AGENT_FILES_CHANGED_DURING_CHECKPOINT" });
await expect(fs.stat(path.join(root, "delta", "checkpoint.json"))).rejects.toMatchObject({ code: "ENOENT" });
});
it("detects mode changes and renames without retransferring unrelated files", async () => {
await fs.writeFile(path.join(live(), "script"), "execute");
await fs.writeFile(path.join(live(), "rename"), "move");
const first = await captureAgentFiles(live());
await fs.chmod(path.join(live(), "script"), 0o755);
await fs.rename(path.join(live(), "rename"), path.join(live(), "renamed"));
const next = await captureAgentFiles(live(), first.manifest, path.join(root, "delta"));
expect(next.stats.copiedFiles).toBe(2);
expect(next.manifest.entries.map(([name]) => name)).toEqual(["renamed", "script"]);
expect((await fs.stat(path.join(root, "delta", "files", "script"))).mode & 0o777).toBe(0o755);
});
it("excludes only the remote transport's reserved root, never user file paths", async () => {
await fs.mkdir(path.join(live(), ".paperclip-runtime"));
await fs.symlink("/missing", path.join(live(), ".paperclip-runtime", "transport-only"));
await fs.writeFile(path.join(live(), "note"), "saved");
await expect(captureAgentFiles(live())).rejects.toThrow("AGENT_FILES_UNSAFE_PATH");
const captured = await captureAgentFiles(live(), undefined, undefined, true, true);
expect(captured.manifest.entries.map(([name]) => name)).toEqual(["note"]);
});
it("does not acknowledge a write made after the captured generation", async () => {
await fs.writeFile(path.join(live(), "note"), "one");
const first = await captureAgentFiles(live(), undefined, path.join(root, "first"));
await fs.writeFile(path.join(live(), "note"), "two");
const second = await captureAgentFiles(live(), first.manifest, path.join(root, "second"));
expect(second.stats.copiedFiles).toBe(1);
expect(await fs.readFile(path.join(root, "first", "files", "note"), "utf8")).toBe("one");
expect(await fs.readFile(path.join(root, "second", "files", "note"), "utf8")).toBe("two");
});
});
@@ -14,8 +14,11 @@ import type { AuthorizationActor } from "./authorization.js";
import type { EnvironmentRuntimeService } from "./environment-runtime.js";
import type { Environment, EnvironmentLease } from "@paperclipai/shared";
import { hasRemoteTerminationReceipt } from "./remote-execution-termination.js";
import { cachedAgentFileManifest, captureAgentFileCheckpoint, checkpointBaseline, checkpointSnapshot, type AgentFileManifest } from "./agent-file-checkpoints.js";
import { logger } from "../middleware/logger.js";
type Copy = typeof copies.$inferSelect;
export class AgentDirectoryReuseInvalidatedError extends Error {}
const completed = new Set(["saved", "unchanged", "resolved", "unavailable"]);
const transports = new Map<string, PreparedAdapterExecutionTargetRuntime>();
const key = (row: Pick<Copy, "companyId" | "runId">) => `${row.companyId}:${row.runId}`;
@@ -32,6 +35,19 @@ function actor(row: Copy): AuthorizationActor {
}
export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string, runId: string) => Promise<Copy | null>, patch: (row: Copy, values: Partial<typeof copies.$inferInsert>) => Promise<Copy>, environmentRuntime?: EnvironmentRuntimeService) {
const store = agentFileStore(db);
const ownerFile = (row: Copy) => path.join(path.dirname(row.localRoot), "owner.json");
async function owns(row: Copy) {
if (row.receipt?.warm !== true) return true;
return await fs.readFile(ownerFile(row), "utf8").then(s => JSON.parse(s).runId === row.runId).catch(error => error.code === "ENOENT" && row.processStoppedAt !== null);
}
async function canReuse(companyId: string, agentId: string, runId: string) {
const row = await get(companyId, runId);
if (!row || row.agentId !== agentId || row.state !== "warm_saved" || row.errorCode || !await owns(row)) return false;
const entry = baseline(row).entries.get(row.entryFile);
if (entry?.kind !== "file" || entry.hash !== row.receipt?.sessionEntryHash) return false;
return store.locked(companyId, agentId, actor(row), false, async (_tx, _agent, root) =>
directorySnapshotSha256(checkpointSnapshot(await cachedAgentFileManifest(root))) === row.baseHash);
}
async function transport(row: Copy, target: AdapterExecutionTarget, recovering: boolean) {
if (target.kind !== "remote" || row.location !== `remote:${target.environmentId ?? ""}`) throw new Error("Agent directory environment changed");
if (recovering && target.transport === "ssh") throw new Error("The original SSH collector is unavailable");
@@ -50,12 +66,50 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
workspaceBaseline: baseline(row), workspaceGitSnapshot: null, workspaceFileMode: "all",
workspaceExclude: [".paperclip-runtime", ".paperclip-runtime/**"] });
}
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string }) {
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; warm?: boolean; reuseRunId?: string; onWarmHandoff?: (copy: Copy) => void }) {
const [agent] = await db.select().from(agents).where(and(eq(agents.id, input.agentId), eq(agents.companyId, input.companyId)));
if (!agent) throw notFound("Agent not found");
if (agentInstructionsBundleMode(agent) !== "managed") return null;
const bound = await resolveInstructionActor(db, { type: "agent", companyId: input.companyId, agentId: input.agentId, runId: input.runId });
const root = resolveManagedInstructionsRoot(agent);
if (input.reuseRunId) {
const prior = await get(input.companyId, input.reuseRunId);
if (!prior || prior.agentId !== input.agentId) throw conflict("Managed warm directory owner changed");
return serial(prior, async current => {
if (current.state !== "warm_saved" || current.errorCode || !await owns(current) || current.location !== (input.target?.kind === "remote" ? `remote:${input.target.environmentId ?? ""}` : "local")) throw conflict("Managed warm directory is not reusable");
const receipt = { ...current.receipt, cleanup: input.target?.kind === "remote" ? { leaseId: input.target.leaseId ?? null, remoteCwd: input.target.remoteCwd } : undefined };
// Insert before taking the agent-row lock: this foreign key needs an
// independent transaction, and must survive post-handoff failures.
const [successor] = await db.insert(copies).values({ runId: input.runId, companyId: input.companyId, agentId: input.agentId, responsibleUserId: bound.onBehalfOfUserId!,
entryFile: current.entryFile, baseHash: current.baseHash, localRoot: current.localRoot, executionRoot: current.executionRoot, location: current.location, state: "prepared", receipt }).returning();
let handedOff = false;
try { return await store.locked(input.companyId, input.agentId, bound, false, async (_tx, lockedAgent, canonical) => {
// The reservation's earlier check is only a hint. Serialize this last
// validation and ownership transfer with every canonical writer.
if (current.entryFile !== deriveBundleState(lockedAgent).entryFile ||
directorySnapshotSha256(checkpointSnapshot(await cachedAgentFileManifest(canonical))) !== current.baseHash) {
throw new AgentDirectoryReuseInvalidatedError("Managed agent files changed before warm directory handoff");
}
// The controller-owned marker fences old callbacks. Attach the new
// collector immediately after it moves, before fallible bookkeeping.
const temporary = `${ownerFile(current)}.next`;
await fs.writeFile(temporary, JSON.stringify({ runId: input.runId }));
await fs.rename(temporary, ownerFile(current));
handedOff = true;
input.onWarmHandoff?.(successor!);
const runtime = transports.get(key(current));
if (runtime) transports.set(key(successor!), runtime);
transports.delete(key(current));
await patch(current, { state: "superseded", receipt: { schema: AGENT_FILES_CONTRACT, warm: true, successorRunId: input.runId } });
return successor!;
}); } catch (error) {
// No process used this proposed owner. Remove it after the canonical
// lock releases so the same run can restore a fresh directory.
if (!handedOff) await db.delete(copies).where(and(eq(copies.companyId, input.companyId), eq(copies.runId, input.runId)));
throw error;
}
});
}
const localRoot = path.join(path.dirname(root), "file-sync", "runs", input.runId, "live");
// Local copies live outside the task cwd. Remote copies use the reserved,
// excluded runtime area inside the provider's confined workspace. They are
@@ -74,7 +128,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
// leaves a recoverable preparing receipt instead of an orphan directory.
const preparing = { entryFile: deriveBundleState(agent).entryFile, baseRevisionId: null, baseHash: "preparing",
localRoot, executionRoot, location, state: "preparing", candidateBase64: null, candidateHash: null,
receipt: { schema: AGENT_FILES_CONTRACT, ...(input.target?.kind === "remote"
receipt: { schema: AGENT_FILES_CONTRACT, ...(input.warm ? { warm: true, directoryRunId: input.runId } : {}), ...(input.target?.kind === "remote"
? { cleanup: { leaseId: input.target.leaseId ?? null, remoteCwd: input.target.remoteCwd } } : {}) }, processStoppedAt: null, attempts: 0,
nextAttemptAt: null, errorCode: null, errorMessage: null };
if (row) row = await patch(row, preparing);
@@ -102,12 +156,15 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
await fs.rm(path.dirname(localRoot), { recursive: true, force: true });
throw error;
});
const loadedEntry = snapshot.entries.get(deriveBundleState(agent).entryFile);
const values = { entryFile: deriveBundleState(agent).entryFile, baseRevisionId: null, baseHash: directorySnapshotSha256(snapshot),
localRoot, executionRoot, location, state: "preparing", candidateBase64: null, candidateHash: null,
receipt: { ...preparing.receipt, baseline: serializeDirectorySnapshot(snapshot), storageWarning },
receipt: { ...preparing.receipt, baseline: serializeDirectorySnapshot(snapshot), storageWarning,
...(input.warm ? { sessionEntryHash: loadedEntry?.kind === "file" ? loadedEntry.hash : null } : {}) },
processStoppedAt: null, attempts: 0, nextAttemptAt: null, errorCode: null, errorMessage: null };
try {
row = await patch(row, values);
if (input.warm) await fs.writeFile(ownerFile(row), JSON.stringify({ runId: row.runId }));
} catch (error) {
await fs.rm(path.dirname(localRoot), { recursive: true, force: true });
throw error;
@@ -150,7 +207,49 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
// ordinary file edits do not change the session configuration digest.
return true;
}
async function checkpoint(row: Copy, target?: AdapterExecutionTarget | null, stopped = false): Promise<Copy> {
if (!await owns(row) || row.state === "superseded") return row;
let failure: unknown;
for (let attempt = 0; attempt < 3; attempt++) {
let captured: Awaited<ReturnType<typeof captureAgentFileCheckpoint>> | undefined;
try {
const previous = row.receipt?.checkpointManifest as AgentFileManifest | undefined ?? checkpointBaseline(baseline(row));
captured = await captureAgentFileCheckpoint({ runId: row.runId, localRoot: row.localRoot, executionRoot: row.executionRoot, previous, target });
const snapshot = checkpointSnapshot(captured.manifest);
const nextHash = directorySnapshotSha256(snapshot);
const changed = nextHash !== row.baseHash;
const { storageWarning } = changed
? await store.apply({ companyId: row.companyId, agentId: row.agentId, sourceDir: path.join(captured.directory, "files"), baseline: baseline(row), checkpoint: captured.manifest }, actor(row))
: await store.locked(row.companyId, row.agentId, actor(row), true, async () => ({
storageWarning: typeof row.receipt?.storageWarning === "string" ? row.receipt.storageWarning : null,
}));
row = await patch(row, { state: stopped ? (changed ? "saved" : "unchanged") : "warm_saved", baseHash: nextHash,
candidateHash: nextHash, errorCode: null, errorMessage: null, attempts: 0, nextAttemptAt: null,
...(stopped ? { processStoppedAt: new Date() } : {}),
receipt: { ...row.receipt, baseline: serializeDirectorySnapshot(snapshot), checkpointManifest: captured.manifest,
checkpointState: changed ? "saved" : "unchanged", checkpointStats: captured.stats, storageWarning } });
return row;
} catch (error) {
failure = error;
logger.warn({ err: error, runId: row.runId, attempt: attempt + 1, stopped }, "Agent file checkpoint failed");
}
finally {
await captured?.cleanup().catch(error => logger.warn({ err: error, runId: row.runId }, "Agent checkpoint scratch cleanup deferred to session retirement"));
}
}
if (!stopped) return patch(row, { errorCode: "AGENT_FILES_CHECKPOINT_UNSTABLE", errorMessage: "Incremental checkpoint could not be validated; collecting after provider stop.", nextAttemptAt: null });
const storageLimit = failure instanceof AgentFileLimitError || String(failure).includes("LIMIT_EXCEEDED");
return patch(row, { state: "unavailable", processStoppedAt: new Date(), errorCode: storageLimit ? "AGENT_FILES_LIMIT_EXCEEDED" : "AGENT_FILES_SAVE_FAILED",
errorMessage: "Agent-file synchronization failed after provider stop. No successful save is claimed.", nextAttemptAt: null,
receipt: { ...row.receipt, storageWarning: storageLimit ? agentStorageWarning(failure instanceof AgentFileLimitError ? failure.message : "Agent folder exceeds a storage limit") : null } });
}
async function collectStopped(row: Copy, target?: AdapterExecutionTarget | null) {
if (!await owns(row) || row.state === "superseded") return row;
if (row.receipt?.warm === true && !completed.has(row.state)) {
row = await checkpoint(row, target, true);
await release(row, target);
return (await get(row.companyId, row.runId))!;
}
if (completed.has(row.state)) { await release(row); return (await get(row.companyId, row.runId))!; }
row = await patch(row, { processStoppedAt: row.processStoppedAt ?? new Date() });
// The stopped working copy is the only temporary tree. Retry transient I/O
@@ -188,6 +287,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
return (await get(row.companyId, row.runId))!;
}
async function release(row: Copy, target?: AdapterExecutionTarget | null) {
if (row.state === "superseded" || !await owns(row) || row.receipt?.warm === true && !row.processStoppedAt) return;
// Collection and environment teardown can both release the same copy.
// A compact successful cleanup receipt is final, even after restart.
if (completed.has(row.state) && row.processStoppedAt && row.receipt?.cleanupPending === false) return;
@@ -195,9 +295,9 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
transports.delete(key(row));
let cleanupPending = false;
await runtime?.cleanupWorkspaceSnapshot?.().catch(() => { cleanupPending = true; });
const cleanupTarget = runtime?.target ?? target;
const cleanupTarget = target ?? runtime?.target;
if (row.processStoppedAt && completed.has(row.state) && cleanupTarget?.kind === "remote") {
const expected = path.posix.join(cleanupTarget.remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, row.runId);
const expected = path.posix.join(cleanupTarget.remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, String(row.receipt?.directoryRunId ?? row.runId));
if (row.executionRoot !== expected) throw new Error("Agent directory cleanup path changed");
const quoted = `'${expected.replaceAll("'", `'"'"'`)}'`;
const remoteCleanupFailed = await runAdapterExecutionTargetShellCommand(row.runId, cleanupTarget, `rm -rf -- ${quoted}`,
@@ -216,6 +316,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
return;
}
await patch(row, { receipt: { schema: AGENT_FILES_CONTRACT, state: row.state, appliedCandidateHash: row.candidateHash, storageWarning: row.receipt?.storageWarning ?? null, cleanupPending,
...(row.receipt?.directoryRunId ? { directoryRunId: row.receipt.directoryRunId } : {}),
...(cleanupPending ? { cleanup: row.receipt?.cleanup } : {}) },
nextAttemptAt: cleanupPending ? new Date(Date.now() + 30_000) : null });
}
@@ -237,7 +338,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
const [environment] = await db.select().from(environments).where(eq(environments.id, lease.environmentId));
if (!environment) return false;
const remoteCwd = cleanup?.remoteCwd ?? lease.metadata?.remoteCwd;
if (typeof remoteCwd !== "string" || row.executionRoot !== path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, row.runId)) return false;
if (typeof remoteCwd !== "string" || row.executionRoot !== path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, String(row.receipt?.directoryRunId ?? row.runId))) return false;
const result = await environmentRuntime.execute({ environment: environment as Environment, lease: lease as EnvironmentLease, command: "rm", args: ["-rf", "--", row.executionRoot],
cwd: remoteCwd, env: {}, timeoutMs: 15_000, bypassSession: true });
return result.exitCode === 0 && !result.timedOut;
@@ -251,7 +352,8 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
return fn(current);
});
}
return { prepare, hasChanges,
return { prepare, hasChanges, canReuse,
checkpointWarm: (row: Copy, target?: AdapterExecutionTarget | null) => serial(row, current => current.receipt?.warm === true ? checkpoint(current, target) : Promise.resolve(current)),
collectStopped: (row: Copy, target?: AdapterExecutionTarget | null) => serial(row, current => collectStopped(current, target)),
release: (row: Copy) => serial(row, release),
};
@@ -0,0 +1,115 @@
import fs from "node:fs/promises";
import path from "node:path";
import { createHash, randomUUID } from "node:crypto";
import { prepareAdapterExecutionTargetRuntime, runAdapterExecutionTargetShellCommand, type AdapterExecutionTarget } from "@paperclipai/adapter-utils/execution-target";
import { type DirectorySnapshot, type SnapshotEntry } from "@paperclipai/adapter-utils/workspace-restore-merge";
import { captureAgentFiles, checkpointPath, type AgentFileManifest, type AgentFileCheckpointStats } from "./scripts/agent-file-checkpoint.mjs";
import { inspectAgentFile, MAX_AGENT_DIRECTORY_BYTES, MAX_AGENT_DIRECTORY_ENTRIES, MAX_AGENT_FILE_BYTES } from "./agent-file-store.js";
export type { AgentFileManifest, AgentFileCheckpointStats };
const quote = (text: string) => `'${text.replaceAll("'", `'"'"'`)}'`;
export function checkpointSnapshot(manifest: AgentFileManifest): DirectorySnapshot {
return { exclude: [], entries: new Map(manifest.entries.map(([name, entry]) => [name,
entry.kind === "dir" ? { kind: "dir" } : { kind: "file", mode: entry.mode!, hash: entry.hash! }] as [string, SnapshotEntry])) };
}
export function checkpointBaseline(snapshot: DirectorySnapshot): AgentFileManifest {
return { version: 1, entries: [...snapshot.entries].map(([name, entry]) => {
if (entry.kind === "symlink") throw new Error("Agent file baseline contains a link");
return [name, entry];
}) };
}
export async function cachedAgentFileManifest(root: string): Promise<AgentFileManifest> {
const filename = path.join(path.dirname(root), "file-sync", "canonical-manifest.json");
const previous = await fs.readFile(filename, "utf8").then(s => JSON.parse(s) as AgentFileManifest).catch(() => undefined);
const { manifest } = await captureAgentFiles(root, previous, undefined, false);
await fs.mkdir(path.dirname(filename), { recursive: true });
const temporary = `${filename}.${randomUUID()}`;
try { await fs.writeFile(temporary, JSON.stringify(manifest)); await fs.rename(temporary, filename); }
finally { await fs.rm(temporary, { force: true }); }
return manifest;
}
/** Host validation does not trust remote hashes, paths, claimed sizes or quotas.
* A missing payload must match the last committed manifest exactly. */
export async function validateAgentFileCheckpoint(directory: string, previous: AgentFileManifest): Promise<AgentFileManifest> {
const filename = path.join(directory, "checkpoint.json");
const info = await fs.lstat(filename);
if (!info.isFile() || info.nlink !== 1 || info.size > 64 * 1024 * 1024) throw new Error("Invalid agent file checkpoint manifest");
const bytes = await fs.readFile(filename);
if (bytes.length > 64 * 1024 * 1024) throw new Error("Agent file checkpoint manifest is too large");
const manifest = JSON.parse(bytes.toString("utf8")) as AgentFileManifest;
if (manifest.version !== 1 || !Array.isArray(manifest.entries) || manifest.entries.length > MAX_AGENT_DIRECTORY_ENTRIES) throw new Error("Invalid agent file checkpoint");
const old = new Map(previous.entries), names = new Set<string>();
let total = 0;
for (const item of manifest.entries) {
if (!Array.isArray(item) || item.length !== 2) throw new Error("Invalid agent file checkpoint entry");
const [name, entry] = item;
checkpointPath(name);
if (names.has(name) || !entry || !["dir", "file"].includes(entry.kind)) throw new Error("Invalid agent file checkpoint entry");
names.add(name);
if (entry.kind === "dir") continue;
if (!Number.isSafeInteger(entry.size) || entry.size! < 0 || entry.size! > MAX_AGENT_FILE_BYTES ||
!Number.isSafeInteger(entry.mode) || !/^[a-f0-9]{64}$/.test(entry.hash ?? "") || typeof entry.stamp !== "string") throw new Error("Invalid agent file checkpoint metadata");
total += entry.size!;
if (total > MAX_AGENT_DIRECTORY_BYTES) throw new Error("Agent file checkpoint exceeds quota");
const before = old.get(name);
if (before?.kind === "file" && before.hash === entry.hash && before.mode === entry.mode) {
if (before.size !== undefined && before.size !== entry.size) throw new Error("Agent file checkpoint size changed without content");
} else {
const file = await inspectAgentFile(path.join(directory, "files"), name, 0);
if (!file || file.hash !== entry.hash || file.size !== entry.size) throw new Error("Agent file checkpoint payload mismatch");
}
}
const entries = new Map(manifest.entries);
for (const name of names) for (let parent = path.posix.dirname(name); parent !== "."; parent = path.posix.dirname(parent)) {
if (entries.get(parent)?.kind !== "dir") throw new Error("Agent file checkpoint parent is missing");
}
return { ...manifest, totalBytes: total };
}
export async function captureAgentFileCheckpoint(input: {
runId: string; localRoot: string; executionRoot: string; previous: AgentFileManifest; target?: AdapterExecutionTarget | null;
}): Promise<{ directory: string; manifest: AgentFileManifest; stats: AgentFileCheckpointStats; cleanup: () => Promise<void> }> {
const directory = path.join(path.dirname(input.localRoot), `checkpoint-${randomUUID()}`);
const target = input.target?.kind === "remote" ? input.target : null;
// Scratch lives inside this session's reserved runtime area, so retirement
// reclaims it even if a transient transport failure interrupted cleanup.
const remoteDirectory = target ? path.posix.join(input.executionRoot, ".paperclip-runtime", path.basename(directory)) : null;
let restore: Awaited<ReturnType<typeof prepareAdapterExecutionTargetRuntime>> | undefined;
let cacheRuntime: Awaited<ReturnType<typeof prepareAdapterExecutionTargetRuntime>> | undefined;
const cleanup = async () => {
try {
await restore?.cleanupWorkspaceSnapshot?.();
} finally {
await fs.rm(directory, { recursive: true, force: true });
if (target && remoteDirectory) {
const paths = [...new Set([remoteDirectory, cacheRuntime?.runtimeRootDir, restore?.runtimeRootDir].filter((p): p is string => Boolean(p)))];
const result = await runAdapterExecutionTargetShellCommand(input.runId, target, `rm -rf -- ${paths.map(quote).join(" ")}`, { cwd: target.remoteCwd, env: {}, timeoutSec: 15 });
if (result.exitCode !== 0 || result.timedOut) throw new Error("Agent checkpoint cleanup failed");
}
}
};
try {
let stats: AgentFileCheckpointStats;
if (!target) ({ stats } = await captureAgentFiles(input.localRoot, input.previous, directory));
else {
const cacheDir = path.join(directory, "cache");
await fs.mkdir(cacheDir, { recursive: true });
const cache = JSON.stringify(input.previous);
await fs.writeFile(path.join(cacheDir, "manifest.json"), cache);
const prepared = cacheRuntime = await prepareAdapterExecutionTargetRuntime({ target, runId: input.runId, adapterKey: `agent-checkpoint-${path.basename(directory)}`,
workspaceLocalDir: directory, workspaceRemoteDir: remoteDirectory!, syncWorkspace: false, assets: [{ key: "cache", localDir: cacheDir, followSymlinks: false }] });
const script = await fs.readFile(new URL("./scripts/agent-file-checkpoint.mjs", import.meta.url), "utf8");
const options = { root: input.executionRoot, output: remoteDirectory, cacheFile: path.posix.join(prepared.assetDirs.cache!, "manifest.json"), cacheHash: createHash("sha256").update(cache).digest("hex") };
const result = await runAdapterExecutionTargetShellCommand(input.runId, target, `node --input-type=module -e ${quote(script)} -- --checkpoint ${quote(JSON.stringify(options))}`,
{ cwd: target.remoteCwd, env: {}, timeoutSec: 120 });
if (result.exitCode !== 0 || result.timedOut) throw new Error(`Agent file checkpoint failed: ${result.stderr.trim()}`);
stats = JSON.parse(result.stdout.trim()) as AgentFileCheckpointStats;
await fs.rm(cacheDir, { recursive: true, force: true });
restore = await prepareAdapterExecutionTargetRuntime({ target, runId: input.runId, adapterKey: `agent-checkpoint-download-${path.basename(directory)}`, workspaceLocalDir: directory,
workspaceRemoteDir: remoteDirectory!, syncWorkspace: true, workspaceInboundMode: "adopt_remote", workspaceBaseline: { exclude: [], entries: new Map() }, workspaceGitSnapshot: null, workspaceFileMode: "all" });
await restore.restoreWorkspace();
}
return { directory, manifest: await validateAgentFileCheckpoint(directory, input.previous), stats, cleanup };
} catch (error) { await cleanup().catch(() => undefined); throw error; }
}
+13 -7
View File
@@ -1,5 +1,6 @@
import { createHash } from "node:crypto";
import fs from "node:fs/promises";
import { cachedAgentFileManifest, checkpointSnapshot, type AgentFileManifest } from "./agent-file-checkpoints.js";
import { constants } from "node:fs";
import { Readable } from "node:stream";
import path from "node:path";
@@ -219,16 +220,20 @@ export function agentFileStore(db: Db) {
await audit(tx, agent, bound, { path: relative, contentHash: incomingHash });
return { contentHash: incomingHash, changed: true };
}),
apply: (input: { companyId: string; agentId: string; sourceDir: string; baseline: DirectorySnapshot }, actor: AuthorizationActor) =>
apply: (input: { companyId: string; agentId: string; sourceDir: string; baseline: DirectorySnapshot; checkpoint?: AgentFileManifest }, actor: AuthorizationActor) =>
locked(input.companyId, input.agentId, actor, true, async (tx, agent, root, bound) => {
const incoming = await snapshotAgentFiles(input.sourceDir);
const entry = await readInstructionBytes(input.sourceDir, deriveBundleState(agent).entryFile);
if (entry === null) throw unprocessable("The configured instruction entry cannot be deleted");
instructionBytes(entry);
const incoming = input.checkpoint ? checkpointSnapshot(input.checkpoint) : await snapshotAgentFiles(input.sourceDir);
const entryPath = deriveBundleState(agent).entryFile;
if (incoming.entries.get(entryPath)?.kind !== "file") throw unprocessable("The configured instruction entry cannot be deleted");
if (!input.checkpoint || JSON.stringify(incoming.entries.get(entryPath)) !== JSON.stringify(input.baseline.entries.get(entryPath))) {
const entry = await readInstructionBytes(input.sourceDir, entryPath);
if (entry === null) throw unprocessable("The configured instruction entry cannot be deleted");
instructionBytes(entry);
}
// Previously saved/imported bytes must remain readable, including when
// over quota. Enforce limits on the incoming and resulting tree so a
// run can delete files to recover instead of being locked out forever.
const current = await snapshotAgentFiles(root, false);
const current = input.checkpoint ? checkpointSnapshot(await cachedAgentFileManifest(root)) : await snapshotAgentFiles(root, false);
// Rebase only the run's changed paths onto the current tree. This makes
// same-file edits/deletions last-sync-wins while untouched files retain
// changes from other runs. Ancestors may need recreating after a writer
@@ -274,7 +279,8 @@ export function agentFileStore(db: Db) {
total += size;
}
assertDirectorySize(total, finalEntries.size);
await mergeDirectoryWithBaseline({ ...input, baseline: applyBaseline, targetDir: root });
await mergeDirectoryWithBaseline({ ...input, baseline: applyBaseline, targetDir: root,
...(input.checkpoint ? { snapshots: { source: incoming, current } } : {}) });
await audit(tx, agent, bound, { sourceRunId: actor.runId, contract: AGENT_FILES_CONTRACT });
return { storageWarning: storageWarning(total, finalEntries.size, fullFile) };
}),
@@ -33,7 +33,7 @@ const liveTargets = new Map<string, AdapterExecutionTarget>();
const targetKey = (companyId: string, runId: string) => `${companyId}:${runId}`;
export function instructionWorkingCopyGuidance(copy: Pick<Copy, "executionRoot" | "entryFile" | "receipt">) {
if (isAgentDirectoryCopy(copy)) return `Your persistent agent directory is ${copy.executionRoot} (AGENT_HOME). Your instruction entry is ${copy.executionRoot}/${copy.entryFile}. Read and write your own files and subfolders there. This directory belongs to this agent across tasks and sessions; task files belong in the task working directory. Paperclip restores this directory before execution and saves changes after the provider stops. Regular files, including binary files, persist; symlinks and special files are unsupported. Check the agent-files save receipt before claiming persistence; Only files you change or delete are synchronized. If another run changes the same file, the last completed synchronization wins. Temporary copies are removed after synchronization; there is no per-run file history. Storage allows 256 MiB per file, 2 GiB total, and 100,000 entries; the instruction entry must remain UTF-8 and at most 1 MiB. Reaching a storage limit never prevents this or future tasks from running. Remove or shrink files to free space; changes that exceed the limits will not be saved.${typeof copy.receipt?.storageWarning === "string" ? `\n\n${copy.receipt.storageWarning}` : ""}`;
if (isAgentDirectoryCopy(copy)) return `Your persistent agent directory is ${copy.executionRoot} (AGENT_HOME). Your instruction entry is ${copy.executionRoot}/${copy.entryFile}. Read and write your own files and subfolders there. This directory belongs to this agent across tasks and sessions; task files belong in the task working directory. Paperclip restores this directory before execution and saves validated changes at turn boundaries. A warm native Codex session keeps the same writable directory between turns. Other sessions collect after the provider stops. Regular files, including binary files, persist; symlinks and special files are unsupported. Check the agent-files save receipt before claiming persistence; Only files you change or delete are synchronized. If another run changes the same file, the last completed synchronization wins. Temporary copies are removed when the owning session stops; there is no per-run file history. Storage allows 256 MiB per file, 2 GiB total, and 100,000 entries; the instruction entry must remain UTF-8 and at most 1 MiB. Reaching a storage limit never prevents this or future tasks from running. Remove or shrink files to free space; changes that exceed the limits will not be saved.${typeof copy.receipt?.storageWarning === "string" ? `\n\n${copy.receipt.storageWarning}` : ""}`;
return `Your editable agent instruction file is ${copy.executionRoot}/${copy.entryFile}. Edit this registered private copy normally. After this run stops, Paperclip saves changed content as a persistent revision if your responsible user still has permission and the baseline has not changed. Check the run's instruction-save receipt before claiming persistence. Use read_agent_instructions, update_agent_instructions, get_agent_instruction_history, and restore_agent_instructions for immediate saves and history. Read first and pin the returned revision. Preserve conflicts; never silently retry against a newer head. Repository instructions, skills, and the loaded prompt are separate and are not collected.`;
}
@@ -71,7 +71,7 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
return { type: "agent", companyId: row.companyId, agentId: row.agentId, runId: row.runId, onBehalfOfUserId: row.responsibleUserId };
}
const directories = agentDirectoryWorkingCopyService(db, get, patch, options.environmentRuntime);
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; legacy?: boolean }) {
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; legacy?: boolean; warm?: boolean; reuseRunId?: string; onWarmHandoff?: (copy: Copy) => void }) {
const existing = await get(input.companyId, input.runId);
if (isAgentDirectoryCopy(existing) || (!existing && !input.legacy)) return directories.prepare(input);
let refreshStoppedCopy = false;
@@ -208,6 +208,10 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
return true;
} catch { return true; }
}
async function checkpointWarm(input: { companyId: string; runId: string; target?: AdapterExecutionTarget | null }) {
const row = await get(input.companyId, input.runId);
return row && isAgentDirectoryCopy(row) ? directories.checkpointWarm(row, input.target) : row;
}
async function acknowledgeExplicitSave(input: { companyId: string; agentId: string; runId: string; entryFile: string; revisionId: string; contentHash: string }) {
const row = await get(input.companyId, input.runId);
@@ -294,10 +298,11 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
async function recoverStopped() {
const pending = await db.select({ copy: copies, runtimeMode: heartbeatRuns.runtimeMode }).from(copies)
.innerJoin(heartbeatRuns, and(eq(heartbeatRuns.companyId, copies.companyId), eq(heartbeatRuns.id, copies.runId)))
.where(and(or(inArray(copies.state, ["prepared", "pending_collection"]),
.where(and(or(inArray(copies.state, ["prepared", "pending_collection", "warm_saved"]),
and(eq(copies.state, "preparing"), sql`${copies.receipt}->>'schema' = 'paperclip.agent-files.v1'`)),
inArray(heartbeatRuns.status, ["succeeded", "failed", "cancelled", "timed_out", "interrupted"]),
lte(copies.attempts, MAX_COLLECTION_ATTEMPTS - 1))).limit(20);
or(isNull(copies.nextAttemptAt), lte(copies.nextAttemptAt, new Date())),
lte(copies.attempts, MAX_COLLECTION_ATTEMPTS - 1))).orderBy(asc(copies.updatedAt)).limit(20);
for (const { copy: row, runtimeMode } of pending) {
if (row.state === "preparing" && isAgentDirectoryCopy(row)) {
// A provider cannot launch until preparation records "prepared". With
@@ -309,13 +314,20 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
const result = await collectStopped({ companyId: row.companyId, runId: row.runId });
if (result && isAgentDirectoryCopy(result) && completed.has(result.state)) await directories.release(result);
} else if (row.location !== "local" && await remoteExecutionHasStopped(db, row.companyId, row.runId)) {
const unavailable = await reportUnavailable(row.companyId, row.runId);
const unavailable = row.receipt?.warm === true
? await patch(row, { state: "unavailable", errorCode: "AGENT_FILES_FINAL_COLLECTION_UNAVAILABLE",
errorMessage: "The last completed checkpoint remains saved. Files changed afterwards could not be collected after remote termination.", nextAttemptAt: null })
: await reportUnavailable(row.companyId, row.runId);
if (unavailable && isAgentDirectoryCopy(unavailable)) {
await directories.release(await patch(unavailable, { processStoppedAt: unavailable.processStoppedAt ?? new Date() }));
}
} else if (row.state === "prepared") {
await patch(row, { state: "pending_collection", errorCode: "INSTRUCTION_STOP_UNCONFIRMED",
errorMessage: "The provider's stop has not been confirmed. Instruction collection is pending; no save is claimed.", nextAttemptAt: null });
} else if (row.state === "warm_saved") {
// Live retained sessions must not occupy every batch and starve stopped
// copies from other agents. This does not permit reads without stop proof.
await patch(row, { nextAttemptAt: new Date(Date.now() + 30_000) });
}
}
return pending.length;
@@ -345,12 +357,12 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
* read a possibly live provider. Already captured candidates remain intact. */
async function reportUnavailable(companyId: string, runId: string) {
const row = await get(companyId, runId);
if (!row || completed.has(row.state) || row.state === "unchanged_turn" || row.candidateBase64 !== null || (isAgentDirectoryCopy(row) && row.candidateHash !== null) || row.state === "conflict") return row;
if (!row || completed.has(row.state) || ["unchanged_turn", "warm_saved", "superseded"].includes(row.state) || row.candidateBase64 !== null || (isAgentDirectoryCopy(row) && row.candidateHash !== null) || row.state === "conflict") return row;
if (isAgentDirectoryCopy(row) && row.state === "unavailable") return row;
return patch(row, { state: "unavailable", errorCode: "INSTRUCTION_COLLECTION_UNAVAILABLE",
errorMessage: "The registered instruction copy could not be retrieved safely before environment release. No instruction save is claimed.", nextAttemptAt: null });
}
async function release(companyId: string, runId: string) { liveTargets.delete(targetKey(companyId, runId)); const row = await get(companyId, runId); if (row && isAgentDirectoryCopy(row)) await directories.release(row); }
return { prepare, get, hasChanges, acknowledgeExplicitSave, collectStopped, recoverCaptured, recoverStopped, list, resolve, reportUnavailable, release };
return { prepare, get, hasChanges, checkpointWarm, canReuseWarm: directories.canReuse, acknowledgeExplicitSave, collectStopped, recoverCaptured, recoverStopped, list, resolve, reportUnavailable, release };
}
+53 -13
View File
@@ -1,6 +1,6 @@
import { isAiAuthenticationBlocked } from "./ai-auth-failure.js";
import { CHAT_COMPLETION_WAKE_REASON, prepareChatCompletionTurn, chatCompletionInstruction, isCompletedOnboardingHandoffWake } from "./chat-completion-delivery.js";
import { isAgentDirectoryCopy } from "./agent-directory-working-copies.js";
import { AgentDirectoryReuseInvalidatedError, isAgentDirectoryCopy } from "./agent-directory-working-copies.js";
import type { PaperclipTurnContext } from "@paperclipai/adapter-utils/server-utils";
import { restoreNativeWorkspaceBestEffort } from "./native-runtime/native-workspace-best-effort.js";
@@ -225,6 +225,7 @@ import {
claimNativeRestartRecoveries,
closeWarmNativeSessionsForEnvironment,
closeIdleWarmNativeSessionsForRestart,
reserveWarmNativeInstructionDirectory,
currentNativeControllerIdentity,
dispatchNativeSessionResumptions,
detachNativeSessionsForRestart,
@@ -20294,6 +20295,7 @@ export function heartbeatService(
Parameters<typeof cleanupGitHubOperationLaunchers>[0] | null = null;
let nativeSessionResumeScheduled = false;
let nativeOwnershipHeld = false;
let nativeInstructionReservation: Awaited<ReturnType<typeof reserveWarmNativeInstructionDirectory>> = null;
let nativeDispatchStarted = false;
let nativeWorkspaceFinalizeScheduled = false;
let nativeWorkspaceSync: Awaited<
@@ -22477,6 +22479,20 @@ export function heartbeatService(
const executionTarget = realizationResult.executionTarget;
let instructionCopy: Awaited<ReturnType<typeof instructionCopies.prepare>> = null;
let instructionSave: Record<string, unknown> | null = null;
const recordInstructionSave = async (saved: NonNullable<Awaited<ReturnType<typeof instructionCopies.get>>>) => {
const receipt = parseObject(saved.receipt);
const storageWarning = readNonEmptyString(receipt.storageWarning);
const state = saved.errorCode === "AGENT_FILES_CHECKPOINT_UNSTABLE" ? "pending_collection"
: saved.state === "warm_saved" ? readNonEmptyString(receipt.checkpointState) ?? "saved" : saved.state;
instructionSave = { state, entryFile: saved.entryFile,
...(isAgentDirectoryCopy(saved) ? { contract: "agent_files", appliedCandidateHash: saved.candidateHash, checkpointStats: receipt.checkpointStats }
: { revisionId: parseObject(receipt.revision).id ?? null }), storageWarning, errorCode: saved.errorCode, errorMessage: saved.errorMessage };
await appendRunEvent(run, { eventType: "instruction_save", stream: "system",
level: !saved.errorCode && !storageWarning && ["saved", "unchanged", "resolved"].includes(state) ? "info" : "warn",
message: storageWarning ?? (state === "saved" ? "Agent files saved."
: state === "unchanged" ? "Instruction working copy is unchanged."
: saved.errorMessage ?? "Instruction edits were not saved."), payload: instructionSave });
};
const collectStoppedInstructions = async () => {
if (!instructionCopy) return;
let saved = await instructionCopies.collectStopped({ companyId: agent.companyId, runId: run.id, target: executionTarget });
@@ -22486,17 +22502,7 @@ export function heartbeatService(
saved = await instructionCopies.collectStopped({ companyId: agent.companyId, runId: run.id, target: executionTarget });
}
if (!saved) return;
const receipt = parseObject(saved.receipt);
const storageWarning = readNonEmptyString(receipt.storageWarning);
instructionSave = { state: saved.state, entryFile: saved.entryFile,
...(isAgentDirectoryCopy(saved) ? { contract: "agent_files", appliedCandidateHash: saved.candidateHash }
: { revisionId: parseObject(receipt.revision).id ?? null }), storageWarning, errorCode: saved.errorCode, errorMessage: saved.errorMessage };
await appendRunEvent(run, { eventType: "instruction_save", stream: "system",
level: !storageWarning && ["saved", "unchanged", "resolved"].includes(saved.state) ? "info" : "warn",
message: storageWarning ?? (saved.state === "saved" ? "Agent files saved."
: saved.state === "unchanged" ? "Instruction working copy is unchanged."
: saved.errorMessage ?? "Instruction edits were not saved. Review the preserved candidate in the agent instruction editor."),
payload: instructionSave });
if (saved.state !== "superseded") await recordInstructionSave(saved);
};
if (managedAiRuntime && aiBinding) {
try { await assertManagedAiProjectAuth({ ...resolvedConfig, cwd: executionWorkspace.cwd }, aiBinding.provider, executionTarget); }
@@ -23321,11 +23327,36 @@ export function heartbeatService(
const savedFileInput = parseObject(parseObject(run.runnerProfileJson).nativeExecutionInput);
const priorFileInput = Object.keys(savedFileInput).length ? savedFileInput : parseObject(parseObject(priorFileRun?.profile).nativeExecutionInput);
const priorWorkingCopy = parseObject(parseObject(parseObject(priorFileInput.runtimeContext).instructions).workingCopy);
instructionCopy = await instructionCopies.prepare({
const warmFiles = nativeRuntimeResolution.kind === "native" && nativeRuntimeResolution.profile.backend === "codex_app_server" &&
(executionTarget?.kind === "remote" && executionTarget.transport === "sandbox"
? executionTarget.runnerLifecyclePolicy?.mode === "warm"
: parseObject(agent.adapterConfig).lifecycleMode === "warm");
if (warmFiles && taskSessionForRun?.lastRunId) {
nativeInstructionReservation = await reserveWarmNativeInstructionDirectory({ companyId: agent.companyId, agentId: agent.id,
previousRunId: taskSessionForRun.lastRunId, runId: run.id, target: executionTarget,
canReuse: () => instructionCopies.canReuseWarm(agent.companyId, agent.id, taskSessionForRun!.lastRunId!),
});
}
const prepareInstructions = (reuseRunId?: string) => instructionCopies.prepare({
companyId: agent.companyId, agentId: agent.id, runId: run.id,
target: executionTarget, cwd: executionWorkspace.cwd,
legacy: Object.keys(priorFileInput).length > 0 && priorWorkingCopy.kind !== "agent_files",
warm: warmFiles, reuseRunId,
onWarmHandoff: copy => {
instructionCopy = copy;
nativeInstructionReservation?.adopt(copy.executionRoot, collectStoppedInstructions);
},
});
try {
instructionCopy = await prepareInstructions(nativeInstructionReservation?.reuseRunId);
} catch (error) {
if (!(error instanceof AgentDirectoryReuseInvalidatedError)) throw error;
// prepare has released its canonical lock. Retirement can now
// collect under that same lock before a fresh restore starts.
await nativeInstructionReservation?.release();
nativeInstructionReservation = null;
instructionCopy = await prepareInstructions();
}
} catch (error) {
if ((error as { status?: number }).status !== 403) throw error;
// Missing write identity must not break a background run's read-only
@@ -24499,6 +24530,13 @@ export function heartbeatService(
onLog,
onEvent: onAdapterEvent,
instructionWorkingCopy: instructionCopy ? {
runId: run.id,
root: instructionCopy.executionRoot,
...(instructionCopy.receipt?.warm === true ? { checkpointWarm: async () => {
const saved = await instructionCopies.checkpointWarm({ companyId: agent.companyId, runId: run.id, target: executionTarget });
if (saved) await recordInstructionSave(saved);
return saved?.state === "warm_saved" && saved.errorCode === null;
} } : {}),
hasChanges: () => instructionCopies.hasChanges({ companyId: agent.companyId, runId: run.id, target: executionTarget }),
collectStopped: collectStoppedInstructions,
} : undefined,
@@ -24981,6 +25019,7 @@ export function heartbeatService(
"failed to revoke heartbeat-run MCP gateway tokens",
);
}
await nativeInstructionReservation?.release();
await instructionCopies.release(agent.companyId, run.id);
}
// Reconcile the referenced-project set against the real remote staging outcome. A referenced
@@ -26289,6 +26328,7 @@ export function heartbeatService(
}
}
} finally {
await nativeInstructionReservation?.release().catch(error => logger.warn({ runId: run.id, err: error }, "Managed warm session preparation cleanup failed"));
if (managedAiRuntime) await managedAiRuntime.cleanup().catch(() => logger.warn({ runId: run.id }, "AI connection refresh or cleanup failed"));
let latestRun = await getRun(run.id).catch(() => null);
try {
@@ -276,6 +276,7 @@ import {
buildNativeHarnessBackupManifest,
cancelNativeSession,
closeWarmNativeSessionsForEnvironment,
reserveWarmNativeInstructionDirectory,
closeIdleWarmNativeSessionsForRestart,
createGovernedWaitEventObservation,
createRemoteRunnerProcessLauncher,
@@ -5918,6 +5919,121 @@ describe("native session same-turn steering", () => {
});
describe("native warm session supervision", () => {
describe("managed directory warm checkpoints", () => {
const result = { result: { summary: "completed" }, terminal: { runTerminalState: "succeeded" },
turnId: "turn", normalizedSessionId: "managed", providerSessionId: "provider", driverKind: "test", driverVersion: "1",
nativeEventCount: 1, highestContiguousSourceSeq: 1, usage: null };
async function start(name: string, checkpoint = true) {
const current = { ...execution, binding: { ...execution.binding, runId: `${name}-one`, executionWorkspaceId: name },
session: { ...execution.session, normalizedSessionId: name, lifecyclePolicy: { mode: "warm" as const, idleTimeoutMs: 60_000 } } } as NativeExecutionInputV1;
const close = vi.fn(async () => undefined);
const session = { close };
const collectStopped = vi.fn(async () => undefined);
const checkpointWarm = vi.fn(async () => checkpoint);
const target = { kind: "remote" as const, transport: "sandbox" as const, environmentId: name, remoteCwd: `/tmp/${name}`, sandboxLeaseAcquisition: { outcome: "created" as const, providerLeaseId: name } };
const copy = { runId: current.binding.runId, root: `/tmp/${name}/home`, collectStopped, checkpointWarm, hasChanges: vi.fn(async () => true) };
state.execute.mockReset().mockImplementationOnce(async options => { await options.onSession?.(session); return result; });
const run = (next: NativeExecutionInputV1, nextCopy = copy) => executePaperclipNativeSession({
db: leaseDb(next), execution: next, runnerInstanceId: name, runnerExecutionTarget: target, instructionWorkingCopy: nextCopy });
await run(current);
const reserve = (canReuse: () => Promise<boolean>, nextTarget = target) => reserveWarmNativeInstructionDirectory({
companyId: current.binding.companyId, agentId: current.binding.agentId, previousRunId: current.binding.runId,
runId: `${name}-two`, target: nextTarget, canReuse });
return { current, copy, close, session, target, run, reserve };
}
afterEach(async () => { await closeIdleWarmNativeSessionsForRestart(); });
it("saves each turn while retaining the provider and collects only the latest directory owner at retirement", async () => {
const f = await start("managed-retained");
expect(f.copy.checkpointWarm).toHaveBeenCalledOnce();
expect(f.copy.hasChanges).not.toHaveBeenCalled();
expect(f.close).not.toHaveBeenCalled();
const reservation = await f.reserve(async () => true);
expect(reservation?.reuseRunId).toBe(f.current.binding.runId);
const secondCopy = { ...f.copy, runId: "managed-retained-two", collectStopped: vi.fn(async () => undefined) };
reservation!.adopt(secondCopy.root, secondCopy.collectStopped);
state.execute.mockImplementationOnce(async options => {
expect(options.existingSession).toBe(f.session);
await options.onSession?.(f.session); return result;
});
await f.run({ ...f.current, binding: { ...f.current.binding, runId: secondCopy.runId } }, secondCopy);
await reservation!.release();
expect(f.close).not.toHaveBeenCalled();
expect(f.copy.checkpointWarm).toHaveBeenCalledTimes(2);
await closeWarmNativeSessionsForEnvironment({ environmentId: f.target.environmentId, reason: "test retirement" });
expect(f.close).toHaveBeenCalledOnce();
expect(f.copy.collectStopped).not.toHaveBeenCalled();
expect(secondCopy.collectStopped).toHaveBeenCalledOnce();
});
it("stops before collecting when a checkpoint cannot stabilize", async () => {
const f = await start("managed-fallback", false);
expect(f.close).toHaveBeenCalledOnce();
expect(f.copy.collectStopped).toHaveBeenCalledOnce();
expect(f.close.mock.invocationCallOrder[0]).toBeLessThan(f.copy.collectStopped.mock.invocationCallOrder[0]!);
});
it.each(["canonical-edit", "replaced-environment", "authorization-error"])("retires before admission for %s", async reason => {
const f = await start(`managed-${reason}`);
const validation = vi.fn(async () => { if (reason === "authorization-error") throw new Error("authorization revoked"); return reason !== "canonical-edit"; });
const attempt = f.reserve(validation, reason === "replaced-environment" ? { ...f.target, remoteCwd: "/different" } : f.target);
if (reason === "authorization-error") await expect(attempt).rejects.toThrow("authorization revoked");
else expect(await attempt).toBeNull();
expect(f.close).toHaveBeenCalledOnce();
expect(f.copy.collectStopped).toHaveBeenCalledOnce();
if (reason === "replaced-environment") expect(validation).not.toHaveBeenCalled();
});
it("retains failed containment for retry and forbids warm reuse until it succeeds", async () => {
const f = await start("managed-stop-failure");
f.close.mockRejectedValueOnce(new Error("containment failed"));
await expect(f.reserve(async () => false)).rejects.toThrow("containment failed");
expect(f.copy.collectStopped).not.toHaveBeenCalled();
const canReuse = vi.fn(async () => true);
expect(await f.reserve(canReuse)).toBeNull();
expect(canReuse).not.toHaveBeenCalled();
expect(f.close).toHaveBeenCalledTimes(2);
expect(f.copy.collectStopped).toHaveBeenCalledOnce();
});
it("cleans up the successor if preparation fails after directory handoff", async () => {
const f = await start("managed-failed-preparation");
const reservation = await f.reserve(async () => true);
const collectSuccessor = vi.fn(async () => undefined);
reservation!.adopt(f.copy.root, collectSuccessor);
await reservation!.release();
await reservation!.release();
expect(f.close).toHaveBeenCalledOnce();
expect(f.copy.collectStopped).not.toHaveBeenCalled();
expect(collectSuccessor).toHaveBeenCalledOnce();
});
it("preserves a handed-off directory when configuration rotates the old provider", async () => {
const f = await start("managed-policy-rotation");
const reservation = await f.reserve(async () => true);
const nextCopy = { ...f.copy, runId: "managed-policy-rotation-two", collectStopped: vi.fn(async () => undefined) };
reservation!.adopt(nextCopy.root, nextCopy.collectStopped);
const replacement = { close: vi.fn(async () => undefined) };
state.execute.mockImplementationOnce(async options => {
expect(options.existingSession).toBeUndefined();
expect(f.close).toHaveBeenCalledOnce();
expect(nextCopy.collectStopped).not.toHaveBeenCalled();
await options.onSession?.(replacement); return result;
});
await f.run({ ...f.current, binding: { ...f.current.binding, runId: nextCopy.runId },
session: { ...f.current.session, lifecyclePolicy: { mode: "warm", idleTimeoutMs: 120_000 } } }, nextCopy);
await reservation!.release();
expect(nextCopy.collectStopped).not.toHaveBeenCalled();
await closeWarmNativeSessionsForEnvironment({ environmentId: f.target.environmentId, reason: "test" });
expect(nextCopy.collectStopped).toHaveBeenCalledOnce();
});
it("cannot reuse a provider with a different registered directory", async () => {
const f = await start("managed-root-replaced");
state.execute.mockImplementationOnce(async options => { expect(options.existingSession).toBeUndefined(); return result; });
await f.run({ ...f.current, binding: { ...f.current.binding, runId: "managed-root-two" } }, { ...f.copy, root: "/new/home" });
expect(f.close).toHaveBeenCalledOnce();
expect(f.copy.collectStopped).toHaveBeenCalledOnce();
});
});
it.each([
{ changed: false, closeFails: false },
{ changed: true, closeFails: false },
@@ -391,6 +391,9 @@ function clearNativeRuntimeRequestResolutions(runId: string): void {
}
type WarmNativeSession = {
agentId: string;
instructionCopy?: { runId: string; root?: string; targetIdentity: string; collectStopped: () => Promise<void> };
preparingRunId?: string;
managedAiCredentialIdentity?: string;
credentialRunId?: string;
githubAccess?: NativeGitHubAccess;
@@ -408,10 +411,13 @@ type WarmNativeSession = {
lastActivityAt: string;
};
async function closeWarmNativeSession(entry: WarmNativeSession, reason: string) {
async function closeWarmNativeSession(entry: WarmNativeSession, reason: string, preserveInstructionsForRunId?: string) {
// Revoke before awaiting process retirement/checkpoint IO.
const stopping = entry.githubAccess?.stop();
try { await entry.session.close({ reason }); }
try {
await entry.session.close({ reason });
if (entry.instructionCopy?.runId !== preserveInstructionsForRunId) await entry.instructionCopy?.collectStopped();
}
finally { await stopping; }
}
@@ -420,6 +426,56 @@ const warmNativeSessions = new Map<string, WarmNativeSession>();
// not inspect or quarantine that owner's state until the save has finished.
const closingWarmNativeSessions = new Map<string, Promise<void>>();
function instructionTargetIdentity(target?: AdapterExecutionTarget | null): string {
return JSON.stringify(target?.kind === "remote" ? {
environmentId: target.environmentId, cwd: target.remoteCwd,
providerLeaseId: target.transport === "sandbox" ? target.sandboxLeaseAcquisition?.providerLeaseId : target.spec,
} : { kind: "local", environmentId: target?.environmentId });
}
/** Reserve the actual live owner before handing its directory to another run.
* A persisted conversation alone is not proof that a working tree is reusable. */
export async function reserveWarmNativeInstructionDirectory(input: {
companyId: string; agentId: string; previousRunId: string; runId: string; target?: AdapterExecutionTarget | null;
canReuse: () => Promise<boolean>;
}): Promise<{ reuseRunId: string; adopt: (root: string, collectStopped: () => Promise<void>) => void; release: () => Promise<void> } | null> {
for (const [id, entry] of warmNativeSessions) {
if (entry.companyId !== input.companyId || entry.agentId !== input.agentId || entry.instructionCopy?.runId !== input.previousRunId) continue;
if (entry.busy || entry.preparingRunId) throw new Error("native_session_supervisor_busy");
if (entry.idleTimer) clearTimeout(entry.idleTimer);
entry.idleTimer = null;
entry.preparingRunId = input.runId;
const retire = async (reason: string) => {
entry.closeOnReleaseReason = reason;
try {
await closeWarmNativeSession(entry, reason);
if (warmNativeSessions.get(id) === entry) warmNativeSessions.delete(id);
} finally {
// Failed containment remains registered and cannot be adopted. The next
// admission retries retirement instead of launching over a live owner.
entry.preparingRunId = undefined;
}
};
let reusable = false;
const targetProven = input.target?.kind !== "remote" || input.target.transport !== "sandbox" || Boolean(input.target.sandboxLeaseAcquisition?.providerLeaseId);
try { reusable = targetProven && !entry.closeOnReleaseReason && entry.instructionCopy.targetIdentity === instructionTargetIdentity(input.target) && await input.canReuse(); }
finally {
if (!reusable) {
await retire("managed agent files changed or environment replaced");
}
}
if (!reusable) return null;
return { reuseRunId: input.previousRunId, adopt: (root, collectStopped) => {
entry.instructionCopy = { runId: input.runId, root, targetIdentity: instructionTargetIdentity(input.target), collectStopped };
}, release: async () => {
if (warmNativeSessions.get(id) === entry && entry.preparingRunId === input.runId) {
await retire("managed warm turn preparation ended without attachment");
}
} };
}
return null;
}
/**
* Close idle native sessions before an operator destroys their remote
* environment. A warm runner owns a long-lived sandbox command stream, so the
@@ -457,7 +513,7 @@ async function closeIdleWarmNativeSessions(input: {
if (input.environmentId !== undefined && entry.environmentId !== input.environmentId) {
continue;
}
if (entry.busy || executingRunnerdSessionScopes.has(sessionId)) {
if (entry.busy || entry.preparingRunId || executingRunnerdSessionScopes.has(sessionId)) {
// A busy turn can complete while another idle session is checkpointing.
// Fence that entry now so its eventual release cannot leave a new idle
// owner behind after the shutdown sweep has already passed it.
@@ -7181,9 +7237,12 @@ export async function executePaperclipNativeSession(input: {
chatAttachmentReadScope?: NativeChatAttachmentReadScope;
onLog?: (stream: "stdout" | "stderr", chunk: string) => Promise<void>;
onEvent?: (event: AdapterRuntimeEvent) => Promise<void>;
/** Only this run's registered private instruction entry.
* Probe at terminal; persist only after owned shutdown. */
/** Only this run's registered directory. Validate a warm checkpoint at the
* terminal boundary, or collect after owned shutdown when it cannot settle. */
instructionWorkingCopy?: {
runId?: string;
root?: string;
checkpointWarm?: () => Promise<boolean>;
hasChanges: () => Promise<boolean>;
collectStopped: () => Promise<void>;
};
@@ -8055,6 +8114,8 @@ async function executePaperclipNativeSessionWithinScope(
if (warmSessionId !== null && warmConfigDigest !== null) {
const entry = warmNativeSessions.get(warmSessionId);
if (entry) {
if (entry.preparingRunId && entry.preparingRunId !== input.execution.binding.runId) throw new Error("native_session_supervisor_busy");
entry.preparingRunId = undefined;
// Old run-scoped environments still require process replacement. A
// session-owned broker can change run authority without replacing it.
const hasBrokerCapability = Boolean(
@@ -8067,6 +8128,7 @@ async function executePaperclipNativeSessionWithinScope(
if (
entry.closeOnReleaseReason !== undefined ||
entry.configDigest !== warmConfigDigest ||
entry.instructionCopy?.root !== input.instructionWorkingCopy?.root ||
entry.managedAiCredentialIdentity !== input.managedAiCredentialIdentity ||
credentialRunChanged ||
Boolean(entry.githubAccess) !== Boolean(input.managedGitHub) ||
@@ -8080,7 +8142,10 @@ async function executePaperclipNativeSessionWithinScope(
if (entry.busy) throw new Error("native_session_supervisor_busy");
if (entry.idleTimer !== null) clearTimeout(entry.idleTimer);
warmNativeSessions.delete(warmSessionId);
await closeWarmNativeSession(entry, "warm native session configuration changed");
// Preparation may already have handed this directory to the incoming
// run. Retiring the old transport must not delete the new run's root;
// its own stop callback owns collection from this point forward.
await closeWarmNativeSession(entry, "warm native session configuration changed", input.execution.binding.runId);
persistedWarmSession = loadWarmNativeCheckpoint(
input.execution,
warmConfigDigest,
@@ -8417,6 +8482,7 @@ async function executePaperclipNativeSessionWithinScope(
existing.session = session;
} else
warmNativeSessions.set(warmSessionId, {
agentId: input.execution.binding.agentId,
managedAiCredentialIdentity: input.managedAiCredentialIdentity,
githubAuthenticationMode:
input.runnerEnvironment?.PAPERCLIP_GITHUB_AUTH_MODE,
@@ -8441,6 +8507,11 @@ async function executePaperclipNativeSessionWithinScope(
idleTimer: null,
lastActivityAt: new Date().toISOString(),
});
const owner = warmNativeSessions.get(warmSessionId);
if (owner && input.instructionWorkingCopy?.runId) owner.instructionCopy = {
runId: input.instructionWorkingCopy.runId, root: input.instructionWorkingCopy.root, targetIdentity: instructionTargetIdentity(input.runnerExecutionTarget),
collectStopped: input.instructionWorkingCopy.collectStopped,
};
} else if (!session && warmSessionId !== null) {
const existing = warmNativeSessions.get(warmSessionId);
// onSession(null) quarantines a transport that can no longer
@@ -9090,7 +9161,8 @@ async function executePaperclipNativeSessionWithinScope(
const instructionCopy = input.instructionWorkingCopy;
const ownedSession = warmNativeSessions.get(warmSessionId);
const collectInstructions = Boolean(instructionCopy && ownedSession?.ownerToken === warmSessionOwnerToken &&
await instructionCopy.hasChanges());
(instructionCopy.checkpointWarm ? !await instructionCopy.checkpointWarm() : await instructionCopy.hasChanges()));
const collectedByOwner = ownedSession?.instructionCopy?.collectStopped === instructionCopy?.collectStopped;
if (collectInstructions && ownedSession) {
// Keep the unchanged warm path intact. A changed private instruction copy
// requires the existing checkpoint-and-close boundary before collection.
@@ -9102,7 +9174,7 @@ async function executePaperclipNativeSessionWithinScope(
lifecyclePolicy.idleTimeoutMs,
false,
);
if (collectInstructions) await instructionCopy!.collectStopped();
if (collectInstructions && !collectedByOwner) await instructionCopy!.collectStopped();
}
const adapterResult: AdapterExecutionResult = {
exitCode: native.terminal.runTerminalState === "succeeded" ? 0 : 1,
@@ -0,0 +1,5 @@
export interface AgentFileManifestEntry { kind: "dir" | "file"; size?: number; mode?: number; hash?: string; stamp?: string }
export interface AgentFileManifest { version?: number; entries: Array<[string, AgentFileManifestEntry]>; totalBytes?: number }
export interface AgentFileCheckpointStats { hashedBytes: number; copiedBytes: number; copiedFiles: number; scannedEntries: number }
export function checkpointPath(name: string): string;
export function captureAgentFiles(root: string, previous?: AgentFileManifest, output?: string, enforceLimits?: boolean, excludeTransportRuntime?: boolean): Promise<{ manifest: AgentFileManifest; stats: AgentFileCheckpointStats }>;
@@ -0,0 +1,138 @@
// This module also runs verbatim in a remote workspace. Keep it dependency-free.
import fs from "node:fs/promises";
import { constants } from "node:fs";
import { createHash } from "node:crypto";
import path from "node:path";
const MAX_FILE = 256 * 1024 * 1024;
const MAX_BYTES = 2 * 1024 * 1024 * 1024;
const MAX_ENTRIES = 100_000;
const stamp = s => [s.dev, s.ino, s.size, s.mode, s.mtimeNs, s.ctimeNs].map(String).join(":");
const sameContent = (a, b) => a?.kind === b?.kind && (a?.kind === "dir" || a?.hash === b?.hash && a?.mode === b?.mode);
const fail = (code) => { throw Object.assign(new Error(code), { code }); };
export function checkpointPath(name) {
if (typeof name !== "string" || !name || name.includes("\\") || name.includes("\0") || name.startsWith("/") ||
name.split("/").some(p => !p || p === "." || p === ".." || p === ".paperclip-runtime") || name === "promptTemplate.legacy.md") fail("AGENT_FILES_UNSAFE_PATH");
return name;
}
async function assertRoot(root) {
let absolute = path.resolve(root);
// Match the managed-file store: macOS owns these aliases. Accepting them
// does not permit a user-created symlink anywhere inside the agent tree.
if (process.platform === "darwin") for (const alias of ["var", "tmp", "etc"]) {
const systemPath = `/${alias}`;
if (absolute.startsWith(`${systemPath}/`)) {
const stat = await fs.lstat(systemPath);
if (stat.isSymbolicLink() && stat.uid === 0 && await fs.realpath(systemPath) === `/private/${alias}`) {
absolute = `/private${absolute}`;
}
}
}
if (await fs.realpath(absolute) !== absolute || !(await fs.lstat(absolute)).isDirectory()) fail("AGENT_FILES_UNSAFE_PATH");
return absolute;
}
/** Metadata is a hash cache, never a dirty-file journal. ctime also detects
* same-size edits whose mtime was restored. Every checkpoint enumerates paths. */
export async function captureAgentFiles(root, previous = { entries: [] }, output, enforceLimits = true, excludeTransportRuntime = false) {
root = await assertRoot(root);
const cache = new Map(previous.entries);
const entries = [];
const observed = new Map();
const stats = { hashedBytes: 0, copiedBytes: 0, copiedFiles: 0, scannedEntries: 0 };
let total = 0;
if (output) await fs.mkdir(path.join(output, "files"), { recursive: true, mode: 0o700 });
const list = async dir => (await fs.readdir(path.join(root, dir)))
.filter(name => !(excludeTransportRuntime && dir === "" && name === ".paperclip-runtime")).sort();
async function walk(dir) {
const names = await list(dir);
observed.set(dir, names);
for (const name of names) {
const relative = checkpointPath(dir ? `${dir}/${name}` : name);
const filename = path.join(root, relative);
const s = await fs.lstat(filename, { bigint: true });
stats.scannedEntries++;
if (enforceLimits && stats.scannedEntries > MAX_ENTRIES) fail("AGENT_FILES_LIMIT_EXCEEDED");
if (s.isDirectory()) {
entries.push([relative, { kind: "dir" }]);
await walk(relative);
} else {
if (!s.isFile() || s.nlink !== 1n) fail("AGENT_FILES_UNSAFE_PATH");
const size = Number(s.size), mode = Number(s.mode), fingerprint = stamp(s);
total += size;
if (enforceLimits && (size > MAX_FILE || total > MAX_BYTES)) fail("AGENT_FILES_LIMIT_EXCEEDED");
const old = cache.get(relative);
let hash = old?.kind === "file" && old.stamp === fingerprint ? old.hash : null;
if (!hash) {
const h = createHash("sha256");
const f = await fs.open(filename, constants.O_RDONLY | constants.O_NOFOLLOW);
try {
if (stamp(await f.stat({ bigint: true })) !== fingerprint) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
let read = 0;
for await (const chunk of f.createReadStream({ autoClose: false })) {
read += chunk.length;
if (read > size) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
h.update(chunk); stats.hashedBytes += chunk.length;
}
if (stamp(await f.stat({ bigint: true })) !== fingerprint) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
hash = h.digest("hex");
} finally { await f.close(); }
}
const entry = { kind: "file", size, mode, hash, stamp: fingerprint };
if (output && !sameContent(old, entry)) {
const target = path.join(output, "files", relative);
await fs.mkdir(path.dirname(target), { recursive: true, mode: 0o700 });
const source = await fs.open(filename, constants.O_RDONLY | constants.O_NOFOLLOW);
let dest;
try {
dest = await fs.open(target, "wx", mode & 0o777);
if (stamp(await source.stat({ bigint: true })) !== fingerprint) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
const h = createHash("sha256");
let copied = 0;
for await (const chunk of source.createReadStream({ autoClose: false })) {
copied += chunk.length;
if (copied > size) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
h.update(chunk);
let offset = 0;
while (offset < chunk.length) {
const { bytesWritten } = await dest.write(chunk, offset, chunk.length - offset);
if (!bytesWritten) fail("AGENT_FILES_CHECKPOINT_WRITE_FAILED");
offset += bytesWritten;
}
stats.copiedBytes += chunk.length;
}
if (h.digest("hex") !== hash || stamp(await source.stat({ bigint: true })) !== fingerprint) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
await dest.sync();
await dest.chmod(mode & 0o777);
stats.copiedFiles++;
} finally { await source.close(); await dest?.close(); }
}
entries.push([relative, entry]);
}
}
}
await walk("");
// Catch deletions, renames, additions and writers racing the scan/copy. These
// are per-file checkpoints, not an atomic snapshot of arbitrary live writers.
for (const [dir, names] of observed) if (JSON.stringify(await list(dir)) !== JSON.stringify(names)) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
for (const [name, entry] of entries) {
const s = await fs.lstat(path.join(root, name), { bigint: true });
if (entry.kind === "dir" ? !s.isDirectory() : stamp(s) !== entry.stamp || !s.isFile() || s.nlink !== 1n) fail("AGENT_FILES_CHANGED_DURING_CHECKPOINT");
}
await assertRoot(root);
const manifest = { version: 1, entries, totalBytes: total };
if (output) await fs.writeFile(path.join(output, "checkpoint.json"), JSON.stringify(manifest), { mode: 0o600, flag: "wx" });
return { manifest, stats };
}
if (process.argv[1] === "--checkpoint") {
try {
const input = JSON.parse(process.argv[2]);
const bytes = await fs.readFile(input.cacheFile);
if (createHash("sha256").update(bytes).digest("hex") !== input.cacheHash) fail("AGENT_FILES_CHECKPOINT_CACHE_MISMATCH");
const result = await captureAgentFiles(input.root, JSON.parse(bytes), input.output, true, true);
console.log(JSON.stringify(result.stats));
} catch (e) { console.error(e.code ?? "AGENT_FILES_CHECKPOINT_FAILED"); process.exitCode = 1; }
}
+7 -1
View File
@@ -108,7 +108,13 @@ canonical Plan revision, capture its pending UI, approve in the browser, and
prove exactly two successful runs. `warm_three_turn` provides exactly two
browser follow-up messages, preserves one project/execution-workspace scope,
verifies host file contents after every turn, and finishes within three
ten-minute turn deadlines. Native turns 1 and 2 include an actionable human review in the completion report's `attentionRequests`. Paperclip creates the review gate from that report. An explicit question-tool wait yields the turn and suppresses its final prose, so it is not interchangeable with this completion-review fixture. Turn 3 reports Done without another review.
ten-minute turn deadlines. The ordinary warm fixture uses managed instructions,
updates AGENT_HOME each turn, and verifies memory, an unchanged 8 MiB binary and
a deletion through public file APIs. Native turns 2 and 3 must copy/hash only the
changed memory file, with a saved receipt and the same provider PID. Journal and
Git stress fixtures retain fixed external bundles as controls. Keep the stable-PID
oracle strict; `instruction-persistence` also covers cold restarts and quota handling.
Native turns 1 and 2 include an actionable human review in the completion report's `attentionRequests`. Paperclip creates the review gate from that report. An explicit question-tool wait yields the turn and suppresses its final prose, so it is not interchangeable with this completion-review fixture. Turn 3 reports Done without another review.
Every selected case runs in its own isolated Paperclip process, and independent
cases may run concurrently. Follow-up turns inside one case retain their shared
+11 -8
View File
@@ -274,9 +274,8 @@ a warning while full and clear it after cleanup. Rejected bytes must not replace
the saved file. This adds at most one 256 MiB saved fixture per isolated agent.
The deadline is twenty minutes per cell, with six expected provider runs;
normal instance/Daytona cleanup, screenshots, evidence, and billing apply. Run with
`pnpm test:e2e:runner -- --suite instruction-persistence`. Managed agent directories
checkpoint and close the provider before collection while retaining conversation
state. The separate `daytona-warm-continuity` suite covers warm runtime behavior.
`pnpm test:e2e:runner -- --suite instruction-persistence`. Managed directories in per-turn sessions collect after provider stop. The separate
`daytona-warm-continuity` suite covers incremental saves while retaining a live native process.
`daytona-warm-continuity` (**Daytona Warm Continuity**) is exactly two paid
cells: legacy Codex and Runner Codex against one reusable warm Daytona
@@ -286,7 +285,12 @@ three browser-driven turns on one issue. Every turn reads and extends the same
nonce file, verifies host copy-back, records scheduler/run/end-to-end timing,
and asserts `created`, `resumed`, `resumed` lease acquisition on one sandbox.
Runner Codex additionally proves stable native session, provider session,
runner instance, PID, and process-start identity. Each turn is bounded to ten
runner instance, PID, and process-start identity. The ordinary warm cell uses
managed instructions and edits AGENT_HOME on every turn: a growing memory file,
an unchanged 8 MiB binary, and a deletion. Public API reads independently verify
the canonical bytes after every turn. Native checkpoint receipts must show only
the changed memory file transferred on turns 2 and 3; the PID oracle remains strict.
The journal stress cell retains fixed external instructions as a control. Each warm turn is bounded to ten
minutes, the cell to thirty minutes, and cleanup explicitly deletes the
sandbox rather than waiting for Daytona's idle timeout.
@@ -310,10 +314,9 @@ over the agent setting. Large copyback plus the next preparation
can exceed the normal five-minute idle window; the PID and process-fingerprint
continuity checks remain strict. The ordinary warm-continuity cell keeps its
existing five-minute policy.
This cell uses a fixed external instruction bundle. Managed agent folders
intentionally checkpoint and stop the provider after every turn for file
collection; external instructions allow this cell to test retained-process
continuity without changing that collection policy.
This Git stress cell uses a fixed external instruction bundle to isolate workspace
transfer from managed agent-file persistence. The ordinary warm cell separately
requires incremental managed-file checkpoints and the same strict process continuity.
It seeds an empty local Git project, creates 60,000 small untracked files through
the real provider, then performs the same three browser-driven review turns.
Each later turn updates all 60,000 generated files to distinct turn-specific
+26 -1
View File
@@ -32,6 +32,29 @@ import {
} from "./selectors.js";
describe("runner E2E catalog", () => {
it.each(["daytona-journal-continuity"])(
"%s retains native processes with fixed external instructions",
(suiteId) => {
const suite = runnerSuites.find(suite => suite.id === suiteId)!;
const input = {
executionId: "warm-instructions", environmentId: "env-1",
environmentFixtureId: "daytona" as const, workspacePath: "/workspace",
secretRefs: { OPENAI_API_KEY: { type: "secret_ref" as const, secretId: "22222222-2222-4222-8222-222222222222", version: "latest" as const } },
};
const profile = suite.profiles.find(profile => profile.id === "runner-codex")!;
const config = profile.buildAgent(input).adapterConfig as Record<string, unknown>;
expect(config).toMatchObject({
instructionsBundleMode: "external", instructionsEntryFile: "AGENTS.md", idleTimeoutMs: 300_000,
});
expect(readFileSync(String(config.instructionsFilePath), "utf8")).toContain("Preserve the workspace files across turns");
expect(suite.definitionMetadata).toMatchObject({ nativeInstructions: "fixed-external" });
expect(runnerProfiles.find(profile => profile.id === "runner-codex")!.buildAgent(input).adapterConfig)
.not.toHaveProperty("instructionsBundleMode", "external");
const legacy = suite.profiles.find(profile => profile.id === "legacy-codex");
if (legacy) expect(legacy.buildAgent(input).adapterConfig).not.toHaveProperty("instructionsBundleMode", "external");
},
);
it("keeps the large Git filename workload explicit-only with three real turns", () => {
const suite = runnerSuites.find(suite => suite.id === "daytona-git-streaming")!;
expect(suite.manualOnly).toBe(true);
@@ -235,7 +258,9 @@ describe("runner E2E catalog", () => {
expect(suite.profiles.map((profile) => profile.id)).toEqual(["runner-codex"]);
expect(daytonaLargeJournalTask.flow).toBe("warm_three_turn");
expect(daytonaLargeJournalTask.buildPrompt("nonce")).toContain("240 separate execution-tool calls");
expect(daytonaLargeJournalTask.buildFollowupMessages!("nonce")).toEqual(daytonaWarmContinuityTask.buildFollowupMessages!("nonce"));
expect(daytonaLargeJournalTask.buildFollowupMessages!("nonce")).toHaveLength(2);
expect(daytonaLargeJournalTask.buildFollowupMessages!("nonce").join("\n")).not.toContain("notes/warm-memory.txt");
expect(daytonaWarmContinuityTask.buildPrompt("nonce")).toContain("notes/warm-memory.txt");
expect(selectRunnerExecutions(parseRunnerSelectors(["--all"]))
.some((cell) => cell.suite.id === suite.id)).toBe(false);
});
+36 -6
View File
@@ -935,6 +935,16 @@ function warmTurnInstructions(turn: 1 | 2 | 3, nonce: string) {
].join("\n");
}
function managedWarmTurnInstructions(turn: 1 | 2 | 3, nonce: string) {
return [warmTurnInstructions(turn, nonce),
"Also update your personal AGENT_HOME with ordinary filesystem tools. Do not edit the loaded AGENTS.md instructions and do not use an API to save these files.",
turn === 1
? `Create notes/warm-memory.txt containing exactly T1-${nonce} followed by a newline. Create notes/unchanged.bin with exactly 8388608 bytes, each byte equal to 93. Create notes/delete-me.txt containing temporary.`
: `Read notes/warm-memory.txt under AGENT_HOME and verify it contains exactly the prior turn lines ${Array.from({ length: turn - 1 }, (_, i) => `T${i + 1}-${nonce}`).join(" | ")}, each followed by a newline. Append exactly T${turn}-${nonce} and a newline. Verify notes/unchanged.bin still has 8388608 bytes, each equal to 93, and leave it unchanged. ${turn === 2 ? "Delete notes/delete-me.txt." : "Verify notes/delete-me.txt is absent."}`,
"Perform these personal-file edits and verification before calling paperclip_finish. Paperclip saves them at the turn boundary.",
].join("\n");
}
export const daytonaWarmContinuityTask: RunnerTaskFixture = {
id: "warm-three-turn",
label: "Warm three-turn workspace continuity",
@@ -947,10 +957,10 @@ export const daytonaWarmContinuityTask: RunnerTaskFixture = {
expectedTerminalState: { issue: "done", run: "succeeded" },
buildTitle: (nonce) => `Runner E2E warm Daytona continuity ${nonce}`,
buildVisibleMarker: (nonce) => warmTurnMarker(3, nonce),
buildPrompt: (nonce) => warmTurnInstructions(1, nonce),
buildPrompt: (nonce) => managedWarmTurnInstructions(1, nonce),
buildFollowupMessages: (nonce) => [
warmTurnInstructions(2, nonce),
warmTurnInstructions(3, nonce),
managedWarmTurnInstructions(2, nonce),
managedWarmTurnInstructions(3, nonce),
],
buildMatchers(nonce, execution) {
// Workspace persistence is the oracle for this story. Exact response text
@@ -980,6 +990,7 @@ export const daytonaLargeJournalTask: RunnerTaskFixture = {
...daytonaWarmContinuityTask,
id: "large-journal-three-turn",
label: "Large journal three-turn workspace continuity",
buildFollowupMessages: nonce => [warmTurnInstructions(2, nonce), warmTurnInstructions(3, nonce)],
buildTitle: (nonce) => `Runner E2E large journal continuity ${nonce}`,
buildPrompt: (nonce) => [
"First exercise ordinary execution history with 240 separate execution-tool calls. In each call, run the Python command below exactly once. Issue the calls one by one. Do not combine them into a shell loop, script, parallel wrapper, or a single tool call: each command must be a separate ordinary execution-tool invocation. Keep a count from 1 through 240. The output is synthetic fixture data and needs no analysis.",
@@ -988,12 +999,30 @@ export const daytonaLargeJournalTask: RunnerTaskFixture = {
warmTurnInstructions(1, nonce),
].join("\n"),
};
export const daytonaGitStreamingTask = createGitStreamingTask(daytonaWarmContinuityTask);
export const daytonaGitStreamingTask = createGitStreamingTask({ ...daytonaWarmContinuityTask, buildPrompt: nonce => warmTurnInstructions(1, nonce), buildFollowupMessages: nonce => [warmTurnInstructions(2, nonce), warmTurnInstructions(3, nonce)] });
const codexContinuityProfiles = runnerProfiles.filter((profile) =>
["legacy-codex", "runner-codex"].includes(profile.id),
);
// The journal stress fixture keeps a fixed external bundle as a control.
// Ordinary warm continuity exercises managed files and incremental checkpoints.
const warmCodexContinuityProfiles = codexContinuityProfiles.map((profile) =>
profile.generation !== "native" ? profile : {
...profile,
buildAgent(input: AgentFixtureBuildInput) {
const agent = profile.buildAgent(input);
return { ...agent, adapterConfig: {
...(agent.adapterConfig as Record<string, unknown>),
instructionsBundleMode: "external",
instructionsRootPath: fileURLToPath(new URL("./fixtures/warm-continuity/", import.meta.url)),
instructionsEntryFile: "AGENTS.md",
instructionsFilePath: fileURLToPath(new URL("./fixtures/warm-continuity/AGENTS.md", import.meta.url)),
} };
},
},
);
export const connectionReviewSuite: RunnerSuiteFixture = {
id: "connection-reviews",
label: "Governed Connection Reviews",
@@ -1292,6 +1321,7 @@ export const runnerSuites: readonly RunnerSuiteFixture[] = [
environments: [daytonaWarmEnvironment],
tasks: [daytonaWarmContinuityTask],
expectedMatrixSize: 2,
definitionMetadata: { version: 3, nativeInstructions: "managed-incremental", managedFileBytes: 8 * 1024 * 1024 },
},
{
id: "daytona-journal-continuity",
@@ -1299,11 +1329,11 @@ export const runnerSuites: readonly RunnerSuiteFixture[] = [
manualOnly: true,
description: "Continue the same native session after separate ordinary tool invocations and their output grow its durable journal beyond 2 MiB.",
groups: ["daytona", "warm"],
profiles: codexContinuityProfiles.filter((profile) => profile.id === "runner-codex"),
profiles: warmCodexContinuityProfiles.filter((profile) => profile.id === "runner-codex"),
environments: [daytonaWarmEnvironment],
tasks: [daytonaLargeJournalTask],
expectedMatrixSize: 1,
definitionMetadata: { version: 4, journalMinimumBytes: 2 * 1024 * 1024, toolInvocations: 240, outputBytesPerInvocation: 65020, scheduling: "explicit-only" },
definitionMetadata: { version: 5, nativeInstructions: "fixed-external", journalMinimumBytes: 2 * 1024 * 1024, toolInvocations: 240, outputBytesPerInvocation: 65020, scheduling: "explicit-only" },
},
{
id: "daytona-git-streaming",
@@ -0,0 +1,6 @@
# Warm continuity fixture agent
Complete the current task and its authorized review revisions using the supplied
task context. Preserve the workspace files across turns and follow the task's
verification and completion instructions. These external instructions are fixed;
do not edit this instruction file or create a personal agent folder.
+11
View File
@@ -1,3 +1,4 @@
import { warmManagedFileEvidence } from "./warm-managed-files.js";
import { gitFinalizationEvidence, gitStreamingEvidence, setupGitStreamingWorkspace } from "./daytona-git-streaming.js";
import { runsCompletionUpdateProbe, completionQualityControls, completionQualityStatus, judgeCompletionQuality, reserveCompletionQuality, type CompletionQualityRecord } from "./completion-quality.js";
import { completionDelivery, type CompletionObservation } from "./completion-updates.js";
@@ -1541,6 +1542,11 @@ for (const execution of executions) {
await writeSanitizedJson(snapshotsDir, `git-finalization-turn-${completedTurn}.json`, finalization, secrets);
expect(finalization.passed, finalization.failures.join("; ")).toBe(true);
}
if (execution.suite.id === "daytona-warm-continuity") {
const evidence = await warmManagedFileEvidence(api, fixtures.agent.id, sortRunsChronologically(waitingState.taskRuns).at(-1)!.id, completedTurn, nonce, execution.profile.generation === "native");
await writeSanitizedJson(snapshotsDir, `managed-warm-turn-${completedTurn}.json`, evidence, secrets);
expect(evidence.passed, JSON.stringify(evidence.checks)).toBe(true);
}
const expectedPrefix = `${Array.from(
{ length: completedTurn },
(_, index) => `T${index + 1}-${nonce}`,
@@ -1823,6 +1829,11 @@ for (const execution of executions) {
};
});
}
if (execution.suite.id === "daytona-warm-continuity") {
const evidence = await warmManagedFileEvidence(api, fixtures.agent.id, selectedRuns.at(-1)!.id, 3, nonce, execution.profile.generation === "native");
await writeSanitizedJson(snapshotsDir, "managed-warm-turn-3.json", evidence, secrets);
expect(evidence.passed, JSON.stringify(evidence.checks)).toBe(true);
}
const finalRun = selectedRuns.at(-1)!;
const run =
selectedRuns.find(
@@ -0,0 +1,11 @@
import { describe, expect, it } from "vitest";
import { gradeWarmManagedFiles } from "./warm-managed-files.js";
describe("managed warm file oracle", () => {
const good = { turn: 2, nonce: "test", content: "T1-test\nT2-test\n", binary: Buffer.alloc(8388608, 93), deletedStatus: 404, native: true,
save: { state: "saved", contract: "agent_files", checkpointStats: { copiedBytes: 16, copiedFiles: 1, hashedBytes: 16 } } };
it("accepts independently saved bytes with an incremental receipt", () => expect(gradeWarmManagedFiles(good).passed).toBe(true));
it.each([
{ content: "T1-test\n" }, { binary: Buffer.from("claimed image") }, { deletedStatus: 200 }, { save: undefined },
{ save: { ...good.save, checkpointStats: { copiedBytes: 8388624, copiedFiles: 2, hashedBytes: 8388624 } } },
])("rejects missing, stale, deleted, or fully recopied data", wrong => expect(gradeWarmManagedFiles({ ...good, ...wrong }).passed).toBe(false));
});
+32
View File
@@ -0,0 +1,32 @@
import { createHash } from "node:crypto";
import type { RunnerApi } from "./api.js";
const record = (v: unknown): Record<string, unknown> => v && typeof v === "object" && !Array.isArray(v) ? v as Record<string, unknown> : {};
export function gradeWarmManagedFiles(input: {
turn: number; nonce: string; content: unknown; binary: Buffer; deletedStatus: number; save: unknown; native: boolean;
}) {
const expected = `${Array.from({ length: input.turn }, (_, i) => `T${i + 1}-${input.nonce}`).join("\n")}\n`;
const expectedBinary = Buffer.alloc(8 * 1024 * 1024, 93);
const save = record(input.save), stats = record(save.checkpointStats);
const checks = [
{ id: "managed-memory-exact", passed: input.content === expected },
{ id: "managed-binary-exact", passed: input.binary.equals(expectedBinary) },
{ id: "managed-deletion", passed: input.deletedStatus === (input.turn === 1 ? 200 : 404) },
{ id: "managed-save-receipt", passed: save.state === "saved" && save.contract === "agent_files" && !save.errorCode },
...(input.native ? [{ id: "incremental-payload", passed: typeof stats.copiedFiles === "number" &&
(input.turn === 1 ? stats.copiedFiles === 3 && stats.copiedBytes === expectedBinary.length + Buffer.byteLength(expected) + 9
: stats.copiedFiles === 1 && stats.copiedBytes === Buffer.byteLength(expected) && stats.hashedBytes === Buffer.byteLength(expected)) }] : []),
];
return { passed: checks.every(c => c.passed), checks, expected, binaryBytes: input.binary.length,
binarySha256: createHash("sha256").update(input.binary).digest("hex"), save };
}
export async function warmManagedFileEvidence(api: RunnerApi, agentId: string, runId: string, turn: number, nonce: string, native: boolean) {
const base = `/api/agents/${agentId}/instructions-bundle/file?path=`;
const note = await api.get<{ content: string }>(`${base}notes%2Fwarm-memory.txt`);
const binary = await api.request.get(`${base}notes%2Funchanged.bin&download=true`);
if (!binary.ok()) throw new Error(`Managed binary download failed: ${binary.status()}`);
const deleted = await api.request.get(`${base}notes%2Fdelete-me.txt&download=true`);
const run = await api.get<{ resultJson?: unknown }>(`/api/heartbeat-runs/${runId}`);
return gradeWarmManagedFiles({ turn, nonce, content: note.content, binary: await binary.body(), deletedStatus: deleted.status(),
save: record(run.resultJson).instructionSave, native });
}