File
Blob: src/client/hooks/use-log-stream.ts
| 1 | import { useEffect, useEffectEvent, useRef, useState } from "react"; |
| 2 | import { RunWsStateMessage, type LogEvent, type RunWsMessage } from "@/contracts"; |
| 3 | import { getApiClient } from "@/client/lib"; |
| 4 | |
| 5 | export interface UseLogStreamOptions { |
| 6 | runId: string; |
| 7 | enabled: boolean; |
| 8 | onEvent: (event: LogEvent) => void; |
| 9 | onStateUpdate?: (msg: RunWsStateMessage) => void; |
| 10 | } |
| 11 | |
| 12 | export type LogStreamStatus = "idle" | "connecting" | "connected" | "reconnecting" | "closed"; |
| 13 | |
| 14 | export const useLogStream = (options: UseLogStreamOptions): LogStreamStatus => { |
| 15 | const { runId, enabled, onEvent, onStateUpdate } = options; |
| 16 | const [status, setStatus] = useState<LogStreamStatus>("idle"); |
| 17 | |
| 18 | const wsRef = useRef<WebSocket | null>(null); |
| 19 | const reconnectTimerRef = useRef<ReturnType<typeof setTimeout> | null>(null); |
| 20 | const backoffRef = useRef(1000); |
| 21 | const emitEvent = useEffectEvent((event: LogEvent) => { |
| 22 | onEvent(event); |
| 23 | }); |
| 24 | const emitStateUpdate = useEffectEvent((message: RunWsStateMessage) => { |
| 25 | onStateUpdate?.(message); |
| 26 | }); |
| 27 | |
| 28 | useEffect(() => { |
| 29 | if (!enabled) { |
| 30 | return; |
| 31 | } |
| 32 | |
| 33 | let active = true; |
| 34 | |
| 35 | const connect = async () => { |
| 36 | setStatus("connecting"); |
| 37 | |
| 38 | try { |
| 39 | const client = getApiClient(); |
| 40 | const { ticket } = await client.getLogStreamTicket(runId); |
| 41 | |
| 42 | if (!active) return; |
| 43 | |
| 44 | const protocol = window.location.protocol === "https:" ? "wss:" : "ws:"; |
| 45 | const wsUrl = `${protocol}//${window.location.host}/api/private/runs/${encodeURIComponent(runId)}/logs?ticket=${encodeURIComponent(ticket)}`; |
| 46 | |
| 47 | const ws = new WebSocket(wsUrl); |
| 48 | wsRef.current = ws; |
| 49 | |
| 50 | ws.onopen = () => { |
| 51 | if (!active) { |
| 52 | ws.close(); |
| 53 | return; |
| 54 | } |
| 55 | setStatus("connected"); |
| 56 | backoffRef.current = 1000; |
| 57 | }; |
| 58 | |
| 59 | ws.onmessage = (event) => { |
| 60 | if (!active) return; |
| 61 | try { |
| 62 | const message = JSON.parse(event.data as string) as RunWsMessage; |
| 63 | if (message.type === "log") { |
| 64 | emitEvent(message.event); |
| 65 | } else if (message.type === "state") { |
| 66 | emitStateUpdate(message); |
| 67 | } |
| 68 | } catch { |
| 69 | // ignore malformed messages |
| 70 | } |
| 71 | }; |
| 72 | |
| 73 | ws.onclose = () => { |
| 74 | wsRef.current = null; |
| 75 | if (!active) return; |
| 76 | setStatus("reconnecting"); |
| 77 | const delay = backoffRef.current; |
| 78 | backoffRef.current = Math.min(delay * 2, 15000); |
| 79 | reconnectTimerRef.current = setTimeout(() => { |
| 80 | reconnectTimerRef.current = null; |
| 81 | if (active) void connect(); |
| 82 | }, delay); |
| 83 | }; |
| 84 | |
| 85 | ws.onerror = () => { |
| 86 | // onclose always follows onerror — reconnection handled there |
| 87 | }; |
| 88 | } catch { |
| 89 | if (!active) return; |
| 90 | setStatus("reconnecting"); |
| 91 | const delay = backoffRef.current; |
| 92 | backoffRef.current = Math.min(delay * 2, 15000); |
| 93 | reconnectTimerRef.current = setTimeout(() => { |
| 94 | reconnectTimerRef.current = null; |
| 95 | if (active) void connect(); |
| 96 | }, delay); |
| 97 | } |
| 98 | }; |
| 99 | |
| 100 | void connect(); |
| 101 | |
| 102 | return () => { |
| 103 | active = false; |
| 104 | if (reconnectTimerRef.current !== null) { |
| 105 | clearTimeout(reconnectTimerRef.current); |
| 106 | reconnectTimerRef.current = null; |
| 107 | } |
| 108 | if (wsRef.current) { |
| 109 | wsRef.current.close(); |
| 110 | wsRef.current = null; |
| 111 | } |
| 112 | }; |
| 113 | }, [enabled, runId]); |
| 114 | |
| 115 | return enabled ? status : "idle"; |
| 116 | }; |