http-bridge.ts 3.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899
  1. /**
  2. * node:http ↔ WHATWG fetch bridge for the /api transport (host side of the
  3. * web carrier; the fetch-shaped handler itself is transport-agnostic).
  4. */
  5. import type { IncomingMessage, ServerResponse } from 'node:http'
  6. /** Default carrier cap for all HTTP RPC bodies: sized for the default
  7. * aggregate image limit (200 MiB) after base64 expansion plus envelope
  8. * headroom (~267.7 MiB required), rounded up for slack. The bridge buffers
  9. * each body in memory, so this cap is also the per-request resident bound. */
  10. export const DEFAULT_MAX_REQUEST_BODY_BYTES = 300 * 1024 * 1024
  11. /** Transport-independent request handler consumed by the Host HTTP bridge. */
  12. export interface FetchHandler {
  13. /**
  14. * Handle one standard Fetch request.
  15. * @param request - request produced by the active transport bridge.
  16. * @returns complete or streaming Fetch response.
  17. */
  18. fetch(request: Request): Promise<Response>
  19. }
  20. /**
  21. * Bridge one node:http request to the fetch-shaped handler (client close
  22. * aborts; response bodies stream out chunk by chunk).
  23. * @param req - incoming node:http request (fully read before dispatch).
  24. * @param res - node:http response the bridge writes and owns to completion.
  25. * @param apiHandler - fetch-shaped API carrier the request is dispatched to.
  26. * @param maxRequestBodyBytes - maximum body bytes buffered before dispatch.
  27. */
  28. export async function bridge(
  29. req: IncomingMessage,
  30. res: ServerResponse,
  31. apiHandler: FetchHandler,
  32. maxRequestBodyBytes = DEFAULT_MAX_REQUEST_BODY_BYTES,
  33. ): Promise<void> {
  34. const abort = new AbortController()
  35. // Client-disconnect detection MUST hang off the response, not the request:
  36. // since Node 16, IncomingMessage 'close' fires as soon as the request body is
  37. // fully consumed (immediately for a bodyless GET), which would abort a
  38. // streaming response right after open. ServerResponse 'close' fires on connection teardown;
  39. // writableEnded distinguishes a normal end() from the client going away.
  40. res.on('close', () => {
  41. if (!res.writableEnded) abort.abort()
  42. })
  43. const declaredLength = req.headers['content-length']
  44. if (declaredLength !== undefined && Number(declaredLength) > maxRequestBodyBytes) {
  45. res.writeHead(413, { connection: 'close' })
  46. res.end()
  47. req.destroy()
  48. return
  49. }
  50. const chunks: Buffer[] = []
  51. let received = 0
  52. for await (const chunk of req) {
  53. const buffer = chunk as Buffer
  54. received += buffer.byteLength
  55. if (received > maxRequestBodyBytes) {
  56. res.writeHead(413, { connection: 'close' })
  57. res.end()
  58. req.destroy()
  59. return
  60. }
  61. chunks.push(buffer)
  62. }
  63. /* v8 ignore next 3 -- `??` arms: node:http always sets url/method on server
  64. requests; the fields are only optional on the client-side IncomingMessage type */
  65. const request = new Request(new URL(req.url ?? '/', 'http://dsh.internal'), {
  66. method: req.method ?? 'GET',
  67. headers: Object.fromEntries(Object.entries(req.headers).filter(([, v]) => typeof v === 'string') as [string, string][]),
  68. ...chunks.length > 0 ? { body: Buffer.concat(chunks) } : {},
  69. signal: abort.signal,
  70. })
  71. const response = await apiHandler.fetch(request)
  72. res.writeHead(response.status, Object.fromEntries(response.headers.entries()))
  73. if (response.body === null) {
  74. res.end()
  75. return
  76. }
  77. for await (const chunk of response.body) {
  78. // Backpressure: a false return means the socket buffer is full — wait for drain
  79. // instead of buffering unboundedly (slow or suspended consumers). 'close' also
  80. // resolves so a mid-wait disconnect can't park this loop forever; the close
  81. // handler above aborts the handler stream, which then ends the iteration.
  82. if (!res.write(chunk)) {
  83. await new Promise<void>((resolve) => {
  84. const done = (): void => {
  85. res.off('drain', done)
  86. res.off('close', done)
  87. resolve()
  88. }
  89. res.once('drain', done)
  90. res.once('close', done)
  91. })
  92. }
  93. }
  94. res.end()
  95. }