Skip to content
File

Blob: src/worker/durable/project-do/repo.ts

typescript171 lines
1import { and, asc, desc, eq, isNotNull, or } from "drizzle-orm";
2 
3import type { ProjectId, RunId } from "@/contracts";
4import type { D1SyncStatus } from "@/worker/contracts";
5import * as projectSchema from "@/worker/db/durable/schema/project-do";
6 
7import type { ProjectConfigRow, ProjectDoContext, ProjectRunRow, ProjectStateRow, ProjectStore } from "./types";
8 
9export const getHighestQueuePosition = (tx: ProjectStore, projectId: ProjectId): number => {
10 const row = tx
11 .select({ position: projectSchema.projectRuns.position })
12 .from(projectSchema.projectRuns)
13 .where(and(eq(projectSchema.projectRuns.projectId, projectId), isNotNull(projectSchema.projectRuns.position)))
14 .orderBy(desc(projectSchema.projectRuns.position))
15 .limit(1)
16 .get();
17 
18 return row?.position ?? 0;
19};
20 
21export const countQueuedRuns = (tx: ProjectStore, projectId: ProjectId): number => {
22 const rows = tx
23 .select({ runId: projectSchema.projectRuns.runId })
24 .from(projectSchema.projectRuns)
25 .where(
26 and(
27 eq(projectSchema.projectRuns.projectId, projectId),
28 or(eq(projectSchema.projectRuns.status, "pending"), eq(projectSchema.projectRuns.status, "executable")),
29 ),
30 )
31 .all();
32 
33 return rows.length;
34};
35 
36export const getProjectStateRow = (tx: ProjectStore, projectId: ProjectId): ProjectStateRow | undefined =>
37 tx
38 .select()
39 .from(projectSchema.projectState)
40 .where(eq(projectSchema.projectState.projectId, projectId))
41 .limit(1)
42 .get();
43 
44export const getProjectConfigRow = (tx: ProjectStore, projectId: ProjectId): ProjectConfigRow | undefined =>
45 tx
46 .select()
47 .from(projectSchema.projectConfig)
48 .where(eq(projectSchema.projectConfig.projectId, projectId))
49 .limit(1)
50 .get();
51 
52export const getRunRow = (tx: ProjectStore, projectId: ProjectId, runId: RunId): ProjectRunRow | undefined =>
53 tx
54 .select()
55 .from(projectSchema.projectRuns)
56 .where(and(eq(projectSchema.projectRuns.projectId, projectId), eq(projectSchema.projectRuns.runId, runId)))
57 .limit(1)
58 .get();
59 
60export const getProjectState = async (
61 context: ProjectDoContext,
62 projectId: ProjectId,
63): Promise<ProjectStateRow | undefined> => {
64 const rows = await context.db
65 .select()
66 .from(projectSchema.projectState)
67 .where(eq(projectSchema.projectState.projectId, projectId))
68 .limit(1);
69 
70 return rows[0];
71};
72 
73export const getProjectConfig = async (
74 context: ProjectDoContext,
75 projectId: ProjectId,
76): Promise<ProjectConfigRow | undefined> => {
77 const rows = await context.db
78 .select()
79 .from(projectSchema.projectConfig)
80 .where(eq(projectSchema.projectConfig.projectId, projectId))
81 .limit(1);
82 
83 return rows[0];
84};
85 
86export const getRunRowByRunId = async (context: ProjectDoContext, runId: RunId): Promise<ProjectRunRow | undefined> => {
87 const rows = await context.db
88 .select()
89 .from(projectSchema.projectRuns)
90 .where(eq(projectSchema.projectRuns.runId, runId))
91 .limit(1);
92 
93 return rows[0];
94};
95 
96export const listProjectRuns = async (context: ProjectDoContext, projectId: ProjectId): Promise<ProjectRunRow[]> =>
97 context.db
98 .select()
99 .from(projectSchema.projectRuns)
100 .where(eq(projectSchema.projectRuns.projectId, projectId))
101 .orderBy(asc(projectSchema.projectRuns.position));
102 
103export const listPendingProjectDetailRows = async (context: ProjectDoContext, projectId: ProjectId) =>
104 context.db
105 .select({
106 runId: projectSchema.projectRuns.runId,
107 branch: projectSchema.projectRuns.branch,
108 queuedAt: projectSchema.projectRuns.createdAt,
109 })
110 .from(projectSchema.projectRuns)
111 .where(
112 and(
113 eq(projectSchema.projectRuns.projectId, projectId),
114 or(eq(projectSchema.projectRuns.status, "pending"), eq(projectSchema.projectRuns.status, "executable")),
115 ),
116 )
117 .orderBy(asc(projectSchema.projectRuns.position));
118 
119export const getOldestCancelRequestedRun = async (
120 context: ProjectDoContext,
121 projectId: ProjectId,
122): Promise<ProjectRunRow | undefined> => {
123 const rows = await context.db
124 .select()
125 .from(projectSchema.projectRuns)
126 .where(
127 and(eq(projectSchema.projectRuns.projectId, projectId), eq(projectSchema.projectRuns.status, "cancel_requested")),
128 )
129 .orderBy(asc(projectSchema.projectRuns.cancelRequestedAt), asc(projectSchema.projectRuns.createdAt))
130 .limit(1);
131 
132 return rows[0];
133};
134 
135export const getOldestRunByD1SyncStatus = async (
136 context: ProjectDoContext,
137 projectId: ProjectId,
138 d1SyncStatus: Extract<D1SyncStatus, "needs_create" | "needs_update" | "needs_terminal_update">,
139): Promise<ProjectRunRow | undefined> => {
140 const rows = await context.db
141 .select()
142 .from(projectSchema.projectRuns)
143 .where(
144 and(eq(projectSchema.projectRuns.projectId, projectId), eq(projectSchema.projectRuns.d1SyncStatus, d1SyncStatus)),
145 )
146 .orderBy(asc(projectSchema.projectRuns.createdAt))
147 .limit(1);
148 
149 return rows[0];
150};
151 
152export const getNextDispatchableRun = async (
153 context: ProjectDoContext,
154 projectId: ProjectId,
155): Promise<ProjectRunRow | undefined> => {
156 const rows = await context.db
157 .select()
158 .from(projectSchema.projectRuns)
159 .where(
160 and(
161 eq(projectSchema.projectRuns.projectId, projectId),
162 eq(projectSchema.projectRuns.status, "executable"),
163 eq(projectSchema.projectRuns.dispatchStatus, "pending"),
164 ),
165 )
166 .orderBy(asc(projectSchema.projectRuns.position))
167 .limit(1);
168 
169 return rows[0];
170};