Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485import { cleanupProbe, runCleanupTimeout, registerCleanupProbe,} from "./fixtures/native-cleanup-probe.mjs";import { NativeFixtureLifetime } from "./fixtures/native-fixture-lifetime.mjs";import * as Effect from "effect/Effect";import assert from "node:assert/strict";import { test } from "node:test";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 { 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";import { diagnosticId } from "../worker/diagnostics.ts";
const root = "/agents/personal-agent/personal";const childPath = (id) => `${root}/sub/conversation/${id}`;async function pollUntil(lifetime, read, predicate) { const deadline = Date.now() + 15000; let value; while (true) { lifetime.signal.throwIfAborted(); value = await lifetime.wait("diagnostic polling", read()); if (predicate(value)) return value; if (Date.now() > deadline) throw new Error(`Timed out: ${JSON.stringify(value)}`); await lifetime.sleep(40); }}test( "native operational export survives restart, excludes content, respects scope/clear/deletion and tolerates storage failure", { timeout: 90000 }, async (t) => { const lifetime = new NativeFixtureLifetime(t, "diagnostics"); const wait = (ms) => lifetime.sleep(ms); const until = (read, predicate) => pollUntil(lifetime, read, predicate); const fetch = (url, options = {}) => { lifetime.signal.throwIfAborted(); return lifetime.wait( "HTTP request", globalThis.fetch(url, { ...options, signal: options.signal ? AbortSignal.any([options.signal, lifetime.signal]) : lifetime.signal, }), ); }; const server = lifetime.own( "port reservation", createServer(), (server) => new Promise((resolve) => server.close(() => resolve())), ); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const { port } = server.address(); await new Promise((resolve) => server.close(resolve)); const origin = `http://127.0.0.1:${port}`; const temporary = await lifetime.start( "directory", () => mkdtemp(join(tmpdir(), "flarebot-diagnostics-")), (path) => rm(path, { recursive: true, force: true }), true, ); const config = JSON.parse( await readFile("dist/release/deployment.json", "utf8"), ); const configPath = join(temporary, "wrangler.json"); await writeFile( configPath, JSON.stringify({ ...config, name: "flarebot-diagnostics-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/think-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const cookie = ( await Effect.runPromise( createOwnerSession( new Secret(customerBindings.FLAREBOT_SESSION_SECRET), { ...installation, runtimeOrigin: origin }, ), ) ).split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; let holdIdentity = false; const identityHeld = Promise.withResolvers(); class AuthenticatedSocket extends WebSocket { emit(event, ...args) { if ( holdIdentity && event === "message" && JSON.parse(String(args[0])).type === "cf_agent_identity" ) { identityHeld.resolve(); return true; } return super.emit(event, ...args); } constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } const clients = []; const connect = async (id) => { lifetime.signal.throwIfAborted(); const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", ...(id ? { basePath: childPath(id).slice(1) } : { name: "personal" }), WebSocket: AuthenticatedSocket, }); lifetime.own("Agent client", client, (client) => client.close()); clients.push(client); await lifetime.ready(client, id ? "conversation" : "owner"); return client; }; const start = () => lifetime.start( "Worker", () => 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: temporary, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }), (worker) => worker.stop(), ); let worker; try { worker = await start(); assert.equal((await fetch(origin + root + "/status")).status, 401); let owner = await connect(); if (cleanupProbe === "diagnostics") { const target = await owner.call("createConversation", [ "Cleanup fixture", ]); holdIdentity = true; const pending = connect(target.id); void pending.catch(() => {}); await identityHeld.promise; await runCleanupTimeout(t, lifetime, pending); return; } const mcp = await owner.call("addMcpConnection", [ { name: "private-mcp-name", endpoint: "https://private-mcp.example.com/mcp", authMode: "headers", }, ]); assert.equal(mcp.ok, true); await until( () => owner.call("listMcpConnections"), (rows) => rows[0]?.state === "authenticating", ); const mcpExport = (await owner.call("exportDiagnostics")).mcpConnections; assert.equal(mcpExport.length, 1); assert.equal(mcpExport[0].id, diagnosticId(mcp.value.id)); assert.equal(mcpExport[0].health.lastContactAt, null); assert.equal(mcpExport[0].health.retryAt, null); assert.deepEqual(Object.keys(mcpExport[0]).sort(), [ "enabled", "health", "id", "lastError", "state", ]);
const target = await owner.call("createConversation", [ "private-mcp-name", "private-mcp.example.com", mcp.value.id, "private-conversation-name", ]); const unrelated = await owner.call("createConversation", [ "private-unrelated-name", ]); await owner.call("updateInstructions", ["private-instructions-sentinel"]); await owner.call("addMemory", ["private-memory-sentinel"]); await owner.call("setProviderKey", [ "anthropic", "private-api-key-sentinel-long-enough", ]); let chat = await connect(target.id); let transport = new WebSocketChatTransport({ agent: chat }); chat.addEventListener("message", (event) => { const frame = JSON.parse(event.data); if (frame.type === MessageType.CF_AGENT_STREAM_RESUMING) transport.handleStreamResuming(frame); }); const history = (id) => fetch(origin + childPath(id) + "/get-messages", { headers }).then((r) => r.json(), ); const send = async (text, abort = new AbortController()) => { const stream = await transport.sendMessages({ chatId: target.id, messages: [ ...(await history(target.id)), { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ], trigger: "submit-message", abortSignal: AbortSignal.any([abort.signal, lifetime.signal]), }); const chunks = []; const done = (async () => { for await (const chunk of stream) chunks.push(chunk); })(); return { done, chunks }; }; await ( await send("tool") ).done; let exported = await until( () => owner.call("exportDiagnostics", [target.id]), (value) => value.conversation.events.some( (e) => e.kind === "turn" && e.phase === "finish", ), ); const events = exported.conversation.events; assert.equal(events.filter((e) => e.kind === "model").length, 2); assert.equal( events.filter((e) => e.kind === "model-attempt" && e.phase === "finish") .length, 2, ); assert.ok( events .filter((e) => e.kind === "model") .every( (e) => e.details.inputTokens === 1 && e.details.outputTokens === 1 && e.details.cost === null && typeof e.durationMs === "number" && Number.isFinite(e.durationMs) && e.durationMs >= 0, ), ); assert.ok( events.some((e) => e.kind === "tool" && e.status === "succeeded"), ); assert.ok( events.some((e) => e.kind === "connection" && e.status === "connected"), ); assert.ok( events .filter((e) => e.kind === "model") .every((e) => events.some( (turn) => turn.kind === "turn" && turn.id === e.requestId, ), ), ); assert.equal(exported.scope.conversationId, diagnosticId(target.id)); assert.equal( ( await owner.call("exportDiagnostics", [unrelated.id]) ).conversation.events.filter((e) => e.kind === "model").length, 0, ); assert.equal((await owner.call("exportDiagnostics")).conversation, null); for (const value of ["../private", {}, 3, crypto.randomUUID()]) await assert.rejects(owner.call("exportDiagnostics", [value])); await assert.rejects(chat.call("diagnosticSnapshot"), /not callable/); await ( await send("activity-failure") ).done; await owner.call("updateModelSettings", [ { provider: "anthropic", model: "claude-haiku-4-5-20251001" }, ]); await assert.rejects( (await send("credential-error")).done, /rejected authentication/, ); await owner.call("updateModelSettings", [ { provider: "workers-ai", model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast", }, ]); await ( await send("activity-progress") ).done; exported = await owner.call("exportDiagnostics", [target.id]); assert.ok( exported.conversation.events.some( (e) => e.kind === "tool" && e.status === "failed", ), ); assert.ok( exported.conversation.events.some( (e) => e.kind === "model-attempt" && e.status === "error" && e.details.httpStatus === 401 && e.details.errorCategory === "authentication" && e.details.gatewayErrorCode === 2009 && e.details.cfRay === "0123456789abcdef-WAW" && e.details.failureStage === "request", ), ); assert.ok( exported.conversation.events.length < 100, "1000 progress updates must not become 1000 diagnostic events", ); const raw = await ( await fetch(`${origin}/__fixture/diagnostics?id=${target.id}`) ).json(); for (const sentinel of [ "private-conversation-name", "private-unrelated-name", "private-instructions-sentinel", "private-memory-sentinel", "private-api-key-sentinel-long-enough", "fixture-private-credential", "fixture-private-output", "fixture-private-progress", "https://user:password", "Provider echoed", "native tool output", ]) { assert.ok(!JSON.stringify(exported).includes(sentinel), sentinel); assert.ok( !JSON.stringify(raw).includes(sentinel), `persisted ${sentinel}`, ); } const task = await owner.call("createTask", [ { id: crypto.randomUUID(), conversationId: target.id, name: "private-task-name", instructions: "second", enabled: true, schedule: { kind: "once", at: new Date( Math.ceil((Date.now() + 2500) / 1000) * 1000, ).toISOString(), }, }, ]); chat.close(); owner.close(); await wait(4500); // No requests or sockets wake the runtime during its native alarm. owner = await connect(); exported = await owner.call("exportDiagnostics", [target.id]); // Does not reconcile task state. const completed = exported.runtime.events.find( (e) => e.kind === "task" && e.status === "completed" && e.details.taskId === diagnosticId(task.id), ); assert.ok( completed, "native scheduled completion projected before an owner read", ); assert.equal(completed.durationMs, null); assert.equal(typeof completed.details.claimToCompletionMs, "number"); assert.ok( Number.isFinite(completed.details.claimToCompletionMs) && completed.details.claimToCompletionMs >= 0, ); assert.ok(!JSON.stringify(exported).includes("private-task-name")); assert.equal( (await owner.call("exportDiagnostics")).runtime.events.filter( (e) => e.kind === "task", ).length, 0, ); const oldIds = exported.conversation.events.map((e) => e.id); for (const client of clients) client.close(); await lifetime.release(worker); worker = await start(); owner = await connect(); exported = await owner.call("exportDiagnostics", [target.id]); assert.ok( exported.conversation.events.some((e) => oldIds.includes(e.id)), ); chat = await connect(target.id); transport = new WebSocketChatTransport({ agent: chat }); const abort = new AbortController(); const slow = await send("slow-first", abort); await until( () => owner.call("exportDiagnostics", [target.id]), (e) => e.conversation.events.some( (row) => row.kind === "model-attempt" && row.phase === "start" && !e.conversation.events.some( (end) => end.id === row.id && end.phase === "finish", ), ), ); void slow.done.catch(() => {}); chat.send(JSON.stringify({ type: MessageType.CF_AGENT_CHAT_CLEAR })); await until( () => history(target.id), (messages) => messages.length === 0, ); abort.abort(); await slow.done.catch(() => {}); await wait(100); const cleared = await owner.call("exportDiagnostics", [target.id]); assert.equal( cleared.conversation.events.length, 0, "late cancellation/finish callbacks cannot repopulate clear", ); await ( await send("second") ).done; assert.ok( ( await owner.call("exportDiagnostics", [target.id]) ).conversation.events.some((e) => e.kind === "model"), ); await fetch(`${origin}/__fixture/diagnostics?id=${target.id}&fail=true`); const afterFailure = await send("after-error"); await afterFailure.done; assert.ok( afterFailure.chunks.some( (c) => c.type === "text-delta" && c.delta.includes("Reply"), ), "diagnostic storage failure cannot fail a turn", ); assert.equal( (await owner.call("exportDiagnostics", [target.id])).conversation .available, false, ); await owner.call("deleteConversation", [target.id]); await assert.rejects( owner.call("exportDiagnostics", [target.id]), /Conversation not found/, ); assert.equal( (await owner.call("exportDiagnostics")).runtime.events.filter( (e) => e.conversationId === diagnosticId(target.id), ).length, 0, ); const inspection = await ( await fetch(origin + "/__fixture/inspect") ).json(); assert.ok(!inspection.facets.includes(target.id)); lifetime.complete(); } finally { await lifetime.cleanup(); } },);
registerCleanupProbe(import.meta.url, "diagnostics");