Skip to content
File

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

typescript299 lines
1import { DurableObject } from "cloudflare:workers";
2import { drizzle, type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite";
3import { migrate } from "drizzle-orm/durable-sqlite/migrator";
4import { ProjectId, type WebhookProvider } from "@/contracts";
5import {
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";
25import {
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";
67import { createLogger } from "@/worker/services/logger";
68import projectMigrations from "../../../drizzle/project-do/migrations.js";
69import * as projectSchema from "@/worker/db/durable/schema/project-do";
70 
71const 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.
75export 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}