Menu
popagent
publicLatest change 98a9a832366ede73ff5e61d826d68af4f775db92 - Add safe autonomous deployment and complete project guidance by AkurAI Build
import type { ChannelEvent } from "../api-types";
import { apiFetch } from "./api";
export function parseChannelEventBlock(block: string): ChannelEvent | undefined {
const data = block.split("\n")
.filter((line) => line.startsWith("data:"))
.map((line) => line.slice(5).trimStart())
.join("\n");
if (!data) return undefined;
try {
const event = JSON.parse(data) as ChannelEvent;
return event?.type && event.channelId ? event : undefined;
} catch {
return undefined;
}
}
export async function consumeChannelEvents(
channelId: string,
signal: AbortSignal,
onEvent: (event: ChannelEvent) => void,
onOpen?: () => void,
): Promise<void> {
const requestController = new AbortController();
const abortRequest = () => requestController.abort();
signal.addEventListener("abort", abortRequest, { once: true });
const response = await apiFetch(`/api/channels/${encodeURIComponent(channelId)}/events`, { signal: requestController.signal });
signal.removeEventListener("abort", abortRequest);
if (!response.ok || !response.body) throw new Error(`Channel stream failed: HTTP ${response.status}`);
if (signal.aborted) { await response.body.cancel(); return; }
onOpen?.();
const reader = response.body.pipeThrough(new TextDecoderStream()).getReader();
const cancelReader = () => void reader.cancel();
signal.addEventListener("abort", cancelReader, { once: true });
let buffer = "";
try {
while (!signal.aborted) {
const { done, value } = await reader.read();
if (done) return;
buffer += value.replaceAll("\r\n", "\n");
let boundary = buffer.indexOf("\n\n");
while (boundary >= 0) {
const event = parseChannelEventBlock(buffer.slice(0, boundary));
if (event) onEvent(event);
buffer = buffer.slice(boundary + 2);
boundary = buffer.indexOf("\n\n");
}
}
} finally {
signal.removeEventListener("abort", cancelReader);
reader.releaseLock();
}
}