import { writable } from 'svelte/store' export interface OikosEvent { id: number ts: string type: string entity_id?: string | null severity: 'info' | 'warning' | 'critical' source: string data?: unknown correlation_id?: string | null } const MAX_BUFFERED = 200 export const liveEvents = writable([]) export const connectionState = writable<'connecting' | 'open' | 'closed'>('connecting') let source: EventSource | null = null let subscriberCount = 0 function connect() { if (source) return connectionState.set('connecting') // The browser's EventSource sends Last-Event-ID automatically on reconnect. source = new EventSource('/api/v1/events/stream') source.onopen = () => connectionState.set('open') source.onmessage = (ev) => { try { const parsed: OikosEvent = JSON.parse(ev.data) liveEvents.update((events) => [parsed, ...events].slice(0, MAX_BUFFERED)) } catch { // skip malformed } } source.onerror = () => { connectionState.set('closed') } } function disconnect() { source?.close() source = null } // Reference-counted: the stream stays open as long as at least one page subscribes. export function subscribeEvents(): () => void { subscriberCount++ if (subscriberCount === 1) connect() return () => { subscriberCount-- if (subscriberCount === 0) disconnect() } }