File
Blob: src/worker/durable/run-do/repo/core.ts
| 1 | import { type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite"; |
| 2 | |
| 3 | import { |
| 4 | LogStream, |
| 5 | OpaqueId, |
| 6 | ProjectId, |
| 7 | RunId, |
| 8 | RunStatus, |
| 9 | StepStatus, |
| 10 | TriggerType, |
| 11 | UnixTimestampMs, |
| 12 | BranchName, |
| 13 | CommitSha, |
| 14 | } from "@/contracts"; |
| 15 | import { |
| 16 | expectTrusted, |
| 17 | nullableTrusted, |
| 18 | PositiveInteger, |
| 19 | type RunLogRecord, |
| 20 | type RunMetaState, |
| 21 | type RunStepState, |
| 22 | } from "@/worker/contracts"; |
| 23 | import * as runSchema from "@/worker/db/durable/schema/run-do"; |
| 24 | |
| 25 | import type { RunStateTransitionConflictReason } from "../state"; |
| 26 | |
| 27 | export type RunDb = DrizzleSqliteDODatabase<typeof runSchema>; |
| 28 | export type RunTx = Parameters<RunDb["transaction"]>[0] extends (tx: infer T) => unknown ? T : never; |
| 29 | export type TryUpdateRunStateResult = |
| 30 | | { |
| 31 | kind: "applied"; |
| 32 | state: RunMetaState; |
| 33 | } |
| 34 | | { |
| 35 | kind: "conflict"; |
| 36 | reason: RunStateTransitionConflictReason; |
| 37 | current: RunMetaState; |
| 38 | }; |
| 39 | |
| 40 | export const toRunMetaState = (row: typeof runSchema.runMeta.$inferSelect): RunMetaState => ({ |
| 41 | runId: expectTrusted(RunId, row.id, "RunId"), |
| 42 | projectId: expectTrusted(ProjectId, row.projectId, "ProjectId"), |
| 43 | status: expectTrusted(RunStatus, row.status, "RunStatus"), |
| 44 | triggerType: expectTrusted(TriggerType, row.triggerType, "TriggerType"), |
| 45 | branch: expectTrusted(BranchName, row.branch, "BranchName"), |
| 46 | commitSha: nullableTrusted(CommitSha, row.commitSha, "CommitSha"), |
| 47 | currentStep: nullableTrusted(PositiveInteger, row.currentStep, "PositiveInteger"), |
| 48 | startedAt: nullableTrusted(UnixTimestampMs, row.startedAt, "UnixTimestampMs"), |
| 49 | finishedAt: nullableTrusted(UnixTimestampMs, row.finishedAt, "UnixTimestampMs"), |
| 50 | exitCode: row.exitCode, |
| 51 | errorMessage: row.errorMessage, |
| 52 | }); |
| 53 | |
| 54 | export const toRunStepState = (row: typeof runSchema.runSteps.$inferSelect): RunStepState => ({ |
| 55 | id: expectTrusted(OpaqueId, row.id, "OpaqueId"), |
| 56 | runId: expectTrusted(RunId, row.runId, "RunId"), |
| 57 | position: expectTrusted(PositiveInteger, row.position, "PositiveInteger"), |
| 58 | name: row.name, |
| 59 | command: row.command, |
| 60 | status: expectTrusted(StepStatus, row.status, "StepStatus"), |
| 61 | startedAt: nullableTrusted(UnixTimestampMs, row.startedAt, "UnixTimestampMs"), |
| 62 | finishedAt: nullableTrusted(UnixTimestampMs, row.finishedAt, "UnixTimestampMs"), |
| 63 | exitCode: row.exitCode, |
| 64 | }); |
| 65 | |
| 66 | export const toRunLogRecord = (row: typeof runSchema.runLogs.$inferSelect): RunLogRecord => ({ |
| 67 | id: expectTrusted(OpaqueId, row.id, "OpaqueId"), |
| 68 | runId: expectTrusted(RunId, row.runId, "RunId"), |
| 69 | seq: expectTrusted(PositiveInteger, row.seq, "PositiveInteger"), |
| 70 | stream: expectTrusted(LogStream, row.stream, "LogStream"), |
| 71 | chunk: row.chunk, |
| 72 | createdAt: expectTrusted(UnixTimestampMs, row.createdAt, "UnixTimestampMs"), |
| 73 | }); |
| 74 | |
| 75 | export const chunkRows = <TRow>(rows: readonly TRow[], maxRowsPerChunk: number): TRow[][] => { |
| 76 | const chunks: TRow[][] = []; |
| 77 | |
| 78 | for (let index = 0; index < rows.length; index += maxRowsPerChunk) { |
| 79 | chunks.push(rows.slice(index, index + maxRowsPerChunk)); |
| 80 | } |
| 81 | |
| 82 | return chunks; |
| 83 | }; |