| 1 | import type { Env } from "./env"; |
| 2 | import { Report, type ReportPayload } from "./report_schema"; |
| 3 | import { |
| 4 | acquireFirebaseGroupLease, |
| 5 | claimFirebaseCrash, |
| 6 | crashStorageMode, |
| 7 | dueFirebaseCrashes, |
| 8 | firebaseGroupState, |
| 9 | firebaseOutboxExists, |
| 10 | firebaseProjectionReceipt, |
| 11 | firebaseStorageReady, |
| 12 | markFirebaseDelivered, |
| 13 | recordFirebaseRetry, |
| 14 | releaseFirebaseGroupLease, |
| 15 | renewFirebaseGroupLease, |
| 16 | type FirebaseGroupLease, |
| 17 | type FirebaseOutboxRow, |
| 18 | } from "./crash_delivery"; |
| 19 | import { firebaseMeta, loadFirebaseGroupMeta } from "./firebase_crash_view"; |
| 20 | import { writeFirebaseCrashGroup } from "./firebase_rtdb"; |
| 21 | |
| 22 | const LATEST_SAMPLES_PER_GROUP = 5; |
| 23 | |
| 24 | export type StoredCrashEvent = { |
| 25 | eventId: string; |
| 26 | fingerprint: string; |
| 27 | receivedAt: string; |
| 28 | keepD1Sample: boolean; |
| 29 | report: ReportPayload; |
| 30 | }; |
| 31 | |
| 32 | export async function deliverCrashEventToFirebase( |
| 33 | env: Env, |
| 34 | event: StoredCrashEvent, |
| 35 | attempts: number, |
| 36 | lease: FirebaseGroupLease, |
| 37 | ): Promise<boolean> { |
| 38 | try { |
| 39 | const [group, receipt, state] = await Promise.all([ |
| 40 | loadFirebaseGroupMeta(env, event.fingerprint), |
| 41 | firebaseProjectionReceipt(env, event.eventId), |
| 42 | firebaseGroupState(env, event.fingerprint), |
| 43 | ]); |
| 44 | if (!group || !receipt || !state || state.sample_state !== "active") { |
| 45 | throw new Error("firebase projection state is missing"); |
| 46 | } |
| 47 | const { installId: _installId, ...publicReport } = event.report; |
| 48 | const stillLatest = Number(receipt.group_count) > Number(group.count) - LATEST_SAMPLES_PER_GROUP; |
| 49 | const renew = () => renewFirebaseGroupLease(env, event.fingerprint, lease); |
| 50 | await writeFirebaseCrashGroup( |
| 51 | env, |
| 52 | firebaseMeta(group), |
| 53 | { ...publicReport, eventId: event.eventId, receivedAt: event.receivedAt }, |
| 54 | stillLatest ? Number(receipt.latest_slot) : null, |
| 55 | Number(receipt.first_sample) === 1, |
| 56 | lease.generation, |
| 57 | Number(state.sample_epoch), |
| 58 | renew, |
| 59 | ); |
| 60 | if (!await renew()) throw new Error("firebase crash group lease was fenced before delivery completion"); |
| 61 | await markFirebaseDelivered(env, event.eventId); |
| 62 | return true; |
| 63 | } catch (error) { |
| 64 | console.error("firebase crash delivery failed", error); |
| 65 | await recordFirebaseRetry(env, event.eventId, "projected", attempts); |
| 66 | return false; |
| 67 | } |
| 68 | } |
| 69 | |
| 70 | function parseStoredCrash(row: FirebaseOutboxRow): StoredCrashEvent | undefined { |
| 71 | try { |
| 72 | const value = JSON.parse(row.payload) as Partial<StoredCrashEvent>; |
| 73 | const report = Report.safeParse(value.report); |
| 74 | if ( |
| 75 | !report.success || value.eventId !== row.event_id || value.fingerprint !== row.fingerprint || |
| 76 | !/^(?:dev:)?[0-9a-f]{64}$/.test(value.fingerprint ?? "") || |
| 77 | typeof value.receivedAt !== "string" || typeof value.keepD1Sample !== "boolean" |
| 78 | ) return undefined; |
| 79 | return { ...value, report: report.data } as StoredCrashEvent; |
| 80 | } catch { |
| 81 | return undefined; |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | export async function drainFirebaseCrashOutbox( |
| 86 | env: Env, |
| 87 | projectCrashEvent: (env: Env, event: StoredCrashEvent) => Promise<void>, |
| 88 | ): Promise<void> { |
| 89 | if (crashStorageMode(env) === "d1" || !firebaseStorageReady(env)) return; |
| 90 | for (const row of await dueFirebaseCrashes(env)) { |
| 91 | const event = parseStoredCrash(row); |
| 92 | if (!event) { |
| 93 | console.error(`firebase crash outbox row ${row.event_id} is invalid`); |
| 94 | await recordFirebaseRetry(env, row.event_id, "queued", row.attempts); |
| 95 | continue; |
| 96 | } |
| 97 | const lease = await acquireFirebaseGroupLease(env, row.fingerprint); |
| 98 | if (!lease) continue; |
| 99 | try { |
| 100 | if (!await firebaseOutboxExists(env, row.event_id)) continue; |
| 101 | if (row.state !== "projected") { |
| 102 | if (!await claimFirebaseCrash(env, row.event_id, new Date().toISOString())) continue; |
| 103 | try { |
| 104 | await projectCrashEvent(env, event); |
| 105 | } catch (error) { |
| 106 | console.error("firebase crash projection failed", error); |
| 107 | await recordFirebaseRetry(env, row.event_id, "queued", row.attempts); |
| 108 | continue; |
| 109 | } |
| 110 | } |
| 111 | if (!await deliverCrashEventToFirebase(env, event, row.attempts, lease)) break; |
| 112 | } finally { |
| 113 | await releaseFirebaseGroupLease(env, row.fingerprint, lease).catch((error) => { |
| 114 | console.error("firebase crash group lease release failed", error); |
| 115 | }); |
| 116 | } |
| 117 | } |
| 118 | } |
| 119 |