Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550import * 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";
// Match the native orchestrator fixture's failure-only response diagnostics.// Never include response bodies, query strings, cookies or request payloads.async function responseJson(response, method) { try { return await response.json(); } catch { const path = new URL(response.url).pathname; const contentType = response.headers.get("Content-Type") ?? "missing"; assert.fail( `${method} ${path} returned invalid JSON (HTTP ${response.status}; Content-Type: ${contentType})`, ); }}
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( "rendered research uses native Chromium, bounded sessions and durable cleanup", { timeout: 240_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-browser-")); 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-browser-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/browser-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/browser-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 }, }); const clients = []; let worker; try { worker = await start(); 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 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) => responseJson( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${id}/get-messages`, { headers }, ), "GET", ); 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/browser-status`).then((r) => responseJson(r, "GET"), ); const fault = (query) => fetch(`${origin}/__fixture/browser-fault?${query}`); const browser = (path, extra = {}) => call("browser_read", { url: `https://browser.fixture.example.com/${path}`, ...extra, }); for (const scenario of [ "late-create-register-fail", "register-delay", "register-fail", "metadata-stall", ]) { const started = Date.now(); const race = await fetch( `${origin}/__fixture/browser-race?scenario=${scenario}`, ).then((r) => responseJson(r, "GET")); assert.equal( race.deletes.length, 1, JSON.stringify({ scenario, race }), ); assert.equal(race.commands, scenario === "metadata-stall" ? 1 : 0); assert.ok(Date.now() - started < 7000); } console.log("Browser races passed"); const staticPage = await call("read_url", { url: "https://browser.fixture.example.com/dynamic", }); assert.equal(staticPage.ok, true); assert.doesNotMatch( staticPage.sources[0].content, /Evidence added by JavaScript/, ); const rendered = await browser("dynamic"); assert.equal(rendered.ok, true, JSON.stringify(rendered)); const source = rendered.sources[0]; assert.equal(source.sourceKind, "browser"); assert.equal(source.title, "Rendered title 🐟"); assert.equal( source.finalUrl, "https://browser.fixture.example.com/rendered-final", ); assert.equal( source.requestedUrl, "https://browser.fixture.example.com/dynamic", ); assert.match(source.content, /Evidence added by JavaScript/); assert.equal(source.links.length, 1); assert.equal( source.links[0].url, "https://browser.fixture.example.com/follow", ); assert.doesNotMatch(JSON.stringify(source), /untrusted.example/); assert.match(source.id, /^web-[a-f0-9]{64}$/); const followed = await call("browser_read", { url: source.links[0].url }); assert.equal(followed.sources[0].title, "Followed source"); assert.equal(followed.sources[0].content, "Follow-up evidence"); const selected = await browser("dynamic", { waitForSelector: "#ready" }); assert.match(selected.sources[0].content, /Evidence added by JavaScript/); const large = await browser("large"); assert.equal(large.sources[0].truncated, true); assert.ok(large.sources[0].content.length <= 16_000); assert.ok(!/[\uD800-\uDBFF]$/.test(large.sources[0].content)); console.log("Rendered evidence passed"); assert.equal((await browser("private")).code, "blocked_url"); assert.equal((await browser("error")).code, "http_error"); assert.equal( (await browser("dynamic", { waitForSelector: "[" })).code, "invalid_selector", ); for (const url of [ "file:///etc/passwd", "http://127.0.0.1/", "https://user:secret@example.com/", "http://[::1]/", "http://localhost/", ]) assert.equal((await call("browser_read", { url })).ok, false); assert.ok((await status()).sessions.every((s) => s.deleted)); assert.deepEqual((await status()).leases, []); // A cleanup failure is a visible failure; the parent's native schedule // retries deletion even though the original tool invocation is over. console.log("Validation passed"); await fault("deletes=1"); assert.equal((await browser("follow")).code, "cleanup_failed"); assert.equal((await status()).leases.length, 1); await waitFor(async () => (await status()).leases.length === 0); assert.ok((await status()).sessions.every((s) => s.deleted));
console.log("Cleanup retry passed"); async function startBrowser(target, id, path = "slow") { const toolId = crypto.randomUUID(); const turn = await send( target, id, "memory:" + JSON.stringify({ tool: "browser_read", input: { url: `https://browser.fixture.example.com/${path}`, waitForSelector: "#ready", }, id: toolId, }), ); turn.done.catch(() => {}); return { ...turn, toolId }; } const count = (await status()).sessions.length; const cancellation = await startBrowser(connection, first.id); await waitFor( async () => (await status()).requests.filter((u) => u.endsWith("/slow")) .length === 1, ); const during = ( await connection.client.call("listToolActivities") ).activities.find((a) => a.toolCallId === cancellation.toolId); assert.equal(during.kind, "browser"); assert.equal(during.status, "running"); assert.ok(during.progress); connection.client.close(); connection = await connect(first.id); assert.equal( (await connection.client.call("listToolActivities")).activities.find( (a) => a.toolCallId === cancellation.toolId, ).status, "running", ); const resumed = await connection.transport.reconnectToStream({ chatId: first.id, }); assert.ok(resumed); const resumeDone = (async () => { for await (const _ of resumed) { } })(); resumeDone.catch(() => {}); connection.transport.cancelActiveServerTurn(); await waitFor(async () => (await status()).sessions[count].deleted); assert.equal( (await connection.client.call("listToolActivities")).activities.find( (a) => a.toolCallId === cancellation.toolId, ).status, "cancelled", ); // Parallel conversations own different session IDs. Cancelling one never // closes its sibling, which remains live until explicitly stopped. console.log("Cancellation passed"); const second = await owner.client.call("createConversation", [ "Parallel browser", ]); const sibling = await connect(second.id); const firstRunning = await startBrowser(connection, first.id); const secondRunning = await startBrowser(sibling, second.id); await waitFor(async () => (await status()).leases.length === 2); connection.transport.cancelActiveServerTurn(); await waitFor(async () => (await status()).leases.length === 1); assert.equal( (await sibling.client.call("listToolActivities")).activities.find( (a) => a.toolCallId === secondRunning.toolId, ).status, "running", ); sibling.transport.cancelActiveServerTurn(); await waitFor(async () => (await status()).leases.length === 0); console.log("Parallel isolation passed"); // Real absolute tool timeout and session deletion, not a fake timer path. assert.equal( (await browser("slow", { waitForSelector: "#never" })).code, "timeout", ); await waitFor(async () => (await status()).leases.length === 0); assert.ok((await status()).sessions.every((s) => s.deleted)); console.log("Timeout passed"); // Native schedule closes a registered browser even without any active call. const seed = await fetch( `${origin}/__fixture/browser-seed?id=${first.id}&ms=1500`, ).then((r) => responseJson(r, "GET")); await waitFor( async () => (await status()).sessions.find((s) => s.id === seed)?.deleted, ); assert.deepEqual((await status()).leases, []); // Conversation deletion occurs before a remote create reply is delivered. const doomed = await owner.client.call("createConversation", [ "Deleted browser", ]); const doomedConnection = await connect(doomed.id); await fault("create=1500"); const beforeLate = (await status()).sessions.length; await startBrowser(doomedConnection, doomed.id); await waitFor(async () => (await status()).sessions.length > beforeLate); const lateSessionId = (await status()).sessions[beforeLate].id; await owner.client.call("deleteConversation", [doomed.id]); await waitFor( async () => ( await fetch( `${origin}/__fixture/browser-probe?id=${lateSessionId}`, ).then((r) => responseJson(r, "GET")) ).status === 404, ); assert.deepEqual((await status()).leases, []); await fault(""); assert.equal( ( await fetch(`${origin}/__fixture/inspect`).then((r) => responseJson(r, "GET"), ) ).facets.includes(doomed.id), false, ); const cleared = await owner.client.call("createConversation", [ "Cleared browser", ]); const clearConnection = await connect(cleared.id); await fault("create=1500"); const beforeClear = (await status()).sessions.length; await startBrowser(clearConnection, cleared.id); await waitFor(async () => (await status()).sessions.length > beforeClear); clearConnection.client.send( JSON.stringify({ type: MessageType.CF_AGENT_CHAT_CLEAR }), ); await waitFor(async () => (await history(cleared.id)).length === 0); await fault(""); const afterClear = await call( "browser_read", { url: source.links[0].url }, crypto.randomUUID(), clearConnection, cleared.id, ); assert.equal(afterClear.ok, true); await waitFor(async () => (await status()).sessions[beforeClear].deleted); assert.equal( (await clearConnection.client.call("listToolActivities")).activities .length, 1, ); const researchId = crypto.randomUUID(); const renderedBatch = await call( "research", { code: `async () => { const rendered = await codemode.browser_read({ url: "https://browser.fixture.example.com/dynamic", waitForSelector: "#ready" }); if (!rendered.ok) return rendered; const followed = await codemode.browser_read({ url: rendered.sources[0].links[0].url }); return { rendered, followed }; }`, }, researchId, ); assert.equal( renderedBatch.result.rendered.ok, true, JSON.stringify(renderedBatch), ); assert.match( renderedBatch.result.rendered.sources[0].content, /Evidence added by JavaScript/, ); assert.match( renderedBatch.result.followed.sources[0].content, /Follow-up evidence/, ); await waitFor(async () => (await status()).leases.length === 0); assert.ok((await status()).sessions.every((session) => session.deleted)); const nested = ( await connection.client.call("listToolActivities") ).activities.filter((activity) => activity.toolCallId.startsWith(`${researchId}:research:`), ); assert.equal(nested.length, 2); assert.ok( nested.every( (activity) => activity.kind === "browser" && activity.status === "succeeded" && activity.progress?.completed === 4, ), ); const batchCancelId = crypto.randomUUID(); const batchCancel = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "research", id: batchCancelId, input: { code: `async () => Promise.all([1, 2].map(() => codemode.browser_read({ url: "https://browser.fixture.example.com/slow", waitForSelector: "#never" })))`, }, }), ); batchCancel.done.catch(() => {}); await waitFor(async () => (await status()).leases.length === 2); connection.transport.cancelActiveServerTurn(); await assert.rejects(batchCancel.done, { name: "AbortError" }); await waitFor(async () => (await status()).leases.length === 0); assert.ok((await status()).sessions.every((session) => session.deleted)); const cancelledBatch = ( await connection.client.call("listToolActivities") ).activities.filter((activity) => activity.toolCallId.startsWith(batchCancelId), ); assert.equal(cancelledBatch.length, 3); assert.ok( cancelledBatch.every((activity) => activity.status === "cancelled"), ); // Durable source/activity snapshots reconnect unchanged. const saved = await history(first.id); const activities = (await connection.client.call("listToolActivities")) .activities; assert.doesNotMatch( JSON.stringify(activities), /fixture.example|Evidence|Rendered title|session_id|secret/, ); assert.ok( activities.some( (a) => a.kind === "browser" && a.status === "succeeded" && a.progress?.completed === 4, ), ); assert.ok( activities.some((a) => a.kind === "browser" && a.status === "failed"), ); assert.ok( activities.some( (a) => a.kind === "browser" && a.status === "cancelled", ), ); // Restart before expiry retains a lease. OnStart closes it before it can // become a reused browser. Local Chromium itself also stops on shutdown; // this verifies durable reconciliation, not remote service survival. await fetch(`${origin}/__fixture/browser-seed?id=${first.id}&ms=30000`); assert.equal((await status()).leases.length, 1); for (const client of clients) client.close(); await worker.stop(); worker = await start(); owner = await connect(); connection = await connect(first.id); assert.deepEqual((await status()).leases, []); assert.deepEqual(await history(first.id), saved); assert.deepEqual( (await connection.client.call("listToolActivities")).activities, activities, ); assert.equal((await browser("follow")).ok, true); } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);