Skip to content
File

Blob: src/worker/durable/run-do.ts

typescript105 lines
1import { DurableObject } from "cloudflare:workers";
2import { drizzle, type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite";
3import { migrate } from "drizzle-orm/durable-sqlite/migrator";
4import { RunId } from "@/contracts";
5import {
6 type AppendRunLogsInput,
7 type EnsureRunInput,
8 RunDetailState,
9 RunMetaState,
10 type UpdateRunStateInput,
11 type UpdateRunStepStateInput,
12 type ReplaceRunStepsInput,
13} from "@/worker/contracts";
14import runMigrations from "../../../drizzle/run-do/migrations.js";
15import * as runSchema from "@/worker/db/durable/schema/run-do";
16import type { TryUpdateRunStateResult } from "@/worker/durable/run-do/repo/core";
17import { createLogger } from "@/worker/services";
18import {
19 appendLogs,
20 listRunLogs,
21 deleteRunData,
22 ensureInitialized,
23 getRunMeta,
24 listRunSteps,
25 repairTerminalState,
26 replaceSteps,
27 tryUpdateRunState,
28 updateRunState,
29 updateStepState,
30 broadcastLogEvents,
31 broadcastStateUpdate,
32 handleRunLogStreamFetch,
33 logRunSocketError,
34} from "@/worker/durable/run-do/index";
35 
36const logger = createLogger("durable.run");
37 
38export class RunDO extends DurableObject {
39 private readonly db: DrizzleSqliteDODatabase<typeof runSchema>;
40 constructor(ctx: DurableObjectState, env: Env) {
41 super(ctx, env);
42 this.db = drizzle(ctx.storage, { schema: runSchema });
43 this.ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong"));
44 ctx.blockConcurrencyWhile(async () => {
45 await migrate(this.db, runMigrations);
46 });
47 }
48 public async fetch(request: Request): Promise<Response> {
49 return handleRunLogStreamFetch(this.ctx, this.db, request);
50 }
51 async ensureInitialized(input: EnsureRunInput): Promise<void> {
52 await ensureInitialized(this.db, input);
53 }
54 async getRunSummary(runId: RunId): Promise<RunMetaState | null> {
55 return getRunMeta(this.db, runId);
56 }
57 async getRunDetail(runId: RunId): Promise<RunDetailState> {
58 const meta = await getRunMeta(this.db, runId);
59 const steps = await listRunSteps(this.db, runId);
60 const recentLogs = await listRunLogs(this.db, runId);
61 return {
62 meta,
63 steps,
64 recentLogs,
65 };
66 }
67 async updateRunState(input: UpdateRunStateInput): Promise<void> {
68 await updateRunState(this.db, input);
69 await broadcastStateUpdate(this.ctx, this.db, logger, input.runId);
70 }
71 async tryUpdateRunState(input: UpdateRunStateInput): Promise<TryUpdateRunStateResult> {
72 const result = await tryUpdateRunState(this.db, input);
73 if (result.kind === "applied") {
74 await broadcastStateUpdate(this.ctx, this.db, logger, input.runId);
75 }
76 return result;
77 }
78 async repairTerminalState(input: UpdateRunStateInput): Promise<void> {
79 await repairTerminalState(this.db, input);
80 await broadcastStateUpdate(this.ctx, this.db, logger, input.runId);
81 }
82 async replaceSteps(input: ReplaceRunStepsInput): Promise<void> {
83 await replaceSteps(this.db, input);
84 await broadcastStateUpdate(this.ctx, this.db, logger, input.runId);
85 }
86 async updateStepState(input: UpdateRunStepStateInput): Promise<void> {
87 await updateStepState(this.db, input);
88 await broadcastStateUpdate(this.ctx, this.db, logger, input.runId);
89 }
90 async appendLogs(input: AppendRunLogsInput): Promise<void> {
91 const appendedLogs = await appendLogs(this.db, input);
92 broadcastLogEvents(this.ctx, logger, input.runId, appendedLogs);
93 }
94 async deleteRunData(runId: RunId): Promise<void> {
95 await deleteRunData(this.db, runId);
96 }
97 webSocketMessage(): void {}
98 webSocketClose(ws: WebSocket, code: number, reason: string, _wasClean: boolean): void {
99 ws.close(code, reason);
100 }
101 webSocketError(ws: WebSocket, error: unknown): void {
102 logRunSocketError(logger, ws, error);
103 }
104}