| 1 | import { batch, createRoot, createSignal } from "solid-js"; |
| 2 | import { LIVE_INTERVAL, type Live, type Series } from "./types/model.ts"; |
| 3 | import { parseResponse } from "hono/client"; |
| 4 | import { api } from "./api.ts"; |
| 5 | |
| 6 | /** Samples per sidebar sparkline: five minutes of ticks. */ |
| 7 | const WINDOW = 150; |
| 8 | |
| 9 | function push(values: number[], value: number | null) { |
| 10 | if (value === null) return values; |
| 11 | const next = values.length >= WINDOW ? values.slice(1) : values.slice(); |
| 12 | next.push(value); |
| 13 | return next; |
| 14 | } |
| 15 | |
| 16 | function seed(series: Series[]) { |
| 17 | return Object.fromEntries(series.map((item) => [item.name, item.v.map((v) => v ?? 0).slice(-WINDOW)])); |
| 18 | } |
| 19 | |
| 20 | /** One shared stream for every view; starts on first use by an admin page and reconnects on its own. */ |
| 21 | export const stream = createRoot(() => { |
| 22 | const [latest, setLatest] = createSignal<Live | null>(null); |
| 23 | const [trails, setTrails] = createSignal<Record<string, { cpu: number[]; memory: number[] }>>({}); |
| 24 | /** Samples appended to the trails so far, counting the seed; sparklines align their buckets to it. */ |
| 25 | const [count, setCount] = createSignal(0); |
| 26 | /** False until the first tick arrives and again while reconnecting, so views can show the numbers are stale. */ |
| 27 | const [live, setLive] = createSignal(false); |
| 28 | let started = false; |
| 29 | let failures = 0; |
| 30 | let stale: ReturnType<typeof setTimeout> | undefined; |
| 31 | |
| 32 | function retry() { |
| 33 | setLive(false); |
| 34 | setTimeout(connect, Math.min(30, 2 ** failures++) * 1000); |
| 35 | } |
| 36 | |
| 37 | // EventSource gives up for good on a non-200 answer, like Vite's 502 while the server restarts. |
| 38 | function connect() { |
| 39 | const events = new EventSource("/api/live"); |
| 40 | events.onmessage = (event) => { |
| 41 | const tick: Live = JSON.parse(event.data); |
| 42 | failures = 0; |
| 43 | clearTimeout(stale); |
| 44 | const remaining = tick.t * 1000 + LIVE_INTERVAL * 3_000 - Date.now(); |
| 45 | stale = setTimeout(() => setLive(false), Math.max(0, remaining)); |
| 46 | batch(() => { |
| 47 | setLive(remaining > 0); |
| 48 | setLatest(tick); |
| 49 | setCount(count() + 1); |
| 50 | setTrails((previous) => Object.fromEntries(Object.entries(tick.services).map(([id, sample]) => { |
| 51 | const trail = previous[id] ?? { cpu: [], memory: [] }; |
| 52 | return [id, { cpu: push(trail.cpu, sample.cpu), memory: push(trail.memory, sample.memory) }]; |
| 53 | }))); |
| 54 | }); |
| 55 | }; |
| 56 | events.onerror = () => { |
| 57 | events.close(); |
| 58 | retry(); |
| 59 | }; |
| 60 | } |
| 61 | |
| 62 | function start() { |
| 63 | if (started) return; |
| 64 | started = true; |
| 65 | connect(); |
| 66 | const history = (metric: "service.cpu" | "service.memory") => |
| 67 | parseResponse(api.metrics[":metric"].$get({ query: { range: String(WINDOW * LIVE_INTERVAL) }, param: { metric } })).catch(() => []); |
| 68 | void Promise.all([history("service.cpu"), history("service.memory")]).then(([cpu, memory]) => { |
| 69 | const cpuTrail = seed(cpu); |
| 70 | const memoryTrail = seed(memory); |
| 71 | batch(() => { |
| 72 | setTrails((previous) => ({ ...previous, ...Object.fromEntries(Object.keys(cpuTrail).map((id) => |
| 73 | [id, { |
| 74 | cpu: [...cpuTrail[id]!, ...(previous[id]?.cpu ?? [])].slice(-WINDOW), |
| 75 | memory: [...(memoryTrail[id] ?? []), ...(previous[id]?.memory ?? [])].slice(-WINDOW), |
| 76 | }])) })); |
| 77 | setCount(count() + (Object.values(cpuTrail)[0]?.length ?? 0)); |
| 78 | }); |
| 79 | }); |
| 80 | } |
| 81 | |
| 82 | return { latest, trails, count, live, start }; |
| 83 | }); |