Files
PaperClipAI/packages/paperclip-runner/devtools/issue-thread/stream-fixture-transport.mjs
Dotta 560e7e48b5 feat(runner): add SDK and developer tooling (#12608)
## Thinking Path

> - Paperclip is the open source app people use to manage AI agents for
work.
> - The runner package already provides the production protocol and
execution spine.
> - Contributors still need stable SDK surfaces, deterministic test
tools, and local inspection tools.
> - Those surfaces share generated contracts and must change as one
package boundary.
> - This pull request adds the package-local SDK, labs, examples, and
drift checks.
> - The benefit is a reviewable developer platform that does not change
application execution selection.

## Linked Issues or Issue Description

**Subsystem affected**

`packages/paperclip-runner` — runner SDK, conformance tools, and
developer tooling.

**Problem or motivation**

The production runner spine is present, but package consumers cannot
build deterministic integrations, inspect sessions, or verify
provider-neutral behavior through supported surfaces.

**Proposed solution**

Add browser, React, standalone, live-session, scenario, conformance, and
evaluation surfaces. Add generated contract inventories and
package-local verification scripts. Keep production application routing
unchanged.

**Alternatives considered**

We considered splitting each generated catalog, SDK surface, and demo
into separate pull requests. Those changes share exports, fixtures, and
drift gates. Splitting them would create intermediate package states
that do not build.

**Roadmap alignment**

No overlapping item appears in `ROADMAP.md`. This work extends the
runner package that is already on `master`.

## What Changed

- Add browser, React, standalone, live-session, and issue-thread SDK
surfaces.
- Add deterministic mock control-plane, scenario, conformance, replay,
and evaluation tools.
- Add bounded Codex, OpenCode, and ACPX development transports and
fixtures.
- Keep deferred managed-provider execution fail-closed. Persisted
compatibility data remains readable.
- Add generated capability inventories with their source files and drift
checks.
- Add examples, package documentation, browser checks, and
clean-consumer checks.
- Preserve the reviewed protocol bounds, replay compatibility aliases,
process environment isolation, and semantic redaction limits.
- Update the ACPX package patch that the existing workspace patch
registry already tracks.
- Do not change `pnpm-lock.yaml`, repository workflows, server runtime
selection, or the application UI.

## Verification

GitHub Actions is the verification authority for this pull request. The
repository CI, package TypeScript and Rust checks, package tests,
generated-output drift checks, browser checks, security scans, and
Greptile review must pass on the exact head.

Local test suites were not run because this series uses parallel GitHub
Actions for verification.

## Risks

This is a large greenfield package change. The main risks are public
export drift, generated-output drift, and optional React consumer
compatibility. Package boundary checks, clean-consumer checks, and
browser tests cover those risks. Production adapter selection and server
execution are outside this pull request.

## Stack

1. **This PR:** runner SDK and developer tooling.
2. [Codex production server
integration](https://github.com/paperclipai/paperclip/pull/12616).
3. [Provider-neutral task-thread
UI](https://github.com/paperclipai/paperclip/pull/12617).

## Model Used

OpenAI Codex, GPT-5, high-reasoning mode, with tool use and code
execution.

## 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 described the issue in-PR following the feature request
template
- [x] I have not referenced internal/instance-local Paperclip issues or
links
- [x] My branch name describes the change and contains no internal
Paperclip ticket id
- [ ] 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 risks above
- [ ] All Paperclip CI gates are green
- [ ] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge
2026-08-31 21:33:11 -05:00

167 lines
5.3 KiB
JavaScript

/**
* Scripted Codex transport for the browser streaming test (track 7Q).
*
* Everything below the provider is the real thing: the same package server
* middleware, the same `CapabilityLiveSession`, the same NDJSON turn stream, and
* the same built browser bundle. Only the provider is scripted, and only so the
* test can decide when a delta arrives — a real Codex process cannot be asked
* to emit exactly four deltas 220 ms apart, and a browser assertion needs that
* to distinguish "streaming" from "arrived all at once".
*
* This module is loaded by `vite.issue-thread-stream.config.ts` only. It is not
* reachable from the shipped server, the deployed surface, or any route: the
* production plugin takes no transport override from a request or an
* environment variable.
*/
const DELTA_INTERVAL_MS = 220;
export const CAPABILITY_STREAM_FIXTURE_PRIVATE_REASONING =
"PRIVATE reasoning text must never reach the browser.";
/** Four deltas: enough for a browser to observe growth twice over, and short. */
export const CAPABILITY_STREAM_FIXTURE_DELTAS = [
"Reading the clean-room issue. ",
"It is blank, with one mock agent and one mock task. ",
"Recording a first status against the mock control plane. ",
"Done — every record stayed in the mock port.",
];
export const CAPABILITY_STREAM_FIXTURE_REPLY = CAPABILITY_STREAM_FIXTURE_DELTAS.join("");
function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
class Notifications {
#values = [];
#waiters = [];
#closed = false;
push(value) {
const waiter = this.#waiters.shift();
if (waiter) waiter({ value, done: false });
else this.#values.push(value);
}
close() {
this.#closed = true;
for (const waiter of this.#waiters.splice(0)) waiter({ value: undefined, done: true });
}
[Symbol.asyncIterator]() {
return {
next: async () => {
const value = this.#values.shift();
if (value) return { value, done: false };
if (this.#closed) return { value: undefined, done: true };
return new Promise((resolve) => this.#waiters.push(resolve));
},
};
}
}
class StreamFixtureTransport {
#queue = new Notifications();
#closed = false;
#turns = 0;
#interrupted = new Set();
async request(method, params) {
if (method === "initialize") return { user: { sessionId: "stream-fixture-session" } };
if (method === "thread/start" || method === "thread/read" || method === "thread/resume") {
return { thread: { id: "stream-fixture-thread", sessionId: "stream-fixture-session" } };
}
if (method === "turn/start") {
this.#turns += 1;
const turnId = `stream-turn-${this.#turns}`;
void this.#runTurn(turnId);
return { turn: { id: turnId, status: "inProgress" } };
}
if (method === "turn/interrupt") {
// A stopped turn keeps whatever it had already streamed, so the surface
// can be checked for a coherent partial reply.
this.#interrupted.add(String(params.turnId));
this.#queue.push({
method: "turn/completed",
params: {
threadId: "stream-fixture-thread",
turn: { id: String(params.turnId), status: "interrupted" },
},
});
return {};
}
throw new Error(`unsupported scripted Codex method ${method}`);
}
notify() {}
notifications() {
return this.#queue;
}
setServerRequestHandler() {}
async close() {
this.#closed = true;
this.#queue.close();
}
processInfo() {
return { pid: 7100, processGroupId: 7100, exited: this.#closed, exitCode: null, signal: null };
}
async #runTurn(turnId) {
this.#queue.push({
method: "turn/started",
params: { threadId: "stream-fixture-thread", turn: { id: turnId, status: "inProgress" } },
});
await sleep(Math.floor(DELTA_INTERVAL_MS / 2));
if (this.#closed || this.#interrupted.has(turnId)) return;
this.#queue.push({
method: "item/reasoning/summaryTextDelta",
params: {
threadId: "stream-fixture-thread",
turnId,
delta: CAPABILITY_STREAM_FIXTURE_PRIVATE_REASONING,
},
});
for (const delta of CAPABILITY_STREAM_FIXTURE_DELTAS) {
await sleep(DELTA_INTERVAL_MS);
if (this.#closed || this.#interrupted.has(turnId)) return;
this.#queue.push({
method: "item/agentMessage/delta",
params: { threadId: "stream-fixture-thread", turnId, delta },
});
}
await sleep(DELTA_INTERVAL_MS);
if (this.#closed || this.#interrupted.has(turnId)) return;
this.#queue.push({
method: "item/completed",
params: {
threadId: "stream-fixture-thread",
turnId,
item: { id: `message-${turnId}`, type: "agentMessage", text: CAPABILITY_STREAM_FIXTURE_REPLY },
},
});
this.#queue.push({
method: "turn/completed",
params: { threadId: "stream-fixture-thread", turn: { id: turnId, status: "completed" } },
});
}
}
export function capabilityStreamFixtureTransportFactory(options = {}) {
const evidence = {
runnerPid: 7100,
runnerProcessGroupId: 7100,
codexPid: 7200,
runnerExited: false,
runnerExitCode: null,
runnerSignal: null,
childEnvironmentKeys: ["CODEX_HOME", "HOME", "PATH"],
diagnostics: [],
};
options.onEvidence?.(evidence);
return { transport: new StreamFixtureTransport(), evidence: () => ({ ...evidence }) };
}