import { useEffect, useState } from "react"; import { api } from "../lib/api"; const LIVE_POLL_MS = 2000; export type RuntimeEvent = Record; export type AmaSessionPhase = "loading" | "ready" | "error"; function sequenceOf(event: RuntimeEvent): number { const seq = Number(event.sequence); return Number.isFinite(seq) ? seq : 0; } function timestampOf(event: RuntimeEvent): string { const value = typeof event.createdAt === "string" ? event.createdAt : typeof event.timestamp === "string" ? event.timestamp : ""; return Number.isFinite(Date.parse(value)) ? value : ""; } function messageIdOf(event: RuntimeEvent): string | null { const payload = event.payload; if (!payload || typeof payload !== "object" || Array.isArray(payload)) return null; const message = (payload as Record).message; if (!message || typeof message !== "object" || Array.isArray(message)) return null; const id = (message as Record).id; return typeof id === "string" && id.length > 0 ? id : null; } function eventKey(event: RuntimeEvent): string { if (typeof event.id === "string" && event.id.length > 0) return `id:${event.id}`; const messageId = messageIdOf(event); if (messageId) return `message:${messageId}`; return `fallback:${sequenceOf(event)}:${timestampOf(event)}:${String(event.type ?? "")}`; } function compareEvents(a: RuntimeEvent, b: RuntimeEvent): number { const aTime = timestampOf(a); const bTime = timestampOf(b); if (aTime && bTime && aTime !== bTime) return aTime.localeCompare(bTime); const sequence = sequenceOf(a) - sequenceOf(b); if (sequence !== 0) return sequence; return eventKey(a).localeCompare(eventKey(b)); } function mergeUnique(base: RuntimeEvent[], incoming: RuntimeEvent[]): RuntimeEvent[] { const seen = new Set(base.map(eventKey)); const added = incoming.filter((event) => { const key = eventKey(event); if (seen.has(key)) return false; seen.add(key); return true; }); if (added.length === 0) return base; return [...base, ...added].sort(compareEvents); } export function useAmaSessionEvents({ taskId, sessionId, taskDone, enabled = true, }: { taskId?: string; sessionId?: string; taskDone: boolean; enabled?: boolean; }) { const [events, setEvents] = useState([]); const [phase, setPhase] = useState("loading"); useEffect(() => { if (!enabled || (!taskId && !sessionId)) { setEvents([]); setPhase("ready"); return; } let active = true; let ws: WebSocket | null = null; let reconnect: ReturnType | null = null; let backfillSeq = 0; setPhase("loading"); setEvents([]); function sendBackfill(socket: WebSocket, cursor?: number) { const requestId = `backfill-${++backfillSeq}`; socket.send( JSON.stringify({ type: "backfill", requestId, limit: 200, ...(cursor !== undefined ? { cursor } : {}), }), ); } const connect = async () => { try { const { url } = sessionId ? await api.sessions.sessionWs(sessionId) : await api.tasks.sessionWs(taskId!); if (!active) return; ws = new WebSocket(url); ws.onopen = () => { if (ws) sendBackfill(ws); }; ws.onmessage = (event) => { try { const frame = JSON.parse(typeof event.data === "string" ? event.data : ""); if (frame?.type === "backfill" && Array.isArray(frame.events)) { setEvents((prev) => mergeUnique(prev, frame.events as RuntimeEvent[])); setPhase("ready"); if (frame.hasMore && typeof frame.nextCursor === "number" && ws?.readyState === WebSocket.OPEN) { sendBackfill(ws, frame.nextCursor); } } else if (frame?.type === "event") { if (!frame.record) return; setEvents((prev) => mergeUnique(prev, [frame.record as RuntimeEvent])); setPhase("ready"); } else if (frame?.type === "error" || frame?.type === "runner_unavailable") { setPhase("error"); } } catch { // ignore malformed runtime frames } }; ws.onclose = () => { ws = null; if (active && !taskDone) reconnect = setTimeout(connect, LIVE_POLL_MS); }; } catch { if (active) { setPhase((prev) => (prev === "loading" ? "error" : prev)); reconnect = setTimeout(connect, LIVE_POLL_MS); } } }; void connect(); return () => { active = false; if (reconnect) clearTimeout(reconnect); ws?.close(); }; }, [enabled, sessionId, taskDone, taskId]); return { events, phase }; }