perf: overlap small repository checkpoint uploads within bounded lanes
This commit is contained in:
parent
bc4e39d2f7
commit
a7108016b3
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -206,11 +206,12 @@ export async function runWithSandboxPerformanceTrace<T>(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; }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 });
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
Loading…
Reference in New Issue