Skip to content
File

Blob: src/worker/db/d1/repositories/runs.ts

typescript92 lines
1import { and, desc, eq, lt, or, sql } from "drizzle-orm";
2 
3import { type D1DbExecutor, runIndex } from "@/worker/db/d1";
4import type { RunId } from "@/contracts";
5import { isTerminalStatus } from "@/worker/contracts";
6 
7export type RunIndexRow = typeof runIndex.$inferSelect;
8export type NewRunIndexRow = typeof runIndex.$inferInsert;
9 
10export interface RunPaginationCursor {
11 queuedAt: number;
12 runId: string;
13}
14 
15export interface RunTerminalValues {
16 status: string;
17 startedAt: number | null;
18 finishedAt: number | null;
19 exitCode: number | null;
20}
21 
22export const insertRunIndex = async (db: D1DbExecutor, row: NewRunIndexRow): Promise<void> => {
23 await db.insert(runIndex).values(row);
24};
25 
26export const upsertRunIndex = async (db: D1DbExecutor, row: NewRunIndexRow): Promise<void> => {
27 const preserveExistingTerminalFields = !isTerminalStatus(row.status);
28 
29 await db
30 .insert(runIndex)
31 .values(row)
32 .onConflictDoUpdate({
33 target: runIndex.id,
34 set: {
35 projectId: row.projectId,
36 triggeredByUserId: row.triggeredByUserId,
37 triggerType: row.triggerType,
38 branch: row.branch,
39 commitSha: sql`coalesce(${row.commitSha}, ${runIndex.commitSha})`,
40 dispatchMode: row.dispatchMode,
41 executionRuntime: row.executionRuntime,
42 status: preserveExistingTerminalFields
43 ? sql`case when ${runIndex.status} in ('passed', 'failed', 'canceled') then ${runIndex.status} else ${row.status} end`
44 : row.status,
45 queuedAt: row.queuedAt,
46 startedAt: preserveExistingTerminalFields
47 ? sql`case when ${runIndex.status} in ('passed', 'failed', 'canceled') then ${runIndex.startedAt} else ${row.startedAt} end`
48 : row.startedAt,
49 finishedAt: preserveExistingTerminalFields
50 ? sql`case when ${runIndex.status} in ('passed', 'failed', 'canceled') then ${runIndex.finishedAt} else ${row.finishedAt} end`
51 : row.finishedAt,
52 exitCode: preserveExistingTerminalFields
53 ? sql`case when ${runIndex.status} in ('passed', 'failed', 'canceled') then ${runIndex.exitCode} else ${row.exitCode} end`
54 : row.exitCode,
55 },
56 });
57};
58 
59export const findRunIndexById = async (db: D1DbExecutor, runId: RunId): Promise<RunIndexRow | undefined> => {
60 const rows = await db.select().from(runIndex).where(eq(runIndex.id, runId)).limit(1);
61 return rows[0];
62};
63 
64export const updateRunIndexTerminal = async (
65 db: D1DbExecutor,
66 runId: string,
67 values: RunTerminalValues,
68): Promise<boolean> => {
69 const rows = await db.update(runIndex).set(values).where(eq(runIndex.id, runId)).returning({ id: runIndex.id });
70 return rows.length === 1;
71};
72 
73export const listProjectRunsPage = async (
74 db: D1DbExecutor,
75 projectId: string,
76 limit: number,
77 cursor?: RunPaginationCursor,
78): Promise<RunIndexRow[]> => {
79 const where =
80 cursor === undefined
81 ? eq(runIndex.projectId, projectId)
82 : and(
83 eq(runIndex.projectId, projectId),
84 or(
85 lt(runIndex.queuedAt, cursor.queuedAt),
86 and(eq(runIndex.queuedAt, cursor.queuedAt), lt(runIndex.id, cursor.runId)),
87 ),
88 );
89 
90 return db.select().from(runIndex).where(where).orderBy(desc(runIndex.queuedAt), desc(runIndex.id)).limit(limit);
91};