137 lines
3.6 KiB
TypeScript
137 lines
3.6 KiB
TypeScript
"use client";
|
|||
|
|
|
||
|
|
import { useEffect, useState } from "react";
|
||
|
|
import type { Config } from "@/lib/config/schema";
|
||
|
|
import type { WidgetResult } from "@/lib/types/widget-result";
|
||
|
|
|
||
|
|
interface Envelope {
|
||
|
|
topic: string;
|
||
|
|
type: string;
|
||
|
|
ts: number;
|
||
|
|
data: unknown;
|
||
|
|
}
|
||
|
|
|
||
|
|
type TopicListener = (envelope: Envelope) => void;
|
||
|
|
|
||
|
|
let sharedSocket: WebSocket | null = null;
|
||
|
|
let refCount = 0;
|
||
|
|
let reconnectAttempts = 0;
|
||
|
|
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
|
||
|
|
const topicListeners = new Map<string, Set<TopicListener>>();
|
||
|
|
const pendingWidgetSubscriptions = new Set<string>();
|
||
|
|
|
||
|
|
function getSocketUrl(): string {
|
||
|
|
const protocol = window.location.protocol === "https:" ? "wss:" : "ws:";
|
||
|
|
return `${protocol}//${window.location.host}/ws`;
|
||
|
|
}
|
||
|
|
|
||
|
|
function scheduleReconnect(): void {
|
||
|
|
if (reconnectTimer || refCount === 0) return;
|
||
|
|
const delay = Math.min(30_000, 1_000 * 2 ** reconnectAttempts);
|
||
|
|
reconnectAttempts += 1;
|
||
|
|
reconnectTimer = setTimeout(() => {
|
||
|
|
reconnectTimer = null;
|
||
|
|
if (refCount > 0) ensureSocket();
|
||
|
|
}, delay);
|
||
|
|
}
|
||
|
|
|
||
|
|
function ensureSocket(): WebSocket {
|
||
|
|
if (sharedSocket && sharedSocket.readyState <= WebSocket.OPEN) {
|
||
|
|
return sharedSocket;
|
||
|
|
}
|
||
|
|
|
||
|
|
const socket = new WebSocket(getSocketUrl());
|
||
|
|
sharedSocket = socket;
|
||
|
|
|
||
|
|
socket.addEventListener("open", () => {
|
||
|
|
reconnectAttempts = 0;
|
||
|
|
for (const widgetId of pendingWidgetSubscriptions) {
|
||
|
|
socket.send(JSON.stringify({ action: "subscribe", widgetId }));
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
socket.addEventListener("message", (event) => {
|
||
|
|
let envelope: Envelope;
|
||
|
|
try {
|
||
|
|
envelope = JSON.parse(event.data);
|
||
|
|
} catch {
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
const listeners = topicListeners.get(envelope.topic);
|
||
|
|
if (!listeners) return;
|
||
|
|
for (const listener of listeners) listener(envelope);
|
||
|
|
});
|
||
|
|
|
||
|
|
socket.addEventListener("close", () => {
|
||
|
|
if (sharedSocket === socket) sharedSocket = null;
|
||
|
|
scheduleReconnect();
|
||
|
|
});
|
||
|
|
|
||
|
|
return socket;
|
||
|
|
}
|
||
|
|
|
||
|
|
function addTopicListener(topic: string, listener: TopicListener): () => void {
|
||
|
|
let listeners = topicListeners.get(topic);
|
||
|
|
if (!listeners) {
|
||
|
|
listeners = new Set();
|
||
|
|
topicListeners.set(topic, listeners);
|
||
|
|
}
|
||
|
|
listeners.add(listener);
|
||
|
|
return () => {
|
||
|
|
listeners?.delete(listener);
|
||
|
|
if (listeners && listeners.size === 0) topicListeners.delete(topic);
|
||
|
|
};
|
||
|
|
}
|
||
|
|
|
||
|
|
export function useWidgetSubscription(widgetId: string): WidgetResult | null {
|
||
|
|
const [result, setResult] = useState<WidgetResult | null>(null);
|
||
|
|
|
||
|
|
useEffect(() => {
|
||
|
|
refCount += 1;
|
||
|
|
pendingWidgetSubscriptions.add(widgetId);
|
||
|
|
const socket = ensureSocket();
|
||
|
|
const topic = `widget:${widgetId}`;
|
||
|
|
|
||
|
|
const removeListener = addTopicListener(topic, (envelope) => {
|
||
|
|
if (envelope.type === "result") setResult(envelope.data as WidgetResult);
|
||
|
|
});
|
||
|
|
|
||
|
|
if (socket.readyState === WebSocket.OPEN) {
|
||
|
|
socket.send(JSON.stringify({ action: "subscribe", widgetId }));
|
||
|
|
}
|
||
|
|
|
||
|
|
return () => {
|
||
|
|
refCount -= 1;
|
||
|
|
removeListener();
|
||
|
|
if (!topicListeners.has(topic)) {
|
||
|
|
pendingWidgetSubscriptions.delete(widgetId);
|
||
|
|
if (socket.readyState === WebSocket.OPEN) {
|
||
|
|
socket.send(JSON.stringify({ action: "unsubscribe", widgetId }));
|
||
|
|
}
|
||
|
|
}
|
||
|
|
};
|
||
|
|
}, [widgetId]);
|
||
|
|
|
||
|
|
return result;
|
||
|
|
}
|
||
|
|
|
||
|
|
export function useConfigSubscription(initial: Config): Config {
|
||
|
|
const [config, setConfig] = useState(initial);
|
||
|
|
|
||
|
|
useEffect(() => {
|
||
|
|
refCount += 1;
|
||
|
|
ensureSocket();
|
||
|
|
|
||
|
|
const removeListener = addTopicListener("config", (envelope) => {
|
||
|
|
if (envelope.type === "config:update") setConfig(envelope.data as Config);
|
||
|
|
});
|
||
|
|
|
||
|
|
return () => {
|
||
|
|
refCount -= 1;
|
||
|
|
removeListener();
|
||
|
|
};
|
||
|
|
}, []);
|
||
|
|
|
||
|
|
return config;
|
||
|
|
}
|