Skip to content
File

Blob: src/workerd/api/tests/http-socket-test.js

javascript565 lines
1// Copyright (c) 2017-2024 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5import { connect, internalNewHttpClient } from 'cloudflare:sockets';
6import { strict as assert } from 'node:assert';
7 
8// Basic connectivity and GET test
9export const oneRequest = {
10 async test(ctrl, env, ctx) {
11 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
12 const httpClient = await internalNewHttpClient(socket);
13 
14 const response = await httpClient.fetch('https://example.com/ping');
15 assert.equal(response.status, 200);
16 const text = await response.text();
17 assert.equal(text, 'pong');
18 },
19};
20 
21export const jsonResponse = {
22 async test(ctrl, env, ctx) {
23 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
24 const httpClient = await internalNewHttpClient(socket);
25 
26 const response = await httpClient.fetch('https://example.com/json');
27 assert.equal(response.status, 200);
28 assert.equal(response.headers.get('content-type'), 'application/json');
29 const data = await response.json();
30 assert.deepEqual(data, { message: 'Hello from HTTP socket server' });
31 },
32};
33 
34export const postRequest = {
35 async test(ctrl, env, ctx) {
36 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
37 const httpClient = await internalNewHttpClient(socket);
38 
39 const postData = 'Hello, world!';
40 const response = await httpClient.fetch('https://example.com/echo', {
41 method: 'POST',
42 body: postData,
43 });
44 assert.equal(response.status, 200);
45 const text = await response.text();
46 assert.equal(text, postData);
47 },
48};
49 
50export const customHeaders = {
51 async test(ctrl, env, ctx) {
52 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
53 const httpClient = await internalNewHttpClient(socket);
54 
55 const response = await httpClient.fetch('https://example.com/headers', {
56 headers: {
57 'X-Custom-Header': 'custom-value',
58 'X-Another-Header': 'another-value',
59 },
60 });
61 assert.equal(response.status, 200);
62 const headers = await response.json();
63 assert.equal(headers['x-custom-header'], 'custom-value');
64 assert.equal(headers['x-another-header'], 'another-value');
65 },
66};
67 
68export const response404 = {
69 async test(ctrl, env, ctx) {
70 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
71 const httpClient = await internalNewHttpClient(socket);
72 
73 const response = await httpClient.fetch('https://example.com/status/404');
74 assert.equal(response.status, 404);
75 const text = await response.text();
76 assert.equal(text, 'Not Found');
77 },
78};
79 
80export const response500 = {
81 async test(ctrl, env, ctx) {
82 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
83 const httpClient = await internalNewHttpClient(socket);
84 
85 const response = await httpClient.fetch('https://example.com/status/500');
86 assert.equal(response.status, 500);
87 const text = await response.text();
88 assert.equal(text, 'Internal Server Error');
89 },
90};
91 
92export const redirect301 = {
93 async test(ctrl, env, ctx) {
94 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
95 const httpClient = await internalNewHttpClient(socket);
96 
97 // TODO(cleanup) when enabling multiple fetches we can only then do redirects.
98 const response = await httpClient.fetch('https://example.com/redirect');
99 assert.equal(response.status, 200);
100 const text = await response.text();
101 assert.equal(text, 'pong');
102 },
103};
104 
105export const multipleRequests = {
106 async test(ctrl, env, ctx) {
107 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
108 const httpClient = await internalNewHttpClient(socket);
109 
110 // First request
111 const response1 = await httpClient.fetch('https://example.com/ping');
112 assert.equal(response1.status, 200);
113 const text1 = await response1.text();
114 assert.equal(text1, 'pong');
115 
116 // Second request on same connection
117 const response2 = await httpClient.fetch('https://example.com/json');
118 assert.equal(response2.status, 200);
119 const data = await response2.json();
120 assert.deepEqual(data, { message: 'Hello from HTTP socket server' });
121 },
122};
123 
124export const multipleConcurrentRequests = {
125 async test(ctrl, env, ctx) {
126 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
127 const httpClient = await internalNewHttpClient(socket);
128 
129 // TODO(cleanup) when multiple fetches are enabled make sure this message is changed
130 await assert.rejects(
131 Promise.all([
132 httpClient.fetch('https://example.com/ping', {
133 signal: AbortSignal.timeout(500),
134 }),
135 httpClient.fetch('https://example.com/json', {
136 signal: AbortSignal.timeout(500),
137 }),
138 ]),
139 {
140 name: 'Error',
141 message: /internal error/,
142 }
143 );
144 },
145};
146 
147// Test that demonstrates socket is unusable after creating HTTP client
148export const socketFetcherUnusable = {
149 async test(ctrl, env, ctx) {
150 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
151 const httpClient = await internalNewHttpClient(socket);
152 
153 // First, successfully use the HTTP client
154 const response = await httpClient.fetch('https://example.com/ping');
155 assert.equal(response.status, 200);
156 
157 // Now try to use the socket directly - this should fail because the streams are detached
158 assert.throws(
159 () => {
160 socket.writable.getWriter();
161 },
162 {
163 name: 'TypeError',
164 message: 'This WritableStream is currently locked to a writer.',
165 }
166 );
167 
168 assert.throws(
169 () => {
170 socket.readable.getReader();
171 },
172 {
173 name: 'TypeError',
174 message: 'This ReadableStream is currently locked to a reader.',
175 }
176 );
177 },
178};
179 
180// Test that closing the socket does not affect the HTTP client
181export const socketCloseThenFetch = {
182 async test(ctrl, env, ctx) {
183 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
184 const httpClient = await internalNewHttpClient(socket);
185 // Now close the socket
186 await socket.close();
187 
188 // Try to send a request with the HTTP client after closing the socket
189 // This should still work since the HTTP client should have its own stream
190 const response1 = await httpClient.fetch('https://example.com/ping');
191 assert.equal(response1.status, 200);
192 },
193};
194 
195// Test that locking a reader before creating the stream has the right error.
196export const lockReaderFail = {
197 async test(ctrl, env, ctx) {
198 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
199 const _reader = socket.readable.getReader();
200 await assert.rejects(internalNewHttpClient(socket), {
201 name: 'TypeError',
202 message: 'The ReadableStream has been locked to a reader.',
203 });
204 },
205};
206 
207// Test that locking a writer before creating the stream has the right error.
208export const lockWriterFail = {
209 async test(ctrl, env, ctx) {
210 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
211 const _writer = socket.writable.getWriter();
212 await assert.rejects(internalNewHttpClient(socket), {
213 name: 'TypeError',
214 message: 'This WritableStream is currently locked to a writer.',
215 });
216 },
217};
218 
219/* TODO(soon) this test does not test the behavior properly
220// Test that queueing up writes before creating the stream throws an error.
221export const writeThenConvertFlushes = {
222 async test(ctrl, env, ctx) {
223 const socket = connect(`localhost:${env.FLUSH_HELLO_SOCKET}`);
224 const writer = socket.writable.getWriter();
225 const encoder = new TextEncoder();
226 await writer.write(encoder.encode('Hello'));
227 writer.releaseLock();
228 const httpClient = await internalNewHttpClient(socket);
229
230 const response = await httpClient.fetch('https://example.com/ping');
231 assert.equal(response.status, 200);
232 const text = await response.text();
233 assert.equal(text, 'pong');
234 },
235};
236*/
237 
238/* TODO(soon) whenever we support inspecting the readable stream for data.
239export const errorRemoteWrites = {
240 async test(ctrl, env, ctx) {
241 const socket = connect(`localhost:${env.SOCKET_PARTIALLY_WRITTEN}`);
242 await sleep(2000);
243 const httpClient = await internalNewHttpClient(socket);
244 await assert.rejects(
245 httpClient.fetch('https://example.com/ping', {
246 signal: AbortSignal.timeout(500),
247 }),
248 {
249 name: 'Error',
250 message: "Readable stream in http protocol starts empty, remote side wrote data that was not consumed",
251 }
252 );
253 },
254};
255*/
256 
257// Test WebSocket connection over the converted Socket
258export const websockets = {
259 async test(ctrl, env, ctx) {
260 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
261 const _httpClient = await internalNewHttpClient(socket);
262 
263 const res = await fetch(
264 `http://localhost:${env.HTTP_SOCKET_SERVER_PORT}/`,
265 {
266 headers: {
267 Connection: 'Upgrade',
268 Upgrade: 'websocket',
269 },
270 }
271 );
272 
273 const ws = res.webSocket;
274 
275 if (!ws) {
276 throw new Error('WebSocket upgrade failed');
277 }
278 
279 ws.accept();
280 ws.send('Hello from test client');
281 
282 // Test the WebSocket connection
283 await new Promise((resolve, reject) => {
284 // Set a timeout for the test
285 const timeout = setTimeout(() => {
286 reject(new Error('WebSocket test timed out after 5 seconds'));
287 }, 5000);
288 
289 ws.addEventListener('message', (event) => {
290 // Verify we got the welcome message
291 if (event.data === 'Welcome to WebSocket server') {
292 // intentionally empty
293 } else if (event.data.startsWith('Echo:')) {
294 clearTimeout(timeout);
295 ws.close();
296 resolve();
297 }
298 });
299 
300 ws.addEventListener('error', (error) => {
301 clearTimeout(timeout);
302 reject(
303 new Error(`WebSocket error: ${error.message || 'Unknown error'}`)
304 );
305 });
306 
307 ws.addEventListener('close', (event) => {
308 if (event && !event.wasClean) {
309 clearTimeout(timeout);
310 reject(
311 new Error(
312 `WebSocket closed abnormally with code ${event.code || 'unknown'}`
313 )
314 );
315 }
316 });
317 });
318 },
319};
320 
321// Test remote end destroys socket.
322export const destroyRemote = {
323 async test(ctrl, env, ctx) {
324 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
325 const httpClient = await internalNewHttpClient(socket);
326 await assert.rejects(httpClient.fetch('https://example.com/destroy'), {
327 name: 'Error',
328 message: 'Network connection lost.',
329 });
330 },
331};
332 
333// Remote destroy first then await internalNewHttpClient
334// Check if connect freezes or if a socket doesn't
335 
336// Test remote end drops socket
337export const dropRemote = {
338 async test(ctrl, env, ctx) {
339 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
340 const httpClient = await internalNewHttpClient(socket);
341 await assert.rejects(
342 httpClient.fetch('https://example.com/drop', {
343 signal: AbortSignal.timeout(500),
344 }),
345 {
346 name: 'TimeoutError',
347 message: 'The operation was aborted due to timeout',
348 }
349 );
350 /* TODO(cleanup) after supporting multiple fetches this should throw a recognizable error
351 response = await httpClient.fetch('https://example.com/ping');
352 */
353 },
354};
355 
356export const connectOnFetcherThrows = {
357 async test(ctrl, env, ctx) {
358 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
359 const httpClient = await internalNewHttpClient(socket);
360 
361 assert.throws(
362 () => {
363 httpClient.connect('localhost:8080');
364 },
365 {
366 name: 'TypeError',
367 message:
368 'connect is not something that can be done on a fetcher converted from a socket',
369 }
370 );
371 },
372};
373 
374export const startTlsSuccess = {
375 async test(ctrl, env, ctx) {
376 const socket = connect(`localhost:${env.STARTTLS_SOCKET}`, {
377 secureTransport: 'starttls',
378 });
379 
380 const writer = socket.writable.getWriter();
381 const reader = socket.readable.getReader();
382 const encoder = new TextEncoder();
383 const decoder = new TextDecoder();
384 
385 try {
386 // Read HELLO greeting
387 const { value: greeting } = await reader.read();
388 const greetingText = decoder.decode(greeting).trim();
389 console.log('startTlsSuccess: Received greeting:', greetingText);
390 
391 if (greetingText === 'HELLO') {
392 console.log('startTlsSuccess: Sending HELLO_BACK');
393 await writer.write(encoder.encode('HELLO_BACK\n'));
394 
395 // Read START_TLS signal
396 const { value: signal } = await reader.read();
397 const signalText = decoder.decode(signal).trim();
398 console.log('startTlsSuccess: Received signal:', signalText);
399 
400 if (signalText === 'START_TLS') {
401 console.log('startTlsSuccess: Received START_TLS, upgrading to TLS');
402 
403 // Release the reader and writer before upgrading
404 reader.releaseLock();
405 writer.releaseLock();
406 
407 // Upgrade to TLS
408 console.log('startTlsSuccess: About to start TLS');
409 const tlsSocket = socket.startTls();
410 console.log('startTlsSuccess: Started TLS');
411 
412 await tlsSocket.opened;
413 console.log(
414 'startTlsSuccess: TLS connection established successfully'
415 );
416 
417 const httpClient = await internalNewHttpClient(tlsSocket);
418 const response = await httpClient.fetch('https://example.com/ping');
419 assert.equal(response.status, 200);
420 const text = await response.text();
421 assert.equal(text, 'pong');
422 }
423 }
424 } catch (err) {
425 console.log('startTlsSuccess: Error:', err.message);
426 throw err;
427 }
428 },
429};
430 
431export const startTlsEarlyConvert = {
432 async test(ctrl, env, ctx) {
433 const socket = connect(`localhost:${env.STARTTLS_SOCKET}`, {
434 secureTransport: 'starttls',
435 });
436 
437 const writer = socket.writable.getWriter();
438 const reader = socket.readable.getReader();
439 const encoder = new TextEncoder();
440 const decoder = new TextDecoder();
441 
442 const { value: greeting } = await reader.read();
443 const greetingText = decoder.decode(greeting).trim();
444 
445 if (greetingText !== 'HELLO') throw 'Wrong Handshake';
446 await writer.write(encoder.encode('HELLO_BACK\n'));
447 
448 // Read START_TLS signal
449 const { value: signal } = await reader.read();
450 const signalText = decoder.decode(signal).trim();
451 if (signalText !== 'START_TLS') throw 'Cannot Start TLS';
452 
453 // Release the reader and writer before upgrading
454 reader.releaseLock();
455 writer.releaseLock();
456 
457 const tlsSocket = socket.startTls();
458 await assert.rejects(internalNewHttpClient(socket), {
459 name: 'TypeError',
460 message: 'This WritableStream is currently locked to a writer.',
461 });
462 const httpClient = await internalNewHttpClient(tlsSocket);
463 await tlsSocket.opened;
464 const response = await httpClient.fetch('https://example.com/ping');
465 assert.equal(response.status, 200);
466 const text = await response.text();
467 assert.equal(text, 'pong');
468 },
469};
470 
471export const startTlsEarlySend = {
472 async test(ctrl, env, ctx) {
473 const socket = connect(`localhost:${env.STARTTLS_SOCKET}`, {
474 secureTransport: 'starttls',
475 });
476 
477 const writer = socket.writable.getWriter();
478 const reader = socket.readable.getReader();
479 const encoder = new TextEncoder();
480 const decoder = new TextDecoder();
481 
482 const { value: greeting } = await reader.read();
483 const greetingText = decoder.decode(greeting).trim();
484 
485 if (greetingText !== 'HELLO') throw 'Wrong Handshake';
486 await writer.write(encoder.encode('HELLO_BACK\n'));
487 
488 // Read START_TLS signal
489 const { value: signal } = await reader.read();
490 const signalText = decoder.decode(signal).trim();
491 if (signalText !== 'START_TLS') throw 'Cannot Start TLS';
492 
493 // Release the reader and writer before upgrading
494 reader.releaseLock();
495 writer.releaseLock();
496 
497 const tlsSocket = socket.startTls();
498 const httpClient = await internalNewHttpClient(tlsSocket);
499 const response = await httpClient.fetch('https://example.com/ping');
500 assert.equal(response.status, 200);
501 const text = await response.text();
502 assert.equal(text, 'pong');
503 },
504};
505 
506export const manualProtocolThenFetcher = {
507 async test(ctrl, env, ctx) {
508 const socket = connect(`localhost:${env.HTTP_SOCKET_SERVER_PORT}`);
509 const writer = socket.writable.getWriter();
510 const reader = socket.readable.getReader();
511 const encoder = new TextEncoder();
512 const decoder = new TextDecoder();
513 
514 // Manually construct and send HTTP request
515 const httpRequest =
516 'GET /ping HTTP/1.1\r\n' +
517 'Host: example.com\r\n' +
518 'Connection: keep-alive\r\n' +
519 '\r\n';
520 
521 await writer.write(encoder.encode(httpRequest));
522 
523 // Read the HTTP response manually
524 let responseData = '';
525 let contentLength = 0;
526 let headersParsed = false;
527 
528 while (true) {
529 const { value, done } = await reader.read();
530 if (done) break;
531 
532 const chunk = decoder.decode(value, { stream: true });
533 responseData += chunk;
534 
535 if (!headersParsed && responseData.includes('\r\n\r\n')) {
536 const [headers, body] = responseData.split('\r\n\r\n', 2);
537 const contentLengthMatch = headers.match(/content-length:\s*(\d+)/i);
538 if (contentLengthMatch) {
539 contentLength = parseInt(contentLengthMatch[1]);
540 }
541 headersParsed = true;
542 
543 // Check if we have all the content
544 if (body.length >= contentLength) {
545 break;
546 }
547 }
548 }
549 
550 // Verify the manual HTTP response
551 assert(responseData.includes('HTTP/1.1 200 OK'));
552 assert(responseData.includes('pong'));
553 
554 // Release locks and convert to HTTP client
555 reader.releaseLock();
556 writer.releaseLock();
557 
558 const httpClient = await internalNewHttpClient(socket);
559 const response = await httpClient.fetch('https://example.com/json');
560 assert.equal(response.status, 200);
561 const data = await response.json();
562 assert.deepEqual(data, { message: 'Hello from HTTP socket server' });
563 },
564};