Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491import type { TaskDiagnostic } from "./diagnostics";import type { Agent } from "agents";import { MAX_TASKS, type TaskDefinition, type TaskRun, type TaskRunPage, type TaskSummary,} from "../shared/tasks";import { nextCron, taskId, taskInput, taskInstant, taskRunPage, taskVersion,} from "./task-validation";
interface TaskRow { id: string; definition: string | null; creation_hash: string | null; conversation_id: string | null; once_consumed_at: string | null;}interface RunRow { payload: string;}
export interface BeginTaskRun { taskId: string; taskVersion: number; source: "scheduled" | "manual"; scheduledFor: string; /** Required for a manual occurrence; retries reuse this owner request ID. */ requestId?: string;}export interface TaskRunProjection { id: string; taskVersion: number; conversationId: string; submissionId: string; status: TaskRun["status"]; startedAt?: string | null; completedAt?: string | null;}export const terminalTaskStatuses = new Set<TaskRun["status"]>([ "completed", "aborted", "skipped", "error",]);
// Application intent and history only. Native Agents schedules and Think// submissions/transcripts own execution; this store never writes SDK tables.export class TaskStore { constructor( private readonly sql: Agent["sql"], private readonly requireConversation: (id: string) => unknown, private readonly storage: Pick<DurableObjectStorage, "transactionSync">, private readonly observe?: (run: TaskDiagnostic) => void, ) {}
private observed(run: TaskRun) { const { id, taskId, conversationId, submissionId, source, status, startedAt, completedAt, } = run; try { this.observe?.({ id, taskId, conversationId, submissionId, source, status, startedAt, completedAt, }); } catch { /* Optional diagnostic projection. */ } }
initialize() { this.sql`CREATE TABLE IF NOT EXISTS flarebot_tasks ( id TEXT PRIMARY KEY, definition TEXT, creation_hash TEXT, conversation_id TEXT, once_consumed_at TEXT )`; this.sql`CREATE INDEX IF NOT EXISTS flarebot_tasks_conversation ON flarebot_tasks(conversation_id)`; this.sql`CREATE TABLE IF NOT EXISTS flarebot_task_runs ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, conversation_id TEXT NOT NULL, submission_id TEXT NOT NULL, created_at TEXT, payload TEXT, recovery_scheduled INTEGER NOT NULL DEFAULT 0 )`; if ( !this.sql<{ name: string }>`PRAGMA table_info(flarebot_task_runs)`.some( (column) => column.name === "recovery_scheduled", ) ) this .sql`ALTER TABLE flarebot_task_runs ADD COLUMN recovery_scheduled INTEGER NOT NULL DEFAULT 0`; this.sql`CREATE INDEX IF NOT EXISTS flarebot_task_runs_history ON flarebot_task_runs(task_id, created_at DESC, id DESC)`; }
private row(id: string): TaskRow | undefined { return this.sql<TaskRow>`SELECT * FROM flarebot_tasks WHERE id = ${id}`[0]; }
private active(id: string): { row: TaskRow; task: TaskDefinition } { const row = this.row(id); if (!row?.definition) throw new Error("Task not found or was deleted"); const task = JSON.parse(row.definition) as TaskDefinition; this.requireConversation(task.conversationId); return { row, task }; }
private summary(row: TaskRow, now = Date.now()): TaskSummary { const task = JSON.parse(row.definition!) as TaskDefinition; const previous = this.sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE task_id = ${task.id} AND payload IS NOT NULL ORDER BY created_at DESC, id DESC LIMIT 1`[0]; return { ...task, previousRun: previous ? (JSON.parse(previous.payload) as TaskRun) : null, nextRunAt: !task.enabled || row.once_consumed_at ? null : task.schedule.kind === "once" ? task.schedule.at : nextCron(task.schedule.expression, now), }; }
get(id: unknown): TaskSummary { return this.summary(this.active(taskId(id)).row); }
list(): TaskSummary[] { const now = Date.now(); return this.sql<TaskRow>`SELECT * FROM flarebot_tasks WHERE definition IS NOT NULL ORDER BY id LIMIT ${MAX_TASKS}`.map( (row) => { this.requireConversation(row.conversation_id!); return this.summary(row, now); }, ); }
async create(value: unknown): Promise<TaskSummary> { const input = taskInput(value, true); const digest = await crypto.subtle.digest( "SHA-256", new TextEncoder().encode(JSON.stringify(input)), ); const hash = [...new Uint8Array(digest)] .map((byte) => byte.toString(16).padStart(2, "0")) .join(""); // Crypto yields. All lifecycle/replay checks and the write occur together // after it, so a concurrent deletion cannot be undone by the pending create. const existing = this.row(input.id); if (existing) { if (!existing.definition) throw new Error("Task was deleted"); if (existing.creation_hash !== hash) throw new Error( "Task creation ID was already used for different input", ); this.requireConversation(input.conversationId); return this.summary(existing); } this.requireConversation(input.conversationId); if ( input.schedule.kind === "once" && Date.parse(input.schedule.at) <= Date.now() ) throw new Error("New one-off task needs a future instant"); const [{ count }] = this.sql<{ count: number; }>`SELECT count(*) AS count FROM flarebot_tasks WHERE definition IS NOT NULL`; if (count >= MAX_TASKS) throw new Error(`Task limit reached (${MAX_TASKS})`); const now = new Date().toISOString(); const task: TaskDefinition = { ...input, version: 1, createdAt: now, updatedAt: now, }; this .sql`INSERT INTO flarebot_tasks (id, definition, creation_hash, conversation_id) VALUES (${task.id}, ${JSON.stringify(task)}, ${hash}, ${task.conversationId})`; return this.get(task.id); }
update(id: unknown, version: unknown, value: unknown): TaskSummary { const key = taskId(id); const expected = taskVersion(version); const input = taskInput(value); const { row, task } = this.active(key); if (task.version !== expected) throw new Error("Task changed. Reload before editing"); const scheduleChanged = JSON.stringify(input.schedule) !== JSON.stringify(task.schedule); if ( input.schedule.kind === "once" && (scheduleChanged || (!task.enabled && input.enabled)) ) { if ( Date.parse(input.schedule.at) <= Date.now() || (!scheduleChanged && row.once_consumed_at) ) throw new Error("Rearming a one-off task needs a new future instant"); } const updated: TaskDefinition = { ...task, ...input, version: task.version + 1, updatedAt: new Date().toISOString(), }; this.sql`UPDATE flarebot_tasks SET definition = ${JSON.stringify(updated)}, once_consumed_at = ${scheduleChanged ? null : row.once_consumed_at} WHERE id = ${key}`; return this.get(key); }
delete(id: unknown, version: unknown): { deleted: true } { return this.storage.transactionSync(() => { const key = taskId(id); const expected = taskVersion(version); const row = this.row(key); if (!row?.definition) { // Reserve even an absent ID: an in-flight create may still be hashing. this.sql`INSERT OR IGNORE INTO flarebot_tasks (id) VALUES (${key})`; return { deleted: true }; } if ((JSON.parse(row.definition) as TaskDefinition).version !== expected) throw new Error("Task changed. Reload before deleting"); this.tombstone(key); return { deleted: true }; }); }
private tombstone(id: string) { // Keep only an ID tombstone. Private nonterminal cleanup references survive // until native cancellation is implemented; no task text/history is retained. this.sql`UPDATE flarebot_tasks SET definition = NULL, creation_hash = NULL, conversation_id = NULL, once_consumed_at = NULL WHERE id = ${id}`; this.sql`DELETE FROM flarebot_task_runs WHERE task_id = ${id} AND payload IS NOT NULL AND json_extract(payload, '$.status') IN ('completed', 'aborted', 'skipped', 'error')`; this.sql`UPDATE flarebot_task_runs SET payload = NULL, created_at = NULL WHERE task_id = ${id} AND payload IS NOT NULL`; }
deleteForConversation(conversationId: string) { return this.storage.transactionSync(() => { for (const { id } of this.sql<{ id: string; }>`SELECT id FROM flarebot_tasks WHERE conversation_id = ${conversationId}`) this.tombstone(id); }); }
history(id: unknown, value?: unknown): TaskRunPage { const key = taskId(id); this.active(key); const { limit, before } = taskRunPage(value); const rows = before ? this.sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE task_id = ${key} AND payload IS NOT NULL AND (created_at < ${before.createdAt} OR (created_at = ${before.createdAt} AND id < ${before.id})) ORDER BY created_at DESC, id DESC LIMIT ${limit + 1}` : this.sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE task_id = ${key} AND payload IS NOT NULL ORDER BY created_at DESC, id DESC LIMIT ${limit + 1}`; const runs = rows .slice(0, limit) .map((row) => JSON.parse(row.payload) as TaskRun); const last = runs.at(-1); return { runs, nextCursor: rows.length > limit && last ? { createdAt: last.createdAt, id: last.id } : null, }; }
// Bounded pages rotate by ID in the execution reconciler. Deleted records // expose only cancellation references, never erased task text/history. unsettled(after = "", limit = 100, taskId: string | null = null) { return this.sql<{ id: string; conversation_id: string; submission_id: string; payload: string | null; }>`SELECT id, conversation_id, submission_id, payload FROM flarebot_task_runs WHERE id > ${after} AND (${taskId} IS NULL OR task_id = ${taskId}) AND (payload IS NULL OR json_extract(payload, '$.status') NOT IN ('completed', 'aborted', 'skipped', 'error')) ORDER BY id LIMIT ${limit}`.map((row) => ({ id: row.id, conversationId: row.conversation_id, submissionId: row.submission_id, run: row.payload ? (JSON.parse(row.payload) as TaskRun) : null, })); }
skipDispatch(run: TaskRun) { const row = this.sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE id = ${run.id} AND payload IS NOT NULL`[0]; if (!row) return null; const current = JSON.parse(row.payload) as TaskRun; // An accepted submission's actual terminal status comes from Think. A stale // callback must not replace that result with a local pre-dispatch skip. return current.status === "dispatching" || current.status === "dispatch_error" ? this.project({ ...run, status: "skipped" }) : current; }
needsDispatchRecovery(id: string) { return ( this.sql`SELECT id FROM flarebot_task_runs WHERE id = ${id} AND payload IS NOT NULL AND recovery_scheduled = 0`.length > 0 ); }
dispatchedRecovery(id: string) { this .sql`UPDATE flarebot_task_runs SET recovery_scheduled = 1 WHERE id = ${id}`; }
finishCleanup(id: string) { this .sql`DELETE FROM flarebot_task_runs WHERE id = ${id} AND payload IS NULL`; }
finishConversationCleanup(id: string) { this .sql`DELETE FROM flarebot_task_runs WHERE conversation_id = ${id} AND payload IS NULL`; }
authorized(run: TaskRun): TaskDefinition | null { try { const task = this.get(run.taskId); const stored = this.sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE id = ${run.id} AND payload IS NOT NULL`[0]; if (!stored) return null; const current = JSON.parse(stored.payload) as TaskRun; return task.version === run.taskVersion && task.conversationId === run.conversationId && current.submissionId === run.submissionId && current.taskVersion === run.taskVersion && (run.source === "manual" || task.enabled) ? task : null; } catch { return null; } }
// Internal occurrence/projection seam for native execution. No owner callable // writes history, and a projection can never insert or revive a missing run. begin(input: BeginTaskRun): TaskRun { return this.storage.transactionSync(() => { const { row, task } = this.active(taskId(input.taskId)); taskVersion(input.taskVersion); const scheduledFor = taskInstant(input.scheduledFor); if (input.source !== "scheduled" && input.source !== "manual") throw new Error("Invalid task run source"); const suffix = input.source === "manual" ? `manual:${taskId(input.requestId)}` : `${input.taskVersion}:${scheduledFor}`; const id = `${task.id}:${suffix}`; const existing = this .sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE id = ${id}`[0]; if (existing?.payload) return JSON.parse(existing.payload) as TaskRun; if ( task.version !== input.taskVersion || (input.source === "scheduled" && !task.enabled) ) throw new Error("Task changed or is disabled"); if ( input.source === "scheduled" && task.schedule.kind === "once" && (row.once_consumed_at || task.schedule.at !== scheduledFor) ) throw new Error("One-off task occurrence is no longer available"); const run: TaskRun = { id, taskId: task.id, taskVersion: task.version, conversationId: task.conversationId, submissionId: `task:${id}`, source: input.source, scheduledFor, createdAt: new Date().toISOString(), status: "dispatching", startedAt: null, completedAt: null, failureCode: null, }; this .sql`INSERT INTO flarebot_task_runs (id, task_id, conversation_id, submission_id, created_at, payload) VALUES (${id}, ${task.id}, ${task.conversationId}, ${run.submissionId}, ${run.createdAt}, ${JSON.stringify(run)})`; if (input.source === "scheduled" && task.schedule.kind === "once") this .sql`UPDATE flarebot_tasks SET once_consumed_at = ${scheduledFor} WHERE id = ${task.id}`; this.observed(run); return run; }); }
project(input: TaskRunProjection): TaskRun | null { const row = this .sql<RunRow>`SELECT payload FROM flarebot_task_runs WHERE id = ${input.id} AND payload IS NOT NULL`[0]; if (!row) return null; const run = JSON.parse(row.payload) as TaskRun; if ( run.taskVersion !== input.taskVersion || run.conversationId !== input.conversationId || run.submissionId !== input.submissionId ) return null; this.active(run.taskId); if (terminalTaskStatuses.has(run.status)) return run; if ( run.status === "running" && (input.status === "pending" || input.status === "dispatch_error" || input.status === "dispatching") ) return run; if ( run.status === "pending" && (input.status === "dispatching" || input.status === "dispatch_error") ) return run; if ( ![ "dispatching", "dispatch_error", "pending", "running", ...terminalTaskStatuses, ].includes(input.status) ) throw new Error("Invalid task run status"); const updated: TaskRun = { ...run, status: input.status, startedAt: input.startedAt ? new Date(input.startedAt).toISOString() : run.startedAt, completedAt: terminalTaskStatuses.has(input.status) ? input.completedAt ? new Date(input.completedAt).toISOString() : new Date().toISOString() : null, failureCode: input.status === "dispatch_error" ? "dispatch_failed" : input.status === "error" ? "turn_failed" : input.status === "aborted" ? "turn_aborted" : input.status === "skipped" ? "turn_skipped" : null, }; this .sql`UPDATE flarebot_task_runs SET payload = ${JSON.stringify(updated)} WHERE id = ${run.id} AND payload IS NOT NULL`; if ( run.status !== updated.status || run.startedAt !== updated.startedAt || run.completedAt !== updated.completedAt ) this.observed(updated); return updated; }}