Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201import fixture, { PersonalAgent as FixturePersonalAgent, Conversation as FixtureConversation,} from "./think-worker";import { getAgentByName, type Schedule } from "agents";import type { Env } from "../../worker/personal-agent";import type { TaskRun } from "../../shared/tasks";import type { ThinkSubmissionInspection } from "@cloudflare/think";import type { TaskSchedulePayload } from "../../worker/task-execution";import type { BeginTaskRun } from "../../worker/task-store";export { Sandbox } from "./think-worker";
interface Fault { holdOutput?: boolean; delay?: number; fail?: number; failOccurrence?: boolean; failedSubmission?: string; lostReply?: boolean; lostReport?: boolean; lateReceipt?: boolean; runningDelay?: number;}export class Conversation extends FixtureConversation { private fault: Fault = {}; private releaseOutput?: () => void;
protected override pauseAfterFirstChunk(text: string, signal?: AbortSignal) { if (!this.fault.holdOutput) return super.pauseAfterFirstChunk(text, signal); return new Promise<void>((resolve, reject) => { const cleanup = () => { clearTimeout(timer); signal?.removeEventListener("abort", finish); this.releaseOutput = undefined; }; const finish = () => { cleanup(); resolve(); }; const timer = setTimeout(() => { cleanup(); reject(new Error("Fixture output was never released")); }, 30_000); this.releaseOutput = finish; signal?.addEventListener("abort", finish, { once: true }); if (signal?.aborted) finish(); }); }
releaseModelOutput() { this.releaseOutput?.(); } async configureExecutionFault(fault: Fault) { this.fault = fault; await this.ctx.storage.put("fixture-execution-fault", fault); } async submitTaskRun(run: TaskRun) { this.fault = (await this.ctx.storage.get<Fault>("fixture-execution-fault")) ?? {}; if (this.fault.delay) { const delay = this.fault.delay; this.fault.delay = 0; await this.ctx.storage.put("fixture-execution-fault", this.fault); await new Promise((resolve) => setTimeout(resolve, delay)); } if (this.fault.failOccurrence) { this.fault.failedSubmission ??= run.submissionId; await this.ctx.storage.put("fixture-execution-fault", this.fault); if (this.fault.failedSubmission === run.submissionId) throw new Error("Fixture occurrence acceptance unavailable"); } if (this.fault.fail) { this.fault.fail--; await this.ctx.storage.put("fixture-execution-fault", this.fault); throw new Error("Fixture acceptance unavailable PRIVATE-ERROR"); } const receipt = await super.submitTaskRun(run); if (this.fault.lostReply) { this.fault.lostReply = false; await this.ctx.storage.put("fixture-execution-fault", this.fault); throw new Error("Fixture lost accepted reply"); } if (this.fault.lateReceipt) { this.fault.lateReceipt = false; await this.ctx.storage.put("fixture-execution-fault", this.fault); await new Promise((resolve) => setTimeout(resolve, 1500)); } return receipt; } protected async onSubmissionStatus(submission: ThinkSubmissionInspection) { if (submission.status === "running" && this.fault.runningDelay) { const delay = this.fault.runningDelay; this.fault.runningDelay = 0; await this.ctx.storage.put("fixture-execution-fault", this.fault); await new Promise((resolve) => setTimeout(resolve, delay)); } if ( this.fault.lostReport && ["completed", "error", "aborted", "skipped"].includes(submission.status) ) return; await super.onSubmissionStatus(submission); const events = (this.env as Env & { TEST_EXECUTION_EVENTS?: string }) .TEST_EXECUTION_EVENTS; if (events && submission.status === "completed") { const response = await fetch(events, { method: "POST", body: JSON.stringify({ conversationId: this.name, submissionId: submission.submissionId, }), }); if (!response.ok) throw new Error("Fixture completion event rejected"); } } async executionSnapshot() { return { outputHeld: !!this.releaseOutput, submissions: await this.listSubmissions({ limit: 100 }), messages: await this.getMessages(), slowSecondAttempts: await this.ctx.storage.get( "fixture-attempts:slow-second", ), recoverAttempts: await this.ctx.storage.get("fixture-attempts:recover"), }; } submitOrdinary(text: string) { return this.submitMessages([ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ]); } clearExecution() { return this.clearMessages(); }}
export class PersonalAgent extends FixturePersonalAgent { async executionFixture(action: string, input: Record<string, unknown>) { if (action === "snapshot") return { tasks: this.tasks.list(), schedules: await this.listSchedules(), runs: this.sql`SELECT * FROM flarebot_task_runs`, facets: this.listSubAgents(FixtureConversation), }; if (action === "begin") return this.beginTaskRun(input as unknown as BeginTaskRun); if (action === "replay") return this.dispatchScheduledTask( input.payload as TaskSchedulePayload, input as unknown as Schedule<TaskSchedulePayload>, ); if (action === "dispatch") return this.dispatchManualTask(input as unknown as TaskRun); if (action === "repair") return this.reconcileTaskExecution(); if (action === "lose-binding") { await this.cancelSchedule(input.id as string); return; } const child = await this.subAgent( Conversation, input.conversationId as string, ); if (action === "fault") return child.configureExecutionFault(input.fault as Fault); if (action === "conversation") return child.executionSnapshot(); if (action === "release-output") return child.releaseModelOutput(); if (action === "ordinary") return child.submitOrdinary(input.text as string); if (action === "clear") return child.clearExecution(); throw new Error("Unknown execution fixture action"); }}
export default { async fetch(request, env, ctx) { const path = new URL(request.url).pathname; if (path.startsWith("/__execution/")) { const personal = await getAgentByName( env.PersonalAgent as unknown as DurableObjectNamespace<PersonalAgent>, "personal", ); try { return Response.json( (await personal.executionFixture( path.split("/").at(-1)!, request.method === "POST" ? await request.json() : {}, )) ?? null, ); } catch (error) { return Response.json({ error: String(error) }, { status: 500 }); } } return fixture.fetch(request, env, ctx); },} satisfies ExportedHandler<Env>;