Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674import 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 { AgentClient } from "agents/client";import { WebSocketChatTransport } from "agents/chat/transport";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";import { MAX_TASKS, MAX_TASK_INSTRUCTIONS } from "../shared/tasks.ts";
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;}
const edit = ({ name, instructions, schedule, enabled }) => ({ name, instructions, schedule, enabled,});const projection = (run, status, extra = {}) => ({ id: run.id, taskVersion: run.taskVersion, conversationId: run.conversationId, submissionId: run.submissionId, status, ...extra,});
test( "durable scheduled task definitions, owner CRUD, UTC forecasts and internal history", { timeout: 120_000 }, async () => { const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const cookie = ( await createOwnerSession( new Secret(customerBindings.FLAREBOT_SESSION_SECRET), { ...installation, runtimeOrigin: origin }, ) ).split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; const persistence = await mkdtemp(join(tmpdir(), "flarebot-tasks-")); const config = JSON.parse( await readFile("dist/release/deployment.json", "utf8"), ); const configPath = join(persistence, "wrangler.json"); await writeFile( configPath, JSON.stringify({ ...config, name: "flarebot-task-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/task-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/task-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: persistence, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }); const clients = []; let worker; try { worker = await start(); 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); await client.ready; return { client, states }; } let owner = await connect(); const call = (name, ...args) => owner.client.call(name, args); const fixture = async (path, body) => { const response = await fetch( `${origin}/__tasks/${path}`, body === undefined ? {} : { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(body), }, ); const result = await response.json(); if (!response.ok) throw new Error(result.error); return result; }; const a = await call("createConversation", "Task target"); const b = await call("createConversation", "Unrelated conversation"); const input = (overrides = {}) => ({ id: crypto.randomUUID(), conversationId: a.id, name: "Private scheduled work", instructions: "TASK-PRIVATE-CONTENT: review the morning report.", schedule: { kind: "once", at: "2096-02-29T11:30:00+02:00" }, enabled: true, ...overrides, }); assert.deepEqual(await call("listTasks"), []); for (const method of [ "beginTaskRun", "projectTaskRun", "fixtureRuns", "inspectTasks", ]) await assert.rejects(call(method, {}), /not callable/); const utc = await fixture("utc"); assert.deepEqual(utc, { offset: 0, monday: "2030-01-07T09:00:00.000Z", leap: "2032-02-29T09:00:00.000Z", });
for (const invalid of [ null, [], {}, input({ unexpected: true }), input({ id: "../bad" }), input({ conversationId: crypto.randomUUID() }), input({ enabled: "true" }), input({ name: "\n" }), input({ name: "x".repeat(121) }), input({ instructions: "\u0000bad" }), input({ instructions: "x".repeat(MAX_TASK_INSTRUCTIONS + 1) }), ]) await assert.rejects(call("createTask", invalid)); for (const at of [ "2096-02-30T09:00:00Z", "2095-02-29T09:00:00Z", "2096-04-31T09:00:00Z", "2096-01-01T24:00:00Z", "2096-01-01T09:00:00", "2096-01-01", "2096-01-01T09:00:00.001Z", "2096-01-01T09:00:00+14:30", "2096-01-01T09:00:00-00:00", "2000-01-01T09:00:00Z", null, 1e99, ]) await assert.rejects( call("createTask", input({ schedule: { kind: "once", at } })), ); for (const expression of [ "* * * * * *", "@daily", "0x 9 * * 1", "1/2 9 * * 1", "0 25 * * *", "0 9 30 2 *", "*/0 * * * *", "0 9 * * monday", "0 9 * * 1,", "x".repeat(121), ]) await assert.rejects( call( "createTask", input({ schedule: { kind: "cron", expression, timezone: "UTC" } }), ), ); for (const schedule of [ { kind: "cron", expression: "0 9 * * 1", timezone: "Europe/Warsaw" }, { kind: "cron", expression: "0 9 * * 1" }, { kind: "once", at: "2096-01-01T09:00:00Z", timezone: "UTC" }, ]) await assert.rejects(call("createTask", input({ schedule })));
const original = input(); let once = await call("createTask", original); assert.equal(once.schedule.at, "2096-02-29T09:30:00.000Z"); assert.equal(once.nextRunAt, once.schedule.at); assert.equal(once.previousRun, null); assert.equal(once.version, 1); assert.deepEqual( await call("createTask", { ...original, name: ` ${original.name} `, schedule: { kind: "once", at: once.schedule.at }, }), once, ); await assert.rejects( call("createTask", { ...original, instructions: "Different initial request", }), /already used/, ); once = await call("updateTask", once.id, once.version, { ...edit(once), instructions: "TASK-EDITED-CONTENT", enabled: false, }); assert.equal(once.nextRunAt, null); assert.deepEqual( await call("createTask", original), once, "replay compares original identity after edits", ); await assert.rejects( call("updateTask", once.id, 1, edit(once)), /changed/, ); await assert.rejects(call("deleteTask", once.id, 1), /changed/); for (const version of [ 0, -1, 1.2, "2", null, Number.MAX_SAFE_INTEGER + 1, ]) await assert.rejects( call("updateTask", once.id, version, edit(once)), /version/, ); await assert.rejects( call("updateTask", once.id, once.version, { ...edit(once), conversationId: b.id, }), /input/, ); once = await call("updateTask", once.id, once.version, { ...edit(once), enabled: true, }); assert.equal(once.nextRunAt, once.schedule.at);
const cronInput = input({ name: "Monday UTC", schedule: { kind: "cron", expression: " 00\t09 * * 01 ", timezone: "UTC", }, }); let cron = await call("createTask", cronInput); assert.equal(cron.schedule.expression, "0 9 * * 1"); assert.equal(new Date(cron.nextRunAt).getUTCDay(), 1); assert.equal(new Date(cron.nextRunAt).getUTCHours(), 9); assert.equal( cron.previousRun, null, "cron's prior calendar date is not history", ); const runInput = (task, overrides = {}) => ({ taskId: task.id, taskVersion: task.version, source: "manual", scheduledFor: "2030-01-01T00:00:00Z", requestId: crypto.randomUUID(), ...overrides, }); const manualInput = runInput(once); const [manual] = await fixture("begin", [manualInput]); assert.equal( (await call("getTask", once.id)).nextRunAt, once.schedule.at, "manual runs do not consume the scheduled one-off", ); assert.deepEqual(await fixture("begin", [manualInput]), [manual]); await fixture("fail", { statement: "UPDATE flarebot_tasks SET once_consumed_at", }); await assert.rejects( fixture("begin", [ runInput(once, { source: "scheduled", scheduledFor: once.schedule.at, }), ]), /Fixture task SQL failure/, ); assert.equal( (await call("getTask", once.id)).nextRunAt, once.schedule.at, ); assert.equal( (await call("listTaskRuns", once.id)).runs.length, 1, "failed once consumption rolls back its inserted occurrence", ); const [scheduled] = await fixture("begin", [ runInput(once, { source: "scheduled", scheduledFor: once.schedule.at }), ]); assert.equal((await call("getTask", once.id)).nextRunAt, null); const complete = await fixture( "project", projection(scheduled, "completed", { startedAt: "2030-01-01T00:00:00Z", completedAt: "2030-01-01T00:00:01Z", }), ); assert.deepEqual( await fixture("project", projection(scheduled, "pending")), complete, "late receipt never regresses completion", ); assert.equal( await fixture("project", { ...projection(manual, "completed"), submissionId: "wrong", }), null, ); once = await call("updateTask", once.id, once.version, { ...edit(once), name: "Completed task renamed", }); assert.equal( once.nextRunAt, null, "metadata edits cannot rearm a consumed one-off", ); assert.equal( (await fixture("project", projection(manual, "aborted"))).failureCode, "turn_aborted", "existing occurrence may report after task edits", ); once = await call("updateTask", once.id, once.version, { ...edit(once), enabled: false, }); await assert.rejects( call("updateTask", once.id, once.version, { ...edit(once), enabled: true, }), /new future instant/, );
const batch = Array.from({ length: 107 }, (_, index) => runInput( cron, index % 2 ? {} : { source: "scheduled", scheduledFor: new Date( Date.parse("2030-01-07T09:00:00Z") + index * 7 * 86400000, ).toISOString(), }, ), ); const runs = await fixture("begin", batch); const errorRun = await fixture( "project", projection(runs[0], "dispatch_error"), ); assert.equal(errorRun.failureCode, "dispatch_failed"); assert.equal( (await fixture("project", projection(runs[0], "running"))).failureCode, null, "uncertain dispatch is repairable", ); assert.equal( (await fixture("project", projection(runs[0], "pending"))).status, "running", ); await fixture("project", projection(runs[0], "error")); const ordered = runs.toSorted((a, b) => a.createdAt < b.createdAt ? 1 : a.createdAt > b.createdAt ? -1 : a.id < b.id ? 1 : -1, ); const seen = []; let before; do { const page = await call("listTaskRuns", cron.id, { limit: 17, ...(before ? { before } : {}), }); seen.push(...page.runs); before = page.nextCursor; } while (before); assert.deepEqual( seen.map((run) => run.id), ordered.map((run) => run.id), ); assert.equal( (await call("getTask", cron.id)).previousRun.id, ordered[0].id, ); assert.equal( (await call("listTaskRuns", cron.id, { limit: 100 })).runs.length, 100, ); for (const options of [ { limit: 101 }, { limit: 0 }, { limit: null }, { limit: 1.1 }, { limit: "3" }, { before: {} }, { before: { createdAt: "2030-02-30T00:00:00.000Z", id: "x" } }, { inject: true }, ]) await assert.rejects(call("listTaskRuns", cron.id, options));
// CAS conflicts are observable on the real owner RPC, even from two clients. const observer = await connect(); const concurrent = await Promise.allSettled([ call("updateTask", cron.id, cron.version, { ...edit(cron), name: "First edit", }), observer.client.call("updateTask", [ cron.id, cron.version, { ...edit(cron), name: "Second edit" }, ]), ]); assert.equal( concurrent.filter((result) => result.status === "fulfilled").length, 1, ); cron = await call("getTask", cron.id); assert.equal(cron.version, 2); assert.doesNotMatch( JSON.stringify(observer.states), /TASK-PRIVATE|TASK-EDITED|scheduled work/, ); assert.deepEqual( await fixture("schedules"), [], "model history seeding does not arm native execution", ); const transcript = async (id) => ( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${id}/get-messages`, { headers }, ) ).json(); assert.deepEqual(await transcript(a.id), []); assert.deepEqual(await transcript(b.id), []);
const conversationClient = await connect(b.id); const transport = new WebSocketChatTransport({ agent: conversationClient.client, }); const stream = await transport.sendMessages({ chatId: b.id, trigger: "submit-message", messages: [ { id: crypto.randomUUID(), role: "user", parts: [ { type: "text", text: "Preserve this existing conversation." }, ], }, ], abortSignal: new AbortController().signal, }); for await (const _chunk of stream) { /* Native Think persists the normal chat. */ } const preservedTranscript = await transcript(b.id); assert.ok(preservedTranscript.length >= 2);
const doomedInput = input({ conversationId: b.id, instructions: "ERASE-TASK-CONTENT", }); const doomed = await call("createTask", doomedInput); const [pending, finished] = await fixture("begin", [ runInput(doomed), runInput(doomed), ]); await fixture("project", projection(finished, "completed")); await fixture("fail", { statement: "DELETE FROM flarebot_task_runs" }); await assert.rejects( call("deleteTask", doomed.id, doomed.version), /Fixture task SQL failure/, ); assert.equal( (await call("getTask", doomed.id)).instructions, doomed.instructions, ); assert.equal( (await call("listTaskRuns", doomed.id)).runs.length, 2, "failed deletion restores task and history together", ); await call("deleteTask", doomed.id, doomed.version); await call("deleteTask", doomed.id, doomed.version); await assert.rejects(call("createTask", doomedInput), /deleted/); await assert.rejects(call("getTask", doomed.id), /deleted/); await assert.rejects(call("listTaskRuns", doomed.id), /deleted/); assert.equal( await fixture("project", projection(pending, "completed")), null, ); assert.ok((await call("listConversations")).some((c) => c.id === b.id)); assert.deepEqual( await transcript(b.id), preservedTranscript, "task deletion retains its shared conversation transcript", ); const storage = await fixture("inspect"); assert.doesNotMatch(JSON.stringify(storage), /ERASE-TASK-CONTENT/); assert.ok( !storage.runs.some((run) => run.id === finished.id), "terminal history is fully erased", ); assert.deepEqual( storage.tasks.find((task) => task.id === doomed.id), { id: doomed.id, definition: null, creation_hash: null, conversation_id: null, once_consumed_at: null, }, ); assert.equal( storage.runs.find((run) => run.id === pending.id).payload, null, ); assert.equal( storage.runs.find((run) => run.id === pending.id).submission_id, pending.submissionId, "private cleanup ID survives content erasure", ); const reserved = input(); await call("deleteTask", reserved.id, 1); await assert.rejects(call("createTask", reserved), /deleted/);
const racing = input(); await Promise.allSettled([ call("createTask", racing), observer.client.call("deleteTask", [racing.id, 1]), ]); await assert.rejects(call("getTask", racing.id), /deleted/); await assert.rejects(call("createTask", racing), /deleted/); const creating = await ( await fetch(`${origin}/__fixture/interrupted-create`) ).json(); await assert.rejects( call("createTask", input({ conversationId: creating })), /Conversation not found/, );
// Conversation deletion gates tasks before native cleanup, including failure. const deletingInput = input({ conversationId: b.id, instructions: "ERASE-CONVERSATION-TASK", }); const deleting = await call("createTask", deletingInput); const [late] = await fixture("begin", [runInput(deleting)]); await fetch(`${origin}/__fixture/fail-delete`); await assert.rejects( call("deleteConversation", b.id), /Fixture deletion failure/, ); await assert.rejects(call("getTask", deleting.id), /deleted/); await assert.rejects( call("createTask", input({ conversationId: b.id })), /Conversation not found/, ); assert.equal( await fixture("project", projection(late, "completed")), null, ); assert.doesNotMatch( JSON.stringify(await fixture("inspect")), /ERASE-CONVERSATION-TASK/, );
const beforeRestart = await call("getTask", once.id); const historyBeforeRestart = await call("listTaskRuns", cron.id, { limit: 100, }); for (const client of clients) client.close(); await worker.stop(); worker = await start(); owner = await connect(); assert.deepEqual(await call("getTask", once.id), beforeRestart); assert.deepEqual( await call("listTaskRuns", cron.id, { limit: 100 }), historyBeforeRestart, ); assert.deepEqual(await call("createTask", original), beforeRestart); await assert.rejects(call("createTask", doomedInput), /deleted/); await assert.rejects(call("createTask", deletingInput), /deleted/); const facets = await (await fetch(`${origin}/__fixture/inspect`)).json(); assert.ok( !facets.facets.includes(b.id), "startup finishes deleting the native facet without recreating it", ); assert.equal(await fixture("project", projection(late, "running")), null); once = await call("updateTask", once.id, beforeRestart.version, { ...edit(beforeRestart), enabled: true, schedule: { kind: "once", at: "2097-01-01T00:00:00Z" }, }); assert.equal(once.nextRunAt, "2097-01-01T00:00:00.000Z"); for (let i = (await call("listTasks")).length; i < MAX_TASKS; i++) await call("createTask", input({ enabled: false, name: `Task ${i}` })); await assert.rejects(call("createTask", input()), /Task limit reached/); assert.equal((await call("listTasks")).length, MAX_TASKS); assert.equal( (await fetch(`${origin}/agents/personal-agent/personal/status`)).status, 401, ); assert.equal( ( await fetch(`${origin}/agents/personal-agent/personal/status`, { headers: { ...headers, Origin: "https://foreign.invalid" }, }) ).status, 403, ); assert.doesNotMatch( await (await fetch(`${origin}/settings`)).text(), /TASK-PRIVATE|TASK-EDITED|Completed task renamed/, ); } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);