Files
ThothII/frontend/src/stream/useSessionStream.ts
T

148 lines
4.1 KiB
TypeScript

import { useLayoutEffect, useRef, useState } from "react";
import { BASE } from "../api/client";
import { joinBackendPath } from "../api/runtime-config";
import { useSessionStore } from "../store/sessionStore";
import type { StreamEvent } from "../api/types";
const STREAM_UPDATE_INTERVAL_MS = 100;
function isStreamingDelta(
event: StreamEvent,
): event is Extract<StreamEvent, { type: "text_delta" | "activity_delta" }> {
return event.type === "text_delta" || event.type === "activity_delta";
}
export function createStreamEventCoalescer(
applyEvent: (event: StreamEvent) => void,
intervalMs = STREAM_UPDATE_INTERVAL_MS,
) {
let pending: Extract<StreamEvent, { type: "text_delta" | "activity_delta" }> | null = null;
let timer: ReturnType<typeof setTimeout> | null = null;
let disposed = false;
const flushPending = () => {
if (!pending) return;
const event = pending;
pending = null;
applyEvent(event);
};
const schedule = () => {
timer = setTimeout(() => {
timer = null;
if (!pending || disposed) return;
flushPending();
schedule();
}, intervalMs);
};
return {
push(event: StreamEvent) {
if (disposed) return;
if (!isStreamingDelta(event)) {
flushPending();
if (timer) clearTimeout(timer);
timer = null;
applyEvent(event);
return;
}
if (!timer) {
applyEvent(event);
schedule();
return;
}
if (pending && pending.type !== event.type) flushPending();
pending = pending
? { ...event, text: pending.text + event.text }
: event;
},
dispose() {
if (disposed) return;
if (timer) clearTimeout(timer);
timer = null;
flushPending();
disposed = true;
},
};
}
export function useSessionStream(
sessionId: string | null,
generation = 0,
cursorResetEpoch = 0,
) {
const [connected, setConnected] = useState(false);
const applyEvent = useSessionStore((s) => s.applyEvent);
const cursor = useRef({
sessionId: null as string | null,
cursorResetEpoch,
lastEventId: "",
});
const activeSource = useRef<object | null>(null);
useLayoutEffect(() => {
if (
cursor.current.sessionId !== sessionId
|| cursor.current.cursorResetEpoch !== cursorResetEpoch
) {
cursor.current = { sessionId, cursorResetEpoch, lastEventId: "" };
}
if (!sessionId) {
activeSource.current = null;
setConnected(false);
return;
}
const query = cursor.current.lastEventId
? `?lastEventId=${encodeURIComponent(cursor.current.lastEventId)}`
: "";
const es = new EventSource(joinBackendPath(BASE, `/sessions/${sessionId}/events${query}`));
const identity = { source: es, sessionId, cursorResetEpoch };
const coalescer = createStreamEventCoalescer(applyEvent);
activeSource.current = identity;
es.onopen = () => {
if (activeSource.current === identity) setConnected(true);
};
es.onerror = () => {
if (activeSource.current === identity) setConnected(false);
};
const handle = (ev: MessageEvent) => {
if (activeSource.current !== identity) return;
if (ev.lastEventId) cursor.current.lastEventId = ev.lastEventId;
try {
coalescer.push(JSON.parse(ev.data) as StreamEvent);
} catch {
/* ignore malformed */
}
};
// onmessage catches unnamed events; addEventListener catches named events
// (the backend sends event: ui_request, event: text_delta, etc.)
es.onmessage = handle;
const namedEvents = [
"ui_request",
"text_delta",
"activity_delta",
"activity_event",
"usage",
"info",
"system_event",
] as const;
for (const name of namedEvents) es.addEventListener(name, handle);
return () => {
for (const name of namedEvents) es.removeEventListener(name, handle);
es.onmessage = null;
es.close();
coalescer.dispose();
if (activeSource.current === identity) {
activeSource.current = null;
setConnected(false);
}
};
}, [sessionId, generation, cursorResetEpoch, applyEvent]);
return { connected };
}