All files / src/utils sse.js

98% Statements 49/50
85.71% Branches 12/14
100% Functions 9/9
100% Lines 45/45

Press n or j to go to the next uncovered block, b, p or k for the previous block.

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 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141                                          11x     11x                               31x 31x 31x       31x 31x   31x   31x                     25x 15x 15x   1x           28x 25x 25x 25x 25x 11x 11x   1x         31x 5x 5x     5x   1x 1x       31x   31x         19x 17x 17x 15x   2x 2x 2x             17x 2x 2x   1x   2x   15x       11x       5x         11x  
/**
 * Server-Sent Events — the push transport for API-only backend access.
 *
 * Ratified by the operator on 2026-08-25 (SHY-0169). The clients hold ~48 live
 * Firestore and RTDB listeners between them, which is precisely how
 * authorization ended up being decided on the phone rather than on the server
 * (EPIC-0006). A request/response endpoint cannot replace a listener, so the
 * server holds it with the Admin SDK and fans out to authorised subscribers
 * over `text/event-stream`.
 *
 * SSE rather than WebSocket because every client→server action in this app is
 * already a mutation, and mutations are ordinary POST/PATCH routes. Nothing
 * needs a duplex channel, and one-way over plain HTTP costs no new
 * infrastructure — it is an ordinary authed route, so it inherits the same
 * server-side authorization as everything else.
 *
 * This module is the transport ONLY. It knows nothing about conversations,
 * rooms or presence, so each migration reuses it instead of inventing its own
 * framing.
 */
 
const log = require('./log');
 
/** Long enough to be cheap, short enough to beat idle proxy timeouts. */
const DEFAULT_HEARTBEAT_MS = 25_000;
 
/**
 * Begin an SSE response.
 *
 * @param {import('express').Request} req
 * @param {import('express').Response} res
 * @param {{ heartbeatMs?: number }} [options]
 * @returns {{
 *   send: (event: string, data: unknown) => boolean,
 *   onClose: (fn: () => void) => void,
 *   close: () => void,
 *   readonly isOpen: boolean,
 * }}
 */
function openStream(req, res, { heartbeatMs = DEFAULT_HEARTBEAT_MS } = {}) {
  res.setHeader('Content-Type', 'text/event-stream');
  res.setHeader('Cache-Control', 'no-cache, no-transform');
  res.setHeader('Connection', 'keep-alive');
  // The one people forget. A buffering proxy in front of Express holds every
  // event until the response ends — and a stream never ends, so the client sees
  // nothing at all and the connection merely looks idle.
  res.setHeader('X-Accel-Buffering', 'no');
  Eif (typeof res.flushHeaders === 'function') res.flushHeaders();
 
  let open = true;
  /** @type {Array<() => void>} */
  const cleanups = [];
 
  /**
   * Run every registered cleanup, even if one throws.
   *
   * The expensive half of a stream is the Firestore listener behind it. If the
   * first cleanup throws — a listener already detached, say — the rest must
   * still run, or a dropped phone leaks a live query for the life of the
   * process.
   */
  function runCleanups() {
    for (const fn of cleanups.splice(0)) {
      try {
        fn();
      } catch (err) {
        log.error('sse', 'Stream cleanup failed', { error: err.message });
      }
    }
  }
 
  function shutdown(endResponse) {
    if (!open) return;
    open = false;
    clearInterval(heartbeat);
    runCleanups();
    if (endResponse) {
      try {
        res.end();
      } catch (err) {
        log.error('sse', 'Failed to end stream', { error: err.message });
      }
    }
  }
 
  const heartbeat = setInterval(() => {
    Iif (!open) return;
    try {
      // A comment line: clients ignore it, but proxies and load balancers see
      // traffic and leave the connection alone.
      res.write(': ping\n\n');
    } catch (err) {
      log.info('sse', 'Heartbeat failed; closing stream', { error: err.message });
      shutdown(false);
    }
  }, heartbeatMs);
 
  req.on('close', () => shutdown(false));
 
  return {
    send(event, data) {
      // Never write to a closed stream. Writing to a dead socket throws
      // asynchronously, which takes the process down — and a fan-out racing a
      // disconnect is the normal case, not an edge case.
      if (!open) return false;
      try {
        res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
        return true;
      } catch (err) {
        log.info('sse', 'Write failed; closing stream', { error: err.message });
        shutdown(false);
        return false;
      }
    },
 
    onClose(fn) {
      // Registering after the stream has already gone runs immediately, so a
      // late registration cannot leak the thing it was meant to release.
      if (!open) {
        try {
          fn();
        } catch (err) {
          log.error('sse', 'Late cleanup failed', { error: err.message });
        }
        return;
      }
      cleanups.push(fn);
    },
 
    close() {
      shutdown(true);
    },
 
    get isOpen() {
      return open;
    },
  };
}
 
module.exports = { openStream, DEFAULT_HEARTBEAT_MS };