// src/codec/sseStream.ts
// Shared SSE stream parser. Used by KbResponseInterceptor + ExecResponseInterceptor.
// Source of truth. Do not duplicate inline.

import { parseEvent, splitEvents, encodeEvent } from "./sseParser.ts";

export async function* parseSseStream(
  body: ReadableStream<Uint8Array>,
  dec: TextDecoder,
): AsyncGenerator<{ rawEvent: string; parsed: Record<string, unknown> | null }> {
  const reader = body.getReader();
  let carry = "";
  try {
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;

      const text = dec.decode(value, { stream: true });
      const { events, remainder } = splitEvents(carry + text);
      carry = remainder;

      for (const rawEvent of events) {
        const event = parseEvent(rawEvent);
        if (event.isDone || event.dataJson === null || typeof event.dataJson !== "object") {
          yield { rawEvent, parsed: null };
          continue;
        }
        const parsed = event.dataJson as Record<string, unknown>;
        const type = typeof parsed.type === "string" ? parsed.type : "message";
        yield { rawEvent: event.eventType ? rawEvent : encodeEvent(type, parsed), parsed };
      }
    }
    if (carry) yield { rawEvent: carry, parsed: null };
  } finally {
    reader.releaseLock();
  }
}

// Re-export for callers that want encodeEvent from the same import:
export { encodeEvent } from "./sseParser.ts";
