Skip to content
File

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

typescript157 lines
1import { eq } from "drizzle-orm";
2 
3import { RunId } from "@/contracts";
4import { isTerminalStatus, type EnsureRunInput, type RunMetaState, type UpdateRunStateInput } from "@/worker/contracts";
5import * as runSchema from "@/worker/db/durable/schema/run-do";
6 
7import { type RunDb, type TryUpdateRunStateResult, toRunMetaState } from "./core";
8import { resolveRunStateUpdate, RunStateTransitionError } from "../state";
9 
10export 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 
19export 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 
69export 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 
82export 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 
118export 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 
152export 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};