import { useEffect, useRef, useState } from "react"; import type { LiveStatus, MapPoint } from "../components/map/types"; export type JobStreamState = "connecting" | "live" | "offline" | "error"; // Subscribes to live points + live position status for a job over the // backend WebSocket. Reconnects with capped exponential backoff; resubscribes // on reconnect. export function useJobStream( jobId: string | null, onPoints: (points: MapPoint[]) => void, onStatus?: (status: LiveStatus) => void, ) { const [connectionState, setConnectionState] = useState("connecting"); const pointsRef = useRef(onPoints); pointsRef.current = onPoints; const statusRef = useRef(onStatus); statusRef.current = onStatus; useEffect(() => { if (!jobId) { setConnectionState("offline"); return; } let socket: WebSocket | null = null; let closed = false; let attempt = 0; let timer: ReturnType | null = null; const connect = () => { let socketErrored = false; setConnectionState(attempt === 0 ? "connecting" : "offline"); const proto = window.location.protocol === "https:" ? "wss" : "ws"; socket = new WebSocket(`${proto}://${window.location.host}/api/ws`); socket.onopen = () => { attempt = 0; setConnectionState("live"); socket?.send( JSON.stringify({ type: "subscribe", channel: `job:${jobId}` }), ); }; socket.onmessage = (event) => { try { const msg = JSON.parse(event.data); if (msg.jobId !== jobId) { return; } if (msg.type === "points") { pointsRef.current(msg.points); } else if (msg.type === "status") { statusRef.current?.(msg as LiveStatus); } } catch { // ignore malformed frames } }; socket.onerror = () => { socketErrored = true; setConnectionState("error"); }; socket.onclose = () => { if (closed) { return; } if (!socketErrored) { setConnectionState("offline"); } attempt += 1; const delay = Math.min(1000 * 2 ** attempt, 15000); timer = setTimeout(connect, delay); }; }; connect(); return () => { closed = true; if (timer) { clearTimeout(timer); } socket?.close(); }; }, [jobId]); return connectionState; }