import { parseConsoleInboxEvent } from "@zhin.js/console-protocol"; const DB_NAME = "zhin-console"; const DB_VERSION = 2; const STORE_INBOX = "inbox"; const STORE_PENDING = "pending"; export type InboxRecord = { id: string; adapter: string; endpoint_id: string; kind: "message" | "request" | "notice"; payload: unknown; updatedAt: number; }; function openDb(): Promise { return new Promise((resolve, reject) => { const req = indexedDB.open(DB_NAME, DB_VERSION); req.onerror = () => reject(req.error); req.onsuccess = () => resolve(req.result); req.onupgradeneeded = () => { const db = req.result; if (!db.objectStoreNames.contains(STORE_INBOX)) { db.createObjectStore(STORE_INBOX, { keyPath: "id" }); } if (!db.objectStoreNames.contains(STORE_PENDING)) { db.createObjectStore(STORE_PENDING, { keyPath: "id" }); } }; }); } export async function idbPutInbox(record: InboxRecord): Promise { const db = await openDb(); await new Promise((resolve, reject) => { const tx = db.transaction(STORE_INBOX, "readwrite"); tx.objectStore(STORE_INBOX).put(record); tx.oncomplete = () => resolve(); tx.onerror = () => reject(tx.error); }); db.close(); } export async function idbListInbox( adapter: string, endpoint_id: string, kind: InboxRecord["kind"], ): Promise { const db = await openDb(); const all = await new Promise((resolve, reject) => { const tx = db.transaction(STORE_INBOX, "readonly"); const req = tx.objectStore(STORE_INBOX).getAll(); req.onsuccess = () => resolve((req.result as InboxRecord[]) ?? []); req.onerror = () => reject(req.error); }); db.close(); return all.filter( (r) => r.adapter === adapter && r.endpoint_id === endpoint_id && r.kind === kind, ); } export async function applyConsoleEvent(event: { type: string; data?: unknown; runtimeId?: string; eventId?: number; timestamp?: number; }): Promise { const parsed = parseConsoleInboxEvent(event); if (!parsed) return; const updatedAt = typeof event.timestamp === 'number' && Number.isFinite(event.timestamp) ? event.timestamp : Date.now(); const resumableId = typeof event.runtimeId === 'string' && event.runtimeId && Number.isSafeInteger(event.eventId) && (event.eventId ?? 0) > 0 ? `${event.runtimeId}:${event.eventId}` : null; await idbPutInbox({ // Resumable Host events are idempotent across history/live redelivery. // Unsequenced sources use a collision-safe append id without shared state. id: resumableId ?? `${parsed.adapter}:${parsed.endpointKey}:${parsed.type}:${updatedAt}:${crypto.randomUUID()}`, adapter: parsed.adapter, endpoint_id: parsed.endpointKey, kind: parsed.kind, payload: parsed.payload, updatedAt, }); }