import { openDB, type IDBPDatabase } from "idb"; import type { DurableSegment, SegmentDurability, SegmentDurabilityPutResult, } from "../chunk-uploader/segment-uploader.js"; export const REPLAY_CHUNK_DB_PREFIX = "proctoring-replay"; const BYTES_KEY = "bytes"; const LAST_SEQUENCE_KEY = "lastSequence"; const DEFAULT_MAX_BYTES = 32 * 1024 * 1024; interface StoredReplayChunk { sequence: number; buffer: ArrayBuffer; mimeType: string; byteSize: number; timingHeaders?: Record; } export interface ReplayChunkStoreOptions { sessionId: string; maxBytes?: number; } export function replayChunkDbName(sessionId: string): string { return `${REPLAY_CHUNK_DB_PREFIX}-${sessionId}`; } /** * Durable, byte-capped replay queue. The newest independently playable * window is always retained; oldest unacknowledged windows are evicted first * when the origin is offline long enough to reach the disk budget. */ export class ReplayChunkStore implements SegmentDurability { private readonly dbName: string; private readonly maxBytes: number; private db: IDBPDatabase | null = null; constructor(options: ReplayChunkStoreOptions) { this.dbName = replayChunkDbName(options.sessionId); this.maxBytes = Math.max(1, options.maxBytes ?? DEFAULT_MAX_BYTES); } name(): string { return this.dbName; } async open(): Promise { if (this.db) return; this.db = await openDB(this.dbName, 1, { upgrade(db) { db.createObjectStore("segments", { keyPath: "sequence" }); db.createObjectStore("meta", { keyPath: "key" }); }, }); } async list(): Promise { const records = (await this.requireDb().getAll("segments")) as StoredReplayChunk[]; records.sort((left, right) => left.sequence - right.sequence); return records.map(toDurableSegment); } async listSequences(): Promise { const keys = await this.requireDb().getAllKeys("segments"); return keys .map(Number) .filter((value) => Number.isSafeInteger(value) && value > 0) .sort((left, right) => left - right); } async get(sequence: number): Promise { const record = (await this.requireDb().get("segments", sequence)) as | StoredReplayChunk | undefined; return record ? toDurableSegment(record) : null; } async put(record: DurableSegment): Promise { const db = this.requireDb(); const buffer = await record.blob.arrayBuffer(); const tx = db.transaction(["segments", "meta"], "readwrite"); const segments = tx.objectStore("segments"); const meta = tx.objectStore("meta"); const prior = (await segments.get(record.sequence)) as StoredReplayChunk | undefined; let total = Number(((await meta.get(BYTES_KEY)) as { value?: number } | undefined)?.value ?? 0); total -= prior?.byteSize ?? 0; total += buffer.byteLength; await segments.put({ sequence: record.sequence, buffer, mimeType: record.blob.type, byteSize: buffer.byteLength, ...(record.timingHeaders ? { timingHeaders: record.timingHeaders } : {}), } satisfies StoredReplayChunk); const evicted: number[] = []; let cursor = await segments.openCursor(); while (cursor && total > this.maxBytes) { const stored = cursor.value as StoredReplayChunk; if (stored.sequence !== record.sequence) { total -= stored.byteSize; evicted.push(stored.sequence); await cursor.delete(); } cursor = await cursor.continue(); } await meta.put({ key: BYTES_KEY, value: Math.max(0, total) }); await tx.done; return { evicted }; } async remove(sequence: number): Promise { const db = this.requireDb(); const tx = db.transaction(["segments", "meta"], "readwrite"); const segments = tx.objectStore("segments"); const meta = tx.objectStore("meta"); const existing = (await segments.get(sequence)) as StoredReplayChunk | undefined; if (existing) await segments.delete(sequence); const total = Number( ((await meta.get(BYTES_KEY)) as { value?: number } | undefined)?.value ?? 0, ); await meta.put({ key: BYTES_KEY, value: Math.max(0, total - (existing?.byteSize ?? 0)), }); await tx.done; } async nextSequence(preferred: number): Promise { const db = this.requireDb(); const tx = db.transaction("meta", "readwrite"); const store = tx.objectStore("meta"); const prior = (await store.get(LAST_SEQUENCE_KEY)) as | { key: string; value: number } | undefined; const sequence = Math.max( Number.isSafeInteger(preferred) && preferred > 0 ? preferred : 1, (prior?.value ?? 0) + 1, ); await store.put({ key: LAST_SEQUENCE_KEY, value: sequence }); await tx.done; return sequence; } async bytes(): Promise { const row = (await this.requireDb().get("meta", BYTES_KEY)) as { value?: number } | undefined; return Number(row?.value ?? 0); } async close(): Promise { this.db?.close(); this.db = null; } private requireDb(): IDBPDatabase { if (!this.db) throw new Error(`Replay chunk store ${this.dbName} is not open`); return this.db; } } function toDurableSegment(record: StoredReplayChunk): DurableSegment { return { sequence: record.sequence, blob: new Blob([record.buffer], { type: record.mimeType || "application/octet-stream", }), ...(record.timingHeaders ? { timingHeaders: record.timingHeaders } : {}), }; }