import * 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 }); } }, );