Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495import * 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<PersonalAgent>, "authorizeTaskRun" | "getModelSettings" | "projectTaskRun">;type SubmissionQueue = Pick< Conversation, "submitMessages" | "cancelSubmission">;
export class ConversationTasks { constructor( private readonly name: () => string, private readonly parent: () => Effect.Effect<TaskParent, AgentFailure>, 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)), ); }); }}