File
Blob: src/worker/durable/project-do.ts
| 1 | import { DurableObject } from "cloudflare:workers"; |
| 2 | import { drizzle, type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite"; |
| 3 | import { migrate } from "drizzle-orm/durable-sqlite/migrator"; |
| 4 | import { ProjectId, type WebhookProvider } from "@/contracts"; |
| 5 | import { |
| 6 | type AcceptManualRunResult, |
| 7 | type AcceptManualRunInput, |
| 8 | type ClaimRunWorkInput, |
| 9 | type ClaimRunWorkResult, |
| 10 | type FinalizeRunExecutionInput, |
| 11 | type FinalizeRunExecutionResult, |
| 12 | type RecoverWorkflowDispatchFailureInput, |
| 13 | type RecoverWorkflowDispatchFailureResult, |
| 14 | type RecordVerifiedWebhookDeliveryInput, |
| 15 | type RecordVerifiedWebhookDeliveryResult, |
| 16 | type ProjectDetailState, |
| 17 | type RecordRunResolvedCommitInput, |
| 18 | type RecordRunResolvedCommitResult, |
| 19 | type RequestRunCancelInput, |
| 20 | type RequestRunCancelResult, |
| 21 | type RunHeartbeatInput, |
| 22 | type RunHeartbeatResult, |
| 23 | expectTrusted, |
| 24 | } from "@/worker/contracts"; |
| 25 | import { |
| 26 | ALARM_FAILURE_RETRY_MS, |
| 27 | PROJECT_ALARM_MAX_ITERATIONS, |
| 28 | acceptManualRun as acceptManualRunCommand, |
| 29 | claimRunWork as claimRunWorkCommand, |
| 30 | finalizeRunExecution as finalizeRunExecutionCommand, |
| 31 | getProjectConfig as getProjectConfigCommand, |
| 32 | getProjectDetailState as getProjectDetailStateCommand, |
| 33 | getProjectExecutionMaterial as getProjectExecutionMaterialCommand, |
| 34 | getProjectWebhookIngressState as getProjectWebhookIngressStateCommand, |
| 35 | getWebhookVerificationMaterial as getWebhookVerificationMaterialCommand, |
| 36 | initializeProject as initializeProjectCommand, |
| 37 | listProjectWebhooks as listProjectWebhooksCommand, |
| 38 | recoverWorkflowDispatchFailure as recoverWorkflowDispatchFailureCommand, |
| 39 | recordRunResolvedCommit as recordRunResolvedCommitCommand, |
| 40 | recordRunHeartbeat as recordRunHeartbeatCommand, |
| 41 | recordVerifiedWebhookDelivery as recordVerifiedWebhookDeliveryCommand, |
| 42 | requestRunCancel as requestRunCancelCommand, |
| 43 | rotateProjectWebhookSecret as rotateProjectWebhookSecretCommand, |
| 44 | runAlarmCycle, |
| 45 | rescheduleAlarm, |
| 46 | scheduleAlarmAt, |
| 47 | scheduleImmediateReconciliation, |
| 48 | type DeleteProjectWebhookInput, |
| 49 | type InitializeProjectInput, |
| 50 | type ProjectConfigState, |
| 51 | type ProjectDoContext, |
| 52 | type ProjectExecutionMaterial, |
| 53 | type ProjectWebhookIngressState, |
| 54 | type RotateProjectWebhookSecretInput, |
| 55 | type StoredProjectWebhook, |
| 56 | type TouchProjectWebhookVersionsInput, |
| 57 | type UpdateProjectConfigInput, |
| 58 | type UpdateProjectConfigResult, |
| 59 | type UpsertProjectWebhookInput, |
| 60 | type UpsertProjectWebhookResult, |
| 61 | updateProjectConfig as updateProjectConfigCommand, |
| 62 | upsertProjectWebhook as upsertProjectWebhookCommand, |
| 63 | deleteProjectWebhook as deleteProjectWebhookCommand, |
| 64 | touchProjectWebhookVersions as touchProjectWebhookVersionsCommand, |
| 65 | type WebhookVerificationMaterial, |
| 66 | } from "@/worker/durable/project-do/index"; |
| 67 | import { createLogger } from "@/worker/services/logger"; |
| 68 | import projectMigrations from "../../../drizzle/project-do/migrations.js"; |
| 69 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 70 | |
| 71 | const logger = createLogger("durable.project"); |
| 72 | |
| 73 | // Stable entrypoint for the ProjectDO RPC surface. The implementation is split across |
| 74 | // `src/worker/durable/project-do/`, but this file keeps the existing export path and class shape. |
| 75 | export class ProjectDO extends DurableObject { |
| 76 | private readonly db: DrizzleSqliteDODatabase<typeof projectSchema>; |
| 77 | private cachedProjectId: ProjectId | null | undefined; |
| 78 | |
| 79 | constructor(ctx: DurableObjectState, env: Env) { |
| 80 | super(ctx, env); |
| 81 | this.db = drizzle(ctx.storage, { schema: projectSchema }); |
| 82 | |
| 83 | ctx.blockConcurrencyWhile(async () => { |
| 84 | await migrate(this.db, projectMigrations); |
| 85 | const projectId = await this.loadStoredProjectId(); |
| 86 | if (!projectId) { |
| 87 | logger.warn("project_state_missing_project_id", { |
| 88 | phase: "constructor", |
| 89 | objectId: this.ctx.id.toString(), |
| 90 | }); |
| 91 | return; |
| 92 | } |
| 93 | const currentAlarm = await this.ctx.storage.getAlarm(); |
| 94 | if (currentAlarm !== null && currentAlarm > Date.now()) { |
| 95 | return; |
| 96 | } |
| 97 | try { |
| 98 | await rescheduleAlarm(this.getProjectContext(), projectId); |
| 99 | } catch (error) { |
| 100 | logger.error("project_constructor_alarm_rearm_failed", { |
| 101 | projectId, |
| 102 | error: error instanceof Error ? error.message : String(error), |
| 103 | }); |
| 104 | } |
| 105 | }); |
| 106 | } |
| 107 | private getProjectContext(): ProjectDoContext { |
| 108 | return { |
| 109 | ctx: this.ctx, |
| 110 | env: this.env, |
| 111 | db: this.db, |
| 112 | logger, |
| 113 | cacheProjectId: (projectId) => { |
| 114 | this.cachedProjectId = projectId; |
| 115 | }, |
| 116 | }; |
| 117 | } |
| 118 | private async loadStoredProjectId(): Promise<ProjectId | null> { |
| 119 | if (this.cachedProjectId !== undefined) { |
| 120 | return this.cachedProjectId; |
| 121 | } |
| 122 | // Do not use ctx.id.name here. Cloudflare does not populate DurableObjectId.name inside |
| 123 | // the Durable Object runtime, even when the Worker created the stub with idFromName(). |
| 124 | const rows = await this.db |
| 125 | .select({ projectId: projectSchema.projectState.projectId }) |
| 126 | .from(projectSchema.projectState) |
| 127 | .limit(1); |
| 128 | const projectId = rows[0]?.projectId; |
| 129 | if (!projectId) { |
| 130 | this.cachedProjectId = null; |
| 131 | return null; |
| 132 | } |
| 133 | try { |
| 134 | const decodedProjectId = expectTrusted(ProjectId, projectId, "ProjectId"); |
| 135 | this.cachedProjectId = decodedProjectId; |
| 136 | return decodedProjectId; |
| 137 | } catch { |
| 138 | this.cachedProjectId = null; |
| 139 | return null; |
| 140 | } |
| 141 | } |
| 142 | public fetch(): Response { |
| 143 | return new Response("ProjectDO not implemented yet.", { status: 501 }); |
| 144 | } |
| 145 | async alarm(alarmInfo?: AlarmInvocationInfo): Promise<void> { |
| 146 | const context = this.getProjectContext(); |
| 147 | const alarmContext = { |
| 148 | objectId: this.ctx.id.toString(), |
| 149 | retryCount: alarmInfo?.retryCount ?? 0, |
| 150 | isRetry: alarmInfo?.isRetry ?? false, |
| 151 | }; |
| 152 | const projectId = await this.loadStoredProjectId(); |
| 153 | if (!projectId) { |
| 154 | logger.warn("project_state_missing_project_id", { |
| 155 | phase: "alarm", |
| 156 | ...alarmContext, |
| 157 | }); |
| 158 | return; |
| 159 | } |
| 160 | logger.info("project_alarm_started", { |
| 161 | projectId, |
| 162 | action: "start", |
| 163 | ...alarmContext, |
| 164 | }); |
| 165 | try { |
| 166 | let progressCount = 0; |
| 167 | while (true) { |
| 168 | if (progressCount >= PROJECT_ALARM_MAX_ITERATIONS) { |
| 169 | logger.warn("project_alarm_iteration_cap_exceeded", { |
| 170 | projectId, |
| 171 | iterationCap: PROJECT_ALARM_MAX_ITERATIONS, |
| 172 | ...alarmContext, |
| 173 | }); |
| 174 | break; |
| 175 | } |
| 176 | const progress = await runAlarmCycle(context, projectId); |
| 177 | if (!progress) { |
| 178 | break; |
| 179 | } |
| 180 | progressCount += 1; |
| 181 | logger.info("project_alarm_progress", { |
| 182 | projectId, |
| 183 | runId: progress.runId, |
| 184 | action: progress.action, |
| 185 | ...alarmContext, |
| 186 | }); |
| 187 | } |
| 188 | if (progressCount === 0) { |
| 189 | logger.info("project_alarm_idle", { |
| 190 | projectId, |
| 191 | action: "idle", |
| 192 | ...alarmContext, |
| 193 | }); |
| 194 | } |
| 195 | await rescheduleAlarm(context, projectId); |
| 196 | const nextAlarmAt = await this.ctx.storage.getAlarm(); |
| 197 | logger.info("project_alarm_rescheduled", { |
| 198 | projectId, |
| 199 | action: "rescheduled", |
| 200 | nextAlarmAt, |
| 201 | progressCount, |
| 202 | ...alarmContext, |
| 203 | }); |
| 204 | } catch (error) { |
| 205 | logger.error("project_alarm_failed", { |
| 206 | projectId, |
| 207 | ...alarmContext, |
| 208 | error: error instanceof Error ? error.message : String(error), |
| 209 | }); |
| 210 | // Keep reconciliation alive beyond the runtime's limited automatic alarm retries. |
| 211 | const retryAt = Date.now() + ALARM_FAILURE_RETRY_MS; |
| 212 | await scheduleAlarmAt(context, retryAt); |
| 213 | logger.info("project_alarm_retry_scheduled", { |
| 214 | projectId, |
| 215 | action: "retry_scheduled", |
| 216 | nextAlarmAt: retryAt, |
| 217 | ...alarmContext, |
| 218 | }); |
| 219 | } |
| 220 | } |
| 221 | // Keep the public RPC method names and signatures stable. Other Worker modules and the queue |
| 222 | // consumer call these methods directly through the ProjectDO stub. |
| 223 | async getProjectDetailState(projectId: ProjectId): Promise<ProjectDetailState> { |
| 224 | return await getProjectDetailStateCommand(this.getProjectContext(), projectId); |
| 225 | } |
| 226 | async initializeProject(input: InitializeProjectInput): Promise<void> { |
| 227 | await initializeProjectCommand(this.getProjectContext(), input); |
| 228 | } |
| 229 | async getProjectConfig(projectId: ProjectId): Promise<ProjectConfigState | null> { |
| 230 | return await getProjectConfigCommand(this.getProjectContext(), projectId); |
| 231 | } |
| 232 | async updateProjectConfig(input: UpdateProjectConfigInput): Promise<UpdateProjectConfigResult> { |
| 233 | return await updateProjectConfigCommand(this.getProjectContext(), input); |
| 234 | } |
| 235 | async getProjectExecutionMaterial(projectId: ProjectId): Promise<ProjectExecutionMaterial | null> { |
| 236 | return await getProjectExecutionMaterialCommand(this.getProjectContext(), projectId); |
| 237 | } |
| 238 | async getProjectWebhookIngressState( |
| 239 | projectId: ProjectId, |
| 240 | provider: WebhookProvider, |
| 241 | ): Promise<ProjectWebhookIngressState | null> { |
| 242 | return await getProjectWebhookIngressStateCommand(this.getProjectContext(), projectId, provider); |
| 243 | } |
| 244 | async acceptManualRun(input: AcceptManualRunInput): Promise<AcceptManualRunResult> { |
| 245 | return await acceptManualRunCommand(this.getProjectContext(), input); |
| 246 | } |
| 247 | async claimRunWork(input: ClaimRunWorkInput): Promise<ClaimRunWorkResult> { |
| 248 | return await claimRunWorkCommand(this.getProjectContext(), input); |
| 249 | } |
| 250 | async finalizeRunExecution(input: FinalizeRunExecutionInput): Promise<FinalizeRunExecutionResult> { |
| 251 | return await finalizeRunExecutionCommand(this.getProjectContext(), input); |
| 252 | } |
| 253 | async recoverWorkflowDispatchFailure( |
| 254 | input: RecoverWorkflowDispatchFailureInput, |
| 255 | ): Promise<RecoverWorkflowDispatchFailureResult> { |
| 256 | return await recoverWorkflowDispatchFailureCommand(this.getProjectContext(), input); |
| 257 | } |
| 258 | async recordRunResolvedCommit(input: RecordRunResolvedCommitInput): Promise<RecordRunResolvedCommitResult> { |
| 259 | return await recordRunResolvedCommitCommand(this.getProjectContext(), input); |
| 260 | } |
| 261 | async requestRunCancel(input: RequestRunCancelInput): Promise<RequestRunCancelResult> { |
| 262 | return await requestRunCancelCommand(this.getProjectContext(), input); |
| 263 | } |
| 264 | async recordRunHeartbeat(input: RunHeartbeatInput): Promise<RunHeartbeatResult> { |
| 265 | return await recordRunHeartbeatCommand(this.getProjectContext(), input); |
| 266 | } |
| 267 | async getWebhookVerificationMaterial( |
| 268 | projectId: ProjectId, |
| 269 | provider: WebhookProvider, |
| 270 | ): Promise<WebhookVerificationMaterial | null> { |
| 271 | return await getWebhookVerificationMaterialCommand(this.getProjectContext(), projectId, provider); |
| 272 | } |
| 273 | async listProjectWebhooks(projectId: ProjectId): Promise<StoredProjectWebhook[]> { |
| 274 | return await listProjectWebhooksCommand(this.getProjectContext(), projectId); |
| 275 | } |
| 276 | async upsertProjectWebhook(input: UpsertProjectWebhookInput): Promise<UpsertProjectWebhookResult> { |
| 277 | return await upsertProjectWebhookCommand(this.getProjectContext(), input); |
| 278 | } |
| 279 | async rotateProjectWebhookSecret(input: RotateProjectWebhookSecretInput): Promise<StoredProjectWebhook | null> { |
| 280 | return await rotateProjectWebhookSecretCommand(this.getProjectContext(), input); |
| 281 | } |
| 282 | async deleteProjectWebhook(input: DeleteProjectWebhookInput): Promise<boolean> { |
| 283 | return await deleteProjectWebhookCommand(this.getProjectContext(), input); |
| 284 | } |
| 285 | async touchProjectWebhookVersions(input: TouchProjectWebhookVersionsInput): Promise<void> { |
| 286 | await touchProjectWebhookVersionsCommand(this.getProjectContext(), input); |
| 287 | } |
| 288 | async recordVerifiedWebhookDelivery( |
| 289 | input: RecordVerifiedWebhookDeliveryInput, |
| 290 | ): Promise<RecordVerifiedWebhookDeliveryResult> { |
| 291 | return await recordVerifiedWebhookDeliveryCommand(this.getProjectContext(), input); |
| 292 | } |
| 293 | async kickReconciliation(): Promise<void> { |
| 294 | // Worker read/write paths call this through executionCtx.waitUntil() to give alarm-driven reconciliation |
| 295 | // a fresh external trigger without making the HTTP response wait on D1 sync or queue delivery. |
| 296 | await scheduleImmediateReconciliation(this.getProjectContext()); |
| 297 | } |
| 298 | } |