返回 DeepSeek-Reasonix
crash_delivery.ts
根目录 / workers / crash-report / src / crash_delivery.ts
1 import type { Env } from "./env";
2 import { firebaseConfigured } from "./firebase_rtdb";
3
4 export const FIREBASE_OUTBOX_LIMIT = 5_000;
5 export const FIREBASE_OUTBOX_BATCH = 200;
6 const FIREBASE_GROUP_LEASE_MS = 60_000;
7 export const FIREBASE_STORAGE_BUDGET_BYTES = 700 * 1024 * 1024;
8 export const FIREBASE_STORAGE_WARNING_BYTES = Math.floor(FIREBASE_STORAGE_BUDGET_BYTES * 0.8);
9 export const FIREBASE_ACTIVE_RESERVATION_BYTES = 640 * 1024;
10 export const FIREBASE_COMPACTED_RESERVATION_BYTES = 128 * 1024;
11 export const FIREBASE_ARCHIVING_RESERVATION_BYTES = 32 * 1024;
12
13 export type FirebaseSampleState = "active" | "compacted" | "archiving" | "archived";
14
15 export type FirebaseGroupState = {
16 fingerprint: string;
17 sample_state: FirebaseSampleState;
18 sample_epoch: number;
19 epoch_first_event_id: string;
20 reserved_bytes: number;
21 last_seen: string;
22 compacted_at: string;
23 archived_at: string;
24 archive_reason: "" | "retention" | "admin";
25 lease_owner: string;
26 lease_generation: number;
27 lease_expires_at: string;
28 };
29
30 export type FirebaseGroupLease = { owner: string; generation: number };
31
32 export type CrashStorageMode = "d1" | "dual" | "firebase";
33
34 export type FirebaseOutboxRow = {
35 event_id: string;
36 fingerprint: string;
37 payload: string;
38 state: "queued" | "processing" | "projected";
39 attempts: number;
40 next_attempt_at: string;
41 created_at: string;
42 updated_at: string;
43 };
44
45 export type FirebaseProjectionReceipt = {
46 group_count: number;
47 latest_slot: number;
48 first_sample: number;
49 };
50
51 export function crashStorageMode(env: Env): CrashStorageMode {
52 const raw = env.CRASH_STORAGE_MODE?.trim().toLowerCase() || "d1";
53 if (raw === "d1" || raw === "dual" || raw === "firebase") return raw;
54 throw new Error(`invalid crash storage mode ${raw}`);
55 }
56
57 export function firebaseStorageReady(env: Env): boolean {
58 return crashStorageMode(env) === "d1" || firebaseConfigured(env);
59 }
60
61 export async function enqueueFirebaseCrash(
62 env: Env,
63 eventId: string,
64 fingerprint: string,
65 payload: string,
66 now: string,
67 ): Promise<"inserted" | "duplicate" | "full"> {
68 const result = await env.DB.prepare(
69 `INSERT OR IGNORE INTO firebase_crash_outbox (
70 event_id, fingerprint, payload, state, attempts, next_attempt_at, created_at, updated_at
71 )
72 SELECT ?1, ?2, ?3, 'queued', 0, ?4, ?4, ?4
73 WHERE (SELECT COUNT(*) FROM firebase_crash_outbox) < ?5
74 AND NOT EXISTS (SELECT 1 FROM firebase_crash_receipts WHERE event_id = ?1)`,
75 ).bind(eventId, fingerprint, payload, now, FIREBASE_OUTBOX_LIMIT).run();
76 if (Number(result.meta?.changes ?? 0) > 0) return "inserted";
77 const existing = await env.DB.prepare(
78 `SELECT event_id FROM firebase_crash_outbox WHERE event_id = ?1
79 UNION ALL SELECT event_id FROM firebase_crash_receipts WHERE event_id = ?1 LIMIT 1`,
80 ).bind(eventId).first<{ event_id: string }>();
81 return existing ? "duplicate" : "full";
82 }
83
84 export async function claimFirebaseCrash(env: Env, eventId: string, now: string): Promise<boolean> {
85 const result = await env.DB.prepare(
86 `UPDATE firebase_crash_outbox SET state = 'processing', updated_at = ?2
87 WHERE event_id = ?1 AND (
88 state = 'queued' OR (state = 'processing' AND datetime(updated_at) < datetime('now', '-10 minutes'))
89 )`,
90 ).bind(eventId, now).run();
91 return Number(result.meta?.changes ?? 0) > 0;
92 }
93
94 export function projectionCompletionStatements(
95 db: D1Database,
96 eventId: string,
97 fingerprint: string,
98 now: string,
99 ): D1PreparedStatement[] {
100 const stateWrite = db.prepare(
101 `UPDATE firebase_crash_group_state
102 SET epoch_first_event_id = CASE WHEN epoch_first_event_id = '' THEN ?1 ELSE epoch_first_event_id END,
103 last_seen = CASE WHEN last_seen < ?2 THEN ?2 ELSE last_seen END
104 WHERE fingerprint = ?3 AND sample_state = 'active'`,
105 ).bind(eventId, now, fingerprint);
106 const receipt = db.prepare(
107 `INSERT OR IGNORE INTO firebase_crash_receipts (
108 event_id, projected_at, group_count, latest_slot, first_sample
109 )
110 SELECT ?1, ?2, groups.count, (groups.count - 1) % 5,
111 CASE WHEN state.epoch_first_event_id = ?1 THEN 1 ELSE 0 END
112 FROM groups
113 JOIN firebase_crash_group_state AS state USING (fingerprint)
114 WHERE groups.fingerprint = ?3`,
115 ).bind(eventId, now, fingerprint);
116 return [
117 stateWrite,
118 receipt,
119 db.prepare(
120 "UPDATE firebase_crash_outbox SET state = 'projected', updated_at = ?2 WHERE event_id = ?1",
121 ).bind(eventId, now),
122 ];
123 }
124
125 export async function firebaseProjectionReceipt(
126 env: Env,
127 eventId: string,
128 ): Promise<FirebaseProjectionReceipt | null> {
129 return env.DB.prepare(
130 "SELECT group_count, latest_slot, first_sample FROM firebase_crash_receipts WHERE event_id = ?1",
131 ).bind(eventId).first<FirebaseProjectionReceipt>();
132 }
133
134 export async function firebaseProjectionExists(env: Env, eventId: string): Promise<boolean> {
135 return Boolean(await env.DB.prepare(
136 "SELECT event_id FROM firebase_crash_receipts WHERE event_id = ?1",
137 ).bind(eventId).first());
138 }
139
140 export async function firebaseOutboxExists(env: Env, eventId: string): Promise<boolean> {
141 return Boolean(await env.DB.prepare(
142 "SELECT event_id FROM firebase_crash_outbox WHERE event_id = ?1",
143 ).bind(eventId).first());
144 }
145
146 export async function firebaseEventExists(env: Env, eventId: string): Promise<boolean> {
147 return Boolean(await env.DB.prepare(
148 `SELECT event_id FROM firebase_crash_outbox WHERE event_id = ?1
149 UNION ALL SELECT event_id FROM firebase_crash_receipts WHERE event_id = ?1 LIMIT 1`,
150 ).bind(eventId).first());
151 }
152
153 export async function firebaseGroupState(
154 env: Env,
155 fingerprint: string,
156 ): Promise<FirebaseGroupState | null> {
157 return env.DB.prepare(
158 `SELECT fingerprint, sample_state, sample_epoch, epoch_first_event_id, reserved_bytes,
159 last_seen, compacted_at, archived_at, archive_reason,
160 lease_owner, lease_generation, lease_expires_at
161 FROM firebase_crash_group_state WHERE fingerprint = ?1`,
162 ).bind(fingerprint).first<FirebaseGroupState>();
163 }
164
165 export async function reserveFirebaseGroup(
166 env: Env,
167 fingerprint: string,
168 lastSeen: string,
169 ): Promise<"reserved" | "full"> {
170 const result = await env.DB.prepare(
171 `INSERT INTO firebase_crash_group_state (
172 fingerprint, sample_state, sample_epoch, epoch_first_event_id,
173 reserved_bytes, last_seen, compacted_at, archived_at, archive_reason
174 )
175 SELECT ?1, 'active', 1, '', ?2, ?3, '', '', ''
176 WHERE (SELECT COALESCE(SUM(reserved_bytes), 0) FROM firebase_crash_group_state) + ?2 <= ?4
177 AND (
178 NOT EXISTS (SELECT 1 FROM groups WHERE fingerprint = ?1) OR
179 EXISTS (SELECT 1 FROM firebase_crash_group_state WHERE fingerprint = ?1)
180 )
181 ON CONFLICT (fingerprint) DO UPDATE SET
182 sample_state = 'active',
183 sample_epoch = CASE
184 WHEN firebase_crash_group_state.sample_state IN ('archiving', 'archived')
185 THEN firebase_crash_group_state.sample_epoch + 1
186 ELSE firebase_crash_group_state.sample_epoch
187 END,
188 epoch_first_event_id = CASE
189 WHEN firebase_crash_group_state.sample_state IN ('archiving', 'archived') THEN ''
190 ELSE firebase_crash_group_state.epoch_first_event_id
191 END,
192 reserved_bytes = ?2,
193 last_seen = CASE WHEN firebase_crash_group_state.last_seen < ?3
194 THEN ?3 ELSE firebase_crash_group_state.last_seen END,
195 compacted_at = '', archived_at = '', archive_reason = ''
196 WHERE firebase_crash_group_state.reserved_bytes >= ?2 OR
197 (SELECT COALESCE(SUM(reserved_bytes), 0) FROM firebase_crash_group_state)
198 - firebase_crash_group_state.reserved_bytes + ?2 <= ?4`,
199 ).bind(fingerprint, FIREBASE_ACTIVE_RESERVATION_BYTES, lastSeen, FIREBASE_STORAGE_BUDGET_BYTES).run();
200 return Number(result.meta?.changes ?? 0) > 0 ? "reserved" : "full";
201 }
202
203 export async function reclaimUnusedFirebaseReservation(env: Env, fingerprint: string): Promise<void> {
204 await env.DB.prepare(
205 `DELETE FROM firebase_crash_group_state
206 WHERE fingerprint = ?1 AND lease_generation = 0
207 AND NOT EXISTS (SELECT 1 FROM groups WHERE fingerprint = ?1)
208 AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1)`,
209 ).bind(fingerprint).run();
210 }
211
212 export async function acquireFirebaseGroupLease(
213 env: Env,
214 fingerprint: string,
215 now = new Date(),
216 ): Promise<FirebaseGroupLease | null> {
217 const owner = crypto.randomUUID();
218 const acquiredAt = now.toISOString();
219 const expiresAt = new Date(now.getTime() + FIREBASE_GROUP_LEASE_MS).toISOString();
220 const acquired = await env.DB.prepare(
221 `UPDATE firebase_crash_group_state
222 SET lease_owner = ?2, lease_generation = lease_generation + 1, lease_expires_at = ?3
223 WHERE fingerprint = ?1 AND (
224 lease_owner = '' OR datetime(lease_expires_at) <= datetime(?4)
225 )
226 RETURNING lease_generation`,
227 ).bind(fingerprint, owner, expiresAt, acquiredAt).first<{ lease_generation: number }>();
228 return acquired ? { owner, generation: Number(acquired.lease_generation) } : null;
229 }
230
231 export async function renewFirebaseGroupLease(
232 env: Env,
233 fingerprint: string,
234 lease: FirebaseGroupLease,
235 now = new Date(),
236 ): Promise<boolean> {
237 const current = now.toISOString();
238 const expiresAt = new Date(now.getTime() + FIREBASE_GROUP_LEASE_MS).toISOString();
239 const result = await env.DB.prepare(
240 `UPDATE firebase_crash_group_state SET lease_expires_at = ?4
241 WHERE fingerprint = ?1 AND lease_owner = ?2 AND lease_generation = ?3
242 AND datetime(lease_expires_at) > datetime(?5)`,
243 ).bind(fingerprint, lease.owner, lease.generation, expiresAt, current).run();
244 return Number(result.meta?.changes ?? 0) > 0;
245 }
246
247 export async function releaseFirebaseGroupLease(
248 env: Env,
249 fingerprint: string,
250 lease: FirebaseGroupLease,
251 ): Promise<void> {
252 await env.DB.prepare(
253 `UPDATE firebase_crash_group_state SET lease_owner = '', lease_expires_at = ''
254 WHERE fingerprint = ?1 AND lease_owner = ?2 AND lease_generation = ?3`,
255 ).bind(fingerprint, lease.owner, lease.generation).run();
256 }
257
258 export async function markFirebaseDelivered(env: Env, eventId: string): Promise<void> {
259 await env.DB.prepare("DELETE FROM firebase_crash_outbox WHERE event_id = ?1").bind(eventId).run();
260 }
261
262 export async function recordFirebaseRetry(
263 env: Env,
264 eventId: string,
265 state: FirebaseOutboxRow["state"],
266 attempts: number,
267 ): Promise<void> {
268 const delaySeconds = Math.min(24 * 60 * 60, 30 * 2 ** Math.min(attempts, 11));
269 const now = new Date();
270 const next = new Date(now.getTime() + delaySeconds * 1000).toISOString();
271 await env.DB.prepare(
272 `UPDATE firebase_crash_outbox
273 SET state = ?2, attempts = attempts + 1, next_attempt_at = ?3, updated_at = ?4
274 WHERE event_id = ?1`,
275 ).bind(eventId, state, next, now.toISOString()).run();
276 }
277
278 export async function dueFirebaseCrashes(env: Env): Promise<FirebaseOutboxRow[]> {
279 const result = await env.DB.prepare(
280 `SELECT event_id, fingerprint, payload, state, attempts, next_attempt_at, created_at, updated_at
281 FROM firebase_crash_outbox
282 WHERE next_attempt_at <= ?1 AND (
283 state IN ('queued', 'projected') OR
284 (state = 'processing' AND datetime(updated_at) < datetime('now', '-10 minutes'))
285 )
286 ORDER BY created_at LIMIT ?2`,
287 ).bind(new Date().toISOString(), FIREBASE_OUTBOX_BATCH).all<FirebaseOutboxRow>();
288 return result.results;
289 }
290
291 export async function purgeFirebaseDeliveryState(env: Env): Promise<void> {
292 await env.DB.batch([
293 env.DB.prepare("DELETE FROM firebase_crash_outbox WHERE datetime(created_at) < datetime('now', '-30 days')"),
294 env.DB.prepare("DELETE FROM firebase_crash_receipts WHERE datetime(projected_at) < datetime('now', '-90 days')"),
295 env.DB.prepare(
296 `UPDATE firebase_crash_group_state SET lease_owner = '', lease_expires_at = ''
297 WHERE lease_owner != '' AND datetime(lease_expires_at) < datetime('now')`,
298 ),
299 ]);
300 }
301
301 lines TYPESCRIPT