File
Blob: src/worker/durable/project-do/transitions/queue.ts
| 1 | import { eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { |
| 4 | type DispatchMode, |
| 5 | type ExecutionRuntime, |
| 6 | type UserId, |
| 7 | type BranchName, |
| 8 | RunId, |
| 9 | UnixTimestampMs, |
| 10 | } from "@/contracts"; |
| 11 | import { |
| 12 | type AcceptManualRunResult, |
| 13 | type AcceptQueuedRunInput, |
| 14 | type ClaimRunWorkInput, |
| 15 | type ClaimRunWorkResult, |
| 16 | type EnsureRunInput, |
| 17 | expectTrusted, |
| 18 | } from "@/worker/contracts"; |
| 19 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 20 | import { generateDurableEntityId } from "@/worker/services"; |
| 21 | |
| 22 | import { countQueuedRuns, getHighestQueuePosition, getProjectStateRow, getRunRow } from "../repo"; |
| 23 | import { ensureProjectState, getSnapshot } from "./shared"; |
| 24 | import type { ProjectDoContext, ProjectStore } from "../types"; |
| 25 | |
| 26 | type AcceptedManualRun = Extract<AcceptManualRunResult, { kind: "accepted" }>; |
| 27 | type RejectedManualRun = Extract<AcceptManualRunResult, { kind: "rejected" }>; |
| 28 | |
| 29 | export interface ResolvedAcceptManualRunInput { |
| 30 | projectId: AcceptQueuedRunInput["projectId"]; |
| 31 | triggeredByUserId: UserId; |
| 32 | branch: BranchName; |
| 33 | repoUrl: string; |
| 34 | configPath: string; |
| 35 | dispatchMode: DispatchMode; |
| 36 | executionRuntime: ExecutionRuntime; |
| 37 | } |
| 38 | |
| 39 | export interface AcceptedQueuedRunTransition extends AcceptedManualRun { |
| 40 | runInitialization: EnsureRunInput; |
| 41 | } |
| 42 | |
| 43 | export type AcceptQueuedRunTransition = AcceptedQueuedRunTransition | RejectedManualRun; |
| 44 | |
| 45 | export const transitionAcceptQueuedRun = ( |
| 46 | context: ProjectDoContext, |
| 47 | tx: ProjectStore, |
| 48 | input: AcceptQueuedRunInput, |
| 49 | currentTime: number, |
| 50 | ): AcceptQueuedRunTransition => { |
| 51 | ensureProjectState(context, tx, input.projectId); |
| 52 | |
| 53 | const queuedCount = countQueuedRuns(tx, input.projectId); |
| 54 | if (queuedCount >= 20) { |
| 55 | return { |
| 56 | kind: "rejected", |
| 57 | reason: "queue_full", |
| 58 | }; |
| 59 | } |
| 60 | |
| 61 | const runId = expectTrusted(RunId, generateDurableEntityId("run", currentTime), "RunId"); |
| 62 | const queuedAt = expectTrusted(UnixTimestampMs, currentTime, "UnixTimestampMs"); |
| 63 | const stateRow = getProjectStateRow(tx, input.projectId); |
| 64 | const executable = stateRow?.activeRunId === null && queuedCount === 0; |
| 65 | const position = getHighestQueuePosition(tx, input.projectId) + 1; |
| 66 | |
| 67 | tx.insert(projectSchema.projectRuns) |
| 68 | .values({ |
| 69 | id: runId, |
| 70 | projectId: input.projectId, |
| 71 | runId, |
| 72 | triggerType: input.triggerType, |
| 73 | triggeredByUserId: input.triggeredByUserId, |
| 74 | branch: input.branch, |
| 75 | commitSha: input.commitSha, |
| 76 | provider: input.provider, |
| 77 | deliveryId: input.deliveryId, |
| 78 | repoUrl: input.repoUrl, |
| 79 | configPath: input.configPath, |
| 80 | dispatchMode: input.dispatchMode, |
| 81 | executionRuntime: input.executionRuntime, |
| 82 | position, |
| 83 | status: executable ? "executable" : "pending", |
| 84 | d1SyncStatus: "needs_create", |
| 85 | dispatchStatus: executable ? "pending" : "blocked", |
| 86 | dispatchAttempts: 0, |
| 87 | lastError: null, |
| 88 | createdAt: currentTime, |
| 89 | cancelRequestedAt: null, |
| 90 | }) |
| 91 | .run(); |
| 92 | |
| 93 | return { |
| 94 | kind: "accepted", |
| 95 | runId, |
| 96 | queuedAt, |
| 97 | executable, |
| 98 | runInitialization: { |
| 99 | runId, |
| 100 | projectId: input.projectId, |
| 101 | triggerType: input.triggerType, |
| 102 | branch: input.branch, |
| 103 | commitSha: input.commitSha, |
| 104 | }, |
| 105 | }; |
| 106 | }; |
| 107 | |
| 108 | export const transitionAcceptManualRun = ( |
| 109 | context: ProjectDoContext, |
| 110 | tx: ProjectStore, |
| 111 | input: ResolvedAcceptManualRunInput, |
| 112 | currentTime: number, |
| 113 | ): AcceptQueuedRunTransition => |
| 114 | transitionAcceptQueuedRun( |
| 115 | context, |
| 116 | tx, |
| 117 | { |
| 118 | projectId: input.projectId, |
| 119 | triggerType: "manual", |
| 120 | triggeredByUserId: input.triggeredByUserId, |
| 121 | branch: input.branch, |
| 122 | commitSha: null, |
| 123 | repoUrl: input.repoUrl, |
| 124 | configPath: input.configPath, |
| 125 | provider: null, |
| 126 | deliveryId: null, |
| 127 | dispatchMode: input.dispatchMode, |
| 128 | executionRuntime: input.executionRuntime, |
| 129 | }, |
| 130 | currentTime, |
| 131 | ); |
| 132 | |
| 133 | export const transitionClaimRunWork = ( |
| 134 | context: ProjectDoContext, |
| 135 | tx: ProjectStore, |
| 136 | input: ClaimRunWorkInput, |
| 137 | ): ClaimRunWorkResult => { |
| 138 | ensureProjectState(context, tx, input.projectId); |
| 139 | |
| 140 | const row = getRunRow(tx, input.projectId, input.runId); |
| 141 | if (!row) { |
| 142 | return { kind: "stale", reason: "run_missing" }; |
| 143 | } |
| 144 | |
| 145 | if (row.status === "canceled") { |
| 146 | return { kind: "stale", reason: "canceled" }; |
| 147 | } |
| 148 | |
| 149 | if (row.status === "active" || row.status === "cancel_requested") { |
| 150 | return { kind: "stale", reason: "run_active" }; |
| 151 | } |
| 152 | |
| 153 | if (row.status === "passed" || row.status === "failed") { |
| 154 | return { kind: "stale", reason: "already_terminal" }; |
| 155 | } |
| 156 | |
| 157 | const stateRow = getProjectStateRow(tx, input.projectId); |
| 158 | if (stateRow?.activeRunId && stateRow.activeRunId !== input.runId) { |
| 159 | return { kind: "stale", reason: "superseded" }; |
| 160 | } |
| 161 | |
| 162 | if (row.status === "executable" && (row.dispatchStatus === "pending" || row.dispatchStatus === "queued")) { |
| 163 | tx.update(projectSchema.projectRuns) |
| 164 | .set({ |
| 165 | status: "active", |
| 166 | position: null, |
| 167 | dispatchStatus: "started", |
| 168 | }) |
| 169 | .where(eq(projectSchema.projectRuns.runId, input.runId)) |
| 170 | .run(); |
| 171 | tx.update(projectSchema.projectState) |
| 172 | .set({ |
| 173 | activeRunId: input.runId, |
| 174 | updatedAt: Date.now(), |
| 175 | }) |
| 176 | .where(eq(projectSchema.projectState.projectId, input.projectId)) |
| 177 | .run(); |
| 178 | |
| 179 | return { |
| 180 | kind: "execute", |
| 181 | snapshot: getSnapshot(row), |
| 182 | }; |
| 183 | } |
| 184 | |
| 185 | return { kind: "stale", reason: "not_currently_executable" }; |
| 186 | }; |