File size: 2,936 Bytes
67d18ac
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
/**
 * SSE keepalive stream generator.
 *
 * Emits pure SSE comment frames (`: keepalive\n\n`) at a fixed interval.
 * Used by the async bridge to keep the client connection alive during
 * ticket-queue wait — comment frames are universally ignored by spec-compliant
 * SSE clients (OpenAI / Anthropic SDKs) but reset their idle-timeout timers.
 *
 * Never emits `data:` frames — those carry semantic events and would confuse
 * strict SDK parsers (e.g. `@ai-sdk/anthropic` zod schemas throw on unknown
 * `type` values, killing the stream — see plan §3.6 anti-pattern #1).
 */

export interface KeepaliveOptions {
  /** Interval between comment frames in ms. */
  intervalMs: number;
  /** Comment text. Default `"keepalive"`. Must not contain newlines. */
  text?: string;
  /** External abort — when fired, the stream closes gracefully. */
  signal?: AbortSignal;
}

/**
 * Build a `ReadableStream<Uint8Array>` that emits `: ${text}\n\n` every
 * `intervalMs`. The stream closes when `signal` aborts. Text encoder is
 * reused across frames; the underlying buffer is fresh per emit.
 *
 * Cadence guarantee: first emit happens after `intervalMs` (NOT immediately)
 * so callers can compose with upstream-output streams without a leading comment.
 */
export function keepaliveStream(opts: KeepaliveOptions): ReadableStream<Uint8Array> {
  const text = (opts.text ?? "keepalive").replace(/[\r\n]/g, " ");
  const frame = new TextEncoder().encode(`: ${text}\n\n`);
  let timer: ReturnType<typeof setTimeout> | undefined;
  let aborted = false;

  return new ReadableStream<Uint8Array>({
    start(controller) {
      if (opts.signal) {
        if (opts.signal.aborted) {
          aborted = true;
          controller.close();
          return;
        }
        opts.signal.addEventListener("abort", () => {
          aborted = true;
          if (timer) {
            clearTimeout(timer);
            timer = undefined;
          }
          try { controller.close(); } catch { /* already closed */ }
        }, { once: true });
      }

      const tick = (): void => {
        if (aborted) return;
        try {
          controller.enqueue(frame);
        } catch {
          // Controller closed by consumer; stop the timer.
          if (timer) {
            clearTimeout(timer);
            timer = undefined;
          }
          return;
        }
        timer = setTimeout(tick, opts.intervalMs);
      };
      timer = setTimeout(tick, opts.intervalMs);
    },
    cancel() {
      aborted = true;
      if (timer) {
        clearTimeout(timer);
        timer = undefined;
      }
    },
  });
}

/**
 * Emit a single immediate keepalive frame as a `Uint8Array`.
 * Useful for flushing one frame before pausing on a long upstream call.
 */
export function keepaliveFrame(text: string = "keepalive"): Uint8Array {
  const clean = text.replace(/[\r\n]/g, " ");
  return new TextEncoder().encode(`: ${clean}\n\n`);
}