Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 287 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import { execFileSync } from "node:child_process";
import { afterEach, describe, expect, test } from "vitest";
import { dockerWarningsSatisfyHarness, runContainerHarness } from "../src/agents/harness/container-launcher.js";import type { ProviderProfile } from "../src/agents/provider-profiles.js";
const roots: string[] = [];
afterEach(async () => { delete process.env.THOUGHTSTREAM_HARNESS_TEST_KEY; delete process.env.HOST_SENTINEL; await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe.skipIf(process.env.THOUGHTSTREAM_RUN_CONTAINER_TESTS !== "1")("workspace-v1 black-box containment", () => { test("fails closed on a host without Docker memory enforcement", async () => { const warnings = JSON.parse(execFileSync("docker", ["info", "--format", "{{json .Warnings}}"], { encoding: "utf8" })) as unknown; if (dockerWarningsSatisfyHarness(warnings)) return;
const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-harness-runtime-")); roots.push(root); const workspaceRoot = path.join(root, "workspace-leases"); const stateRoot = path.join(root, "state-leases"); const workspace = path.join(workspaceRoot, "lease-a"); const state = path.join(stateRoot, "lease-a"); await Promise.all([fs.mkdir(workspace, { recursive: true }), fs.mkdir(state, { recursive: true })]); const receipt = await runContainerHarness({ runId: "run-runtime-admission", image: "thoughtstream/pi-coding-harness:local", trustedImageId: fixtureImageId(), workspaceLeaseRoot: workspaceRoot, workspacePath: workspace, workspaceIdentity: "workspace-a", stateLeaseRoot: stateRoot, statePath: state, session: { mode: "new" }, systemPrompt: "Do not run.", prompt: "Do not run.", tools: [], providerProfile: fixtureProfile(), model: { id: "fixture-model", reasoning: false, contextWindow: 32_000, maxOutputTokens: 1_000 }, maxProviderRequests: 1, maxTotalProviderRequestBytes: 64 * 1024, maxTotalProviderResponseBytes: 64 * 1024, timeoutMs: 5_000, }); expect(receipt).toMatchObject({ status: "failed", errorCode: "runtime-unavailable" }); });
test("runs Pi tools while denying host secrets, sibling files, network, and root writes", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-harness-test-")); roots.push(root); const workspaceRoot = path.join(root, "workspace-leases"); const stateRoot = path.join(root, "state-leases"); const workspace = path.join(workspaceRoot, "lease-a"); const state = path.join(stateRoot, "lease-a"); await Promise.all([ fs.mkdir(workspace, { recursive: true }), fs.mkdir(state, { recursive: true }), ]); await fs.writeFile(path.join(workspaceRoot, "host-canary"), "HOST_CANARY_VALUE"); await fs.symlink("../host-canary", path.join(workspace, "escape-link")); await fs.mkdir(path.join(workspace, ".pi", "extensions"), { recursive: true }); await fs.writeFile( path.join(workspace, ".pi", "extensions", "hostile.mjs"), "import fs from 'node:fs'; fs.writeFileSync('/workspace/extension-ran', 'unsafe');\n", ); const dockerPath = await createDockerTestWrapper(root); process.env.THOUGHTSTREAM_HARNESS_TEST_KEY = "BROKER_SECRET_VALUE"; process.env.HOST_SENTINEL = "HOST_ENV_SECRET_VALUE";
let providerCalls = 0; const fetchImpl: typeof fetch = async (_input, init) => { providerCalls += 1; expect(new Headers(init?.headers).get("authorization")).toBe("Bearer BROKER_SECRET_VALUE"); if (providerCalls === 1) return sseResponse(toolCall("bash", { command: [ "set +e", "printf 'env=' > containment.txt", "env | grep -E 'BROKER_SECRET_VALUE|HOST_ENV_SECRET_VALUE' >> containment.txt", "printf '\\nprocess_env=' >> containment.txt", "cat /proc/1/environ | tr '\\0' '\\n' | grep -E 'BROKER_SECRET_VALUE|HOST_ENV_SECRET_VALUE' >> containment.txt", "printf '\\ncanary=' >> containment.txt", "cat escape-link >> containment.txt 2>/dev/null", "printf '\\nroot_write=' >> containment.txt", "touch /escaped-host-file 2>> containment.txt", "printf '\\nnetwork=' >> containment.txt", "node -e \"fetch('http://127.0.0.1:9',{signal:AbortSignal.timeout(500)}).then(()=>console.log('reachable')).catch(()=>console.log('blocked'))\" >> containment.txt 2>&1", "printf '\\nexternal_network=' >> containment.txt", "node -e \"fetch('http://1.1.1.1',{signal:AbortSignal.timeout(500)}).then(()=>console.log('reachable')).catch(()=>console.log('blocked'))\" >> containment.txt 2>&1", "printf '\\ndocker=' >> containment.txt", "test -S /var/run/docker.sock && echo present >> containment.txt || echo absent >> containment.txt", "printf '\\ndata_kib=' >> containment.txt", "ulimit -d >> containment.txt", "exit 0", ].join("; "), timeout: 5, })); return sseResponse(textChunk("containment probe complete")); };
const receipt = await runContainerHarness({ runId: "run-container-containment", image: "thoughtstream/pi-coding-harness:local", trustedImageId: fixtureImageId(), dockerPath, workspaceLeaseRoot: workspaceRoot, workspacePath: workspace, workspaceIdentity: "workspace-a", stateLeaseRoot: stateRoot, statePath: state, session: { mode: "new" }, systemPrompt: "Execute the requested probe with the bash tool, then report completion.", prompt: "Run the containment probe.", tools: ["bash"], providerProfile: fixtureProfile(), model: { id: "fixture-model", reasoning: false, contextWindow: 32_000, maxOutputTokens: 1_000 }, maxProviderRequests: 4, maxTotalProviderRequestBytes: 256 * 1024, maxTotalProviderResponseBytes: 256 * 1024, timeoutMs: 30_000, fetchImpl, });
expect(receipt.status, JSON.stringify(receipt)).toBe("completed"); expect(receipt.result).toMatchObject({ status: "completed", finalText: "containment probe complete" }); expect(receipt.providerUsage?.requests).toBe(2); expect(providerCalls).toBe(2); const proof = await fs.readFile(path.join(workspace, "containment.txt"), "utf8"); expect(proof).not.toContain("BROKER_SECRET_VALUE"); expect(proof).not.toContain("HOST_ENV_SECRET_VALUE"); expect(proof).not.toContain("HOST_CANARY_VALUE"); expect(proof).toContain("Read-only file system"); expect(proof).toContain("network=blocked"); expect(proof).toContain("external_network=blocked"); expect(proof).toContain("docker=absent"); expect(proof).toContain("data_kib=1048576"); await expect(fs.stat(path.join(workspace, "extension-ran"))).rejects.toThrow(); await expect(fs.stat("/escaped-host-file")).rejects.toThrow(); }, 45_000);
test("persists only compatible sessions and rejects image or workspace identity swaps", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-harness-session-")); roots.push(root); const workspaceRoot = path.join(root, "workspace-leases"); const stateRoot = path.join(root, "state-leases"); const workspace = path.join(workspaceRoot, "lease-a"); const state = path.join(stateRoot, "lease-a"); await Promise.all([fs.mkdir(workspace, { recursive: true }), fs.mkdir(state, { recursive: true })]); const dockerPath = await createDockerTestWrapper(root); process.env.THOUGHTSTREAM_HARNESS_TEST_KEY = "fixture-secret"; const imageId = fixtureImageId();
const common = { image: "thoughtstream/pi-coding-harness:local", trustedImageId: imageId, dockerPath, workspaceLeaseRoot: workspaceRoot, workspacePath: workspace, workspaceIdentity: "workspace-a", stateLeaseRoot: stateRoot, statePath: state, systemPrompt: "Reply concisely.", prompt: "Continue.", tools: [] as [], providerProfile: fixtureProfile(), model: { id: "fixture-model", reasoning: false, contextWindow: 32_000, maxOutputTokens: 1_000 }, maxProviderRequests: 2, maxTotalProviderRequestBytes: 128 * 1024, maxTotalProviderResponseBytes: 128 * 1024, timeoutMs: 30_000, }; const first = await runContainerHarness({ ...common, runId: "run-session-first", session: { mode: "new" }, fetchImpl: async () => sseResponse(textChunk("first")), }); expect(first.status, JSON.stringify(first)).toBe("completed"); if (!first.result || first.result.status !== "completed") throw new Error("Missing first session result");
const resumed = await runContainerHarness({ ...common, runId: "run-session-resume", session: { mode: "resume", sessionId: first.result.sessionId }, fetchImpl: async () => sseResponse(textChunk("second")), }); expect(resumed.result).toMatchObject({ status: "completed", sessionId: first.result.sessionId, finalText: "second" });
let mismatchedProviderCalls = 0; const mismatch = await runContainerHarness({ ...common, runId: "run-session-mismatch", workspaceIdentity: "workspace-b", session: { mode: "resume", sessionId: first.result.sessionId }, fetchImpl: async () => { mismatchedProviderCalls += 1; return sseResponse(textChunk("must not run")); }, }); expect(mismatch).toMatchObject({ status: "failed", errorCode: "session-policy-rejected" }); expect(mismatchedProviderCalls).toBe(0);
const untrustedImage = await runContainerHarness({ ...common, runId: "run-image-mismatch", image: "node:22-bookworm-slim", session: { mode: "new" }, fetchImpl: async () => { throw new Error("broker must not start"); }, }); expect(untrustedImage).toMatchObject({ status: "failed", errorCode: "image-rejected" }); }, 45_000);});
function fixtureProfile(): ProviderProfile { return { id: "fixture", provider: "openai-compatible", baseUrl: "https://fixture.invalid/v1", route: "/chat/completions", apiKeyEnv: "THOUGHTSTREAM_HARNESS_TEST_KEY", allowedModels: new Set(["fixture-model"]), imageInputModels: new Set(), jsonObjectResponseFormat: false, jsonSchemaResponseFormat: false, requestTimeoutMs: 5_000, maxRequestBytes: 64 * 1024, maxResponseBytes: 64 * 1024, };}
function fixtureImageId(): string { return execFileSync("docker", ["image", "inspect", "--format", "{{.Id}}", "thoughtstream/pi-coding-harness:local"], { encoding: "utf8", }).trim();}
async function createDockerTestWrapper(root: string): Promise<string> { const wrapper = path.join(root, "docker-test-wrapper"); await fs.writeFile(wrapper, [ "#!/bin/sh", "if [ \"$1\" = info ]; then", " printf '%s\\n' '[]'", " exit 0", "fi", "exec /usr/bin/docker \"$@\"", "", ].join("\n"), { mode: 0o700 }); return wrapper;}
function toolCall(name: string, args: Record<string, unknown>): string { return [ chunk({ role: "assistant", tool_calls: [{ index: 0, id: "call_fixture", type: "function", function: { name, arguments: JSON.stringify(args) }, }], }, null), chunk({}, "tool_calls"), ].map((value) => `data: ${JSON.stringify(value)}\n\n`).join("") + "data: [DONE]\n\n";}
function textChunk(text: string): string { return [ chunk({ role: "assistant", content: text }, null), chunk({}, "stop"), ].map((value) => `data: ${JSON.stringify(value)}\n\n`).join("") + "data: [DONE]\n\n";}
function sseResponse(body: string): Response { return new Response(body, { status: 200, headers: { "content-type": "text/event-stream" } });}
function chunk(delta: Record<string, unknown>, finishReason: string | null): Record<string, unknown> { return { id: "chatcmpl-harness-fixture", object: "chat.completion.chunk", created: 1, model: "fixture-model", choices: [{ index: 0, delta, finish_reason: finishReason }], };}