File
Blob: src/worker/durable/project-do/project-config.ts
| 1 | import { eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { DispatchMode, ProjectId, type WebhookProvider } from "@/contracts"; |
| 4 | import { expectTrusted } from "@/worker/contracts"; |
| 5 | import { getWebhookProviderCatalogEntry, validateWebhookConfigForUpsert } from "@/lib/webhooks"; |
| 6 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 7 | import type { EncryptedSecret } from "@/worker/security/secrets"; |
| 8 | |
| 9 | import { getProjectConfig, getProjectConfigRow } from "./repo"; |
| 10 | import { rescheduleAlarmInTransaction } from "./sidecar-state"; |
| 11 | import { ensureProjectState } from "./transitions/shared"; |
| 12 | import type { |
| 13 | InitializeProjectInput, |
| 14 | ProjectConfigRow, |
| 15 | ProjectConfigState, |
| 16 | ProjectDoContext, |
| 17 | ProjectExecutionMaterial, |
| 18 | ProjectStore, |
| 19 | ProjectWebhookIngressState, |
| 20 | UpdateProjectConfigInput, |
| 21 | UpdateProjectConfigResult, |
| 22 | } from "./types"; |
| 23 | import { |
| 24 | getProjectWebhookRow, |
| 25 | listProjectWebhookRows, |
| 26 | parseStoredWebhook, |
| 27 | parseWebhookVerificationMaterial, |
| 28 | } from "./webhooks"; |
| 29 | import { touchProjectWebhookRows } from "./webhooks/repo"; |
| 30 | |
| 31 | const getNextProjectUpdatedAt = (now: number, previousUpdatedAt?: number): number => |
| 32 | previousUpdatedAt === undefined ? now : Math.max(now, previousUpdatedAt + 1); |
| 33 | |
| 34 | const toEncryptedRepoToken = ( |
| 35 | row: Pick<ProjectConfigRow, "repoTokenCiphertext" | "repoTokenKeyVersion" | "repoTokenNonce">, |
| 36 | ): EncryptedSecret | null => |
| 37 | row.repoTokenCiphertext !== null && row.repoTokenKeyVersion !== null && row.repoTokenNonce !== null |
| 38 | ? { |
| 39 | ciphertext: row.repoTokenCiphertext, |
| 40 | keyVersion: row.repoTokenKeyVersion, |
| 41 | nonce: row.repoTokenNonce, |
| 42 | } |
| 43 | : null; |
| 44 | |
| 45 | const toProjectConfigState = (row: ProjectConfigRow): ProjectConfigState => ({ |
| 46 | projectId: expectTrusted(ProjectId, row.projectId, "ProjectId"), |
| 47 | name: row.name, |
| 48 | repoUrl: row.repoUrl, |
| 49 | defaultBranch: row.defaultBranch, |
| 50 | configPath: row.configPath, |
| 51 | dispatchMode: expectTrusted(DispatchMode, row.dispatchMode, "DispatchMode"), |
| 52 | createdAt: row.createdAt, |
| 53 | updatedAt: row.updatedAt, |
| 54 | }); |
| 55 | |
| 56 | const validateProjectRepoUrlAgainstConfiguredWebhooks = ( |
| 57 | tx: ProjectStore, |
| 58 | projectId: ProjectId, |
| 59 | nextRepoUrl: string, |
| 60 | ): UpdateProjectConfigResult | null => { |
| 61 | const conflictingProviders: WebhookProvider[] = []; |
| 62 | const webhookRows = listProjectWebhookRows(tx, projectId); |
| 63 | |
| 64 | for (const webhookRow of webhookRows) { |
| 65 | const storedWebhook = parseStoredWebhook(webhookRow, []); |
| 66 | const result = validateWebhookConfigForUpsert({ |
| 67 | provider: storedWebhook.provider, |
| 68 | projectRepoUrl: nextRepoUrl, |
| 69 | incomingConfig: undefined, |
| 70 | existingConfig: storedWebhook.config, |
| 71 | creating: false, |
| 72 | }); |
| 73 | |
| 74 | if (!result.ok) { |
| 75 | conflictingProviders.push(storedWebhook.provider); |
| 76 | } |
| 77 | } |
| 78 | |
| 79 | if (conflictingProviders.length === 0) { |
| 80 | return null; |
| 81 | } |
| 82 | |
| 83 | const providerNames = conflictingProviders.map((provider) => getWebhookProviderCatalogEntry(provider).displayName); |
| 84 | return { |
| 85 | kind: "invalid", |
| 86 | status: 400, |
| 87 | code: "project_repo_url_conflicts_with_webhook", |
| 88 | message: `Project repoUrl conflicts with configured webhook providers: ${providerNames.join(", ")}. Update or delete those webhooks first.`, |
| 89 | details: { |
| 90 | providers: conflictingProviders, |
| 91 | }, |
| 92 | }; |
| 93 | }; |
| 94 | |
| 95 | const upsertProjectConfigRow = (tx: ProjectStore, input: InitializeProjectInput): void => { |
| 96 | const updatedAt = getNextProjectUpdatedAt(input.updatedAt); |
| 97 | const existingRow = getProjectConfigRow(tx, input.projectId); |
| 98 | |
| 99 | if (!existingRow) { |
| 100 | tx.insert(projectSchema.projectConfig) |
| 101 | .values({ |
| 102 | projectId: input.projectId, |
| 103 | name: input.name, |
| 104 | repoUrl: input.repoUrl, |
| 105 | defaultBranch: input.defaultBranch, |
| 106 | configPath: input.configPath, |
| 107 | repoTokenCiphertext: input.encryptedRepoToken?.ciphertext ?? null, |
| 108 | repoTokenKeyVersion: input.encryptedRepoToken?.keyVersion ?? null, |
| 109 | repoTokenNonce: input.encryptedRepoToken?.nonce ?? null, |
| 110 | dispatchMode: input.dispatchMode, |
| 111 | executionRuntime: input.executionRuntime, |
| 112 | createdAt: input.createdAt, |
| 113 | updatedAt, |
| 114 | }) |
| 115 | .run(); |
| 116 | return; |
| 117 | } |
| 118 | |
| 119 | tx.update(projectSchema.projectConfig) |
| 120 | .set({ |
| 121 | name: input.name, |
| 122 | repoUrl: input.repoUrl, |
| 123 | defaultBranch: input.defaultBranch, |
| 124 | configPath: input.configPath, |
| 125 | repoTokenCiphertext: input.encryptedRepoToken?.ciphertext ?? null, |
| 126 | repoTokenKeyVersion: input.encryptedRepoToken?.keyVersion ?? null, |
| 127 | repoTokenNonce: input.encryptedRepoToken?.nonce ?? null, |
| 128 | dispatchMode: input.dispatchMode, |
| 129 | executionRuntime: input.executionRuntime, |
| 130 | createdAt: input.createdAt, |
| 131 | updatedAt, |
| 132 | }) |
| 133 | .where(eq(projectSchema.projectConfig.projectId, input.projectId)) |
| 134 | .run(); |
| 135 | }; |
| 136 | |
| 137 | export const initializeProject = async (context: ProjectDoContext, input: InitializeProjectInput): Promise<void> => { |
| 138 | context.db.transaction((tx) => { |
| 139 | ensureProjectState(context, tx, input.projectId); |
| 140 | upsertProjectConfigRow(tx, input); |
| 141 | tx.update(projectSchema.projectState) |
| 142 | .set({ |
| 143 | projectIndexSyncStatus: "current", |
| 144 | }) |
| 145 | .where(eq(projectSchema.projectState.projectId, input.projectId)) |
| 146 | .run(); |
| 147 | }); |
| 148 | }; |
| 149 | |
| 150 | export const getProjectConfigState = async ( |
| 151 | context: ProjectDoContext, |
| 152 | projectId: ProjectId, |
| 153 | ): Promise<ProjectConfigState | null> => { |
| 154 | const row = await getProjectConfig(context, projectId); |
| 155 | return row ? toProjectConfigState(row) : null; |
| 156 | }; |
| 157 | |
| 158 | export const getProjectExecutionMaterial = async ( |
| 159 | context: ProjectDoContext, |
| 160 | projectId: ProjectId, |
| 161 | ): Promise<ProjectExecutionMaterial | null> => { |
| 162 | const row = await getProjectConfig(context, projectId); |
| 163 | if (!row) { |
| 164 | return null; |
| 165 | } |
| 166 | |
| 167 | return { |
| 168 | projectId, |
| 169 | encryptedRepoToken: toEncryptedRepoToken(row), |
| 170 | }; |
| 171 | }; |
| 172 | |
| 173 | export const getProjectWebhookIngressState = async ( |
| 174 | context: ProjectDoContext, |
| 175 | projectId: ProjectId, |
| 176 | provider: WebhookProvider, |
| 177 | ): Promise<ProjectWebhookIngressState | null> => { |
| 178 | const configRow = await getProjectConfig(context, projectId); |
| 179 | if (!configRow) { |
| 180 | return null; |
| 181 | } |
| 182 | |
| 183 | const webhookRow = getProjectWebhookRow(context.db, projectId, provider); |
| 184 | if (!webhookRow) { |
| 185 | return null; |
| 186 | } |
| 187 | |
| 188 | return { |
| 189 | projectId, |
| 190 | repoUrl: configRow.repoUrl, |
| 191 | defaultBranch: configRow.defaultBranch, |
| 192 | configPath: configRow.configPath, |
| 193 | webhook: parseWebhookVerificationMaterial(webhookRow), |
| 194 | }; |
| 195 | }; |
| 196 | |
| 197 | export const updateProjectConfig = async ( |
| 198 | context: ProjectDoContext, |
| 199 | input: UpdateProjectConfigInput, |
| 200 | ): Promise<UpdateProjectConfigResult> => |
| 201 | context.ctx.storage.transaction(async (txn) => { |
| 202 | const result = context.db.transaction((tx): UpdateProjectConfigResult => { |
| 203 | ensureProjectState(context, tx, input.projectId); |
| 204 | |
| 205 | const existingRow = getProjectConfigRow(tx, input.projectId); |
| 206 | if (!existingRow) { |
| 207 | return { |
| 208 | kind: "not_found", |
| 209 | }; |
| 210 | } |
| 211 | |
| 212 | const nextRepoUrl = input.repoUrl ?? existingRow.repoUrl; |
| 213 | const repoConflict = validateProjectRepoUrlAgainstConfiguredWebhooks(tx, input.projectId, nextRepoUrl); |
| 214 | if (repoConflict) { |
| 215 | return repoConflict; |
| 216 | } |
| 217 | |
| 218 | const nextDefaultBranch = input.defaultBranch ?? existingRow.defaultBranch; |
| 219 | const nextConfigPath = input.configPath ?? existingRow.configPath; |
| 220 | const nextDispatchMode = |
| 221 | input.dispatchMode ?? expectTrusted(DispatchMode, existingRow.dispatchMode, "DispatchMode"); |
| 222 | const updatedAt = getNextProjectUpdatedAt(input.now, existingRow.updatedAt); |
| 223 | const webhookHandlingChanged = |
| 224 | nextRepoUrl !== existingRow.repoUrl || |
| 225 | nextDefaultBranch !== existingRow.defaultBranch || |
| 226 | nextConfigPath !== existingRow.configPath; |
| 227 | |
| 228 | tx.update(projectSchema.projectConfig) |
| 229 | .set({ |
| 230 | name: input.name ?? existingRow.name, |
| 231 | repoUrl: nextRepoUrl, |
| 232 | defaultBranch: nextDefaultBranch, |
| 233 | configPath: nextConfigPath, |
| 234 | dispatchMode: nextDispatchMode, |
| 235 | repoTokenCiphertext: |
| 236 | input.encryptedRepoToken === undefined |
| 237 | ? existingRow.repoTokenCiphertext |
| 238 | : (input.encryptedRepoToken?.ciphertext ?? null), |
| 239 | repoTokenKeyVersion: |
| 240 | input.encryptedRepoToken === undefined |
| 241 | ? existingRow.repoTokenKeyVersion |
| 242 | : (input.encryptedRepoToken?.keyVersion ?? null), |
| 243 | repoTokenNonce: |
| 244 | input.encryptedRepoToken === undefined |
| 245 | ? existingRow.repoTokenNonce |
| 246 | : (input.encryptedRepoToken?.nonce ?? null), |
| 247 | updatedAt, |
| 248 | }) |
| 249 | .where(eq(projectSchema.projectConfig.projectId, input.projectId)) |
| 250 | .run(); |
| 251 | |
| 252 | if (webhookHandlingChanged) { |
| 253 | touchProjectWebhookRows(tx, { |
| 254 | projectId: input.projectId, |
| 255 | now: updatedAt, |
| 256 | }); |
| 257 | } |
| 258 | |
| 259 | tx.update(projectSchema.projectState) |
| 260 | .set({ |
| 261 | projectIndexSyncStatus: "needs_update", |
| 262 | }) |
| 263 | .where(eq(projectSchema.projectState.projectId, input.projectId)) |
| 264 | .run(); |
| 265 | |
| 266 | return { |
| 267 | kind: "applied", |
| 268 | config: { |
| 269 | projectId: input.projectId, |
| 270 | name: input.name ?? existingRow.name, |
| 271 | repoUrl: nextRepoUrl, |
| 272 | defaultBranch: nextDefaultBranch, |
| 273 | configPath: nextConfigPath, |
| 274 | dispatchMode: nextDispatchMode, |
| 275 | createdAt: existingRow.createdAt, |
| 276 | updatedAt, |
| 277 | }, |
| 278 | }; |
| 279 | }); |
| 280 | |
| 281 | if (result.kind === "applied") { |
| 282 | await rescheduleAlarmInTransaction(context, txn, input.projectId); |
| 283 | } |
| 284 | |
| 285 | return result; |
| 286 | }); |