Skip to content
File

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

typescript273 lines
1import { and, asc, desc, eq, lt } from "drizzle-orm";
2 
3import {
4 type ProjectId,
5 ProjectId as ProjectIdCodec,
6 UnixTimestampMs,
7 WebhookProviderConfig as WebhookProviderConfigCodec,
8 WebhookProvider as WebhookProviderCodec,
9 type WebhookProviderConfig,
10 type WebhookProvider,
11 WebhookDeliveryOutcome,
12 WebhookEventKind,
13} from "@/contracts";
14import {
15 type RecordVerifiedWebhookDeliveryInput,
16 type RecordVerifiedWebhookDeliveryResult,
17 expectTrusted,
18} from "@/worker/contracts";
19import * as projectSchema from "@/worker/db/durable/schema/project-do";
20import { generateDurableEntityId } from "@/worker/services";
21 
22import type { ProjectStore } from "../types";
23import type {
24 DeleteProjectWebhookInput,
25 ParsedWebhookDeliveryRow,
26 RotateProjectWebhookSecretInput,
27 StoredProjectWebhook,
28 TouchProjectWebhookVersionsInput,
29 WebhookVerificationMaterial,
30} from "./types";
31 
32const parseWebhookConfig = (configJson: string | null): WebhookProviderConfig | null => {
33 if (configJson === null) {
34 return null;
35 }
36 
37 return WebhookProviderConfigCodec.assertDecode(JSON.parse(configJson) as unknown);
38};
39 
40export const serializeWebhookConfig = (config: WebhookProviderConfig | null): string | null =>
41 config === null ? null : JSON.stringify(config);
42 
43const getNextWebhookUpdatedAt = (now: number, previousUpdatedAt?: number): number =>
44 previousUpdatedAt === undefined ? now : Math.max(now, previousUpdatedAt + 1);
45 
46const parseWebhookRowState = (row: typeof projectSchema.projectWebhooks.$inferSelect) => ({
47 id: row.id,
48 projectId: expectTrusted(ProjectIdCodec, row.projectId, "ProjectId"),
49 provider: expectTrusted(WebhookProviderCodec, row.provider, "WebhookProvider"),
50 enabled: row.enabled !== 0,
51 config: parseWebhookConfig(row.configJson),
52 updatedAt: expectTrusted(UnixTimestampMs, row.updatedAt, "UnixTimestampMs"),
53});
54 
55export const getProjectWebhookRow = (
56 tx: ProjectStore,
57 projectId: ProjectId,
58 provider: WebhookProvider,
59): typeof projectSchema.projectWebhooks.$inferSelect | undefined =>
60 tx
61 .select()
62 .from(projectSchema.projectWebhooks)
63 .where(
64 and(eq(projectSchema.projectWebhooks.projectId, projectId), eq(projectSchema.projectWebhooks.provider, provider)),
65 )
66 .limit(1)
67 .get();
68 
69export const listProjectWebhookRows = (
70 tx: ProjectStore,
71 projectId: ProjectId,
72): Array<typeof projectSchema.projectWebhooks.$inferSelect> =>
73 tx
74 .select()
75 .from(projectSchema.projectWebhooks)
76 .where(eq(projectSchema.projectWebhooks.projectId, projectId))
77 .orderBy(asc(projectSchema.projectWebhooks.provider))
78 .all();
79 
80export const listRecentWebhookDeliveryRows = (
81 tx: ProjectStore,
82 projectId: ProjectId,
83 provider: WebhookProvider,
84 limit: number,
85): Array<typeof projectSchema.projectWebhookDeliveries.$inferSelect> =>
86 tx
87 .select()
88 .from(projectSchema.projectWebhookDeliveries)
89 .where(
90 and(
91 eq(projectSchema.projectWebhookDeliveries.projectId, projectId),
92 eq(projectSchema.projectWebhookDeliveries.provider, provider),
93 ),
94 )
95 .orderBy(desc(projectSchema.projectWebhookDeliveries.receivedAt), desc(projectSchema.projectWebhookDeliveries.id))
96 .limit(limit)
97 .all();
98 
99export const getWebhookDeliveryRow = (
100 tx: ProjectStore,
101 projectId: ProjectId,
102 provider: WebhookProvider,
103 deliveryId: string,
104): typeof projectSchema.projectWebhookDeliveries.$inferSelect | undefined =>
105 tx
106 .select()
107 .from(projectSchema.projectWebhookDeliveries)
108 .where(
109 and(
110 eq(projectSchema.projectWebhookDeliveries.projectId, projectId),
111 eq(projectSchema.projectWebhookDeliveries.provider, provider),
112 eq(projectSchema.projectWebhookDeliveries.deliveryId, deliveryId),
113 ),
114 )
115 .limit(1)
116 .get();
117 
118export const updateWebhookDeliveryRow = (
119 tx: ProjectStore,
120 rowId: string,
121 input: RecordVerifiedWebhookDeliveryInput,
122 outcome: RecordVerifiedWebhookDeliveryResult["outcome"],
123 runId: string | null,
124 receivedAt: number,
125): void => {
126 tx.update(projectSchema.projectWebhookDeliveries)
127 .set({
128 eventKind: input.payload.eventKind,
129 eventName: input.payload.eventName,
130 outcome,
131 repoUrl: input.payload.repoUrl,
132 ref: input.payload.ref,
133 branch: input.payload.branch,
134 commitSha: input.payload.commitSha,
135 beforeSha: input.payload.beforeSha,
136 runId,
137 receivedAt,
138 })
139 .where(eq(projectSchema.projectWebhookDeliveries.id, rowId))
140 .run();
141};
142 
143export const pruneWebhookDeliveries = (tx: ProjectStore, projectId: ProjectId, receivedBefore: number): void => {
144 tx.delete(projectSchema.projectWebhookDeliveries)
145 .where(
146 and(
147 eq(projectSchema.projectWebhookDeliveries.projectId, projectId),
148 lt(projectSchema.projectWebhookDeliveries.receivedAt, receivedBefore),
149 ),
150 )
151 .run();
152};
153 
154export const insertProjectWebhookRow = (
155 tx: ProjectStore,
156 input: {
157 projectId: ProjectId;
158 provider: WebhookProvider;
159 enabled: boolean;
160 config: WebhookProviderConfig | null;
161 encryptedSecret: RotateProjectWebhookSecretInput["encryptedSecret"];
162 now: number;
163 },
164): void => {
165 const updatedAt = getNextWebhookUpdatedAt(input.now);
166 tx.insert(projectSchema.projectWebhooks)
167 .values({
168 id: generateDurableEntityId("whk", input.now),
169 projectId: input.projectId,
170 provider: input.provider,
171 configJson: serializeWebhookConfig(input.config),
172 secretCiphertext: input.encryptedSecret.ciphertext,
173 secretKeyVersion: input.encryptedSecret.keyVersion,
174 secretNonce: input.encryptedSecret.nonce,
175 enabled: input.enabled ? 1 : 0,
176 createdAt: input.now,
177 updatedAt,
178 })
179 .run();
180};
181 
182export const updateProjectWebhookSettingsRow = (
183 tx: ProjectStore,
184 input: {
185 rowId: string;
186 previousUpdatedAt: number;
187 enabled: boolean;
188 config: WebhookProviderConfig | null | undefined;
189 now: number;
190 },
191): void => {
192 const updatedAt = getNextWebhookUpdatedAt(input.now, input.previousUpdatedAt);
193 tx.update(projectSchema.projectWebhooks)
194 .set({
195 enabled: input.enabled ? 1 : 0,
196 updatedAt,
197 ...(input.config === undefined ? {} : { configJson: serializeWebhookConfig(input.config) }),
198 })
199 .where(eq(projectSchema.projectWebhooks.id, input.rowId))
200 .run();
201};
202 
203export const rotateProjectWebhookSecretRow = (tx: ProjectStore, input: RotateProjectWebhookSecretInput): boolean => {
204 const existingRow = getProjectWebhookRow(tx, input.projectId, input.provider);
205 if (!existingRow) {
206 return false;
207 }
208 
209 const updatedAt = getNextWebhookUpdatedAt(input.now, existingRow.updatedAt);
210 tx.update(projectSchema.projectWebhooks)
211 .set({
212 secretCiphertext: input.encryptedSecret.ciphertext,
213 secretKeyVersion: input.encryptedSecret.keyVersion,
214 secretNonce: input.encryptedSecret.nonce,
215 updatedAt,
216 })
217 .where(eq(projectSchema.projectWebhooks.id, existingRow.id))
218 .run();
219 
220 return true;
221};
222 
223export const deleteProjectWebhookRow = (tx: ProjectStore, input: DeleteProjectWebhookInput): boolean => {
224 const existingRow = getProjectWebhookRow(tx, input.projectId, input.provider);
225 if (!existingRow) {
226 return false;
227 }
228 
229 tx.delete(projectSchema.projectWebhooks).where(eq(projectSchema.projectWebhooks.id, existingRow.id)).run();
230 return true;
231};
232 
233export const touchProjectWebhookRows = (tx: ProjectStore, input: TouchProjectWebhookVersionsInput): void => {
234 const rows = listProjectWebhookRows(tx, input.projectId);
235 for (const row of rows) {
236 tx.update(projectSchema.projectWebhooks)
237 .set({
238 updatedAt: getNextWebhookUpdatedAt(input.now, row.updatedAt),
239 })
240 .where(eq(projectSchema.projectWebhooks.id, row.id))
241 .run();
242 }
243};
244 
245export const parseStoredWebhook = (
246 row: typeof projectSchema.projectWebhooks.$inferSelect,
247 recentDeliveries: Array<typeof projectSchema.projectWebhookDeliveries.$inferSelect>,
248): StoredProjectWebhook => ({
249 ...parseWebhookRowState(row),
250 createdAt: row.createdAt,
251 recentDeliveries: recentDeliveries.map(parseWebhookDeliveryRow),
252});
253 
254export const parseWebhookVerificationMaterial = (
255 row: typeof projectSchema.projectWebhooks.$inferSelect,
256): WebhookVerificationMaterial => ({
257 ...parseWebhookRowState(row),
258 encryptedSecret: {
259 ciphertext: row.secretCiphertext,
260 keyVersion: row.secretKeyVersion,
261 nonce: row.secretNonce,
262 },
263});
264 
265export const parseWebhookDeliveryRow = (
266 row: typeof projectSchema.projectWebhookDeliveries.$inferSelect,
267): ParsedWebhookDeliveryRow => ({
268 ...row,
269 provider: expectTrusted(WebhookProviderCodec, row.provider, "WebhookProvider"),
270 eventKind: expectTrusted(WebhookEventKind, row.eventKind, "WebhookEventKind"),
271 outcome: expectTrusted(WebhookDeliveryOutcome, row.outcome, "WebhookDeliveryOutcome"),
272});