import type { ThinkSubmissionInspection } from "@cloudflare/think"; import type { Schedule, SubAgentStub } from "agents"; import * as Effect from "effect/Effect"; import type { TaskRun, TaskSummary } from "../shared/tasks"; import { agentCall, AgentFailure } from "./agent-io"; import type { Conversation } from "./conversation"; 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: Pick< PersonalAgent, | "listSchedules" | "cancelSchedule" | "schedule" | "scheduleEvery" | "getScheduleById" >, private readonly store: TaskStore, private readonly conversation: ( id: string, ) => Effect.Effect | null, AgentFailure>, ) {} private matches(schedule: Schedule, task: TaskSummary) { return ( schedule.callback === "dispatchScheduledTask" && JSON.stringify(schedule.payload) === JSON.stringify(taskSchedulePayload(task)) ); } schedules() { return Effect.gen({ self: this }, function* () { const tasks = this.store.list(); const allSchedules = yield* agentCall(() => 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)) yield* agentCall(() => 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)) yield* agentCall(() => 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 = yield* agentCall(() => 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) yield* agentCall(() => this.host.cancelSchedule(native.id)); } if ( this.store.list().some((task) => task.nextRunAt) || this.store.unsettled().length ) yield* agentCall(() => this.host.scheduleEvery(30, "reconcileTaskExecution"), ); else for (const schedule of allSchedules.filter( (s) => s.callback === "reconcileTaskExecution", )) yield* agentCall(() => this.host.cancelSchedule(schedule.id)); }); } summary(id: unknown): Effect.Effect { return Effect.gen({ self: this }, function* () { const task = this.store.get(id); if (!task.nextRunAt) return { ...task, schedulingError: null }; const binding = (yield* agentCall(() => this.host.listSchedules())).find( (s) => this.matches(s, task), ); const native = binding ? yield* agentCall(() => this.host.getScheduleById(binding.id)) : undefined; const current = this.store.get(id); if (current.version !== task.version) return yield* this.summary(id); return { ...current, nextRunAt: current.nextRunAt && native ? new Date(native.time * 1000).toISOString() : null, schedulingError: current.nextRunAt && !native ? "schedule_unavailable" : null, }; }); } dispatch(run: TaskRun) { return Effect.gen({ self: this }, function* () { if (terminalTaskStatuses.has(run.status)) return run; return yield* Effect.gen({ self: this }, function* () { if (!this.store.authorized(run)) return this.skip(run); const child = yield* this.conversation(run.conversationId); if (!child || !this.store.authorized(run)) return this.skip(run); const receipt = yield* agentCall(() => 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)) yield* agentCall(() => child.cancelTaskRun(run.submissionId)); return receipt ? this.store.project(submissionProjection(run, receipt)) : this.skip(run); }).pipe( Effect.catchCause(() => Effect.gen({ self: this }, function* () { 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. return yield* Effect.fail( new AgentFailure({ message: "Task dispatch failed" }), ); }), ), ); }); } private skip(run: TaskRun) { return this.store.skipDispatch(run); } repair(taskId: string | null = null, limit = 100) { return Effect.gen({ self: this }, function* () { 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; yield* Effect.gen({ self: this }, function* () { const child = yield* this.conversation(row.conversationId); if (!child) return; // Conversation teardown owns facet cancellation. const allowed = row.run && this.store.authorized(row.run); if (!allowed) yield* agentCall(() => child.cancelTaskRun(row.submissionId)); const status = yield* agentCall(() => 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. yield* agentCall(() => this.host.schedule(1, "dispatchRecoveredTask", row.run, { idempotent: true, retry: { maxAttempts: 3, baseDelayMs: 500, maxDelayMs: 2000 }, }), ); this.store.dispatchedRecovery(row.id); } }).pipe( Effect.catchCause(() => Effect.void), Effect.timeoutOrElse({ duration: Math.max(1, deadline - Date.now()), orElse: () => Effect.void, }), ); if (Date.now() >= deadline) return; } if (!taskId && rows.length < limit) this.cursor = ""; }); } }