const HOP_BY_HOP = new Set([
  "host","connection","keep-alive","proxy-authenticate","proxy-authorization",
  "te","trailer","transfer-encoding","upgrade","content-length",
]);

export interface ForwardOpts {
  upstreamUrl: string;
  method: string;
  path: string;
  headers: Headers;
  body: RequestInit["body"];
  signal?: AbortSignal;
  onStreamError?: (err: unknown) => void;
}

export async function forwardToUpstream(opts: ForwardOpts): Promise<Response> {
  const out = new Headers();
  opts.headers.forEach((v, k) => {
    if (!HOP_BY_HOP.has(k.toLowerCase())) out.set(k, v);
  });
  const url = opts.upstreamUrl.replace(/\/$/, "") + opts.path;
  const init: RequestInit = { method: opts.method, headers: out };
  if (opts.body != null && opts.method !== "GET" && opts.method !== "HEAD") {
    init.body = opts.body ?? null;
  }
  if (opts.signal) init.signal = opts.signal;
  const res = await fetch(url, init);
  const cleanHeaders = new Headers(res.headers);
  cleanHeaders.delete("content-encoding");
  cleanHeaders.delete("content-length");
  const body = res.body ? wrapStreamWithErrorHandler(res.body, opts.onStreamError) : null;
  return new Response(body, { status: res.status, statusText: res.statusText, headers: cleanHeaders });
}

function wrapStreamWithErrorHandler(
  src: ReadableStream<Uint8Array>,
  onError?: (err: unknown) => void,
): ReadableStream<Uint8Array> {
  if (!onError) return src; // bypass: caller wraps errors downstream
  let cancelled = false;
  // eslint-disable-next-line @typescript-eslint/no-explicit-any
  let reader: any;
  return new ReadableStream<Uint8Array>({
    async start(controller) {
      reader = src.getReader();
      try {
        while (true) {
          const { done, value } = await reader.read();
          if (done) break;
          controller.enqueue(value);
        }
        if (!cancelled) controller.close();
      } catch (err) {
        if (!cancelled) {
          onError?.(err);
          try { controller.error(err); } catch { /* already closed */ }
        }
      } finally {
        try { reader.releaseLock(); } catch { /* ignore */ }
      }
    },
    async cancel(reason) {
      cancelled = true;
      try { await reader.cancel(reason); } catch { /* ignore */ }
    },
  });
}
