Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175import * as Effect from "effect/Effect";import { requestedWebReader, RESEARCH_INSTRUCTIONS, WEB_INSTRUCTIONS,} from "../shared/web.ts";import { SHELL_INSTRUCTIONS } from "../shared/shell.ts";import assert from "node:assert/strict";import { createHmac } from "node:crypto";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 { DEFAULT_INSTRUCTIONS, MAX_INSTRUCTIONS_LENGTH,} from "../shared/instructions.ts";import { Secret } from "../configuration/secrets.ts";import { createOwnerSession } from "../worker/session.ts";import { customerBindings, installation } from "./fixtures/config.mjs";
const root = "/agents/personal-agent/personal";const pathFor = (id) => `${root}/sub/conversation/${id}`;const textOf = (history) => history .flatMap((m) => m.parts.filter((p) => p.type === "text").map((p) => p.text)) .join("\n");
test("explicit URL reads and pseudo-call retries require a real reader tool", () => { const message = (role, text) => ({ role, content: [{ type: "text", text }], }); assert.equal( requestedWebReader([ message( "user", "Visit https://news.ycombinator.com/item?id=49580164 and return the response in markdown.", ), ]), "read_url", ); assert.equal( requestedWebReader([ message( "user", "https://developers.cloudflare.com/fundamentals/oauth/create-an-oauth-client/ what's this telling me?", ), ]), "read_url", ); assert.equal( requestedWebReader([ message("user", "Render https://example.com in a browser."), ]), "browser_read", ); assert.equal( requestedWebReader([ message("user", "Visit https://example.com"), message( "assistant", 'I will try that. [read_url(url="https://example.com")]', ), message("user", "Doesn't seem to have worked"), ]), "read_url", ); assert.equal( requestedWebReader([message("user", "Visit https://example.com")], true), undefined, "tool-result continuations must be allowed to answer", ); assert.equal( requestedWebReader([ message("user", "Format https://example.com as a Markdown link"), ]), undefined, );});async function waitFor(predicate) { const deadline = Date.now() + 20_000; while (!(await predicate())) { if (Date.now() > deadline) throw new Error("Timed out waiting for Think"); await new Promise((resolve) => setTimeout(resolve, 30)); }}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;}
function deniedSocket(url, headers, status) { return new Promise((resolve, reject) => { const socket = new WebSocket(url, { headers, handshakeTimeout: 5_000 }); socket.on("open", () => { socket.close(); reject(new Error("Unauthorized socket accepted")); }); socket.on("error", () => {}); socket.on("unexpected-response", (_request, response) => { response.resume(); socket.terminate(); try { assert.equal(response.statusCode, status); resolve(); } catch (error) { reject(error); } }); });}
test( "persistent conversations support CRUD, isolated native streaming, deletion, reconnect and recovery", { timeout: 180_000 }, async (t) => { const logs = []; for (const method of ["log", "info", "warn", "error"]) { const original = console[method]; t.mock.method(console, method, (...args) => { logs.push(args.map(String).join(" ")); original(...args); }); } const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const secret = new Secret(customerBindings.FLAREBOT_SESSION_SECRET); const cookie = ( await Effect.runPromise( createOwnerSession(secret, { ...installation, runtimeOrigin: origin, }), ) ).split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; const persistence = await mkdtemp(join(tmpdir(), "flarebot-think-")); 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-think-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/think-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/think-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 = []; const ownerFrames = []; let worker; try { worker = await start(); class AuthenticatedSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } async function connectOwner() { const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", name: "personal", WebSocket: AuthenticatedSocket, }); clients.push(client); client.addEventListener("message", (event) => ownerFrames.push(event.data), ); await client.ready; return client; } let owner = await connectOwner(); const inspect = async () => (await fetch(`${origin}/__fixture/inspect`)).json(); const create = async (name) => owner.call("createConversation", name === undefined ? [] : [name]); assert.deepEqual(await owner.call("listConversations"), []); for (const name of [ null, 42, {}, "", " ", "a".repeat(121), "line\nbreak", "bad\u0000name", ]) { await assert.rejects(create(name), /Conversation name/); } assert.deepEqual(await owner.call("listConversations"), []); const history = async (id) => { const response = await fetch(origin + pathFor(id) + "/get-messages", { headers, }); assert.equal(response.status, 200, await response.clone().text()); assert.equal(response.headers.get("cache-control"), "no-store"); return response.json(); }; const firstMetadata = await create(" Research "); const firstId = firstMetadata.id; const secondMetadata = await create(); const secondId = secondMetadata.id; assert.equal(firstMetadata.name, "Research"); assert.equal(secondMetadata.name, "New conversation"); assert.deepEqual(Object.keys(firstMetadata).sort(), [ "createdAt", "id", "name", "updatedAt", ]); assert.ok(Number.isFinite(Date.parse(firstMetadata.createdAt))); for (const id of [ null, 1, {}, "", "../../personal", firstId.toUpperCase(), `${firstId}/get-messages`, ]) { await assert.rejects( owner.call("renameConversation", [id, "Renamed"]), /Invalid conversation ID/, ); await assert.rejects( owner.call("deleteConversation", [id]), /Invalid conversation ID/, ); } await assert.rejects( owner.call("renameConversation", [crypto.randomUUID(), "Renamed"]), /Conversation not found/, ); await assert.rejects( owner.call("renameConversation", [firstId, ""]), /Conversation name/, ); const renamed = await owner.call("renameConversation", [ firstId, "Renamed research", ]); assert.equal(renamed.name, "Renamed research"); assert.equal(renamed.createdAt, firstMetadata.createdAt); assert.ok(renamed.updatedAt >= firstMetadata.updatedAt); await assert.rejects( owner.call("subAgent", ["Conversation", crypto.randomUUID()]), /not callable/, ); await assert.rejects( owner.call("prepareConversation", [firstId]), /not callable/, ); for (const id of [firstId, secondId]) { assert.deepEqual(await history(id), []); for (const suffix of ["", "/get-messages"]) { assert.equal( (await fetch(origin + pathFor(id) + suffix)).status, 401, ); assert.equal( ( await fetch(origin + pathFor(id) + suffix, { headers: { ...headers, Origin: "https://attacker.example.com" }, }) ).status, 403, ); } await deniedSocket( origin.replace("http", "ws") + pathFor(id), { Origin: origin }, 401, ); } for (const path of [ pathFor(crypto.randomUUID()), `${pathFor(firstId)}/sub/conversation/${secondId}`, `/agents/conversation/${firstId}`, `${pathFor(firstId)}/arbitrary`, ]) { assert.equal((await fetch(origin + path, { headers })).status, 404); await deniedSocket(origin.replace("http", "ws") + path, headers, 404); } async function connect(id) { const frames = []; const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", basePath: pathFor(id).slice(1), WebSocket: AuthenticatedSocket, minReconnectionDelay: 50, maxReconnectionDelay: 250, }); const transport = new WebSocketChatTransport({ agent: client }); client.addEventListener("message", (event) => { const frame = JSON.parse(event.data); frames.push(frame); 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(); }); clients.push(client); await Promise.race([ client.ready, new Promise((_, reject) => { const timer = setTimeout( () => reject( new Error( `Facet readiness timed out: ${JSON.stringify(frames)}`, ), ), 10_000, ); timer.unref(); }), ]); return { client, transport, frames }; } const first = await connect(firstId); const observer = await connect(firstId); const second = await connect(secondId); async function send(connection, id, text, body) { const messages = [ ...(await history(id)), { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ]; const chunks = []; const stream = await connection.transport.sendMessages({ chatId: id, messages, body, trigger: "submit-message", abortSignal: new AbortController().signal, }); const done = (async () => { for await (const chunk of stream) chunks.push(chunk); })(); return { chunks, done }; } const nameOf = async (id) => (await owner.call("listConversations")).find((item) => item.id === id) ?.name; const titleCalls = async () => (await inspect()).titleCalls; await assert.rejects( owner.call("nameConversationFromFirstMessage", [secondId, "Injected"]), /not callable/, ); for (const [message, expected] of [ ["title:Configure Flarebot locally", "Configure Flarebot locally"], ["title:failure", "New conversation"], ["title:persistence-failure", "New conversation"], ["title:long-output", "New conversation"], ["title:unsafe-output", "New conversation"], ["title:bad-output", "New conversation"], ["title:opaque-output", "New conversation"], ["title:prefixed-output", "New conversation"], ["title:", "New conversation"], ["title:" + "a".repeat(121), "New conversation"], ["title:line\nbreak", "New conversation"], ['title:read_url(url="private")', "New conversation"], [ "title:Visit https://user:password@example.com/?token=private", "New conversation", ], ["title:API key private-credential", "New conversation"], ]) { const id = (await create()).id; assert.equal(await nameOf(id), "New conversation"); const connection = await connect(id); const before = (await titleCalls()).length; await ( await send(connection, id, message) ).done; await waitFor(async () => (await titleCalls()).length === before + 1); if (expected !== "New conversation") await waitFor(async () => (await nameOf(id)) === expected); const originalMessages = await history(id); const duplicate = await connection.transport.sendMessages({ chatId: id, messages: originalMessages.filter((item) => item.role === "user"), trigger: "submit-message", abortSignal: new AbortController().signal, }); for await (const _chunk of duplicate) { /* drain native duplicate replay */ } await ( await send(connection, id, "second") ).done; assert.equal(await nameOf(id), expected); assert.equal((await titleCalls()).length, before + 1); connection.client.close(); await owner.call("deleteConversation", [id]); } assert.ok( ownerFrames.some( (frame) => JSON.parse(frame).type === "conversations-changed", ), ); for (const customName of ["Custom name", "New conversation"]) { const id = (await create(customName)).id; const connection = await connect(id); const before = (await titleCalls()).length; await ( await send(connection, id, "title:Configure Flarebot locally") ).done; assert.equal(await nameOf(id), customName); assert.equal((await titleCalls()).length, before); connection.client.close(); await owner.call("deleteConversation", [id]); } const raceId = (await create()).id; const race = await connect(raceId); const beforeRace = (await titleCalls()).length; const raceTurn = await send(race, raceId, "title:slow"); await raceTurn.done; await waitFor(async () => (await titleCalls()).length === beforeRace + 1); assert.equal( await nameOf(raceId), "New conversation", "chat completes before title inference", ); await owner.call("renameConversation", [raceId, "Manual name"]); await owner.call("renameConversation", [raceId, "New conversation"]); await new Promise((resolve) => setTimeout(resolve, 1400)); await ( await send(race, raceId, "second") ).done; assert.equal( await nameOf(raceId), "New conversation", "manual rename back to default still wins", ); assert.equal((await titleCalls()).length, beforeRace + 1); race.client.close(); await owner.call("deleteConversation", [raceId]); await ( await send(second, secondId, "title:Understand Cloudflare OAuth") ).done; await waitFor( async () => (await nameOf(secondId)) === "Understand Cloudflare OAuth", ); const defaults = await owner.call("getInstructionSettings"); assert.deepEqual(defaults, { instructions: DEFAULT_INSTRUCTIONS, defaultInstructions: DEFAULT_INSTRUCTIONS, customized: false, updatedAt: null, }); await assert.rejects(owner.call("readInstructions"), /not callable/); const assertInstructions = async (connection, id, expected) => { const turn = await send(connection, id, "instructions"); await turn.done; const reply = textOf([(await history(id)).at(-1)]); const system = JSON.parse( reply.slice("Reply ".length, -" complete".length), ); assert.equal(system.length, 1); const prefix = expected + WEB_INSTRUCTIONS + RESEARCH_INSTRUCTIONS + SHELL_INSTRUCTIONS; assert.equal( system[0].slice(0, prefix.length), prefix, "model sees current effective instructions and existing tool guidance, no stale frozen prompt", ); assert.match( system[0].slice(prefix.length), /^\n\nScheduling:\nUse createSchedule/, ); assert.match( system[0], /\nCurrent UTC time: \d{4}-\d\d-\d\dT\d\d:\d\d:\d\d\.\d{3}Z\n\nAvailable skills\./, ); assert.equal(system[0].split("Available skills.").length - 1, 1); assert.match(system[0], /- summarize-report: /); assert.match(system[0], /- prepare-decision: /); }; await assertInstructions(first, firstId, DEFAULT_INSTRUCTIONS); const customInstructions = "Use concise Polish answers.\nKeep my personal preferences private.\tAsk when uncertain."; const customSettings = await owner.call("updateInstructions", [ customInstructions, ]); assert.equal(customSettings.instructions, customInstructions); assert.equal(customSettings.customized, true); assert.ok(Number.isFinite(Date.parse(customSettings.updatedAt))); for (const value of [ null, false, 42, {}, "", " \n\t", "x".repeat(MAX_INSTRUCTIONS_LENGTH + 1), "bad\u0000text", "bad\u007ftext", ]) { await assert.rejects( owner.call("updateInstructions", [value]), /Instructions must/, ); assert.deepEqual( await owner.call("getInstructionSettings"), customSettings, ); } await Promise.all([ assertInstructions(first, firstId, customInstructions), assertInstructions(second, secondId, customInstructions), ]); const instructionId = (await create()).id; const instructionConnection = await connect(instructionId); await assertInstructions( instructionConnection, instructionId, customInstructions, ); instructionConnection.client.close(); await owner.call("deleteConversation", [instructionId]); const resetSettings = await owner.call("resetInstructions"); assert.equal(resetSettings.customized, false); assert.equal(resetSettings.instructions, DEFAULT_INSTRUCTIONS); await assertInstructions(first, firstId, DEFAULT_INSTRUCTIONS); await owner.call("updateInstructions", [customInstructions]); assert.ok( !JSON.stringify(await owner.call("getStatus")).includes( customInstructions, ), ); const defaultSettings = await owner.call("getModelSettings"); assert.deepEqual(defaultSettings, { configuration: { provider: "workers-ai", model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast", }, providers: { "workers-ai": true, anthropic: true, openrouter: false, }, providerSetupUrl: null, credentials: { anthropic: "missing", openrouter: "missing" }, }); const catalog = await owner.call("getModelCatalog"); const external = { provider: "anthropic", model: catalog.anthropic[0] }; const apiKey = "sk-ant-fixture-original-key-never-return"; const replacementKey = "sk-ant-fixture-replacement-key-never-return"; await assert.rejects( owner.call("readModelConfiguration"), /not callable/, ); await owner.call("updateModelSettings", [external]); const missingKeyTurn = await send(first, firstId, "configuration"); await assert.rejects(missingKeyTurn.done, /Add an Anthropic API key/); const configured = await owner.call("setProviderKey", [ "anthropic", apiKey, ]); assert.deepEqual(configured, { ...defaultSettings, configuration: external, credentials: { anthropic: "configured", openrouter: "missing" }, }); for (const invalid of [ null, {}, { ...external, apiKey }, { ...external, model: "arbitrary" }, { provider: "workers-ai", model: external.model }, ]) { await assert.rejects( owner.call("updateModelSettings", [invalid]), /supported provider and model/, ); assert.deepEqual(await owner.call("getModelSettings"), configured); } await assert.rejects( owner.call("setProviderKey", ["anthropic", "bad\nkey"]), /API key/, ); assert.deepEqual(await owner.call("getModelSettings"), configured); const assertConfiguration = async (connection, id, configuration) => { const turn = await send(connection, id, "configuration"); await turn.done; const calls = await connection.client.call("fixtureModelCalls"); assert.equal(calls.at(-1).maxOutputTokens, 16384); assert.equal( textOf([(await history(id)).at(-1)]), `Reply ${configuration.provider}/${configuration.model} complete`, ); }; // The composer sends configuration in the native request, not a separate // settings mutation. Each step sees that same model and effort. const scopedId = (await create("Scoped model and effort")).id; const scoped = await connect(scopedId); const selected = { ...external, effort: "low" }; assert.equal( (await scoped.client.call("getConversationModelSettings")).override, null, ); await ( await send(scoped, scopedId, "tool", { modelConfiguration: selected }) ).done; const calls = await scoped.client.call("fixtureModelCalls"); assert.ok(calls.length >= 2, "tool follow-up must invoke the model"); for (const call of calls) { assert.deepEqual(call.configuration, selected); assert.equal(call.maxOutputTokens, 16384); assert.deepEqual(call.options.anthropic, { effort: "low", thinking: { type: "adaptive" }, }); } assert.deepEqual( (await scoped.client.call("getConversationModelSettings")).override, selected, ); assert.equal( (await second.client.call("getConversationModelSettings")).override, null, ); assert.deepEqual( (await owner.call("getModelSettings")).configuration, external, ); const invalidEffort = await send(scoped, scopedId, "configuration", { modelConfiguration: { ...selected, effort: "off" }, }); await assert.rejects(invalidEffort.done, /supported effort/); assert.deepEqual( (await scoped.client.call("getConversationModelSettings")).override, selected, ); scoped.client.close(); const reopened = await connect(scopedId); assert.deepEqual( (await reopened.client.call("getConversationModelSettings")).override, selected, ); await assertConfiguration(reopened, scopedId, selected); // Reset applies an explicit snapshot to this turn while clearing the // conversation override for future default-following turns. await ( await send(reopened, scopedId, "configuration", { modelConfiguration: defaultSettings.configuration, useDefaultModel: true, }) ).done; assert.equal( (await reopened.client.call("getConversationModelSettings")).override, null, ); reopened.client.close(); await owner.call("deleteConversation", [scopedId]); const openRouterConfiguration = { provider: "openrouter", model: catalog["openrouter"][0], }; const openRouterKey = "openrouter-fixture-private-key-never-return"; await assert.rejects( owner.call("updateModelSettings", [openRouterConfiguration]), /Enable OpenRouter in Settings/, ); const enabled = await fetch(`${origin}/__fixture/enable-openrouter`, { method: "POST", }); assert.equal(enabled.status, 204); await enabled.arrayBuffer(); await owner.call("updateModelSettings", [openRouterConfiguration]); const missingOpenRouterKey = await send(first, firstId, "configuration"); await assert.rejects( missingOpenRouterKey.done, /Add an OpenRouter API key/, ); const openRouterSettings = await owner.call("setProviderKey", [ "openrouter", openRouterKey, ]); assert.deepEqual(openRouterSettings.credentials, { anthropic: "configured", openrouter: "configured", }); await assertConfiguration(first, firstId, openRouterConfiguration); assert.ok( !JSON.stringify((await inspect()).credentialsStorage).includes( openRouterKey, ), ); await owner.call("setProviderKey", ["openrouter", null]); assert.deepEqual((await owner.call("getModelSettings")).credentials, { anthropic: "configured", openrouter: "missing", }); // A removal only affects its own provider. Keep OpenRouter configured to verify // independent ciphertext survives the later legacy Anthropic migration. await owner.call("setProviderKey", ["openrouter", openRouterKey]); await owner.call("updateModelSettings", [external]); await Promise.all([ assertConfiguration(first, firstId, external), assertConfiguration(second, secondId, external), ]); const newId = (await create("New settings conversation")).id; const newlyCreated = await connect(newId); await assertConfiguration(newlyCreated, newId, external); newlyCreated.client.close(); await owner.call("deleteConversation", [newId]); const providerFailure = await send(first, firstId, "credential-error"); await assert.rejects( providerFailure.done, /Anthropic rejected authentication/, ); await owner.call("setProviderKey", ["anthropic", replacementKey]); const secondWorkersModel = { provider: "workers-ai", model: catalog["workers-ai"][1], }; await owner.call("updateModelSettings", [secondWorkersModel]); await assertConfiguration(first, firstId, secondWorkersModel); await owner.call("updateModelSettings", [defaultSettings.configuration]); const firstTurn = await send(first, firstId, "slow-first"); await waitFor(() => firstTurn.chunks.some((c) => c.type === "text-delta"), ); assert.ok( !firstTurn.chunks.some((c) => c.type === "finish"), "stream yields before turn finishes", ); const secondTurn = await send(second, secondId, "second"); await secondTurn.done; await firstTurn.done; await waitFor(async () => /Reply slow-first complete/.test(textOf(await history(firstId))), ); await waitFor(async () => /Reply second complete/.test(textOf(await history(secondId))), ); assert.ok(!textOf(await history(secondId)).includes("slow-first")); await waitFor(() => JSON.stringify(observer.frames).includes("slow-first"), ); const toolTurn = await send(first, firstId, "tool"); await toolTurn.done; const toolHistory = await history(firstId); assert.ok( toolHistory .flatMap((m) => m.parts) .some( (p) => p.type === "tool-fixtureEcho" && p.state === "output-available" && p.output.echoed === "native tool output", ), ); assert.match(textOf(toolHistory), /Reply tool complete/);
const cancelTurn = await send(first, firstId, "slow-cancel"); await waitFor(() => cancelTurn.chunks.some((c) => c.type === "text-delta"), ); assert.equal(first.transport.cancelActiveServerTurn(), true); await assert.rejects(cancelTurn.done, { name: "AbortError" }); await waitFor( async () => (await history(firstId)).at(-1)?.role === "assistant", ); const cancelled = (await history(firstId)).at(-1); assert.equal(textOf([cancelled]), "Reply "); const errorTurn = await send(first, firstId, "error"); await assert.rejects(errorTurn.done, /Workers AI request failed/); const afterError = await send(first, firstId, "after-error"); await afterError.done; assert.match( textOf(await history(firstId)), /Reply after-error complete/, );
// Detaching a tab preserves the durable server turn; another native transport // reattaches to the buffered stream without submitting the user message twice. const detachedTurn = await send(first, firstId, "slow-detach"); detachedTurn.done.catch(() => {}); await waitFor(() => detachedTurn.chunks.some((c) => c.type === "text-delta"), ); first.client.close(); const reconnect = await connect(firstId); const resumed = await reconnect.transport.reconnectToStream({ chatId: firstId, }); assert.ok(resumed); const resumedChunks = []; for await (const chunk of resumed) resumedChunks.push(chunk); assert.ok(resumedChunks.some((c) => c.type === "text-delta")); assert.equal( (await history(firstId)).filter( (m) => m.role === "user" && textOf([m]) === "slow-detach", ).length, 1, ); const stableSecondHistory = await history(secondId);
// Delete a live stream while observers, history requests and stale frames // still address it. Native parent bridging must not recreate the facet. const deletedId = (await create("Delete while running")).id; const liveDeleted = await connect(deletedId); const deletedObserver = await connect(deletedId); const deletedTurn = await send(liveDeleted, deletedId, "slow-delete"); deletedTurn.done.catch(() => {}); await waitFor(() => deletedTurn.chunks.some((c) => c.type === "text-delta"), ); const closedCodes = []; for (const connection of [liveDeleted, deletedObserver]) connection.client.addEventListener("close", (event) => closedCodes.push(event.code), ); const staleFrames = setInterval( () => liveDeleted.client.send( JSON.stringify({ type: MessageType.CF_AGENT_STREAM_RESUME }), ), 5, ); try { await Promise.all([ owner.call("deleteConversation", [deletedId]), owner.call("deleteConversation", [deletedId]), ...Array.from({ length: 5 }, async () => { const response = await fetch( origin + pathFor(deletedId) + "/get-messages", { headers }, ); assert.ok([200, 404, 503].includes(response.status)); }), ]); await waitFor(() => closedCodes.length >= 2); } finally { clearInterval(staleFrames); liveDeleted.client.close(); deletedObserver.client.close(); } assert.deepEqual(closedCodes.slice(0, 2), [4004, 4004]); await waitFor(async () => !(await inspect()).facets.includes(deletedId)); await assert.rejects( owner.call("renameConversation", [deletedId, "Resurrect"]), /Conversation not found/, ); for (const id of [deletedId, crypto.randomUUID()]) { await owner.call("deleteConversation", [id]); assert.equal( (await fetch(origin + pathFor(id) + "/get-messages", { headers })) .status, 404, ); await deniedSocket( origin.replace("http", "ws") + pathFor(id), headers, 404, ); } const stableMetadata = await owner.call("listConversations"); assert.deepEqual( stableMetadata.map((item) => item.id).sort(), [firstId, secondId].sort(), ); assert.deepEqual( stableMetadata.find((item) => item.id === firstId), renamed, ); assert.deepEqual((await inspect()).facets, [firstId, secondId].sort());
// Failed teardown stays inaccessible and is retried on the next full wake. const failedDeleteId = (await create("Failed teardown")).id; await fetch(`${origin}/__fixture/fail-delete`, { method: "POST" }); await assert.rejects( owner.call("deleteConversation", [failedDeleteId]), /Conversation deletion unavailable/, ); const interruptedCreateId = await ( await fetch(`${origin}/__fixture/interrupted-create`, { method: "POST", }) ).json(); for (const id of [failedDeleteId, interruptedCreateId]) { assert.equal( (await fetch(origin + pathFor(id) + "/get-messages", { headers })) .status, 404, ); await deniedSocket( origin.replace("http", "ws") + pathFor(id), headers, 404, ); } assert.deepEqual(await owner.call("listConversations"), stableMetadata); await owner.call("updateModelSettings", [external]); const beforeRestartTitleCalls = await titleCalls(); const beforeRestartSettings = await owner.call("getModelSettings"); const beforeRestartStorage = (await inspect()).settingsStorage; const beforeRestartCredentials = (await inspect()).credentialsStorage; const beforeRestartEnabledProviders = (await inspect()) .enabledProvidersStorage; assert.ok(!JSON.stringify(beforeRestartStorage).includes(replacementKey)); const recoveryTurn = await send(reconnect, firstId, "recover"); recoveryTurn.done.catch(() => {}); await waitFor(() => recoveryTurn.chunks.some((c) => c.type === "text-delta"), ); // Allow Think's native buffered stream checkpoint to reach SQLite. await new Promise((resolve) => setTimeout(resolve, 400)); // Emulate the previous schema at rest after this turn has read its key. // Migration runs on the actual Worker restart below. const retiredResponse = await fetch( `${origin}/__fixture/retired-opencode`, { method: "POST" }, ); assert.equal(retiredResponse.status, 204); await retiredResponse.arrayBuffer(); const legacyResponse = await fetch(`${origin}/__fixture/legacy-key`, { method: "POST", }); assert.equal(legacyResponse.status, 204); await legacyResponse.arrayBuffer(); const stagedRetired = await inspect(); assert.equal( JSON.parse(stagedRetired.settingsStorage[0].configuration).provider, "opencode-go", ); assert.ok( stagedRetired.credentialsStorage.some( (row) => row.provider === "opencode-go", ), ); assert.ok( stagedRetired.enabledProvidersStorage.some( (row) => row.provider === "opencode-go", ), ); for (const client of clients) client.close(); await worker.stop(); worker = undefined; worker = await start(); owner = await connectOwner(); const migratedSettings = await owner.call("getModelSettings"); assert.deepEqual(migratedSettings, { ...beforeRestartSettings, configuration: defaultSettings.configuration, }); assert.equal( (await owner.call("getInstructionSettings")).instructions, customInstructions, ); assert.equal( (await owner.call("getInstructionSettings")).customized, true, ); const migrated = await inspect(); assert.deepEqual(migrated.settingsStorage, [ { ...beforeRestartStorage[0], configuration: JSON.stringify(defaultSettings.configuration), }, ]); assert.deepEqual(migrated.credentialsStorage, beforeRestartCredentials); assert.deepEqual( migrated.enabledProvidersStorage, beforeRestartEnabledProviders, ); assert.ok( !migrated.credentialsStorage.some( (row) => row.provider === "opencode-go", ), ); assert.ok( !migrated.enabledProvidersStorage.some( (row) => row.provider === "opencode-go", ), ); assert.equal(migratedSettings.providers.openrouter, true); assert.equal(migratedSettings.credentials.openrouter, "configured"); assert.equal(migratedSettings.credentials.anthropic, "configured");
// A second full wake must leave the migrated rows unchanged. owner.close(); await worker.stop(); worker = undefined; worker = await start(); owner = await connectOwner(); assert.deepEqual(await owner.call("getModelSettings"), migratedSettings); const migratedAgain = await inspect(); assert.deepEqual(migratedAgain.settingsStorage, migrated.settingsStorage); assert.deepEqual( migratedAgain.credentialsStorage, migrated.credentialsStorage, ); assert.deepEqual( migratedAgain.enabledProvidersStorage, migrated.enabledProvidersStorage, ); await owner.call("updateModelSettings", [external]); assert.deepEqual(await owner.call("listConversations"), stableMetadata); const restarted = await inspect(); assert.deepEqual(restarted.facets, [firstId, secondId].sort()); assert.deepEqual( restarted.metadata, [firstId, secondId].sort().map((id) => ({ id, status: "active" })), ); assert.equal( ( await fetch(origin + pathFor(deletedId) + "/get-messages", { headers, }) ).status, 404, ); await deniedSocket( origin.replace("http", "ws") + pathFor(deletedId), headers, 404, ); const recovered = await connect(firstId); await waitFor(async () => /Recovered recover complete/.test(textOf(await history(firstId))), ); assert.deepEqual(await history(secondId), stableSecondHistory); assert.equal( (await history(firstId)).filter( (m) => m.role === "user" && textOf([m]) === "recover", ).length, 1, ); assert.equal( textOf([(await history(firstId)).find((m) => m.id === cancelled.id)]), "Reply ", "cancelled turn is not resumed after restart", ); assert.ok(!JSON.stringify(recovered.frames).includes(secret.reveal())); await assertInstructions(recovered, firstId, customInstructions); await assertConfiguration(recovered, firstId, external); const secondReconnected = await connect(secondId); await assertConfiguration(secondReconnected, secondId, external); assert.deepEqual(await titleCalls(), beforeRestartTitleCalls); assert.equal(await nameOf(secondId), "Understand Cloudflare OAuth"); await owner.call("setProviderKey", ["anthropic", null]); assert.equal( (await owner.call("getModelSettings")).credentials.anthropic, "missing", ); assert.equal((await inspect()).settingsStorage[0].anthropic_key, null); assert.ok( !(await inspect()).credentialsStorage.some( (row) => row.provider === "anthropic", ), ); const removedKeyTurn = await send(recovered, firstId, "configuration"); await assert.rejects(removedKeyTurn.done, /Add an Anthropic API key/); await owner.call("updateModelSettings", [defaultSettings.configuration]); await assertConfiguration( recovered, firstId, defaultSettings.configuration, ); const stateFrames = ownerFrames .map((frame) => JSON.parse(frame)) .filter((frame) => frame.type === "cf_agent_state"); assert.ok(stateFrames.length > 0); assert.ok( !JSON.stringify(stateFrames).includes(customInstructions), "custom text is never a generic native state broadcast", ); const surfaces = JSON.stringify({ ownerFrames, logs, first: first.frames, second: second.frames, recovered: recovered.frames, firstHistory: await history(firstId), secondHistory: await history(secondId), status: await owner.call("getStatus"), }); assert.ok( logs.some((line) => line.includes("Anthropic rejected authentication")), "native error logs were captured", ); for (const sensitive of [ apiKey, replacementKey, openRouterKey, secret.reveal(), ]) assert.ok( !surfaces.includes(sensitive), "credentials absent from native protocol, snapshots, history and logs", );
const [cookieName, signedCookie] = cookie.split("="); const claims = JSON.parse( Buffer.from(signedCookie.split(".")[0], "base64url").toString(), ); const now = Math.floor(Date.now() / 1000); const shortBody = Buffer.from( JSON.stringify({ ...claims, issuedAt: now, expiresAt: now + 2 }), ).toString("base64url"); const expiryCookie = `${cookieName}=${shortBody}.${createHmac("sha256", secret.reveal()).update(shortBody).digest("base64url")}`; const expiring = new WebSocket( origin.replace("http", "ws") + pathFor(secondId), { headers: { Origin: origin, Cookie: expiryCookie }, closeTimeout: 100, }, ); const closed = await new Promise((resolve, reject) => { const timeout = setTimeout(() => { expiring.terminate(); reject(new Error("Facet session did not expire")); }, 10_000); expiring.on("error", reject); expiring.on("close", (code) => { clearTimeout(timeout); resolve(code); }); }); assert.equal(closed, 4001); } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);