Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262import 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<SubAgentStub<Conversation> | null, AgentFailure>, ) {}
private matches(schedule: Schedule<unknown>, 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<TaskSummary, AgentFailure> { 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 = ""; }); }}