Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365import 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_dev } from "wrangler";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";
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) => { response.resume(); socket.terminate(); try { assert.equal(response.statusCode, status, url); resolve(); } catch (error) { reject(error); } }); });}
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 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 = (overrides = {}) => unstable_dev("dist/release/worker/index.js", { config: configPath, vars: { ...customerBindings, FLAREBOT_ENV: "development", FLAREBOT_DEV_OVERRIDES: JSON.stringify({ runtimeOrigin: origin }), ...overrides, }, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: persistence, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }); 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 fetch(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"]) { const response = await fetch(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 }), }; assert.equal( ( await fetch(origin + root, { method: "POST", headers: requestHeaders, }) ).status, 403, ); await deniedSocket(wsOrigin + root, requestHeaders, 403); } assert.equal( ( await fetch(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 fetch(origin + path, { headers })).status, 404, path, ); await deniedSocket(wsOrigin + path, headers, 404); } const response = await fetch(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 fetch(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 fetch(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 createOwnerSession(secret, { ...effectiveInstallation, ...changed, }) ).split(";")[0]; const denied = await fetch(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 }); } },);