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 { 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 }); } }, );