AkurAI Build
Menu

popagent

public

Latest 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();
  }
}