Skip to content
File

Blob: src/worker/durable/project-do/reconciliation/d1-sync.ts

typescript297 lines
1import { eq } from "drizzle-orm";
2 
3import { type ProjectId, RunId } from "@/contracts";
4import { D1SyncStatus, expectTrusted } from "@/worker/contracts";
5import { createD1Db } from "@/worker/db/d1";
6import { updateProjectIndexReplicaById, upsertRunIndex } from "@/worker/db/d1/repositories";
7import * as projectSchema from "@/worker/db/durable/schema/project-do";
8 
9import { D1_RETRY_DELAYS_MS } from "../constants";
10import { getOldestRunByD1SyncStatus, getProjectConfig, getProjectState, getRunRowByRunId } from "../repo";
11import { ensureRunInitialized, setRunDoTerminal } from "../run-do-sync";
12import {
13 getD1RetryState,
14 getProjectIndexRetryState,
15 setD1RetryState,
16 setProjectIndexRetryState,
17 shouldWaitForD1Retry,
18} from "../sidecar-state";
19import {
20 buildLatestQueuedRunIndexTruth,
21 getNextD1RetryAttempt,
22 getRetryDelay,
23 getRunStub,
24 resolvePostQueuedSyncD1Status,
25 type ProjectRunTableRow,
26} from "./shared";
27import { isTerminalStatus } from "../transitions/shared";
28import type { ProjectDoContext } from "../types";
29 
30const getNextProjectIndexRetryAttempt = (retryState: Awaited<ReturnType<typeof getProjectIndexRetryState>>): number =>
31 retryState ? retryState.attempt + 1 : 1;
32 
33export const reconcileProjectIndexReplicaSync = async (
34 context: ProjectDoContext,
35 projectId: ProjectId,
36): Promise<ProjectId | null> => {
37 const stateRow = await getProjectState(context, projectId);
38 if (!stateRow || stateRow.projectIndexSyncStatus !== "needs_update") {
39 return null;
40 }
41 
42 const retryState = await getProjectIndexRetryState(context);
43 if (retryState && retryState.nextAt > Date.now()) {
44 return null;
45 }
46 
47 try {
48 const configRow = await getProjectConfig(context, projectId);
49 if (!configRow) {
50 throw new Error(`Project config ${projectId} is missing during read replica sync.`);
51 }
52 
53 const db = createD1Db(context.env.DB.withSession("first-primary"));
54 const updated = await updateProjectIndexReplicaById(db, projectId, {
55 name: configRow.name,
56 repoUrl: configRow.repoUrl,
57 defaultBranch: configRow.defaultBranch,
58 configPath: configRow.configPath,
59 updatedAt: configRow.updatedAt,
60 });
61 if (!updated) {
62 throw new Error(`Project index ${projectId} is missing during read replica sync.`);
63 }
64 
65 await context.db
66 .update(projectSchema.projectState)
67 .set({ projectIndexSyncStatus: "current" })
68 .where(eq(projectSchema.projectState.projectId, projectId));
69 await setProjectIndexRetryState(context, null);
70 return projectId;
71 } catch (error) {
72 const attempt = getNextProjectIndexRetryAttempt(retryState);
73 const nextAt = Date.now() + getRetryDelay(attempt, D1_RETRY_DELAYS_MS);
74 await setProjectIndexRetryState(context, {
75 attempt,
76 nextAt,
77 phase: "project_index",
78 });
79 context.logger.error("project_index_sync_failed", {
80 projectId,
81 attempt,
82 error: error instanceof Error ? error.message : String(error),
83 });
84 return null;
85 }
86};
87 
88export const syncRunMetadataToD1 = async (
89 context: ProjectDoContext,
90 projectId: ProjectId,
91 runId: RunId,
92 rowOverride?: ProjectRunTableRow,
93): Promise<RunId | null> => {
94 const row = rowOverride ?? (await getRunRowByRunId(context, runId));
95 if (!row) {
96 return null;
97 }
98 
99 const retryState = await getD1RetryState(context, runId);
100 
101 try {
102 const truth = await buildLatestQueuedRunIndexTruth(context, runId, row);
103 if (!truth) {
104 return null;
105 }
106 
107 const db = createD1Db(context.env.DB.withSession("first-primary"));
108 await upsertRunIndex(db, truth.truthRow);
109 
110 let latestRow = await getRunRowByRunId(context, runId);
111 if (!latestRow) {
112 throw new Error(`Run ${runId} disappeared after D1 metadata sync.`);
113 }
114 
115 const latestD1SyncStatus = expectTrusted(D1SyncStatus, latestRow.d1SyncStatus, "D1SyncStatus");
116 let preserveMetadataRetry = false;
117 if (latestD1SyncStatus === "needs_update") {
118 try {
119 await ensureRunInitialized(context, latestRow);
120 latestRow = (await getRunRowByRunId(context, runId)) ?? latestRow;
121 } catch (error) {
122 latestRow = (await getRunRowByRunId(context, runId)) ?? latestRow;
123 preserveMetadataRetry = true;
124 
125 const attempt = getNextD1RetryAttempt(retryState, "metadata");
126 const nextAt = Date.now() + getRetryDelay(attempt, D1_RETRY_DELAYS_MS);
127 await setD1RetryState(context, runId, { attempt, nextAt, phase: "metadata" });
128 context.logger.error("run_do_metadata_sync_failed", {
129 projectId,
130 runId,
131 attempt,
132 error: error instanceof Error ? error.message : String(error),
133 });
134 }
135 }
136 
137 const nextD1SyncStatus = resolvePostQueuedSyncD1Status(latestRow, truth.truthRow, preserveMetadataRetry);
138 
139 await context.db
140 .update(projectSchema.projectRuns)
141 .set({ d1SyncStatus: nextD1SyncStatus })
142 .where(eq(projectSchema.projectRuns.runId, row.runId));
143 
144 if (nextD1SyncStatus !== "needs_update") {
145 await setD1RetryState(context, runId, null);
146 }
147 
148 return runId;
149 } catch (error) {
150 const attempt = getNextD1RetryAttempt(retryState, "metadata");
151 const nextAt = Date.now() + getRetryDelay(attempt, D1_RETRY_DELAYS_MS);
152 await setD1RetryState(context, runId, { attempt, nextAt, phase: "metadata" });
153 context.logger.error("d1_metadata_sync_failed", {
154 projectId,
155 runId,
156 attempt,
157 error: error instanceof Error ? error.message : String(error),
158 });
159 return null;
160 }
161};
162 
163export const reconcileAcceptedRunD1Sync = async (
164 context: ProjectDoContext,
165 projectId: ProjectId,
166): Promise<RunId | null> => {
167 const row = await getOldestRunByD1SyncStatus(context, projectId, "needs_create");
168 if (!row) {
169 return null;
170 }
171 
172 const runId = expectTrusted(RunId, row.runId, "RunId");
173 const retryState = await getD1RetryState(context, runId);
174 if (shouldWaitForD1Retry(retryState, "create")) {
175 return null;
176 }
177 
178 try {
179 const truth = await buildLatestQueuedRunIndexTruth(context, runId, row);
180 if (!truth) {
181 return null;
182 }
183 
184 const db = createD1Db(context.env.DB.withSession("first-primary"));
185 await upsertRunIndex(db, truth.truthRow);
186 
187 const latestRow = await getRunRowByRunId(context, runId);
188 if (!latestRow) {
189 throw new Error(`Run ${runId} disappeared after D1 create sync.`);
190 }
191 
192 await context.db
193 .update(projectSchema.projectRuns)
194 .set({ d1SyncStatus: resolvePostQueuedSyncD1Status(latestRow, truth.truthRow) })
195 .where(eq(projectSchema.projectRuns.runId, row.runId));
196 await setD1RetryState(context, runId, null);
197 return runId;
198 } catch (error) {
199 const attempt = getNextD1RetryAttempt(retryState, "create");
200 const nextAt = Date.now() + getRetryDelay(attempt, D1_RETRY_DELAYS_MS);
201 await setD1RetryState(context, runId, { attempt, nextAt, phase: "create" });
202 context.logger.error("d1_create_sync_failed", {
203 projectId,
204 runId,
205 attempt,
206 error: error instanceof Error ? error.message : String(error),
207 });
208 return null;
209 }
210};
211 
212export const reconcileRunMetadataD1Sync = async (
213 context: ProjectDoContext,
214 projectId: ProjectId,
215): Promise<RunId | null> => {
216 const row = await getOldestRunByD1SyncStatus(context, projectId, "needs_update");
217 if (!row) {
218 return null;
219 }
220 
221 const runId = expectTrusted(RunId, row.runId, "RunId");
222 const retryState = await getD1RetryState(context, runId);
223 if (shouldWaitForD1Retry(retryState, "metadata")) {
224 return null;
225 }
226 
227 return await syncRunMetadataToD1(context, projectId, runId, row);
228};
229 
230export const reconcileTerminalRunD1Sync = async (
231 context: ProjectDoContext,
232 projectId: ProjectId,
233): Promise<RunId | null> => {
234 const row = await getOldestRunByD1SyncStatus(context, projectId, "needs_terminal_update");
235 if (!row) {
236 return null;
237 }
238 
239 const runId = expectTrusted(RunId, row.runId, "RunId");
240 const retryState = await getD1RetryState(context, runId);
241 if (shouldWaitForD1Retry(retryState, "terminal")) {
242 return null;
243 }
244 
245 try {
246 const runStub = getRunStub(context, runId);
247 let runMeta = await runStub.getRunSummary(runId);
248 
249 if (
250 (!runMeta || !isTerminalStatus(runMeta.status) || runMeta.finishedAt === null) &&
251 isTerminalStatus(row.status)
252 ) {
253 await setRunDoTerminal(context, row, row.status, row.lastError);
254 runMeta = await runStub.getRunSummary(runId);
255 }
256 
257 if (!runMeta || !isTerminalStatus(runMeta.status) || runMeta.finishedAt === null) {
258 throw new Error(`Run ${runId} is missing terminal RunDO metadata.`);
259 }
260 
261 const db = createD1Db(context.env.DB.withSession("first-primary"));
262 await upsertRunIndex(db, {
263 id: row.runId,
264 projectId: row.projectId,
265 triggeredByUserId: row.triggeredByUserId,
266 triggerType: row.triggerType,
267 branch: row.branch,
268 commitSha: row.commitSha,
269 dispatchMode: row.dispatchMode,
270 executionRuntime: row.executionRuntime,
271 status: runMeta.status,
272 queuedAt: row.createdAt,
273 startedAt: runMeta.startedAt,
274 finishedAt: runMeta.finishedAt,
275 exitCode: runMeta.exitCode,
276 });
277 
278 await context.db
279 .update(projectSchema.projectRuns)
280 .set({ d1SyncStatus: "done" })
281 .where(eq(projectSchema.projectRuns.runId, row.runId));
282 await setD1RetryState(context, runId, null);
283 return runId;
284 } catch (error) {
285 const attempt = getNextD1RetryAttempt(retryState, "terminal");
286 const nextAt = Date.now() + getRetryDelay(attempt, D1_RETRY_DELAYS_MS);
287 await setD1RetryState(context, runId, { attempt, nextAt, phase: "terminal" });
288 context.logger.error("d1_terminal_sync_failed", {
289 projectId,
290 runId,
291 attempt,
292 error: error instanceof Error ? error.message : String(error),
293 });
294 return null;
295 }
296};