File
Blob: src/worker/durable/run-do.ts
| 1 | import { DurableObject } from "cloudflare:workers"; |
| 2 | import { drizzle, type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite"; |
| 3 | import { migrate } from "drizzle-orm/durable-sqlite/migrator"; |
| 4 | import { RunId } from "@/contracts"; |
| 5 | import { |
| 6 | type AppendRunLogsInput, |
| 7 | type EnsureRunInput, |
| 8 | RunDetailState, |
| 9 | RunMetaState, |
| 10 | type UpdateRunStateInput, |
| 11 | type UpdateRunStepStateInput, |
| 12 | type ReplaceRunStepsInput, |
| 13 | } from "@/worker/contracts"; |
| 14 | import runMigrations from "../../../drizzle/run-do/migrations.js"; |
| 15 | import * as runSchema from "@/worker/db/durable/schema/run-do"; |
| 16 | import type { TryUpdateRunStateResult } from "@/worker/durable/run-do/repo/core"; |
| 17 | import { createLogger } from "@/worker/services"; |
| 18 | import { |
| 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 | |
| 36 | const logger = createLogger("durable.run"); |
| 37 | |
| 38 | export 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 | } |