refactor(adapter-utils): remove the retired duplex_v1 sandbox bridge transport (#12171)

## Thinking Path

> - Paperclip is the open source app people use to manage AI agents for
work
> - The adapter utilities provide sandbox transport paths for agent
execution
> - The retired `duplex_v1` path remains in host, gateway, and test code
after `http2_v1` replaced it
> - Retired transport code adds maintenance cost and leaves an unsafe
fallback for unknown gateway modes
> - This pull request removes the retired path, moves shared `http2_v1`
contracts to a leaf module, and closes mode dispatch to a fixed
allowlist
> - The benefit is a smaller transport surface and explicit failure for
unsupported modes

## Linked Issues or Issue Description

Refs #12120

The `http2_v1` transport replaced `duplex_v1`, but the retired broker,
gateway, constants, and tests remain in the adapter utilities. An
unknown bridge mode can also fall through to the queue gateway when a
queue directory exists. This change removes the retired code and rejects
unsupported modes before gateway selection.

## What Changed

- Delete the host `duplex_v1` broker and its transport-only tests.
- Delete the in-sandbox duplex gateway and retired mode constants.
- Move shared `http2_v1` symbols into `bridge-transport-contract.ts`.
- Update the remaining importers and repair their focused tests.
- Validate bridge modes against `http2_v1` and `queue_v1` before queue
lookup.
- Keep `queue_v1`, `duplex-frame-codec.ts`, and duplex telemetry
dimensions unchanged.

## Verification

- [x] `npx tsc --noEmit -p packages/adapter-utils` passes.
- [x] `npx vitest run packages/adapter-utils/src` passes: 48 files and
968 tests pass, with 4 pre-existing platform skips.
- [x] Full CI is green on this pull request.
- [x] Greptile review is complete and every finding is resolved.

## Risks

The change removes an internal transport that no host path selects. The
main risk is an overlooked import or test dependency. Targeted typecheck
and tests cover the adapter utility package. Full CI must confirm
workspace-wide compatibility.

## Model Used

Anthropic Claude Sonnet 5 assisted with the implementation, as recorded
in the commit. The commit does not record a context-window size or
reasoning mode.

## 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:
Nicky Leach 2026-08-25 08:47:43 -07:00 committed by GitHub
parent 02a984068c
commit 2862e18484
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
12 changed files with 276 additions and 5158 deletions

View File

@ -36,8 +36,6 @@ import {
type AcpxEngineExecutorOptions,
} from "./execute.js";
import { runChildProcess } from "../server-utils.js";
import { createDuplexBridgeBroker } from "../duplex-bridge-broker.js";
import type { CommandManagedDuplexChannel } from "../command-managed-runtime.js";
import {
getActiveStepContext,
runWithRuntimeParent,
@ -5734,56 +5732,55 @@ describe("ACPX engine run lifecycle corrections (F3: one teardown error policy)"
});
});
describe("ACPX engine sandbox duplex run-disposition seam (fail-closed)", () => {
describe("ACPX engine sandbox bridge run-disposition seam (fail-closed)", () => {
beforeEach(() => {
vi.clearAllMocks();
});
// A minimal in-memory duplex channel. The test drives a channel exit to latch a
// real loss in a real broker, so the seam reads a real run disposition.
function createFakeDuplexChannel() {
let exitListener: ((exit: { exitCode: number | null }) => void) | null = null;
const channel: CommandManagedDuplexChannel = {
write(): void {},
onData(): void {},
onExit(listener: (exit: { exitCode: number | null }) => void): void {
exitListener = listener;
},
stop(): void {},
close(): Promise<void> {
return Promise.resolve();
},
};
return {
channel,
emitExit: (exit: { exitCode: number | null }) => exitListener?.(exit),
};
}
// Wrap a started real broker in a paperclip bridge handle. The handle exposes
// the same run-disposition surface the sandbox bridge exposes, so the seam runs
// against the real latch and the real orderly-completion mark.
async function bridgeOverBroker(fake: ReturnType<typeof createFakeDuplexChannel>) {
const broker = await createDuplexBridgeBroker({
channel: fake.channel,
forwardRequest: async () => ({ status: 200 }),
/**
* A minimal run-disposition latch. It reproduces the same ordering rule the
* bridge transport applies: the first ordered loss or orderly completion
* latches the terminal disposition, and a later call never overrides it.
* The test drives the latch directly, with no channel and no live process.
*/
function createFakeBridgeHandle() {
let lossOrdered = false;
let lossReason: string | null = null;
let completionOrdered = false;
const readDisposition = () => ({ failed: lossOrdered, lossReason });
const markOrderlyCompletion = vi.fn(() => {
if (completionOrdered || lossOrdered) return;
completionOrdered = true;
});
const settleRunDisposition = vi.fn(() => {
markOrderlyCompletion();
return readDisposition();
});
broker.start();
const markOrderlyCompletion = vi.fn(() => broker.markOrderlyCompletion());
const settleRunDisposition = vi.fn(() => broker.settleRunDisposition());
const stop = vi.fn(async () => {});
const handle = {
env: {
PAPERCLIP_API_URL: "http://127.0.0.1:1",
PAPERCLIP_API_KEY: "bridge-token",
PAPERCLIP_API_BRIDGE_MODE: "duplex_v1",
PAPERCLIP_API_BRIDGE_MODE: "http2_v1",
},
readRunDisposition: () => broker.runDisposition,
readRunDisposition: () => readDisposition(),
settleRunDisposition,
markOrderlyCompletion,
stop,
};
return { broker, handle, markOrderlyCompletion, settleRunDisposition };
return {
handle,
markOrderlyCompletion,
settleRunDisposition,
readDisposition,
// Record the first ordered loss. A loss ordered after a completion, or a
// second loss, is a no-op — the same rule the real transport applies.
emitLoss: (reason: string) => {
if (lossOrdered || completionOrdered) return;
lossOrdered = true;
lossReason = reason;
},
};
}
// A runtime whose one turn completes cleanly. The `beforeResult` hook runs at
@ -5888,14 +5885,13 @@ describe("ACPX engine sandbox duplex run-disposition seam (fail-closed)", () =>
} as never);
}
it("fails a completed run when the duplex channel was lost before the completion", async () => {
it("fails a completed run when the bridge channel was lost before the completion", async () => {
const sandbox = await setupRemoteSandbox();
const fake = createFakeDuplexChannel();
const { broker, handle, settleRunDisposition } = await bridgeOverBroker(fake);
const fake = createFakeBridgeHandle();
// Latch the loss before the ACP terminal resolves.
const runtime = runtimeWithControlledResult(() => fake.emitExit({ exitCode: 1 }));
const runtime = runtimeWithControlledResult(() => fake.emitLoss("provider_exit"));
const result = await runRemote(handle, runtime, sandbox);
const result = await runRemote(fake.handle, runtime, sandbox);
// The lost channel overrides the nominally completed terminal to a failure.
expect(result.exitCode).not.toBe(0);
@ -5905,65 +5901,62 @@ describe("ACPX engine sandbox duplex run-disposition seam (fail-closed)", () =>
expect(result.resultJson).toMatchObject({ status: "failed" });
// The seam read the disposition through the atomic settle step, and the
// latched loss kept the failure, so no orderly completion ordered.
expect(settleRunDisposition).toHaveBeenCalledTimes(1);
expect(broker.runDisposition.failed).toBe(true);
expect(fake.settleRunDisposition).toHaveBeenCalledTimes(1);
expect(fake.readDisposition().failed).toBe(true);
});
it("keeps a completed run a success when the channel stays live, and a later teardown loss is benign", async () => {
const sandbox = await setupRemoteSandbox();
const fake = createFakeDuplexChannel();
const { broker, handle, settleRunDisposition } = await bridgeOverBroker(fake);
const fake = createFakeBridgeHandle();
// No loss before the completion.
const runtime = runtimeWithControlledResult();
const result = await runRemote(handle, runtime, sandbox);
const result = await runRemote(fake.handle, runtime, sandbox);
expect(result.exitCode).toBe(0);
expect(result.errorCode ?? null).toBeNull();
// The atomic settle step marked the orderly completion for the
// success-eligible terminal.
expect(settleRunDisposition).toHaveBeenCalledTimes(1);
expect(fake.settleRunDisposition).toHaveBeenCalledTimes(1);
// A teardown loss ordered after the orderly completion is a normal teardown,
// so the run disposition stays a success.
fake.emitExit({ exitCode: 0 });
expect(broker.runDisposition.failed).toBe(false);
fake.emitLoss("provider_exit");
expect(fake.readDisposition().failed).toBe(false);
});
it("does not let a later completion or activity clear the loss latch", async () => {
const sandbox = await setupRemoteSandbox();
const fake = createFakeDuplexChannel();
const { broker, handle } = await bridgeOverBroker(fake);
const fake = createFakeBridgeHandle();
// Latch the loss before the ACP terminal resolves.
const runtime = runtimeWithControlledResult(() => fake.emitExit({ exitCode: 1 }));
const runtime = runtimeWithControlledResult(() => fake.emitLoss("provider_exit"));
const result = await runRemote(handle, runtime, sandbox);
const result = await runRemote(fake.handle, runtime, sandbox);
expect(result.exitCode).not.toBe(0);
expect(result.errorCode).toBe("duplex_channel_lost");
// A later orderly-completion mark and further channel activity cannot clear
// the latched loss.
broker.markOrderlyCompletion();
fake.emitExit({ exitCode: 0 });
expect(broker.runDisposition.failed).toBe(true);
expect(broker.runDisposition.lossReason).toBe("provider_exit");
fake.markOrderlyCompletion();
fake.emitLoss("provider_exit");
expect(fake.readDisposition().failed).toBe(true);
expect(fake.readDisposition().lossReason).toBe("provider_exit");
});
it("marks an orderly completion on a failed terminal so the teardown loss emits no false loss", async () => {
const sandbox = await setupRemoteSandbox();
const fake = createFakeDuplexChannel();
const { broker, handle, markOrderlyCompletion } = await bridgeOverBroker(fake);
const fake = createFakeBridgeHandle();
// The turn fails, and no channel loss ordered before the finalization.
const runtime = runtimeWithFailedResult();
const result = await runRemote(handle, runtime, sandbox);
const result = await runRemote(fake.handle, runtime, sandbox);
// The failed terminal stays a failure, but not a duplex loss.
// The failed terminal stays a failure, but not a bridge-channel loss.
expect(result.exitCode).not.toBe(0);
expect(result.errorCode).not.toBe("duplex_channel_lost");
// The non-success-eligible terminal marked the orderly completion, so the
// teardown channel_exit orders after the mark and does not latch a loss.
expect(markOrderlyCompletion).toHaveBeenCalledTimes(1);
fake.emitExit({ exitCode: 0 });
expect(broker.runDisposition.failed).toBe(false);
// teardown loss orders after the mark and does not latch a loss.
expect(fake.markOrderlyCompletion).toHaveBeenCalledTimes(1);
fake.emitLoss("provider_exit");
expect(fake.readDisposition().failed).toBe(false);
});
});

View File

@ -34,7 +34,7 @@ import {
type SandboxAdditionalSource,
} from "@paperclipai/adapter-utils/execution-target";
import type { DuplexLossReason } from "../duplex-observability.js";
import { DUPLEX_CHANNEL_LOST_ERROR_CODE } from "../duplex-bridge-broker.js";
import { DUPLEX_CHANNEL_LOST_ERROR_CODE } from "../bridge-transport-contract.js";
import {
DEFAULT_PAPERCLIP_AGENT_PROMPT_TEMPLATE,
applyPaperclipWorkspaceEnv,

View File

@ -0,0 +1,61 @@
/**
* Shared contract for the sandbox callback bridge transports.
*
* The retired duplex_v1 broker first defined these symbols. The host
* broker is gone, but the http2_v1 transport and the ACPX engine
* run-disposition seam still use them. This leaf module holds the
* survivors, so a caller of the run-disposition seam does not import the
* whole HTTP/2 bridge server module graph to reach one error code.
*/
import type { DuplexLossReason } from "./duplex-observability.js";
/**
* The typed error code the host reports when the bridge control channel
* died before an orderly completion. Both the ACP lane and the CLI lane
* report this one code, so the run disposition is identical across the two
* lanes.
*/
export const DUPLEX_CHANNEL_LOST_ERROR_CODE = "duplex_channel_lost";
/**
* The terminal run disposition a bridge transport computes from its ordered
* lifecycle. A `failed` disposition means a terminal loss ordered before an
* orderly completion, so the run must not report success. The typed loss
* reason names the cause; it is `null` for a success.
*/
export interface DuplexBrokerRunDisposition {
/** True when a terminal loss ordered before an orderly completion. */
failed: boolean;
/** The typed, closed loss reason on a failure; `null` on a success. */
lossReason: DuplexLossReason | null;
}
/** The nested timeout budgets. Each inner budget is smaller than its outer budget. */
export interface DuplexBrokerBudgets {
/** The deadline for one forward call, in milliseconds. */
forwardTimeoutMs: number;
/** The deadline for the broker to send one response frame, in milliseconds. */
responseBudgetMs: number;
/** The deadline the in-sandbox gateway waits for the response frame, in milliseconds. */
gatewayWaitMs: number;
}
/** The default nested budgets: forward 30 s, response 32 s, gateway wait 35 s. */
export const DEFAULT_DUPLEX_BROKER_BUDGETS: DuplexBrokerBudgets = {
forwardTimeoutMs: 30_000,
responseBudgetMs: 32_000,
gatewayWaitMs: 35_000,
};
/**
* The safe HTTP methods. RFC 7231 section 4.2.1 defines this set. A safe method
* does not change host state, so the host applies no mutation for it. A caller
* can retry a safe method after a forward failure without a double-apply risk.
*/
const SAFE_BRIDGE_METHODS = new Set(["GET", "HEAD", "OPTIONS", "TRACE"]);
/** Report whether the method is safe, so a forward failure stays retryable. */
export function isSafeBridgeMethod(method: string): boolean {
return SAFE_BRIDGE_METHODS.has(method.trim().toUpperCase());
}

View File

@ -1,384 +0,0 @@
import { afterEach, describe, expect, it } from "vitest";
import {
createDuplexBridgeBroker,
type DuplexBridgeBroker,
type DuplexBrokerForwardResult,
} from "./duplex-bridge-broker.js";
import {
DUPLEX_CHANNEL_AGGREGATE_BYTES_EXCEEDED,
DuplexAggregateByteLedger,
} from "./duplex-aggregate-byte-ledger.js";
import {
DUPLEX_FRAME_VERSION,
DuplexFrameDecoder,
encodeDuplexFrame,
type DuplexRequestFrame,
type DuplexResponseFrame,
} from "./duplex-frame-codec.js";
import { splitBodyIntoChunkFrames } from "./duplex-body-spool.js";
import type { CommandManagedDuplexChannel } from "./command-managed-runtime.js";
/**
* Regression harness for the broker aggregate byte ledger charging.
*
* The harness feeds request frames straight into the broker through an in-memory
* channel, the same as a provider that controls the transport. A test controls
* when each forward settles, so it can assert the ledger charge while the forward
* is still in flight and after it settles.
*/
/** One pending forward a test settles by hand. */
interface PendingForward {
id: string;
resolve: (result: DuplexBrokerForwardResult) => void;
reject: (error: Error) => void;
settled: boolean;
}
/**
* One reassembled response the broker wrote back. The broker writes a response as
* one envelope frame that carries `bodyByteCount`, then the `body_chunk` frames
* that carry the body. The harness reassembles the body and exposes it as `body`.
*/
type ReassembledResponse = DuplexResponseFrame & { body: string };
/** One request the harness feeds: the envelope frame plus its raw body text. */
interface RequestInput {
frame: DuplexRequestFrame;
bodyText: string;
}
/** The in-memory channel plus the levers a test uses to drive the broker. */
interface FakeChannelHarness {
channel: CommandManagedDuplexChannel;
feed: (input: RequestInput) => void;
exit: () => void;
responses: ReassembledResponse[];
forwards: PendingForward[];
resolveForward: (id: string, body?: string) => void;
}
/** Build one valid request: the envelope carries `bodyByteCount`, the body rides body_chunk frames. */
function requestFrame(id: string, method = "POST"): RequestInput {
const bodyText = JSON.stringify({ id });
return {
frame: {
version: DUPLEX_FRAME_VERSION,
type: "request",
id,
method,
path: `/api/issues/${id}`,
query: "",
headers: { "content-type": "application/json" },
bodyByteCount: Buffer.byteLength(bodyText, "utf8"),
},
bodyText,
};
}
function createFakeChannelHarness(): FakeChannelHarness {
const responses: ReassembledResponse[] = [];
const forwards: PendingForward[] = [];
const writtenDecoder = new DuplexFrameDecoder();
const responseAssembly = new Map<
string,
{ frame: DuplexResponseFrame; received: number; chunks: Buffer[] }
>();
let dataListener: ((chunk: Uint8Array) => void) | null = null;
let exitListener: ((exit: { exitCode: number | null }) => void) | null = null;
const channel: CommandManagedDuplexChannel = {
write: (data) => {
for (const result of writtenDecoder.push(data)) {
if (!result.ok) continue;
const frame = result.frame;
if (frame.type === "response") {
if (frame.bodyByteCount === 0) {
responses.push({ ...frame, body: "" });
} else {
responseAssembly.set(frame.id, { frame, received: 0, chunks: [] });
}
} else if (frame.type === "body_chunk") {
const assembly = responseAssembly.get(frame.id);
if (!assembly) continue;
const decoded = Buffer.from(frame.data, "base64");
assembly.received += decoded.length;
assembly.chunks.push(decoded);
if (assembly.received >= assembly.frame.bodyByteCount) {
responseAssembly.delete(frame.id);
responses.push({
...assembly.frame,
body: Buffer.concat(assembly.chunks).toString("utf8"),
});
}
}
}
},
onData: (listener) => {
dataListener = listener;
},
onExit: (listener) => {
exitListener = listener;
},
stop: () => undefined,
close: () => Promise.resolve(),
};
return {
channel,
feed: ({ frame, bodyText }) => {
if (!dataListener) throw new Error("The broker did not bind the data listener.");
dataListener(Buffer.from(encodeDuplexFrame(frame), "utf8"));
if (bodyText.length > 0) {
for (const chunk of splitBodyIntoChunkFrames(
frame.id,
Buffer.from(bodyText, "utf8"),
DUPLEX_FRAME_VERSION,
)) {
dataListener(Buffer.from(encodeDuplexFrame(chunk), "utf8"));
}
}
},
exit: () => {
if (!exitListener) throw new Error("The broker did not bind the exit listener.");
exitListener({ exitCode: 0 });
},
responses,
forwards,
resolveForward: (id, body = JSON.stringify({ ok: true })) => {
const forward = [...forwards].reverse().find((entry) => entry.id === id && !entry.settled);
if (!forward) throw new Error(`No unsettled forward for id ${id}.`);
forward.settled = true;
forward.resolve({ status: 200, headers: { "content-type": "application/json" }, body });
},
};
}
/** A forward handler that hands each call to the harness and never auto-resolves. */
function controllableForward(harness: FakeChannelHarness) {
return (request: DuplexRequestFrame): Promise<DuplexBrokerForwardResult> =>
new Promise<DuplexBrokerForwardResult>((resolve, reject) => {
harness.forwards.push({ id: request.id, resolve, reject, settled: false });
});
}
/**
* Wait for the microtasks and one macrotask to settle, so the broker reassembles
* a request body, dispatches its forward, and runs each settled promise handler.
*/
async function flush(): Promise<void> {
await new Promise((resolve) => setTimeout(resolve, 0));
}
describe("duplex bridge broker aggregate byte ledger", () => {
const brokers: DuplexBridgeBroker[] = [];
afterEach(async () => {
while (brokers.length > 0) {
const broker = brokers.pop();
if (broker) await broker.close();
}
});
it("charges request-frame, request-payload, and seen-id tokens, then releases the request tokens on forward settlement", async () => {
const harness = createFakeChannelHarness();
const ledger = new DuplexAggregateByteLedger({ ceilingBytes: 1_000_000 });
const broker = await createDuplexBridgeBroker({
channel: harness.channel,
forwardRequest: controllableForward(harness),
duplexAggregateByteLedger: ledger,
});
brokers.push(broker);
broker.start();
harness.feed(requestFrame("req-1"));
// The dispatch retains three tokens at the request envelope: the raw frame, the
// normalized payload, and the no-replay set entry. The reservation is
// synchronous, so the charge holds before the body reassembles.
expect(ledger.liveTokenCount).toBe(3);
expect(ledger.bytesInUse).toBeGreaterThan(0);
const chargedInFlight = ledger.bytesInUse;
// The broker reassembles the request body, then dispatches the forward.
await flush();
harness.resolveForward("req-1");
await flush();
// The forward settled. The finally owner released the request-frame and the
// request-payload tokens. The seen-id token stays charged for the channel
// lifetime, so exactly one token remains.
expect(ledger.liveTokenCount).toBe(1);
expect(ledger.bytesInUse).toBeGreaterThan(0);
expect(ledger.bytesInUse).toBeLessThan(chargedInFlight);
// A close releases the seen-id token, so the ledger returns to zero.
await broker.close();
expect(ledger.bytesInUse).toBe(0);
expect(ledger.liveTokenCount).toBe(0);
});
it("keeps request tokens charged when the response timer answers before the forward settles, and releases them on settlement", async () => {
const harness = createFakeChannelHarness();
const ledger = new DuplexAggregateByteLedger({ ceilingBytes: 1_000_000 });
const broker = await createDuplexBridgeBroker({
channel: harness.channel,
forwardRequest: controllableForward(harness),
// Squeeze the nested budgets so the response timer fires quickly while the
// forward stays unsettled.
budgets: { forwardTimeoutMs: 5, responseBudgetMs: 10, gatewayWaitMs: 20 },
duplexAggregateByteLedger: ledger,
});
brokers.push(broker);
broker.start();
harness.feed(requestFrame("req-1", "GET"));
expect(ledger.liveTokenCount).toBe(3);
// Wait past the response budget so the backstop answers the gateway while the
// forward is still in flight.
await new Promise((resolve) => setTimeout(resolve, 40));
expect(harness.responses.some((frame) => frame.id === "req-1")).toBe(true);
// The gateway got its answer, but the forward has not settled. The request
// tokens and the seen-id token stay charged: nothing released yet.
expect(ledger.liveTokenCount).toBe(3);
// Settle the orphaned forward. Its finally owner now releases the two request
// tokens; the seen-id token still stays for the channel lifetime.
harness.resolveForward("req-1");
await flush();
expect(ledger.liveTokenCount).toBe(1);
await broker.close();
expect(ledger.bytesInUse).toBe(0);
});
it("refuses a one-byte-over dispatch with the fixed marker and retains nothing", async () => {
const harness = createFakeChannelHarness();
// Measure the exact per-request charge with a throwaway broker, so the test
// can size a ceiling one byte below it and force the reservation to fail.
const probeHarness = createFakeChannelHarness();
const measured = new DuplexAggregateByteLedger({ ceilingBytes: 1_000_000 });
const measuringBroker = await createDuplexBridgeBroker({
channel: probeHarness.channel,
forwardRequest: controllableForward(probeHarness),
duplexAggregateByteLedger: measured,
});
measuringBroker.start();
probeHarness.feed(requestFrame("req-1"));
const perRequestBytes = measured.bytesInUse;
await measuringBroker.close();
const ledger = new DuplexAggregateByteLedger({ ceilingBytes: perRequestBytes - 1 });
const broker = await createDuplexBridgeBroker({
channel: harness.channel,
forwardRequest: controllableForward(harness),
duplexAggregateByteLedger: ledger,
});
brokers.push(broker);
broker.start();
harness.feed(requestFrame("req-1"));
await flush();
// The broker retained nothing: no token, no forward, no seen id.
expect(ledger.liveTokenCount).toBe(0);
expect(ledger.bytesInUse).toBe(0);
expect(harness.forwards.length).toBe(0);
// The refusal carries the fixed marker and is a bounded terminal response.
const refusal = harness.responses.find((frame) => frame.id === "req-1");
expect(refusal).toBeTruthy();
expect(refusal?.status).toBe(503);
expect(refusal?.body).toContain(DUPLEX_CHANNEL_AGGREGATE_BYTES_EXCEEDED);
// The refusal did not retain the id, so a resend after pressure eases is
// admitted: raise the ceiling and re-feed.
const roomyLedger = new DuplexAggregateByteLedger({ ceilingBytes: 1_000_000 });
const roomyHarness = createFakeChannelHarness();
const roomyBroker = await createDuplexBridgeBroker({
channel: roomyHarness.channel,
forwardRequest: controllableForward(roomyHarness),
duplexAggregateByteLedger: roomyLedger,
});
brokers.push(roomyBroker);
roomyBroker.start();
roomyHarness.feed(requestFrame("req-1"));
await flush();
expect(roomyHarness.forwards.length).toBe(1);
});
it("releases every retained token on a terminal channel loss", async () => {
const harness = createFakeChannelHarness();
const ledger = new DuplexAggregateByteLedger({ ceilingBytes: 1_000_000 });
const broker = await createDuplexBridgeBroker({
channel: harness.channel,
forwardRequest: controllableForward(harness),
duplexAggregateByteLedger: ledger,
});
brokers.push(broker);
broker.start();
harness.feed(requestFrame("req-1"));
harness.feed(requestFrame("req-2"));
// The reservation is synchronous, so both dispatches charge three tokens each
// at the request envelope, before the bodies reassemble.
expect(ledger.liveTokenCount).toBe(6);
// Let both requests reassemble and start their forwards, so the loss below
// transfers live forwards to the orphan registry.
await flush();
// The channel exits mid-flight. The broker latches the loss, transfers the
// in-flight forwards to the orphan registry, and releases the seen-id tokens.
harness.exit();
expect(broker.runDisposition.failed).toBe(true);
// The seen-id tokens released at loss; the two request tokens per forward stay
// charged until each aborted forward settles.
expect(ledger.liveTokenCount).toBe(4);
// The aborted forwards settle. Each finally owner releases its two request
// tokens, so the ledger returns to zero.
harness.forwards[0]?.reject(new Error("aborted"));
harness.forwards[1]?.reject(new Error("aborted"));
await flush();
expect(ledger.bytesInUse).toBe(0);
expect(ledger.liveTokenCount).toBe(0);
});
it("does not double-release a token when a forward settles after the channel closed", async () => {
const harness = createFakeChannelHarness();
let underflow = 0;
const ledger = new DuplexAggregateByteLedger({
ceilingBytes: 1_000_000,
telemetry: {
setBytesInUse() {},
recordReservationRejection() {},
recordAccountingUnderflow() {
underflow += 1;
},
},
});
const broker = await createDuplexBridgeBroker({
channel: harness.channel,
forwardRequest: controllableForward(harness),
duplexAggregateByteLedger: ledger,
});
brokers.push(broker);
broker.start();
harness.feed(requestFrame("req-1"));
// Let the request reassemble and start its forward before the close, so the
// close orphans a live forward that settles afterward.
await flush();
await broker.close();
// The close aborted the in-flight forward and orphaned it. The forward now
// settles after the close: its finally owner releases the request tokens once.
harness.forwards[0]?.reject(new Error("aborted"));
await flush();
expect(ledger.bytesInUse).toBe(0);
expect(ledger.liveTokenCount).toBe(0);
// No accounting defect: every token released exactly one time.
expect(underflow).toBe(0);
});
});

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@ -1,734 +0,0 @@
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
import net from "node:net";
import { mkdtemp, rm, writeFile } from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { Readable } from "node:stream";
import { afterEach, describe, expect, it } from "vitest";
import {
authorizeSandboxCallbackBridgeRequestWithRoutes,
getSandboxCallbackBridgeServerSource,
SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE,
} from "./sandbox-callback-bridge.js";
import {
DuplexFrameDecoder,
DUPLEX_BODY_CHUNK_RAW_BYTES,
DUPLEX_FRAME_VERSION,
encodeDuplexFrame,
type DuplexFrame,
type DuplexReadyFrame,
type DuplexRequestFrame,
} from "./duplex-frame-codec.js";
import {
createDuplexBridgeBroker,
type DuplexBridgeBroker,
type DuplexBrokerForwardResult,
} from "./duplex-bridge-broker.js";
import type { ReassembledBody } from "./duplex-body-spool.js";
import type { CommandManagedDuplexChannel } from "./command-managed-runtime.js";
/**
* Local real-process end-to-end harness for the composed duplex path.
*
* The harness spawns the real generated gateway with plain `node` from a
* temporary file, attaches the real host broker to the child stdin and stdout,
* and forwards each request to a local fake API server on a real HTTP call. It
* uses no provider credentials.
*
* The child stdio pipes are kernel pipes. They give partial reads, split a
* multi-byte UTF-8 sequence across two chunks, apply backpressure, and give a
* true end of file when the child dies. The harness proves the composed path
* handles each of these real conditions.
*/
/** One recorded call the fake API server received. */
interface FakeApiRequest {
method: string;
url: string;
auth: string | null;
runId: string | null;
body: string;
}
/** The responder the fake API server calls for each request. */
type FakeApiResponder = (req: IncomingMessage, res: ServerResponse, body: string) => void;
/** The handle for one fake API server. */
interface FakeApiServer {
origin: string;
requests: FakeApiRequest[];
setResponder: (responder: FakeApiResponder) => void;
/** The count of sockets the server holds open right now. */
openSocketCount: () => number;
/** True while the server still accepts connections. */
listening: () => boolean;
close: () => Promise<void>;
}
/**
* Start a local fake API server. The forward handler targets it, so one real
* HTTP call proves the composed path end to end. The server tracks each open
* socket, so a test asserts teardown leaves no open socket handle.
*/
async function startFakeApiServer(): Promise<FakeApiServer> {
const requests: FakeApiRequest[] = [];
const sockets = new Set<net.Socket>();
let responder: FakeApiResponder = (req, res) => {
const url = new URL(req.url ?? "/", "http://127.0.0.1");
res.writeHead(200, { "content-type": "application/json" });
res.end(JSON.stringify({ ok: true, path: url.pathname }));
};
const server = createServer((req, res) => {
const chunks: Buffer[] = [];
req.on("data", (chunk: Buffer) => chunks.push(chunk));
req.on("end", () => {
const body = Buffer.concat(chunks).toString("utf8");
requests.push({
method: req.method ?? "GET",
url: req.url ?? "/",
auth: req.headers.authorization ?? null,
runId:
typeof req.headers["x-paperclip-run-id"] === "string"
? req.headers["x-paperclip-run-id"]
: null,
body,
});
responder(req, res, body);
});
});
server.on("connection", (socket) => {
sockets.add(socket);
socket.on("close", () => sockets.delete(socket));
});
await new Promise<void>((resolve, reject) => {
server.once("error", reject);
server.listen(0, "127.0.0.1", () => resolve());
});
const address = server.address();
if (!address || typeof address === "string") {
throw new Error("The fake API server did not expose a TCP port.");
}
return {
origin: `http://127.0.0.1:${address.port}`,
requests,
setResponder: (next) => {
responder = next;
},
openSocketCount: () => sockets.size,
listening: () => server.listening,
close: () =>
new Promise<void>((resolve) => {
// Destroy each open socket first. A keep-alive client socket keeps the
// server open, so `server.close` alone could stall. The destroy makes the
// close callback fire and drives the open socket count to zero.
for (const socket of sockets) socket.destroy();
server.close(() => resolve());
}),
};
}
/** The channel view over the spawned child, plus the frames the child sent host-ward. */
interface ChildDuplexChannel {
channel: CommandManagedDuplexChannel;
observedFrames: DuplexFrame[];
stderr: () => string;
}
/**
* Wrap the spawned child as a {@link CommandManagedDuplexChannel}. The broker
* writes to the child stdin, reads the child stdout, and learns of the child
* exit through this channel.
*
* The host end forwards each raw stdout chunk unchanged; the channel carries
* bytes, so no decode step sits between the pipe and the broker.
*/
function attachChildDuplexChannel(child: ChildProcessWithoutNullStreams): ChildDuplexChannel {
const observed = new DuplexFrameDecoder();
const observedFrames: DuplexFrame[] = [];
let dataListener: ((chunk: Uint8Array) => void) | null = null;
let exitListener: ((exit: { exitCode: number | null }) => void) | null = null;
let pendingBytes: Buffer = Buffer.alloc(0);
let pendingExit: { exitCode: number | null } | null = null;
let stderrText = "";
child.stdout.on("data", (buffer: Buffer) => {
for (const result of observed.push(buffer)) {
if (result.ok) observedFrames.push(result.frame);
}
if (buffer.length === 0) return;
if (dataListener) dataListener(buffer);
else pendingBytes = pendingBytes.length === 0 ? buffer : Buffer.concat([pendingBytes, buffer]);
});
child.stderr.on("data", (buffer: Buffer) => {
stderrText += buffer.toString("utf8");
});
child.on("exit", (code) => {
const exit = { exitCode: code };
if (exitListener) exitListener(exit);
else pendingExit = exit;
});
// Swallow a stdin EPIPE. The broker can write one more frame while the child
// exits; the write fails and the broker records the loss on its own path.
child.stdin.on("error", () => undefined);
const channel: CommandManagedDuplexChannel = {
write: (data) => {
child.stdin.write(data);
},
onData: (listener) => {
dataListener = listener;
if (pendingBytes.length > 0) {
const replay = pendingBytes;
pendingBytes = Buffer.alloc(0);
listener(replay);
}
},
onExit: (listener) => {
exitListener = listener;
if (pendingExit) {
const exit = pendingExit;
pendingExit = null;
listener(exit);
}
},
stop: () => {
child.kill("SIGKILL");
},
close: () =>
new Promise<void>((resolve) => {
// Close the write side. The child stdin reaches end of file, so the
// gateway sees a real EOF.
child.stdin.end(() => resolve());
}),
};
return { channel, observedFrames, stderr: () => stderrText };
}
/** The forward-handler mode. `proxy` calls the fake API; `hang` blocks until abort. */
type ForwardMode = "proxy" | "hang";
/** The options for one harness. */
interface HarnessOptions {
hostApiToken?: string;
runId?: string;
maxBodyBytes?: number;
lossExitGraceMs?: number;
}
/** The full harness handle for one composed duplex path. */
interface DuplexE2EHarness {
baseUrl: string;
bridgeToken: string;
hostApiToken: string;
runId: string;
nonce: string;
broker: DuplexBridgeBroker;
api: FakeApiServer;
child: ChildProcessWithoutNullStreams;
observedFrames: DuplexFrame[];
forwardedRequests: DuplexRequestFrame[];
setForwardMode: (mode: ForwardMode) => void;
stderr: () => string;
killChild: () => Promise<void>;
waitFor: (predicate: () => boolean, message: string, timeoutMs?: number) => Promise<void>;
teardown: () => Promise<void>;
}
/** Reserve one free loopback port. The host assigns it to the gateway. */
async function reserveLoopbackPort(): Promise<number> {
return new Promise<number>((resolve, reject) => {
const probe = net.createServer();
probe.once("error", reject);
probe.listen(0, "127.0.0.1", () => {
const address = probe.address();
if (!address || typeof address === "string") {
probe.close(() => reject(new Error("Could not reserve a loopback port.")));
return;
}
const reserved = address.port;
probe.close(() => resolve(reserved));
});
});
}
/**
* Build the composed duplex path: a spawned gateway child, the real broker on
* the child stdio, and a local fake API server as the forward target.
*/
async function createHarness(options: HarnessOptions = {}): Promise<DuplexE2EHarness> {
const hostApiToken = options.hostApiToken ?? "real-run-jwt";
const runId = options.runId ?? "run-e2e";
const maxBodyBytes = options.maxBodyBytes ?? 1_000_000;
const bridgeToken = "duplex-e2e-bridge-token";
const nonce = "e2e112233445566778899aabbccddeeff";
let forwardMode: ForwardMode = "proxy";
const forwardedRequests: DuplexRequestFrame[] = [];
const api = await startFakeApiServer();
const assignedPort = await reserveLoopbackPort();
const tmpDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-duplex-e2e-"));
const entrypoint = path.join(tmpDir, "gateway.mjs");
await writeFile(entrypoint, getSandboxCallbackBridgeServerSource(), "utf8");
const child = spawn(process.execPath, [entrypoint], {
stdio: ["pipe", "pipe", "pipe"],
env: {
...process.env,
PAPERCLIP_API_BRIDGE_MODE: SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE,
PAPERCLIP_BRIDGE_HOST: "127.0.0.1",
PAPERCLIP_BRIDGE_PORT: String(assignedPort),
PAPERCLIP_BRIDGE_NONCE: nonce,
PAPERCLIP_BRIDGE_TOKEN: bridgeToken,
PAPERCLIP_BRIDGE_MAX_BODY_BYTES: String(maxBodyBytes),
...(options.lossExitGraceMs != null
? { PAPERCLIP_BRIDGE_LOSS_EXIT_GRACE_MS: String(options.lossExitGraceMs) }
: {}),
},
}) as ChildProcessWithoutNullStreams;
const { channel, observedFrames, stderr } = attachChildDuplexChannel(child);
const forwardRequest = async (
request: DuplexRequestFrame,
opts: { signal: AbortSignal; body: ReassembledBody },
): Promise<DuplexBrokerForwardResult> => {
forwardedRequests.push(request);
// Keep the real route allowlist on the forward seam. The broker forwards
// only an allowed route; it denies any other route with a 403.
const denial = authorizeSandboxCallbackBridgeRequestWithRoutes(request);
if (denial) {
return {
status: 403,
headers: { "content-type": "application/json" },
body: JSON.stringify({ error: denial }),
};
}
if (forwardMode === "hang") {
// Block until the broker aborts the forward. This lets a test hold an
// outstanding request open while it forces a loss.
await new Promise<never>((_resolve, reject) => {
const onAbort = () => reject(new Error("The forward call was aborted."));
if (opts.signal.aborted) {
onAbort();
return;
}
opts.signal.addEventListener("abort", onAbort, { once: true });
});
}
const method = request.method.trim().toUpperCase() || "GET";
const headers = new Headers();
for (const [key, value] of Object.entries(request.headers)) {
if (value.trim().length === 0) continue;
headers.set(key, value);
}
// Apply the real host token and the run id, the same as the file bridge path.
headers.set("authorization", `Bearer ${hostApiToken}`);
headers.set("x-paperclip-run-id", runId);
const target = new URL(`${request.path}${request.query ?? ""}`, api.origin);
// Stream the reassembled request body to the host, the same as the real
// forward. A streamed body needs `duplex: "half"`.
const forwardInit: RequestInit & { duplex?: "half" } = { method, headers, signal: opts.signal };
if (method !== "GET" && method !== "HEAD") {
forwardInit.body = Readable.toWeb(
opts.body.createReadStream(),
) as unknown as ReadableStream<Uint8Array>;
forwardInit.duplex = "half";
}
const response = await fetch(target, forwardInit);
const body = await response.text();
const outHeaders: Record<string, string> = {};
response.headers.forEach((value, key) => {
if (key.toLowerCase() === "content-length") return;
outHeaders[key] = value;
});
return { status: response.status, headers: outHeaders, body };
};
const broker = await createDuplexBridgeBroker({
channel,
forwardRequest,
logger: () => undefined,
});
broker.start();
const waitFor = async (
predicate: () => boolean,
message: string,
timeoutMs = 5000,
): Promise<void> => {
const deadline = Date.now() + timeoutMs;
for (;;) {
if (predicate()) return;
if (Date.now() > deadline) {
throw new Error(`Timed out waiting for ${message}. stderr: ${stderr()}`);
}
await new Promise<void>((resolve) => {
const timer = setTimeout(resolve, 20);
timer.unref?.();
});
}
};
const waitForExit = (): Promise<void> =>
new Promise<void>((resolve) => {
if (child.exitCode !== null || child.signalCode !== null) {
resolve();
return;
}
child.once("exit", () => resolve());
});
const killChild = async (): Promise<void> => {
child.kill("SIGKILL");
await waitForExit();
};
let torndown = false;
const teardown = async (): Promise<void> => {
if (torndown) return;
torndown = true;
child.kill("SIGKILL");
await waitForExit();
child.stdin.destroy();
child.stdout.destroy();
child.stderr.destroy();
await api.close();
await rm(tmpDir, { recursive: true, force: true });
};
return {
baseUrl: `http://127.0.0.1:${assignedPort}`,
bridgeToken,
hostApiToken,
runId,
nonce,
broker,
api,
child,
observedFrames,
forwardedRequests,
setForwardMode: (mode) => {
forwardMode = mode;
},
stderr,
killChild,
waitFor,
teardown,
};
}
/**
* Build a large body of multi-byte UTF-8 characters. "€" uses three UTF-8 bytes
* and "😀" uses four. A pipe read boundary at a 65536-byte multiple lands inside
* a "€" sequence, because 65536 is not a multiple of three. This guarantees a
* multi-byte character split across two pipe chunks.
*/
function buildMultiByteBody(targetBytes: number): string {
const marker = "😀-start-😀";
const filler = "€".repeat(Math.ceil(targetBytes / 3));
return `${marker}${filler}${marker}`;
}
describe("duplex bridge local end-to-end harness", () => {
const harnesses: DuplexE2EHarness[] = [];
afterEach(async () => {
while (harnesses.length > 0) {
const harness = harnesses.pop();
if (harness) await harness.teardown();
}
});
const readyPredicate = (harness: DuplexE2EHarness) => () =>
harness.observedFrames.some((frame) => frame.type === "ready");
it("delivers a valid READY frame from the spawned gateway child to the attached broker", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "a READY frame from the gateway child");
const ready = harness.observedFrames.find(
(frame) => frame.type === "ready",
) as DuplexReadyFrame;
expect(ready.version).toBe(2);
expect(ready.nonce).toBe(harness.nonce);
// READY is liveness only. It carries no address data.
expect((ready as unknown as Record<string, unknown>).address).toBeUndefined();
expect((ready as unknown as Record<string, unknown>).port).toBeUndefined();
// The broker read the READY frame and stayed open with no loss.
expect(harness.broker.state).toBe("open");
expect(harness.broker.lossRecord).toBeNull();
}, 20000);
it("carries one real HTTP request through the child and the broker to the fake API unchanged", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
harness.api.setResponder((req, res) => {
const url = new URL(req.url ?? "/", "http://127.0.0.1");
res.writeHead(200, { "content-type": "application/json" });
res.end(JSON.stringify({ ok: true, echoedPath: url.pathname }));
});
const response = await fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
expect(response.status).toBe(200);
await expect(response.json()).resolves.toEqual({ ok: true, echoedPath: "/api/agents/me" });
// The fake API saw one request with the real host token and the run id. The
// broker replaced the bridge token, so the fake API never saw it.
expect(harness.api.requests).toHaveLength(1);
expect(harness.api.requests[0]).toMatchObject({
method: "GET",
url: "/api/agents/me",
auth: `Bearer ${harness.hostApiToken}`,
runId: harness.runId,
});
}, 20000);
it("reassembles a large response that spans many pipe chunks", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
// The body is larger than the pipe buffer, so the response frame crosses the
// child stdin pipe in many chunks. An ASCII body isolates the chunk-span
// behavior from the multi-byte behavior the next test covers.
const largeBody = "x".repeat(256 * 1024);
expect(Buffer.byteLength(largeBody, "utf8")).toBeGreaterThan(200 * 1024);
harness.api.setResponder((_req, res) => {
res.writeHead(200, { "content-type": "text/plain; charset=utf-8" });
res.end(largeBody);
});
const response = await fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
expect(response.status).toBe(200);
const received = await response.text();
// The whole body returns complete after it crossed the pipe in many chunks.
expect(received.length).toBe(largeBody.length);
expect(received).toBe(largeBody);
}, 20000);
it("decodes a multi-byte UTF-8 sequence split across chunk borders", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
const multiByteBody = buildMultiByteBody(256 * 1024);
expect(Buffer.byteLength(multiByteBody, "utf8")).toBeGreaterThan(200 * 1024);
harness.api.setResponder((_req, res) => {
res.writeHead(200, { "content-type": "text/plain; charset=utf-8" });
res.end(multiByteBody);
});
const response = await fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
expect(response.status).toBe(200);
const received = await response.text();
// A multi-byte character split across a pipe chunk border decodes to the
// same character, so the whole body returns unchanged.
expect(received.length).toBe(multiByteBody.length);
expect(received).toBe(multiByteBody);
}, 20000);
it("matches each concurrent request through one child to its own response", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
harness.api.setResponder((req, res) => {
const url = new URL(req.url ?? "/", "http://127.0.0.1");
res.writeHead(200, { "content-type": "application/json" });
res.end(JSON.stringify({ path: url.pathname }));
});
const ids = Array.from({ length: 8 }, (_unused, index) => `issue-${index}`);
const responses = await Promise.all(
ids.map((id) =>
fetch(`${harness.baseUrl}/api/issues/${id}`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
}).then(async (response) => ({ status: response.status, body: await response.json() })),
),
);
responses.forEach((result, index) => {
expect(result.status).toBe(200);
expect(result.body).toEqual({ path: `/api/issues/${ids[index]}` });
});
expect(harness.api.requests).toHaveLength(8);
}, 20000);
it("answers 409 to outstanding and 503 to new requests when the broker closes its write side", async () => {
const harness = await createHarness({ lossExitGraceMs: 3000 });
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
// Hold the forward open, so the broker sends no response frame. The gateway
// keeps the request outstanding until the loss.
harness.setForwardMode("hang");
const outstanding = fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
void outstanding.catch(() => undefined);
await harness.waitFor(
() => harness.forwardedRequests.length >= 1,
"the broker to receive the forwarded request",
);
// The broker closes its write side, so the child stdin reaches a real EOF.
await harness.broker.close();
const lossResponse = await outstanding;
expect(lossResponse.status).toBe(409);
expect(lossResponse.headers.get("x-paperclip-bridge-outcome")).toBe("indeterminate");
await expect(lossResponse.json()).resolves.toEqual({ error: "outcome_indeterminate" });
const afterLoss = await fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
expect(afterLoss.status).toBe(503);
await expect(afterLoss.json()).resolves.toEqual({ error: "bridge_unavailable" });
}, 20000);
it("moves the broker to lost on a child kill, dispatches nothing more, and leaks no handle after teardown", async () => {
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
// Hold one forward open, so a request is in flight when the child dies.
harness.setForwardMode("hang");
const outstanding = fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
void outstanding.catch(() => undefined);
await harness.waitFor(
() => harness.forwardedRequests.length >= 1,
"the broker to receive the forwarded request",
);
const forwardedBeforeKill = harness.forwardedRequests.length;
// A real process kill closes the stdout pipe and ends the child.
await harness.killChild();
await harness.waitFor(
() => harness.broker.state === "lost",
"the broker to record the channel loss",
);
expect(harness.broker.lossRecord?.reason).toBe("channel_exit");
// The broker dispatches nothing more after the loss.
await new Promise<void>((resolve) => {
const timer = setTimeout(resolve, 100);
timer.unref?.();
});
expect(harness.forwardedRequests.length).toBe(forwardedBeforeKill);
// The outstanding fetch fails because the child died before it answered.
await expect(outstanding).rejects.toThrow();
await harness.teardown();
// Teardown left no open pipe or socket handle.
expect(harness.child.stdin.destroyed).toBe(true);
expect(harness.child.stdout.destroyed).toBe(true);
expect(harness.child.stderr.destroyed).toBe(true);
expect(harness.api.listening()).toBe(false);
expect(harness.api.openSocketCount()).toBe(0);
}, 20000);
it("fails a request with a local 502 when the host sends a malformed response chunk", async () => {
const RAW = DUPLEX_BODY_CHUNK_RAW_BYTES;
// Each case declares one malformed body_chunk a broken or hostile host could
// send. The gateway must reject it at once with a local 502. It must not grow
// its response reassembly buffer until the response timeout. `bodyByteCount`
// sets the declared body size on the response envelope. `data` is the base64
// payload of the one injected chunk.
const cases: Array<{ name: string; bodyByteCount: number; data: string; error: string }> = [
{ name: "an empty chunk", bodyByteCount: RAW, data: "", error: "duplex response body_chunk is empty" },
{
name: "an undersized non-final chunk",
bodyByteCount: RAW + 8,
data: Buffer.alloc(4, 1).toString("base64"),
error: "duplex response body_chunk has the wrong size",
},
{
name: "an oversized chunk",
bodyByteCount: RAW + 8,
data: Buffer.alloc(RAW + 4, 1).toString("base64"),
error: "duplex response body_chunk has the wrong size",
},
{
name: "an overrun past the declared size",
bodyByteCount: 10,
data: Buffer.alloc(200, 1).toString("base64"),
error: "duplex response body overruns the declared size",
},
{
name: "a non-canonical base64 chunk",
bodyByteCount: 10,
data: "AB==",
error: "duplex response body_chunk is not canonical base64",
},
];
const harness = await createHarness();
harnesses.push(harness);
await harness.waitFor(readyPredicate(harness), "the gateway to become ready");
// Hold the broker forward open, so only the injected frames answer each
// request. The gateway reassembles the response body itself, so a raw
// malformed chunk on its stdin drives the reject path under test.
harness.setForwardMode("hang");
for (const testCase of cases) {
const forwardedBefore = harness.forwardedRequests.length;
const pendingFetch = fetch(`${harness.baseUrl}/api/agents/me`, {
headers: { authorization: `Bearer ${harness.bridgeToken}` },
});
void pendingFetch.catch(() => undefined);
await harness.waitFor(
() => harness.forwardedRequests.length > forwardedBefore,
`the broker to receive the forwarded request for ${testCase.name}`,
);
const id = harness.forwardedRequests[harness.forwardedRequests.length - 1].id;
// Inject the response envelope, then the one malformed body_chunk, straight
// to the gateway stdin. This bypasses the broker response encoder, so the
// gateway sees the exact bytes a broken or hostile host could send.
harness.child.stdin.write(
encodeDuplexFrame({
version: DUPLEX_FRAME_VERSION,
type: "response",
id,
status: 200,
headers: {},
bodyByteCount: testCase.bodyByteCount,
outcome: "completed",
}),
);
harness.child.stdin.write(
encodeDuplexFrame({
version: DUPLEX_FRAME_VERSION,
type: "body_chunk",
id,
seq: 0,
data: testCase.data,
}),
);
const response = await pendingFetch;
expect(response.status).toBe(502);
await expect(response.json()).resolves.toEqual({ error: testCase.error });
}
}, 20000);
});

