Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530import assert from "node:assert/strict";import { createHash, generateKeyPairSync, randomBytes, sign,} from "node:crypto";import { execFileSync } from "node:child_process";import { cp, mkdir, mkdtemp, readFile, readdir, rm, writeFile,} from "node:fs/promises";import { createServer, request as httpRequest } from "node:http";import { tmpdir } from "node:os";import { join, resolve } from "node:path";import { test } from "node:test";import { parse } from "jsonc-parser";import { AgentClient } from "agents/client";import WebSocket from "ws";import { unstable_dev } from "wrangler";import { customerBindings, installation } from "./fixtures/config.mjs";import { LOGIN_PURPOSE } from "../shared/bridge.ts";
const opaque = () => randomBytes(32).toString("base64url");const sha256 = (bytes) => createHash("sha256").update(bytes).digest("hex");const personalPath = "/agents/personal-agent/personal";const releases = ["0.1.0-fixture.1", "0.1.0-fixture.2"];const delay = (ms) => new Promise((done) => setTimeout(done, ms));async function until(read, accepts) { const deadline = Date.now() + 30_000; while (Date.now() < deadline) { const value = await read(); if (accepts(value)) return value; await delay(100); } assert.fail("Native conversation/task did not complete before deadline");}async function listen(server) { await new Promise((done) => server.listen(0, "127.0.0.1", done)); return server.address().port;}async function close(server) { server.closeAllConnections(); await new Promise((done) => server.close(done));}const cookie = (response, name) => response.headers .getSetCookie() .find((value) => value.startsWith(name + "=")) ?.split(";")[0];
// These archives are test identities, never entries in the publisher catalog.// Compile twice, then start workerd with no_bundle so replacement uses exactly// the recorded bytes; changing an environment variable is not this proof.async function buildArtifact( temporary, base, version, sourceRoot = process.cwd(),) { const directory = join(temporary, version); const assets = join(directory, "assets"); await mkdir(directory); await cp(resolve("dist/client"), assets, { recursive: true }); await writeFile(join(assets, "upgrade-fixture.txt"), version + "\n"); const config = { ...base, name: "flarebot-upgrade-state-fixture", main: resolve(sourceRoot, "tests/fixtures/upgrade-customer-worker.ts"), assets: { ...base.assets, directory: assets }, define: { __UPGRADE_FIXTURE_RELEASE__: JSON.stringify(version) }, }; const configPath = join(directory, "build.json"); await writeFile(configPath, JSON.stringify(config)); execFileSync( process.execPath, [ resolve("node_modules/wrangler/bin/wrangler.js"), "deploy", "--dry-run", "--config", configPath, "--outdir", join(directory, "worker"), ], { stdio: "pipe" }, ); const main = join(directory, "worker/upgrade-customer-worker.js"); assert.ok((await readFile(main, "utf8")).includes(version)); const inventory = []; async function walk(directory, prefix = "") { for (const item of (await readdir(directory, { withFileTypes: true })).sort( (a, b) => a.name.localeCompare(b.name), )) { const path = prefix + item.name; if (item.isDirectory()) await walk(join(directory, item.name), path + "/"); else if (!path.endsWith(".map") && !path.endsWith("README.md")) inventory.push({ path, sha256: sha256(await readFile(join(directory, item.name))), }); } } await walk(join(directory, "worker"), "worker/"); await walk(assets, "assets/"); const manifest = { version, files: inventory }; await writeFile(join(directory, "manifest.json"), JSON.stringify(manifest)); delete config.define; config.main = main; config.no_bundle = true; const runtimeConfig = join(directory, "runtime.json"); await writeFile(runtimeConfig, JSON.stringify(config)); return { version, main, config: runtimeConfig, manifest, digest: sha256(JSON.stringify(manifest)), };}
test( "distinct compiled customer artifacts preserve native SQLite, facets, BYOK, owner login and scheduled execution", { timeout: 240_000 }, async () => { const temporary = await mkdtemp(join(tmpdir(), "flarebot-upgrade-state-")); const base = parse(await readFile("wrangler.jsonc", "utf8")); const key = generateKeyPairSync("ed25519"); const pin = { keyId: "upgrade-fixture", publicKey: key.publicKey.export({ format: "jwk" }).x, }; const installed = { ...installation, bridge: pin }; const codes = new Map(); let exchangeCount = 0; // The control-plane assertion issuer is the explicit network seam. The // customer performs its real challenge, PKCE exchange, assertion validation // and session-cookie issuance. OAuth/provider upload contracts are separate. const bridge = createServer(async (req, res) => { try { const chunks = []; for await (const chunk of req) chunks.push(chunk); const body = new URLSearchParams(Buffer.concat(chunks).toString()); const expected = codes.get(body.get("code")); assert.ok(expected); assert.equal(body.get("state"), expected.state); assert.equal(body.get("installationId"), installed.installationId); assert.equal(body.get("audience"), installed.runtimeOrigin); assert.equal( createHash("sha256").update(body.get("verifier")).digest("base64url"), expected.challenge, ); codes.delete(body.get("code")); const now = Math.floor(Date.now() / 1000); const claims = { iss: installed.controlPlaneOrigin, aud: installed.runtimeOrigin, sub: installed.ownerSubject, installationId: installed.installationId, purpose: LOGIN_PURPOSE, ...expected, iat: now, exp: now + 60, jti: opaque(), }; const payload = [{ alg: "EdDSA", kid: pin.keyId, typ: "JWT" }, claims] .map((part) => Buffer.from(JSON.stringify(part)).toString("base64url"), ) .join("."); exchangeCount++; res.writeHead(200, { "Content-Type": "application/json" }); res.end( JSON.stringify({ assertion: payload + "." + sign(null, Buffer.from(payload), key.privateKey).toString( "base64url", ), }), ); } catch { res.writeHead(403); res.end(); } }); let worker, ownerCookie; const clients = []; try { const artifacts = []; for (const [index, version] of releases.entries()) artifacts.push( await buildArtifact( temporary, base, version, index === 0 ? process.env.FLAREBOT_UPGRADE_BASELINE : undefined, ), ); assert.notEqual(artifacts[0].digest, artifacts[1].digest); for (const path of [ "worker/upgrade-customer-worker.js", "assets/upgrade-fixture.txt", ]) assert.notEqual( artifacts[0].manifest.files.find((f) => f.path === path).sha256, artifacts[1].manifest.files.find((f) => f.path === path).sha256, ); const bridgePort = await listen(bridge); const reservation = createServer(); const port = await listen(reservation); await close(reservation); const origin = `http://127.0.0.1:${port}`; const headers = () => ({ Cookie: ownerCookie, Origin: installed.runtimeOrigin, }); const request = (path, init = {}) => new Promise((resolve, reject) => { // Node fetch overwrites Sec-Fetch-Mode with cors; preserve actual // navigation headers for production bridge routing. const outgoing = httpRequest( origin + path, { method: init.method ?? "GET", headers: init.headers, }, (incoming) => { const chunks = []; incoming.on("data", (chunk) => chunks.push(chunk)); incoming.on("end", () => { const headers = new Headers(); for (let i = 0; i < incoming.rawHeaders.length; i += 2) headers.append( incoming.rawHeaders[i], incoming.rawHeaders[i + 1], ); resolve( new Response(Buffer.concat(chunks), { status: incoming.statusCode, headers, }), ); }); }, ); outgoing.on("error", reject); outgoing.end(init.body); }); async function login() { const begin = await request("/auth/login", { headers: { "Sec-Fetch-Mode": "navigate" }, }); assert.equal(begin.status, 303); const redirect = new URL(begin.headers.get("location")); assert.equal(redirect.origin, installed.controlPlaneOrigin); const code = opaque(); const state = redirect.searchParams.get("state"); codes.set(code, { state, challenge: redirect.searchParams.get("challenge"), }); const callback = `/auth/callback?${new URLSearchParams({ code, state })}`; const completed = await request(callback, { headers: { "Sec-Fetch-Mode": "navigate", Cookie: cookie(begin, "__Host-flarebot-login"), }, }); assert.equal(completed.status, 303); assert.equal(completed.headers.get("location"), "/"); const session = cookie(completed, "__Host-flarebot-session"); assert.ok(session); assert.match( completed.headers .getSetCookie() .find((c) => c.startsWith("__Host-flarebot-session=")), /HttpOnly; Secure; SameSite=Lax/, ); assert.equal( ( await request(callback, { headers: { "Sec-Fetch-Mode": "navigate", Cookie: cookie(begin, "__Host-flarebot-login"), }, }) ).status, 403, ); return session; } async function fixture(action, body) { const response = await request("/__upgrade__/" + action, { headers: { ...headers(), "Content-Type": "application/json" }, ...(body ? { method: "POST", body: JSON.stringify(body) } : {}), }); assert.equal(response.status, 200, action); return response.json(); } async function connect() { class OwnerSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers: headers(), closeTimeout: 100 }); } } const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", name: "personal", WebSocket: OwnerSocket, }); clients.push(client); await client.ready; return client; } const start = (artifact) => unstable_dev(artifact.main, { config: artifact.config, vars: { ...customerBindings, FLAREBOT_INSTALLATION: JSON.stringify({ ...installed, release: { version: artifact.version, artifactDigest: artifact.digest, operationId: artifact === artifacts[0] ? "1".repeat(32) : "2".repeat(32), }, }), TEST_BRIDGE_PORT: String(bridgePort), }, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: join(temporary, "same-native-storage"), logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }); worker = await start(artifacts[0]); assert.equal((await request(personalPath)).status, 401); assert.equal((await request("/__upgrade__/snapshot")).status, 401); ownerCookie = await login(); let owner = await connect(); let call = (method, ...args) => owner.call(method, args); assert.deepEqual(await fixture("identity"), { release: releases[0] }); assert.equal( await (await request("/upgrade-fixture.txt")).text(), releases[0] + "\n", ); const first = await call("createConversation", "Preserved research"); const second = await call("createConversation", "Preserved planning"); const instructions = await call( "updateInstructions", "UPGRADE-PRIVATE-INSTRUCTIONS: use primary sources.", ); const memory = await call( "addMemory", "UPGRADE-PRIVATE-MEMORY: preferred timezone is Europe/Warsaw.", ); await call( "setProviderKey", "anthropic", "sk-ant-upgrade-preserved-private-sentinel", ); const model = await call("updateModelSettings", { provider: "anthropic", model: "claude-sonnet-5", }); for (const conversation of [first, second]) await fixture("submit", { id: conversation.id, text: "second" }); await until( () => fixture("snapshot"), (s) => s.conversations.length === 2 && s.conversations.every( (c) => c.messages.length === 2 && c.submissions[0]?.status === "completed", ), ); const future = await call("createTask", { id: crypto.randomUUID(), conversationId: second.id, name: "Preserved future native schedule", instructions: "second", enabled: true, schedule: { kind: "once", at: new Date( Math.ceil(Date.now() / 1000) * 1000 + 3_600_000, ).toISOString(), }, }); const executed = await call("createTask", { id: crypto.randomUUID(), conversationId: first.id, name: "Native alarm creates preserved history", instructions: "second", enabled: true, schedule: { kind: "once", at: new Date( Math.ceil(Date.now() / 1000) * 1000 + 3000, ).toISOString(), }, }); await until( () => call("listTaskRuns", executed.id), (page) => page.runs[0]?.status === "completed", ); const history = await call("listTaskRuns", executed.id); const before = await fixture("snapshot"); assert.equal(before.parentName, "personal"); assert.match(before.parentId, /^[a-f0-9]{64}$/); assert.equal(before.facets.length, 2); assert.deepEqual( before.facets.map((facet) => facet.className), ["Conversation", "Conversation"], ); assert.deepEqual( before.facets.map((facet) => facet.name).sort(), [first.id, second.id].sort(), ); assert.equal(before.runs.length, 1); assert.equal(JSON.parse(before.runs[0].payload).status, "completed"); assert.equal(before.keyDecrypts, true); assert.equal(before.settings[0].anthropic_key, null); assert.ok(before.credentials[0].encrypted_key); assert.ok(!JSON.stringify(before.credentials).includes("sk-ant-upgrade")); assert.equal(before.schedules.length, 1); assert.equal(before.schedules[0].payload.taskId, future.id); const previousCookie = ownerCookie; for (const client of clients) client.close(); await worker.stop(); worker = undefined; worker = await start(artifacts[1]); assert.deepEqual(await fixture("identity"), { release: releases[1] }); assert.equal( await (await request("/upgrade-fixture.txt")).text(), releases[1] + "\n", ); assert.equal( (await request(personalPath, { headers: headers() })).status, 200, "old owner cookie authorizes new artifact", ); owner = await connect(); call = (method, ...args) => owner.call(method, args); const after = await fixture("snapshot"); assert.deepEqual( after, before, "same parent/facet IDs, timestamps, messages, native submissions, ciphertext, schedules and history", ); assert.deepEqual(await call("getInstructionSettings"), instructions); assert.deepEqual(await call("listMemories"), [memory]); assert.deepEqual(await call("getModelSettings"), model); assert.deepEqual(await call("listTaskRuns", executed.id), history); ownerCookie = await login(); assert.equal( exchangeCount, 2, "fresh post-upgrade login completes another real PKCE exchange", ); assert.equal( (await request(personalPath, { headers: headers() })).status, 200, ); assert.equal( (await request(personalPath, { headers: { Cookie: previousCookie } })) .status, 200, ); assert.equal((await request(personalPath)).status, 401); // Existing facets and native task machinery remain writable in the new // artifact, beyond a successful read of old SQLite rows. await fixture("submit", { id: second.id, text: "after-error" }); await until( () => fixture("snapshot"), (s) => s.conversations.find((c) => c.name === second.id).messages.length === 4 && s.conversations .find((c) => c.name === second.id) .submissions.every( (submission) => submission.status === "completed", ), ); await call("runTaskNow", future.id, crypto.randomUUID()); await until( () => call("listTaskRuns", future.id), (page) => page.runs[0]?.status === "completed", ); const final = await fixture("snapshot"); assert.equal(final.parentId, before.parentId); assert.deepEqual(final.facets, before.facets); assert.deepEqual(final.credentials, before.credentials); assert.equal(final.keyDecrypts, true); assert.deepEqual(final.schedules, before.schedules); assert.equal(final.runs.length, 2); // Re-read checksums to guard against accidental rebuilding/mutation of the // already accepted old/new fixture artifacts during native replacement. for (const artifact of artifacts) for (const file of artifact.manifest.files) assert.equal( sha256( await readFile(join(temporary, artifact.version, file.path)), ), file.sha256, ); } finally { for (const client of clients) client.close(); await worker?.stop(); await close(bridge); await rm(temporary, { recursive: true, force: true }); } },);