File
Blob: src/worker/durable/project-do/commands.ts
| 1 | import { and, eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { |
| 4 | BranchName, |
| 5 | DispatchMode, |
| 6 | ExecutionRuntime, |
| 7 | type ProjectId, |
| 8 | RunId, |
| 9 | TriggerType, |
| 10 | UnixTimestampMs, |
| 11 | type WebhookProvider, |
| 12 | } from "@/contracts"; |
| 13 | import { |
| 14 | type AcceptManualRunResult, |
| 15 | type AcceptManualRunInput, |
| 16 | type ClaimRunWorkInput, |
| 17 | type ClaimRunWorkResult, |
| 18 | D1SyncStatus, |
| 19 | nullableTrusted, |
| 20 | type PendingRunState, |
| 21 | type ProjectDetailState, |
| 22 | ProjectRunStatus, |
| 23 | type FinalizeRunExecutionInput, |
| 24 | type FinalizeRunExecutionResult, |
| 25 | type RecoverWorkflowDispatchFailureInput, |
| 26 | type RecoverWorkflowDispatchFailureResult, |
| 27 | type RecordRunResolvedCommitResult, |
| 28 | type RequestRunCancelInput, |
| 29 | type RecordRunResolvedCommitInput, |
| 30 | type RequestRunCancelResult, |
| 31 | type RunHeartbeatInput, |
| 32 | type RunHeartbeatResult, |
| 33 | expectTrusted, |
| 34 | } from "@/worker/contracts"; |
| 35 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 36 | |
| 37 | import { DISPATCH_RETRY_DELAYS_MS, SANDBOX_CLEANUP_RETRY_DELAYS_MS } from "./constants"; |
| 38 | import { getProjectConfigRow, getRunRow, listPendingProjectDetailRows } from "./repo"; |
| 39 | import { |
| 40 | getProjectConfigState, |
| 41 | getProjectExecutionMaterial as getProjectExecutionMaterialState, |
| 42 | getProjectWebhookIngressState as getProjectWebhookIngressStateState, |
| 43 | initializeProject as initializeProjectConfig, |
| 44 | updateProjectConfig as updateProjectConfigState, |
| 45 | } from "./project-config"; |
| 46 | import { getRetryDelay, syncRunMetadataToD1 } from "./reconciliation/index"; |
| 47 | import { |
| 48 | ensureRunInitialized, |
| 49 | ensureRunInitializedWithPayload, |
| 50 | setRunDoTerminal, |
| 51 | updateRunDoCancelRequested, |
| 52 | } from "./run-do-sync"; |
| 53 | import { |
| 54 | armReconciliation, |
| 55 | rescheduleAlarmInTransaction, |
| 56 | sandboxCleanupRetryKey, |
| 57 | setDispatchRetryAt, |
| 58 | setHeartbeatAt, |
| 59 | } from "./sidecar-state"; |
| 60 | import { |
| 61 | ensureProjectState, |
| 62 | nextTerminalD1SyncStatus, |
| 63 | promoteNextPendingRun, |
| 64 | transitionAcceptManualRun, |
| 65 | transitionClaimRunWork, |
| 66 | transitionFinalizeRunExecution, |
| 67 | transitionRecordRunResolvedCommit, |
| 68 | transitionRequestRunCancel, |
| 69 | } from "./transitions"; |
| 70 | import type { |
| 71 | InitializeProjectInput, |
| 72 | ProjectConfigState, |
| 73 | ProjectDoContext, |
| 74 | ProjectExecutionMaterial, |
| 75 | ProjectWebhookIngressState, |
| 76 | UpdateProjectConfigInput, |
| 77 | UpdateProjectConfigResult, |
| 78 | } from "./types"; |
| 79 | |
| 80 | const now = (): UnixTimestampMs => expectTrusted(UnixTimestampMs, Date.now(), "UnixTimestampMs"); |
| 81 | const runBestEffortSidecar = async ( |
| 82 | context: ProjectDoContext, |
| 83 | operation: string, |
| 84 | projectId: ProjectId, |
| 85 | runId: RunId | null, |
| 86 | effect: () => Promise<void>, |
| 87 | ): Promise<void> => { |
| 88 | try { |
| 89 | await effect(); |
| 90 | } catch (error) { |
| 91 | context.logger.warn("project_do_sidecar_write_failed", { |
| 92 | operation, |
| 93 | projectId, |
| 94 | runId, |
| 95 | error: error instanceof Error ? error.message : String(error), |
| 96 | }); |
| 97 | } |
| 98 | }; |
| 99 | const runCriticalReconciliationMutation = async <T>( |
| 100 | context: ProjectDoContext, |
| 101 | projectId: ProjectId, |
| 102 | mutation: (txn: DurableObjectTransaction) => Promise<T> | T, |
| 103 | ): Promise<T> => |
| 104 | context.ctx.storage.transaction(async (txn) => { |
| 105 | const result = await mutation(txn); |
| 106 | await rescheduleAlarmInTransaction(context, txn, projectId); |
| 107 | return result; |
| 108 | }); |
| 109 | |
| 110 | // Keep these handlers aligned with the public ProjectDO RPC contract. HTTP handlers and the queue |
| 111 | // consumer call the shell methods by name, so this module only moves implementation behind that surface. |
| 112 | export const getProjectDetailState = async ( |
| 113 | context: ProjectDoContext, |
| 114 | projectId: ProjectId, |
| 115 | ): Promise<ProjectDetailState> => { |
| 116 | context.db.transaction((tx) => { |
| 117 | ensureProjectState(context, tx, projectId); |
| 118 | }); |
| 119 | |
| 120 | const stateRows = await context.db |
| 121 | .select() |
| 122 | .from(projectSchema.projectState) |
| 123 | .where(eq(projectSchema.projectState.projectId, projectId)) |
| 124 | .limit(1); |
| 125 | const pendingRows = await listPendingProjectDetailRows(context, projectId); |
| 126 | |
| 127 | const pendingRuns: PendingRunState[] = pendingRows.map((pendingRun) => ({ |
| 128 | runId: expectTrusted(RunId, pendingRun.runId, "RunId"), |
| 129 | branch: expectTrusted(BranchName, pendingRun.branch, "BranchName"), |
| 130 | queuedAt: expectTrusted(UnixTimestampMs, pendingRun.queuedAt, "UnixTimestampMs"), |
| 131 | })); |
| 132 | |
| 133 | return { |
| 134 | activeRunId: nullableTrusted(RunId, stateRows[0]?.activeRunId ?? null, "RunId"), |
| 135 | pendingRuns, |
| 136 | }; |
| 137 | }; |
| 138 | |
| 139 | export const acceptManualRun = async ( |
| 140 | context: ProjectDoContext, |
| 141 | input: AcceptManualRunInput, |
| 142 | ): Promise<AcceptManualRunResult> => { |
| 143 | const currentTime = Date.now(); |
| 144 | const transition = await context.ctx.storage.transaction(async (txn) => { |
| 145 | const projectConfigRow = getProjectConfigRow(context.db, input.projectId); |
| 146 | if (!projectConfigRow) { |
| 147 | throw new Error(`Project config ${input.projectId} is missing during manual run acceptance.`); |
| 148 | } |
| 149 | |
| 150 | const result = transitionAcceptManualRun( |
| 151 | context, |
| 152 | context.db, |
| 153 | { |
| 154 | projectId: input.projectId, |
| 155 | triggeredByUserId: input.triggeredByUserId, |
| 156 | branch: input.branch ?? expectTrusted(BranchName, projectConfigRow.defaultBranch, "BranchName"), |
| 157 | repoUrl: projectConfigRow.repoUrl, |
| 158 | configPath: projectConfigRow.configPath, |
| 159 | dispatchMode: expectTrusted(DispatchMode, projectConfigRow.dispatchMode, "DispatchMode"), |
| 160 | executionRuntime: expectTrusted(ExecutionRuntime, projectConfigRow.executionRuntime, "ExecutionRuntime"), |
| 161 | }, |
| 162 | currentTime, |
| 163 | ); |
| 164 | if (result.kind === "accepted") { |
| 165 | await rescheduleAlarmInTransaction(context, txn, input.projectId); |
| 166 | } |
| 167 | |
| 168 | return result; |
| 169 | }); |
| 170 | |
| 171 | if (transition.kind === "rejected") { |
| 172 | return transition; |
| 173 | } |
| 174 | |
| 175 | try { |
| 176 | await ensureRunInitializedWithPayload(context, transition.runInitialization); |
| 177 | } catch (error) { |
| 178 | context.logger.error("run_do_initialize_failed", { |
| 179 | projectId: input.projectId, |
| 180 | runId: transition.runId, |
| 181 | error: error instanceof Error ? error.message : String(error), |
| 182 | }); |
| 183 | } |
| 184 | |
| 185 | return { |
| 186 | kind: "accepted", |
| 187 | runId: transition.runId, |
| 188 | queuedAt: transition.queuedAt, |
| 189 | executable: transition.executable, |
| 190 | }; |
| 191 | }; |
| 192 | |
| 193 | export const claimRunWork = async ( |
| 194 | context: ProjectDoContext, |
| 195 | input: ClaimRunWorkInput, |
| 196 | ): Promise<ClaimRunWorkResult> => { |
| 197 | const result = await runCriticalReconciliationMutation(context, input.projectId, () => |
| 198 | transitionClaimRunWork(context, context.db, input), |
| 199 | ); |
| 200 | |
| 201 | if (result.kind === "execute") { |
| 202 | await runBestEffortSidecar(context, "set_heartbeat", input.projectId, result.snapshot.runId, async () => { |
| 203 | await setHeartbeatAt(context, result.snapshot.runId, Date.now()); |
| 204 | }); |
| 205 | } |
| 206 | |
| 207 | return result; |
| 208 | }; |
| 209 | |
| 210 | export const finalizeRunExecution = async ( |
| 211 | context: ProjectDoContext, |
| 212 | input: FinalizeRunExecutionInput, |
| 213 | ): Promise<FinalizeRunExecutionResult> => { |
| 214 | const result = await runCriticalReconciliationMutation(context, input.projectId, async (txn) => { |
| 215 | const transition = transitionFinalizeRunExecution(context, context.db, input); |
| 216 | |
| 217 | if (input.sandboxDestroyed) { |
| 218 | await txn.delete(sandboxCleanupRetryKey(input.runId)); |
| 219 | } else { |
| 220 | await txn.put(sandboxCleanupRetryKey(input.runId), { |
| 221 | attempt: 1, |
| 222 | nextAt: Date.now() + SANDBOX_CLEANUP_RETRY_DELAYS_MS[0], |
| 223 | }); |
| 224 | } |
| 225 | |
| 226 | return transition; |
| 227 | }); |
| 228 | |
| 229 | await runBestEffortSidecar(context, "clear_heartbeat", input.projectId, input.runId, async () => { |
| 230 | await setHeartbeatAt(context, input.runId, null); |
| 231 | }); |
| 232 | return result; |
| 233 | }; |
| 234 | |
| 235 | export const recoverWorkflowDispatchFailure = async ( |
| 236 | context: ProjectDoContext, |
| 237 | input: RecoverWorkflowDispatchFailureInput, |
| 238 | ): Promise<RecoverWorkflowDispatchFailureResult> => { |
| 239 | const result = await runCriticalReconciliationMutation(context, input.projectId, async () => { |
| 240 | const row = getRunRow(context.db, input.projectId, input.runId); |
| 241 | if (!row) { |
| 242 | return { |
| 243 | kind: "stale" as const, |
| 244 | }; |
| 245 | } |
| 246 | |
| 247 | if (row.status === "active" || row.status === "cancel_requested") { |
| 248 | return { |
| 249 | kind: "already_active" as const, |
| 250 | }; |
| 251 | } |
| 252 | |
| 253 | if (row.status !== "executable" || row.dispatchStatus !== "queued") { |
| 254 | return { |
| 255 | kind: "stale" as const, |
| 256 | }; |
| 257 | } |
| 258 | |
| 259 | const nextAttempt = row.dispatchAttempts + 1; |
| 260 | if (nextAttempt > DISPATCH_RETRY_DELAYS_MS.length) { |
| 261 | await context.db |
| 262 | .update(projectSchema.projectRuns) |
| 263 | .set({ |
| 264 | status: "failed", |
| 265 | position: null, |
| 266 | dispatchStatus: "terminal", |
| 267 | dispatchAttempts: nextAttempt, |
| 268 | d1SyncStatus: nextTerminalD1SyncStatus(expectTrusted(D1SyncStatus, row.d1SyncStatus, "D1SyncStatus")), |
| 269 | lastError: "dispatch_failed", |
| 270 | }) |
| 271 | .where(eq(projectSchema.projectRuns.runId, row.runId)); |
| 272 | promoteNextPendingRun(context.db, input.projectId); |
| 273 | |
| 274 | return { |
| 275 | kind: "terminal" as const, |
| 276 | row, |
| 277 | }; |
| 278 | } |
| 279 | |
| 280 | await context.db |
| 281 | .update(projectSchema.projectRuns) |
| 282 | .set({ |
| 283 | dispatchStatus: "pending", |
| 284 | dispatchAttempts: nextAttempt, |
| 285 | lastError: input.errorMessage, |
| 286 | }) |
| 287 | .where(eq(projectSchema.projectRuns.runId, row.runId)); |
| 288 | |
| 289 | return { |
| 290 | kind: "rearmed" as const, |
| 291 | nextAt: Date.now() + getRetryDelay(nextAttempt, DISPATCH_RETRY_DELAYS_MS), |
| 292 | }; |
| 293 | }); |
| 294 | |
| 295 | if (result.kind === "stale" || result.kind === "already_active") { |
| 296 | return result; |
| 297 | } |
| 298 | |
| 299 | if (result.kind === "terminal") { |
| 300 | await runBestEffortSidecar(context, "clear_dispatch_retry", input.projectId, input.runId, async () => { |
| 301 | await setDispatchRetryAt(context, input.runId, null); |
| 302 | }); |
| 303 | try { |
| 304 | await setRunDoTerminal(context, result.row, "failed", "dispatch_failed"); |
| 305 | } catch (error) { |
| 306 | context.logger.error("workflow_dispatch_failure_terminalize_failed", { |
| 307 | projectId: input.projectId, |
| 308 | runId: input.runId, |
| 309 | error: error instanceof Error ? error.message : String(error), |
| 310 | }); |
| 311 | } |
| 312 | return { |
| 313 | kind: "terminal", |
| 314 | }; |
| 315 | } |
| 316 | |
| 317 | await runBestEffortSidecar(context, "set_dispatch_retry", input.projectId, input.runId, async () => { |
| 318 | await setDispatchRetryAt(context, input.runId, result.nextAt); |
| 319 | }); |
| 320 | return { |
| 321 | kind: "rearmed", |
| 322 | }; |
| 323 | }; |
| 324 | |
| 325 | export const recordRunResolvedCommit = async ( |
| 326 | context: ProjectDoContext, |
| 327 | input: RecordRunResolvedCommitInput, |
| 328 | ): Promise<RecordRunResolvedCommitResult> => { |
| 329 | const transition = await runCriticalReconciliationMutation(context, input.projectId, () => |
| 330 | transitionRecordRunResolvedCommit(context, context.db, input), |
| 331 | ); |
| 332 | |
| 333 | if (transition.kind === "stale") { |
| 334 | return transition; |
| 335 | } |
| 336 | |
| 337 | await syncRunMetadataToD1(context, input.projectId, input.runId, transition.row); |
| 338 | return { |
| 339 | kind: "applied", |
| 340 | }; |
| 341 | }; |
| 342 | |
| 343 | export const requestRunCancel = async ( |
| 344 | context: ProjectDoContext, |
| 345 | input: RequestRunCancelInput, |
| 346 | ): Promise<RequestRunCancelResult> => { |
| 347 | const transition = await runCriticalReconciliationMutation(context, input.projectId, () => |
| 348 | transitionRequestRunCancel(context, context.db, input, now()), |
| 349 | ); |
| 350 | |
| 351 | if (transition.runDoAction === "canceled") { |
| 352 | await runBestEffortSidecar(context, "clear_dispatch_retry", input.projectId, input.runId, async () => { |
| 353 | await setDispatchRetryAt(context, input.runId, null); |
| 354 | }); |
| 355 | await runBestEffortSidecar(context, "clear_heartbeat", input.projectId, input.runId, async () => { |
| 356 | await setHeartbeatAt(context, input.runId, null); |
| 357 | }); |
| 358 | try { |
| 359 | await setRunDoTerminal(context, transition.row, "canceled", null); |
| 360 | } catch (error) { |
| 361 | context.logger.error("pending_cancel_run_do_terminalize_failed", { |
| 362 | projectId: input.projectId, |
| 363 | runId: input.runId, |
| 364 | error: error instanceof Error ? error.message : String(error), |
| 365 | }); |
| 366 | } |
| 367 | } else if (transition.runDoAction === "cancel_requested") { |
| 368 | try { |
| 369 | await updateRunDoCancelRequested(context, transition.row); |
| 370 | } catch (error) { |
| 371 | context.logger.error("active_cancel_run_do_update_failed", { |
| 372 | projectId: input.projectId, |
| 373 | runId: input.runId, |
| 374 | error: error instanceof Error ? error.message : String(error), |
| 375 | }); |
| 376 | } |
| 377 | } |
| 378 | |
| 379 | return { |
| 380 | runId: input.runId, |
| 381 | status: transition.runStatus, |
| 382 | cancelRequestedAt: transition.cancelRequestedAt, |
| 383 | }; |
| 384 | }; |
| 385 | |
| 386 | export const recordRunHeartbeat = async ( |
| 387 | context: ProjectDoContext, |
| 388 | input: RunHeartbeatInput, |
| 389 | ): Promise<RunHeartbeatResult> => { |
| 390 | const row = context.db.transaction((tx) => { |
| 391 | ensureProjectState(context, tx, input.projectId); |
| 392 | return tx |
| 393 | .select() |
| 394 | .from(projectSchema.projectRuns) |
| 395 | .where( |
| 396 | and(eq(projectSchema.projectRuns.projectId, input.projectId), eq(projectSchema.projectRuns.runId, input.runId)), |
| 397 | ) |
| 398 | .limit(1) |
| 399 | .get(); |
| 400 | }); |
| 401 | if (!row) { |
| 402 | return null; |
| 403 | } |
| 404 | |
| 405 | const control = { |
| 406 | runId: input.runId, |
| 407 | status: expectTrusted(ProjectRunStatus, row.status, "ProjectRunStatus"), |
| 408 | cancelRequestedAt: nullableTrusted(UnixTimestampMs, row.cancelRequestedAt, "UnixTimestampMs"), |
| 409 | }; |
| 410 | |
| 411 | if (row.status === "active" || row.status === "cancel_requested") { |
| 412 | await runBestEffortSidecar(context, "set_heartbeat", input.projectId, input.runId, async () => { |
| 413 | await setHeartbeatAt(context, input.runId, Date.now()); |
| 414 | }); |
| 415 | await runBestEffortSidecar(context, "arm_reconciliation", input.projectId, input.runId, async () => { |
| 416 | await armReconciliation(context, input.projectId); |
| 417 | }); |
| 418 | } |
| 419 | |
| 420 | return control; |
| 421 | }; |
| 422 | |
| 423 | export const initializeProject = async (context: ProjectDoContext, input: InitializeProjectInput): Promise<void> => { |
| 424 | await initializeProjectConfig(context, input); |
| 425 | }; |
| 426 | |
| 427 | export const getProjectConfig = async ( |
| 428 | context: ProjectDoContext, |
| 429 | projectId: ProjectId, |
| 430 | ): Promise<ProjectConfigState | null> => { |
| 431 | return await getProjectConfigState(context, projectId); |
| 432 | }; |
| 433 | |
| 434 | export const updateProjectConfig = async ( |
| 435 | context: ProjectDoContext, |
| 436 | input: UpdateProjectConfigInput, |
| 437 | ): Promise<UpdateProjectConfigResult> => { |
| 438 | return await updateProjectConfigState(context, input); |
| 439 | }; |
| 440 | |
| 441 | export const getProjectExecutionMaterial = async ( |
| 442 | context: ProjectDoContext, |
| 443 | projectId: ProjectId, |
| 444 | ): Promise<ProjectExecutionMaterial | null> => { |
| 445 | return await getProjectExecutionMaterialState(context, projectId); |
| 446 | }; |
| 447 | |
| 448 | export const getProjectWebhookIngressState = async ( |
| 449 | context: ProjectDoContext, |
| 450 | projectId: ProjectId, |
| 451 | provider: WebhookProvider, |
| 452 | ): Promise<ProjectWebhookIngressState | null> => { |
| 453 | return await getProjectWebhookIngressStateState(context, projectId, provider); |
| 454 | }; |