import log from "~/utils/log"; import type { HeadscaleClient } from "./api"; /** * Defines a resource that can be fetched and polled by the live store. */ export interface ResourceDefinition { /** * A unique key to identify the resource */ readonly key: string; /** * How often to poll for changes (in milliseconds) */ readonly pollInterval: number; /** * A callback to fire to get the latest data for this resource */ readonly fetch: (client: HeadscaleClient) => Promise; } /** * Helper function to define a resource with proper typing to be used * as a keying input for the live store. * @param key A unique key to identify the resource * @param config The resource configuration */ export function defineResource( key: string, config: Omit, "key">, ): ResourceDefinition { return { key, ...config }; } export const nodesResource = defineResource("nodes", { pollInterval: 5_000, fetch: (api) => api.nodes.list(), }); export const usersResource = defineResource("users", { pollInterval: 15_000, fetch: (api) => api.users.list(), }); interface Snapshot { data: T; version: string; fetchedAt: number; } type ChangeListener = (resourceKey: string, version: string) => void; export interface LiveStore { get(resource: ResourceDefinition, apiClient: HeadscaleClient): Promise>; refresh(resource: ResourceDefinition, apiClient: HeadscaleClient): Promise; getVersions(): Record; subscribe(listener: ChangeListener): () => void; dispose(): void; } export function createLiveStore(resources: ResourceDefinition[]): LiveStore { const snapshots = new Map>(); const serializedCache = new Map(); const listeners = new Set(); const intervals = new Map>(); let storedApiClient: HeadscaleClient | undefined; let versionCounter = 0; function notifyListeners(resourceKey: string, version: string) { for (const listener of listeners) { listener(resourceKey, version); } } async function fetchResource( resource: ResourceDefinition, apiClient: HeadscaleClient, ): Promise { const data = await resource.fetch(apiClient); const json = JSON.stringify(data); const previousJson = serializedCache.get(resource.key); if (previousJson === json) { log.debug("api", "Live store: %s unchanged", resource.key); return; } const version = String(++versionCounter); serializedCache.set(resource.key, json); const snapshot: Snapshot = { data, version, fetchedAt: Date.now(), }; snapshots.set(resource.key, snapshot); log.debug("api", "Live store: %s updated (v%s)", resource.key, version); if (previousJson !== undefined) { notifyListeners(resource.key, version); } } function ensurePolling(resource: ResourceDefinition) { if (intervals.has(resource.key)) { return; } const interval = setInterval(async () => { if (!storedApiClient) { return; } try { await fetchResource(resource, storedApiClient); } catch (error) { log.error("api", "Live store: failed to poll %s", resource.key, error); } }, resource.pollInterval); intervals.set(resource.key, interval); log.debug( "api", "Live store: started polling %s every %dms", resource.key, resource.pollInterval, ); } function findResource(key: string): ResourceDefinition | undefined { return resources.find((r) => r.key === key); } return { async get( resource: ResourceDefinition, apiClient: HeadscaleClient, ): Promise> { storedApiClient = apiClient; const def = findResource(resource.key); if (!def) { throw new Error(`LiveStore: unknown resource "${resource.key}"`); } if (!snapshots.has(resource.key)) { await fetchResource(def, apiClient); } ensurePolling(def); return snapshots.get(resource.key) as Snapshot; }, async refresh(resource: ResourceDefinition, apiClient: HeadscaleClient): Promise { storedApiClient = apiClient; const def = findResource(resource.key); if (!def) { throw new Error(`LiveStore: unknown resource "${resource.key}"`); } await fetchResource(def, apiClient); }, getVersions(): Record { const versions: Record = {}; for (const [key, snapshot] of snapshots) { versions[key] = snapshot.version; } return versions; }, subscribe(listener) { listeners.add(listener); return () => { listeners.delete(listener); }; }, dispose() { for (const interval of intervals.values()) { clearInterval(interval); } intervals.clear(); snapshots.clear(); serializedCache.clear(); listeners.clear(); storedApiClient = undefined; }, }; }