1import { batch, createRoot, createSignal } from "solid-js";
2import { LIVE_INTERVAL, type Live, type Series } from "./types/model.ts";
3import { parseResponse } from "hono/client";
4import { api } from "./api.ts";
5
6/** Samples per sidebar sparkline: five minutes of ticks. */
7const WINDOW = 150;
8
9function 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
16function 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. */
21export 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});