Skip to content
File

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

typescript72 lines
1import { asc, eq } from "drizzle-orm";
2 
3import { RunId } from "@/contracts";
4import { type ReplaceRunStepsInput, type RunStepState, type UpdateRunStepStateInput } from "@/worker/contracts";
5import * as runSchema from "@/worker/db/durable/schema/run-do";
6 
7import { assertStepTransition } from "../state";
8import { chunkRows, type RunDb, type RunTx, toRunStepState } from "./core";
9 
10const MAX_SQL_BOUND_PARAMETERS = 100;
11const RUN_STEP_INSERT_COLUMN_COUNT = 9;
12const MAX_RUN_STEP_ROWS_PER_INSERT = Math.floor(MAX_SQL_BOUND_PARAMETERS / RUN_STEP_INSERT_COLUMN_COUNT);
13 
14export 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 
24export 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 
50export 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};