Skip to content
File

Blob: src/worker/durable/project-do/sidecar-state.ts

typescript335 lines
1import { type ProjectId, RunId } from "@/contracts";
2import { expectTrusted } from "@/worker/contracts";
3 
4import {
5 CANCEL_RECONCILE_RETRY_MS,
6 HEARTBEAT_STALE_AFTER_MS,
7 PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS,
8} from "./constants";
9import { getProjectState, listProjectRuns } from "./repo";
10import { isTerminalStatus } from "./transitions";
11import type {
12 D1RetryPhase,
13 D1RetryState,
14 ProjectDoContext,
15 ProjectIndexRetryState,
16 SandboxCleanupRetryState,
17} from "./types";
18 
19type SidecarStorage = Pick<DurableObjectStorage, "get" | "getAlarm" | "setAlarm" | "deleteAlarm">;
20 
21export const dispatchRetryKey = (runId: RunId): string => `run:${runId}:dispatch-retry-at`;
22 
23export const d1RetryKey = (runId: RunId): string => `run:${runId}:d1-retry`;
24 
25export const heartbeatKey = (runId: RunId): string => `run:${runId}:heartbeat-at`;
26export const sandboxCleanupRetryKey = (runId: RunId): string => `run:${runId}:sandbox-cleanup-retry`;
27 
28export const projectIndexRetryKey = (): string => "project:index-sync-retry";
29 
30const isD1RetryPhase = (value: unknown): value is D1RetryPhase =>
31 value === "create" || value === "metadata" || value === "terminal";
32 
33const isProjectIndexRetryPhase = (value: unknown): value is "project_index" => value === "project_index";
34 
35export const getD1RetryAtForPhase = (retryState: D1RetryState | null, phase: D1RetryPhase): number | null => {
36 if (!retryState) {
37 return null;
38 }
39 
40 return retryState.phase === phase ? retryState.nextAt : null;
41};
42 
43export const shouldWaitForD1Retry = (retryState: D1RetryState | null, phase: D1RetryPhase): boolean => {
44 const retryAt = getD1RetryAtForPhase(retryState, phase);
45 return retryAt !== null && retryAt > Date.now();
46};
47 
48const getDispatchRetryAtFromStorage = async (
49 storage: Pick<DurableObjectStorage, "get">,
50 runId: RunId,
51): Promise<number | null> => {
52 const value = await storage.get<number>(dispatchRetryKey(runId));
53 return typeof value === "number" ? value : null;
54};
55 
56export const getDispatchRetryAt = async (context: ProjectDoContext, runId: RunId): Promise<number | null> =>
57 getDispatchRetryAtFromStorage(context.ctx.storage, runId);
58 
59export const setDispatchRetryAt = async (
60 context: ProjectDoContext,
61 runId: RunId,
62 value: number | null,
63): Promise<void> => {
64 if (value === null) {
65 await context.ctx.storage.delete(dispatchRetryKey(runId));
66 return;
67 }
68 
69 await context.ctx.storage.put(dispatchRetryKey(runId), value);
70};
71 
72const getD1RetryStateFromStorage = async (
73 storage: Pick<DurableObjectStorage, "get">,
74 runId: RunId,
75): Promise<D1RetryState | null> => {
76 const value = await storage.get<D1RetryState>(d1RetryKey(runId));
77 if (!value || typeof value !== "object") {
78 return null;
79 }
80 
81 if (!("attempt" in value) || !("nextAt" in value)) {
82 return null;
83 }
84 
85 if (!("phase" in value) || !isD1RetryPhase(value.phase)) {
86 return null;
87 }
88 
89 return {
90 attempt: Number(value.attempt),
91 nextAt: Number(value.nextAt),
92 phase: value.phase,
93 };
94};
95 
96export const getD1RetryState = async (context: ProjectDoContext, runId: RunId): Promise<D1RetryState | null> =>
97 getD1RetryStateFromStorage(context.ctx.storage, runId);
98 
99export const setD1RetryState = async (
100 context: ProjectDoContext,
101 runId: RunId,
102 value: D1RetryState | null,
103): Promise<void> => {
104 if (value === null) {
105 await context.ctx.storage.delete(d1RetryKey(runId));
106 return;
107 }
108 
109 await context.ctx.storage.put(d1RetryKey(runId), value);
110};
111 
112const getProjectIndexRetryStateFromStorage = async (
113 storage: Pick<DurableObjectStorage, "get">,
114): Promise<ProjectIndexRetryState | null> => {
115 const value = await storage.get<ProjectIndexRetryState>(projectIndexRetryKey());
116 if (!value || typeof value !== "object") {
117 return null;
118 }
119 
120 if (!("attempt" in value) || !("nextAt" in value)) {
121 return null;
122 }
123 
124 if (!("phase" in value) || !isProjectIndexRetryPhase(value.phase)) {
125 return null;
126 }
127 
128 return {
129 attempt: Number(value.attempt),
130 nextAt: Number(value.nextAt),
131 phase: "project_index",
132 };
133};
134 
135export const getProjectIndexRetryState = async (context: ProjectDoContext): Promise<ProjectIndexRetryState | null> =>
136 getProjectIndexRetryStateFromStorage(context.ctx.storage);
137 
138export const setProjectIndexRetryState = async (
139 context: ProjectDoContext,
140 value: ProjectIndexRetryState | null,
141): Promise<void> => {
142 if (value === null) {
143 await context.ctx.storage.delete(projectIndexRetryKey());
144 return;
145 }
146 
147 await context.ctx.storage.put(projectIndexRetryKey(), value);
148};
149 
150const getHeartbeatAtFromStorage = async (
151 storage: Pick<DurableObjectStorage, "get">,
152 runId: RunId,
153): Promise<number | null> => {
154 const value = await storage.get<number>(heartbeatKey(runId));
155 return typeof value === "number" ? value : null;
156};
157 
158export const getHeartbeatAt = async (context: ProjectDoContext, runId: RunId): Promise<number | null> =>
159 getHeartbeatAtFromStorage(context.ctx.storage, runId);
160 
161export const setHeartbeatAt = async (context: ProjectDoContext, runId: RunId, value: number | null): Promise<void> => {
162 if (value === null) {
163 await context.ctx.storage.delete(heartbeatKey(runId));
164 return;
165 }
166 
167 await context.ctx.storage.put(heartbeatKey(runId), value);
168};
169 
170const getSandboxCleanupRetryStateFromStorage = async (
171 storage: Pick<DurableObjectStorage, "get">,
172 runId: RunId,
173): Promise<SandboxCleanupRetryState | null> => {
174 const value = await storage.get<SandboxCleanupRetryState>(sandboxCleanupRetryKey(runId));
175 if (!value || typeof value !== "object") {
176 return null;
177 }
178 
179 if (!("attempt" in value) || !("nextAt" in value)) {
180 return null;
181 }
182 
183 return {
184 attempt: Number(value.attempt),
185 nextAt: Number(value.nextAt),
186 };
187};
188 
189export const getSandboxCleanupRetryState = async (
190 context: ProjectDoContext,
191 runId: RunId,
192): Promise<SandboxCleanupRetryState | null> => getSandboxCleanupRetryStateFromStorage(context.ctx.storage, runId);
193 
194export const setSandboxCleanupRetryState = async (
195 context: ProjectDoContext,
196 runId: RunId,
197 value: SandboxCleanupRetryState | null,
198): Promise<void> => {
199 if (value === null) {
200 await context.ctx.storage.delete(sandboxCleanupRetryKey(runId));
201 return;
202 }
203 
204 await context.ctx.storage.put(sandboxCleanupRetryKey(runId), value);
205};
206 
207const scheduleAlarmAtOnStorage = async (
208 storage: Pick<DurableObjectStorage, "getAlarm" | "setAlarm" | "deleteAlarm">,
209 timestamp: number | null,
210): Promise<void> => {
211 if (timestamp === null) {
212 await storage.deleteAlarm();
213 return;
214 }
215 
216 const existing = await storage.getAlarm();
217 const currentTime = Date.now();
218 const nextTimestamp = timestamp <= currentTime ? currentTime : timestamp;
219 
220 // Wrangler local dev can preserve a past-due alarm timestamp across restarts without
221 // delivering it. Treat overdue alarms as stale and re-arm them so reconciliation resumes.
222 if (existing === null || existing <= currentTime || nextTimestamp < existing) {
223 await storage.setAlarm(nextTimestamp);
224 }
225};
226 
227export const scheduleAlarmAt = async (context: ProjectDoContext, timestamp: number | null): Promise<void> => {
228 await scheduleAlarmAtOnStorage(context.ctx.storage, timestamp);
229};
230 
231export const scheduleImmediateReconciliation = async (context: ProjectDoContext): Promise<void> => {
232 await scheduleAlarmAt(context, Date.now());
233};
234 
235export const armReconciliation = async (context: ProjectDoContext, projectId: ProjectId): Promise<void> => {
236 try {
237 await rescheduleAlarm(context, projectId);
238 } catch (error) {
239 context.logger.error("project_reconciliation_arm_failed", {
240 projectId,
241 error: error instanceof Error ? error.message : String(error),
242 });
243 
244 try {
245 await context.ctx.storage.setAlarm(Date.now());
246 } catch (fallbackError) {
247 context.logger.error("project_reconciliation_fallback_arm_failed", {
248 projectId,
249 error: fallbackError instanceof Error ? fallbackError.message : String(fallbackError),
250 });
251 }
252 }
253};
254 
255const computeNextAlarmAt = async (
256 context: ProjectDoContext,
257 storage: SidecarStorage,
258 projectId: ProjectId,
259): Promise<number | null> => {
260 const stateRow = await getProjectState(context, projectId);
261 const rows = await listProjectRuns(context, projectId);
262 const currentTime = Date.now();
263 let nextAt: number | null = null;
264 const hasNonTerminalRows = rows.some((row) => !isTerminalStatus(row.status));
265 
266 const consider = (candidate: number): void => {
267 nextAt = nextAt === null ? candidate : Math.min(nextAt, candidate);
268 };
269 
270 if (stateRow?.activeRunId) {
271 const activeRunId = expectTrusted(RunId, stateRow.activeRunId, "RunId");
272 const heartbeatAt = await getHeartbeatAtFromStorage(storage, activeRunId);
273 // If claimRunWork() committed but the sidecar heartbeat write has not happened yet,
274 // fall back to the active-claim timestamp instead of treating the run as stale immediately.
275 consider((heartbeatAt ?? stateRow.updatedAt) + HEARTBEAT_STALE_AFTER_MS);
276 }
277 
278 if (rows.some((row) => row.status === "cancel_requested")) {
279 consider(currentTime + CANCEL_RECONCILE_RETRY_MS);
280 }
281 
282 if (stateRow?.projectIndexSyncStatus === "needs_update") {
283 consider((await getProjectIndexRetryStateFromStorage(storage))?.nextAt ?? currentTime);
284 }
285 
286 const hasExecutable = rows.some((row) => row.status === "executable");
287 const hasPending = rows.some((row) => row.status === "pending");
288 if (!stateRow?.activeRunId && !hasExecutable && hasPending) {
289 consider(currentTime);
290 }
291 
292 for (const row of rows) {
293 const runId = expectTrusted(RunId, row.runId, "RunId");
294 const sandboxCleanupRetryState = await getSandboxCleanupRetryStateFromStorage(storage, runId);
295 
296 if (sandboxCleanupRetryState) {
297 consider(sandboxCleanupRetryState.nextAt);
298 }
299 
300 if (row.status === "executable" && row.dispatchStatus === "pending") {
301 consider((await getDispatchRetryAtFromStorage(storage, runId)) ?? currentTime);
302 }
303 
304 if (
305 row.d1SyncStatus === "needs_create" ||
306 row.d1SyncStatus === "needs_update" ||
307 row.d1SyncStatus === "needs_terminal_update"
308 ) {
309 const phase: D1RetryPhase =
310 row.d1SyncStatus === "needs_create" ? "create" : row.d1SyncStatus === "needs_update" ? "metadata" : "terminal";
311 consider(getD1RetryAtForPhase(await getD1RetryStateFromStorage(storage, runId), phase) ?? currentTime);
312 }
313 }
314 
315 if (nextAt === null && hasNonTerminalRows) {
316 consider(currentTime + PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS);
317 }
318 
319 return nextAt;
320};
321 
322export const rescheduleAlarmInTransaction = async (
323 context: ProjectDoContext,
324 txn: DurableObjectTransaction,
325 projectId: ProjectId,
326): Promise<void> => {
327 const nextAt = await computeNextAlarmAt(context, txn, projectId);
328 await scheduleAlarmAtOnStorage(txn, nextAt);
329};
330 
331export const rescheduleAlarm = async (context: ProjectDoContext, projectId: ProjectId): Promise<void> => {
332 const nextAt = await computeNextAlarmAt(context, context.ctx.storage, projectId);
333 await scheduleAlarmAtOnStorage(context.ctx.storage, nextAt);
334};