codekingpro/portable-devtools
114k
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 