Skip to content
File

Blob: src/worker/durable/project-do/transitions/shared.ts

typescript102 lines
1import { and, asc, eq } from "drizzle-orm";
2 
3import {
4 BranchName,
5 CommitSha,
6 DispatchMode,
7 ExecutionRuntime,
8 ProjectId,
9 RunId,
10 TriggerType,
11 UnixTimestampMs,
12 UserId,
13} from "@/contracts";
14import { D1SyncStatus, isTerminalStatus, nullableTrusted, expectTrusted } from "@/worker/contracts";
15import * as projectSchema from "@/worker/db/durable/schema/project-do";
16 
17import { getProjectStateRow } from "../repo";
18import type { ProjectDoContext, ProjectRunRow, ProjectStore } from "../types";
19 
20export 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 
33export 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 
47export { isTerminalStatus } from "@/worker/contracts";
48 
49export 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 
57export const nextMetadataD1SyncStatus = (current: D1SyncStatus): D1SyncStatus => {
58 if (current === "done" || current === "needs_terminal_update") {
59 return current;
60 }
61 
62 return "needs_update";
63};
64 
65export 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};