* Creates a readable stream from a WebSocket server, allowing for data to be read from the WebSocket. * @param {import("@cloudflare/workers-types").WebSocket} webSocketServer The WebSocket server to create the readable stream from. * @param {string} earlyDataHeader The header containing early data
(webSocketServer, earlyDataHeader, log)
| 317 | * @returns {ReadableStream} A readable stream that can be used to read data from the WebSocket. |
| 318 | */ |
| 319 | function makeReadableWebSocketStream(webSocketServer, earlyDataHeader, log) { |
| 320 | let readableStreamCancel = false; |
| 321 | const stream = new ReadableStream({ |
| 322 | start(controller) { |
| 323 | webSocketServer.addEventListener('message', (event) => { |
| 324 | const message = event.data; |
| 325 | controller.enqueue(message); |
| 326 | }); |
| 327 | |
| 328 | webSocketServer.addEventListener('close', () => { |
| 329 | safeCloseWebSocket(webSocketServer); |
| 330 | controller.close(); |
| 331 | }); |
| 332 | |
| 333 | webSocketServer.addEventListener('error', (err) => { |
| 334 | log('webSocketServer has error'); |
| 335 | controller.error(err); |
| 336 | }); |
| 337 | const { earlyData, error } = base64ToArrayBuffer(earlyDataHeader); |
| 338 | if (error) { |
| 339 | controller.error(error); |
| 340 | } else if (earlyData) { |
| 341 | controller.enqueue(earlyData); |
| 342 | } |
| 343 | }, |
| 344 | |
| 345 | pull(controller) { |
| 346 | // if ws can stop read if stream is full, we can implement backpressure |
| 347 | // https://streams.spec.whatwg.org/#example-rs-push-backpressure |
| 348 | }, |
| 349 | |
| 350 | cancel(reason) { |
| 351 | log(`ReadableStream was canceled, due to ${reason}`) |
| 352 | readableStreamCancel = true; |
| 353 | safeCloseWebSocket(webSocketServer); |
| 354 | } |
| 355 | }); |
| 356 | |
| 357 | return stream; |
| 358 | } |
| 359 | |
| 360 | // https://xtls.github.io/development/protocols/vless.html |
| 361 | // https://github.com/zizifn/excalidraw-backup/blob/main/v2ray-protocol.excalidraw |
no outgoing calls
no test coverage detected