Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534import * as Effect from "effect/Effect";import assert from "node:assert/strict";import { createHmac } from "node:crypto";import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";import { request as httpRequest } from "node:http";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 WebSocket from "ws";import { unstable_getMiniflareWorkerOptions } from "wrangler";import { Miniflare, convertV4MiniflareOptions, Log, LogLevel } from "miniflare";import { Secret } from "../configuration/secrets.ts";import { createOwnerSession, SESSION_COOKIE } from "../worker/session.ts";import { customerBindings, installation } from "./fixtures/config.mjs";
const root = "/agents/personal-agent/personal";const skillRoutes = [ "/api/skills/manage", "/api/skills/review-upload", "/api/skills/review-url", "/api/skills/install",];const mcpAuthorization = "/api/mcp/mcp-00000000-0000-4000-8000-000000000000/authorize";const mcpCallback = `${root}/mcp-callback?code=private-code&state=forged.mcp-00000000-0000-4000-8000-000000000000`;
async function runtimeFetch(site, url, options) { try { // Local workerd can close after rejecting an unread POST body without a // close header. Do not reuse that socket for the next fixture request. const headers = new Headers(options?.headers); headers.set("Connection", "close"); return await fetch(url, { ...options, headers }); } catch (cause) { throw new Error( `Runtime request failed at ${site}: ${options?.method ?? "GET"} ${new URL(url).pathname}`, { cause }, ); }}
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;}
async function waitFor(predicate) { const deadline = Date.now() + 10_000; while (!predicate()) { if (Date.now() > deadline) throw new Error("Timed out waiting for native Agent event"); await new Promise((resolve) => setTimeout(resolve, 20)); }}
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) => { let bodyPrefix = ""; if (response.statusCode !== status) { response.setEncoding("utf8"); response.on("data", (chunk) => { bodyPrefix = (bodyPrefix + chunk).slice(0, 1024); }); } response.once("error", reject); response.once("end", () => { socket.terminate(); try { assert.equal( response.statusCode, status, JSON.stringify({ url, origin: headers.Origin ?? "(missing)", contentType: response.headers["content-type"] ?? null, bodyPrefix, }), ); resolve(); } catch (error) { reject(error); } }); response.resume(); }); });}
test( "private native Agent routing, state sync and restart persistence", { timeout: 120_000 }, async () => { const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const wsOrigin = `ws://127.0.0.1:${port}`; const effectiveInstallation = { ...installation, runtimeOrigin: origin }; const secret = new Secret(customerBindings.FLAREBOT_SESSION_SECRET); const cookieHeader = await Effect.runPromise( createOwnerSession(secret, effectiveInstallation), ); assert.match( cookieHeader, /; Path=\/; HttpOnly; Secure; SameSite=Lax; Max-Age=28800$/, ); const cookie = cookieHeader.split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; const persistence = await mkdtemp(join(tmpdir(), "flarebot-runtime-")); const deployment = JSON.parse( await readFile("dist/release/deployment.json", "utf8"), ); const configPath = join(persistence, "wrangler.json"); // The account-neutral artifact intentionally omits a Worker name. Give the // local installation a stable name so restarts address the same namespace. await writeFile( configPath, JSON.stringify({ ...deployment, name: "flarebot-runtime-test", main: resolve("dist/release/worker/index.js"), assets: { ...deployment.assets, directory: resolve("dist/release/assets"), }, }), ); const start = async (overrides = {}) => { const { workerOptions, main, externalWorkers } = unstable_getMiniflareWorkerOptions(configPath, undefined, { overrides: { enableContainers: false }, }); const runtime = new Miniflare( convertV4MiniflareOptions({ host: "127.0.0.1", port, resourcePersistencePath: persistence, log: new Log(LogLevel.ERROR), workers: [ { ...workerOptions, name: "flarebot-runtime-test", modules: true, modulesRoot: resolve("dist/release/worker"), // The release contains one already-bundled JavaScript module. modulesRules: undefined, scriptPath: main, bindings: { ...workerOptions.bindings, ...customerBindings, FLAREBOT_ENV: "development", FLAREBOT_DEV_OVERRIDES: JSON.stringify({ runtimeOrigin: origin, }), ...overrides, }, }, ...externalWorkers, ], }), ); try { await runtime.ready; return { stop: () => runtime.dispose() }; } catch (error) { await runtime.dispose(); throw error; } }; let worker; const clients = []; try { worker = await start(); const localLogin = await new Promise((resolve, reject) => { const request = httpRequest(origin + "/auth/login", (response) => { response.resume(); response.on("end", () => resolve(response)); }); request.on("error", reject); request.end(); }); assert.equal(localLogin.statusCode, 303); assert.equal(localLogin.headers.location, "/"); const localCookie = localLogin.headers["set-cookie"]?.[0]; assert.match( localCookie, /^__Host-flarebot-session=.+; Path=\/; HttpOnly; Secure; SameSite=Lax; Max-Age=28800$/, ); assert.equal( ( await runtimeFetch( "development login status", origin + `${root}/status`, { headers: { Cookie: localCookie.split(";")[0] }, }, ) ).status, 200, ); const [body] = cookie.slice(SESSION_COOKIE.length + 1).split("."); const claims = JSON.parse(Buffer.from(body, "base64url").toString()); const sign = (changes) => { const encoded = Buffer.from( JSON.stringify({ ...claims, ...changes }), ).toString("base64url"); return `${SESSION_COOKIE}=${encoded}.${createHmac("sha256", secret.reveal()).update(encoded).digest("base64url")}`; }; const invalidCookies = [ "", `${cookie}tampered`, `${SESSION_COOKIE}=not-a-token`, `${cookie}; ${cookie}`, sign({ subject: "another-owner" }), sign({ installationId: "f".repeat(32) }), sign({ audience: "https://other.example.com" }), sign({ issuedAt: claims.issuedAt - 100, expiresAt: claims.issuedAt - 1, }), sign({ issuedAt: claims.issuedAt + 100 }), ]; for (const invalid of invalidCookies) { for (const path of [ root, `${root}/get-messages`, "/api/status", mcpAuthorization, mcpCallback, ...skillRoutes, ]) { const response = await runtimeFetch( "invalid cookie route", origin + path, { headers: { Cookie: invalid }, }, ); assert.equal(response.status, 401, path); assert.equal(response.headers.get("cache-control"), "no-store"); assert.equal(await response.text(), "Authentication required"); } await deniedSocket( wsOrigin + root, { Cookie: invalid, Origin: origin }, 401, ); } for (const badOrigin of [ undefined, "null", "https://attacker.example.com", ]) { const requestHeaders = { Cookie: cookie, ...(badOrigin === undefined ? {} : { Origin: badOrigin }), }; const deniedRoot = await runtimeFetch( "invalid origin root", origin + root, { method: "POST", headers: requestHeaders, }, ); assert.equal(deniedRoot.status, 403); assert.equal(await deniedRoot.text(), "Forbidden"); const deniedAuthorization = await runtimeFetch( "invalid origin MCP authorization", origin + mcpAuthorization, { method: "POST", headers: requestHeaders, redirect: "manual", body: new URLSearchParams({ revision: "1" }), }, ); assert.equal(deniedAuthorization.status, 403); assert.equal(await deniedAuthorization.text(), "Forbidden"); // Repeat body-bearing denials to catch reuse of a socket closed after // an early response. Every request must retain its original 403 body. for (let round = 0; round < 100; round++) { for (const path of skillRoutes) { const deniedSkill = await runtimeFetch( `invalid origin Skill route (round ${round})`, origin + path, { method: "POST", headers: requestHeaders, body: "{}", }, ); assert.equal(deniedSkill.status, 403); assert.equal(await deniedSkill.text(), "Forbidden"); } } await deniedSocket(wsOrigin + root, requestHeaders, 403); } assert.equal( ( await runtimeFetch("cross-origin read", origin + root, { headers: { ...headers, Origin: "https://attacker.example.com" }, }) ).status, 403, ); for (const path of [ "/agents/PersonalAgent/personal", "/agents/personal-agent/other", "/agents/other/personal", `${root}/sub/conversation/unknown`, `${root}/get-messages`, `${root}/anything`, "/agents/personal-agent/%70ersonal", ]) { assert.equal( ( await runtimeFetch("unknown native route", origin + path, { headers, }) ).status, 404, path, ); await deniedSocket(wsOrigin + path, headers, 404); } const invalidMcpCallback = await runtimeFetch( "invalid MCP callback", origin + mcpCallback, { headers: { Cookie: cookie }, redirect: "manual", }, ); assert.equal(invalidMcpCallback.status, 303); assert.equal( invalidMcpCallback.headers.get("Location"), "/extensions?mcpAuth=failed", ); assert.equal(invalidMcpCallback.headers.get("Cache-Control"), "no-store"); assert.equal( invalidMcpCallback.headers.get("Referrer-Policy"), "no-referrer", ); assert.doesNotMatch( await invalidMcpCallback.text(), /private-code|forged/, ); assert.equal( ( await runtimeFetch( "MCP authorization method", origin + mcpAuthorization, { headers, redirect: "manual", }, ) ).status, 405, ); const response = await runtimeFetch( "initial native status", origin + `${root}/status`, { headers }, ); assert.equal(response.status, 200); assert.equal(response.headers.get("cache-control"), "no-store"); const initial = await response.json(); assert.deepEqual(Object.keys(initial).sort(), [ "createdAt", "schemaVersion", ]); assert.equal(initial.schemaVersion, 1); assert.ok(Number.isFinite(Date.parse(initial.createdAt))); assert.equal( ( await runtimeFetch("native method rejection", origin + root, { method: "POST", headers, }) ).status, 405, );
// Inject browser-equivalent cookies/Origin in Node only. The native client // owns identity, RPC, state protocol and reconnect/backoff in both contexts. class AuthenticatedSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers }); } } function connect() { const states = []; const errors = []; const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", name: "personal", WebSocket: AuthenticatedSocket, minReconnectionDelay: 50, maxReconnectionDelay: 250, onStateUpdate: (state) => states.push(state), onStateUpdateError: (error) => errors.push(error), }); clients.push(client); return { client, states, errors }; } const first = connect(); const second = connect(); await Promise.all([first.client.ready, second.client.ready]); await waitFor(() => first.states.length && second.states.length); assert.deepEqual(first.states.at(-1), initial); assert.deepEqual(second.states.at(-1), initial); assert.deepEqual(await first.client.call("getStatus"), initial); await assert.rejects( first.client.call("setState", [{ schemaVersion: 999 }]), ); first.client.setState({ schemaVersion: 999, createdAt: "overwritten", secret: "injected-secret", }); await waitFor(() => first.errors.length); assert.deepEqual(await first.client.call("getStatus"), initial); assert.deepEqual(second.states.at(-1), initial);
first.client.close(); const reconnect = connect(); await reconnect.client.ready; await waitFor(() => reconnect.states.length); assert.deepEqual(reconnect.states.at(-1), initial);
const beforeRestart = second.states.length; await worker.stop(); worker = undefined; worker = await start(); await waitFor(() => second.states.length > beforeRestart); await second.client.ready; assert.deepEqual(second.states.at(-1), initial); assert.deepEqual(await second.client.call("getStatus"), initial); assert.deepEqual( await ( await runtimeFetch("status after restart", origin + root, { headers }) ).json(), initial, ); assert.ok(!JSON.stringify(second.states).includes(secret.reveal())); const now = Math.floor(Date.now() / 1000); const shortCookie = sign({ issuedAt: now, expiresAt: now + 2 }); const expiring = new WebSocket(wsOrigin + root, { // Local forwarding may leave TCP half-open after the close frame. // Bound Node ws transport cleanup; the server must still send 4001. closeTimeout: 100, headers: { Origin: origin, Cookie: shortCookie }, }); try { const closed = await new Promise((resolve, reject) => { const timeout = setTimeout( () => reject(new Error("Session socket did not expire")), 10_000, ); expiring.on("error", reject); expiring.on("close", (code, reason) => { clearTimeout(timeout); resolve({ code, reason: reason.toString() }); }); }); assert.deepEqual(closed, { code: 4001, reason: "Session expired" }); await deniedSocket( wsOrigin + root, { Origin: origin, Cookie: shortCookie }, 401, ); } finally { expiring.terminate(); } for (const client of clients) client.close(); // A configuration typo must never reassign an existing personal database // to another owner or installation, even with a valid newly signed cookie. for (const changed of [ { ownerSubject: "replacement-owner" }, { installationId: "f".repeat(32) }, ]) { await worker.stop(); worker = undefined; worker = await start({ FLAREBOT_INSTALLATION: JSON.stringify({ ...installation, ...changed, }), }); const changedCookie = ( await Effect.runPromise( createOwnerSession(secret, { ...effectiveInstallation, ...changed, }), ) ).split(";")[0]; const denied = await runtimeFetch( "changed installation identity", origin + root, { headers: { Origin: origin, Cookie: changedCookie }, }, ); assert.equal(denied.status, 503); assert.equal(await denied.text(), "Personal runtime unavailable"); await deniedSocket( wsOrigin + root, { Origin: origin, Cookie: changedCookie }, 503, ); } } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);