Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
receiver.js491 linesDownload Raw Back to websocket
1'use strict'
2
3const { Writable } = require('node:stream')
4const assert = require('node:assert')
5const { parserStates, opcodes, states, emptyBuffer, sentCloseFrameState } = require('./constants')
6const { kReadyState, kSentClose, kResponse, kReceivedClose } = require('./symbols')
7const { channels } = require('../../core/diagnostics')
8const {
9  isValidStatusCode,
10  isValidOpcode,
11  failWebsocketConnection,
12  websocketMessageReceived,
13  utf8Decode,
14  isControlFrame,
15  isTextBinaryFrame,
16  isContinuationFrame
17} = require('./util')
18const { WebsocketFrameSend } = require('./frame')
19const { closeWebSocketConnection } = require('./connection')
20const { PerMessageDeflate } = require('./permessage-deflate')
21const { MessageSizeExceededError } = require('../../core/errors')
22
23// This code was influenced by ws released under the MIT license.
24// Copyright (c) 2011 Einar Otto Stangvik <einaros@gmail.com>
25// Copyright (c) 2013 Arnout Kazemier and contributors
26// Copyright (c) 2016 Luigi Pinca and contributors
27
28class ByteParser extends Writable {
29  #buffers = []
30  #fragmentsBytes = 0
31  #byteOffset = 0
32  #loop = false
33
34  #state = parserStates.INFO
35
36  #info = {}
37  #fragments = []
38
39  /** @type {Map<string, PerMessageDeflate>} */
40  #extensions
41
42  /** @type {number} */
43  #maxPayloadSize
44
45  /**
46   * @param {import('./websocket').WebSocket} ws
47   * @param {Map<string, string>|null} extensions
48   * @param {{ maxPayloadSize?: number }} [options]
49   */
50  constructor (ws, extensions, options = {}) {
51    super()
52
53    this.ws = ws
54    this.#extensions = extensions == null ? new Map() : extensions
55    this.#maxPayloadSize = options.maxPayloadSize ?? 0
56
57    if (this.#extensions.has('permessage-deflate')) {
58      this.#extensions.set('permessage-deflate', new PerMessageDeflate(extensions, options))
59    }
60  }
61
62  /**
63   * @param {Buffer} chunk
64   * @param {() => void} callback
65   */
66  _write (chunk, _, callback) {
67    this.#buffers.push(chunk)
68    this.#byteOffset += chunk.length
69    this.#loop = true
70
71    this.run(callback)
72  }
73
74  #validatePayloadLength () {
75    if (
76      this.#maxPayloadSize > 0 &&
77      !isControlFrame(this.#info.opcode) &&
78      this.#info.payloadLength > this.#maxPayloadSize
79    ) {
80      failWebsocketConnection(this.ws, 'Payload size exceeds maximum allowed size')
81      return false
82    }
83
84    return true
85  }
86
87  /**
88   * Runs whenever a new chunk is received.
89   * Callback is called whenever there are no more chunks buffering,
90   * or not enough bytes are buffered to parse.
91   */
92  run (callback) {
93    while (this.#loop) {
94      if (this.#state === parserStates.INFO) {
95        // If there aren't enough bytes to parse the payload length, etc.
96        if (this.#byteOffset < 2) {
97          return callback()
98        }
99
100        const buffer = this.consume(2)
101        const fin = (buffer[0] & 0x80) !== 0
102        const opcode = buffer[0] & 0x0F
103        const masked = (buffer[1] & 0x80) === 0x80
104
105        const fragmented = !fin && opcode !== opcodes.CONTINUATION
106        const payloadLength = buffer[1] & 0x7F
107
108        const rsv1 = buffer[0] & 0x40
109        const rsv2 = buffer[0] & 0x20
110        const rsv3 = buffer[0] & 0x10
111
112        if (!isValidOpcode(opcode)) {
113          failWebsocketConnection(this.ws, 'Invalid opcode received')
114          return callback()
115        }
116
117        if (masked) {
118          failWebsocketConnection(this.ws, 'Frame cannot be masked')
119          return callback()
120        }
121
122        // MUST be 0 unless an extension is negotiated that defines meanings
123        // for non-zero values.  If a nonzero value is received and none of
124        // the negotiated extensions defines the meaning of such a nonzero
125        // value, the receiving endpoint MUST _Fail the WebSocket
126        // Connection_.
127        // This document allocates the RSV1 bit of the WebSocket header for
128        // PMCEs and calls the bit the "Per-Message Compressed" bit.  On a
129        // WebSocket connection where a PMCE is in use, this bit indicates
130        // whether a message is compressed or not.
131        if (rsv1 !== 0 && !this.#extensions.has('permessage-deflate')) {
132          failWebsocketConnection(this.ws, 'Expected RSV1 to be clear.')
133          return
134        }
135
136        if (rsv2 !== 0 || rsv3 !== 0) {
137          failWebsocketConnection(this.ws, 'RSV1, RSV2, RSV3 must be clear')
138          return
139        }
140
141        if (fragmented && !isTextBinaryFrame(opcode)) {
142          // Only text and binary frames can be fragmented
143          failWebsocketConnection(this.ws, 'Invalid frame type was fragmented.')
144          return
145        }
146
147        // If we are already parsing a text/binary frame and do not receive either
148        // a continuation frame or close frame, fail the connection.
149        if (isTextBinaryFrame(opcode) && this.#fragments.length > 0) {
150          failWebsocketConnection(this.ws, 'Expected continuation frame')
151          return
152        }
153
154        if (this.#info.fragmented && fragmented) {
155          // A fragmented frame can't be fragmented itself
156          failWebsocketConnection(this.ws, 'Fragmented frame exceeded 125 bytes.')
157          return
158        }
159
160        // "All control frames MUST have a payload length of 125 bytes or less
161        // and MUST NOT be fragmented."
162        if ((payloadLength > 125 || fragmented) && isControlFrame(opcode)) {
163          failWebsocketConnection(this.ws, 'Control frame either too large or fragmented')
164          return
165        }
166
167        if (isContinuationFrame(opcode) && this.#fragments.length === 0 && !this.#info.compressed) {
168          failWebsocketConnection(this.ws, 'Unexpected continuation frame')
169          return
170        }
171
172        if (payloadLength <= 125) {
173          this.#info.payloadLength = payloadLength
174          this.#state = parserStates.READ_DATA
175
176          if (!this.#validatePayloadLength()) {
177            return
178          }
179        } else if (payloadLength === 126) {
180          this.#state = parserStates.PAYLOADLENGTH_16
181        } else if (payloadLength === 127) {
182          this.#state = parserStates.PAYLOADLENGTH_64
183        }
184
185        if (isTextBinaryFrame(opcode)) {
186          this.#info.binaryType = opcode
187          this.#info.compressed = rsv1 !== 0
188        }
189
190        this.#info.opcode = opcode
191        this.#info.masked = masked
192        this.#info.fin = fin
193        this.#info.fragmented = fragmented
194      } else if (this.#state === parserStates.PAYLOADLENGTH_16) {
195        if (this.#byteOffset < 2) {
196          return callback()
197        }
198
199        const buffer = this.consume(2)
200
201        this.#info.payloadLength = buffer.readUInt16BE(0)
202        this.#state = parserStates.READ_DATA
203
204        if (!this.#validatePayloadLength()) {
205          return
206        }
207      } else if (this.#state === parserStates.PAYLOADLENGTH_64) {
208        if (this.#byteOffset < 8) {
209          return callback()
210        }
211
212        const buffer = this.consume(8)
213        const upper = buffer.readUInt32BE(0)
214        const lower = buffer.readUInt32BE(4)
215
216        // 2^31 is the maximum bytes an arraybuffer can contain
217        // on 32-bit systems. Although, on 64-bit systems, this is
218        // 2^53-1 bytes.
219        // https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Errors/Invalid_array_length
220        // https://source.chromium.org/chromium/chromium/src/+/main:v8/src/common/globals.h;drc=1946212ac0100668f14eb9e2843bdd846e510a1e;bpv=1;bpt=1;l=1275
221        // https://source.chromium.org/chromium/chromium/src/+/main:v8/src/objects/js-array-buffer.h;l=34;drc=1946212ac0100668f14eb9e2843bdd846e510a1e
222        if (upper !== 0 || lower > 2 ** 31 - 1) {
223          failWebsocketConnection(this.ws, 'Received payload length > 2^31 bytes.')
224          return
225        }
226
227        this.#info.payloadLength = lower
228        this.#state = parserStates.READ_DATA
229
230        if (!this.#validatePayloadLength()) {
231          return
232        }
233      } else if (this.#state === parserStates.READ_DATA) {
234        if (this.#byteOffset < this.#info.payloadLength) {
235          return callback()
236        }
237
238        const body = this.consume(this.#info.payloadLength)
239
240        if (isControlFrame(this.#info.opcode)) {
241          this.#loop = this.parseControlFrame(body)
242          this.#state = parserStates.INFO
243        } else {
244          if (!this.#info.compressed) {
245            this.writeFragments(body)
246
247            if (this.#maxPayloadSize > 0 && this.#fragmentsBytes > this.#maxPayloadSize) {
248              failWebsocketConnection(this.ws, new MessageSizeExceededError().message)
249              return
250            }
251
252            // If the frame is not fragmented, a message has been received.
253            // If the frame is fragmented, it will terminate with a fin bit set
254            // and an opcode of 0 (continuation), therefore we handle that when
255            // parsing continuation frames, not here.
256            if (!this.#info.fragmented && this.#info.fin) {
257              websocketMessageReceived(this.ws, this.#info.binaryType, this.consumeFragments())
258            }
259
260            this.#state = parserStates.INFO
261          } else {
262            this.#extensions.get('permessage-deflate').decompress(
263              body,
264              this.#info.fin,
265              (error, data) => {
266                if (error) {
267                  failWebsocketConnection(this.ws, error.message)
268                  return
269                }
270
271                this.writeFragments(data)
272
273                if (this.#maxPayloadSize > 0 && this.#fragmentsBytes > this.#maxPayloadSize) {
274                  failWebsocketConnection(this.ws, new MessageSizeExceededError().message)
275                  return
276                }
277
278                if (!this.#info.fin) {
279                  this.#state = parserStates.INFO
280                  this.#loop = true
281                  this.run(callback)
282                  return
283                }
284
285                websocketMessageReceived(this.ws, this.#info.binaryType, this.consumeFragments())
286
287                this.#loop = true
288                this.#state = parserStates.INFO
289                this.run(callback)
290              }
291            )
292
293            this.#loop = false
294            break
295          }
296        }
297      }
298    }
299  }
300
301  /**
302   * Take n bytes from the buffered Buffers
303   * @param {number} n
304   * @returns {Buffer}
305   */
306  consume (n) {
307    if (n > this.#byteOffset) {
308      throw new Error('Called consume() before buffers satiated.')
309    } else if (n === 0) {
310      return emptyBuffer
311    }
312
313    if (this.#buffers[0].length === n) {
314      this.#byteOffset -= this.#buffers[0].length
315      return this.#buffers.shift()
316    }
317
318    const buffer = Buffer.allocUnsafe(n)
319    let offset = 0
320
321    while (offset !== n) {
322      const next = this.#buffers[0]
323      const { length } = next
324
325      if (length + offset === n) {
326        buffer.set(this.#buffers.shift(), offset)
327        break
328      } else if (length + offset > n) {
329        buffer.set(next.subarray(0, n - offset), offset)
330        this.#buffers[0] = next.subarray(n - offset)
331        break
332      } else {
333        buffer.set(this.#buffers.shift(), offset)
334        offset += next.length
335      }
336    }
337
338    this.#byteOffset -= n
339
340    return buffer
341  }
342
343  writeFragments (fragment) {
344    this.#fragmentsBytes += fragment.length
345    this.#fragments.push(fragment)
346  }
347
348  consumeFragments () {
349    const fragments = this.#fragments
350
351    if (fragments.length === 1) {
352      this.#fragmentsBytes = 0
353      return fragments.shift()
354    }
355
356    const output = Buffer.concat(fragments, this.#fragmentsBytes)
357    this.#fragments = []
358    this.#fragmentsBytes = 0
359
360    return output
361  }
362
363  parseCloseBody (data) {
364    assert(data.length !== 1)
365
366    // https://datatracker.ietf.org/doc/html/rfc6455#section-7.1.5
367    /** @type {number|undefined} */
368    let code
369
370    if (data.length >= 2) {
371      // _The WebSocket Connection Close Code_ is
372      // defined as the status code (Section 7.4) contained in the first Close
373      // control frame received by the application
374      code = data.readUInt16BE(0)
375    }
376
377    if (code !== undefined && !isValidStatusCode(code)) {
378      return { code: 1002, reason: 'Invalid status code', error: true }
379    }
380
381    // https://datatracker.ietf.org/doc/html/rfc6455#section-7.1.6
382    /** @type {Buffer} */
383    let reason = data.subarray(2)
384
385    // Remove BOM
386    if (reason[0] === 0xEF && reason[1] === 0xBB && reason[2] === 0xBF) {
387      reason = reason.subarray(3)
388    }
389
390    try {
391      reason = utf8Decode(reason)
392    } catch {
393      return { code: 1007, reason: 'Invalid UTF-8', error: true }
394    }
395
396    return { code, reason, error: false }
397  }
398
399  /**
400   * Parses control frames.
401   * @param {Buffer} body
402   */
403  parseControlFrame (body) {
404    const { opcode, payloadLength } = this.#info
405
406    if (opcode === opcodes.CLOSE) {
407      if (payloadLength === 1) {
408        failWebsocketConnection(this.ws, 'Received close frame with a 1-byte body.')
409        return false
410      }
411
412      this.#info.closeInfo = this.parseCloseBody(body)
413
414      if (this.#info.closeInfo.error) {
415        const { code, reason } = this.#info.closeInfo
416
417        closeWebSocketConnection(this.ws, code, reason, reason.length)
418        failWebsocketConnection(this.ws, reason)
419        return false
420      }
421
422      if (this.ws[kSentClose] !== sentCloseFrameState.SENT) {
423        // If an endpoint receives a Close frame and did not previously send a
424        // Close frame, the endpoint MUST send a Close frame in response.  (When
425        // sending a Close frame in response, the endpoint typically echos the
426        // status code it received.)
427        let body = emptyBuffer
428        if (this.#info.closeInfo.code) {
429          body = Buffer.allocUnsafe(2)
430          body.writeUInt16BE(this.#info.closeInfo.code, 0)
431        }
432        const closeFrame = new WebsocketFrameSend(body)
433
434        this.ws[kResponse].socket.write(
435          closeFrame.createFrame(opcodes.CLOSE),
436          (err) => {
437            if (!err) {
438              this.ws[kSentClose] = sentCloseFrameState.SENT
439            }
440          }
441        )
442      }
443
444      // Upon either sending or receiving a Close control frame, it is said
445      // that _The WebSocket Closing Handshake is Started_ and that the
446      // WebSocket connection is in the CLOSING state.
447      this.ws[kReadyState] = states.CLOSING
448      this.ws[kReceivedClose] = true
449
450      return false
451    } else if (opcode === opcodes.PING) {
452      // Upon receipt of a Ping frame, an endpoint MUST send a Pong frame in
453      // response, unless it already received a Close frame.
454      // A Pong frame sent in response to a Ping frame must have identical
455      // "Application data"
456
457      if (!this.ws[kReceivedClose]) {
458        const frame = new WebsocketFrameSend(body)
459
460        this.ws[kResponse].socket.write(frame.createFrame(opcodes.PONG))
461
462        if (channels.ping.hasSubscribers) {
463          channels.ping.publish({
464            payload: body
465          })
466        }
467      }
468    } else if (opcode === opcodes.PONG) {
469      // A Pong frame MAY be sent unsolicited.  This serves as a
470      // unidirectional heartbeat.  A response to an unsolicited Pong frame is
471      // not expected.
472
473      if (channels.pong.hasSubscribers) {
474        channels.pong.publish({
475          payload: body
476        })
477      }
478    }
479
480    return true
481  }
482
483  get closingInfo () {
484    return this.#info.closeInfo
485  }
486}
487
488module.exports = {
489  ByteParser
490}
491 
codekingpro/portable-devtools · Team Ai