| 1 | import type { Env } from "./env"; |
| 2 | import { |
| 3 | FIREBASE_ARCHIVING_RESERVATION_BYTES, |
| 4 | FIREBASE_ACTIVE_RESERVATION_BYTES, |
| 5 | FIREBASE_COMPACTED_RESERVATION_BYTES, |
| 6 | FIREBASE_OUTBOX_LIMIT, |
| 7 | FIREBASE_STORAGE_BUDGET_BYTES, |
| 8 | acquireFirebaseGroupLease, |
| 9 | firebaseGroupState, |
| 10 | releaseFirebaseGroupLease, |
| 11 | renewFirebaseGroupLease, |
| 12 | type FirebaseGroupLease, |
| 13 | type FirebaseSampleState, |
| 14 | } from "./crash_delivery"; |
| 15 | import { firebaseMeta, loadFirebaseGroupMeta } from "./firebase_crash_view"; |
| 16 | import { |
| 17 | deleteFirebaseCrashGroupConditional, |
| 18 | writeFirebaseGroupMeta, |
| 19 | writeFirebaseSampleMarkers, |
| 20 | } from "./firebase_rtdb"; |
| 21 | |
| 22 | const LIFECYCLE_BATCH = 20; |
| 23 | |
| 24 | export type FirebaseStorageSummary = { |
| 25 | active: number; |
| 26 | compacted: number; |
| 27 | archiving: number; |
| 28 | archived: number; |
| 29 | reservedBytes: number; |
| 30 | budgetBytes: number; |
| 31 | outboxCount: number; |
| 32 | oldestOutboxSeconds: number; |
| 33 | stuckArchiving: number; |
| 34 | }; |
| 35 | |
| 36 | type LifecycleCandidate = { |
| 37 | fingerprint: string; |
| 38 | sample_state: FirebaseSampleState; |
| 39 | }; |
| 40 | |
| 41 | function renewer(env: Env, fingerprint: string, lease: FirebaseGroupLease) { |
| 42 | return () => renewFirebaseGroupLease(env, fingerprint, lease); |
| 43 | } |
| 44 | |
| 45 | async function pendingOutbox(env: Env, fingerprint: string): Promise<boolean> { |
| 46 | return Boolean(await env.DB.prepare( |
| 47 | "SELECT event_id FROM firebase_crash_outbox WHERE fingerprint = ?1 LIMIT 1", |
| 48 | ).bind(fingerprint).first()); |
| 49 | } |
| 50 | |
| 51 | async function compactGroup(env: Env, fingerprint: string, lease: FirebaseGroupLease, now: string): Promise<void> { |
| 52 | const [group, state] = await Promise.all([ |
| 53 | loadFirebaseGroupMeta(env, fingerprint), firebaseGroupState(env, fingerprint), |
| 54 | ]); |
| 55 | if (!group || !state || state.sample_state !== "active" || await pendingOutbox(env, fingerprint)) return; |
| 56 | const renew = renewer(env, fingerprint, lease); |
| 57 | await writeFirebaseSampleMarkers( |
| 58 | env, fingerprint, Number(group.count), lease.generation, Number(state.sample_epoch), |
| 59 | "compacted", false, renew, |
| 60 | ); |
| 61 | await writeFirebaseGroupMeta( |
| 62 | env, fingerprint, firebaseMeta(group), lease.generation, Number(state.sample_epoch), |
| 63 | "compacted", renew, |
| 64 | ); |
| 65 | await env.DB.prepare( |
| 66 | `UPDATE firebase_crash_group_state |
| 67 | SET sample_state = 'compacted', reserved_bytes = ?4, compacted_at = ?5 |
| 68 | WHERE fingerprint = ?1 AND sample_state = 'active' |
| 69 | AND lease_owner = ?2 AND lease_generation = ?3 |
| 70 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1)`, |
| 71 | ).bind(fingerprint, lease.owner, lease.generation, FIREBASE_COMPACTED_RESERVATION_BYTES, now).run(); |
| 72 | } |
| 73 | |
| 74 | async function beginArchive( |
| 75 | env: Env, |
| 76 | fingerprint: string, |
| 77 | lease: FirebaseGroupLease, |
| 78 | now: string, |
| 79 | reason: "retention" | "admin", |
| 80 | ): Promise<void> { |
| 81 | const [group, state] = await Promise.all([ |
| 82 | loadFirebaseGroupMeta(env, fingerprint), firebaseGroupState(env, fingerprint), |
| 83 | ]); |
| 84 | if (!group || !state || !["active", "compacted"].includes(state.sample_state)) { |
| 85 | if (reason === "admin") throw new Error("firebase crash group lifecycle state is missing"); |
| 86 | return; |
| 87 | } |
| 88 | if (reason === "retention" && await pendingOutbox(env, fingerprint)) return; |
| 89 | const renew = renewer(env, fingerprint, lease); |
| 90 | await writeFirebaseSampleMarkers( |
| 91 | env, fingerprint, Number(group.count), lease.generation, Number(state.sample_epoch), |
| 92 | "archiving", true, renew, |
| 93 | ); |
| 94 | await writeFirebaseGroupMeta( |
| 95 | env, fingerprint, firebaseMeta(group), lease.generation, Number(state.sample_epoch), |
| 96 | "archiving", renew, |
| 97 | ); |
| 98 | if (reason === "admin") { |
| 99 | await env.DB.batch([ |
| 100 | env.DB.prepare("DELETE FROM reports WHERE fingerprint = ?1").bind(fingerprint), |
| 101 | env.DB.prepare("DELETE FROM report_daily WHERE fingerprint = ?1").bind(fingerprint), |
| 102 | env.DB.prepare("DELETE FROM report_installations WHERE fingerprint = ?1").bind(fingerprint), |
| 103 | env.DB.prepare("DELETE FROM report_event_dimensions WHERE fingerprint = ?1").bind(fingerprint), |
| 104 | env.DB.prepare( |
| 105 | "DELETE FROM firebase_crash_outbox WHERE fingerprint = ?1 AND datetime(created_at) <= datetime(?2)", |
| 106 | ).bind(fingerprint, now), |
| 107 | env.DB.prepare("DELETE FROM groups WHERE fingerprint = ?1").bind(fingerprint), |
| 108 | env.DB.prepare( |
| 109 | `UPDATE firebase_crash_group_state |
| 110 | SET sample_state = CASE WHEN EXISTS ( |
| 111 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 112 | ) THEN 'active' ELSE 'archiving' END, |
| 113 | sample_epoch = sample_epoch + CASE WHEN EXISTS ( |
| 114 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 115 | ) THEN 1 ELSE 0 END, |
| 116 | epoch_first_event_id = CASE WHEN EXISTS ( |
| 117 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 118 | ) THEN '' ELSE epoch_first_event_id END, |
| 119 | reserved_bytes = CASE WHEN EXISTS ( |
| 120 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 121 | ) THEN ?4 ELSE ?5 END, |
| 122 | archived_at = CASE WHEN EXISTS ( |
| 123 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 124 | ) THEN '' ELSE ?6 END, |
| 125 | archive_reason = CASE WHEN EXISTS ( |
| 126 | SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1 |
| 127 | ) THEN '' ELSE 'admin' END |
| 128 | WHERE fingerprint = ?1 AND lease_owner = ?2 AND lease_generation = ?3`, |
| 129 | ).bind( |
| 130 | fingerprint, lease.owner, lease.generation, |
| 131 | FIREBASE_ACTIVE_RESERVATION_BYTES, FIREBASE_ARCHIVING_RESERVATION_BYTES, now, |
| 132 | ), |
| 133 | ]); |
| 134 | return; |
| 135 | } |
| 136 | await env.DB.prepare( |
| 137 | `UPDATE firebase_crash_group_state |
| 138 | SET sample_state = 'archiving', reserved_bytes = ?4, archived_at = ?5, archive_reason = 'retention' |
| 139 | WHERE fingerprint = ?1 AND sample_state IN ('active', 'compacted') |
| 140 | AND lease_owner = ?2 AND lease_generation = ?3 |
| 141 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1)`, |
| 142 | ).bind(fingerprint, lease.owner, lease.generation, FIREBASE_ARCHIVING_RESERVATION_BYTES, now).run(); |
| 143 | } |
| 144 | |
| 145 | async function finishArchive(env: Env, fingerprint: string, lease: FirebaseGroupLease): Promise<void> { |
| 146 | const state = await firebaseGroupState(env, fingerprint); |
| 147 | if (!state || state.sample_state !== "archiving" || await pendingOutbox(env, fingerprint)) return; |
| 148 | await deleteFirebaseCrashGroupConditional(env, fingerprint, renewer(env, fingerprint, lease)); |
| 149 | if (state.archive_reason === "admin") { |
| 150 | await env.DB.prepare( |
| 151 | `DELETE FROM firebase_crash_group_state |
| 152 | WHERE fingerprint = ?1 AND sample_state = 'archiving' |
| 153 | AND lease_owner = ?2 AND lease_generation = ?3 |
| 154 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1)`, |
| 155 | ).bind(fingerprint, lease.owner, lease.generation).run(); |
| 156 | } else { |
| 157 | await env.DB.prepare( |
| 158 | `UPDATE firebase_crash_group_state |
| 159 | SET sample_state = 'archived', reserved_bytes = 0, lease_owner = '', lease_expires_at = '' |
| 160 | WHERE fingerprint = ?1 AND sample_state = 'archiving' |
| 161 | AND lease_owner = ?2 AND lease_generation = ?3 |
| 162 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = ?1)`, |
| 163 | ).bind(fingerprint, lease.owner, lease.generation).run(); |
| 164 | } |
| 165 | } |
| 166 | |
| 167 | async function withLease( |
| 168 | env: Env, |
| 169 | candidate: LifecycleCandidate, |
| 170 | operation: (lease: FirebaseGroupLease) => Promise<void>, |
| 171 | ): Promise<boolean> { |
| 172 | const lease = await acquireFirebaseGroupLease(env, candidate.fingerprint); |
| 173 | if (!lease) return false; |
| 174 | try { |
| 175 | await operation(lease); |
| 176 | return true; |
| 177 | } finally { |
| 178 | await releaseFirebaseGroupLease(env, candidate.fingerprint, lease); |
| 179 | } |
| 180 | } |
| 181 | |
| 182 | export async function archiveFirebaseGroupForAdmin(env: Env, fingerprint: string): Promise<void> { |
| 183 | const lease = await acquireFirebaseGroupLease(env, fingerprint); |
| 184 | if (!lease) throw new Error("firebase crash group is busy"); |
| 185 | try { |
| 186 | await beginArchive(env, fingerprint, lease, new Date().toISOString(), "admin"); |
| 187 | } finally { |
| 188 | await releaseFirebaseGroupLease(env, fingerprint, lease); |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | export async function runFirebaseCrashLifecycle(env: Env, now = new Date()): Promise<void> { |
| 193 | const iso = now.toISOString(); |
| 194 | let remaining = LIFECYCLE_BATCH; |
| 195 | const queries = [ |
| 196 | `SELECT fingerprint, sample_state FROM firebase_crash_group_state |
| 197 | WHERE sample_state = 'archiving' AND datetime(archived_at) <= datetime(?1, '-24 hours') |
| 198 | ORDER BY archived_at LIMIT ?2`, |
| 199 | `SELECT state.fingerprint, state.sample_state FROM firebase_crash_group_state AS state |
| 200 | JOIN groups USING (fingerprint) |
| 201 | WHERE state.sample_state IN ('active', 'compacted') AND groups.status IN ('resolved', 'ignored') |
| 202 | AND datetime(groups.last_seen) <= datetime(?1, '-60 days') |
| 203 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = state.fingerprint) |
| 204 | ORDER BY groups.last_seen LIMIT ?2`, |
| 205 | `SELECT state.fingerprint, state.sample_state FROM firebase_crash_group_state AS state |
| 206 | JOIN groups USING (fingerprint) |
| 207 | WHERE state.sample_state = 'active' AND groups.status IN ('resolved', 'ignored') |
| 208 | AND datetime(groups.last_seen) <= datetime(?1, '-30 days') |
| 209 | AND NOT EXISTS (SELECT 1 FROM firebase_crash_outbox WHERE fingerprint = state.fingerprint) |
| 210 | ORDER BY groups.last_seen LIMIT ?2`, |
| 211 | ]; |
| 212 | for (let phase = 0; phase < queries.length && remaining > 0; phase++) { |
| 213 | const candidates = await env.DB.prepare(queries[phase]).bind(iso, remaining).all<LifecycleCandidate>(); |
| 214 | for (const candidate of candidates.results) { |
| 215 | const processed = await withLease(env, candidate, async (lease) => { |
| 216 | if (phase === 0) await finishArchive(env, candidate.fingerprint, lease); |
| 217 | else if (phase === 1) await beginArchive(env, candidate.fingerprint, lease, iso, "retention"); |
| 218 | else await compactGroup(env, candidate.fingerprint, lease, iso); |
| 219 | }); |
| 220 | if (processed) remaining--; |
| 221 | if (remaining === 0) break; |
| 222 | } |
| 223 | } |
| 224 | } |
| 225 | |
| 226 | export async function firebaseStorageSummary(env: Env): Promise<FirebaseStorageSummary> { |
| 227 | const [state, outbox] = await Promise.all([ |
| 228 | env.DB.prepare( |
| 229 | `SELECT |
| 230 | SUM(CASE WHEN sample_state = 'active' THEN 1 ELSE 0 END) AS active, |
| 231 | SUM(CASE WHEN sample_state = 'compacted' THEN 1 ELSE 0 END) AS compacted, |
| 232 | SUM(CASE WHEN sample_state = 'archiving' THEN 1 ELSE 0 END) AS archiving, |
| 233 | SUM(CASE WHEN sample_state = 'archived' THEN 1 ELSE 0 END) AS archived, |
| 234 | COALESCE(SUM(reserved_bytes), 0) AS reserved_bytes, |
| 235 | SUM(CASE WHEN sample_state = 'archiving' AND datetime(archived_at) < datetime('now', '-48 hours') |
| 236 | THEN 1 ELSE 0 END) AS stuck_archiving |
| 237 | FROM firebase_crash_group_state`, |
| 238 | ).first<Record<string, number>>(), |
| 239 | env.DB.prepare( |
| 240 | `SELECT COUNT(*) AS outbox_count, |
| 241 | COALESCE(CAST((julianday('now') - julianday(MIN(created_at))) * 86400 AS INTEGER), 0) |
| 242 | AS oldest_outbox_seconds |
| 243 | FROM firebase_crash_outbox`, |
| 244 | ).first<Record<string, number>>(), |
| 245 | ]); |
| 246 | return { |
| 247 | active: Number(state?.active ?? 0), |
| 248 | compacted: Number(state?.compacted ?? 0), |
| 249 | archiving: Number(state?.archiving ?? 0), |
| 250 | archived: Number(state?.archived ?? 0), |
| 251 | reservedBytes: Number(state?.reserved_bytes ?? 0), |
| 252 | budgetBytes: FIREBASE_STORAGE_BUDGET_BYTES, |
| 253 | outboxCount: Number(outbox?.outbox_count ?? 0), |
| 254 | oldestOutboxSeconds: Number(outbox?.oldest_outbox_seconds ?? 0), |
| 255 | stuckArchiving: Number(state?.stuck_archiving ?? 0), |
| 256 | }; |
| 257 | } |
| 258 | |
| 259 | export const FIREBASE_OUTBOX_WARNING = Math.floor(FIREBASE_OUTBOX_LIMIT * 0.8); |
| 260 |