Skip to content
File

Blob: src/node/internal/internal_http_incoming.ts

typescript582 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4// Copyright Joyent and Node contributors. All rights reserved. MIT license.
5 
6import { EventEmitter } from 'node-internal:events';
7import { Readable } from 'node-internal:streams_readable';
8import { isIPv4, Socket } from 'node-internal:internal_net';
9import type {
10 IncomingMessage as _IncomingMessage,
11 IncomingHttpHeaders,
12} from 'node:http';
13const kHeaders = Symbol('kHeaders');
14const kHeadersDistinct = Symbol('kHeadersDistinct');
15const kHeadersCount = Symbol('kHeadersCount');
16 
17export let setIncomingMessageFetchResponse: (
18 incoming: IncomingMessage,
19 response: Response,
20 resetTimers?: (opts: { finished: boolean }) => void
21) => void;
22 
23export let setIncomingMessageSocket: (
24 incoming: IncomingMessage,
25 options: {
26 headers: Headers;
27 localPort: number;
28 }
29) => void;
30 
31export let setIncomingRequestBody: (
32 incoming: IncomingMessage,
33 body: ReadableStream | null
34) => void;
35 
36export class IncomingMessage extends Readable implements _IncomingMessage {
37 #response?: Response;
38 #reader?: ReadableStreamDefaultReader<Uint8Array>;
39 #reading = false;
40 #socket: unknown;
41 #stream: ReadableStream | null = null;
42 
43 override aborted = false;
44 url: string = '';
45 // @ts-expect-error TS2416 Type-inconsistencies
46 method: string | null = null;
47 // @ts-expect-error TS2416 Type-inconsistencies
48 statusCode: number | null = null;
49 // @ts-expect-error TS2416 Type-inconsistencies
50 statusMessage: string | null = null;
51 httpVersionMajor = 1;
52 httpVersionMinor = 1;
53 httpVersion: string = '1.1';
54 complete = false;
55 rawHeaders: string[] = [];
56 joinDuplicateHeaders = false;
57 
58 // The cloudflare property is currently only used on the server-side
59 // to access properties like `req.cf`, and the `env` and `ctx`
60 // objects.
61 cloudflare: {
62 // Technically, the type should be IncomingRequestCfProperties but
63 // we don't have that type in the workerd runtime at the moment.
64 cf?: Record<string, unknown> | undefined;
65 env?: unknown;
66 ctx?: unknown;
67 } = { cf: undefined, env: undefined, ctx: undefined };
68 
69 [kHeaders]: IncomingHttpHeaders | null = null;
70 [kHeadersDistinct]: Record<string, string[]> | null = null;
71 [kHeadersCount]: number = 0;
72 
73 // Flag for when we decide that this message cannot possibly be
74 // read by the user, so there's no point continuing to handle it.
75 _dumped = false;
76 _consuming = false;
77 _paused = false;
78 
79 static {
80 setIncomingMessageFetchResponse = (
81 incoming: IncomingMessage,
82 response: Response
83 ): void => {
84 incoming.#setFetchResponse(response);
85 };
86 
87 // This method sets the socket property of the IncomingMessage object.
88 // Please rest assured that this method implements a subset of Socket since
89 // in Node.js it's net.Socket which isn't possible to implement within our own
90 // implementation since our implementation is based on Request and Response objects.
91 setIncomingMessageSocket = (
92 incoming: IncomingMessage,
93 { headers, localPort }: { headers: Headers; localPort: number }
94 ): void => {
95 const connectingIp = headers.get('cf-connecting-ip');
96 const isConnectingIpIpv4 = connectingIp ? isIPv4(connectingIp) : true;
97 // Return a port number between 2^15 and 2^16.
98 const remotePort = (Math.random() * 0x8000) | 0x8000;
99 
100 // Some libraries such as on-finished (which Express.js depends on)
101 // Ref: https://github.com/jshttp/on-finished/blob/d2974f5a18f468ea56f58acb2f6d402f4b5142f0/index.js
102 // calls EventEmitter events on socket attribute.
103 const socket = new EventEmitter();
104 
105 Object.defineProperties(socket, {
106 encrypted: {
107 value: headers.get('x-forwarded-proto') === 'https',
108 writable: false,
109 configurable: true,
110 },
111 readable: {
112 get: () => {
113 return incoming.readable;
114 },
115 configurable: true,
116 },
117 remoteFamily: {
118 get: () => {
119 if (incoming.destroyed) {
120 return undefined;
121 }
122 return isConnectingIpIpv4 ? 'IPv4' : 'IPv6';
123 },
124 configurable: true,
125 },
126 remoteAddress: {
127 get: () => {
128 // This is defined in production, and will fallback to localhost on local development
129 // where request headers does not contain cf-connecting-ip.
130 return incoming.destroyed
131 ? undefined
132 : (connectingIp ?? '127.0.0.1');
133 },
134 configurable: true,
135 },
136 remotePort: {
137 get: () => {
138 // Return a port in the ephemeral range (32768-65535) as clients would use,
139 // and undefined if the socket is destroyed.
140 return incoming.destroyed ? undefined : remotePort;
141 },
142 configurable: true,
143 },
144 localAddress: {
145 // Host will have a value like "my-worker.yagiz.workers.dev",
146 value: headers.get('host') ?? '127.0.0.1',
147 writable: false,
148 configurable: true,
149 },
150 localPort: {
151 // This is the port defined by the `server.listen(port)` call.
152 value: localPort,
153 writable: false,
154 configurable: true,
155 },
156 destroy: {
157 value: (err: Error | undefined): IncomingMessage =>
158 incoming.destroy(err),
159 writable: false,
160 configurable: true,
161 },
162 });
163 
164 incoming.#socket = socket;
165 };
166 
167 setIncomingRequestBody = (
168 incoming: IncomingMessage,
169 stream: ReadableStream | null
170 ): void => {
171 incoming.#stream = stream;
172 };
173 }
174 
175 constructor() {
176 super({});
177 this._readableState.readingMore = true;
178 }
179 
180 #setFetchResponse(response: Response): void {
181 this[kHeaders] = {};
182 this[kHeadersDistinct] = {};
183 for (const header of response.headers.keys()) {
184 const value = response.headers.get(header) as string;
185 this[kHeaders][header] = value;
186 this[kHeadersDistinct][header] = [value];
187 this[kHeadersCount]++;
188 }
189 
190 this.#response = response;
191 this._readableState.readingMore = true;
192 
193 this.url = response.url;
194 this.statusCode = response.status;
195 this.statusMessage = response.statusText;
196 
197 this.once('end', () => {
198 // We need to emit close in a queueMicrotask because
199 // this is the only way we can ensure that the close event is emitted after destroy.
200 queueMicrotask(() => this.emit('close'));
201 });
202 
203 this.on('timeout', () => {
204 this._consuming = false;
205 });
206 
207 this.#stream = this.#response.body;
208 }
209 
210 async #tryRead(): Promise<void> {
211 if (this.#stream == null || this.#reading) return;
212 
213 this.#reading = true;
214 
215 try {
216 this.#reader ??= this.#stream.getReader();
217 
218 while (!this.destroyed) {
219 const data = await this.#reader.read();
220 if (data.done) {
221 this.complete = true;
222 this.push(null);
223 break;
224 }
225 
226 // Backpressure - stop reading until _read() is called again
227 if (!this.push(data.value)) {
228 break;
229 }
230 }
231 } catch (e) {
232 this.destroy(e as Error);
233 } finally {
234 this.#reading = false;
235 this.#reader?.releaseLock();
236 }
237 }
238 
239 // As this is an implementation of stream.Readable, we provide a _read()
240 // function that pumps the next chunk out of the underlying ReadableStream.
241 override _read(_n: number): void {
242 if (!this._consuming) {
243 this._readableState.readingMore = false;
244 this._consuming = true;
245 }
246 
247 // Difference from Node.js -
248 // The Node.js implementation will already have its internal buffer
249 // filled by the parserOnBody function.
250 // For our implementation, we use the ReadableStream instance.
251 if (this.#stream == null) {
252 // For GET and HEAD requests, the stream would be empty.
253 // Simply signal that we're done.
254 this.complete = true;
255 this.push(null);
256 return;
257 }
258 
259 this.#tryRead(); // eslint-disable-line @typescript-eslint/no-floating-promises
260 }
261 
262 #onError(error: Error | null, cb: (err?: Error | null) => void): void {
263 // This is to keep backward compatible behavior.
264 // An error is emitted only if there are listeners attached to the event.
265 if (this.listenerCount('error') === 0) {
266 cb();
267 } else {
268 cb(error);
269 }
270 }
271 
272 override _destroy(
273 error: Error | null,
274 callback: (error?: Error | null) => void
275 ): void {
276 if (!this.readableEnded || !this.complete) {
277 this.aborted = true;
278 this.emit('aborted');
279 }
280 
281 queueMicrotask(() => {
282 this.#onError(error, callback);
283 });
284 }
285 
286 // Add the given (field, value) pair to the message
287 //
288 // Per RFC2616, section 4.2 it is acceptable to join multiple instances of the
289 // same header with a ', ' if the header in question supports specification of
290 // multiple values this way. The one exception to this is the Cookie header,
291 // which has multiple values joined with a '; ' instead. If a header's values
292 // cannot be joined in either of these ways, we declare the first instance the
293 // winner and drop the second. Extended header fields (those beginning with
294 // 'x-') are always joined.
295 _addHeaderLine(
296 field: string,
297 value: string,
298 dest: IncomingHttpHeaders
299 ): void {
300 field = matchKnownFields(field);
301 const flag = field.charCodeAt(0);
302 if (flag === 0 || flag === 2) {
303 field = field.slice(1);
304 // Make a delimited list
305 if (typeof dest[field] === 'string') {
306 // eslint-disable-next-line @typescript-eslint/restrict-plus-operands
307 dest[field] += (flag === 0 ? ', ' : '; ') + value;
308 } else {
309 dest[field] = value;
310 }
311 } else if (flag === 1) {
312 // Array header -- only Set-Cookie at the moment
313 if (dest['set-cookie'] !== undefined) {
314 dest['set-cookie'].push(value);
315 } else {
316 dest['set-cookie'] = [value];
317 }
318 } else if (this.joinDuplicateHeaders) {
319 // RFC 9110 https://www.rfc-editor.org/rfc/rfc9110#section-5.2
320 // https://github.com/nodejs/node/issues/45699
321 // allow authorization multiple fields
322 // Make a delimited list
323 if (dest[field] === undefined) {
324 dest[field] = value;
325 } else {
326 // eslint-disable-next-line @typescript-eslint/restrict-plus-operands
327 dest[field] += ', ' + value;
328 }
329 } else if (dest[field] === undefined) {
330 // Drop duplicates
331 dest[field] = value;
332 }
333 }
334 
335 _addHeaderLines(headers: string[] | null, n: number): void {
336 if (Array.isArray(headers) && !this.complete) {
337 this.rawHeaders = headers;
338 this[kHeadersCount] = n;
339 
340 if (this[kHeaders]) {
341 for (let i = 0; i < n; i += 2) {
342 this._addHeaderLine(
343 headers[i] as string,
344 headers[i + 1] as string,
345 this[kHeaders]
346 );
347 }
348 }
349 }
350 }
351 
352 get headers(): Record<string, string | string[] | undefined> {
353 if (!this[kHeaders]) {
354 this[kHeaders] = {};
355 
356 const src = this.rawHeaders;
357 const dst = this[kHeaders];
358 
359 for (let n = 0; n < this[kHeadersCount]; n += 2) {
360 this._addHeaderLine(src[n] as string, src[n + 1] as string, dst);
361 }
362 }
363 return this[kHeaders];
364 }
365 
366 set headers(val: IncomingHttpHeaders) {
367 this[kHeaders] = val;
368 }
369 
370 get headersDistinct(): Record<string, string[]> {
371 if (!this[kHeadersDistinct]) {
372 this[kHeadersDistinct] = {};
373 
374 const src = this.rawHeaders;
375 const dst = this[kHeadersDistinct];
376 
377 for (let n = 0; n < this[kHeadersCount]; n += 2) {
378 this._addHeaderLineDistinct(
379 src[n] as string,
380 src[n + 1] as string,
381 dst
382 );
383 }
384 }
385 return this[kHeadersDistinct];
386 }
387 
388 set headersDistinct(val: Record<string, string[]>) {
389 this[kHeadersDistinct] = val;
390 }
391 
392 get trailers(): Record<string, string | undefined> {
393 return {};
394 }
395 
396 set trailers(_val: NodeJS.Dict<string>) {
397 // Workerd doesn't support trailers.
398 }
399 
400 get trailersDistinct(): Record<string, string[]> {
401 return {};
402 }
403 
404 set trailersDistinct(_val: Record<string, string[]>) {
405 // Workerd doesn't support trailers.
406 }
407 
408 _addHeaderLineDistinct(
409 field: string,
410 value: string,
411 dest: Record<string, string[]>
412 ): void {
413 field = field.toLowerCase();
414 if (!dest[field]) {
415 dest[field] = [value];
416 } else {
417 dest[field]?.push(value);
418 }
419 }
420 
421 // Call this instead of resume() if we want to just
422 // dump all the data to /dev/null
423 _dump(): void {
424 if (!this._dumped) {
425 this._dumped = true;
426 // If there is buffered data, it may trigger 'data' events.
427 // Remove 'data' event listeners explicitly.
428 this.removeAllListeners('data');
429 this.resume();
430 }
431 }
432 
433 setTimeout(_msecs: number, callback?: () => void): this {
434 if (callback) {
435 this.on('timeout', callback);
436 }
437 return this;
438 }
439 
440 override pipe<T extends NodeJS.WritableStream>(
441 destination: T,
442 options?: { end?: boolean }
443 ): T {
444 const shouldEnd = options?.end !== false;
445 
446 // Handle the piping manually for better control
447 this.on('data', (chunk: string | Uint8Array) => {
448 destination.write(chunk);
449 });
450 
451 this.once('end', () => {
452 if (shouldEnd) {
453 destination.end();
454 }
455 });
456 
457 this.once('error', (err: unknown) => {
458 destination.emit('error', err);
459 });
460 
461 // Always ensure reading starts - call resume to trigger the stream
462 this.resume();
463 
464 return destination;
465 }
466 
467 set connection(value: unknown) {
468 this.#socket = value;
469 }
470 
471 get connection(): Socket {
472 return this.#socket as Socket;
473 }
474 
475 get socket(): Socket {
476 return this.#socket as Socket;
477 }
478}
479 
480// This function is used to help avoid the lowercasing of a field name if it
481// matches a 'traditional cased' version of a field name. It then returns the
482// lowercased name to both avoid calling toLowerCase() a second time and to
483// indicate whether the field was a 'no duplicates' field. If a field is not a
484// 'no duplicates' field, a `0` byte is prepended as a flag. The one exception
485// to this is the Set-Cookie header which is indicated by a `1` byte flag, since
486// it is an 'array' field and thus is treated differently in _addHeaderLines().
487// TODO: perhaps http_parser could be returning both raw and lowercased versions
488// of known header names to avoid us having to call toLowerCase() for those
489// headers.
490function matchKnownFields(field: string, lowercased: boolean = false): string {
491 switch (field.length) {
492 case 3:
493 if (field === 'Age' || field === 'age') return 'age';
494 break;
495 case 4:
496 if (field === 'Host' || field === 'host') return 'host';
497 if (field === 'From' || field === 'from') return 'from';
498 if (field === 'ETag' || field === 'etag') return 'etag';
499 if (field === 'Date' || field === 'date') return '\u0000date';
500 if (field === 'Vary' || field === 'vary') return '\u0000vary';
501 break;
502 case 6:
503 if (field === 'Server' || field === 'server') return 'server';
504 if (field === 'Cookie' || field === 'cookie') return '\u0002cookie';
505 if (field === 'Origin' || field === 'origin') return '\u0000origin';
506 if (field === 'Expect' || field === 'expect') return '\u0000expect';
507 if (field === 'Accept' || field === 'accept') return '\u0000accept';
508 break;
509 case 7:
510 if (field === 'Referer' || field === 'referer') return 'referer';
511 if (field === 'Expires' || field === 'expires') return 'expires';
512 if (field === 'Upgrade' || field === 'upgrade') return '\u0000upgrade';
513 break;
514 case 8:
515 if (field === 'Location' || field === 'location') return 'location';
516 if (field === 'If-Match' || field === 'if-match') return '\u0000if-match';
517 break;
518 case 10:
519 if (field === 'User-Agent' || field === 'user-agent') return 'user-agent';
520 if (field === 'Set-Cookie' || field === 'set-cookie') return '\u0001';
521 if (field === 'Connection' || field === 'connection')
522 return '\u0000connection';
523 break;
524 case 11:
525 if (field === 'Retry-After' || field === 'retry-after')
526 return 'retry-after';
527 break;
528 case 12:
529 if (field === 'Content-Type' || field === 'content-type')
530 return 'content-type';
531 if (field === 'Max-Forwards' || field === 'max-forwards')
532 return 'max-forwards';
533 break;
534 case 13:
535 if (field === 'Authorization' || field === 'authorization')
536 return 'authorization';
537 if (field === 'Last-Modified' || field === 'last-modified')
538 return 'last-modified';
539 if (field === 'Cache-Control' || field === 'cache-control')
540 return '\u0000cache-control';
541 if (field === 'If-None-Match' || field === 'if-none-match')
542 return '\u0000if-none-match';
543 break;
544 case 14:
545 if (field === 'Content-Length' || field === 'content-length')
546 return 'content-length';
547 break;
548 case 15:
549 if (field === 'Accept-Encoding' || field === 'accept-encoding')
550 return '\u0000accept-encoding';
551 if (field === 'Accept-Language' || field === 'accept-language')
552 return '\u0000accept-language';
553 if (field === 'X-Forwarded-For' || field === 'x-forwarded-for')
554 return '\u0000x-forwarded-for';
555 break;
556 case 16:
557 if (field === 'Content-Encoding' || field === 'content-encoding')
558 return '\u0000content-encoding';
559 if (field === 'X-Forwarded-Host' || field === 'x-forwarded-host')
560 return '\u0000x-forwarded-host';
561 break;
562 case 17:
563 if (field === 'If-Modified-Since' || field === 'if-modified-since')
564 return 'if-modified-since';
565 if (field === 'Transfer-Encoding' || field === 'transfer-encoding')
566 return '\u0000transfer-encoding';
567 if (field === 'X-Forwarded-Proto' || field === 'x-forwarded-proto')
568 return '\u0000x-forwarded-proto';
569 break;
570 case 19:
571 if (field === 'Proxy-Authorization' || field === 'proxy-authorization')
572 return 'proxy-authorization';
573 if (field === 'If-Unmodified-Since' || field === 'if-unmodified-since')
574 return 'if-unmodified-since';
575 break;
576 }
577 if (lowercased) {
578 return '\u0000' + field;
579 }
580 return matchKnownFields(field.toLowerCase(), true);
581}