File
Blob: src/workerd/api/tests/http-socket-server.js
| 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 | |
| 5 | // This file is used as a sidecar for the http-socket tests. |
| 6 | // It creates an HTTP server that will respond to requests from the convertSocketToFetcher API. |
| 7 | const http = require('node:http'); |
| 8 | const crypto = require('node:crypto'); |
| 9 | const net = require('node:net'); |
| 10 | const tls = require('node:tls'); |
| 11 | const assert = require('node:assert'); |
| 12 | |
| 13 | // Handle upgrade requests to switch from HTTP to WebSocket protocol |
| 14 | function upgradeToWebSocketConnection(req, socket, head) { |
| 15 | // Check if it's a WebSocket upgrade request |
| 16 | if (req.headers['upgrade'] !== 'websocket') { |
| 17 | socket.end('HTTP/1.1 400 Bad Request'); |
| 18 | return; |
| 19 | } |
| 20 | |
| 21 | //console.log('WebSocket upgrade request received'); |
| 22 | //console.log('WebSocket headers:', req.headers); |
| 23 | |
| 24 | // Get the WebSocket key from the client |
| 25 | const webSocketKey = req.headers['sec-websocket-key']; |
| 26 | const GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'; // WebSocket protocol GUID |
| 27 | |
| 28 | // Create the accept key by concatenating the key and GUID, then hashing |
| 29 | const acceptKey = crypto |
| 30 | .createHash('sha1') |
| 31 | .update(webSocketKey + GUID) |
| 32 | .digest('base64'); |
| 33 | |
| 34 | // Write WebSocket handshake response headers |
| 35 | socket.write( |
| 36 | 'HTTP/1.1 101 Switching Protocols\r\n' + |
| 37 | 'Upgrade: websocket\r\n' + |
| 38 | 'Connection: Upgrade\r\n' + |
| 39 | `Sec-WebSocket-Accept: ${acceptKey}\r\n` + |
| 40 | '\r\n' |
| 41 | ); |
| 42 | |
| 43 | // Socket is now a WebSocket connection |
| 44 | handleWebSocketConnection(socket); |
| 45 | } |
| 46 | |
| 47 | // Function to send a message to the client |
| 48 | function sendMessage(socket, message) { |
| 49 | const payload = Buffer.from(message); |
| 50 | const payloadLength = payload.length; |
| 51 | |
| 52 | // Create frame header |
| 53 | let header; |
| 54 | |
| 55 | if (payloadLength <= 125) { |
| 56 | header = Buffer.alloc(2); |
| 57 | header[1] = payloadLength; |
| 58 | } else if (payloadLength <= 65535) { |
| 59 | header = Buffer.alloc(4); |
| 60 | header[1] = 126; |
| 61 | header.writeUInt16BE(payloadLength, 2); |
| 62 | } else { |
| 63 | header = Buffer.alloc(10); |
| 64 | header[1] = 127; |
| 65 | // Write length as 64-bit integer (simplified here) |
| 66 | header.writeBigUInt64BE(BigInt(payloadLength), 2); |
| 67 | } |
| 68 | |
| 69 | // Set the first byte: FIN bit (0x80) + opcode 0x01 for text data |
| 70 | header[0] = 0x81; |
| 71 | |
| 72 | // Combine header and payload |
| 73 | const frame = Buffer.concat([header, payload]); |
| 74 | |
| 75 | // Send the frame |
| 76 | socket.write(frame); |
| 77 | } |
| 78 | |
| 79 | // Parse WebSocket frames to extract message data |
| 80 | function parseWebSocketFrame(buffer) { |
| 81 | if (buffer.length < 2) return null; |
| 82 | |
| 83 | const isFinalFrame = !!(buffer[0] & 0x80); |
| 84 | const opcode = buffer[0] & 0x0f; |
| 85 | const isMasked = !!(buffer[1] & 0x80); |
| 86 | let payloadLength = buffer[1] & 0x7f; |
| 87 | |
| 88 | let maskingKeyOffset = 2; |
| 89 | if (payloadLength === 126) { |
| 90 | payloadLength = buffer.readUInt16BE(2); |
| 91 | maskingKeyOffset = 4; |
| 92 | } else if (payloadLength === 127) { |
| 93 | payloadLength = Number(buffer.readBigUInt64BE(2)); |
| 94 | maskingKeyOffset = 10; |
| 95 | } |
| 96 | |
| 97 | if (!isMasked) { |
| 98 | // According to the spec, client messages must be masked |
| 99 | return { |
| 100 | opcode, |
| 101 | payload: Buffer.alloc(0), |
| 102 | isControl: (opcode & 0x8) !== 0, |
| 103 | }; |
| 104 | } |
| 105 | |
| 106 | const maskingKey = buffer.slice(maskingKeyOffset, maskingKeyOffset + 4); |
| 107 | const payloadOffset = maskingKeyOffset + 4; |
| 108 | |
| 109 | if (buffer.length < payloadOffset + payloadLength) { |
| 110 | return null; // Not enough data |
| 111 | } |
| 112 | |
| 113 | const payload = Buffer.alloc(payloadLength); |
| 114 | for (let i = 0; i < payloadLength; i++) { |
| 115 | payload[i] = buffer[payloadOffset + i] ^ maskingKey[i % 4]; |
| 116 | } |
| 117 | |
| 118 | return { |
| 119 | opcode, |
| 120 | payload, |
| 121 | isControl: (opcode & 0x8) !== 0, |
| 122 | }; |
| 123 | } |
| 124 | |
| 125 | // Create HTTP server |
| 126 | const server = http.createServer((req, res) => { |
| 127 | //console.log(`Received request: ${req.method} ${req.url}`); |
| 128 | |
| 129 | if (req.url === '/ping') { |
| 130 | res.writeHead(200, { 'Content-Type': 'text/plain' }); |
| 131 | res.end('pong'); |
| 132 | } else if (req.url === '/json') { |
| 133 | res.writeHead(200, { 'Content-Type': 'application/json' }); |
| 134 | res.end(JSON.stringify({ message: 'Hello from HTTP socket server' })); |
| 135 | } else if (req.url === '/echo' && req.method === 'POST') { |
| 136 | let body = ''; |
| 137 | req.on('data', (chunk) => { |
| 138 | body += chunk.toString(); |
| 139 | }); |
| 140 | req.on('end', () => { |
| 141 | res.writeHead(200, { 'Content-Type': 'text/plain' }); |
| 142 | res.end(body); |
| 143 | }); |
| 144 | } else if (req.url === '/headers') { |
| 145 | const headers = {}; |
| 146 | for (const [key, value] of Object.entries(req.headers)) { |
| 147 | headers[key] = value; |
| 148 | } |
| 149 | res.writeHead(200, { 'Content-Type': 'application/json' }); |
| 150 | res.end(JSON.stringify(headers)); |
| 151 | } else if (req.url === '/status/404') { |
| 152 | res.writeHead(404, { 'Content-Type': 'text/plain' }); |
| 153 | res.end('Not Found'); |
| 154 | } else if (req.url === '/status/500') { |
| 155 | res.writeHead(500, { 'Content-Type': 'text/plain' }); |
| 156 | res.end('Internal Server Error'); |
| 157 | } else if (req.url === '/drop') { |
| 158 | res.socket.drop(); |
| 159 | } else if (req.url === '/destroy') { |
| 160 | res.socket.destroy(); |
| 161 | } else if (req.url === '/redirect') { |
| 162 | res.writeHead(301, { Location: '/ping' }); |
| 163 | res.end('Moved Permanently'); |
| 164 | } else { |
| 165 | res.writeHead(404, { 'Content-Type': 'text/plain' }); |
| 166 | res.end('Not Found'); |
| 167 | } |
| 168 | }); |
| 169 | |
| 170 | // Handle WebSocket upgrade requests |
| 171 | server.on('upgrade', upgradeToWebSocketConnection); |
| 172 | |
| 173 | // Function to handle WebSocket connections |
| 174 | function handleWebSocketConnection(socket) { |
| 175 | //console.log('WebSocket connection established'); |
| 176 | |
| 177 | // Send a welcome message |
| 178 | sendMessage(socket, 'Welcome to WebSocket server'); |
| 179 | |
| 180 | // Handle incoming data |
| 181 | let buffer = Buffer.alloc(0); |
| 182 | socket.on('data', (data) => { |
| 183 | buffer = Buffer.concat([buffer, data]); |
| 184 | |
| 185 | // Process frames until we can't anymore |
| 186 | let frame; |
| 187 | while ((frame = parseWebSocketFrame(buffer))) { |
| 188 | // Handle different frame types |
| 189 | if (frame.isControl) { |
| 190 | if (frame.opcode === 0x8) { |
| 191 | // Close frame |
| 192 | socket.end(); |
| 193 | return; |
| 194 | } |
| 195 | // Skip other control frames |
| 196 | continue; |
| 197 | } |
| 198 | |
| 199 | // For text/binary frames, echo the message back |
| 200 | if (frame.opcode === 0x1) { |
| 201 | // Text frame |
| 202 | const message = frame.payload.toString('utf8'); |
| 203 | //console.log('Received WebSocket message:', message); |
| 204 | |
| 205 | // Echo the message back |
| 206 | sendMessage(socket, `Echo: ${message}`); |
| 207 | } |
| 208 | |
| 209 | // Remove processed frame from buffer |
| 210 | const frameSize = |
| 211 | frame.payload.length + (frame.payload.length > 125 ? 4 : 2) + 4; // header + masking key + payload |
| 212 | buffer = buffer.slice(frameSize); |
| 213 | } |
| 214 | }); |
| 215 | |
| 216 | // Handle socket close |
| 217 | socket.on('end', () => { |
| 218 | //console.log('WebSocket connection closed'); |
| 219 | }); |
| 220 | |
| 221 | socket.on('error', (error) => { |
| 222 | console.error('WebSocket error:', error); |
| 223 | }); |
| 224 | } |
| 225 | |
| 226 | server.listen(process.env.HTTP_SOCKET_SERVER_PORT, () => { |
| 227 | console.info( |
| 228 | `HTTP Socket test server listening on port ${server.address().port}` |
| 229 | ); |
| 230 | }); |
| 231 | |
| 232 | // This socket grabs connections and immediately drop them |
| 233 | const dropServer = net.createServer((socket) => { |
| 234 | socket.on('error', (err) => { |
| 235 | console.log('DROP: ' + err.name); |
| 236 | console.log('DROP: ' + err.message); |
| 237 | }); |
| 238 | var ready = true; |
| 239 | // Repeatedly send a page of data till the socket is borked |
| 240 | const repeatedString = 'A'.repeat(4_096); |
| 241 | while (ready) { |
| 242 | ready = socket.write(repeatedString + '\n'); |
| 243 | } |
| 244 | socket.write(repeatedString + '\n'); |
| 245 | }); |
| 246 | |
| 247 | dropServer.listen(process.env.SOCKET_PARTIALLY_WRITTEN, () => { |
| 248 | console.info( |
| 249 | `Drop Socket test server listening on port ${dropServer.address().port}` |
| 250 | ); |
| 251 | }); |
| 252 | |
| 253 | // Flush Hello Socket server that checks for hello message and responds with HTTP pong |
| 254 | const flushHelloServer = net.createServer((socket) => { |
| 255 | let receivedHello = false; |
| 256 | let buffer = Buffer.alloc(0); |
| 257 | |
| 258 | socket.on('data', (data) => { |
| 259 | buffer = Buffer.concat([buffer, data]); |
| 260 | const message = buffer.toString().trim(); |
| 261 | |
| 262 | if (!receivedHello && message.includes('Hello')) { |
| 263 | receivedHello = true; |
| 264 | // Clear the buffer after processing hello |
| 265 | buffer = Buffer.alloc(0); |
| 266 | return; |
| 267 | } |
| 268 | |
| 269 | if (receivedHello) { |
| 270 | // Respond with HTTP pong response |
| 271 | const httpResponse = |
| 272 | 'HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: 4\r\n\r\npong'; |
| 273 | socket.write(httpResponse); |
| 274 | socket.end(); |
| 275 | } |
| 276 | }); |
| 277 | |
| 278 | socket.on('error', (err) => { |
| 279 | console.log('FLUSH_HELLO error:', err.message); |
| 280 | }); |
| 281 | }); |
| 282 | |
| 283 | flushHelloServer.listen(process.env.FLUSH_HELLO_SOCKET, () => { |
| 284 | console.info( |
| 285 | `Flush Hello Socket server listening on port ${flushHelloServer.address().port}` |
| 286 | ); |
| 287 | }); |
| 288 | |
| 289 | // Create a self-signed certificate for TLS with proper SAN extension |
| 290 | function createSelfSignedCert() { |
| 291 | const key = `-----BEGIN PRIVATE KEY----- |
| 292 | MIIEvAIBADANBgkqhkiG9w0BAQEFAASCBKYwggSiAgEAAoIBAQDIWfGy2tRsqANt |
| 293 | J1F/52bIDzMDxlmSkDpu3U3Ehq6TmH2hNBcLOWuWLvG8Np9artnzk8QnodfN8yEJ |
| 294 | 0HRzZ6mRjVIUHJOb3+L1+0ePOM8dtWvG0AOd95K0T0imJRLPR18UjHl5OLE7mMS9 |
| 295 | CGHa6mDTGKzTcdxtpkjiyoNgfdKKSKzLplga5if36leGJ2+mEhOAc/cV1kqOx+hP |
| 296 | VGAOz5p3OnkgolC8hZ3WTsAFEYMU0QoPNs7jVCVNGH9t3qPiWV2/7XaNQoMnMHgq |
| 297 | yyljcvDo0O0dw0HKtBVIMhz5xHHGZqJNM/R36MeHzO+cIYhJz+ncknu62+IlQOCY |
| 298 | eyHuxvpRAgMBAAECggEAB3SXYqsze+6doAZ2SS7se3XbVWDgbOyKjB0Wm4FShkIG |
| 299 | rMTCNcP3fbF2A+W5dNesWxzM0Be87thFCrcz2iaJoBW07/QnRwXkDXjCDzGTPX0G |
| 300 | i3GqrMptbmHD55DaHBYBEwPuMkVajQfwjEM/VvThUQGqTrz+MaNeM3hLPsA34Tbl |
| 301 | wigHK+4tyAFLkYkvsxXYHs1F23ey2ubFUyBI8gvfayOvr4MOVfEbQZsmz3IB5jjk |
| 302 | oK5EaMJuhty65pb9Pi6ncSbfVQ2aciNgHjZs/is/WQfQPB2jYBPpmaFQbZc2pseZ |
| 303 | 8zTLvF+GKZng4hQE1F+F6DUf9sxb54Eu1XqzqGlrdQKBgQDkQc5fJj72rF9RkiBN |
| 304 | zeEK0Ngihm9Jzsb544i7DT/4jPWwa1+dt+kNVqguDbcxi4tCi+P8BwTnfX5M31d2 |
| 305 | /Si9nDBpP+WLLRHjFq2AQf05JvFc1/7s4KvWBrfKbwiN+uxrstB2lDuFp386B20U |
| 306 | stsCBu/nawjt5lYY7S0zPokkJQKBgQDgs9rTJNnRuaaEpUxqRdR6zR3z0Qe1B1Ga |
| 307 | y6z8OqiX3N+wM9uAcrTzGOtPKvB4glcJZrrQp0NSxkJVsFBzPBCfUfjIV/IiOqS1 |
| 308 | nE/rrEKEG7ZVkxDUsdeS8KiVKBJjLch/hrT0udgv4vndXqeaJuNzbD1yfiKbPnY5 |
| 309 | yGC78uqvvQKBgGUGMx6twMRQekeSEzYcXvP4hxCQy4SxPiOvbv7K2HtbeApDG6ik |
| 310 | k0NSDVGExIXrKxGi9J7BRIxoYJQJbZ6+YV+6VzreCuxUYExP5y6TBk5bTAw5lRym |
| 311 | O6eYhZPVHMYqPqVUGSvCY629+nNmggLdPk1hYKDeIK+aeJTDtHOvw+b5AoGAZI74 |
| 312 | wf8+34WWyMv026ZuhZpf6ipEqbYhxgWaX7KcmoHFNWSvudcbtaMUQ3Sy8ytZaiKo |
| 313 | PhJspZGGRDTIfBmIUtRrYrVA7iKSbZgLiCuqBNcmDTvoj1cbY24B8+Zf/DSUAsY1 |
| 314 | G0RERIHuUiw3E1yN86yf/yoFsLYOUKOk7teyQX0CgYBFUzUeVszODI/1NeKuGdoR |
| 315 | 1uJk2gt7hH4t7GdX4G78h3i6P5pa7qSyUXqBm53nXrYhEFRiVgDh5BoRnoKKWomj |
| 316 | IDidHZX3fMoQvdKAT8sUbT/Q3KKWPFqGv7frdpQe+JM0DwwjjPrMtUWhuve455XW |
| 317 | bCKgoxoaqEPUzY9CyI+RZg== |
| 318 | -----END PRIVATE KEY-----`; |
| 319 | |
| 320 | const cert = `-----BEGIN CERTIFICATE----- |
| 321 | MIIDmTCCAoGgAwIBAgIUFdaq1aN7zHoSsmLp6ItxnM1VOPEwDQYJKoZIhvcNAQEL |
| 322 | BQAwTjELMAkGA1UEBhMCVVMxDTALBgNVBAgMBFRlc3QxDTALBgNVBAcMBFRlc3Qx |
| 323 | DTALBgNVBAoMBFRlc3QxEjAQBgNVBAMMCWxvY2FsaG9zdDAeFw0yNTA3MjYwNjI3 |
| 324 | MDZaFw0yNjA3MjYwNjI3MDZaME4xCzAJBgNVBAYTAlVTMQ0wCwYDVQQIDARUZXN0 |
| 325 | MQ0wCwYDVQQHDARUZXN0MQ0wCwYDVQQKDARUZXN0MRIwEAYDVQQDDAlsb2NhbGhv |
| 326 | c3QwggEiMA0GCSqGSIb3DQEBAQUAA4IBDwAwggEKAoIBAQDIWfGy2tRsqANtJ1F/ |
| 327 | 52bIDzMDxlmSkDpu3U3Ehq6TmH2hNBcLOWuWLvG8Np9artnzk8QnodfN8yEJ0HRz |
| 328 | Z6mRjVIUHJOb3+L1+0ePOM8dtWvG0AOd95K0T0imJRLPR18UjHl5OLE7mMS9CGHa |
| 329 | 6mDTGKzTcdxtpkjiyoNgfdKKSKzLplga5if36leGJ2+mEhOAc/cV1kqOx+hPVGAO |
| 330 | z5p3OnkgolC8hZ3WTsAFEYMU0QoPNs7jVCVNGH9t3qPiWV2/7XaNQoMnMHgqyylj |
| 331 | cvDo0O0dw0HKtBVIMhz5xHHGZqJNM/R36MeHzO+cIYhJz+ncknu62+IlQOCYeyHu |
| 332 | xvpRAgMBAAGjbzBtMB0GA1UdDgQWBBQ6R0NLwjufqWjT75cFdzHW9bh2DDAfBgNV |
| 333 | HSMEGDAWgBQ6R0NLwjufqWjT75cFdzHW9bh2DDAPBgNVHRMBAf8EBTADAQH/MBoG |
| 334 | A1UdEQQTMBGCCWxvY2FsaG9zdIcEfwAAATANBgkqhkiG9w0BAQsFAAOCAQEAToR4 |
| 335 | CaI9HAfSSXE+6fPthp+qrwPmfx3rW0RskpjKuqemZmIK7ydU9pcYuGtc6CPen614 |
| 336 | RfFoaWtPltbNU0KV79P4zTRNYxqTKcEyhsjyAGbLA+bJtJE3hlDrfPGVyepZETXE |
| 337 | 7Ig4XPyXi4M+WmvLboAF2dHC+H1XoWp3agIN45VRnr5uPVNX19dTbr0gc3WxLEUH |
| 338 | N839sKGVB9GVaQhF/4Z8ia0bluirf+6SNAaN/veJA40ixGEkHN3gqX4ZTZWl5rji |
| 339 | 3cLdth83wmueKxiBp8ov78ubmdPiBsyIVrsb8jEsxRPAKX8gEx09S/yIIs1ZEGi9 |
| 340 | wG73FlfFCc09zjmYww== |
| 341 | -----END CERTIFICATE-----`; |
| 342 | |
| 343 | return { key, cert }; |
| 344 | } |
| 345 | |
| 346 | // STARTTLS Socket server that implements proper handshake protocol |
| 347 | const startTlsSocketServer = net.createServer((s) => { |
| 348 | console.log('STARTTLS_SOCKET: New connection received'); |
| 349 | |
| 350 | // Send initial greeting |
| 351 | s.write('HELLO\n'); |
| 352 | |
| 353 | // Wait for one response then upgrade to TLS |
| 354 | s.once('data', (data) => { |
| 355 | const response = data.toString().trim(); |
| 356 | console.log('STARTTLS_SOCKET: Received response:', response); |
| 357 | |
| 358 | if (response === 'HELLO_BACK') { |
| 359 | s.write('START_TLS\n', () => { |
| 360 | console.log('STARTTLS_SOCKET: Sent START_TLS, upgrading to TLS'); |
| 361 | |
| 362 | // Small delay to ensure START_TLS is sent |
| 363 | console.log('STARTTLS_SOCKET: Creating TLS socket'); |
| 364 | const tlsSocket = new tls.TLSSocket(s, { |
| 365 | isServer: true, |
| 366 | server: startTlsSocketServer, |
| 367 | secureContext: tls.createSecureContext(createSelfSignedCert()), |
| 368 | requestCert: false, |
| 369 | SNICallback: (hostname, callback) => { |
| 370 | console.log( |
| 371 | 'STARTTLS_SOCKET: SNI callback for hostname:', |
| 372 | hostname |
| 373 | ); |
| 374 | callback(null, null); |
| 375 | }, |
| 376 | }); |
| 377 | |
| 378 | console.log('STARTTLS_SOCKET: Setting up TLS event handlers'); |
| 379 | |
| 380 | tlsSocket.on('secure', () => { |
| 381 | console.log('STARTTLS_SOCKET: TLS handshake complete'); |
| 382 | |
| 383 | // Handle TLS data |
| 384 | tlsSocket.on('data', (data) => { |
| 385 | const message = data.toString().trim(); |
| 386 | console.log('STARTTLS_SOCKET: Received TLS message:', message); |
| 387 | |
| 388 | if (message === 'ping') { |
| 389 | console.log('STARTTLS_SOCKET: Sending pong response'); |
| 390 | tlsSocket.write('pong\n', (err) => { |
| 391 | if (err) { |
| 392 | console.log('STARTTLS_SOCKET: Error writing pong:', err); |
| 393 | } else { |
| 394 | console.log('STARTTLS_SOCKET: Pong sent successfully'); |
| 395 | } |
| 396 | }); |
| 397 | } else if (message.includes('GET') || message.includes('POST')) { |
| 398 | // Handle HTTP requests over TLS |
| 399 | console.log('STARTTLS_SOCKET: TLS HTTP request:', message); |
| 400 | const httpResponse = |
| 401 | 'HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: 4\r\n\r\npong'; |
| 402 | tlsSocket.write(httpResponse); |
| 403 | } |
| 404 | }); |
| 405 | }); |
| 406 | |
| 407 | tlsSocket.on('error', (err) => { |
| 408 | console.log('STARTTLS_SOCKET TLS error:', err.message); |
| 409 | }); |
| 410 | |
| 411 | tlsSocket.on('close', () => { |
| 412 | console.log('STARTTLS_SOCKET: TLS socket closed'); |
| 413 | }); |
| 414 | |
| 415 | console.log( |
| 416 | 'STARTTLS_SOCKET: TLS socket created, waiting for handshake' |
| 417 | ); |
| 418 | |
| 419 | // The TLS handshake should start when the client initiates it |
| 420 | }); |
| 421 | } |
| 422 | }); |
| 423 | |
| 424 | s.on('error', (err) => { |
| 425 | console.log('STARTTLS_SOCKET socket error:', err.message); |
| 426 | }); |
| 427 | }); |
| 428 | |
| 429 | startTlsSocketServer.listen(process.env.STARTTLS_SOCKET, () => { |
| 430 | console.info( |
| 431 | `STARTTLS Socket server listening on port ${startTlsSocketServer.address().port}` |
| 432 | ); |
| 433 | }); |