File
Blob: src/worker/durable/project-do/transitions/shared.ts
| 1 | import { and, asc, eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { |
| 4 | BranchName, |
| 5 | CommitSha, |
| 6 | DispatchMode, |
| 7 | ExecutionRuntime, |
| 8 | ProjectId, |
| 9 | RunId, |
| 10 | TriggerType, |
| 11 | UnixTimestampMs, |
| 12 | UserId, |
| 13 | } from "@/contracts"; |
| 14 | import { D1SyncStatus, isTerminalStatus, nullableTrusted, expectTrusted } from "@/worker/contracts"; |
| 15 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 16 | |
| 17 | import { getProjectStateRow } from "../repo"; |
| 18 | import type { ProjectDoContext, ProjectRunRow, ProjectStore } from "../types"; |
| 19 | |
| 20 | export const ensureProjectState = (context: ProjectDoContext, tx: ProjectStore, projectId: ProjectId): void => { |
| 21 | context.cacheProjectId(projectId); |
| 22 | tx.insert(projectSchema.projectState) |
| 23 | .values({ |
| 24 | projectId, |
| 25 | activeRunId: null, |
| 26 | projectIndexSyncStatus: "current", |
| 27 | updatedAt: Date.now(), |
| 28 | }) |
| 29 | .onConflictDoNothing() |
| 30 | .run(); |
| 31 | }; |
| 32 | |
| 33 | export const getSnapshot = (row: ProjectRunRow) => ({ |
| 34 | runId: expectTrusted(RunId, row.runId, "RunId"), |
| 35 | projectId: expectTrusted(ProjectId, row.projectId, "ProjectId"), |
| 36 | triggerType: expectTrusted(TriggerType, row.triggerType, "TriggerType"), |
| 37 | triggeredByUserId: nullableTrusted(UserId, row.triggeredByUserId, "UserId"), |
| 38 | branch: expectTrusted(BranchName, row.branch, "BranchName"), |
| 39 | commitSha: nullableTrusted(CommitSha, row.commitSha, "CommitSha"), |
| 40 | repoUrl: row.repoUrl, |
| 41 | configPath: row.configPath, |
| 42 | dispatchMode: expectTrusted(DispatchMode, row.dispatchMode, "DispatchMode"), |
| 43 | executionRuntime: expectTrusted(ExecutionRuntime, row.executionRuntime, "ExecutionRuntime"), |
| 44 | queuedAt: expectTrusted(UnixTimestampMs, row.createdAt, "UnixTimestampMs"), |
| 45 | }); |
| 46 | |
| 47 | export { isTerminalStatus } from "@/worker/contracts"; |
| 48 | |
| 49 | export const nextTerminalD1SyncStatus = (current: D1SyncStatus): D1SyncStatus => { |
| 50 | if (current === "done") { |
| 51 | return "done"; |
| 52 | } |
| 53 | |
| 54 | return current === "needs_create" ? "needs_create" : "needs_terminal_update"; |
| 55 | }; |
| 56 | |
| 57 | export const nextMetadataD1SyncStatus = (current: D1SyncStatus): D1SyncStatus => { |
| 58 | if (current === "done" || current === "needs_terminal_update") { |
| 59 | return current; |
| 60 | } |
| 61 | |
| 62 | return "needs_update"; |
| 63 | }; |
| 64 | |
| 65 | export const promoteNextPendingRun = (tx: ProjectStore, projectId: ProjectId): RunId | null => { |
| 66 | const stateRow = getProjectStateRow(tx, projectId); |
| 67 | if (stateRow?.activeRunId) { |
| 68 | return null; |
| 69 | } |
| 70 | |
| 71 | const existingExecutable = tx |
| 72 | .select({ runId: projectSchema.projectRuns.runId }) |
| 73 | .from(projectSchema.projectRuns) |
| 74 | .where(and(eq(projectSchema.projectRuns.projectId, projectId), eq(projectSchema.projectRuns.status, "executable"))) |
| 75 | .limit(1) |
| 76 | .get(); |
| 77 | if (existingExecutable) { |
| 78 | return null; |
| 79 | } |
| 80 | |
| 81 | const nextRow = tx |
| 82 | .select({ runId: projectSchema.projectRuns.runId }) |
| 83 | .from(projectSchema.projectRuns) |
| 84 | .where(and(eq(projectSchema.projectRuns.projectId, projectId), eq(projectSchema.projectRuns.status, "pending"))) |
| 85 | .orderBy(asc(projectSchema.projectRuns.position)) |
| 86 | .limit(1) |
| 87 | .get(); |
| 88 | if (!nextRow) { |
| 89 | return null; |
| 90 | } |
| 91 | |
| 92 | tx.update(projectSchema.projectRuns) |
| 93 | .set({ |
| 94 | status: "executable", |
| 95 | dispatchStatus: "pending", |
| 96 | }) |
| 97 | .where(eq(projectSchema.projectRuns.runId, nextRow.runId)) |
| 98 | .run(); |
| 99 | |
| 100 | return expectTrusted(RunId, nextRow.runId, "RunId"); |
| 101 | }; |