File
Blob: src/worker/durable/run-do/repo/meta.ts
| 1 | import { eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { RunId } from "@/contracts"; |
| 4 | import { isTerminalStatus, type EnsureRunInput, type RunMetaState, type UpdateRunStateInput } from "@/worker/contracts"; |
| 5 | import * as runSchema from "@/worker/db/durable/schema/run-do"; |
| 6 | |
| 7 | import { type RunDb, type TryUpdateRunStateResult, toRunMetaState } from "./core"; |
| 8 | import { resolveRunStateUpdate, RunStateTransitionError } from "../state"; |
| 9 | |
| 10 | export const getRunMeta = async (db: RunDb, runId: RunId): Promise<RunMetaState | null> => { |
| 11 | const rows = await db.select().from(runSchema.runMeta).where(eq(runSchema.runMeta.id, runId)).limit(1); |
| 12 | if (rows.length === 0) { |
| 13 | return null; |
| 14 | } |
| 15 | |
| 16 | return toRunMetaState(rows[0]); |
| 17 | }; |
| 18 | |
| 19 | export const ensureInitialized = async (db: RunDb, input: EnsureRunInput): Promise<void> => { |
| 20 | const payload = input; |
| 21 | const existing = await getRunMeta(db, payload.runId); |
| 22 | if (existing) { |
| 23 | if ( |
| 24 | existing.projectId !== payload.projectId || |
| 25 | existing.triggerType !== payload.triggerType || |
| 26 | existing.branch !== payload.branch |
| 27 | ) { |
| 28 | throw new Error(`Run ${payload.runId} is already initialized with different immutable fields.`); |
| 29 | } |
| 30 | |
| 31 | if (existing.commitSha === payload.commitSha) { |
| 32 | return; |
| 33 | } |
| 34 | |
| 35 | if (existing.commitSha === null && payload.commitSha !== null) { |
| 36 | await db |
| 37 | .update(runSchema.runMeta) |
| 38 | .set({ commitSha: payload.commitSha }) |
| 39 | .where(eq(runSchema.runMeta.id, payload.runId)); |
| 40 | return; |
| 41 | } |
| 42 | |
| 43 | if (existing.commitSha !== null && payload.commitSha === null) { |
| 44 | return; |
| 45 | } |
| 46 | |
| 47 | if (existing.commitSha !== payload.commitSha) { |
| 48 | throw new Error(`Run ${payload.runId} is already initialized with different immutable fields.`); |
| 49 | } |
| 50 | |
| 51 | return; |
| 52 | } |
| 53 | |
| 54 | await db.insert(runSchema.runMeta).values({ |
| 55 | id: payload.runId, |
| 56 | projectId: payload.projectId, |
| 57 | status: "queued", |
| 58 | triggerType: payload.triggerType, |
| 59 | branch: payload.branch, |
| 60 | commitSha: payload.commitSha, |
| 61 | currentStep: null, |
| 62 | startedAt: null, |
| 63 | finishedAt: null, |
| 64 | exitCode: null, |
| 65 | errorMessage: null, |
| 66 | }); |
| 67 | }; |
| 68 | |
| 69 | export const updateRunState = async (db: RunDb, input: UpdateRunStateInput): Promise<void> => { |
| 70 | const payload = input; |
| 71 | const current = await getRunMeta(db, payload.runId); |
| 72 | if (!current) { |
| 73 | throw new Error(`Run ${payload.runId} is not initialized.`); |
| 74 | } |
| 75 | |
| 76 | await db |
| 77 | .update(runSchema.runMeta) |
| 78 | .set(resolveRunStateUpdate(current, payload)) |
| 79 | .where(eq(runSchema.runMeta.id, payload.runId)); |
| 80 | }; |
| 81 | |
| 82 | export const tryUpdateRunState = async (db: RunDb, input: UpdateRunStateInput): Promise<TryUpdateRunStateResult> => { |
| 83 | const payload = input; |
| 84 | const current = await getRunMeta(db, payload.runId); |
| 85 | if (!current) { |
| 86 | throw new Error(`Run ${payload.runId} is not initialized.`); |
| 87 | } |
| 88 | |
| 89 | try { |
| 90 | const resolved = resolveRunStateUpdate(current, payload); |
| 91 | await db.update(runSchema.runMeta).set(resolved).where(eq(runSchema.runMeta.id, payload.runId)); |
| 92 | |
| 93 | return { |
| 94 | kind: "applied", |
| 95 | state: { |
| 96 | ...current, |
| 97 | status: resolved.status, |
| 98 | currentStep: resolved.currentStep === undefined ? current.currentStep : resolved.currentStep, |
| 99 | startedAt: resolved.startedAt, |
| 100 | finishedAt: resolved.finishedAt, |
| 101 | exitCode: resolved.exitCode === undefined ? current.exitCode : resolved.exitCode, |
| 102 | errorMessage: resolved.errorMessage === undefined ? current.errorMessage : resolved.errorMessage, |
| 103 | }, |
| 104 | }; |
| 105 | } catch (error) { |
| 106 | if (!(error instanceof RunStateTransitionError)) { |
| 107 | throw error; |
| 108 | } |
| 109 | |
| 110 | return { |
| 111 | kind: "conflict", |
| 112 | reason: error.reason, |
| 113 | current, |
| 114 | }; |
| 115 | } |
| 116 | }; |
| 117 | |
| 118 | export const repairTerminalState = async (db: RunDb, input: UpdateRunStateInput): Promise<void> => { |
| 119 | const payload = input; |
| 120 | const current = await getRunMeta(db, payload.runId); |
| 121 | if (!current) { |
| 122 | throw new Error(`Run ${payload.runId} is not initialized.`); |
| 123 | } |
| 124 | |
| 125 | if (!isTerminalStatus(payload.status)) { |
| 126 | throw new Error(`Run ${payload.runId} cannot be repaired to non-terminal status ${payload.status}.`); |
| 127 | } |
| 128 | |
| 129 | const startedAt = payload.startedAt ?? current.startedAt; |
| 130 | if (payload.status === "passed" && startedAt === null) { |
| 131 | throw new Error(`Run ${payload.runId} cannot be repaired to terminal status ${payload.status} without startedAt.`); |
| 132 | } |
| 133 | |
| 134 | const finishedAt = payload.finishedAt ?? current.finishedAt; |
| 135 | if (finishedAt === null) { |
| 136 | throw new Error(`Run ${payload.runId} cannot be repaired to terminal status ${payload.status} without finishedAt.`); |
| 137 | } |
| 138 | |
| 139 | await db |
| 140 | .update(runSchema.runMeta) |
| 141 | .set({ |
| 142 | status: payload.status, |
| 143 | currentStep: payload.currentStep ?? null, |
| 144 | startedAt, |
| 145 | finishedAt, |
| 146 | exitCode: payload.exitCode, |
| 147 | errorMessage: payload.errorMessage, |
| 148 | }) |
| 149 | .where(eq(runSchema.runMeta.id, payload.runId)); |
| 150 | }; |
| 151 | |
| 152 | export const deleteRunData = async (db: RunDb, runId: RunId): Promise<void> => { |
| 153 | await db.delete(runSchema.runLogs).where(eq(runSchema.runLogs.runId, runId)); |
| 154 | await db.delete(runSchema.runSteps).where(eq(runSchema.runSteps.runId, runId)); |
| 155 | await db.delete(runSchema.runMeta).where(eq(runSchema.runMeta.id, runId)); |
| 156 | }; |