import 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([ "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, 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`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`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`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 { 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`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`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`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`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`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`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; } }