View File

@ -19,9 +19,9 @@
* HTTP/2 is the preferred transport. `queue_v1` is the soft-deprecated fallback.
* The `http2_v1` host readiness gate imports {@link decodeDuplexLine} from this
* file to read the one READY line every gateway sends, so this file's READY
* frame path stays live for both transports. `duplex-bridge-broker.ts` and
* `duplex-body-spool.ts` still import the request, response, and body-chunk
* frame types from this file, so this phase keeps every frame type here.
* frame path stays live for both transports. `duplex-body-spool.ts` still
* imports the request, response, and body-chunk frame types from this file, so
* this file keeps every frame type here.
*/
import {

File diff suppressed because it is too large Load Diff

View File

@ -51,7 +51,7 @@ import {
DUPLEX_CHANNEL_LOST_ERROR_CODE,
isSafeBridgeMethod,
type DuplexBrokerRunDisposition,
} from "./duplex-bridge-broker.js";
} from "./bridge-transport-contract.js";
import { decodeDuplexLine, DEFAULT_MAX_DUPLEX_FRAME_BYTES } from "./duplex-frame-codec.js";
import type { ReassembledBody } from "./duplex-body-spool.js";
import {

View File

@ -2941,6 +2941,71 @@ describe("sandbox callback bridge", () => {
expect(stderr).toContain("EADDRINUSE");
}, 15_000);
it("exits nonzero for the retired duplex_v1 mode instead of starting the queue gateway", async () => {
// The closed mode allowlist rejects `duplex_v1` before the queue-directory
// check, so a stale `duplex_v1` launch environment fails startup instead
// of silently falling through to the queue gateway.
const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-bridge-mode-duplex-"));
cleanupDirs.push(rootDir);
const entrypoint = path.join(rootDir, "paperclip-bridge-server.mjs");
await writeFile(entrypoint, getSandboxCallbackBridgeServerSource(), "utf8");
const queueDir = path.join(rootDir, "queue");
await mkdir(queueDir, { recursive: true });
const child = spawn(process.execPath, [entrypoint], {
env: {
...process.env,
PAPERCLIP_API_BRIDGE_MODE: "duplex_v1",
PAPERCLIP_BRIDGE_QUEUE_DIR: queueDir,
PAPERCLIP_BRIDGE_TOKEN: "test-token",
PAPERCLIP_BRIDGE_PORT: "0",
},
stdio: ["ignore", "ignore", "pipe"],
});
let stderr = "";
child.stderr.on("data", (chunk) => {
stderr += chunk;
});
const exitCode = await new Promise<number | null>((resolve) => {
child.on("close", resolve);
});
expect(exitCode).not.toBe(0);
expect(stderr).toContain("Unsupported PAPERCLIP_API_BRIDGE_MODE: duplex_v1");
}, 15_000);
it("exits nonzero for an unknown bridge mode instead of starting the queue gateway", async () => {
// The closed mode allowlist rejects every value it does not name, not
// only the retired duplex transport.
const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-bridge-mode-unknown-"));
cleanupDirs.push(rootDir);
const entrypoint = path.join(rootDir, "paperclip-bridge-server.mjs");
await writeFile(entrypoint, getSandboxCallbackBridgeServerSource(), "utf8");
const queueDir = path.join(rootDir, "queue");
await mkdir(queueDir, { recursive: true });
const child = spawn(process.execPath, [entrypoint], {
env: {
...process.env,
PAPERCLIP_API_BRIDGE_MODE: "totally_unknown_mode",
PAPERCLIP_BRIDGE_QUEUE_DIR: queueDir,
PAPERCLIP_BRIDGE_TOKEN: "test-token",
PAPERCLIP_BRIDGE_PORT: "0",
},
stdio: ["ignore", "ignore", "pipe"],
});
let stderr = "";
child.stderr.on("data", (chunk) => {
stderr += chunk;
});
const exitCode = await new Promise<number | null>((resolve) => {
child.on("close", resolve);
});
expect(exitCode).not.toBe(0);
expect(stderr).toContain("Unsupported PAPERCLIP_API_BRIDGE_MODE: totally_unknown_mode");
}, 15_000);
it("test_http2_gateway_writes_no_frame_between_ready_and_the_preface", async () => {
// Spawn the real generated gateway in http2_v1 mode and read its raw
// stdout bytes. The only frame-codec write on this path is the READY

View File

@ -73,35 +73,16 @@ const SANDBOX_EXEC_CHANNEL_ENV = "PAPERCLIP_SANDBOX_EXEC_CHANNEL";
const SANDBOX_EXEC_CHANNEL_BRIDGE = "bridge";
// The bridge modes the generated gateway supports. The file mode polls a
// request/response queue on disk. The retired duplex mode forwarded one
// request frame to stdout and resolved one response frame from stdin; the
// generated gateway still defines it, but no mode dispatch selects it anymore
// — the http2 mode replaced it as the active non-file transport. The http2
// mode runs one Node HTTP/2 client session directly on stdin/stdout, after it
// sends the one READY line the host readiness gate expects. The generated
// `.mjs` selects the mode from `PAPERCLIP_API_BRIDGE_MODE`.
// request/response queue on disk. The http2 mode runs one Node HTTP/2 client
// session directly on stdin/stdout, after it sends the one READY line the
// host readiness gate expects. The generated `.mjs` selects the mode from
// `PAPERCLIP_API_BRIDGE_MODE`. The generated gateway rejects every other
// value with a fixed startup error, including the retired duplex transport.
// HTTP/2 is the preferred transport. `queue_v1` is the soft-deprecated fallback.
const SANDBOX_CALLBACK_BRIDGE_FILE_MODE = "queue_v1";
export const SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE = "duplex_v1";
/** The active non-file transport mode. It replaced {@link SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE}
* in the mode-selection path. */
/** The active non-file transport mode. */
export const SANDBOX_CALLBACK_BRIDGE_HTTP2_MODE = "http2_v1";
// The duplex gateway HTTP wait budget default. The gateway waits this long for a
// response frame before it answers the local caller with a 502 timeout. The
// value stays configurable through `PAPERCLIP_BRIDGE_RESPONSE_TIMEOUT_MS`, the
// same environment key the file mode reads.
const DEFAULT_DUPLEX_GATEWAY_WAIT_BUDGET_MS = 35_000;
// How often the duplex gateway writes a heartbeat frame to stdout.
const DEFAULT_DUPLEX_GATEWAY_HEARTBEAT_INTERVAL_MS = 5_000;
// The duplex gateway treats the channel as lost when no inbound frame arrives on
// stdin within this window.
const DEFAULT_DUPLEX_GATEWAY_HEARTBEAT_TIMEOUT_MS = 20_000;
// After a loss, the duplex gateway answers each new local request with a 503,
// then exits when this grace ends. The grace gives the local caller time to read
// the 409 for an outstanding request and a 503 for a new request.
const DEFAULT_DUPLEX_GATEWAY_LOSS_EXIT_GRACE_MS = 1_000;
/** Span name that wraps one Paperclip-API callback request — read the request,
* write the response, and remove the request file. */
const CALLBACK_BRIDGE_RELAY_REQUEST_SPAN = "sandbox.callbackBridge.relayRequest";
@ -1847,9 +1828,8 @@ export async function startSandboxCallbackBridgeServer(input: {
//
// This gateway turns each local loopback request into one HTTP/2 stream to
// the host server in `http2-bridge-server.ts`. It sits beside the file-mode
// and duplex-mode gateways above; it changes neither of them, and no mode
// dispatch selects it yet — a later phase wires it into the generated
// in-sandbox entrypoint and the transport-selection path.
// gateway above; it changes neither of them. The generated in-sandbox
// entrypoint and the transport-selection path already select it.
// ---------------------------------------------------------------------------
/**
@ -1998,7 +1978,7 @@ function forwardOneHttp2Request(
* local request the caller hands it (already checked against the bridge
* token — see {@link SandboxHttp2BridgeGatewayRequest.receivedToken}) as one
* HTTP/2 stream. It keeps the header allowlist on the sandbox side, exactly
* as the file-mode and duplex-mode gateways do.
* as the file-mode gateway does.
*/
export function createSandboxHttp2BridgeGateway(
options: CreateSandboxHttp2BridgeGatewayOptions,
@ -2306,36 +2286,18 @@ const bridgeToken = process.env.PAPERCLIP_BRIDGE_TOKEN;
const host = process.env.PAPERCLIP_BRIDGE_HOST || "127.0.0.1";
const port = Number(process.env.PAPERCLIP_BRIDGE_PORT || "0");
// The host assigns the loopback port and passes it through the launch
// environment. The duplex gateway binds exactly this port; it never selects a
// environment. The gateway binds exactly this port; it never selects a
// different one. The host also passes one random per-open nonce here. The
// gateway echoes it in the READY frame so the host correlates READY with this
// channel open. The nonce is a liveness signal, not authentication.
const bridgeNonce = process.env.PAPERCLIP_BRIDGE_NONCE || "";
const pollIntervalMs = Number(process.env.PAPERCLIP_BRIDGE_POLL_INTERVAL_MS || "100");
const responseTimeoutMs = Number(
process.env.PAPERCLIP_BRIDGE_RESPONSE_TIMEOUT_MS ||
(bridgeMode === "${SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE}"
? "${DEFAULT_DUPLEX_GATEWAY_WAIT_BUDGET_MS}"
: "${DEFAULT_BRIDGE_RESPONSE_TIMEOUT_MS}"),
process.env.PAPERCLIP_BRIDGE_RESPONSE_TIMEOUT_MS || "${DEFAULT_BRIDGE_RESPONSE_TIMEOUT_MS}",
);
const maxQueueDepth = Number(process.env.PAPERCLIP_BRIDGE_MAX_QUEUE_DEPTH || "${DEFAULT_BRIDGE_MAX_QUEUE_DEPTH}");
const maxBodyBytes = Number(process.env.PAPERCLIP_BRIDGE_MAX_BODY_BYTES || "${DEFAULT_BRIDGE_MAX_BODY_BYTES}");
// The host passes the separate sandbox-process raw-decoder cap here. The in-sandbox
// decoder enforces it locally under the "sandbox_process" scope; it never shares the
// host aggregate byte ledger.
const maxDuplexDecoderBytes = Number(
process.env.PAPERCLIP_BRIDGE_MAX_DUPLEX_DECODER_BYTES || "${DEFAULT_BRIDGE_MAX_DUPLEX_DECODER_BYTES}",
);
const heartbeatIntervalMs = Number(
process.env.PAPERCLIP_BRIDGE_HEARTBEAT_INTERVAL_MS || "${DEFAULT_DUPLEX_GATEWAY_HEARTBEAT_INTERVAL_MS}",
);
const heartbeatTimeoutMs = Number(
process.env.PAPERCLIP_BRIDGE_HEARTBEAT_TIMEOUT_MS || "${DEFAULT_DUPLEX_GATEWAY_HEARTBEAT_TIMEOUT_MS}",
);
const lossExitGraceMs = Number(
process.env.PAPERCLIP_BRIDGE_LOSS_EXIT_GRACE_MS || "${DEFAULT_DUPLEX_GATEWAY_LOSS_EXIT_GRACE_MS}",
);
// The header allowlist. Both the file gateway and the duplex gateway strip an
// The header allowlist. Both the file gateway and the http2 gateway strip an
// inbound request to these headers before they forward it. One copy serves both
// modes. The route allowlist stays on the host: both modes forward a request to
// the host, and the host enforces the same route allowlist for each.
@ -2344,11 +2306,19 @@ const allowedHeaders = new Set(${JSON.stringify([...DEFAULT_SANDBOX_CALLBACK_BRI
if (!bridgeToken) {
throw new Error("PAPERCLIP_BRIDGE_TOKEN is required.");
}
// Closed allowlist for the bridge mode. The generated gateway supports exactly
// two transports: http2 and the file-mode queue. Every other value, including
// the retired duplex transport, fails startup at once instead of falling
// through to a mode that never ran. This check runs before the queue-directory
// check below, so an unsupported mode never reaches a state where a missing
// queue directory masks the real problem.
if (
bridgeMode !== "${SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE}" &&
bridgeMode !== "${SANDBOX_CALLBACK_BRIDGE_HTTP2_MODE}" &&
!queueDir
bridgeMode !== "${SANDBOX_CALLBACK_BRIDGE_FILE_MODE}"
) {
throw new Error("Unsupported PAPERCLIP_API_BRIDGE_MODE: " + bridgeMode);
}
if (bridgeMode !== "${SANDBOX_CALLBACK_BRIDGE_HTTP2_MODE}" && !queueDir) {
throw new Error("PAPERCLIP_BRIDGE_QUEUE_DIR and PAPERCLIP_BRIDGE_TOKEN are required.");
}
@ -2598,350 +2568,6 @@ async function runFileGateway() {
});
}
// Split one whole body buffer into body_chunk frames. The gateway buffers the
// whole source body once, then splits it here. Each frame carries one fixed raw
// slice as base64 text, except the final frame, which carries the remaining
// bytes. A zero-length body yields no frame. Send-side true streaming is a
// separate goal; this whole-body split is acceptable for this gateway.
function splitDuplexBodyIntoChunks(id, body) {
const frames = [];
let seq = 0;
for (let offset = 0; offset < body.length; offset += DUPLEX_BODY_CHUNK_RAW_BYTES) {
const slice = body.subarray(offset, offset + DUPLEX_BODY_CHUNK_RAW_BYTES);
frames.push({
version: DUPLEX_FRAME_VERSION,
type: "body_chunk",
id: id,
seq: seq,
data: slice.toString("base64"),
});
seq += 1;
}
return frames;
}
function runDuplexGateway() {
// One outstanding local request per id. Each entry holds the HTTP resolver and
// the wait-budget timer.
const pending = new Map();
// One in-flight response reassembly per id. The host returns a response as an
// envelope frame that carries bodyByteCount, then the body_chunk frames that
// carry the body. The gateway reassembles the response body in memory here; it
// does not spill, because a production response body stays small.
const responseAssembly = new Map();
let unavailable = false;
let lossTriggered = false;
let lastInboundAt = Date.now();
function diag(message) {
// Diagnostics go to stderr only. Stdout carries frames.
process.stderr.write("[paperclip-bridge] " + message + "\\n");
}
function writeFrame(frame) {
process.stdout.write(encodeDuplexFrame(frame));
}
function stopTimers() {
clearInterval(heartbeatSendTimer);
clearInterval(heartbeatWatchTimer);
}
// Loss behavior. On stdin end, a close frame, or a heartbeat timeout, answer
// every outstanding request with a non-retryable 409 (the request may have
// reached the host), answer every new request with a retryable 503 (the
// request never left the gateway), then exit. The gateway never replays a
// request.
function triggerLoss(reason) {
if (lossTriggered) return;
lossTriggered = true;
unavailable = true;
diag("duplex channel lost: " + reason);
stopTimers();
for (const entry of pending.values()) {
clearTimeout(entry.timer);
entry.resolve({
status: 409,
headers: {
"content-type": "application/json",
"x-paperclip-bridge-outcome": "indeterminate",
},
body: JSON.stringify({ error: "outcome_indeterminate" }),
});
}
pending.clear();
responseAssembly.clear();
const exitTimer = setTimeout(() => process.exit(0), lossExitGraceMs);
if (typeof exitTimer.unref === "function") exitTimer.unref();
}
// Fail one outstanding request with a bounded local 502. The gateway calls it
// when a response reassembly breaks: a reordered seq, a base64 that is not
// canonical, or a total that overruns the declared body size. The host is the
// response peer, so this is a defensive local error, not a channel loss.
function failRequest(id, message) {
responseAssembly.delete(id);
const entry = pending.get(id);
if (!entry) return;
pending.delete(id);
clearTimeout(entry.timer);
entry.resolve({
status: 502,
headers: { "content-type": "application/json" },
body: JSON.stringify({ error: message }),
});
}
function handleInboundFrame(frame) {
if (frame.type === "response") {
const entry = pending.get(frame.id);
if (!entry) return;
// Map an indeterminate outcome to a non-retryable 409, the same contract
// the file gateway applies through the outcome header.
const statusCode =
frame.outcome === "indeterminate" ? 409 : typeof frame.status === "number" ? frame.status : 200;
const headers = frame.headers || {};
if (frame.bodyByteCount === 0) {
// A zero-length body accepts no body_chunk, so the response completes now.
pending.delete(frame.id);
responseAssembly.delete(frame.id);
clearTimeout(entry.timer);
entry.resolve({ status: statusCode, headers: headers, body: "" });
return;
}
// The body rides body_chunk frames that share this id. Record the envelope
// and wait for the chunks.
responseAssembly.set(frame.id, {
status: statusCode,
headers: headers,
bodyByteCount: frame.bodyByteCount,
received: 0,
nextSeq: 0,
chunks: [],
});
return;
}
if (frame.type === "body_chunk") {
const asm = responseAssembly.get(frame.id);
if (!asm) return;
if (frame.seq !== asm.nextSeq) {
failRequest(frame.id, "duplex response body_chunk seq is out of order");
return;
}
const decoded = Buffer.from(frame.data, "base64");
if (decoded.toString("base64") !== frame.data) {
failRequest(frame.id, "duplex response body_chunk is not canonical base64");
return;
}
if (decoded.length === 0) {
failRequest(frame.id, "duplex response body_chunk is empty");
return;
}
if (asm.received + decoded.length > asm.bodyByteCount) {
failRequest(frame.id, "duplex response body overruns the declared size");
return;
}
const isFinal = asm.received + decoded.length === asm.bodyByteCount;
if (
decoded.length > DUPLEX_BODY_CHUNK_RAW_BYTES ||
(!isFinal && decoded.length !== DUPLEX_BODY_CHUNK_RAW_BYTES)
) {
failRequest(frame.id, "duplex response body_chunk has the wrong size");
return;
}
asm.nextSeq += 1;
asm.received += decoded.length;
asm.chunks.push(decoded);
if (asm.received === asm.bodyByteCount) {
const entry = pending.get(frame.id);
responseAssembly.delete(frame.id);
if (!entry) return;
pending.delete(frame.id);
clearTimeout(entry.timer);
entry.resolve({
status: asm.status,
headers: asm.headers,
body: Buffer.concat(asm.chunks).toString("utf8"),
});
}
return;
}
if (frame.type === "close") {
triggerLoss("host sent a close frame");
return;
}
// A heartbeat, ready, error, or request frame proves liveness but needs no
// local action here.
}
const decoder = new DuplexFrameDecoder({ maxAggregateBytes: maxDuplexDecoderBytes });
process.stdin.on("data", (chunk) => {
lastInboundAt = Date.now();
for (const result of decoder.push(chunk)) {
if (!result.ok) {
diag("dropped an inbound frame: " + result.error.code);
continue;
}
handleInboundFrame(result.frame);
}
});
process.stdin.on("end", () => triggerLoss("stdin reached end of file"));
process.stdin.on("error", (error) =>
triggerLoss("stdin error: " + (error && error.message ? error.message : String(error))),
);
const heartbeatSendTimer = setInterval(() => {
writeFrame({ version: DUPLEX_FRAME_VERSION, type: "heartbeat" });
}, heartbeatIntervalMs);
const heartbeatWatchTimer = setInterval(() => {
if (Date.now() - lastInboundAt > heartbeatTimeoutMs) {
triggerLoss("no inbound frame within the heartbeat timeout");
}
}, Math.max(200, Math.floor(heartbeatTimeoutMs / 4)));
const server = createServer(async (req, res) => {
try {
const auth = req.headers.authorization || "";
const receivedToken = auth.startsWith("Bearer ") ? auth.slice("Bearer ".length) : "";
if (!tokensMatch(receivedToken)) {
writeJsonResponse(res, 401, { error: "Invalid bridge token." });
return;
}
if (unavailable) {
writeJsonResponse(res, 503, { error: "bridge_unavailable" });
return;
}
if (pending.size >= maxQueueDepth) {
writeJsonResponse(res, 503, { error: "Bridge request queue is full." });
return;
}
const url = new URL(req.url || "/", "http://127.0.0.1");
const contentType = typeof req.headers["content-type"] === "string" ? req.headers["content-type"] : "";
if (req.method && req.method !== "GET" && req.method !== "HEAD" && !/json/i.test(contentType)) {
writeJsonResponse(res, 415, { error: "Bridge only accepts JSON request bodies." });
return;
}
const requestId = randomUUID();
const requestBodyBuffer = Buffer.from(await readBody(req), "utf8");
if (unavailable) {
writeJsonResponse(res, 503, { error: "bridge_unavailable" });
return;
}
// The request body rides body_chunk frames. The envelope carries only the
// raw byte count, so the envelope stays small and the body splits into
// fixed-size slices that each stay under the frame bound.
const requestFrame = {
version: DUPLEX_FRAME_VERSION,
type: "request",
id: requestId,
method: req.method || "GET",
path: url.pathname,
query: url.search,
headers: normalizeHeaders(req.headers),
bodyByteCount: requestBodyBuffer.length,
};
const encodedRequest = encodeDuplexFrameChecked(requestFrame);
if (!encodedRequest.ok) {
writeJsonResponse(res, 413, { error: "request_too_large" });
return;
}
// Pre-encode every body_chunk frame and enforce the frame size bound on
// each. A fixed raw slice never exceeds the bound, so this guard is
// defensive. On a rejection, fail this one local request with a clean 413.
// The frames never leave the gateway, so no other in-flight request is
// affected and the channel stays open.
const chunkFrames = splitDuplexBodyIntoChunks(requestId, requestBodyBuffer);
const encodedChunks = [];
let chunkTooLarge = false;
for (const chunk of chunkFrames) {
const encodedChunk = encodeDuplexFrameChecked(chunk);
if (!encodedChunk.ok) {
chunkTooLarge = true;
break;
}
encodedChunks.push(encodedChunk.line);
}
if (chunkTooLarge) {
writeJsonResponse(res, 413, { error: "request_too_large" });
return;
}
const response = await new Promise((resolve) => {
const timer = setTimeout(() => {
responseAssembly.delete(requestId);
if (pending.delete(requestId)) {
resolve({
status: 502,
headers: { "content-type": "application/json" },
body: JSON.stringify({ error: "Timed out waiting for host bridge response." }),
});
}
}, responseTimeoutMs);
pending.set(requestId, { resolve: resolve, timer: timer });
process.stdout.write(encodedRequest.line);
for (const line of encodedChunks) process.stdout.write(line);
});
res.statusCode = typeof response.status === "number" ? response.status : 200;
for (const [key, value] of Object.entries(response.headers || {})) {
if (typeof value !== "string" || key.toLowerCase() === "content-length") continue;
res.setHeader(key, value);
}
res.end(typeof response.body === "string" ? response.body : "");
} catch (error) {
writeJsonResponse(res, 502, { error: error instanceof Error ? error.message : String(error) });
}
});
process.on("SIGINT", () => {
try {
server.close();
} catch (error) {
diag("server close error: " + (error && error.message ? error.message : String(error)));
}
process.exit(0);
});
process.on("SIGTERM", () => {
try {
server.close();
} catch (error) {
diag("server close error: " + (error && error.message ? error.message : String(error)));
}
process.exit(0);
});
// Bind-or-exit. The host assigns a positive loopback port. The gateway binds
// exactly that port. On a non-positive assigned port or a bind failure it
// writes a diagnostic and exits with a nonzero code. It never selects a
// different port, so no untrusted workload can steer the endpoint.
if (!Number.isInteger(port) || port <= 0) {
diag("duplex gateway requires a positive assigned PAPERCLIP_BRIDGE_PORT; got " + String(port));
process.exit(1);
}
server.on("error", (error) => {
diag("duplex gateway could not bind port " + String(port) + ": " + (error && error.message ? error.message : String(error)));
process.exit(1);
});
server.listen(port, host, () => {
const address = server.address();
if (!address || typeof address === "string") {
diag("duplex gateway did not expose a TCP address");
process.exit(1);
return;
}
// The gateway sends READY only after the listener binds. READY is a liveness
// signal: it carries the frame version and the echoed nonce, and no address
// data. The host builds the endpoint from its own stored port. Stdout carries
// only frames.
writeFrame({
version: DUPLEX_FRAME_VERSION,
type: "ready",
nonce: bridgeNonce,
});
// READY is on the wire, so the host will adopt this process. From here on
// an uncaught fault must not kill the listener.
gatewayReady = true;
});
}
// ---------------------------------------------------------------------------
// http2_v1: run one Node HTTP/2 client session directly on stdin/stdout.
//
@ -3160,15 +2786,14 @@ function runHttp2Gateway() {
});
}
// The startup check above already rejected every value except http2 and
// queue, so this dispatch names both modes explicitly and never falls
// through to the queue gateway for an unsupported mode.
if (bridgeMode === "${SANDBOX_CALLBACK_BRIDGE_HTTP2_MODE}") {
runHttp2Gateway();
} else if (bridgeMode === "${SANDBOX_CALLBACK_BRIDGE_DUPLEX_MODE}") {
// No host selection path sets this mode anymore (http2_v1 replaced it), but
// the generated gateway keeps the mode reachable: it stays defined here,
// unchanged, so nothing that still spawns the gateway directly with this
// mode name breaks.
runDuplexGateway();
} else {
} else if (bridgeMode === "${SANDBOX_CALLBACK_BRIDGE_FILE_MODE}") {
await runFileGateway();
} else {
throw new Error("Unsupported PAPERCLIP_API_BRIDGE_MODE: " + bridgeMode);
}`;
}