diff --git a/doc/sandbox-work-folders.md b/doc/sandbox-work-folders.md index 4f96f4c85e..7e20f54294 100644 --- a/doc/sandbox-work-folders.md +++ b/doc/sandbox-work-folders.md @@ -247,10 +247,12 @@ either observes completion or fails visibly and retains the working copy for the next run's existing intent reconciliation. This does not make arbitrary sandbox commands or repository mutations retryable. -Repository checkpoints transfer at most four distinct content-addressed blobs -concurrently, avoiding duplicate uploads for identical files. Small-file reads -are grouped into at most 1 MiB and 64 files per remote command, with at most -four read batches cached per checkpoint. Larger files stream independently. +Repository checkpoints transfer up to sixteen distinct batch-readable blobs +of at most 1 MiB concurrently, plus at most four larger streaming blobs. +Transports without batched reads retain the four-stream limit. Identical files +share one content-addressed upload. Small-file reads are grouped into at most +1 MiB and 64 files per remote command, with at most four read batches cached +per checkpoint and sixteen additional batches held by active readers. Retries bypass that cache and reopen the actual file. All active transfers must settle, and a second filesystem scan must match, before the complete checkpoint reference can advance. Scoped-file retry receipts and diff --git a/server/src/__tests__/heartbeat-list.test.ts b/server/src/__tests__/heartbeat-list.test.ts index e45e4e958a..0c1df4da95 100644 --- a/server/src/__tests__/heartbeat-list.test.ts +++ b/server/src/__tests__/heartbeat-list.test.ts @@ -6,6 +6,8 @@ import { startEmbeddedPostgresTestDatabase, } from "./helpers/embedded-postgres.js"; import { boundHeartbeatRunEventPayloadForStorage, heartbeatService } from "../services/heartbeat.ts"; +import { measureSandboxOperation, runWithSandboxPerformanceTrace, type SandboxPerformanceRecord } from "../services/sandbox-performance.js"; +import { redactEventPayload } from "../redaction.js"; const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport(); const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe : describe.skip; @@ -283,6 +285,23 @@ describeEmbeddedPostgres("heartbeat list", () => { }); describe("heartbeat run event payload bounding", () => { + it("preserves every performance record through the actual storage bounds and redaction", async () => { + const original: SandboxPerformanceRecord[] = []; + const stored: SandboxPerformanceRecord[] = []; + await runWithSandboxPerformanceTrace({ runId: "trace-batch-boundary", enabled: true, + onBatch: async (batch) => { + original.push(...batch.records); + const payload = redactEventPayload(boundHeartbeatRunEventPayloadForStorage(batch)); + stored.push(...payload!.records as SandboxPerformanceRecord[]); + } }, async () => { + for (let i = 0; i < 350; i++) await measureSandboxOperation("sandbox.read", { fileIndex: i }, async () => undefined); + }); + expect(stored).toHaveLength(351); + expect(stored).toEqual(original.map((record) => ({ ...record, traceId: record.traceId ?? null }))); + expect(new Set(stored.map((record) => record.id)).size).toBe(351); + expect(stored.find((record) => record.name === "sandbox.run")?.attributes).toMatchObject({ recordCount: 351, dropped: 0 }); + }); + it("truncates oversized adapter metadata before storage", () => { const payload = boundHeartbeatRunEventPayloadForStorage({ adapterType: "codex_local", diff --git a/server/src/__tests__/sandbox-performance.test.ts b/server/src/__tests__/sandbox-performance.test.ts index 0fa7d045a0..152fb88ec1 100644 --- a/server/src/__tests__/sandbox-performance.test.ts +++ b/server/src/__tests__/sandbox-performance.test.ts @@ -93,7 +93,7 @@ describe("sandbox performance trace", () => { for (let i = 0; i < 350; i++) await measureSandboxOperation("sandbox.read", { fileIndex: i }, async () => undefined); finished = true; }); - expect(batches.map((batch) => batch.records.length)).toEqual([250, 50]); + expect(batches.map((batch) => batch.records.length)).toEqual([50, 50, 50, 50, 50, 50]); expect(batches.every((batch) => batch.dropped === 51)).toBe(true); const root = spans.find((span) => span.name === "sandbox.run")!; expect(root.attributes["paperclip.sandbox.recordCount"]).toBe(300); diff --git a/server/src/__tests__/work-folder-repositories.test.ts b/server/src/__tests__/work-folder-repositories.test.ts index 34d853397c..577c4d60f1 100644 --- a/server/src/__tests__/work-folder-repositories.test.ts +++ b/server/src/__tests__/work-folder-repositories.test.ts @@ -36,7 +36,7 @@ describe("bounded repository checkpoint transfers", () => { afterEach(() => vi.restoreAllMocks()); afterAll(async () => { await database?.cleanup(); }); - async function fixture() { + async function fixture(batched = false) { const taskId = randomUUID(); await db.insert(issues).values({ id: taskId, companyId, title: "Checkpoint" }); const [binding] = await db.insert(taskRepositoryBindings).values({ companyId, taskId, @@ -47,7 +47,8 @@ describe("bounded repository checkpoint transfers", () => { let entries: WorkTreeEntry[] = []; const scan = vi.fn(async () => entries); const transport: WorkFolderTransport = { - home: async () => "/home/runner", scan, readBatch: undefined, + home: async () => "/home/runner", scan, + readBatch: batched ? vi.fn(async (_root, entries) => entries.map((entry) => contents.get(entry.path)!)) : undefined, read: vi.fn((_root, filePath) => { const source = Readable.from([contents.get(filePath)!]); sources.push(source); @@ -186,14 +187,61 @@ describe("bounded repository checkpoint transfers", () => { expect(f.sources.every((source) => source.destroyed)).toBe(true); }); - it("drains in-flight PUTs after failure without scheduling more blobs or replacing the protected checkpoint", async () => { - const f = await fixture(); + it("overlaps sixteen batch-readable blobs while keeping large streams bounded to four", async () => { + const f = await fixture(true); + const entries = Array.from({ length: 24 }, (_, i) => f.file(`small-${i}`)); + entries.splice(1, 0, ...Array.from({ length: 6 }, (_, i) => f.file(`large-${i}`, `${i}`.repeat(1024 * 1024 + 1)))); + const largeKeys = new Set(entries.filter((entry) => entry.byteSize > 1024 * 1024).map((entry) => entry.sha256)); + f.setEntries(entries); + const headGate = gate(), putGate = gate(); + const activeHeads = [0, 0], activePuts = [0, 0], maxHeads = [0, 0], maxPuts = [0, 0]; + const lane = (key: string) => largeKeys.has(key.split("/").at(-1)!) ? 1 : 0; + f.headObject.mockImplementation(async ({ objectKey }) => { + const index = lane(objectKey); + activeHeads[index]!++; + maxHeads[index] = Math.max(maxHeads[index]!, activeHeads[index]!); + try { await headGate.promise; return { exists: false }; } finally { activeHeads[index]!--; } + }); + const put = f.putObject.getMockImplementation()!; + f.putObject.mockImplementation(async (input) => { + if (!input.objectKey.includes("/blobs/")) return put(input); + const index = lane(input.objectKey); + activePuts[index]!++; + maxPuts[index] = Math.max(maxPuts[index]!, activePuts[index]!); + try { await putGate.promise; await put(input); } finally { activePuts[index]!--; } + }); + const saving = f.save(); + try { + await vi.waitFor(() => expect(activeHeads).toEqual([16, 4])); + expect(f.headObject).toHaveBeenCalledTimes(20); + headGate.release(); + await vi.waitFor(() => expect(activePuts).toEqual([16, 4])); + expect(f.headObject).toHaveBeenCalledTimes(20); + expect((await f.current()).checkpointKey).toBeNull(); + } finally { headGate.release(); putGate.release(); } + await saving; + expect(maxHeads).toEqual([16, 4]); + expect(maxPuts).toEqual([16, 4]); + expect(f.headObject).toHaveBeenCalledTimes(30); + const manifest = JSON.parse(f.objects.get(f.binding.checkpointKey!)!.toString("utf8")); + expect(manifest.files.map(({ objectKey: _key, ...entry }: WorkTreeEntry & { objectKey: string }) => entry)).toEqual(entries); + expect(f.transport.read).toHaveBeenCalledTimes(6); + expect(f.sources.every((source) => source.destroyed)).toBe(true); + for (const [, group] of vi.mocked(f.transport.readBatch!).mock.calls) { + expect(group.reduce((size, entry) => size + entry.byteSize, 0)).toBeLessThanOrEqual(1024 * 1024); + expect(group).toHaveLength(24); + } + }); + + it.each([false, true])("drains in-flight PUTs after failure without replacing the checkpoint (batched=%s)", async (batched) => { + const f = await fixture(batched); + const parallelism = batched ? 16 : 4; const original = f.file("saved"); f.setEntries([original]); await f.save(); const previous = await f.current(); const failed = f.file("fail"), queued = f.file("must-not-start"); - f.setEntries([original, failed, f.file("held-a"), f.file("held-b"), f.file("held-c"), queued]); + f.setEntries([original, failed, ...Array.from({ length: parallelism - 1 }, (_, i) => f.file(`held-${i}`)), queued]); const failedKey = `${companyId}/task-repositories/${f.binding.id}/blobs/${failed.sha256}`; const queuedKey = `${companyId}/task-repositories/${f.binding.id}/blobs/${queued.sha256}`; const fail = gate(), held = gate(); @@ -214,9 +262,9 @@ describe("bounded repository checkpoint transfers", () => { const outcome = f.save().then(() => ({ error: null }), (failure: unknown) => ({ error: failure })) .finally(() => { settled = true; }); try { - await vi.waitFor(() => expect(active).toBe(4)); + await vi.waitFor(() => expect(active).toBe(parallelism)); fail.release(); - await vi.waitFor(() => expect(active).toBe(3)); + await vi.waitFor(() => expect(active).toBe(parallelism - 1)); await setImmediate(); expect(settled).toBe(false); expect((await f.current()).checkpointKey).toBe(previous.checkpointKey); @@ -231,13 +279,14 @@ describe("bounded repository checkpoint transfers", () => { const tracked = await db.select().from(workFolderObjects).where(eq(workFolderObjects.repositoryBindingId, f.binding.id)); expect(tracked.find((object) => object.objectKey === previous.checkpointKey)!.deleteAfter).toBeNull(); expect(tracked.find((object) => object.objectKey.endsWith(`/blobs/${original.sha256}`))!.deleteAfter).toBeNull(); - expect(tracked.filter((object) => object.deleteAfter !== null)).toHaveLength(4); + expect(tracked.filter((object) => object.deleteAfter !== null)).toHaveLength(parallelism); }); - it("drains failed concurrent HEADs without opening more file streams", async () => { - const f = await fixture(); + it.each([false, true])("drains failed concurrent HEADs without opening streams (batched=%s)", async (batched) => { + const f = await fixture(batched); + const parallelism = batched ? 16 : 4; const first = f.file("fail-head"); - f.setEntries([first, f.file("two"), f.file("three"), f.file("four"), f.file("queued")]); + f.setEntries([first, ...Array.from({ length: parallelism - 1 }, (_, i) => f.file(`held-${i}`)), f.file("queued")]); const fail = gate(), held = gate(); const error = new Error("HEAD unavailable"); let active = 0, settled = false; @@ -252,16 +301,16 @@ describe("bounded repository checkpoint transfers", () => { const outcome = f.save().then(() => ({ error: null }), (failure: unknown) => ({ error: failure })) .finally(() => { settled = true; }); try { - await vi.waitFor(() => expect(active).toBe(4)); + await vi.waitFor(() => expect(active).toBe(parallelism)); fail.release(); - await vi.waitFor(() => expect(active).toBe(3)); + await vi.waitFor(() => expect(active).toBe(parallelism - 1)); await setImmediate(); expect(settled).toBe(false); expect(f.sources).toHaveLength(0); } finally { fail.release(); held.release(); } expect((await outcome).error).toBe(error); expect(active).toBe(0); - expect(f.headObject).toHaveBeenCalledTimes(4); + expect(f.headObject).toHaveBeenCalledTimes(parallelism); expect(f.sources).toHaveLength(0); expect(f.putObject).not.toHaveBeenCalled(); expect((await f.current()).checkpointKey).toBeNull(); diff --git a/server/src/services/sandbox-performance.ts b/server/src/services/sandbox-performance.ts index 3be762b871..da7c120caa 100644 --- a/server/src/services/sandbox-performance.ts +++ b/server/src/services/sandbox-performance.ts @@ -206,11 +206,12 @@ export async function runWithSandboxPerformanceTrace(input: { finally { trace.closed = true; // No per-file DB writes, no synchronous exporter call and no unbounded - // promise queue. Persist at most 250 records per batch after measured work. + // promise queue. The run-log sanitizer allows fifty array items. Keep each + // batch within that bound so retained records are not silently truncated. // The SDK independently exports ended spans through its batch processor. - for (let offset = 0; offset < trace.records.length; offset += 250) { + for (let offset = 0; offset < trace.records.length; offset += 50) { try { await input.onBatch?.({ schema: "paperclip.sandbox-performance.v1", runHash: trace.runHash, - records: trace.records.slice(offset, offset + 250), dropped: trace.dropped }); } catch { break; } + records: trace.records.slice(offset, offset + 50), dropped: trace.dropped }); } catch { break; } } } } diff --git a/server/src/services/work-folder-read-cache.ts b/server/src/services/work-folder-read-cache.ts index ed0edfafa7..cbcfaa8df7 100644 --- a/server/src/services/work-folder-read-cache.ts +++ b/server/src/services/work-folder-read-cache.ts @@ -2,14 +2,14 @@ import { captureSandboxPerformanceContext, measureSandboxOperation, measureSandb import { Readable } from "node:stream"; import type { WorkTreeEntry } from "./work-folder-transport.js"; -const MAX_BATCH_BYTES = 1024 * 1024; +export const WORK_FOLDER_READ_BATCH_MAX_BYTES = 1024 * 1024; const MAX_BATCH_ENTRIES = 64; const MAX_CACHED_BATCHES = 4; /** * A checkpoint-local cache, not a snapshot or a retry source. The caller limits - * concurrent readers to four. At most four 1 MiB batches stay cached; evicted - * batches held by those active readers can add at most another 4 MiB. Larger + * concurrent batch readers to sixteen. At most four 1 MiB batches stay cached; + * evicted batches held by active readers can add at most another 16 MiB. Larger * files use the separately bounded streaming fallback. Metadata is O(entries). */ export function createWorkFolderReadCache( @@ -21,8 +21,8 @@ export function createWorkFolderReadCache( let batch: WorkTreeEntry[] = []; let batchBytes = 0; for (const entry of entries) { - if (entry.kind !== "file" || entry.linkTarget || entry.byteSize > MAX_BATCH_BYTES) continue; - if (batch.length >= MAX_BATCH_ENTRIES || batchBytes + entry.byteSize > MAX_BATCH_BYTES) { + if (entry.kind !== "file" || entry.linkTarget || entry.byteSize > WORK_FOLDER_READ_BATCH_MAX_BYTES) continue; + if (batch.length >= MAX_BATCH_ENTRIES || batchBytes + entry.byteSize > WORK_FOLDER_READ_BATCH_MAX_BYTES) { batch = []; batchBytes = 0; } locations.set(entry.path, { batch, index: batch.length }); diff --git a/server/src/services/work-folder-repositories.ts b/server/src/services/work-folder-repositories.ts index 2157cbff8f..9523f8b098 100644 --- a/server/src/services/work-folder-repositories.ts +++ b/server/src/services/work-folder-repositories.ts @@ -1,5 +1,5 @@ import { measureSandboxOperation, measureSandboxStream } from "./sandbox-performance.js"; -import { createWorkFolderReadCache } from "./work-folder-read-cache.js"; +import { createWorkFolderReadCache, WORK_FOLDER_READ_BATCH_MAX_BYTES } from "./work-folder-read-cache.js"; import { prefetchWorkFiles } from "./work-folder-transfer.js"; import { createHash, randomUUID } from "node:crypto"; import { and, eq, inArray, isNull } from "drizzle-orm"; @@ -58,32 +58,42 @@ export function workFolderRepositoryService(db: Db, storage: StorageProvider, tr const readCache = transport.readBatch ? createWorkFolderReadCache([...unknown.values()], (entries) => transport.readBatch!(root, entries), (entry) => transport.read(root, entry.path, entry.byteSize)) : undefined; - let next = 0; + // Small, batch-readable objects mostly wait for object-store round trips. + // Give them a wider bounded lane without multiplying large remote streams. + // With no batch transport, all objects retain the four-stream limit. + const indexedObjects = objects.map((object, fileIndex) => ({ object, fileIndex })); + const lanes = [ + { parallelism: 16, objects: indexedObjects.filter(({ object: [, entry] }) => readCache && entry.byteSize <= WORK_FOLDER_READ_BATCH_MAX_BYTES) }, + { parallelism: 4, objects: indexedObjects.filter(({ object: [, entry] }) => !readCache || entry.byteSize > WORK_FOLDER_READ_BATCH_MAX_BYTES) }, + ]; let failure: { error: unknown } | undefined; try { - await Promise.all(Array.from({ length: Math.min(4, objects.length) }, async () => { - while (!failure) { - const fileIndex = next++; - const object = objects[fileIndex]; - if (!object) return; - const [objectKey, entry] = object; - try { - await measureSandboxOperation("work_folder.repository.object_intent", { fileIndex }, () => registerWorkFolderObject(db, storage, { objectKey, companyId: binding.companyId, repositoryBindingId: binding.id })); - if (failure) return; - const { exists } = await measureSandboxOperation("work_folder.repository.object_head", { fileIndex, bytes: entry.byteSize, requestCount: 1 }, async (span) => { - const result = await storage.headObject({ objectKey }); span.set({ exists: result.exists }); return result; - }); - if (!exists && !failure) { - await measureSandboxOperation("work_folder.repository.object_upload", { fileIndex, bytes: entry.byteSize, parallelism: 4 }, () => uploadWorkFolderObject(storage, { objectKey, contentType: "application/octet-stream", - contentLength: entry.byteSize, sha256: entry.sha256!, - createSource: () => readCache ? readCache.read(entry) : transport.read(root, entry.path, entry.byteSize) })); + await Promise.all(lanes.map(async (lane) => { + let next = 0; + await Promise.all(Array.from({ length: Math.min(lane.parallelism, lane.objects.length) }, async () => { + while (!failure) { + const item = lane.objects[next++]; + if (!item) return; + const { object, fileIndex } = item; + const [objectKey, entry] = object; + try { + await measureSandboxOperation("work_folder.repository.object_intent", { fileIndex }, () => registerWorkFolderObject(db, storage, { objectKey, companyId: binding.companyId, repositoryBindingId: binding.id })); + if (failure) return; + const { exists } = await measureSandboxOperation("work_folder.repository.object_head", { fileIndex, bytes: entry.byteSize, requestCount: 1 }, async (span) => { + const result = await storage.headObject({ objectKey }); span.set({ exists: result.exists }); return result; + }); + if (!exists && !failure) { + await measureSandboxOperation("work_folder.repository.object_upload", { fileIndex, bytes: entry.byteSize, parallelism: lane.parallelism }, () => uploadWorkFolderObject(storage, { objectKey, contentType: "application/octet-stream", + contentLength: entry.byteSize, sha256: entry.sha256!, + createSource: () => readCache ? readCache.read(entry) : transport.read(root, entry.path, entry.byteSize) })); + } + } catch (error) { + // Stop scheduling after the first error, but drain the other workers + // before returning. Their streaming PUTs must not outlive this save. + failure ??= { error }; } - } catch (error) { - // Stop scheduling after the first error, but drain the other workers - // before returning. Their streaming PUTs must not outlive this save. - failure ??= { error }; } - } + })); })); } finally { readCache?.clear(); } if (failure) throw failure.error;