File
Blob: src/worker/durable/run-do/repo/steps.ts
| 1 | import { asc, eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { RunId } from "@/contracts"; |
| 4 | import { type ReplaceRunStepsInput, type RunStepState, type UpdateRunStepStateInput } from "@/worker/contracts"; |
| 5 | import * as runSchema from "@/worker/db/durable/schema/run-do"; |
| 6 | |
| 7 | import { assertStepTransition } from "../state"; |
| 8 | import { chunkRows, type RunDb, type RunTx, toRunStepState } from "./core"; |
| 9 | |
| 10 | const MAX_SQL_BOUND_PARAMETERS = 100; |
| 11 | const RUN_STEP_INSERT_COLUMN_COUNT = 9; |
| 12 | const MAX_RUN_STEP_ROWS_PER_INSERT = Math.floor(MAX_SQL_BOUND_PARAMETERS / RUN_STEP_INSERT_COLUMN_COUNT); |
| 13 | |
| 14 | export const listRunSteps = async (db: RunDb, runId: RunId): Promise<RunStepState[]> => { |
| 15 | const rows = await db |
| 16 | .select() |
| 17 | .from(runSchema.runSteps) |
| 18 | .where(eq(runSchema.runSteps.runId, runId)) |
| 19 | .orderBy(asc(runSchema.runSteps.position)); |
| 20 | |
| 21 | return rows.map(toRunStepState); |
| 22 | }; |
| 23 | |
| 24 | export const replaceSteps = async (db: RunDb, input: ReplaceRunStepsInput): Promise<void> => { |
| 25 | const payload = input; |
| 26 | const stepRows = payload.steps.map((step) => ({ |
| 27 | id: `${payload.runId}:step:${step.position}`, |
| 28 | runId: payload.runId, |
| 29 | position: step.position, |
| 30 | name: step.name, |
| 31 | command: step.command, |
| 32 | status: "queued" as const, |
| 33 | startedAt: null, |
| 34 | finishedAt: null, |
| 35 | exitCode: null, |
| 36 | })); |
| 37 | |
| 38 | db.transaction((tx: RunTx) => { |
| 39 | tx.delete(runSchema.runSteps).where(eq(runSchema.runSteps.runId, payload.runId)).run(); |
| 40 | if (stepRows.length === 0) { |
| 41 | return; |
| 42 | } |
| 43 | |
| 44 | for (const batch of chunkRows(stepRows, MAX_RUN_STEP_ROWS_PER_INSERT)) { |
| 45 | tx.insert(runSchema.runSteps).values(batch).run(); |
| 46 | } |
| 47 | }); |
| 48 | }; |
| 49 | |
| 50 | export const updateStepState = async (db: RunDb, input: UpdateRunStepStateInput): Promise<void> => { |
| 51 | const payload = input; |
| 52 | const stepId = `${payload.runId}:step:${payload.position}`; |
| 53 | const rows = await db.select().from(runSchema.runSteps).where(eq(runSchema.runSteps.id, stepId)).limit(1); |
| 54 | const currentRow = rows[0]; |
| 55 | if (!currentRow) { |
| 56 | throw new Error(`Run step ${stepId} was not found.`); |
| 57 | } |
| 58 | |
| 59 | const current = toRunStepState(currentRow); |
| 60 | assertStepTransition(current.status, payload.status); |
| 61 | |
| 62 | await db |
| 63 | .update(runSchema.runSteps) |
| 64 | .set({ |
| 65 | status: payload.status, |
| 66 | startedAt: payload.startedAt, |
| 67 | finishedAt: payload.finishedAt, |
| 68 | exitCode: payload.exitCode, |
| 69 | }) |
| 70 | .where(eq(runSchema.runSteps.id, stepId)); |
| 71 | }; |