Skip to content
File

Blob: src/client/hooks/use-log-stream.ts

typescript117 lines
1import { useEffect, useEffectEvent, useRef, useState } from "react";
2import { RunWsStateMessage, type LogEvent, type RunWsMessage } from "@/contracts";
3import { getApiClient } from "@/client/lib";
4 
5export interface UseLogStreamOptions {
6 runId: string;
7 enabled: boolean;
8 onEvent: (event: LogEvent) => void;
9 onStateUpdate?: (msg: RunWsStateMessage) => void;
10}
11 
12export type LogStreamStatus = "idle" | "connecting" | "connected" | "reconnecting" | "closed";
13 
14export 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};