Files
PaperClipAI/packages/adapter-utils/src/sync-operation-schedule.test.ts
T
Nicky LeachandPaperclip 233be4b36c feat: parallelize sandbox file-sync behind a provider opt-in capability (#11736)
## Thinking Path

> - Paperclip runs AI agents through local and remote execution
adapters.
> - Sandbox providers move workspace and asset files before and after
agent runs.
> - Serial file transfers delay startup and teardown when several
operations do not depend on each other.
> - Providers need an opt-in contract so existing providers keep their
serial behavior.
> - This pull request adds a bounded scheduler and routes inbound and
outbound sync operations through it.
> - The benefit is shorter sandbox setup and teardown with stable
errors, clear telemetry, and a safe opt-in path.

## Linked Issues or Issue Description

**Subsystem affected**

Cross-cutting (multiple of the above): packages/shared,
packages/adapter-utils, packages/plugins, and server.

**Problem or motivation**

Sandbox sync processes the workspace, assets, and referenced projects in
series. This adds avoidable wait time to agent startup and teardown.

**Proposed solution**

Add a fail-closed provider capability named concurrentSyncOperations.
Use a bounded scheduler with a limit of four operations. Preserve
operation order for error reporting. Keep non-opted-in providers on the
serial path.

**Alternatives considered**

Increase the serial transfer speed or add provider-specific schedulers.
Those options do not provide one shared contract or stable behavior
across providers.

**Roadmap alignment**

ROADMAP.md lists cloud and sandbox agents as a product area. This change
improves sandbox execution without changing the control-plane contract.

**Additional context**

The Daytona provider opts in. Board trials on this commit showed overlap
for inbound sync and outbound restore, with no referenced-project
staging failures.

## What Changed

- Add the concurrentSyncOperations sandbox capability and fail-closed
parsing.
- Add a bounded settle-all scheduler with stable input-order errors.
- Parallelize inbound workspace, asset, and referenced-project sync
operations when the provider opts in.
- Parallelize outbound workspace and asset restore operations when the
provider opts in.
- Surface referenced-project failure text in run logs and server
telemetry.
- Add Daytona sync spans and the capability declaration.
- Preserve in-flight upload scratch tarballs during workspace wipe.
- Add unit and regression tests for the scheduler, coordinators,
provider behavior, telemetry, and wipe race.

## Verification

- Run the adapter-utils and server type checks.
- Run the targeted adapter-utils, server, and Daytona test suites.
- Run the full automated sweep.
- Review six cold Daytona trials, with three serial and three parallel
runs.
- Confirm that parallel trials show inbound overlap and outbound restore
overlap.
- Confirm that providers without the capability keep serial behavior.

## Risks

- Providers must opt in only when their file operations can run safely
at the same time.
- A provider that declares the capability incorrectly can expose
transfer races.
- The scheduler keeps a limit of four to bound resource use.
- Providers without the capability keep the prior serial behavior.

## Model Used

OpenAI GPT-5 in the Codex runtime. The model used tool calls, code
inspection, and GitHub workflow support. The model did not author the
implementation commits.

## Checklist

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

---------

Co-authored-by: Paperclip <noreply@paperclip.ing>
2026-08-19 12:34:11 -07:00

167 lines
5.6 KiB
TypeScript

import { describe, expect, it } from "vitest";
import {
SYNC_OPERATION_CONCURRENCY_LIMIT,
scheduleSyncOperations,
} from "./sync-operation-schedule.js";
// A deferred task fake. The `task` thunk records that it started and returns a
// promise that the test resolves or rejects by hand. The fake lets a test hold
// a task open and check the exact moment a later task starts.
interface DeferredTask<T> {
readonly task: () => Promise<T>;
resolve(value: T): void;
reject(reason: unknown): void;
started(): boolean;
}
function makeDeferred<T>(): DeferredTask<T> {
let started = false;
let resolveFn!: (value: T) => void;
let rejectFn!: (reason: unknown) => void;
const promise = new Promise<T>((resolve, reject) => {
resolveFn = resolve;
rejectFn = reject;
});
return {
task: () => {
started = true;
return promise;
},
resolve: (value: T) => resolveFn(value),
reject: (reason: unknown) => rejectFn(reason),
started: () => started,
};
}
// Drain the microtask and timer queues so every pending worker step runs.
function flush(): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, 0));
}
describe("scheduleSyncOperations", () => {
it("exposes the shared bound constant as 4", () => {
expect(SYNC_OPERATION_CONCURRENCY_LIMIT).toBe(4);
});
it("keeps at most `bound` tasks active in concurrent mode", async () => {
const deferreds = [makeDeferred<number>(), makeDeferred<number>(), makeDeferred<number>()];
const tasks = deferreds.map((deferred) => deferred.task);
const scheduled = scheduleSyncOperations(tasks, true, 2);
await flush();
// The bound is 2, so only the first two tasks start.
expect(deferreds[0].started()).toBe(true);
expect(deferreds[1].started()).toBe(true);
expect(deferreds[2].started()).toBe(false);
// One active task settles. A worker frees, so the third task starts.
deferreds[0].resolve(0);
await flush();
expect(deferreds[2].started()).toBe(true);
deferreds[1].resolve(1);
deferreds[2].resolve(2);
const results = await scheduled;
expect(results).toEqual([
{ status: "fulfilled", value: 0 },
{ status: "fulfilled", value: 1 },
{ status: "fulfilled", value: 2 },
]);
});
it("runs one task at a time in input order in serial mode", async () => {
const deferreds = [makeDeferred<number>(), makeDeferred<number>(), makeDeferred<number>()];
const tasks = deferreds.map((deferred) => deferred.task);
const scheduled = scheduleSyncOperations(tasks, false);
await flush();
// Serial mode holds one task active, so only the first task starts.
expect(deferreds[0].started()).toBe(true);
expect(deferreds[1].started()).toBe(false);
expect(deferreds[2].started()).toBe(false);
deferreds[0].resolve(0);
await flush();
// The next task starts only after the previous task settles.
expect(deferreds[1].started()).toBe(true);
expect(deferreds[2].started()).toBe(false);
deferreds[1].resolve(1);
await flush();
expect(deferreds[2].started()).toBe(true);
deferreds[2].resolve(2);
const results = await scheduled;
expect(results).toEqual([
{ status: "fulfilled", value: 0 },
{ status: "fulfilled", value: 1 },
{ status: "fulfilled", value: 2 },
]);
});
it("waits for every started task to settle when one rejects", async () => {
const first = makeDeferred<number>();
const second = makeDeferred<number>();
const tasks = [first.task, second.task];
const scheduled = scheduleSyncOperations(tasks, true, 2);
let settled = false;
void scheduled.then(() => {
settled = true;
});
await flush();
// Both tasks are active under the bound of 2.
expect(first.started()).toBe(true);
expect(second.started()).toBe(true);
// The first task rejects, but the second task stays open.
first.reject(new Error("first failed"));
await flush();
// The scheduler does not return before every started task settles.
expect(settled).toBe(false);
second.resolve(2);
const results = await scheduled;
expect(settled).toBe(true);
// The results keep input order, and the rejection carries its reason.
expect(results[0]).toEqual({ status: "rejected", reason: new Error("first failed") });
expect(results[1]).toEqual({ status: "fulfilled", value: 2 });
});
// One table proves settle-all and input-order for both call modes. Both future
// call sites share these proven cases.
const MODE_TABLE = [
{ name: "serial mode", concurrent: false, bound: SYNC_OPERATION_CONCURRENCY_LIMIT },
{ name: "concurrent mode", concurrent: true, bound: 2 },
];
for (const mode of MODE_TABLE) {
it(`returns settled results in input order in ${mode.name}`, async () => {
// The tasks settle out of order: index 2 first, then index 0, then a
// rejection at index 1. The result array must still keep input order.
const failure = new Error("index one failed");
const tasks: Array<() => Promise<string>> = [
() => Promise.resolve("zero"),
() => Promise.reject(failure),
() => Promise.resolve("two"),
];
const results = await scheduleSyncOperations(tasks, mode.concurrent, mode.bound);
expect(results).toEqual([
{ status: "fulfilled", value: "zero" },
{ status: "rejected", reason: failure },
{ status: "fulfilled", value: "two" },
]);
});
}
it("returns an empty result list for no tasks", async () => {
const results = await scheduleSyncOperations<number>([], true, 4);
expect(results).toEqual([]);
});
});