Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598import * 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() + 20_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( "web tools execute with bounded AI search sources, cancellation and durable history", { timeout: 180_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-web-")); 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-web-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/web-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/web-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) => ( 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/web-status`).then((r) => r.json()); for (const scenario of [ "success", "pre-abort", "late-create", "late-connect", "cancel-command", "timeout-command", "cleanup-fail", ]) { const race = await fetch( `${origin}/__fixture/browser-race?scenario=${scenario}`, ).then((r) => r.json()); assert.equal( race.error, scenario === "success" ? "" : scenario === "cleanup-fail" ? "cleanup_failed" : scenario === "pre-abort" || scenario === "cancel-command" ? "cancelled" : "timeout", JSON.stringify({ scenario, race }), ); assert.equal(race.creates, scenario === "pre-abort" ? 0 : 1); assert.equal(race.deletes.length, race.creates); assert.equal( race.connects, ["pre-abort", "late-create"].includes(scenario) ? 0 : 1, ); if (race.connects) assert.equal(race.closed, 1); if (["late-create", "late-connect", "pre-abort"].includes(scenario)) assert.equal(race.commands, 0); } const nativeTimeout = await fetch( `${origin}/__fixture/native-fetch-timeout`, ).then((r) => r.json()); assert.equal(nativeTimeout.result.code, "timeout"); assert.equal( nativeTimeout.bodyAborts, 1, "SDK timeout stays active after headers", ); const page = await call("read_url", { url: "https://source.fixture.example.com/redirect", }); assert.equal(page.ok, true); assert.equal(page.sources.length, 1); const source = page.sources[0]; assert.equal(source.finalUrl, "https://other.fixture.example.com/html"); assert.equal( source.requestedUrl, "https://source.fixture.example.com/redirect", ); assert.equal(source.title, "Fish & Chips 🐟"); assert.match(source.content, /A < B & C. Café 🐟/); assert.doesNotMatch(source.content, /script-secret|navigation-secret/); assert.equal(source.sourceKind, "page"); assert.equal(source.truncated, false); assert.match(source.id, /^web-[a-f0-9]{64}$/); assert.ok(Number.isFinite(Date.parse(source.fetchedAt))); const same = await call("read_url", { url: source.finalUrl + "#section", }); assert.equal(same.sources[0].id, source.id); for (const url of [ "file:///etc/passwd", "https://name:secret@example.com", "http://127.0.0.1/secret", "http://[::ffff:127.0.0.1]/secret", "http://localhost/", "http://0x7f000001/", ]) assert.equal((await call("read_url", { url })).ok, false); for (const [path, code] of [ ["private", "blocked_url"], ["loop", "request_failed"], ["binary", "unsupported_content"], ["error", "http_error"], ]) { const failure = await call("read_url", { url: `https://source.fixture.example.com/${path}`, }); assert.equal(failure.code, code); assert.doesNotMatch(JSON.stringify(failure), /private-error-marker/); } assert.ok( !(await status()).calls.some((url) => url.includes("127.0.0.1")), ); const large = await call("read_url", { url: "https://source.fixture.example.com/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)); const search = await call("web_search", { query: "normal", limit: 5 }); assert.equal(search.ok, true, JSON.stringify(search)); assert.equal( search.sources.length, 1, "duplicate citations coalesce and unsafe destinations are excluded", ); assert.equal(search.sources[0].sourceKind, "search"); assert.equal(search.sources[0].content, ""); assert.deepEqual(search.generatedSummary, { text: "Generated search summary, not page text.", truncated: false, }); assert.equal(search.sources[0].title, "Fixture & research"); assert.equal( search.sources[0].finalUrl, "https://source.fixture.example.com/html", ); const searchRequest = (await status()).searches[0]; assert.equal(searchRequest.model, "openai/gpt-5.4-mini"); assert.equal(searchRequest.input.input, "normal"); assert.deepEqual(searchRequest.input.tools, [ { type: "web_search", search_context_size: "low" }, ]); assert.deepEqual(searchRequest.input.tool_choice, { type: "web_search" }); assert.equal(searchRequest.input.max_tool_calls, 1); assert.equal(searchRequest.input.max_output_tokens, 12000); assert.equal(searchRequest.input.store, false); assert.deepEqual(searchRequest.input.include, [ "web_search_call.action.sources", ]); assert.deepEqual(searchRequest.options, { gateway: { id: "default", skipCache: true, collectLog: false }, returnRawResponse: true, signal: true, }); const preAborted = await fetch( `${origin}/__fixture/search-pre-abort`, ).then((r) => r.json()); assert.equal(preAborted.code, "cancelled"); assert.equal( (await status()).searches.length, 1, "pre-abort makes no billed request", ); for (const [query, code] of [ ["http-402", "search_billing"], ["http-401", "search_authentication"], ["http-403", "search_authentication"], ["http-429", "search_rate_limited"], ["http-503", "search_unavailable"], ["throws", "search_unavailable"], ["refusal", "search_unavailable"], ...[ "unexpected", "incomplete", "failed-search", "searching-only", "mixed-empty", "missing-sources", "malformed", "oversized", "unsafe-only", ].map((query) => [query, "invalid_search_response"]), ]) { const before = (await status()).searches.length; const failure = await call("web_search", { query, limit: 5 }); assert.equal(failure.code, code, query); assert.doesNotMatch(JSON.stringify(failure), /private-error-marker/); if (query.startsWith("http-")) assert.equal(failure.status, Number(query.slice(5))); assert.equal( (await status()).searches.length, before + 1, "no automatic billed retry", ); } assert.deepEqual(await call("web_search", { query: "empty", limit: 5 }), { ok: true, sources: [], }); const mixed = await call("web_search", { query: "mixed-search", limit: 5, }); assert.equal(mixed.ok, true, JSON.stringify(mixed)); assert.deepEqual( mixed.sources.map((source) => source.finalUrl), ["https://source.fixture.example.com/html"], "unfinished search metadata must not become source evidence", ); assert.equal( mixed.generatedSummary.text, "Generated search summary, not page text.", ); for (const limit of [1, 3, 5]) { const many = await call("web_search", { query: "many", limit }); assert.equal(many.sources.length, limit); } const longSearch = await call("web_search", { query: "long", limit: 5 }); assert.equal(longSearch.generatedSummary.truncated, true); assert.equal(longSearch.generatedSummary.text.length, 6000); assert.equal(longSearch.sources[0].title.length, 240); assert.ok(!/[\uD800-\uDBFF]$/.test(longSearch.generatedSummary.text)); const cancelId = crypto.randomUUID(); const cancelled = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "read_url", input: { url: "https://source.fixture.example.com/slow" }, id: cancelId, }), ); cancelled.done.catch(() => {}); await waitFor( async () => (await status()).calls.filter((url) => url.endsWith("/slow")) .length === 2, ); connection.transport.cancelActiveServerTurn(); await assert.rejects(cancelled.done, { name: "AbortError" }); await waitFor(async () => (await status()).bodyAborts === 2); const activities = (await connection.client.call("listToolActivities")) .activities; assert.ok(activities.some((a) => a.status === "failed")); assert.ok(activities.some((a) => a.status === "succeeded")); assert.ok(activities.some((a) => a.status === "cancelled")); assert.ok(activities.every((a) => a.kind === "web")); assert.doesNotMatch( JSON.stringify(activities), /fixture.example|snippet|Fish|secret|127\.0/, ); const searchCancelId = crypto.randomUUID(); const searchCancel = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "web_search", input: { query: "slow", limit: 5 }, id: searchCancelId, }), ); searchCancel.done.catch(() => {}); await waitFor(async () => (await status()).searches.some((s) => s.input.input === "slow"), ); const inspect = () => fetch(`${origin}/__fixture/inspect`).then((r) => r.json()); assert.equal( (await inspect()).browserLeases.length, 0, "search needs no browser lease", ); connection.transport.cancelActiveServerTurn(); await assert.rejects(searchCancel.done, { name: "AbortError" }); await waitFor(async () => (await status()).searchAborts === 1); const bodyCancel = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "web_search", input: { query: "slow-body", limit: 5 }, id: crypto.randomUUID(), }), ); bodyCancel.done.catch(() => {}); await waitFor(async () => (await status()).searches.some((s) => s.input.input === "slow-body"), ); connection.transport.cancelActiveServerTurn(); await assert.rejects(bodyCancel.done, { name: "AbortError" }); await waitFor(async () => (await status()).searchBodyCancels === 1); const timedOut = await call("web_search", { query: "timeout", limit: 5 }); assert.ok(timedOut, JSON.stringify((await history(first.id)).slice(-2))); assert.equal(timedOut.code, "timeout"); assert.equal( (await status()).searchAborts, 2, "deadline aborts the binding request", ); const savedActivities = ( await connection.client.call("listToolActivities") ).activities; assert.equal( savedActivities.find((a) => a.toolCallId === searchCancelId).status, "cancelled", ); const researchId = crypto.randomUUID(); const modelCallsBeforeResearch = ( await connection.client.call("fixtureModelCalls") ).length; const research = await call( "research", { code: `async () => { const search = await codemode.web_search({ query: "normal", limit: 3 }); if (!search.ok) return search; return Promise.all(search.sources.map(source => codemode.read_url({ url: source.finalUrl }))); }`, }, researchId, ); assert.ok(research.result.length > 0, JSON.stringify(research)); assert.ok( research.result.every( (page) => page.ok && page.sources[0].sourceKind === "page", ), ); assert.equal( (await connection.client.call("fixtureModelCalls")).length, modelCallsBeforeResearch + 2, "one inference plans search and reads; the second answers from their combined evidence", ); assert.match(research.result[0].sources[0].content, /Research title/); const nested = ( await connection.client.call("listToolActivities") ).activities.filter((activity) => activity.toolCallId.startsWith(`${researchId}:research:`), ); assert.equal(nested.length, research.result.length + 1); assert.ok(nested.every((activity) => activity.status === "succeeded")); assert.equal( (await history(first.id)) .flatMap((message) => message.parts) .filter((part) => part.toolCallId === researchId).length, 1, "search and dependent reads use one model-facing tool call", ); const validation = await call("research", { code: `async () => { try { await codemode.web_search({ query: "" }); return false; } catch { return true; } }`, }); assert.equal( validation.result, true, "nested calls validate their schemas", ); const isolation = await call("research", { code: `async () => { try { await fetch("https://example.com"); return false; } catch { try { await codemode.shell({ command: "true" }); return false; } catch { return true; } } }`, }); assert.equal( isolation.result, true, "sandbox cannot bypass tools with outbound fetch", ); const bounded = await call("research", { code: `async () => { let completed = 0; try { for (let i = 0; i < 13; i++) { await codemode.read_url({ url: "https://source.fixture.example.com/plain" }); completed++; } } catch { return completed; } }`, }); assert.equal(bounded.result, 12); const beforeBatchCancel = await status(); const batchCancelId = crypto.randomUUID(); const batchCancel = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "research", id: batchCancelId, input: { code: `async () => Promise.all(Array.from({ length: 5 }, () => codemode.read_url({ url: "https://source.fixture.example.com/slow" })))`, }, }), ); batchCancel.done.catch(() => {}); await waitFor( async () => (await status()).calls.length === beforeBatchCancel.calls.length + 3, ); connection.transport.cancelActiveServerTurn(); await assert.rejects(batchCancel.done, { name: "AbortError" }); await waitFor( async () => (await status()).bodyAborts === beforeBatchCancel.bodyAborts + 3, ); assert.equal( (await status()).calls.length, beforeBatchCancel.calls.length + 3, "queued calls cannot start after cancellation", ); const batchActivities = ( await connection.client.call("listToolActivities") ).activities.filter((activity) => activity.toolCallId.startsWith(batchCancelId), ); assert.equal(batchActivities.length, 4); assert.ok( batchActivities.every((activity) => activity.status === "cancelled"), ); const researchActivities = ( await connection.client.call("listToolActivities") ).activities; const saved = await history(first.id); connection.client.close(); connection = await connect(first.id); assert.deepEqual(await history(first.id), saved); for (const client of clients) client.close(); await worker.stop(); worker = await start(); owner = await connect(); connection = await connect(first.id); assert.deepEqual(await history(first.id), saved); assert.deepEqual( (await connection.client.call("listToolActivities")).activities, researchActivities, ); const after = await call("read_url", { url: "https://source.fixture.example.com/plain", }); assert.equal(after.sources[0].content, "Plain evidence"); } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);