MCPcopy Create free account
hub / github.com/Ridter/cf_workers_proxy / makeReadableWebSocketStream

Function makeReadableWebSocketStream

js/worker.js:210–268  ·  view source on GitHub ↗

* * @param {import("@cloudflare/workers-types").WebSocket} webSocketServer * @param {string} earlyDataHeader for ws 0rtt * @param {(info: string)=> void} log for ws 0rtt

(webSocketServer, earlyDataHeader, log)

Source from the content-addressed store, hash-verified

208 * @param {(info: string)=> void} log for ws 0rtt
209 */
210function 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

Callers 1

vlessOverWSHandlerFunction · 0.85

Calls

no outgoing calls

Tested by

no test coverage detected