Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217import type { ThinkSubmissionInspection } from "@cloudflare/think";import type { Schedule } from "agents";import type { TaskRun, TaskSummary } from "../shared/tasks";import type { PersonalAgent } from "./personal-agent";import { terminalTaskStatuses, type TaskStore } from "./task-store";
export interface TaskSchedulePayload { taskId: string; taskVersion: number; // Date matching in the native scheduler ignores the date. Include it in the // canonical payload; cron identity instead includes the fired schedule.time. at: string | null;}
export function taskSchedulePayload(task: TaskSummary): TaskSchedulePayload { return { taskId: task.id, taskVersion: task.version, at: task.schedule.kind === "once" ? task.schedule.at : null, };}
export function submissionProjection( run: TaskRun, status: ThinkSubmissionInspection,) { return { id: run.id, taskVersion: run.taskVersion, conversationId: run.conversationId, submissionId: status.submissionId, status: status.status, startedAt: status.startedAt ? new Date(status.startedAt).toISOString() : null, completedAt: status.completedAt ? new Date(status.completedAt).toISOString() : null, };}
// Native schedules are the binding registry. Only task intent and safe run// projections live in application SQL; there is no copied timer or turn queue.export class TaskExecution { private cursor = ""; constructor( private readonly host: PersonalAgent, private readonly store: TaskStore, ) {}
private matches(schedule: Schedule<unknown>, task: TaskSummary) { return ( schedule.callback === "dispatchScheduledTask" && JSON.stringify(schedule.payload) === JSON.stringify(taskSchedulePayload(task)) ); }
async schedules() { const tasks = this.store.list(); const allSchedules = await this.host.listSchedules(); const schedules = allSchedules.filter( (s) => s.callback === "dispatchScheduledTask", ); for (const schedule of allSchedules.filter( (s) => s.callback === "dispatchManualTask" || s.callback === "dispatchRecoveredTask", )) { if (!this.store.authorized(schedule.payload as TaskRun)) await this.host.cancelSchedule(schedule.id); } for (const schedule of schedules) { // Re-read after each await: cancellation cannot stop a callback snapshot. let task: TaskSummary | undefined; try { task = this.store.get((schedule.payload as TaskSchedulePayload).taskId); } catch { /* Deleted task binding. */ } if (!task?.enabled || !this.matches(schedule, task)) await this.host.cancelSchedule(schedule.id); } for (const snapshot of tasks) { let task: TaskSummary; try { task = this.store.get(snapshot.id); } catch { continue; } if (!task.enabled || !task.nextRunAt) continue; const native = await this.host.schedule( task.schedule.kind === "once" ? new Date(task.schedule.at) : task.schedule.expression, "dispatchScheduledTask", taskSchedulePayload(task), { idempotent: true, retry: { maxAttempts: 3, baseDelayMs: 500, maxDelayMs: 2000 }, }, ); let current: TaskSummary | undefined; try { current = this.store.get(task.id); } catch { /* Deleted during arming. */ } if (!current?.enabled || current.version !== task.version) await this.host.cancelSchedule(native.id); } if ( this.store.list().some((task) => task.nextRunAt) || this.store.unsettled().length ) await this.host.scheduleEvery(30, "reconcileTaskExecution"); else for (const schedule of allSchedules.filter( (s) => s.callback === "reconcileTaskExecution", )) await this.host.cancelSchedule(schedule.id); }
async summary(id: unknown): Promise<TaskSummary> { const task = this.store.get(id); if (!task.nextRunAt) return { ...task, schedulingError: null }; const binding = (await this.host.listSchedules()).find((s) => this.matches(s, task), ); const native = binding ? await this.host.getScheduleById(binding.id) : undefined; const current = this.store.get(id); if (current.version !== task.version) return this.summary(id); return { ...current, nextRunAt: current.nextRunAt && native ? new Date(native.time * 1000).toISOString() : null, schedulingError: current.nextRunAt && !native ? "schedule_unavailable" : null, }; }
async dispatch(run: TaskRun) { if (terminalTaskStatuses.has(run.status)) return run; try { if (!this.store.authorized(run)) return this.skip(run); const child = await this.host.taskConversation(run.conversationId); if (!child || !this.store.authorized(run)) return this.skip(run); const receipt = await child.submitTaskRun(run); // Receipt may be older than an already delivered final status report. if (receipt) this.store.project(submissionProjection(run, receipt)); if (!this.store.authorized(run)) await child.cancelTaskRun(run.submissionId); return receipt ? this.store.project(submissionProjection(run, receipt)) : this.skip(run); } catch { if (!this.store.authorized(run)) return this.skip(run); this.store.project({ ...run, status: "dispatch_error" }); // SDK callback retry reuses exactly the same logical submission identity. throw new Error("Task dispatch failed"); } }
private skip(run: TaskRun) { return this.store.skipDispatch(run); }
async repair(taskId: string | null = null, limit = 100) { const rows = this.store.unsettled(taskId ? "" : this.cursor, limit, taskId); const deadline = Date.now() + 5000; if (!taskId) this.cursor = ""; for (const row of rows) { if (!taskId) this.cursor = row.id; try { const child = await this.host.taskConversation(row.conversationId); if (!child) continue; // Conversation teardown owns facet cancellation. const allowed = row.run && this.store.authorized(row.run); if (!allowed) await child.cancelTaskRun(row.submissionId); const status = await child.inspectTaskRun(row.submissionId); if (row.run && status) this.store.project(submissionProjection(row.run, status)); else if ( !row.run && (!status || terminalTaskStatuses.has(status.status)) ) this.store.finishCleanup(row.id); else if (row.run && !allowed) this.skip(row.run); else if ( row.run && !status && this.store.needsDispatchRecovery(row.id) && Date.now() - Date.parse(row.run.createdAt) >= 10_000 ) { // One bounded native recovery batch closes a parent crash between the // occurrence write and acceptance. Always reuse the SAME submission. // A distinct callback cannot dedupe against an executing original // manual row that the SDK will remove when its callback returns. await this.host.schedule(1, "dispatchRecoveredTask", row.run, { idempotent: true, retry: { maxAttempts: 3, baseDelayMs: 500, maxDelayMs: 2000 }, }); this.store.dispatchedRecovery(row.id); } } catch { // Safe projection stays pending/error until the next native reconciliation. // Never manufacture another accepted turn because an observer was lost. } if (Date.now() >= deadline) return; } if (!taskId && rows.length < limit) this.cursor = ""; }}