File
Blob: src/worker/durable/project-do/sidecar-state.ts
| 1 | import { type ProjectId, RunId } from "@/contracts"; |
| 2 | import { expectTrusted } from "@/worker/contracts"; |
| 3 | |
| 4 | import { |
| 5 | CANCEL_RECONCILE_RETRY_MS, |
| 6 | HEARTBEAT_STALE_AFTER_MS, |
| 7 | PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS, |
| 8 | } from "./constants"; |
| 9 | import { getProjectState, listProjectRuns } from "./repo"; |
| 10 | import { isTerminalStatus } from "./transitions"; |
| 11 | import type { |
| 12 | D1RetryPhase, |
| 13 | D1RetryState, |
| 14 | ProjectDoContext, |
| 15 | ProjectIndexRetryState, |
| 16 | SandboxCleanupRetryState, |
| 17 | } from "./types"; |
| 18 | |
| 19 | type SidecarStorage = Pick<DurableObjectStorage, "get" | "getAlarm" | "setAlarm" | "deleteAlarm">; |
| 20 | |
| 21 | export const dispatchRetryKey = (runId: RunId): string => `run:${runId}:dispatch-retry-at`; |
| 22 | |
| 23 | export const d1RetryKey = (runId: RunId): string => `run:${runId}:d1-retry`; |
| 24 | |
| 25 | export const heartbeatKey = (runId: RunId): string => `run:${runId}:heartbeat-at`; |
| 26 | export const sandboxCleanupRetryKey = (runId: RunId): string => `run:${runId}:sandbox-cleanup-retry`; |
| 27 | |
| 28 | export const projectIndexRetryKey = (): string => "project:index-sync-retry"; |
| 29 | |
| 30 | const isD1RetryPhase = (value: unknown): value is D1RetryPhase => |
| 31 | value === "create" || value === "metadata" || value === "terminal"; |
| 32 | |
| 33 | const isProjectIndexRetryPhase = (value: unknown): value is "project_index" => value === "project_index"; |
| 34 | |
| 35 | export 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 | |
| 43 | export const shouldWaitForD1Retry = (retryState: D1RetryState | null, phase: D1RetryPhase): boolean => { |
| 44 | const retryAt = getD1RetryAtForPhase(retryState, phase); |
| 45 | return retryAt !== null && retryAt > Date.now(); |
| 46 | }; |
| 47 | |
| 48 | const 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 | |
| 56 | export const getDispatchRetryAt = async (context: ProjectDoContext, runId: RunId): Promise<number | null> => |
| 57 | getDispatchRetryAtFromStorage(context.ctx.storage, runId); |
| 58 | |
| 59 | export 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 | |
| 72 | const 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 | |
| 96 | export const getD1RetryState = async (context: ProjectDoContext, runId: RunId): Promise<D1RetryState | null> => |
| 97 | getD1RetryStateFromStorage(context.ctx.storage, runId); |
| 98 | |
| 99 | export 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 | |
| 112 | const 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 | |
| 135 | export const getProjectIndexRetryState = async (context: ProjectDoContext): Promise<ProjectIndexRetryState | null> => |
| 136 | getProjectIndexRetryStateFromStorage(context.ctx.storage); |
| 137 | |
| 138 | export 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 | |
| 150 | const 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 | |
| 158 | export const getHeartbeatAt = async (context: ProjectDoContext, runId: RunId): Promise<number | null> => |
| 159 | getHeartbeatAtFromStorage(context.ctx.storage, runId); |
| 160 | |
| 161 | export 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 | |
| 170 | const 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 | |
| 189 | export const getSandboxCleanupRetryState = async ( |
| 190 | context: ProjectDoContext, |
| 191 | runId: RunId, |
| 192 | ): Promise<SandboxCleanupRetryState | null> => getSandboxCleanupRetryStateFromStorage(context.ctx.storage, runId); |
| 193 | |
| 194 | export 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 | |
| 207 | const 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 | |
| 227 | export const scheduleAlarmAt = async (context: ProjectDoContext, timestamp: number | null): Promise<void> => { |
| 228 | await scheduleAlarmAtOnStorage(context.ctx.storage, timestamp); |
| 229 | }; |
| 230 | |
| 231 | export const scheduleImmediateReconciliation = async (context: ProjectDoContext): Promise<void> => { |
| 232 | await scheduleAlarmAt(context, Date.now()); |
| 233 | }; |
| 234 | |
| 235 | export 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 | |
| 255 | const 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 | |
| 322 | export 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 | |
| 331 | export 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 | }; |