返回源码地图

packages/host/webserver/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 15;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * @deepseek-ai/dsh-host-webserver — node:http route registration with optional
3 * gzip, index injection, and one fallback seat. It knows no harness concepts
4 * and serves no files; the composing application owns dist serving. Electron
5 * uses file:// plus IPC instead, and this package never prints the URL.
6 * Route handlers retain direct response ownership.
7 */
8
9import { createServer } from 'node:http'
10import type { IncomingMessage, ServerResponse, Server } from 'node:http'
11import type { AddressInfo } from 'node:net'
12import type { Duplex } from 'node:stream'
13import { Context, Service } from '@deepseek-ai/cordis'
14import z from '@deepseek-ai/schemastery'
15import compressionMiddleware from 'compression'
16import Negotiator from 'negotiator'
17import { renderIndexInjections, type IndexInjection } from './injections.ts'
18
19export { renderIndexInjections } from './injections.ts'
20export type { IndexInjection, IndexInjectionPlacement } from './injections.ts'
21
22declare module '@deepseek-ai/cordis' {
23 interface Context {
24 webServer: WebServer
25 }
26 interface Events {
27 /**
28 * Collect the structured index injection table. Emitted on every index
29 * render and every worker boot-payload request; listeners push their
30 * current rows, so a row's data is read fresh at emit time.
31 * @param table - Mutable row table; listeners append in activation order.
32 * @mode emit
33 */
34 'webserver/index-inject'(table: IndexInjection[]): void
35 }
36}
37
38/** Route match kind: 'exact' matches the pathname verbatim; 'prefix' p matches p and p/<anything>. */
39export type WebRouteKind = 'exact' | 'prefix'
40
41/** One named route registration. */
42export interface WebRoute {
43 kind: WebRouteKind
44 /** Absolute pathname, no trailing slash. */
45 path: string
46 /** Owns the full response lifecycle (may hold the response open, e.g. SSE). */
47 handler: (req: IncomingMessage, res: ServerResponse) => void | Promise<void>
48}
49
50/** One exact-path HTTP upgrade registration. */
51export interface WebUpgradeRoute {
52 /** Absolute pathname, no trailing slash. */
53 path: string
54 /** Owns protocol negotiation and the upgraded socket after dispatch. */
55 handler: (req: IncomingMessage, socket: Duplex, head: Buffer) => void | Promise<void>
56}
57
58/** Web server listen and response-compression config. */
59export interface Config {
60 /** Listen host; the two supported values are loopback and all-interfaces. */
61 host: '127.0.0.1' | '0.0.0.0'
62 /** Listen port; zero requests an OS-assigned port. */
63 port: number
64 /** Response compression for socket-backed HTTP requests. @default 'none' */
65 compression?: 'none' | 'gzip'
66 /** Gzip DEFLATE level from 0 through 9. @default 1 */
67 compressionLevel?: number
68 /** Minimum known response length eligible for gzip; unknown-length streams are eligible. @default 1024 */
69 compressionThresholdBytes?: number
70}
71
72const DEFAULT_COMPRESSION = 'none' as const
73const DEFAULT_COMPRESSION_LEVEL = 1
74const DEFAULT_COMPRESSION_THRESHOLD_BYTES = 1024
75
76interface ResolvedConfig extends Config {
77 compression: 'none' | 'gzip'
78 compressionLevel: number
79 compressionThresholdBytes: number
80}
81
82type NodeMiddleware = (
83 req: IncomingMessage,
84 res: ServerResponse,
85 next: () => void,
86) => void
87
88function createGzipMiddleware(config: ResolvedConfig): NodeMiddleware {
89 // `compression` is typed for Express, but its runtime uses only the
90 // node:http request and response members supplied here.
91 const middleware = compressionMiddleware({
92 level: config.compressionLevel,
93 threshold: config.compressionThresholdBytes,
94 filter(request, response) {
95 if (response.getHeader('content-range') !== undefined) return false
96 const contentType = response.getHeader('content-type')
97 if (typeof contentType === 'string' && contentType.toLowerCase().startsWith('text/event-stream')) return false
98 if (typeof contentType === 'string' && /^multipart\/form-data(?:;|$)/i.test(contentType)) return true
99 return compressionMiddleware.filter(request, response)
100 },
101 }) as NodeMiddleware
102
103 return (req, res, next) => {
104 // The Web Worker tunnel has no socket and transfers identity bytes.
105 if ((res as { socket?: unknown }).socket === undefined) {
106 next()
107 return
108 }
109 const encoding = new Negotiator(req).encoding(['gzip', 'identity'])
110 const gzipRequest = Object.create(req) as IncomingMessage
111 Object.defineProperty(gzipRequest, 'headers', {
112 value: { ...req.headers, 'accept-encoding': encoding === 'gzip' ? 'gzip' : 'identity' },
113 })
114 middleware(gzipRequest, res, next)
115 }
116}
117
118/**
119 * The browser HTTP carrier service. Activation listens immediately. Route
120 * registration order does not affect requests because configured named routes
121 * must be distinct, and the fallback handler answers anything not yet claimed
122 * during startup with 404 until its owner registers. A listen failure rejects
123 * initialization, and the boot process reports the failed fiber.
124 */
125export class WebServer extends Service {
126 static Config: z<Config> = z.object({
127 host: z.union([z.const('127.0.0.1'), z.const('0.0.0.0')]).required(),
128 port: z.natural().max(65535).required(),
129 compression: z.union([z.const('none'), z.const('gzip')]).default(DEFAULT_COMPRESSION),
130 compressionLevel: z.number().step(1).min(0).max(9).default(DEFAULT_COMPRESSION_LEVEL),
131 compressionThresholdBytes: z.natural().default(DEFAULT_COMPRESSION_THRESHOLD_BYTES),
132 })
133
134 private readonly exact = new Map<string, WebRoute>()
135 private readonly prefixes = new Map<string, WebRoute>()
136 private readonly upgrades = new Map<string, WebUpgradeRoute>()
137 private readonly upgradedSockets = new Set<Duplex>()
138 private readonly indexTaps: ((html: string) => string)[] = []
139 private fallback: WebRoute['handler'] | undefined
140 private server!: Server
141 private listenedPort!: number
142 private readonly gzip: NodeMiddleware | undefined
143
144 constructor(ctx: Context, private config: Config) {
145 super(ctx, 'webServer')
146 const resolved = config as ResolvedConfig
147 this.gzip = resolved.compression === 'gzip' ? createGzipMiddleware(resolved) : undefined
148 }
149
150 /** The listening port (the OS-assigned value when config.port is 0). */
151 get port(): number {
152 return this.listenedPort
153 }
154
155 /** The configured bind host (the loopback or all-interfaces literal). */
156 get host(): Config['host'] {
157 return this.config.host
158 }
159
160 /**
161 * Register a named route. Duplicate (kind, path) throws — route patterns are
162 * a composition-level contract, so a collision is a misconfiguration.
163 * @param route - kind, path, and the owning handler.
164 * @returns the disposer removing the route.
165 */
166 register(route: WebRoute): () => void {
167 const table = route.kind === 'exact' ? this.exact : this.prefixes
168 if (table.has(route.path)) {
169 throw new Error(`webserver: duplicate ${route.kind} route "${route.path}"`)
170 }
171 table.set(route.path, route)
172 return () => { table.delete(route.path) }
173 }
174
175 /**
176 * Register an exact-path HTTP upgrade route. Duplicate paths throw because
177 * one socket can have only one protocol owner.
178 * @param route - pathname and handler owning negotiation plus socket use.
179 * @returns the disposer removing the route.
180 */
181 registerUpgrade(route: WebUpgradeRoute): () => void {
182 if (this.upgrades.has(route.path)) {
183 throw new Error(`webserver: duplicate upgrade route "${route.path}"`)
184 }
185 this.upgrades.set(route.path, route)
186 return () => { this.upgrades.delete(route.path) }
187 }
188
189 /**
190 * Claim the fallback seat: the handler answering every request no named
191 * route matches (the SPA dist server in the shipped Web composition). One
192 * owner only — a second registration throws, because two fallbacks cannot
193 * compose.
194 * @param handler - owns the full response lifecycle of unmatched requests.
195 * @returns the disposer releasing the seat.
196 */
197 registerFallback(handler: WebRoute['handler']): () => void {
198 if (this.fallback !== undefined) {
199 throw new Error('webserver: fallback already registered')
200 }
201 this.fallback = handler
202 return () => { this.fallback = undefined }
203 }
204
205 /**
206 * Register a raw-HTML index transform, the escape hatch for markup no
207 * {@link IndexInjection} row expresses: {@link renderIndex} applies taps in
208 * registration order after rendering the structured rows.
209 * @param transform - pure html-to-html function.
210 * @returns the disposer removing the transform.
211 */
212 tapIndex(transform: (html: string) => string): () => void {
213 this.indexTaps.push(transform)
214 return () => {
215 const at = this.indexTaps.indexOf(transform)
216 if (at !== -1) this.indexTaps.splice(at, 1)
217 }
218 }
219
220 /** Listen; resolves once the socket is bound (rejection = FAILED fiber). */
221 async [Service.init](): Promise<void> {
222 const handle = async (req: IncomingMessage, res: ServerResponse): Promise<void> => {
223 /* v8 ignore next -- `?? '/'` arm: node:http always sets url on server
224 requests; the field is only optional on the client-side IncomingMessage type */
225 const rawPath = new URL(req.url ?? '/', 'http://x').pathname
226 const route = this.match(rawPath)
227 if (route !== undefined) {
228 await route.handler(req, res)
229 return
230 }
231 const fallback = this.fallback
232 if (fallback === undefined) {
233 res.writeHead(404)
234 res.end()
235 return
236 }
237 await fallback(req, res)
238 }
239 // Last-resort guard: handle() rejecting would otherwise be an unhandled
240 // rejection killing the process on one malformed request (bad %-escape,
241 // client dropping mid-body). Per-request failures log and answer 400 —
242 // never a process exit.
243 this.server = createServer((req, res) => {
244 const next = (): void => {
245 void handle(req, res).catch((err: unknown) => {
246 this.ctx.logger.warn(err instanceof Error ? err : new Error(String(err)))
247 if (res.headersSent) {
248 res.destroy()
249 return
250 }
251 res.writeHead(400)
252 res.end()
253 })
254 }
255 if (this.gzip === undefined) next()
256 else this.gzip(req, res, next)
257 })
258 this.server.on('upgrade', (req, socket, head) => {
259 const onError = (error: Error): void => {
260 this.ctx.logger.warn(error)
261 socket.destroy()
262 }
263 socket.on('error', onError)
264 socket.once('close', () => {
265 socket.off('error', onError)
266 this.upgradedSockets.delete(socket)
267 })
268 let route: WebUpgradeRoute | undefined
269 try {
270 /* v8 ignore next -- node:http always sets url on server requests. */
271 route = this.upgrades.get(new URL(req.url ?? '/', 'http://x').pathname)
272 } catch (error) {
273 this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
274 socket.destroy()
275 return
276 }
277 if (route === undefined) {
278 socket.destroy()
279 return
280 }
281 this.upgradedSockets.add(socket)
282 try {
283 Promise.resolve(route.handler(req, socket, head)).catch((error: unknown) => {
284 this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
285 socket.destroy()
286 })
287 } catch (error) {
288 this.ctx.logger.warn(error instanceof Error ? error : new Error(String(error)))
289 socket.destroy()
290 }
291 })
292
293 await new Promise<void>((resolve, reject) => {
294 this.server.once('error', reject)
295 this.server.listen(this.config.port, this.config.host, () => {
296 this.server.off('error', reject)
297 this.server.on('error', (err) => { this.ctx.logger.error(err) })
298 this.listenedPort = (this.server.address() as AddressInfo).port
299 resolve()
300 })
301 })
302
303 // Node does not include upgraded sockets in closeAllConnections(). The service
304 // owns them with the other connections, so it tracks and destroys them explicitly.
305 this.ctx.effect(() => async () => {
306 const serverClosed = new Promise<void>((resolve) => {
307 this.server.close(() => { resolve() })
308 })
309 this.server.closeAllConnections()
310 const upgradedClosed = [...this.upgradedSockets].map(socket => new Promise<void>((resolve) => {
311 socket.once('close', () => { resolve() })
312 socket.destroy()
313 }))
314 await Promise.all([serverClosed, ...upgradedClosed])
315 }, 'webServer.listen')
316 }
317
318 /** Longest-prefix-wins over the prefix table after an exact-table miss. */
319 private match(pathname: string): WebRoute | undefined {
320 const exact = this.exact.get(pathname)
321 if (exact !== undefined) return exact
322 let best: WebRoute | undefined
323 for (const [prefix, route] of this.prefixes) {
324 if (pathname !== prefix && !pathname.startsWith(`${prefix}/`)) continue
325 if (best === undefined || prefix.length > best.path.length) best = route
326 }
327 return best
328 }
329
330 /**
331 * Run an index.html body through the registered taps in registration order
332 * — called by the fallback owner on every index response it renders.
333 * @param html - the raw index.html body.
334 * @returns the transformed body.
335 */
336 applyIndexTaps(html: string): string {
337 let out = html
338 for (const transform of this.indexTaps) out = transform(out)
339 return out
340 }
341
342 /**
343 * Gather the structured injection table: one `webserver/index-inject` emit,
344 * every subscriber pushes its current rows. Fresh per call, so subscribers
345 * read live state (module graph, theme preference) at emit time.
346 * @returns rows in subscriber activation order.
347 */
348 collectIndexInjections(): IndexInjection[] {
349 const table: IndexInjection[] = []
350 this.ctx.emit('webserver/index-inject', table)
351 return table
352 }
353
354 /**
355 * Render one index.html body: the structured injection table first, then
356 * the raw `tapIndex` transforms over the result.
357 * @param html - the raw index.html body.
358 * @returns the transformed body.
359 */
360 renderIndex(html: string): string {
361 return this.applyIndexTaps(renderIndexInjections(html, this.collectIndexInjections()))
362 }
363}
364
365export default WebServer