File
Blob: src/worker/db/d1/repositories/runs.ts
| 1 | import { and, desc, eq, lt, or, sql } from "drizzle-orm"; |
| 2 | |
| 3 | import { type D1DbExecutor, runIndex } from "@/worker/db/d1"; |
| 4 | import type { RunId } from "@/contracts"; |
| 5 | import { isTerminalStatus } from "@/worker/contracts"; |
| 6 | |
| 7 | export type RunIndexRow = typeof runIndex.$inferSelect; |
| 8 | export type NewRunIndexRow = typeof runIndex.$inferInsert; |
| 9 | |
| 10 | export interface RunPaginationCursor { |
| 11 | queuedAt: number; |
| 12 | runId: string; |
| 13 | } |
| 14 | |
| 15 | export interface RunTerminalValues { |
| 16 | status: string; |
| 17 | startedAt: number | null; |
| 18 | finishedAt: number | null; |
| 19 | exitCode: number | null; |
| 20 | } |
| 21 | |
| 22 | export const insertRunIndex = async (db: D1DbExecutor, row: NewRunIndexRow): Promise<void> => { |
| 23 | await db.insert(runIndex).values(row); |
| 24 | }; |
| 25 | |
| 26 | export 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 | |
| 59 | export 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 | |
| 64 | export 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 | |
| 73 | export 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 | }; |