index.ts 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267
  1. /**
  2. * @deepseek-ai/dsh-host-webserver — Web route-registration plugin: a node:http
  3. * server plus the `httpServer` service (HTTP and upgrade route registries,
  4. * index transform taps, and the single fallback seat for everything no route
  5. * claims). Knows no harness concepts and serves no files; the composing
  6. * application's frontend plugin owns dist serving through the fallback seam.
  7. * Web shape only — Electron loads dist over file:// and carries fetch over an
  8. * IPC bridge. This package never prints: the URL line belongs to the shell.
  9. */
  10. import { createServer } from 'node:http'
  11. import type { IncomingMessage, ServerResponse, Server } from 'node:http'
  12. import type { AddressInfo } from 'node:net'
  13. import type { Duplex } from 'node:stream'
  14. import { Context, Service } from 'cordis'
  15. import z from 'schemastery'
  16. declare module 'cordis' {
  17. interface Context {
  18. httpServer: HttpServerService
  19. }
  20. }
  21. /** Route match kind: 'exact' matches the pathname verbatim; 'prefix' p matches p and p/<anything>. */
  22. export type WebRouteKind = 'exact' | 'prefix'
  23. /** One named route registration. */
  24. export interface WebRoute {
  25. kind: WebRouteKind
  26. /** Absolute pathname, no trailing slash. */
  27. path: string
  28. /** Owns the full response lifecycle (may hold the response open, e.g. SSE). */
  29. handler: (req: IncomingMessage, res: ServerResponse) => void | Promise<void>
  30. }
  31. /** One exact-path HTTP upgrade registration. */
  32. export interface WebUpgradeRoute {
  33. /** Absolute pathname, no trailing slash. */
  34. path: string
  35. /** Owns protocol negotiation and the upgraded socket after dispatch. */
  36. handler: (req: IncomingMessage, socket: Duplex, head: Buffer) => void | Promise<void>
  37. }
  38. /** Gateway config: the listen address. */
  39. export interface Config {
  40. /** Listen host; the two supported values are loopback and all-interfaces. */
  41. host: '127.0.0.1' | '0.0.0.0'
  42. /** Listen port; zero requests an OS-assigned port. */
  43. port: number
  44. }
  45. /**
  46. * The web-shape HTTP carrier service. Activation listens immediately (route
  47. * registration order carries no request-facing semantics: named routes are
  48. * composed to be disjoint, and the fallback seat answers anything not yet
  49. * claimed during the boot window — 404 until its owner registers). A listen
  50. * failure throws out of init — a FAILED fiber the boot's fail-loud sweep
  51. * reports.
  52. */
  53. export class HttpServerService extends Service {
  54. static Config: z<Config> = z.object({
  55. host: z.union([z.const('127.0.0.1'), z.const('0.0.0.0')]).required(),
  56. port: z.natural().max(65535).required(),
  57. })
  58. private readonly exact = new Map<string, WebRoute>()
  59. private readonly prefixes = new Map<string, WebRoute>()
  60. private readonly upgrades = new Map<string, WebUpgradeRoute>()
  61. private readonly upgradedSockets = new Set<Duplex>()
  62. private readonly indexTaps: ((html: string) => string)[] = []
  63. private fallback: WebRoute['handler'] | undefined
  64. private server!: Server
  65. private listenedPort!: number
  66. constructor(ctx: Context, private config: Config) {
  67. super(ctx, 'httpServer')
  68. }
  69. /** The listening port (the OS-assigned value when config.port is 0). */
  70. get port(): number {
  71. return this.listenedPort
  72. }
  73. /** The configured bind host (the loopback or all-interfaces literal). */
  74. get host(): Config['host'] {
  75. return this.config.host
  76. }
  77. /**
  78. * Register a named route. Duplicate (kind, path) throws — route patterns are
  79. * a composition-level contract, so a collision is a misconfiguration.
  80. * @param route - kind, path, and the owning handler.
  81. * @returns the disposer removing the route.
  82. */
  83. register(route: WebRoute): () => void {
  84. const table = route.kind === 'exact' ? this.exact : this.prefixes
  85. if (table.has(route.path)) {
  86. throw new Error(`webserver: duplicate ${route.kind} route "${route.path}"`)
  87. }
  88. table.set(route.path, route)
  89. return () => { table.delete(route.path) }
  90. }
  91. /**
  92. * Register an exact-path HTTP upgrade route. Duplicate paths throw because
  93. * one socket can have only one protocol owner.
  94. * @param route - pathname and handler owning negotiation plus socket use.
  95. * @returns the disposer removing the route.
  96. */
  97. registerUpgrade(route: WebUpgradeRoute): () => void {
  98. if (this.upgrades.has(route.path)) {
  99. throw new Error(`webserver: duplicate upgrade route "${route.path}"`)
  100. }
  101. this.upgrades.set(route.path, route)
  102. return () => { this.upgrades.delete(route.path) }
  103. }
  104. /**
  105. * Claim the fallback seat: the handler answering every request no named
  106. * route matches (the SPA dist server in the shipped Web composition). One
  107. * owner only — a second registration throws, because two fallbacks cannot
  108. * compose.
  109. * @param handler - owns the full response lifecycle of unmatched requests.
  110. * @returns the disposer releasing the seat.
  111. */
  112. registerFallback(handler: WebRoute['handler']): () => void {
  113. if (this.fallback !== undefined) {
  114. throw new Error('webserver: fallback already registered')
  115. }
  116. this.fallback = handler
  117. return () => { this.fallback = undefined }
  118. }
  119. /**
  120. * Register an index.html transform, applied by the fallback owner to every
  121. * index response ({@link applyIndexTaps}) in registration order.
  122. * @param transform - pure html-to-html function.
  123. * @returns the disposer removing the transform.
  124. */
  125. tapIndex(transform: (html: string) => string): () => void {
  126. this.indexTaps.push(transform)
  127. return () => {
  128. const at = this.indexTaps.indexOf(transform)
  129. if (at !== -1) this.indexTaps.splice(at, 1)
  130. }
  131. }
  132. /** Listen; resolves once the socket is bound (rejection = FAILED fiber). */
  133. async [Service.init](): Promise<void> {
  134. const handle = async (req: IncomingMessage, res: ServerResponse): Promise<void> => {
  135. /* v8 ignore next -- `?? '/'` arm: node:http always sets url on server
  136. requests; the field is only optional on the client-side IncomingMessage type */
  137. const rawPath = new URL(req.url ?? '/', 'http://x').pathname
  138. const route = this.match(rawPath)
  139. if (route !== undefined) {
  140. await route.handler(req, res)
  141. return
  142. }
  143. const fallback = this.fallback
  144. if (fallback === undefined) {
  145. res.writeHead(404)
  146. res.end()
  147. return
  148. }
  149. await fallback(req, res)
  150. }
  151. // Last-resort guard: handle() rejecting would otherwise be an unhandled
  152. // rejection killing the process on one malformed request (bad %-escape,
  153. // client dropping mid-body). Per-request failures log and answer 400 —
  154. // never a process exit.
  155. this.server = createServer((req, res) => {
  156. handle(req, res).catch((err: unknown) => {
  157. this.ctx.logger.warn(err instanceof Error ? err : new Error(String(err)))
  158. if (res.headersSent) {
  159. res.destroy()
  160. return
  161. }
  162. res.writeHead(400)
  163. res.end()
  164. })
  165. })
  166. this.server.on('upgrade', (req, socket, head) => {
  167. const onError = (error: Error): void => {
  168. this.ctx.logger.warn(error)
  169. socket.destroy()
  170. }
  171. socket.on('error', onError)
  172. socket.once('close', () => {
  173. socket.off('error', onError)
  174. this.upgradedSockets.delete(socket)
  175. })
  176. let route: WebUpgradeRoute | undefined
  177. try {
  178. /* v8 ignore next -- node:http always sets url on server requests. */
  179. route = this.upgrades.get(new URL(req.url ?? '/', 'http://x').pathname)
  180. } catch (error) {
  181. this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
  182. socket.destroy()
  183. return
  184. }
  185. if (route === undefined) {
  186. socket.destroy()
  187. return
  188. }
  189. this.upgradedSockets.add(socket)
  190. try {
  191. Promise.resolve(route.handler(req, socket, head)).catch((error: unknown) => {
  192. this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
  193. socket.destroy()
  194. })
  195. } catch (error) {
  196. this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
  197. socket.destroy()
  198. }
  199. })
  200. await new Promise<void>((resolve, reject) => {
  201. this.server.once('error', reject)
  202. this.server.listen(this.config.port, this.config.host, () => {
  203. this.server.off('error', reject)
  204. this.server.on('error', (err) => { this.ctx.logger.error(err) })
  205. this.listenedPort = (this.server.address() as AddressInfo).port
  206. resolve()
  207. })
  208. })
  209. // Node does not include upgraded sockets in closeAllConnections(), so the
  210. // service tracks and destroys them as part of the same ownership boundary.
  211. this.ctx.effect(() => async () => {
  212. const serverClosed = new Promise<void>((resolve) => {
  213. this.server.close(() => { resolve() })
  214. })
  215. this.server.closeAllConnections()
  216. const upgradedClosed = [...this.upgradedSockets].map(socket => new Promise<void>((resolve) => {
  217. socket.once('close', () => { resolve() })
  218. socket.destroy()
  219. }))
  220. await Promise.all([serverClosed, ...upgradedClosed])
  221. }, 'httpServer.listen')
  222. }
  223. /** Longest-prefix-wins over the prefix table after an exact-table miss. */
  224. private match(pathname: string): WebRoute | undefined {
  225. const exact = this.exact.get(pathname)
  226. if (exact !== undefined) return exact
  227. let best: WebRoute | undefined
  228. for (const [prefix, route] of this.prefixes) {
  229. if (pathname !== prefix && !pathname.startsWith(`${prefix}/`)) continue
  230. if (best === undefined || prefix.length > best.path.length) best = route
  231. }
  232. return best
  233. }
  234. /**
  235. * Run an index.html body through the registered taps in registration order
  236. * — called by the fallback owner on every index response it renders.
  237. * @param html - the raw index.html body.
  238. * @returns the transformed body.
  239. */
  240. applyIndexTaps(html: string): string {
  241. let out = html
  242. for (const transform of this.indexTaps) out = transform(out)
  243. return out
  244. }
  245. }
  246. export default HttpServerService