Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574import * as Effect from "effect/Effect";import assert from "node:assert/strict";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 { test } from "node:test";import { chromium } from "playwright";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";
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));async function until(read, predicate, label) { const deadline = Date.now() + 20_000; while (Date.now() < deadline) { const value = await read(); if (predicate(value)) return value; await sleep(25); } assert.fail(`Timed out: ${label}`);}const editable = ({ name, instructions, schedule, enabled }) => ({ name, instructions, schedule, enabled,});const request = (name) => ({ name, instructions: "Research important Cloudflare releases and summarize the sources.", schedule: { kind: "cron", expression: "0 9 * * 1", timezone: "UTC" },});
test( "conversation scheduling uses native actions, durable tasks and the real task editor", { timeout: 180_000 }, async (t) => { 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)); const origin = `http://127.0.0.1:${port}`; const temporary = await mkdtemp( join(tmpdir(), "flarebot-schedule-action-"), ); 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-schedule-action-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/schedule-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 }; const start = () => unstable_dev("tests/fixtures/schedule-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 }, }); const clients = []; let worker, browser, owner, connection, first, second, saved; class OwnerSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } async function connect(id) { const states = []; const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", ...(id ? { basePath: `agents/personal-agent/personal/sub/conversation/${id}`, } : { name: "personal" }), WebSocket: OwnerSocket, onStateUpdate: (state) => states.push(state), }); clients.push(client); const transport = new WebSocketChatTransport({ agent: client }); client.addEventListener("message", (event) => { const frame = JSON.parse(event.data); if (frame.type === MessageType.CF_AGENT_STREAM_RESUMING) transport.handleStreamResuming(frame); if (frame.type === MessageType.CF_AGENT_STREAM_RESUME_NONE) transport.handleStreamResumeNone(frame); if (frame.type === MessageType.CF_AGENT_STREAM_PENDING) transport.handleStreamPending(); }); await client.ready; return { client, transport, states }; } const call = (method, ...args) => owner.client.call(method, args); const list = () => call("listTasks"); const history = async (id) => ( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${id}/get-messages`, { headers }, ) ).json(); async function fixture(namespace, action, input = {}) { const response = await fetch(`${origin}/__${namespace}/${action}`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(input), }); const value = await response.json(); assert.equal(response.ok, true, JSON.stringify(value)); return value; } const fault = (mode) => fixture("schedule", "fault", { mode }); const snapshot = () => fixture("execution", "snapshot"); async function send(text, target = connection, conversationId = first.id) { const stream = await target.transport.sendMessages({ chatId: conversationId, trigger: "submit-message", messages: [ ...(await history(conversationId)), { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ], abortSignal: new AbortController().signal, }); const chunks = []; const done = (async () => { for await (const chunk of stream) chunks.push(chunk); })(); done.catch(() => {}); return { chunks, done }; } async function action( input, id = crypto.randomUUID(), target = connection, conversationId = first.id, ) { const turn = await send( "schedule:" + JSON.stringify({ tool: "createSchedule", input, id }), target, conversationId, ); await turn.done; const part = (await history(conversationId)) .flatMap((m) => m.parts) .findLast( (p) => p.toolCallId === id && ("output" in p || "errorText" in p), ); return { output: part?.output, part, chunks: turn.chunks, id }; } try { worker = await start(); owner = await connect(); first = await call("createConversation", "Cloudflare research"); second = await call("createConversation", "Unrelated conversation"); connection = await connect(first.id);
await t.test( "native model selects the action; saved task opens and edits in packaged Kumo UI", async () => { for (const method of [ "createTaskForConversation", "authorizeTaskRun", ]) await assert.rejects(call(method), /not callable/); const result = await action(request("Monday Cloudflare releases")); saved = result.output; assert.equal(saved.conversationId, first.id); assert.equal(saved.status, "scheduled"); assert.equal(saved.url, `/tasks/${saved.id}`); assert.deepEqual(saved.schedule, request("").schedule); assert.ok(Date.parse(saved.nextRunAt) > Date.now()); const native = (await snapshot()).schedules.filter( (s) => s.callback === "dispatchScheduledTask" && s.payload.taskId === saved.id, ); assert.equal(native.length, 1); assert.equal( new Date(native[0].time * 1000).toISOString(), saved.nextRunAt, ); assert.equal((await history(second.id)).length, 0); const activity = ( await connection.client.call("listToolActivities") ).activities.find((a) => a.toolCallId === result.id); assert.equal(activity.kind, "schedule"); assert.equal(activity.status, "succeeded"); assert.ok( result.chunks.some((c) => c.type === "tool-output-available"), ); browser = await chromium.launch({ headless: true }); const context = await browser.newContext(); await context.addCookies([ { name: cookie.split("=")[0], value: cookie.slice(cookie.indexOf("=") + 1), url: origin.replace("http:", "https:"), secure: true, httpOnly: true, sameSite: "Strict", }, ]); const page = await context.newPage(); const errors = []; page.on("pageerror", (e) => errors.push(e.message)); await page.goto(`${origin}/conversations/${first.id}`); const viewTask = page .locator(".chat-tool") .getByRole("link", { name: "View scheduled task", exact: true }); assert.equal(await viewTask.getAttribute("href"), saved.url); assert.equal( await viewTask.getAttribute("data-kumo-component"), "LinkButton", ); await viewTask.click(); await page.waitForURL(origin + saved.url); await page .getByRole("heading", { name: saved.name, exact: true }) .waitFor(); await page .getByRole("button", { name: "Edit task", exact: true }) .click(); const dialog = page.getByRole("dialog"); await dialog .getByRole("textbox", { name: "Task name", exact: true }) .fill("Tuesday Cloudflare releases"); await dialog .getByRole("textbox", { name: "Cron expression (UTC)", exact: true, }) .fill("30 10 * * 2"); await dialog .getByRole("button", { name: "Save task", exact: true }) .click(); await dialog.waitFor({ state: "hidden" }); saved = await call("getTask", saved.id); assert.equal(saved.version, 2); assert.equal(saved.name, "Tuesday Cloudflare releases"); assert.equal(saved.schedule.expression, "30 10 * * 2"); assert.equal(saved.conversationId, first.id); await page.reload(); await page .getByRole("heading", { name: saved.name, exact: true }) .waitFor(); assert.deepEqual(errors, []); await browser.close(); browser = undefined; }, );
await t.test( "one-off normalization and rejected unresolved or invalid schedules", async () => { const instant = new Date( Math.ceil(Date.now() / 1000) * 1000 + 86_400_000, ).toISOString(); const once = await action({ ...request("Tomorrow research"), schedule: { kind: "once", at: instant.replace(".000Z", "+00:00") }, }); assert.equal(once.output.schedule.at, instant); const count = (await list()).length; for (const schedule of [ { kind: "cron", expression: "0 9 * * 1", timezone: "Europe/Warsaw", }, { kind: "cron", expression: "0 9 * * 1" }, { kind: "cron", expression: "every Monday morning", timezone: "UTC", }, { kind: "cron", expression: "0garbage 9 * * 1", timezone: "UTC" }, { kind: "once", at: "2030-01-01T09:00:00" }, { kind: "once", at: "2030-02-30T09:00:00Z" }, ]) { const result = await action({ ...request("Invalid schedule"), schedule, }); assert.ok( result.output?.error || result.part?.errorText, JSON.stringify(result), ); } const forged = await action({ ...request("Wrong conversation"), conversationId: second.id, }); assert.ok(forged.part?.errorText || forged.output?.error); assert.equal((await list()).length, count); }, );
let replayId, replayInput, replayTask; await t.test( "lost accepted reply replays the same task after edits; native ledger settles once", async () => { await fault("lost-reply"); replayId = crypto.randomUUID(); replayInput = request("Accepted before lost reply"); assert.ok((await action(replayInput, replayId)).output.error); replayTask = (await list()).find( (task) => task.name === replayInput.name, ); assert.ok(replayTask); replayTask = await call( "updateTask", replayTask.id, replayTask.version, { ...editable(replayTask), name: "Edited after accepted reply", schedule: { kind: "cron", expression: "15 11 * * 3", timezone: "UTC", }, }, ); // Leave the native action unsettled across a full Worker restart. for (const client of clients) client.close(); await worker.stop(); worker = await start(); owner = await connect(); connection = await connect(first.id); assert.equal( (await history(first.id)) .flatMap((m) => m.parts) .filter((p) => p.toolCallId === replayId && p.output?.error) .length, 1, ); const retry = await action(replayInput, replayId); assert.equal(retry.output.id, replayTask.id); assert.equal(retry.output.version, 2); assert.equal(retry.output.name, replayTask.name); assert.deepEqual(retry.output.schedule, replayTask.schedule); const settled = await action(replayInput, replayId); assert.deepEqual(settled.output, retry.output); assert.equal( (await list()).filter((task) => task.id === replayTask.id).length, 1, ); assert.equal( (await snapshot()).schedules.filter( (s) => s.callback === "dispatchScheduledTask" && s.payload.taskId === replayTask.id, ).length, 1, ); const conflict = await action( { ...replayInput, name: "Conflicting key reuse" }, replayId, ); assert.ok(conflict.output.error); assert.equal( (await list()).filter( (task) => task.name === "Conflicting key reuse", ).length, 0, ); }, );
await t.test( "deleted accepted task cannot be resurrected by an unsettled replay", async () => { await fault("lost-reply"); const id = crypto.randomUUID(), input = request("Deleted before replay"); assert.ok((await action(input, id)).output.error); const task = (await list()).find((task) => task.name === input.name); await call("deleteTask", task.id, task.version); assert.ok((await action(input, id)).output.error); assert.ok(!(await list()).some((item) => item.id === task.id)); }, );
await t.test( "native binding failure returns honest saved-but-unavailable status", async () => { await fault("unavailable"); const result = await action(request("Native schedule unavailable")); assert.equal(result.output.status, "unavailable"); assert.equal(result.output.nextRunAt, null); assert.equal(result.output.schedulingError, "schedule_unavailable"); assert.equal(result.output.url, `/tasks/${result.output.id}`); const activities = ( await connection.client.call("listToolActivities") ).activities; assert.equal( activities.find((a) => a.toolCallId === result.id).status, "failed", ); assert.doesNotMatch( JSON.stringify(connection.states), /PRIVATE-SCHEDULE-ERROR|Research important|Accepted before/, ); await fault(""); assert.ok((await call("getTask", result.output.id)).nextRunAt); }, );
await t.test( "stop and clear during parent acquisition prevent a later write", async () => { for (const mode of ["stop", "clear"]) { await fixture("schedule", "delay-acquisition", { conversationId: first.id, }); const id = crypto.randomUUID(), input = request(`Interrupted ${mode}`); const turn = await send( "schedule:" + JSON.stringify({ tool: "createSchedule", input, id }), ); await until( () => fixture("schedule", "acquiring", { conversationId: first.id }), Boolean, "parent acquisition", ); if (mode === "stop") connection.transport.cancelActiveServerTurn(); else await fixture("execution", "clear", { conversationId: first.id }); await turn.done.catch(() => {}); await sleep(1800); assert.ok(!(await list()).some((task) => task.name === input.name)); } }, );
await t.test( "deletion during acquisition or before parent mutation cannot restore a task or facet", async () => { for (const mode of ["acquisition", "late-parent"]) { const conversation = await call( "createConversation", `Delete during ${mode}`, ); const target = await connect(conversation.id); if (mode === "acquisition") await fixture("schedule", "delay-acquisition", { conversationId: conversation.id, }); else await fault("late-parent"); const id = crypto.randomUUID(), input = request(`Deleted ${mode}`); const turn = await send( "schedule:" + JSON.stringify({ tool: "createSchedule", input, id }), target, conversation.id, ); await until( () => fixture( "schedule", mode === "acquisition" ? "acquiring" : "waiting", { conversationId: conversation.id }, ), Boolean, mode, ); await call("deleteConversation", conversation.id); target.client.close(); await sleep(1800); assert.ok( !(await list()).some( (task) => task.conversationId === conversation.id, ), ); assert.ok( !(await snapshot()).facets.some( (facet) => facet.name === conversation.id, ), ); } }, );
await t.test( "complete native instructions retain settings, other tools, clarification and current UTC", async () => { await call("updateInstructions", "CUSTOM-SCHEDULE-INSTRUCTIONS"); await call( "addMemory", "Instructions for Cloudflare research use primary sources.", ); const before = Date.now(); const turn = await send("instructions"); await turn.done; const text = (await history(first.id)) .at(-1) .parts.filter((p) => p.type === "text") .map((p) => p.text) .join(""); assert.match(text, /CUSTOM-SCHEDULE-INSTRUCTIONS/); assert.match(text, /untrusted user data/); assert.match(text, /read_url/); assert.match(text, /shell/); assert.match(text, /Every Monday morning/); assert.match(text, /never invent 09:00/); assert.match( text, /execution instructions are not requests to create another task/, ); const instant = /Current UTC time: ([0-9TZ:.-]+)/.exec(text)?.[1]; assert.ok( Date.parse(instant) >= before && Date.parse(instant) <= Date.now(), instant, ); assert.doesNotMatch( JSON.stringify( (await connection.client.call("listToolActivities")).activities, ), /PRIVATE-SCHEDULE-ERROR|primary sources|Cloudflare releases/, ); }, ); } finally { for (const client of clients) client.close(); await browser?.close(); await worker?.stop(); await rm(temporary, { recursive: true, force: true }); } },);