import * as Effect from "effect/Effect"; import type { ThinkSubmissionInspection } from "@cloudflare/think"; import type { TaskRun } from "../shared/tasks"; import type { PersonalAgent } from "./personal-agent"; import type { Conversation } from "./conversation"; import { agentCall, type AgentFailure } from "./agent-io"; import { submissionProjection } from "./task-execution"; type TaskParent = Pick< DurableObjectStub, "authorizeTaskRun" | "getModelSettings" | "projectTaskRun" >; type SubmissionQueue = Pick< Conversation, "submitMessages" | "cancelSubmission" >; export class ConversationTasks { constructor( private readonly name: () => string, private readonly parent: () => Effect.Effect, private readonly queue: SubmissionQueue, ) {} submit(run: TaskRun) { return Effect.gen({ self: this }, function* () { const parent = yield* this.parent(); const task = yield* agentCall(() => parent.authorizeTaskRun(run)); if (!task || task.conversationId !== this.name()) return null; const { configuration } = yield* agentCall(() => parent.getModelSettings(), ); const receipt = yield* agentCall(() => this.queue.submitMessages( [ { id: run.submissionId, role: "user", parts: [{ type: "text", text: task.instructions }], metadata: { scheduledModelConfiguration: configuration }, }, ], { submissionId: run.submissionId, idempotencyKey: run.submissionId, metadata: { taskRun: run, modelConfiguration: configuration }, }, ), ); if (!(yield* agentCall(() => parent.authorizeTaskRun(run)))) yield* agentCall(() => this.queue.cancelSubmission( run.submissionId, "Task changed or was deleted", ), ); return receipt; }); } observe(submission: ThinkSubmissionInspection) { return Effect.gen({ self: this }, function* () { const run = submission.metadata?.taskRun as TaskRun | undefined; if ( !run || run.submissionId !== submission.submissionId || run.conversationId !== this.name() ) return; if (submission.status === "pending" || submission.status === "running") { const allowed = yield* this.parent().pipe( Effect.flatMap((parent) => agentCall(() => parent.authorizeTaskRun(run)), ), Effect.map(Boolean), Effect.catchTag("AgentFailure", () => Effect.succeed(false)), ); if (!allowed) { // Think emits running before applying queued messages. Cancel explicitly // so a changed task cannot enter the queue through that start race. yield* agentCall(() => this.queue.cancelSubmission( run.submissionId, "Task changed or unavailable", ), ); return; } } const parent = yield* this.parent(); yield* agentCall(() => parent.projectTaskRun(submissionProjection(run, submission)), ); }); } }