Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466import * as Effect from "effect/Effect";import assert from "node:assert/strict";import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";import { createServer } from "node:net";import { tmpdir } from "node:os";import { join, resolve } from "node:path";import { test } from "node:test";import { AgentClient } from "agents/client";import { WebSocketChatTransport } from "agents/chat/transport";import { MessageType } from "agents/chat";import WebSocket from "ws";import { unstable_dev } from "wrangler";import { Secret } from "../configuration/secrets.ts";import { createOwnerSession } from "../worker/session.ts";import { customerBindings, installation } from "./fixtures/config.mjs";
async function waitFor(predicate) { const deadline = Date.now() + 40_000; while (!(await predicate())) { if (Date.now() > deadline) throw new Error("Timed out waiting for tool activity"); await new Promise((resolve) => setTimeout(resolve, 25)); }}async function freePort() { const server = createServer(); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const { port } = server.address(); await new Promise((resolve) => server.close(resolve)); return port;}
test( "native Sandbox shell streams, isolates and destroys temporary workspaces", { timeout: 360_000 }, async () => { const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const cookie = ( await Effect.runPromise( createOwnerSession( new Secret(customerBindings.FLAREBOT_SESSION_SECRET), { ...installation, runtimeOrigin: origin }, ), ) ).split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; const persistence = await mkdtemp(join(tmpdir(), "flarebot-shell-")); const config = JSON.parse( await readFile("dist/release/deployment.json", "utf8"), ); const configPath = join(persistence, "wrangler.json"); await writeFile( configPath, JSON.stringify({ ...config, name: "flarebot-shell-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/shell-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/shell-worker.ts", { config: configPath, vars: { ...customerBindings, FLAREBOT_ENV: "development", FLAREBOT_DEV_OVERRIDES: JSON.stringify({ runtimeOrigin: origin }), }, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: persistence, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false, enableContainers: true, }, }); const clients = []; let worker; try { worker = await start(); await worker.fetch("/"); console.log("Shell fixture ready"); class OwnerSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } async function connect(id) { const states = []; const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", ...(id ? { basePath: `agents/personal-agent/personal/sub/conversation/${id}`, } : { name: "personal" }), WebSocket: OwnerSocket, onStateUpdate: (state) => states.push(state), }); clients.push(client); const transport = new WebSocketChatTransport({ agent: client }); client.addEventListener("message", (event) => { const frame = JSON.parse(event.data); if (frame.type === MessageType.CF_AGENT_STREAM_RESUMING) transport.handleStreamResuming(frame); if (frame.type === MessageType.CF_AGENT_STREAM_RESUME_NONE) transport.handleStreamResumeNone(frame); if (frame.type === MessageType.CF_AGENT_STREAM_PENDING) transport.handleStreamPending(); }); await waitFor(() => states.length > 0); await client.ready; return { client, transport, states }; } let owner = await connect(); const first = await owner.client.call("createConversation", [ "Remembering", ]); let connection = await connect(first.id); const history = async (id) => ( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${id}/get-messages`, { headers }, ) ).json(); async function send(target, id, text) { const stream = await target.transport.sendMessages({ chatId: id, trigger: "submit-message", messages: [ ...(await history(id)), { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ], abortSignal: new AbortController().signal, }); const chunks = []; const done = (async () => { for await (const chunk of stream) chunks.push(chunk); })(); return { done, chunks }; } async function call( tool, input, id = crypto.randomUUID(), target = connection, conversationId = first.id, ) { const turn = await send( target, conversationId, "memory:" + JSON.stringify({ tool, input, id }), ); await turn.done; return (await history(conversationId)) .flatMap((message) => message.parts) .findLast((part) => part.toolCallId === id && "output" in part) ?.output; } const status = () => fetch(`${origin}/__fixture/shell-status`).then((r) => r.json()); const fault = (query) => fetch(`${origin}/__fixture/shell-fault?${query}`); const shell = (command, extra = {}) => call("shell", { command, ...extra }); const firstCall = crypto.randomUUID(); const streaming = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "shell", id: firstCall, input: { command: "printf 'first-out\\n'; sleep 1; printf 'first-err\\n' >&2; sleep 1; printf 'last-out\\n'", }, }), ); await waitFor(() => streaming.chunks.some( (chunk) => chunk.type === "tool-output-available" && chunk.preliminary && chunk.output.stdout.includes("first-out"), ), ); assert.ok( !streaming.chunks.some( (chunk) => chunk.type === "tool-output-available" && chunk.output.stdout?.includes("last-out"), ), "first output arrives before final output", ); await streaming.done; const result = streaming.chunks.findLast( (chunk) => chunk.type === "tool-output-available" && !chunk.preliminary, )?.output; assert.equal( result?.status, "succeeded", JSON.stringify(streaming.chunks), ); assert.equal(result.stdout, "first-out\nlast-out\n"); assert.equal(result.stderr, "first-err\n"); assert.equal(result.exitCode, 0); assert.equal(result.cleanup, "closed"); assert.equal((await status()).leases.length, 0); console.log("Native shell streaming passed"); const noContainers = async () => { await waitFor(async () => { const current = await status(); return ( current.leases.length === 0 && current.containers.every(({ state }) => state.status.startsWith("stopped"), ) ); }); }; const failed = await shell("printf 'useful-error' >&2; exit 7"); assert.equal(failed.exitCode, 7, JSON.stringify(failed)); assert.equal(failed.stderr.trim(), "useful-error"); assert.equal(failed.status, "failed"); const files = await shell(`cat > /workspace/data.txt <<'DATA'34DATAcat > /workspace/process.js <<'SCRIPT'const fs = require('fs'); console.log(fs.readFileSync('/workspace/data.txt', 'utf8').trim().split('\\n').map(Number).reduce((a,b)=>a+b,0));SCRIPTnode /workspace/process.js`); assert.equal(files.stdout.trim(), "7"); assert.equal(files.status, "succeeded"); const fresh = await shell( "test ! -e /workspace/data.txt && test ! -e /workspace/process.js && printf fresh", ); assert.equal(fresh.stdout.trim(), "fresh"); assert.equal( (await shell("env | sort")).stdout.includes( customerBindings.FLAREBOT_SESSION_SECRET, ), false, ); await noContainers();
const outputParts = (turn) => turn.chunks.filter((chunk) => chunk.type === "tool-output-available"); const begin = ( command, timeoutMs = 30000, target = connection, id = first.id, ) => send( target, id, "memory:" + JSON.stringify({ tool: "shell", id: crypto.randomUUID(), input: { command, timeoutMs }, }), ); const cancelled = await begin("printf 'retained-partial\\n'; sleep 30"); await waitFor(() => outputParts(cancelled).some((part) => part.output.stdout.includes("retained-partial"), ), ); const stopStarted = Date.now(); connection.transport.cancelActiveServerTurn(); await cancelled.done.catch((error) => assert.equal(error.name, "AbortError"), ); await noContainers(); assert.ok( Date.now() - stopStarted < 8000, "cancellation kills a silent running command promptly", ); assert.ok( (await history(first.id)) .flatMap((message) => message.parts) .some((part) => part.output?.stdout?.includes("retained-partial")), "native history retains partial shell output", ); const reopened = await connect(first.id); assert.ok( (await history(first.id)) .flatMap((message) => message.parts) .some((part) => part.output?.stdout?.includes("retained-partial")), ); reopened.client.close();
const timeout = await shell("printf 'before-timeout\\n'; sleep 30", { timeoutMs: 2000, }); assert.equal(timeout.error, "timeout", JSON.stringify(timeout)); assert.equal(timeout.stdout.trim(), "before-timeout"); await noContainers(); const flooded = await shell( "node -e 'setInterval(()=>process.stdout.write(\"🐟\".repeat(4096)), 10)'", ); assert.equal(flooded.error, "output_limit"); assert.equal(flooded.truncated, true); assert.ok(Buffer.byteLength(flooded.stdout + flooded.stderr) <= 32768); await noContainers(); console.log("Shell files, cancellation and limits passed");
const second = await owner.client.call("createConversation", [ "Separate workspace", ]); const sibling = await connect(second.id); const isolatedA = await begin( "printf a > /workspace/a-marker; printf 'a-ready\\n'; sleep 3; test ! -e /workspace/b-marker && printf 'a-isolated'", ); await waitFor(() => outputParts(isolatedA).some((part) => part.output.stdout.includes("a-ready"), ), ); const isolatedB = await begin( "printf b > /workspace/b-marker; test ! -e /workspace/a-marker && printf 'b-isolated'; sleep 1", 30000, sibling, second.id, ); await Promise.all([isolatedA.done, isolatedB.done]); assert.equal( outputParts(isolatedA).at(-1).output.stdout.trim(), "a-ready\na-isolated", ); assert.equal( outputParts(isolatedB).at(-1).output.stdout.trim(), "b-isolated", ); await noContainers();
await fault("destroy=4"); const pending = await shell("printf cleanup"); assert.equal(pending.cleanup, "pending"); assert.ok((await status()).leases.length > 0); await noContainers(); console.log("Shell isolation and durable cleanup retries passed");
const deleted = await owner.client.call("createConversation", [ "Delete during launch", ]); const deleting = await connect(deleted.id); await fault("launch=2000"); const late = await begin( "printf late; sleep 30", 30000, deleting, deleted.id, ); late.done.catch(() => {}); await waitFor(async () => (await status()).leases.some( (lease) => lease.conversation_id === deleted.id && lease.status === "launching", ), ); const lateId = (await status()).leases.find( (lease) => lease.conversation_id === deleted.id, ).id; await owner.client.call("deleteConversation", [deleted.id]); await noContainers(); const lateEvents = (await status()).events.filter( (event) => event.id === lateId, ); assert.ok( lateEvents.filter((event) => event.event === "destroy-settled") .length >= 2, JSON.stringify(lateEvents), ); assert.equal( ( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${deleted.id}/get-messages`, { headers }, ) ).status, 404, ); console.log("Deletion during shell launch passed"); const restarting = await owner.client.call("createConversation", [ "Restart during native launch", ]); const restartConnection = await connect(restarting.id); await fault("nativeLaunch=2500"); const restartTurn = await begin( "printf 'after-restart'; sleep 30", 30000, restartConnection, restarting.id, ); restartTurn.done.catch(() => {}); await waitFor(async () => (await status()).events.some( (event) => event.event === "native-launch-waiting", ), ); const restartingId = (await status()).leases.find( (lease) => lease.conversation_id === restarting.id, ).id; await fetch(`${origin}/__fixture/shell-restart`); owner = await connect(); await waitFor( async () => !(await status()).leases.some((lease) => lease.id === restartingId), ); await waitFor(async () => (await status()).containers .find((container) => container.id === restartingId) ?.state.status.startsWith("stopped"), ); await owner.client.call("deleteConversation", [restarting.id]); await noContainers(); console.log("Parent restart during native shell launch passed"); const persisted = await history(first.id); for (const client of clients) client.close(); await worker.stop(); worker = await start(); await worker.fetch("/"); owner = await connect(); connection = await connect(first.id); assert.deepEqual(await history(first.id), persisted); assert.equal( ( await shell("test ! -e /workspace/a-marker && printf restarted") ).stdout.trim(), "restarted", ); await noContainers(); const activity = await connection.client.call("listToolActivities"); assert.ok(activity.activities.some((entry) => entry.kind === "shell")); assert.doesNotMatch( JSON.stringify(activity), /retained-partial|useful-error|process.js|first-out/, ); } finally { for (const client of clients) client.close(); if (worker) await worker.stop(); await rm(persistence, { recursive: true, force: true }); } },);