Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927import 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 WebSocket from "ws";import { WebSocketChatTransport } from "agents/chat/transport";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));const edit = ({ name, instructions, schedule, enabled }) => ({ name, instructions, schedule, enabled,});const instant = (ms) => new Date(Math.ceil((Date.now() + ms) / 1000) * 1000).toISOString();
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;}
test( "native unattended schedules, durable Think acceptance, lifecycle gates and repair", { timeout: 360_000 }, async (t) => { 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-execution-")); 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-execution-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/execution-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/execution-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, owner; class OwnerSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } async function connect(id) { 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, }); clients.push(client); await client.ready; if (!id) owner = client; return client; } const close = () => { for (const client of clients.splice(0)) client.close(); }; const call = (method, ...args) => owner.call(method, args); async function fixture(action, input = {}) { const response = await fetch(`${origin}/__execution/${action}`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(input), }); const result = await response.json(); if (!response.ok) throw new Error(result.error); return result; } const snapshot = (id) => fixture("conversation", { conversationId: id }); const fault = (id, value) => fixture("fault", { conversationId: id, fault: value }); async function until(fn, predicate, message, timeout = 15_000) { const deadline = Date.now() + timeout; let value; while (Date.now() < deadline) { value = await fn(); if (predicate(value)) return value; await sleep(150); } assert.fail(`${message}: ${JSON.stringify(value)}`); } const scenario = (name, body) => t.test( name, { skip: process.env.FLAREBOT_EXECUTION_CASE ? !name.includes(process.env.FLAREBOT_EXECUTION_CASE) : false, }, async () => { try { await body(); } finally { if (!clients.includes(owner)) await connect(); } }, ); try { worker = await start(); await connect(); const target = await call("createConversation", "Unattended target"); const unrelated = await call("createConversation", "Unrelated"); const input = (conversationId = target.id, overrides = {}) => ({ id: crypto.randomUUID(), conversationId, name: "Scheduled check", instructions: "second", enabled: true, schedule: { kind: "once", at: instant(5000) }, ...overrides, }); const create = (id, overrides) => call("createTask", input(id, overrides)); const runs = (id) => call("listTaskRuns", id); const complete = (id) => until( () => runs(id), (page) => page.runs[0]?.status === "completed", "Task completion", ); for (const method of [ "dispatchScheduledTask", "dispatchManualTask", "dispatchRecoveredTask", "reconcileTaskExecution", "authorizeTaskRun", "taskConversation", ]) await assert.rejects(call(method, {}), /not callable/);
await scenario( "one-off fires with every client closed and no requests until after deadline", async () => { const task = await create(); const native = (await fixture("snapshot")).schedules.find( (s) => s.payload?.taskId === task.id, ); assert.equal( task.nextRunAt, new Date(native.time * 1000).toISOString(), ); close(); await sleep(Date.parse(task.nextRunAt) - Date.now() + 8000); // Read native transcript FIRST: an owner task read must not manufacture execution. const peekAt = Date.now(); const result = await snapshot(target.id); assert.ok( result.submissions[0].completedAt < peekAt, "Native completion preceded the first request", ); assert.equal( result.messages.filter((m) => m.role === "user").length, 1, ); assert.equal(result.submissions[0].status, "completed"); assert.equal((await snapshot(unrelated.id)).messages.length, 0); await connect(); const page = await complete(task.id); assert.equal( page.runs[0].submissionId, result.submissions[0].submissionId, ); assert.equal((await call("getTask", task.id)).nextRunAt, null); await fixture("replay", native); await fixture("replay", native); assert.equal((await snapshot(target.id)).submissions.length, 1); }, );
await scenario( "full worker restart before due retains native unattended alarm", async () => { const conversation = await call( "createConversation", "Restart alarm", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(12_000) }, }); close(); await worker.stop(); worker = await start(); assert.ok( Date.now() < Date.parse(task.nextRunAt), "Worker restarted before deadline", ); await sleep(Date.parse(task.nextRunAt) - Date.now() + 8000); const peekAt = Date.now(); const result = await snapshot(conversation.id); assert.ok( result.submissions[0].completedAt < peekAt, "Native completion preceded the first request", ); assert.equal( result.messages.filter((m) => m.role === "user").length, 1, ); assert.equal(result.submissions[0].status, "completed"); await connect(); await complete(task.id); }, );
await scenario( "manual dedupe, lost acceptance reply and late receipt preserve one normal result", async () => { const conversation = await call( "createConversation", "Manual idempotency", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); const requestId = crypto.randomUUID(); await fault(conversation.id, { lostReply: true }); const first = await call("runTaskNow", task.id, requestId); const second = await call("runTaskNow", task.id, requestId); assert.equal(first.id, second.id); await complete(task.id); assert.equal( (await call("getTask", task.id)).nextRunAt, task.nextRunAt, ); await call("runTaskNow", task.id, requestId); await sleep(1500); assert.equal( (await snapshot(conversation.id)).messages.filter( (m) => m.role === "user", ).length, 1, ); await fault(conversation.id, { lateReceipt: true }); const late = await fixture("begin", { taskId: task.id, taskVersion: task.version, source: "manual", requestId: crypto.randomUUID(), scheduledFor: instant(0), }); await fixture("dispatch", late); const raw = await fixture("snapshot"); assert.equal( raw.tasks.find((t) => t.id === task.id).previousRun.status, "completed", "Final native report preceded the late pending receipt", ); await call("deleteTask", task.id, task.version); }, );
await scenario( "slow original manual acceptance retains a distinct native recovery callback", async () => { const conversation = await call( "createConversation", "Slow acceptance recovery", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); await fault(conversation.id, { delay: 14_000 }); const run = await call("runTaskNow", task.id, crypto.randomUUID()); await sleep(11_500); await call("getTask", task.id); // Read-triggered public inspection sees no acceptance yet. const state = await fixture("snapshot"); const original = state.schedules.find( (s) => s.callback === "dispatchManualTask" && s.payload?.id === run.id, ); const recovery = state.schedules.find( (s) => s.callback === "dispatchRecoveredTask" && s.payload?.id === run.id, ); assert.ok( original, "The original native callback is still executing", ); assert.ok(recovery, "Recovery has a separate native row"); assert.notEqual(recovery.id, original.id); await complete(task.id); const result = await snapshot(conversation.id); assert.equal(result.submissions.length, 1); assert.equal(result.submissions[0].submissionId, run.submissionId); assert.equal( result.messages.filter((m) => m.role === "user").length, 1, ); await call("deleteTask", task.id, task.version); }, );
await scenario( "lost observer report repaired by native inspection without another submission", async () => { const conversation = await call( "createConversation", "Lost status report", ); await fault(conversation.id, { lostReport: true }); const task = await create(conversation.id); await until( () => snapshot(conversation.id), (s) => s.submissions[0]?.status === "completed", "Native completion", ); const before = await fixture("snapshot"); assert.equal( before.tasks.find((t) => t.id === task.id).previousRun.status, "running", ); await complete(task.id); assert.equal((await snapshot(conversation.id)).submissions.length, 1); }, );
await scenario( "edit, disable and task deletion gate callbacks paused before native acceptance", async () => { for (const operation of ["edit", "disable", "delete"]) { const conversation = await call("createConversation", operation); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); await fault(conversation.id, { delay: 1800 }); await call("runTaskNow", task.id, crypto.randomUUID()); await sleep(1100); if (operation === "delete") await call("deleteTask", task.id, task.version); else await call("updateTask", task.id, task.version, { ...edit(task), ...(operation === "disable" ? { enabled: false } : { instructions: "edited" }), }); await sleep(1800); assert.equal( (await snapshot(conversation.id)).messages.length, 0, operation, ); if (operation !== "delete") await call("deleteTask", task.id, task.version + 1); } }, );
await scenario( "queued and running cancellation targets task submission only", async () => { const conversation = await call( "createConversation", "Queued cancel", ); const client = await connect(conversation.id); const transport = new WebSocketChatTransport({ agent: client }); const stream = await transport.sendMessages({ chatId: conversation.id, trigger: "submit-message", messages: [ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "recover" }], }, ], abortSignal: new AbortController().signal, }); const ordinary = (async () => { for await (const chunk of stream) { /* Native ordinary chat continues. */ } })(); ordinary.catch(() => {}); await until( () => snapshot(conversation.id), (s) => s.messages.some((m) => m.role === "user"), "Ordinary chat running", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); const run = await call("runTaskNow", task.id, crypto.randomUUID()); await until( () => snapshot(conversation.id), (s) => s.submissions.some( (r) => r.submissionId === run.submissionId && ["pending", "running"].includes(r.status), ), "Task queued", ); await call("updateTask", task.id, task.version, { ...edit(task), enabled: false, }); const result = await snapshot(conversation.id); assert.equal( result.submissions.find((r) => r.submissionId === run.submissionId) .status, "aborted", ); await ordinary; assert.match( JSON.stringify((await snapshot(conversation.id)).messages), /Reply recover complete/, ); await call("deleteTask", task.id, task.version + 1);
const runningConversation = await call( "createConversation", "Running cancel", ); const runningTask = await create(runningConversation.id, { instructions: "slow-second", schedule: { kind: "once", at: instant(180_000) }, }); const running = await call( "runTaskNow", runningTask.id, crypto.randomUUID(), ); await until( () => snapshot(runningConversation.id), (s) => s.submissions[0]?.status === "running" && s.slowSecondAttempts === 1, "Task inference running", ); await call("updateTask", runningTask.id, runningTask.version, { ...edit(runningTask), enabled: false, }); assert.equal( (await snapshot(runningConversation.id)).submissions.find( (s) => s.submissionId === running.submissionId, ).status, "aborted", ); assert.ok( ( await call("exportDiagnostics", runningConversation.id) ).runtime.events.some( (event) => event.kind === "task" && event.status === "aborted", ), ); await call("deleteTask", runningTask.id, runningTask.version + 1); }, );
await scenario( "clearing native queued submissions projects skipped without replay", async () => { const conversation = await call( "createConversation", "Clear pending submissions", ); const slow = await create(conversation.id, { instructions: "recover", schedule: { kind: "once", at: instant(180_000) }, }); const queued = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); const first = await fixture("begin", { taskId: slow.id, taskVersion: slow.version, source: "manual", requestId: crypto.randomUUID(), scheduledFor: instant(0), }); await fixture("dispatch", first); await until( () => snapshot(conversation.id), (s) => s.recoverAttempts === 1, "First native submission executing", ); const second = await fixture("begin", { taskId: queued.id, taskVersion: queued.version, source: "manual", requestId: crypto.randomUUID(), scheduledFor: instant(0), }); await fixture("dispatch", second); assert.equal( (await snapshot(conversation.id)).submissions.find( (s) => s.submissionId === second.submissionId, ).status, "pending", ); await fixture("clear", { conversationId: conversation.id }); const page = await runs(queued.id); assert.equal(page.runs[0].status, "skipped"); assert.equal(page.runs[0].failureCode, "turn_skipped"); await fixture("dispatch", second); assert.equal((await snapshot(conversation.id)).messages.length, 0); await call("deleteTask", slow.id, slow.version); await call("deleteTask", queued.id, queued.version); }, );
await scenario( "running observer gate prevents queued-start message application after edit", async () => { const conversation = await call("createConversation", "Running gate"); await fault(conversation.id, { runningDelay: 2000 }); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(180_000) }, }); await call("runTaskNow", task.id, crypto.randomUUID()); await until( () => snapshot(conversation.id), (s) => s.submissions[0]?.status === "running", "Claimed queued submission", ); await call("updateTask", task.id, task.version, { ...edit(task), enabled: false, }); await sleep(2300); assert.equal((await snapshot(conversation.id)).messages.length, 0); await call("deleteTask", task.id, task.version + 1); }, );
await scenario( "native inference failure is safe terminal history", async () => { const conversation = await call( "createConversation", "Inference failure", ); const task = await create(conversation.id, { instructions: "error" }); const page = await until( () => runs(task.id), (page) => page.runs[0]?.status === "error", "Native terminal error", ); assert.equal(page.runs[0].failureCode, "turn_failed"); assert.ok( ( await call("exportDiagnostics", conversation.id) ).runtime.events.some( (event) => event.kind === "task" && event.status === "error", ), ); assert.doesNotMatch( JSON.stringify(page), /PRIVATE-ERROR|Fixture model|unavailable/, ); }, );
await scenario( "scheduled turns keep shared instructions, memory and native tools", async () => { await call("updateInstructions", "Be brief. EXECUTION-INSTRUCTIONS."); const memory = await call( "addMemory", "executionsentinel is a saved explicit fact.", ); const contextConversation = await call( "createConversation", "Scheduled context", ); const toolConversation = await call( "createConversation", "Scheduled tool", ); await fault(contextConversation.id, { fail: 1 }); const contextTask = await create(contextConversation.id, { instructions: "memory-context executionsentinel", }); const toolTask = await create(toolConversation.id, { instructions: "tool", }); await complete(contextTask.id); await complete(toolTask.id); const context = JSON.stringify( (await snapshot(contextConversation.id)).messages, ); assert.match(context, /EXECUTION-INSTRUCTIONS/); assert.match(context, /executionsentinel is a saved explicit fact/); const tool = JSON.stringify( (await snapshot(toolConversation.id)).messages, ); assert.match(tool, /tool-fixtureEcho/); assert.match(tool, /output-available/); await call("deleteMemory", memory.id, memory.version); await call("resetInstructions"); }, );
await scenario( "scheduled turns ignore the chat override and retain their default snapshot", async () => { const defaults = await call("getModelSettings"); const conversation = await call( "createConversation", "Scheduled model isolation", ); const client = await connect(conversation.id); const selection = { provider: "workers-ai", model: "thinkingmachines/inkling-256k", effort: "high", }; const transport = new WebSocketChatTransport({ agent: client }); const stream = await transport.sendMessages({ chatId: conversation.id, trigger: "submit-message", messages: [ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "configuration" }], }, ], body: { modelConfiguration: selection }, abortSignal: new AbortController().signal, }); for await (const _ of stream) { /* Persist the real chat override. */ } const taskDefault = { ...selection, effort: "low" }; await call("updateModelSettings", taskDefault); const task = await create(conversation.id, { instructions: "tool", schedule: { kind: "once", at: instant(180_000) }, }); const run = await call("runTaskNow", task.id, crypto.randomUUID()); await complete(task.id); const calls = await client.call("fixtureModelCalls"); assert.deepEqual(calls[0].configuration, selection); assert.ok(calls.length >= 3); for (const call of calls.slice(1)) { assert.deepEqual(call.configuration, taskDefault); assert.deepEqual(call.options.anthropic, { effort: "low" }); } assert.deepEqual( (await client.call("getConversationModelSettings")).override, selection, ); const submission = (await snapshot(conversation.id)).submissions.find( (submission) => submission.submissionId === run.submissionId, ); assert.deepEqual(submission.metadata.modelConfiguration, taskDefault); const followUp = await transport.sendMessages({ chatId: conversation.id, trigger: "submit-message", messages: [ ...(await snapshot(conversation.id)).messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "configuration" }], }, ], body: { modelConfiguration: selection }, abortSignal: new AbortController().signal, }); for await (const _ of followUp) { /* Scheduled metadata must not leak to chat. */ } assert.deepEqual( (await client.call("fixtureModelCalls")).at(-1).configuration, selection, ); await call("updateModelSettings", defaults.configuration); await call("deleteTask", task.id, task.version); }, );
await scenario( "accepted scheduled turn reaches native recovery or safe terminal error after restart", async () => { const conversation = await call( "createConversation", "Scheduled recovery", ); const task = await create(conversation.id, { instructions: "recover", schedule: { kind: "once", at: instant(2000) }, }); const before = await until( () => snapshot(conversation.id), (s) => s.submissions[0]?.status === "running" && s.recoverAttempts === 1, "Accepted turn before restart", ); await sleep(800); // Native buffered stream checkpoint. close(); await worker.stop(); worker = await start(); await connect(); // Recovery is explicitly woken here; the separate alarm tests prove // unattended deadlines. This test checks native accepted-turn continuity. const recovered = await until( () => snapshot(conversation.id), (s) => ["completed", "error"].includes(s.submissions[0]?.status), "Native scheduled turn recovery", 20_000, ); assert.equal( recovered.submissions[0].submissionId, before.submissions[0].submissionId, ); assert.equal( recovered.messages.filter( (m) => m.role === "user" && m.parts.some((p) => p.type === "text" && p.text === "recover"), ).length, 1, ); t.diagnostic( `Accepted scheduled restart native outcome: ${recovered.submissions[0].status}`, ); if (recovered.submissions[0].status === "completed") assert.match(JSON.stringify(recovered.messages), /Recovered/); const history = await runs(task.id); assert.equal(history.runs[0].status, recovered.submissions[0].status); assert.equal(history.runs.length, 1); assert.equal( history.runs[0].submissionId, before.submissions[0].submissionId, ); }, );
await scenario( "parent restart after placeholder recovers same unaccepted occurrence unattended", async () => { const conversation = await call( "createConversation", "Acceptance crash gap", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(3000) }, }); const run = await fixture("begin", { taskId: task.id, taskVersion: task.version, source: "scheduled", scheduledFor: task.schedule.at, }); const native = (await fixture("snapshot")).schedules.find( (s) => s.payload?.taskId === task.id, ); await fixture("lose-binding", { id: native.id }); close(); await worker.stop(); worker = await start(); await sleep(55_000); // Includes an early (<10s) maintenance tick and its next 30s tick. const peekAt = Date.now(); const result = await snapshot(conversation.id); await connect(); assert.ok( result.submissions[0].completedAt < peekAt, "Recovery completed before the first request", ); assert.equal(result.submissions.length, 1); assert.equal(result.submissions[0].submissionId, run.submissionId); assert.equal(result.submissions[0].status, "completed"); await complete(task.id); }, );
await scenario( "conversation deletion and restart erase cleanup references without facet resurrection", async () => { const conversation = await call( "createConversation", "Deleted target", ); const task = await create(conversation.id, { schedule: { kind: "once", at: instant(4000) }, }); const native = (await fixture("snapshot")).schedules.find( (s) => s.payload?.taskId === task.id, ); await call("deleteConversation", conversation.id); close(); await worker.stop(); worker = await start(); await sleep(5000); await fixture("replay", native); const state = await fixture("snapshot"); assert.ok(!state.facets.some((f) => f.name === conversation.id)); assert.ok(!state.tasks.some((t) => t.id === task.id)); assert.ok( !state.runs.some((r) => r.conversation_id === conversation.id), ); await connect(); }, );
// Isolate real UTC cron callbacks from the short race-test timing windows. const recurringConversation = await call( "createConversation", "Real recurrence", ); const recurring = await create(recurringConversation.id, { schedule: { kind: "cron", expression: "* * * * *", timezone: "UTC" }, }); const recurringDue = Date.parse(recurring.nextRunAt); const failingConversation = await call( "createConversation", "Recurring dispatch failure", ); await fault(failingConversation.id, { failOccurrence: true }); const failingRecurring = await create(failingConversation.id, { schedule: { kind: "cron", expression: "* * * * *", timezone: "UTC" }, });
await scenario( "two real UTC cron occurrences and next recurrence after exhausted dispatch failure", async () => { // No requests wake either conversation during the remaining alarm window. close(); await sleep(Math.max(0, recurringDue + 60_000 + 8000 - Date.now())); const state = await snapshot(recurringConversation.id); assert.ok( state.submissions.filter((s) => s.status === "completed").length >= 2, ); assert.equal( state.messages.filter((m) => m.role === "user").length, state.submissions.length, ); await connect(); const page = await runs(recurring.id); assert.ok(new Set(page.runs.map((r) => r.scheduledFor)).size >= 2); const failurePage = await runs(failingRecurring.id); assert.ok(failurePage.runs.some((r) => r.status === "completed")); assert.ok( failurePage.runs.some((r) => r.status === "dispatch_error"), "Exhausted native dispatch remains visible beside the next successful recurrence", ); const diagnosticRuns = ( await call("exportDiagnostics", failingConversation.id) ).runtime.events.filter((event) => event.kind === "task"); assert.ok( diagnosticRuns.some((event) => event.status === "dispatch_error"), ); assert.ok( diagnosticRuns.some((event) => event.status === "completed"), ); const failureState = await snapshot(failingConversation.id); assert.ok( failureState.submissions.every((s) => s.status === "completed"), ); }, );
for (const task of await call("listTasks")) await call("deleteTask", task.id, task.version); await fixture("repair"); assert.ok( !(await fixture("snapshot")).schedules.some((s) => [ "dispatchScheduledTask", "dispatchManualTask", "dispatchRecoveredTask", "reconcileTaskExecution", ].includes(s.callback), ), "No orphan task maintenance alarm", ); } finally { close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);