* * @param {import("@cloudflare/workers-types").WebSocket} webSocketServer * @param {string} earlyDataHeader for ws 0rtt * @param {(info: string)=> void} log for ws 0rtt
(webSocketServer, earlyDataHeader, log)
| 208 | * @param {(info: string)=> void} log for ws 0rtt |
| 209 | */ |
| 210 | function makeReadableWebSocketStream(webSocketServer, earlyDataHeader, log) { |
| 211 | let readableStreamCancel = false; |
| 212 | const stream = new ReadableStream({ |
| 213 | start(controller) { |
| 214 | webSocketServer.addEventListener('message', (event) => { |
| 215 | if (readableStreamCancel) { |
| 216 | return; |
| 217 | } |
| 218 | const message = event.data; |
| 219 | controller.enqueue(message); |
| 220 | }); |
| 221 | |
| 222 | // The event means that the client closed the client -> server stream. |
| 223 | // However, the server -> client stream is still open until you call close() on the server side. |
| 224 | // The WebSocket protocol says that a separate close message must be sent in each direction to fully close the socket. |
| 225 | webSocketServer.addEventListener('close', () => { |
| 226 | // client send close, need close server |
| 227 | // if stream is cancel, skip controller.close |
| 228 | safeCloseWebSocket(webSocketServer); |
| 229 | if (readableStreamCancel) { |
| 230 | return; |
| 231 | } |
| 232 | controller.close(); |
| 233 | } |
| 234 | ); |
| 235 | webSocketServer.addEventListener('error', (err) => { |
| 236 | log('webSocketServer has error'); |
| 237 | controller.error(err); |
| 238 | } |
| 239 | ); |
| 240 | // for ws 0rtt |
| 241 | const { earlyData, error } = base64ToArrayBuffer(earlyDataHeader); |
| 242 | if (error) { |
| 243 | controller.error(error); |
| 244 | } else if (earlyData) { |
| 245 | controller.enqueue(earlyData); |
| 246 | } |
| 247 | }, |
| 248 | |
| 249 | pull(controller) { |
| 250 | // if ws can stop read if stream is full, we can implement backpressure |
| 251 | // https://streams.spec.whatwg.org/#example-rs-push-backpressure |
| 252 | }, |
| 253 | cancel(reason) { |
| 254 | // 1. pipe WritableStream has error, this cancel will called, so ws handle server close into here |
| 255 | // 2. if readableStream is cancel, all controller.close/enqueue need skip, |
| 256 | // 3. but from testing controller.error still work even if readableStream is cancel |
| 257 | if (readableStreamCancel) { |
| 258 | return; |
| 259 | } |
| 260 | log(`ReadableStream was canceled, due to ${reason}`) |
| 261 | readableStreamCancel = true; |
| 262 | safeCloseWebSocket(webSocketServer); |
| 263 | } |
| 264 | }); |
| 265 | |
| 266 | return stream; |
| 267 |
no outgoing calls
no test coverage detected