Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710import assert from "node:assert/strict";import { generateKeyPairSync } from "node:crypto";import { mkdtemp, readFile, writeFile, rm, mkdir } from "node:fs/promises";import { createServer } 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 { chromium } from "playwright";import { unstable_dev } from "wrangler";import { oauthBindings, oauthConfig } from "./fixtures/oauth-config.mjs";import { compileGolden, writeAccepted, acceptedConfiguration, goldenProxy, listen, close, freePort, requestBytes, until, sha,} from "./fixtures/golden-harness.mjs";import "./fixtures/config.mjs";
const CP = "https://publisher.test";test( "v0.1 connected golden path: native install, tools, unattended task and persisted history", { timeout: 420000 }, async (t) => { const temporary = await mkdtemp(join(tmpdir(), "flarebot-golden-")); const cpPort = await freePort(), customerPort = await freePort(); const base = parse(await readFile("wrangler.control-plane.jsonc", "utf8")); let cp, customer, browser, context, proxy, runner, runnerPort, accepted, customerOrigin, customerConfig, completion, disconnectedAt; let acceptedUploads = 0, acceptedApplications = 0, runnerError; const completions = []; const assets = new Map(); let assetManifest; const keys = generateKeyPairSync("ed25519"); const bridge = { keyId: "golden-key", publicKey: keys.publicKey.export({ format: "jwk" }).x, }; let artifact; let failed = false; const checkpoint = async (name, action) => { let failure; await t.test(name, async () => { try { await action(); } catch (error) { failure = error; failed = true; throw error; } }); if (failure) throw failure; }; const options = (config, port, vars, containers = false) => ({ config, vars, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: join(temporary, "state", String(port)), logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false, enableContainers: containers, }, }); const startCustomer = () => unstable_dev( join(temporary, "accepted/worker/index.js"), options( customerConfig, customerPort, { ...accepted.variables, TEST_CP_PORT: String(cpPort), TEST_RUNNER_PORT: String(runnerPort), }, true, ), ); t.after(async () => { const failures = []; for (const cleanup of [ () => browser?.close(), () => proxy?.close(), () => customer?.stop(), () => cp?.stop(), () => close(runner), () => rm(temporary, { recursive: true, force: true }), ]) { try { await cleanup(); } catch (error) { failures.push(error); t.diagnostic(`Cleanup failed: ${error.message}`); } } if (failures.length && !failed) throw new AggregateError(failures, "Golden cleanup failed"); }); artifact = await compileGolden(temporary); runner = createServer(async (req, res) => { try { if (req.url === "/artifact") { res.writeHead(200, { "Content-Type": "application/json" }); res.end( JSON.stringify({ identity: artifact.identity, files: Object.fromEntries( Object.entries(artifact.files).map(([path, bytes]) => [ path, bytes.toString("base64"), ]), ), }), ); return; } if (req.url === "/completed") { const event = JSON.parse(await requestBytes(req)); const received = { ...event, receivedAt: Date.now() }; if (!completion) { assert.ok( disconnectedAt, "Native completion must follow disconnection", ); assert.equal( proxy.evidence.customerSockets.size, 0, "Customer socket remained open at native completion", ); assert.equal( proxy.evidence.customerRequests, disconnectedAt.requests, "An owner request preceded native completion", ); completion = received; } else { // A replayed terminal notification after reconnect is not new work. assert.equal(event.submissionId, completion.submissionId); assert.equal(event.run.id, completion.run.id); assert.equal(event.run.taskId, completion.run.taskId); } completions.push(received); res.writeHead(200); res.end(); return; } assert.equal(req.url, "/accepted"); const path = req.headers["x-provider-path"]; const body = await requestBytes(req); if (path.endsWith("/containers/applications")) { const deployment = JSON.parse(artifact.files["deployment.json"]), container = deployment.containers[0]; const workerName = `flarebot-${accepted.installation.installationId}`; assert.deepEqual(JSON.parse(body), { name: `flarebot-shell-${accepted.installation.installationId}`, instances: 0, configuration: { image: container.image, instance_type: container.instance_type, }, max_instances: container.max_instances, scheduling_policy: "default", durable_objects: { namespace_id: sha(workerName + "Sandbox").slice(0, 32), }, }); acceptedApplications++; } else if (path.endsWith("/assets-upload-session")) assetManifest = JSON.parse(body).manifest; else { const form = await new Request("http://local/", { method: "POST", headers: { "Content-Type": req.headers["content-type"] }, body, }).formData(); if (path.endsWith("/workers/assets/upload")) { for (const [hash, file] of form) assets.set(hash, Buffer.from(await file.text(), "base64")); } else { assert.equal(req.headers["x-provider-method"], "PUT"); assert.equal( acceptedUploads++, 0, "Fresh install must upload once", ); const metadata = JSON.parse(await form.get("metadata").text()); const modules = new Map(); for (const [name, file] of form) if (name !== "metadata") modules.set(name, Buffer.from(await file.arrayBuffer())); await writeAccepted( join(temporary, "accepted"), artifact, metadata, modules, assetManifest, assets, ); const deployment = JSON.parse(artifact.files["deployment.json"]); const { variables, installation } = acceptedConfiguration( metadata, deployment, path.split("/").at(-1), ); assert.equal( installation.release.artifactDigest, artifact.identity.artifactDigest, ); assert.equal(installation.controlPlaneOrigin, CP); assert.deepEqual(installation.bridge, bridge); customerOrigin = installation.runtimeOrigin; accepted = { variables, installation, metadata }; customerConfig = join(temporary, "customer.json"); // Local-only adapters: filesystem asset/module paths, native local // resource allocation and loopback external boundaries. All executable // binding/runtime settings above must equal the accepted provider PUT. await writeFile( customerConfig, JSON.stringify({ ...deployment, name: path.split("/").at(-1), main: join(temporary, "accepted/worker/index.js"), no_bundle: true, assets: { directory: join(temporary, "accepted/assets"), binding: "ASSETS", }, }), ); customer = await startCustomer(); } } res.writeHead(200); res.end(); } catch (error) { runnerError = error; res.writeHead(500); res.end("Golden runner boundary failed"); } }); runnerPort = await listen(runner); const cpConfig = join(temporary, "cp.json"); await writeFile( cpConfig, JSON.stringify({ ...base, name: "flarebot-golden-cp", main: resolve("tests/fixtures/golden-control-worker.ts"), assets: { directory: resolve("dist/control-plane/client"), binding: "ASSETS", }, durable_objects: { bindings: [ ...base.durable_objects.bindings, { name: "PROVIDER", class_name: "Provider" }, ], }, exports: { ...base.exports, Provider: { type: "durable-object", storage: "sqlite" }, }, }), ); cp = await unstable_dev( "tests/fixtures/golden-control-worker.ts", options(cpConfig, cpPort, { ...oauthBindings(CP), TEST_RUNNER_PORT: String(runnerPort), TEST_CUSTOMER_PORT: String(customerPort), FLAREBOT_CONTROL_PLANE: JSON.stringify({ ...oauthConfig(CP), bridge }), FLAREBOT_BRIDGE_SIGNING_KEY: keys.privateKey .export({ type: "pkcs8", format: "der" }) .toString("base64url"), }), ); proxy = await goldenProxy(temporary, cpPort, customerPort); browser = await chromium.launch({ headless: true, args: [ "--no-proxy-server", `--host-resolver-rules=MAP dash.cloudflare.com 127.0.0.1:${proxy.port}, MAP publisher.test 127.0.0.1:${proxy.port}, MAP *.fixture-account.workers.dev 127.0.0.1:${proxy.port}`, ], }); context = await browser.newContext({ ignoreHTTPSErrors: true }); let page = await context.newPage(); page.setDefaultTimeout(30000); let actionError; page.on("response", async (response) => { if ( response.request().method() === "POST" && new URL(response.url()).pathname.startsWith("/api/") && !response.ok() ) { const data = await response.json().catch(() => ({})); actionError = new Error( `Installation POST ${new URL(response.url()).pathname}: ${response.status()} ${data.error ?? "unknown"}`, ); } }); const errors = []; page.on("pageerror", (error) => errors.push(error.message)); let firstId, secondId, taskId, storage; const newConversation = async () => { const previous = page.url(); await page .getByRole("button", { name: "New conversation", exact: true }) .click(); await page.waitForURL( (url) => url.href !== previous && url.pathname.startsWith("/conversations/"), ); await page .getByRole("textbox", { name: "Message", exact: true }) .waitFor(); }; const send = async (text, response) => { await page .getByRole("textbox", { name: "Message", exact: true }) .fill(text); await page .getByRole("button", { name: "Send message", exact: true }) .click(); await page.getByText(response, { exact: false }).last().waitFor(); await page .getByRole("button", { name: "Send message", exact: true }) .waitFor(); }; await checkpoint( "1. Fresh installation passes real signed health and owner login", async () => { await page.goto(CP + "/connect"); await page .getByRole("button", { name: "Connect Cloudflare", exact: true }) .click(); await page .getByRole("button", { name: "Select Personal account", exact: true }) .click(); await page .getByText( "Account selected. New installations will use this account.", ) .waitFor(); await page .getByRole("button", { name: "Install Flarebot", exact: true }) .click(); await until( async () => { if (runnerError) throw runnerError; if (actionError) throw actionError; const value = await page.evaluate( async () => (await (await fetch("/api/installations")).json()) .installations[0], ); if (value?.status === "failed") throw new Error(`Installation failed: ${value.errorCode}`); return value; }, (value) => value?.status === "ready", "production installation health", 180000, ); assert.equal(acceptedUploads, 1); assert.equal(acceptedApplications, 1); assert.ok(customer); const record = await page.evaluate( async () => (await (await fetch("/api/installations")).json()).installations[0], ); assert.equal(record.status, "ready"); assert.equal( record.installedRelease.artifactDigest, artifact.identity.artifactDigest, ); assert.equal( record.operationId, accepted.installation.release.operationId, ); await page .getByRole("link", { name: "Open Flarebot", exact: true }) .click(); await page.waitForURL(customerOrigin + "/"); await page .getByRole("button", { name: "New conversation", exact: true }) .waitFor(); const cookie = (await context.cookies(customerOrigin)).find( (c) => c.name === "__Host-flarebot-session", ); assert.ok(cookie?.secure && cookie.httpOnly); assert.equal(cookie.sameSite, "Lax"); assert.ok( proxy.evidence.navigations.some( (n) => n.path === "/auth/login" && n.host !== "publisher.test", ), ); assert.ok( proxy.evidence.navigations.some( (n) => n.path === "/auth/callback" && n.host !== "publisher.test", ), ); assert.ok(proxy.evidence.oauthVisits >= 1); }, ); await checkpoint( "2. Configure and persist a nondefault model through Settings", async () => { await page.getByRole("link", { name: "Settings", exact: true }).click(); const section = page.getByRole("region", { name: "Model and provider", exact: true, }); await section .getByRole("combobox", { name: "Model", exact: true }) .click(); await page .getByRole("option", { name: "@cf/meta/llama-4-scout-17b-16e-instruct · Free allowance", exact: true, }) .click(); await section .getByRole("button", { name: "Save model", exact: true }) .click(); await section.getByText("Model saved.", { exact: false }).waitFor(); await page.reload(); assert.match( await section .getByRole("combobox", { name: "Model", exact: true }) .innerText(), /llama-4-scout/, ); }, ); await checkpoint("3. Ask the agent to remember a fact", async () => { await newConversation(); await send( "Please remember that My favorite orchard fruit is quince.", "Saved: My favorite orchard fruit is quince.", ); firstId = new URL(page.url()).pathname.split("/").at(-1); }); await checkpoint( "4. Recall real memory in a new conversation", async () => { await newConversation(); await send( "What is my favorite orchard fruit?", "Memory answer: My favorite orchard fruit is quince.", ); secondId = new URL(page.url()).pathname.split("/").at(-1); assert.notEqual(firstId, secondId); }, ); await checkpoint( "5. Research through native web tools and cite read page evidence", async () => { await send( "Research the garden report on the web and cite the page.", "Web evidence:", ); await page .getByText("17 cedar trees", { exact: false }) .last() .waitFor(); assert.ok( await page .locator('a[href="https://source.golden.example.com/report"]') .count(), ); }, ); await checkpoint( "6. Research JavaScript-rendered evidence in native Chromium", async () => { await send( "Read the live garden in a browser after its JavaScript renders and cite it.", "Browser evidence:", ); await page .getByText("23 birch trees", { exact: false }) .last() .waitFor(); assert.ok( await page .locator('a[href="https://browser.golden.example.com/report"]') .count(), ); }, ); await checkpoint( "7. Execute a real small Node.js task in Sandbox Docker", async () => { await send( "Use the shell to compute the sum of 3, 5 and 8 with Node.js.", "Shell result: 16", ); }, ); await checkpoint("8. Ask the agent to schedule a future task", async () => { const at = new Date( Math.ceil((Date.now() + 30000) / 1000) * 1000, ).toISOString(); await send( `At ${at} report my saved favorite orchard fruit.`, "Scheduled orchard reminder for", ); const href = await page .getByRole("link", { name: "View task", exact: true }) .last() .getAttribute("href"); taskId = href.split("/").at(-1); assert.match(taskId, /^[a-f0-9-]{36}$/); assert.ok(Date.parse(at) > Date.now() + 10000); }); await checkpoint( "9. Disconnect every customer client before the task is due", async () => { storage = await context.storageState(); await context.close(); context = null; await until( () => proxy.evidence.customerSockets.size, (n) => n === 0, "all customer WebSockets closed", ); disconnectedAt = { time: Date.now(), requests: proxy.evidence.customerRequests, }; }, ); await checkpoint( "10. Native scheduled work completes before any reconnect or owner read", async () => { await until( () => { if (runnerError) throw runnerError; return completion; }, Boolean, "passive native scheduled completion", 90000, ); assert.equal(completion.run.taskId, taskId); assert.equal(completion.run.conversationId, secondId); assert.equal(completion.run.source, "scheduled"); assert.equal(completion.status, "completed"); assert.ok(completion.receivedAt > disconnectedAt.time); }, ); await checkpoint( "11. Reconnect using the owner cookie issued by the real bridge", async () => { // Stop/restart only after passive completion, preserving native identity/storage. await customer.stop(); customer = await startCustomer(); context = await browser.newContext({ ignoreHTTPSErrors: true, storageState: storage, }); page = await context.newPage(); page.setDefaultTimeout(30000); await page.goto(`${customerOrigin}/conversations/${secondId}`); await page .getByText("Scheduled result: My favorite orchard fruit is quince.", { exact: false, }) .last() .waitFor(); }, ); await checkpoint( "12. Persisted conversation history, sources and exact scheduled result remain intact", async () => { const snapshot = await page.evaluate(async () => { const response = await fetch("/__golden__/snapshot"); if (!response.ok) throw new Error(`Snapshot ${response.status}`); return response.json(); }); assert.equal(snapshot.memories.length, 1); assert.equal(snapshot.tasks.length, 1); assert.equal(snapshot.runs.length, 1); assert.equal(snapshot.runs[0].id, completion.run.id); assert.equal(JSON.parse(snapshot.runs[0].payload).status, "completed"); assert.equal(snapshot.runs[0].submission_id, completion.submissionId); assert.equal(snapshot.browserLeases.length, 0); assert.equal(snapshot.shellLeases.length, 0); const first = snapshot.conversations.find((c) => c.id === firstId), second = snapshot.conversations.find((c) => c.id === secondId); assert.ok( JSON.stringify(first.messages).includes( "Saved: My favorite orchard fruit is quince.", ), ); const history = JSON.stringify(second.messages); const scheduled = second.submissions.filter( (submission) => submission.metadata?.taskRun, ); assert.equal(scheduled.length, 1); assert.equal(scheduled[0].submissionId, completion.submissionId); assert.equal(scheduled[0].status, "completed"); assert.equal(scheduled[0].metadata.taskRun.id, completion.run.id); // Native hook delivery may repeat; execution and transcript must not. assert.ok(completions.length >= 1); assert.ok( completions.every( (event) => event.submissionId === completion.submissionId && event.run.id === completion.run.id && event.run.taskId === taskId, ), ); const textOf = (message) => message.parts .filter((part) => part.type === "text") .map((part) => part.text) .join(""); const scheduledPrompt = "Report my saved favorite orchard fruit now."; const scheduledAnswer = "Scheduled result: My favorite orchard fruit is quince."; assert.equal( second.messages.filter( (message) => message.role === "user" && textOf(message) === scheduledPrompt, ).length, 1, ); assert.equal( second.messages.filter( (message) => message.role === "assistant" && textOf(message) === scheduledAnswer, ).length, 1, ); assert.ok( first.messages.every( (message) => ![scheduledPrompt, scheduledAnswer].includes(textOf(message)), ), ); assert.equal( first.submissions.filter((submission) => submission.metadata?.taskRun) .length, 0, ); for (const tool of [ "web_search", "read_url", "browser_read", "shell", "createSchedule", ]) { const activities = second.activities.filter( (activity) => activity.toolName === tool, ); assert.equal(activities.length, 1, tool); assert.equal(activities[0].status, "succeeded", tool); } for (const evidence of [ "Memory answer:", "17 cedar trees", "23 birch trees", "Shell result: 16", "Scheduled result:", "source.golden.example.com", "browser.golden.example.com", ]) assert.ok(history.includes(evidence), evidence); assert.ok( !second.messages.some( (m) => m.role === "user" && JSON.stringify(m).includes("Please remember that"), ), ); await page.goto(`${customerOrigin}/tasks/${taskId}`); await page.getByText("Completed", { exact: true }).last().waitFor(); assert.deepEqual(errors, []); assert.equal(runnerError, undefined); }, ); },);