Files
PaperClipAI/server/src/__tests__/plugin-host-services-span.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

332 lines
12 KiB
TypeScript

import { beforeEach, describe, expect, it, vi } from "vitest";
import { createHostClientHandlers } from "../../../packages/plugins/sdk/src/host-client-factory.js";
import type { WorkerHostCallContext } from "../../../packages/plugins/sdk/src/protocol.js";
import { SANDBOX_STARTUP_SPAN_ATTRS as A } from "@paperclipai/adapter-utils/acpx-engine/startup-timing";
import {
buildHostServices,
clampProviderSpanAttributes,
parseTraceparent,
} from "../services/plugin-host-services.js";
// Capture every span the host trust boundary hands to the real tracer.
const mockRecordSpan = vi.hoisted(() => vi.fn());
vi.mock("../instrumentation.js", () => ({
recordProviderPluginSpan: mockRecordSpan,
traceparentFromContextToken: () => undefined,
}));
function createEventBusStub() {
return {
forPlugin() {
return { emit: vi.fn(), subscribe: vi.fn() };
},
} as never;
}
// A well-formed W3C traceparent (the host mints it; the handler validates it).
const VALID_TRACEPARENT = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01";
function servicesFor() {
return buildHostServices({} as never, "plugin-record-id", "daytona", createEventBusStub());
}
function handlersFor(capabilities: readonly string[]) {
return createHostClientHandlers({
pluginId: "daytona",
capabilities: capabilities as never,
services: servicesFor(),
});
}
describe("plugin provider span host handler", () => {
beforeEach(() => {
mockRecordSpan.mockReset();
});
it("records a span with the clamped name and the allowlisted attributes", async () => {
const services = servicesFor();
await services.tracer.record(
{
name: "pack",
attributes: {
[A.provider]: "daytona",
[A.packWallMs]: 12,
},
},
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
expect(mockRecordSpan).toHaveBeenCalledTimes(1);
const call = mockRecordSpan.mock.calls[0]![0] as {
name: string;
parent: { traceId: string; spanId: string; traceFlags: number };
attributes: Record<string, unknown>;
};
expect(call.name).toBe("sandbox.daytona.pack");
expect(call.parent).toEqual({
traceId: "0af7651916cd43dd8448eb211c80319c",
spanId: "b7ad6b7169203331",
traceFlags: 1,
});
expect(call.attributes[A.provider]).toBe("daytona");
expect(call.attributes[A.packWallMs]).toBe(12);
});
it("drops every forbidden attribute before the span reaches the tracer", async () => {
const services = servicesFor();
await services.tracer.record(
{
name: "transfer",
attributes: {
[A.transferGuardCount]: 2,
// Forbidden fields that must never ride a span.
[A.execCommand]: "bash",
command: "rm -rf /",
args: "--force",
stdout: "secret output",
stderr: "secret error",
path: "/etc/passwd",
extra: "leak",
},
},
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
expect(mockRecordSpan).toHaveBeenCalledTimes(1);
const attributes = (mockRecordSpan.mock.calls[0]![0] as { attributes: Record<string, unknown> })
.attributes;
// Only the allowlisted key survives; every forbidden key is dropped.
expect(attributes).toEqual({ [A.transferGuardCount]: 2 });
for (const forbidden of [A.execCommand, "command", "args", "stdout", "stderr", "path", "extra"]) {
expect(attributes).not.toHaveProperty(forbidden);
}
});
it("drops a status message (it could carry standard-stream text)", async () => {
const services = servicesFor();
await services.tracer.record(
{ name: "transfer", status: { code: 2, message: "secret error text" } },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
const call = mockRecordSpan.mock.calls[0]![0] as { status?: { code: number; message?: string } };
expect(call.status).toEqual({ code: 2 });
expect(call.status).not.toHaveProperty("message");
});
it("clamps an unknown span name to sandbox.daytona.other", async () => {
const services = servicesFor();
await services.tracer.record(
{ name: "rm -rf / --no-preserve-root" },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
expect((mockRecordSpan.mock.calls[0]![0] as { name: string }).name).toBe(
"sandbox.daytona.other",
);
});
it("admits each per-round-trip span name to sandbox.daytona.<name>", async () => {
const services = servicesFor();
for (const name of [
"ensureDirectory",
"checkSymlinkEscape",
"promote",
"extractTarball",
"postUploadCommand",
]) {
await services.tracer.record(
{ name },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
}
const recorded = mockRecordSpan.mock.calls.map((c) => (c[0] as { name: string }).name);
expect(recorded).toEqual([
"sandbox.daytona.ensureDirectory",
"sandbox.daytona.checkSymlinkEscape",
"sandbox.daytona.promote",
"sandbox.daytona.extractTarball",
"sandbox.daytona.postUploadCommand",
]);
});
it("admits the session open and close span names", async () => {
const services = servicesFor();
for (const name of ["session.open", "session.close"]) {
await services.tracer.record(
{ name },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
}
const recorded = mockRecordSpan.mock.calls.map((c) => (c[0] as { name: string }).name);
expect(recorded).toEqual([
"sandbox.daytona.session.open",
"sandbox.daytona.session.close",
]);
});
it("forwards a valid start-time and end-time pair to the recorder", async () => {
const services = servicesFor();
const startTimeMs = Date.now() - 4500;
const endTimeMs = startTimeMs + 4500;
await services.tracer.record(
{ name: "ensureDirectory", startTimeMs, endTimeMs },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
const call = mockRecordSpan.mock.calls[0]![0] as {
startTimeMs?: number;
endTimeMs?: number;
};
expect(call.startTimeMs).toBe(startTimeMs);
expect(call.endTimeMs).toBe(endTimeMs);
});
it("accepts a pair whose end is a small skew ahead of the host clock", async () => {
const services = servicesFor();
// The worker clock leads the host clock by a few seconds. This small skew
// is within the allowed bound, so the host keeps the native width.
const startTimeMs = Date.now() + 5000;
const endTimeMs = startTimeMs + 1000;
await services.tracer.record(
{ name: "ensureDirectory", startTimeMs, endTimeMs },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
const call = mockRecordSpan.mock.calls[0]![0] as {
startTimeMs?: number;
endTimeMs?: number;
};
expect(call.startTimeMs).toBe(startTimeMs);
expect(call.endTimeMs).toBe(endTimeMs);
});
it("drops an invalid timestamp pair so the synchronous path runs", async () => {
const services = servicesFor();
const now = Date.now();
const invalidPairs: Array<{ startTimeMs?: unknown; endTimeMs?: unknown; why: string }> = [
{ startTimeMs: now, endTimeMs: now - 1000, why: "reversed order" },
{ startTimeMs: Number.NaN, endTimeMs: now, why: "non-finite start" },
{ startTimeMs: now, endTimeMs: Number.POSITIVE_INFINITY, why: "non-finite end" },
{ startTimeMs: now, endTimeMs: now + 11 * 60 * 1000, why: "over-ceiling duration" },
{ startTimeMs: now - 2 * 60 * 60 * 1000, endTimeMs: now - 2 * 60 * 60 * 1000 + 10, why: "over-age start" },
{ startTimeMs: now, endTimeMs: now + 2 * 60 * 1000, why: "end far in the future" },
{ startTimeMs: now + 5 * 60 * 1000, endTimeMs: now + 5 * 60 * 1000 + 10, why: "start and end in the future" },
];
for (const pair of invalidPairs) {
mockRecordSpan.mockReset();
await services.tracer.record(
{ name: "ensureDirectory", startTimeMs: pair.startTimeMs, endTimeMs: pair.endTimeMs } as never,
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
// The span still records (the synchronous path), but without a timestamp.
const call = mockRecordSpan.mock.calls[0]![0] as {
startTimeMs?: number;
endTimeMs?: number;
};
expect(call.startTimeMs, pair.why).toBeUndefined();
expect(call.endTimeMs, pair.why).toBeUndefined();
}
});
it("records the synchronous path when the timestamp pair is absent", async () => {
const services = servicesFor();
await services.tracer.record(
{ name: "pack" },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
const call = mockRecordSpan.mock.calls[0]![0] as {
startTimeMs?: number;
endTimeMs?: number;
};
expect(call.startTimeMs).toBeUndefined();
expect(call.endTimeMs).toBeUndefined();
});
it("rejects a malformed traceparent — no span is recorded", async () => {
const services = servicesFor();
for (const bad of [
undefined,
"not-a-traceparent",
"00-xyz-b7ad6b7169203331-01",
"00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331", // missing flags
"00-00000000000000000000000000000000-b7ad6b7169203331-01", // all-zero trace id
]) {
await services.tracer.record(
{ name: "pack" },
{ traceparent: bad } as WorkerHostCallContext,
);
}
expect(mockRecordSpan).not.toHaveBeenCalled();
});
it("rejects a span from a plugin that lacks the environment-driver capability", async () => {
const handlers = handlersFor([]);
await expect(
handlers["span.record"](
{ name: "pack" },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
),
).rejects.toThrow(/capabilit/i);
expect(mockRecordSpan).not.toHaveBeenCalled();
});
it("admits a span from a plugin that holds the environment-driver capability", async () => {
const handlers = handlersFor(["environment.drivers.register"]);
await handlers["span.record"](
{ name: "pack", attributes: { [A.provider]: "daytona" } },
{ traceparent: VALID_TRACEPARENT } as WorkerHostCallContext,
);
expect(mockRecordSpan).toHaveBeenCalledTimes(1);
});
});
describe("parseTraceparent", () => {
it("accepts a well-formed traceparent and returns the parts", () => {
expect(parseTraceparent(VALID_TRACEPARENT)).toEqual({
traceId: "0af7651916cd43dd8448eb211c80319c",
spanId: "b7ad6b7169203331",
traceFlags: 1,
});
});
it("rejects malformed, all-zero, and forbidden-version traceparents", () => {
expect(parseTraceparent(undefined)).toBeNull();
expect(parseTraceparent("garbage")).toBeNull();
expect(parseTraceparent("00-0af7651916cd43dd8448eb211c80319c-0000000000000000-01")).toBeNull();
expect(parseTraceparent("ff-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01")).toBeNull();
});
});
describe("clampProviderSpanAttributes", () => {
it("keeps only allowlisted keys and normalizes the provider family", () => {
expect(
clampProviderSpanAttributes({
[A.provider]: "some-operator-key",
[A.packWallMs]: 7,
[A.transferWallMs]: Number.NaN,
[A.execCommand]: "bash",
}),
).toEqual({
// An unknown provider key maps to `plugin`, never the raw key.
[A.provider]: "plugin",
[A.packWallMs]: 7,
// A non-finite number yields no attribute; `exec.command` is not allowed.
});
});
it("keeps the transfer direction for each closed-set value", () => {
expect(clampProviderSpanAttributes({ [A.transferDirection]: "inbound" })).toEqual({
[A.transferDirection]: "inbound",
});
expect(clampProviderSpanAttributes({ [A.transferDirection]: "outbound" })).toEqual({
[A.transferDirection]: "outbound",
});
});
it("drops a transfer direction outside the closed set or of the wrong type", () => {
// A free-form string, an empty string, and a non-string all yield no
// attribute, so the direction stays bounded and low-cardinality.
expect(clampProviderSpanAttributes({ [A.transferDirection]: "sideways" })).toEqual({});
expect(clampProviderSpanAttributes({ [A.transferDirection]: "" })).toEqual({});
expect(clampProviderSpanAttributes({ [A.transferDirection]: 1 })).toEqual({});
});
});