From 140a2b770f10df5703c935312c54455f4a9049cc Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Fri, 1 May 2026 16:15:23 +0530 Subject: [PATCH 01/10] ongoing changes --- index.js | 110 ++++++++++++++++++++++++++----------------------------- 1 file changed, 51 insertions(+), 59 deletions(-) diff --git a/index.js b/index.js index 4d3a703..58c232a 100644 --- a/index.js +++ b/index.js @@ -3,11 +3,16 @@ const net = require('net') const EventEmitter = require('events') function connPiper(connection, _dst, opts = {}, stats = {}) { - const logger = opts.logger || { log: () => {} } - logger.log({ type: 1, msg: 'Starting TCP connection piper' }) + const logger = opts.logger || { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {} + } + logger.info('Starting TCP connection piper') const loc = _dst() if (loc === null) { - logger.log({ type: 2, msg: 'Connection rejected (null destination)' }) + logger.warn('Connection rejected (null destination)') connection.destroy() // don't return rejection error if (!stats.rejectCnt) { stats.rejectCnt = 0 @@ -24,39 +29,34 @@ function connPiper(connection, _dst, opts = {}, stats = {}) { stats.locCnt++ stats.remCnt++ let destroyed = false + + // loc is the remote stream experimentation here loc.on('data', (d) => { - logger.log({ type: 0, msg: `Data from local to remote: ${d.length} bytes` }) + logger.debug(`Data from local to remote: ${d.length} bytes`) connection.write(d) }) connection.on('data', (d) => { - logger.log({ type: 0, msg: `Data from remote to local: ${d.length} bytes` }) + logger.debug(`Data from remote to local: ${d.length} bytes`) loc.write(d) }) loc.on('error', destroy).on('close', destroy) connection.on('error', destroy).on('close', destroy) loc.on('end', () => { - logger.log({ type: 0, msg: 'Local end, ending connection' }) + logger.debug('Local end, ending connection') connection.end() }) connection.on('end', () => { - logger.log({ type: 0, msg: 'Connection end, ending local' }) + logger.debug('Connection end, ending local') loc.end() }) loc.on('connect', (err) => { - if (err) { - logger.log({ type: 3, msg: err.message }) - } else { - logger.log({ type: 1, msg: 'Connected' }) - } + if (err) logger.error(err.message) + else logger.info('Connected') }) function destroy(err) { - if (destroyed) { - return - } - logger.log({ - type: 1, - msg: `Destroying connection piper${err ? ` due to error: ${err.message}` : ''}` - }) + if (destroyed) return + if (err) logger.info(`Destroying connection piper due to error: ${err.message}`) + else logger.info('Destroying connection piper') stats.locCnt-- stats.remCnt-- destroyed = true @@ -80,10 +80,8 @@ class UdpSocket { this.event = new EventEmitter() this.rinfo = null this.connect() - this.logger.log({ - type: 1, - msg: `UDP socket created${opts.bind ? `, binding to ${opts.host}:${opts.port}` : ''}` - }) + if (opts.bind) this.logger.info(`UDP socket created, binding to ${opts.host}:${opts.port}`) + else this.logger.info('UDP socket created') } connect() { @@ -91,28 +89,26 @@ class UdpSocket { this.server.bind(this.opts.port, this.opts.host) } this.server.on('message', (msg, rinfo) => { - this.logger.log({ - type: 0, - msg: `UDP message from server: ${msg.length} bytes from ${rinfo.address}:${rinfo.port}` - }) + this.logger.debug( + `UDP message from server: ${msg.length} bytes from ${rinfo.address}:${rinfo.port}` + ) this.event.emit('message', msg, rinfo) this.rinfo = rinfo }) this.client.on('message', (response, rinfo) => { - this.logger.log({ - type: 0, - msg: `UDP message from client: ${response.length} bytes from ${rinfo.address}:${rinfo.port}` - }) + this.logger.debug( + `UDP message from client: ${response.length} bytes from ${rinfo.address}:${rinfo.port}` + ) this.event.emit('message', response) }) this.client.on('error', (err) => { - this.logger.log({ type: 3, msg: `UDP error: ${err.stack}` }) + this.logger.error(`UDP error: ${err.stack}`) this.client.close() }) } write(msg) { - this.logger.log({ type: 0, msg: `Writing UDP message: ${msg.length} bytes` }) + this.logger.debug(`Writing UDP message: ${msg.length} bytes`) if (this.rinfo) { this.server.send(msg, 0, msg.length, this.rinfo.port, this.rinfo.address) } else { @@ -132,10 +128,8 @@ class UdpConnPiper { this.destroyed = false this._bindListeners() this.connect() - this.logger.log({ - type: 1, - msg: `Starting UDP connection piper${opts.client ? ' (client mode)' : ' (server mode)'}` - }) + if (opts.client) this.logger.info('Starting UDP connection piper (client mode)') + else this.logger.info('Starting UDP connection piper (server mode)') } _bindListeners() { @@ -151,12 +145,12 @@ class UdpConnPiper { connect() { if (this.destroyed) return - this.logger.log({ type: 0, msg: 'Connecting UDP piper' }) + this.logger.debug('Connecting UDP piper') this.removeListeners() this.localStream = typeof this.local === 'function' ? this.local() : this.local this.remoteStream = typeof this.remote === 'function' ? this.remote() : this.remote if (!this.localStream || !this.remoteStream) { - this.logger.log({ type: 2, msg: 'UDP connect failed (missing streams)' }) + this.logger.warn('UDP connect failed (missing streams)') this.destroy() return } @@ -170,7 +164,7 @@ class UdpConnPiper { this.remoteStream.on('message', this.bound.onConnectionMessage) this.remoteStream.on('error', this.bound.onConnectionError) this.remoteStream.on('close', this.bound.onConnectionClose) - this.logger.log({ type: 0, msg: 'UDP listeners attached' }) + this.logger.debug('UDP listeners attached') } removeListeners() { @@ -184,35 +178,33 @@ class UdpConnPiper { this.remoteStream.off('error', this.bound.onConnectionError) this.remoteStream.off('close', this.bound.onConnectionClose) } - this.logger.log({ type: 0, msg: 'UDP listeners removed' }) + this.logger.debug('UDP listeners removed') } onLocMessage(msg, rinfo) { - this.logger.log({ type: 0, msg: `UDP message from local: ${msg.length} bytes` }) + this.logger.debug(`UDP message from local: ${msg.length} bytes`) if (this.remoteStream && !this.destroyed) { this.remoteStream.trySend?.(msg) } } onConnectionMessage(msg) { - this.logger.log({ type: 0, msg: `UDP message from connection: ${msg.length} bytes` }) + this.logger.debug(`UDP message from connection: ${msg.length} bytes`) if (this.localStream && !this.destroyed) { this.localStream.write?.(msg) } } _handleError(err) { - this.logger.log({ type: 2, msg: `UDP error: ${err ? err.message : 'close'}` }) + this.logger.warn(`UDP error: ${err ? err.message : 'close'}`) this.destroy(err) } destroy(err) { if (this.destroyed) return this.destroyed = true - this.logger.log({ - type: 1, - msg: `Destroying UDP piper${err ? ` due to error: ${err.message}` : ''}` - }) + if (err) this.logger.info(`Destroying UDP piper due to error: ${err.message}`) + else this.logger.info('Destroying UDP piper') this.removeListeners() try { this.localStream?.destroy?.(err) @@ -221,7 +213,7 @@ class UdpConnPiper { this.remoteStream?.close?.(err) } catch (e) {} if (this.client) { - this.logger.log({ type: 1, msg: `Scheduling retry in ${this.retryDelay}ms` }) + this.logger.info(`Scheduling retry in ${this.retryDelay}ms`) setTimeout(() => { this.destroyed = false this.connect() @@ -270,13 +262,13 @@ function createUdpFramedProxy(listenOpts, connectRemote, logger, onBind) { const clients = new Map() proxySocket.on('error', (err) => { - logger.log({ type: 3, msg: `Proxy socket error: ${err.stack}` }) + logger.error(`Proxy socket error: ${err.stack}`) proxySocket.close() }) proxySocket.on('message', (msg, rinfo) => { const clientId = `${rinfo.address}:${rinfo.port}` - logger.log({ type: 0, msg: `UDP message from ${clientId}: ${msg.length} bytes` }) + logger.debug(`UDP message from ${clientId}: ${msg.length} bytes`) let client = clients.get(clientId) if (!client) { const remoteStream = connectRemote() @@ -288,20 +280,20 @@ function createUdpFramedProxy(listenOpts, connectRemote, logger, onBind) { const len = client.buffer.readUInt32BE(0) if (client.buffer.length < 4 + len) break const response = client.buffer.slice(4, 4 + len) - logger.log({ type: 0, msg: `UDP response for ${clientId}: ${response.length} bytes` }) + logger.debug(`UDP response for ${clientId}: ${response.length} bytes`) proxySocket.send(response, 0, response.length, rinfo.port, rinfo.address, (err) => { - if (err) logger.log({ type: 3, msg: `Send error to ${clientId}: ${err.stack}` }) + if (err) logger.error(`Send error to ${clientId}: ${err.stack}`) }) client.buffer = client.buffer.slice(4 + len) } }) remoteStream.on('error', (err) => { - logger.log({ type: 3, msg: `Remote error for ${clientId}: ${err.stack}` }) + logger.error(`Remote error for ${clientId}: ${err.stack}`) clients.delete(clientId) remoteStream.destroy() }) remoteStream.on('close', () => { - logger.log({ type: 0, msg: `Remote close for ${clientId}` }) + logger.debug(`Remote close for ${clientId}`) clients.delete(clientId) }) } @@ -319,24 +311,24 @@ function pipeUdpFramedServer(remoteStream, localOpts, logger, stats) { const localSocket = dgram.createSocket('udp4') let buffer = Buffer.alloc(0) localSocket.on('error', (err) => { - logger.log({ type: 3, msg: `Local UDP socket error: ${err.stack}` }) + logger.error(`Local UDP socket error: ${err.stack}`) remoteStream.destroy(err) }) localSocket.on('message', (msg) => { - logger.log({ type: 0, msg: `Data from local to remote: ${msg.length} bytes` }) + logger.debug(`Data from local to remote: ${msg.length} bytes`) const lenBuf = Buffer.alloc(4) lenBuf.writeUInt32BE(msg.length, 0) remoteStream.write(Buffer.concat([lenBuf, msg])) }) remoteStream.on('data', (chunk) => { - logger.log({ type: 0, msg: `Data from remote to local: ${chunk.length} bytes` }) + logger.debug(`Data from remote to local: ${chunk.length} bytes`) buffer = Buffer.concat([buffer, chunk]) while (buffer.length >= 4) { const len = buffer.readUInt32BE(0) if (buffer.length < 4 + len) break const msg = buffer.slice(4, 4 + len) localSocket.send(msg, 0, msg.length, +localOpts.port, localOpts.host, (err) => { - if (err) logger.log({ type: 3, msg: `Send error to local: ${err.stack}` }) + if (err) logger.error(`Send error to local: ${err.stack}`) }) buffer = buffer.slice(4 + len) } From 2511db6117edd3950e5e97076ed51dd3a4bacfe4 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Fri, 22 May 2026 09:25:07 +0530 Subject: [PATCH 02/10] refactor --- README.md | 6 +- index.js | 344 +--------------------------------------------- lib/tcp-piper.js | 69 ++++++++++ lib/udp-framed.js | 105 ++++++++++++++ package.json | 10 +- 5 files changed, 182 insertions(+), 352 deletions(-) create mode 100644 lib/tcp-piper.js create mode 100644 lib/udp-framed.js diff --git a/README.md b/README.md index a1f0c17..a9823a2 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # Hyper Cmd Lib Net -Network Library to Interface with Hyperswarm and local connections. Supports both UDP and TCP connections. +Network library to pipe local and hypderdht udp/tcp streams. ## Install @@ -10,10 +10,6 @@ Network Library to Interface with Hyperswarm and local connections. Supports bot #### `connPiper` -#### `udpPiper` - -#### `udpConnect` - #### `createTcpProxy` #### `pipeTcpServer` diff --git a/index.js b/index.js index 58c232a..5342184 100644 --- a/index.js +++ b/index.js @@ -1,348 +1,8 @@ -const dgram = require('bare-dgram') -const net = require('net') -const EventEmitter = require('events') - -function connPiper(connection, _dst, opts = {}, stats = {}) { - const logger = opts.logger || { - debug: () => {}, - info: () => {}, - warn: () => {}, - error: () => {} - } - logger.info('Starting TCP connection piper') - const loc = _dst() - if (loc === null) { - logger.warn('Connection rejected (null destination)') - connection.destroy() // don't return rejection error - if (!stats.rejectCnt) { - stats.rejectCnt = 0 - } - stats.rejectCnt++ - return - } - if (!stats.locCnt) { - stats.locCnt = 0 - } - if (!stats.remCnt) { - stats.remCnt = 0 - } - stats.locCnt++ - stats.remCnt++ - let destroyed = false - - // loc is the remote stream experimentation here - loc.on('data', (d) => { - logger.debug(`Data from local to remote: ${d.length} bytes`) - connection.write(d) - }) - connection.on('data', (d) => { - logger.debug(`Data from remote to local: ${d.length} bytes`) - loc.write(d) - }) - loc.on('error', destroy).on('close', destroy) - connection.on('error', destroy).on('close', destroy) - loc.on('end', () => { - logger.debug('Local end, ending connection') - connection.end() - }) - connection.on('end', () => { - logger.debug('Connection end, ending local') - loc.end() - }) - loc.on('connect', (err) => { - if (err) logger.error(err.message) - else logger.info('Connected') - }) - function destroy(err) { - if (destroyed) return - if (err) logger.info(`Destroying connection piper due to error: ${err.message}`) - else logger.info('Destroying connection piper') - stats.locCnt-- - stats.remCnt-- - destroyed = true - loc.end() - connection.end() - loc.destroy(err) - connection.destroy(err) - if (opts.onDestroy) { - opts.onDestroy(err) - } - } - return {} -} - -class UdpSocket { - constructor(opts) { - this.opts = opts - this.logger = this.opts.logger || { log: () => {} } - this.server = dgram.createSocket('udp4') - this.client = dgram.createSocket('udp4') - this.event = new EventEmitter() - this.rinfo = null - this.connect() - if (opts.bind) this.logger.info(`UDP socket created, binding to ${opts.host}:${opts.port}`) - else this.logger.info('UDP socket created') - } - - connect() { - if (this.opts.bind) { - this.server.bind(this.opts.port, this.opts.host) - } - this.server.on('message', (msg, rinfo) => { - this.logger.debug( - `UDP message from server: ${msg.length} bytes from ${rinfo.address}:${rinfo.port}` - ) - this.event.emit('message', msg, rinfo) - this.rinfo = rinfo - }) - this.client.on('message', (response, rinfo) => { - this.logger.debug( - `UDP message from client: ${response.length} bytes from ${rinfo.address}:${rinfo.port}` - ) - this.event.emit('message', response) - }) - this.client.on('error', (err) => { - this.logger.error(`UDP error: ${err.stack}`) - this.client.close() - }) - } - - write(msg) { - this.logger.debug(`Writing UDP message: ${msg.length} bytes`) - if (this.rinfo) { - this.server.send(msg, 0, msg.length, this.rinfo.port, this.rinfo.address) - } else { - this.client.send(msg, 0, msg.length, this.opts.port, this.opts.host) - } - } -} - -class UdpConnPiper { - constructor(remote, local, opts = {}) { - this.opts = opts - this.logger = this.opts.logger || { log: () => {} } - this.remote = remote - this.local = local - this.client = opts.client - this.retryDelay = opts.retryDelay || 2000 - this.destroyed = false - this._bindListeners() - this.connect() - if (opts.client) this.logger.info('Starting UDP connection piper (client mode)') - else this.logger.info('Starting UDP connection piper (server mode)') - } - - _bindListeners() { - this.bound = { - onLocMessage: this.onLocMessage.bind(this), - onConnectionMessage: this.onConnectionMessage.bind(this), - onLocError: this._handleError.bind(this), - onLocClose: this._handleError.bind(this), - onConnectionError: this._handleError.bind(this), - onConnectionClose: this._handleError.bind(this) - } - } - - connect() { - if (this.destroyed) return - this.logger.debug('Connecting UDP piper') - this.removeListeners() - this.localStream = typeof this.local === 'function' ? this.local() : this.local - this.remoteStream = typeof this.remote === 'function' ? this.remote() : this.remote - if (!this.localStream || !this.remoteStream) { - this.logger.warn('UDP connect failed (missing streams)') - this.destroy() - return - } - this.attachListeners() - } - - attachListeners() { - this.localStream.event.on('message', this.bound.onLocMessage) - this.localStream.server.on('error', this.bound.onLocError) - this.localStream.server.on('close', this.bound.onLocClose) - this.remoteStream.on('message', this.bound.onConnectionMessage) - this.remoteStream.on('error', this.bound.onConnectionError) - this.remoteStream.on('close', this.bound.onConnectionClose) - this.logger.debug('UDP listeners attached') - } - - removeListeners() { - if (this.localStream) { - this.localStream.event.off('message', this.bound.onLocMessage) - this.localStream.server.off('error', this.bound.onLocError) - this.localStream.server.off('close', this.bound.onLocClose) - } - if (this.remoteStream) { - this.remoteStream.off('message', this.bound.onConnectionMessage) - this.remoteStream.off('error', this.bound.onConnectionError) - this.remoteStream.off('close', this.bound.onConnectionClose) - } - this.logger.debug('UDP listeners removed') - } - - onLocMessage(msg, rinfo) { - this.logger.debug(`UDP message from local: ${msg.length} bytes`) - if (this.remoteStream && !this.destroyed) { - this.remoteStream.trySend?.(msg) - } - } - - onConnectionMessage(msg) { - this.logger.debug(`UDP message from connection: ${msg.length} bytes`) - if (this.localStream && !this.destroyed) { - this.localStream.write?.(msg) - } - } - - _handleError(err) { - this.logger.warn(`UDP error: ${err ? err.message : 'close'}`) - this.destroy(err) - } - - destroy(err) { - if (this.destroyed) return - this.destroyed = true - if (err) this.logger.info(`Destroying UDP piper due to error: ${err.message}`) - else this.logger.info('Destroying UDP piper') - this.removeListeners() - try { - this.localStream?.destroy?.(err) - } catch (e) {} - try { - this.remoteStream?.close?.(err) - } catch (e) {} - if (this.client) { - this.logger.info(`Scheduling retry in ${this.retryDelay}ms`) - setTimeout(() => { - this.destroyed = false - this.connect() - }, this.retryDelay) - } - } -} - -function udpConnect(opts, callback) { - const socket = new UdpSocket(opts) - if (typeof callback === 'function') { - callback(socket) - } else { - return socket - } -} - -function udpPiper(connection, _dst, opts) { - return new UdpConnPiper(connection, _dst, opts) -} - -function createTcpProxy(listenOpts, connectRemote, piperOpts, stats, onListen) { - const proxy = net.createServer({ allowHalfOpen: true }, (c) => { - connPiper(c, connectRemote, piperOpts, stats) - }) - proxy.listen(listenOpts.port, listenOpts.host, onListen) - return proxy -} - -function pipeTcpServer(remoteStream, localOpts, piperOpts, stats) { - connPiper( - remoteStream, - () => - net.connect({ - port: +localOpts.port, - host: localOpts.host, - allowHalfOpen: true - }), - piperOpts, - stats - ) -} - -function createUdpFramedProxy(listenOpts, connectRemote, logger, onBind) { - const proxySocket = dgram.createSocket('udp4') - const clients = new Map() - - proxySocket.on('error', (err) => { - logger.error(`Proxy socket error: ${err.stack}`) - proxySocket.close() - }) - - proxySocket.on('message', (msg, rinfo) => { - const clientId = `${rinfo.address}:${rinfo.port}` - logger.debug(`UDP message from ${clientId}: ${msg.length} bytes`) - let client = clients.get(clientId) - if (!client) { - const remoteStream = connectRemote() - client = { remoteStream, rinfo, buffer: Buffer.alloc(0) } - clients.set(clientId, client) - remoteStream.on('data', (chunk) => { - client.buffer = Buffer.concat([client.buffer, chunk]) - while (client.buffer.length >= 4) { - const len = client.buffer.readUInt32BE(0) - if (client.buffer.length < 4 + len) break - const response = client.buffer.slice(4, 4 + len) - logger.debug(`UDP response for ${clientId}: ${response.length} bytes`) - proxySocket.send(response, 0, response.length, rinfo.port, rinfo.address, (err) => { - if (err) logger.error(`Send error to ${clientId}: ${err.stack}`) - }) - client.buffer = client.buffer.slice(4 + len) - } - }) - remoteStream.on('error', (err) => { - logger.error(`Remote error for ${clientId}: ${err.stack}`) - clients.delete(clientId) - remoteStream.destroy() - }) - remoteStream.on('close', () => { - logger.debug(`Remote close for ${clientId}`) - clients.delete(clientId) - }) - } - const lenBuf = Buffer.alloc(4) - lenBuf.writeUInt32BE(msg.length, 0) - client.remoteStream.write(Buffer.concat([lenBuf, msg])) - }) - - proxySocket.bind(listenOpts.port, listenOpts.host, onBind) - - return { proxySocket, clients } -} - -function pipeUdpFramedServer(remoteStream, localOpts, logger, stats) { - const localSocket = dgram.createSocket('udp4') - let buffer = Buffer.alloc(0) - localSocket.on('error', (err) => { - logger.error(`Local UDP socket error: ${err.stack}`) - remoteStream.destroy(err) - }) - localSocket.on('message', (msg) => { - logger.debug(`Data from local to remote: ${msg.length} bytes`) - const lenBuf = Buffer.alloc(4) - lenBuf.writeUInt32BE(msg.length, 0) - remoteStream.write(Buffer.concat([lenBuf, msg])) - }) - remoteStream.on('data', (chunk) => { - logger.debug(`Data from remote to local: ${chunk.length} bytes`) - buffer = Buffer.concat([buffer, chunk]) - while (buffer.length >= 4) { - const len = buffer.readUInt32BE(0) - if (buffer.length < 4 + len) break - const msg = buffer.slice(4, 4 + len) - localSocket.send(msg, 0, msg.length, +localOpts.port, localOpts.host, (err) => { - if (err) logger.error(`Send error to local: ${err.stack}`) - }) - buffer = buffer.slice(4 + len) - } - }) - remoteStream.on('end', () => localSocket.close()) - localSocket.on('close', () => remoteStream.end()) - remoteStream.on('error', (err) => localSocket.close()) - localSocket.on('error', (err) => remoteStream.destroy(err)) -} +const { connPiper, createTcpProxy, pipeTcpServer } = require('./lib/tcp-piper.js') +const { createUdpFramedProxy, pipeUdpFramedServer } = require('./lib/udp-framed.js') module.exports = { connPiper, - udpPiper, - udpConnect, createTcpProxy, pipeTcpServer, createUdpFramedProxy, diff --git a/lib/tcp-piper.js b/lib/tcp-piper.js new file mode 100644 index 0000000..3c9d6c4 --- /dev/null +++ b/lib/tcp-piper.js @@ -0,0 +1,69 @@ +const net = require('net') + +function connPiper(a, bFactory, opts = {}) { + const logger = opts.logger || { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {} + } + + logger.info('Starting TCP connection piper') + + let b + + try { + b = bFactory() + } catch (err) { + logger.error(`b factory error: ${err.message}`) + a.destroy() + } + + if (b === null) { + logger.warn('Connection rejected (null b)') + a.destroy() + return + } + + let destroyed = false + const destroy = (err) => { + if (destroyed) return + destroyed = true + if (err) logger.info(`Destroying piper due to error: ${err.message}`) + else logger.info('Destroying piper') + a.destroy(err) + b.destroy(err) + opts.onDestroy?.(err) + } + + a.pipe(b) + b.pipe(a) + + a.on('error', destroy).on('close', () => destroy()) + b.on('error', destroy).on('close', () => destroy()) + a.on('data', (d) => logger.debug(`a→b: ${d.length} bytes`)) + b.on('data', (d) => logger.debug(`b→a: ${d.length} bytes`)) + a.on('end', () => logger.debug('a ended (peer-a sent FIN); pipe will end b')) + b.on('end', () => logger.debug('b ended (peer-b sent FIN); pipe will end a')) + b.on('connect', () => logger.info('Connected')) +} + +function createTcpProxy(remoteTunnel, opts, onListen) { + const proxy = net.createServer({ allowHalfOpen: true }, (localStream) => { + connPiper(localStream, remoteTunnel, opts) + }) + proxy.listen(opts.port, opts.host, onListen) + return proxy +} + +function pipeTcpServer(remoteStream, leftover, opts) { + const sock = net.connect({ + port: opts.port, + host: opts.host, + allowHalfOpen: true + }) + if (leftover && leftover.length) sock.write(leftover) + return connPiper(remoteStream, () => sock, opts) +} + +module.exports = { connPiper, createTcpProxy, pipeTcpServer } diff --git a/lib/udp-framed.js b/lib/udp-framed.js new file mode 100644 index 0000000..ea065c1 --- /dev/null +++ b/lib/udp-framed.js @@ -0,0 +1,105 @@ +const dgram = require('bare-dgram') + +function createUdpFramedProxy(createTunnel, opts = {}, onBind) { + const port = opts.port + const host = opts.host + const logger = opts.logger || { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {} + } + + const proxySocket = dgram.createSocket('udp4') + const clients = new Map() + + proxySocket.on('error', (err) => { + logger.error(`Proxy socket error: ${err.stack}`) + proxySocket.close() + }) + + proxySocket.on('message', (msg, rinfo) => { + const clientId = `${rinfo.address}:${rinfo.port}` + logger.debug(`UDP message from ${clientId}: ${msg.length} bytes`) + let client = clients.get(clientId) + if (!client) { + const stream = createTunnel() + client = { stream, rinfo, buffer: Buffer.alloc(0) } + clients.set(clientId, client) + stream.on('data', (chunk) => { + client.buffer = Buffer.concat([client.buffer, chunk]) + while (client.buffer.length >= 4) { + const len = client.buffer.readUInt32BE(0) + if (client.buffer.length < 4 + len) break + const response = client.buffer.slice(4, 4 + len) + logger.debug(`UDP response for ${clientId}: ${response.length} bytes`) + proxySocket.send(response, 0, response.length, rinfo.port, rinfo.address, (err) => { + if (err) logger.error(`Socket error while sending packet to ${clientId}: ${err.stack}`) + }) + client.buffer = client.buffer.slice(4 + len) + } + }) + stream.on('error', (err) => { + logger.error(`Remote error for ${clientId}: ${err.stack}`) + clients.delete(clientId) + stream.destroy() + }) + stream.on('close', () => { + logger.debug(`Remote close for ${clientId}`) + clients.delete(clientId) + }) + } + const lenBuf = Buffer.alloc(4) + lenBuf.writeUInt32BE(msg.length, 0) + client.stream.write(Buffer.concat([lenBuf, msg])) + }) + + proxySocket.bind(port, host, onBind) + + return { proxySocket, clients } +} + +function pipeUdpFramedServer(stream, leftover, opts = {}) { + const port = opts.port + const host = opts.host + const logger = opts.logger + + const localSocket = dgram.createSocket('udp4') + let buffer = leftover && leftover.length ? Buffer.from(leftover) : Buffer.alloc(0) + + const drain = () => { + while (buffer.length >= 4) { + const len = buffer.readUInt32BE(0) + if (buffer.length < 4 + len) break + const msg = buffer.slice(4, 4 + len) + localSocket.send(msg, 0, msg.length, port, host, (err) => { + if (err) logger.error(`Error sending packet to local socket: ${err.stack}`) + }) + buffer = buffer.slice(4 + len) + } + } + + localSocket.on('error', (err) => { + logger.error(`Local UDP socket error: ${err.stack}`) + stream.destroy(err) + }) + localSocket.on('message', (msg) => { + logger.debug(`Data from local to remote: ${msg.length} bytes`) + const lenBuf = Buffer.alloc(4) + lenBuf.writeUInt32BE(msg.length, 0) + stream.write(Buffer.concat([lenBuf, msg])) + }) + stream.on('data', (chunk) => { + logger.debug(`Data from remote to local: ${chunk.length} bytes`) + buffer = Buffer.concat([buffer, chunk]) + drain() + }) + stream.on('end', () => localSocket.close()) + localSocket.on('close', () => stream.end()) + stream.on('error', () => localSocket.close()) + localSocket.on('error', (err) => stream.destroy(err)) + + if (buffer.length) drain() +} + +module.exports = { createUdpFramedProxy, pipeUdpFramedServer } diff --git a/package.json b/package.json index d12c52a..f77df72 100644 --- a/package.json +++ b/package.json @@ -1,11 +1,11 @@ { "name": "@holesail/hyper-cmd-lib-net", "version": "1.1.2", - "description": "Network Library to Interface with Hyperswarm and local connections", + "description": "Network library to pipe local and hypderdht udp/tcp streams", "main": "index.js", "scripts": { "test": "npx prettier . --check", - "lint": "npx prettier . --write" + "format": "npx prettier . --write" }, "repository": { "type": "git", @@ -25,11 +25,11 @@ "homepage": "https://github.com/holesail/hyper-cmd-lib-net/", "dependencies": { "bare-dgram": "^1.0.1", - "bare-net": "^2.2.0", - "net": "npm:bare-net@^2.2.0" + "bare-net": "^2.3.1", + "net": "npm:bare-net@^2.3.1" }, "devDependencies": { - "prettier": "^3.6.2", + "prettier": "^3.8.3", "prettier-config-holepunch": "^2.0.0" } } From 21914afa729077cb306516dd8ea4af557b94f0d7 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Fri, 22 May 2026 09:40:51 +0530 Subject: [PATCH 03/10] disable Nagle's algorithm --- lib/tcp-piper.js | 2 ++ 1 file changed, 2 insertions(+) diff --git a/lib/tcp-piper.js b/lib/tcp-piper.js index 3c9d6c4..2be6a65 100644 --- a/lib/tcp-piper.js +++ b/lib/tcp-piper.js @@ -50,6 +50,7 @@ function connPiper(a, bFactory, opts = {}) { function createTcpProxy(remoteTunnel, opts, onListen) { const proxy = net.createServer({ allowHalfOpen: true }, (localStream) => { + localStream.setNoDelay(true) connPiper(localStream, remoteTunnel, opts) }) proxy.listen(opts.port, opts.host, onListen) @@ -62,6 +63,7 @@ function pipeTcpServer(remoteStream, leftover, opts) { host: opts.host, allowHalfOpen: true }) + sock.setNoDelay(true) if (leftover && leftover.length) sock.write(leftover) return connPiper(remoteStream, () => sock, opts) } From 9607632fd01d25191ba9476968f912bdd71e4761 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Fri, 22 May 2026 09:45:49 +0530 Subject: [PATCH 04/10] minor fix --- lib/tcp-piper.js | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/tcp-piper.js b/lib/tcp-piper.js index 2be6a65..7fa03fb 100644 --- a/lib/tcp-piper.js +++ b/lib/tcp-piper.js @@ -17,6 +17,7 @@ function connPiper(a, bFactory, opts = {}) { } catch (err) { logger.error(`b factory error: ${err.message}`) a.destroy() + return } if (b === null) { From d676e67c5c85935ff0d5f6fec5ea415b901d6a80 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Fri, 22 May 2026 15:08:55 +0530 Subject: [PATCH 05/10] fix cleanup and add a max packet size --- lib/udp-framed.js | 52 ++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 42 insertions(+), 10 deletions(-) diff --git a/lib/udp-framed.js b/lib/udp-framed.js index ea065c1..d22980e 100644 --- a/lib/udp-framed.js +++ b/lib/udp-framed.js @@ -1,5 +1,7 @@ const dgram = require('bare-dgram') +const DEFAULT_MAX_FRAME_SIZE = 65535 + function createUdpFramedProxy(createTunnel, opts = {}, onBind) { const port = opts.port const host = opts.host @@ -9,6 +11,7 @@ function createUdpFramedProxy(createTunnel, opts = {}, onBind) { warn: () => {}, error: () => {} } + const maxFrameSize = opts.maxFrameSize ?? DEFAULT_MAX_FRAME_SIZE const proxySocket = dgram.createSocket('udp4') const clients = new Map() @@ -30,13 +33,19 @@ function createUdpFramedProxy(createTunnel, opts = {}, onBind) { client.buffer = Buffer.concat([client.buffer, chunk]) while (client.buffer.length >= 4) { const len = client.buffer.readUInt32BE(0) + if (len > maxFrameSize) { + logger.error(`Frame too large from ${clientId}: ${len} > ${maxFrameSize}`) + clients.delete(clientId) + stream.destroy(new Error(`Frame too large: ${len}`)) + return + } if (client.buffer.length < 4 + len) break - const response = client.buffer.slice(4, 4 + len) + const response = client.buffer.subarray(4, 4 + len) logger.debug(`UDP response for ${clientId}: ${response.length} bytes`) proxySocket.send(response, 0, response.length, rinfo.port, rinfo.address, (err) => { if (err) logger.error(`Socket error while sending packet to ${clientId}: ${err.stack}`) }) - client.buffer = client.buffer.slice(4 + len) + client.buffer = client.buffer.subarray(4 + len) } }) stream.on('error', (err) => { @@ -62,26 +71,49 @@ function createUdpFramedProxy(createTunnel, opts = {}, onBind) { function pipeUdpFramedServer(stream, leftover, opts = {}) { const port = opts.port const host = opts.host - const logger = opts.logger + const logger = opts.logger || { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {} + } + const maxFrameSize = opts.maxFrameSize ?? DEFAULT_MAX_FRAME_SIZE const localSocket = dgram.createSocket('udp4') let buffer = leftover && leftover.length ? Buffer.from(leftover) : Buffer.alloc(0) + let destroyed = false + const cleanup = (err) => { + if (destroyed) return + destroyed = true + if (err) logger.info(`Destroying pipe due to error: ${err.message}`) + else logger.info('Destroying pipe') + try { + localSocket.close() + } catch {} + if (err) stream.destroy(err) + else stream.end() + } + const drain = () => { while (buffer.length >= 4) { const len = buffer.readUInt32BE(0) + if (len > maxFrameSize) { + cleanup(new Error(`Frame too large: ${len} > ${maxFrameSize}`)) + return + } if (buffer.length < 4 + len) break - const msg = buffer.slice(4, 4 + len) + const msg = buffer.subarray(4, 4 + len) localSocket.send(msg, 0, msg.length, port, host, (err) => { if (err) logger.error(`Error sending packet to local socket: ${err.stack}`) }) - buffer = buffer.slice(4 + len) + buffer = buffer.subarray(4 + len) } } localSocket.on('error', (err) => { logger.error(`Local UDP socket error: ${err.stack}`) - stream.destroy(err) + cleanup(err) }) localSocket.on('message', (msg) => { logger.debug(`Data from local to remote: ${msg.length} bytes`) @@ -89,15 +121,15 @@ function pipeUdpFramedServer(stream, leftover, opts = {}) { lenBuf.writeUInt32BE(msg.length, 0) stream.write(Buffer.concat([lenBuf, msg])) }) + localSocket.on('close', () => cleanup()) stream.on('data', (chunk) => { logger.debug(`Data from remote to local: ${chunk.length} bytes`) buffer = Buffer.concat([buffer, chunk]) drain() }) - stream.on('end', () => localSocket.close()) - localSocket.on('close', () => stream.end()) - stream.on('error', () => localSocket.close()) - localSocket.on('error', (err) => stream.destroy(err)) + stream.on('end', () => cleanup()) + stream.on('close', () => cleanup()) + stream.on('error', (err) => cleanup(err)) if (buffer.length) drain() } From 1c8136f01a05c9313cbe68521a59de5138d5abf2 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Thu, 4 Jun 2026 10:01:31 +0530 Subject: [PATCH 06/10] make sure port is numeric --- lib/tcp-piper.js | 2 +- package.json | 4 ++++ 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/lib/tcp-piper.js b/lib/tcp-piper.js index 7fa03fb..b8d4da9 100644 --- a/lib/tcp-piper.js +++ b/lib/tcp-piper.js @@ -60,7 +60,7 @@ function createTcpProxy(remoteTunnel, opts, onListen) { function pipeTcpServer(remoteStream, leftover, opts) { const sock = net.connect({ - port: opts.port, + port: +opts.port, host: opts.host, allowHalfOpen: true }) diff --git a/package.json b/package.json index f77df72..c8692db 100644 --- a/package.json +++ b/package.json @@ -15,6 +15,10 @@ "events": { "bare": "bare-events", "default": "events" + }, + "net": { + "bare": "bare-net", + "default": "net" } }, "author": "supersuryaansh", From 7fc785d1d320e39dfb07e66e4563e1ddf6a8d2c9 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Sat, 25 Jul 2026 10:18:43 +0530 Subject: [PATCH 07/10] add test suite --- .github/workflows/test-bare.yml | 24 +++++ .github/workflows/test-node.yml | 2 +- lib/tcp-piper.js | 11 ++ package.json | 12 ++- test/connection-piper.js | 185 ++++++++++++++++++++++++++++++++ test/helpers.js | 132 +++++++++++++++++++++++ test/index.js | 5 + test/tcp-proxy.js | 116 ++++++++++++++++++++ test/tcp-server.js | 66 ++++++++++++ test/udp-framed-proxy.js | 162 ++++++++++++++++++++++++++++ test/udp-framed-server.js | 145 +++++++++++++++++++++++++ 11 files changed, 855 insertions(+), 5 deletions(-) create mode 100644 .github/workflows/test-bare.yml create mode 100644 test/connection-piper.js create mode 100644 test/helpers.js create mode 100644 test/index.js create mode 100644 test/tcp-proxy.js create mode 100644 test/tcp-server.js create mode 100644 test/udp-framed-proxy.js create mode 100644 test/udp-framed-server.js diff --git a/.github/workflows/test-bare.yml b/.github/workflows/test-bare.yml new file mode 100644 index 0000000..b9dc344 --- /dev/null +++ b/.github/workflows/test-bare.yml @@ -0,0 +1,24 @@ +name: Test Bare Status +on: + push: + branches: + - main + pull_request: + branches: + - main +jobs: + build: + strategy: + matrix: + node-version: [lts/*] + os: [ubuntu-latest, macos-latest, windows-latest] + runs-on: ${{ matrix.os }} + steps: + - uses: actions/checkout@v3 + - name: Use Node.js ${{ matrix.node-version }} + uses: actions/setup-node@v3 + with: + node-version: ${{ matrix.node-version }} + - run: npm install -g bare + - run: npm install + - run: npm run test-bare diff --git a/.github/workflows/test-node.yml b/.github/workflows/test-node.yml index a6f145d..fc03e8a 100644 --- a/.github/workflows/test-node.yml +++ b/.github/workflows/test-node.yml @@ -1,4 +1,4 @@ -name: Test Status +name: Test Node Status on: push: branches: diff --git a/lib/tcp-piper.js b/lib/tcp-piper.js index b8d4da9..0d4a119 100644 --- a/lib/tcp-piper.js +++ b/lib/tcp-piper.js @@ -50,11 +50,22 @@ function connPiper(a, bFactory, opts = {}) { } function createTcpProxy(remoteTunnel, opts, onListen) { + const sockets = new Set() + const proxy = net.createServer({ allowHalfOpen: true }, (localStream) => { localStream.setNoDelay(true) + sockets.add(localStream) + localStream.on('close', () => sockets.delete(localStream)) connPiper(localStream, remoteTunnel, opts) }) proxy.listen(opts.port, opts.host, onListen) + + const close = proxy.close.bind(proxy) + proxy.close = (cb) => { + for (const sock of sockets) sock.destroy() + return close(cb) + } + return proxy } diff --git a/package.json b/package.json index c8692db..79eae22 100644 --- a/package.json +++ b/package.json @@ -4,7 +4,8 @@ "description": "Network library to pipe local and hypderdht udp/tcp streams", "main": "index.js", "scripts": { - "test": "npx prettier . --check", + "test": "prettier . --check && node test/index.js", + "test-bare": "prettier . --check && bare test/index.js", "format": "npx prettier . --write" }, "repository": { @@ -29,11 +30,14 @@ "homepage": "https://github.com/holesail/hyper-cmd-lib-net/", "dependencies": { "bare-dgram": "^1.0.1", - "bare-net": "^2.3.1", - "net": "npm:bare-net@^2.3.1" + "bare-net": "^2.3.2", + "bare-tcp": "^2.5.3", + "net": "npm:bare-net@^2.3.2" }, "devDependencies": { + "brittle": "^4.0.0", "prettier": "^3.8.3", - "prettier-config-holepunch": "^2.0.0" + "prettier-config-holepunch": "^2.0.0", + "streamx": "^2.28.0" } } diff --git a/test/connection-piper.js b/test/connection-piper.js new file mode 100644 index 0000000..c69db45 --- /dev/null +++ b/test/connection-piper.js @@ -0,0 +1,185 @@ +const test = require('brittle') +const net = require('net') +const { connPiper } = require('../lib/tcp-piper.js') +const { spyLogger, tick, listenTcp, tcpPair, closePair, readOnce } = require('./helpers.js') + +test('connPiper - pipes data from a to b', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + + connPiper(A.local, () => B.local) + const received = readOnce(B.remote) + A.remote.write('hello-a-to-b') + + t.is(await received, 'hello-a-to-b') +}) + +test('connPiper - pipes data from b to a', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + + connPiper(A.local, () => B.local) + const received = readOnce(A.remote) + B.remote.write('hello-b-to-a') + + t.is(await received, 'hello-b-to-a') +}) + +test('connPiper - bFactory throwing destroys a and logs the error', async function (t) { + const A = await tcpPair() + t.teardown(() => closePair(A)) + const { logger, calls } = spyLogger() + + connPiper( + A.local, + () => { + throw new Error('boom') + }, + { logger } + ) + await tick() + + t.ok(A.local.destroyed) + t.is(calls.error.length, 1) + t.ok(/b factory error: boom/.test(calls.error[0][0])) +}) + +test('connPiper - bFactory returning null destroys a and logs a warning', async function (t) { + const A = await tcpPair() + t.teardown(() => closePair(A)) + const { logger, calls } = spyLogger() + + connPiper(A.local, () => null, { logger }) + await tick() + + t.ok(A.local.destroyed) + t.is(calls.warn.length, 1) + t.ok(/Connection rejected \(null b\)/.test(calls.warn[0][0])) +}) + +test('connPiper - works with the default no-op logger', async function (t) { + const A = await tcpPair() + t.teardown(() => closePair(A)) + + connPiper(A.local, () => null) + await tick() + t.ok(A.local.destroyed) +}) + +test('connPiper - a erroring destroys b exactly once and notifies onDestroy', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + const destroys = [] + + connPiper(A.local, () => B.local, { onDestroy: (err) => destroys.push(err) }) + A.local.destroy(new Error('kaboom')) + await tick() + + t.ok(B.local.destroyed) + t.is(destroys.length, 1) + t.is(destroys[0].message, 'kaboom') +}) + +test('connPiper - b erroring destroys a exactly once and notifies onDestroy', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + const destroys = [] + + connPiper(A.local, () => B.local, { onDestroy: (err) => destroys.push(err) }) + B.local.destroy(new Error('kaboom-b')) + await tick() + + t.ok(A.local.destroyed) + t.is(destroys.length, 1) + t.is(destroys[0].message, 'kaboom-b') +}) + +test('connPiper - destroy is idempotent across both error and close events', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + const destroys = [] + + connPiper(A.local, () => B.local, { onDestroy: (err) => destroys.push(err) }) + // destroy(err) on a real socket emits both 'error' and 'close', and both are + // wired to the same internal destroy() closure — this alone exercises the guard. + A.local.destroy(new Error('once-only')) + await tick() + + t.is(destroys.length, 1) +}) + +test('connPiper - logs debug messages describing byte counts in both directions', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + const { logger, calls } = spyLogger() + + connPiper(A.local, () => B.local, { logger }) + A.remote.write('abc') + B.remote.write('de') + await tick() + + t.ok(calls.debug.some((args) => /a.b: 3 bytes/.test(args[0]))) + t.ok(calls.debug.some((args) => /b.a: 2 bytes/.test(args[0]))) +}) + +test('connPiper - logs Connected when b emits a connect event', async function (t) { + const A = await tcpPair() + t.teardown(() => closePair(A)) + + const bServer = net.createServer({ allowHalfOpen: true }) + t.teardown(() => bServer.close()) + const addr = await listenTcp(bServer) + + const { logger, calls } = spyLogger() + const b = net.connect(addr.port, '127.0.0.1') + t.teardown(() => b.destroy()) + // net.connect() hasn't fired 'connect' yet here — connPiper's own listener + // attaches before that happens, same as pipeTcpServer does in production. + connPiper(A.local, () => b, { logger }) + + await new Promise((resolve) => bServer.once('connection', resolve)) + await tick() + + t.ok(calls.info.some((args) => args[0] === 'Connected')) +}) + +test('connPiper - a ending its readable side ends b in turn', async function (t) { + const A = await tcpPair() + const B = await tcpPair() + t.teardown(() => { + closePair(A) + closePair(B) + }) + const { logger, calls } = spyLogger() + + connPiper(A.local, () => B.local, { logger }) + const bFinished = new Promise((resolve) => B.local.once('finish', resolve)) + A.remote.end() + await bFinished + + t.ok(calls.debug.some((args) => /a ended/.test(args[0]))) +}) diff --git a/test/helpers.js b/test/helpers.js new file mode 100644 index 0000000..39e7a8a --- /dev/null +++ b/test/helpers.js @@ -0,0 +1,132 @@ +const { Duplex } = require('streamx') +const net = require('net') +const dgram = require('bare-dgram') + +function fakeSocket() { + const written = [] + const s = new Duplex({ + read() {}, + write(chunk, cb) { + written.push(chunk) + cb(null) + } + }) + s.written = written + return s +} + +// A stream that echoes back whatever is written to it, simulating a remote +// peer that reflects data (used as the "far end" of a tunnel/proxy). +function echoStream() { + return new Duplex({ + read() {}, + write(chunk, cb) { + this.push(chunk) + cb(null) + } + }) +} + +// Similar to echoStream(), but pushes each write back in two separate chunks +function splitEchoStream() { + return new Duplex({ + read() {}, + write(chunk, cb) { + const mid = Math.floor(chunk.length / 2) || 1 + this.push(chunk.subarray(0, mid)) + setImmediate(() => { + this.push(chunk.subarray(mid)) + cb(null) + }) + } + }) +} + +function spyLogger() { + const calls = { debug: [], info: [], warn: [], error: [] } + const logger = { + debug: (...a) => calls.debug.push(a), + info: (...a) => calls.info.push(a), + warn: (...a) => calls.warn.push(a), + error: (...a) => calls.error.push(a) + } + return { logger, calls } +} + +function frame(payload) { + const len = Buffer.alloc(4) + len.writeUInt32BE(payload.length, 0) + return Buffer.concat([len, payload]) +} + +function tick(ms = 20) { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +function bindUdpClient() { + return new Promise((resolve, reject) => { + const sock = dgram.createSocket('udp4') + sock.on('error', reject) + sock.bind(0, '127.0.0.1', () => resolve(sock)) + }) +} + +function listenTcp(server) { + return new Promise((resolve) => { + server.listen(0, '127.0.0.1', () => resolve(server.address())) + }) +} + +// Two real, independently connected loopback TCP sockets. `local` is the +// server-accepted side — what gets handed to connPiper as `a` or `b` — and +// `remote` is the test's own hook into "the other end" of that connection, +// used to inject inbound bytes (remote.write) and observe whatever connPiper +// forwarded out to it (remote.on('data', ...)). +function tcpPair() { + return new Promise((resolve) => { + const server = net.createServer({ allowHalfOpen: true }) + server.listen(0, '127.0.0.1', () => { + const addr = server.address() + const remote = net.connect(addr.port, '127.0.0.1') + server.once('connection', (local) => resolve({ local, remote, server })) + }) + }) +} + +function closePair(pair) { + pair.remote.destroy() + pair.local.destroy() + pair.server.close() +} + +function readOnce(sock) { + return new Promise((resolve) => sock.once('data', (d) => resolve(d.toString()))) +} + +async function withEchoUdpServer(fn) { + const server = dgram.createSocket('udp4') + server.on('message', (msg, rinfo) => { + server.send(msg, 0, msg.length, rinfo.port, rinfo.address) + }) + await new Promise((resolve) => server.bind(0, '127.0.0.1', resolve)) + try { + await fn(server.address()) + } finally { + server.close() + } +} + +module.exports = { + fakeSocket, + echoStream, + splitEchoStream, + spyLogger, + frame, + tick, + bindUdpClient, + listenTcp, + tcpPair, + closePair, + readOnce, + withEchoUdpServer +} diff --git a/test/index.js b/test/index.js new file mode 100644 index 0000000..c5cd179 --- /dev/null +++ b/test/index.js @@ -0,0 +1,5 @@ +require('./connection-piper.js') +require('./tcp-proxy.js') +require('./tcp-server.js') +require('./udp-framed-proxy.js') +require('./udp-framed-server.js') diff --git a/test/tcp-proxy.js b/test/tcp-proxy.js new file mode 100644 index 0000000..3ee994b --- /dev/null +++ b/test/tcp-proxy.js @@ -0,0 +1,116 @@ +const test = require('brittle') +const net = require('net') +const { createTcpProxy } = require('../lib/tcp-piper.js') +const { echoStream, tick, listenTcp } = require('./helpers.js') + +test('createTcpProxy - proxies a real TCP client round trip through the tunnel factory', async function (t) { + const addr = await new Promise((resolve) => { + const proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => resolve(proxy.address()) + ) + t.teardown(() => proxy.close()) + }) + + const client = net.connect(addr.port, '127.0.0.1') + t.teardown(() => client.destroy()) + await new Promise((resolve) => client.once('connect', resolve)) + + const reply = new Promise((resolve) => client.once('data', (d) => resolve(d.toString()))) + client.write('ping-through-proxy') + + t.is(await reply, 'ping-through-proxy') +}) + +test('createTcpProxy - calls onListen once the server is bound', async function (t) { + let called = false + const proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => { + called = true + } + ) + t.teardown(() => proxy.close()) + + await new Promise((resolve) => proxy.once('listening', resolve)) + t.ok(called) +}) + +test('createTcpProxy - invokes the tunnel factory once per incoming connection', async function (t) { + let calls = 0 + const proxy = createTcpProxy( + () => { + calls++ + return echoStream() + }, + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.close()) + + const addr = await listenTcp(proxy) + const c1 = net.connect(addr.port, '127.0.0.1') + const c2 = net.connect(addr.port, '127.0.0.1') + t.teardown(() => { + c1.destroy() + c2.destroy() + }) + await Promise.all([ + new Promise((resolve) => c1.once('connect', resolve)), + new Promise((resolve) => c2.once('connect', resolve)) + ]) + await tick() + + t.is(calls, 2) +}) + +test('createTcpProxy - two connections stay isolated from each other', async function (t) { + const proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.close()) + + const addr = await listenTcp(proxy) + const c1 = net.connect(addr.port, '127.0.0.1') + const c2 = net.connect(addr.port, '127.0.0.1') + t.teardown(() => { + c1.destroy() + c2.destroy() + }) + await Promise.all([ + new Promise((resolve) => c1.once('connect', resolve)), + new Promise((resolve) => c2.once('connect', resolve)) + ]) + + const r1 = new Promise((resolve) => c1.once('data', (d) => resolve(d.toString()))) + const r2 = new Promise((resolve) => c2.once('data', (d) => resolve(d.toString()))) + c1.write('from-c1') + c2.write('from-c2') + + t.is(await r1, 'from-c1') + t.is(await r2, 'from-c2') +}) + +test('createTcpProxy - closing the proxy destroys any still-open connections', async function (t) { + const proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => {} + ) + + const addr = await listenTcp(proxy) + const client = net.connect(addr.port, '127.0.0.1') + client.on('error', () => {}) // the far end destroys mid-flight; ECONNRESET is expected here + t.teardown(() => client.destroy()) + await new Promise((resolve) => client.once('connect', resolve)) + + const closed = new Promise((resolve) => client.once('close', resolve)) + await new Promise((resolve) => proxy.close(resolve)) + await closed + + t.pass('client socket was closed by proxy.close()') +}) diff --git a/test/tcp-server.js b/test/tcp-server.js new file mode 100644 index 0000000..1a5181b --- /dev/null +++ b/test/tcp-server.js @@ -0,0 +1,66 @@ +const test = require('brittle') +const net = require('net') +const { pipeTcpServer } = require('../lib/tcp-piper.js') +const { fakeSocket, tick, listenTcp } = require('./helpers.js') + +test('pipeTcpServer - writes leftover bytes to the local connection before piping', async function (t) { + const echoServer = net.createServer((sock) => sock.pipe(sock)) + t.teardown(() => echoServer.close()) + const addr = await listenTcp(echoServer) + + const remoteStream = fakeSocket() + t.teardown(() => remoteStream.destroy()) + pipeTcpServer(remoteStream, Buffer.from('LEFTOVER'), { port: addr.port, host: '127.0.0.1' }) + await tick(30) + + t.alike(remoteStream.written, [Buffer.from('LEFTOVER')]) +}) + +test('pipeTcpServer - accepts a string port and pipes further data after the leftover', async function (t) { + const echoServer = net.createServer((sock) => sock.pipe(sock)) + t.teardown(() => echoServer.close()) + const addr = await listenTcp(echoServer) + + const remoteStream = fakeSocket() + t.teardown(() => remoteStream.destroy()) + pipeTcpServer(remoteStream, Buffer.from('LEFTOVER'), { + port: String(addr.port), + host: '127.0.0.1' + }) + await tick(30) + remoteStream.push(Buffer.from('MORE-DATA')) + await tick(30) + + t.alike(remoteStream.written, [Buffer.from('LEFTOVER'), Buffer.from('MORE-DATA')]) +}) + +test('pipeTcpServer - works with no leftover at all', async function (t) { + const echoServer = net.createServer((sock) => sock.pipe(sock)) + t.teardown(() => echoServer.close()) + const addr = await listenTcp(echoServer) + + const remoteStream = fakeSocket() + t.teardown(() => remoteStream.destroy()) + pipeTcpServer(remoteStream, null, { port: addr.port, host: '127.0.0.1' }) + await tick(20) + t.is(remoteStream.written.length, 0, 'nothing sent yet without leftover or data') + + remoteStream.push(Buffer.from('first-word')) + await tick(30) + t.alike(remoteStream.written, [Buffer.from('first-word')]) +}) + +test('pipeTcpServer - a refused TCP connection destroys the remote stream', async function (t) { + // Bind a server, close it immediately, then connect: the port is free but + // nothing is listening, so the connect attempt fails with ECONNREFUSED. + const throwaway = net.createServer(() => {}) + const addr = await listenTcp(throwaway) + await new Promise((resolve) => throwaway.close(resolve)) + + const remoteStream = fakeSocket() + const closed = new Promise((resolve) => remoteStream.once('close', resolve)) + pipeTcpServer(remoteStream, null, { port: addr.port, host: '127.0.0.1' }) + await closed + + t.ok(remoteStream.destroyed) +}) diff --git a/test/udp-framed-proxy.js b/test/udp-framed-proxy.js new file mode 100644 index 0000000..64ffe6f --- /dev/null +++ b/test/udp-framed-proxy.js @@ -0,0 +1,162 @@ +const test = require('brittle') +const { createUdpFramedProxy } = require('../lib/udp-framed.js') +const { echoStream, splitEchoStream, spyLogger, tick, bindUdpClient } = require('./helpers.js') + +test('createUdpFramedProxy - round-trips a datagram through the tunnel factory', async function (t) { + const proxy = createUdpFramedProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const client = await bindUdpClient() + t.teardown(() => client.close()) + + const reply = new Promise((resolve) => client.once('message', (msg) => resolve(msg.toString()))) + const payload = Buffer.from('hello-udp') + client.send(payload, 0, payload.length, addr.port, '127.0.0.1') + + t.is(await reply, 'hello-udp') +}) + +test('createUdpFramedProxy - reuses the same tunnel for repeated datagrams from one client', async function (t) { + let calls = 0 + const proxy = createUdpFramedProxy( + () => { + calls++ + return echoStream() + }, + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const client = await bindUdpClient() + t.teardown(() => client.close()) + + const send = (msg) => + new Promise((resolve) => { + client.once('message', (m) => resolve(m.toString())) + client.send(Buffer.from(msg), 0, msg.length, addr.port, '127.0.0.1') + }) + + t.is(await send('first'), 'first') + t.is(await send('second'), 'second') + t.is(calls, 1) + t.is(proxy.clients.size, 1) +}) + +test('createUdpFramedProxy - gives distinct clients independent tunnels and replies', async function (t) { + let calls = 0 + const proxy = createUdpFramedProxy( + () => { + calls++ + return echoStream() + }, + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const c1 = await bindUdpClient() + const c2 = await bindUdpClient() + t.teardown(() => { + c1.close() + c2.close() + }) + + const r1 = new Promise((resolve) => c1.once('message', (m) => resolve(m.toString()))) + const r2 = new Promise((resolve) => c2.once('message', (m) => resolve(m.toString()))) + c1.send(Buffer.from('from-c1'), 0, 7, addr.port, '127.0.0.1') + c2.send(Buffer.from('from-c2'), 0, 7, addr.port, '127.0.0.1') + + t.is(await r1, 'from-c1') + t.is(await r2, 'from-c2') + t.is(calls, 2) +}) + +test('createUdpFramedProxy - oversized reply frame destroys the tunnel and evicts the client', async function (t) { + let calls = 0 + const { logger, calls: logs } = spyLogger() + const proxy = createUdpFramedProxy( + () => { + calls++ + return echoStream() + }, + { port: 0, host: '127.0.0.1', maxFrameSize: 5, logger }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const client = await bindUdpClient() + t.teardown(() => client.close()) + + const oversized = Buffer.from('0123456789') // 10 bytes > maxFrameSize(5) + client.send(oversized, 0, oversized.length, addr.port, '127.0.0.1') + await tick(30) + + t.is(proxy.clients.size, 0, 'client entry removed after oversized frame') + t.ok(logs.error.some((args) => /Frame too large/.test(args[0]))) + + const small = Buffer.from('small') + client.send(small, 0, small.length, addr.port, '127.0.0.1') + await tick(30) + + t.is(calls, 2, 'a fresh tunnel is created for the same client after eviction') +}) + +test('createUdpFramedProxy - removes the client entry when its tunnel closes', async function (t) { + let tunnel + const proxy = createUdpFramedProxy( + () => { + tunnel = echoStream() + return tunnel + }, + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const client = await bindUdpClient() + t.teardown(() => client.close()) + + const payload = Buffer.from('hi') + client.send(payload, 0, payload.length, addr.port, '127.0.0.1') + await tick(30) + t.is(proxy.clients.size, 1) + + tunnel.emit('close') + await tick() + t.is(proxy.clients.size, 0) +}) + +test('createUdpFramedProxy - reassembles a reply frame split across multiple stream chunks', async function (t) { + const proxy = createUdpFramedProxy( + () => splitEchoStream(), + { port: 0, host: '127.0.0.1' }, + () => {} + ) + t.teardown(() => proxy.proxySocket.close()) + await tick() + + const addr = proxy.proxySocket.address() + const client = await bindUdpClient() + t.teardown(() => client.close()) + + const reply = new Promise((resolve) => client.once('message', (msg) => resolve(msg.toString()))) + const payload = Buffer.from('reassemble-me-please') + client.send(payload, 0, payload.length, addr.port, '127.0.0.1') + + t.is(await reply, 'reassemble-me-please') +}) diff --git a/test/udp-framed-server.js b/test/udp-framed-server.js new file mode 100644 index 0000000..3ed32fb --- /dev/null +++ b/test/udp-framed-server.js @@ -0,0 +1,145 @@ +const test = require('brittle') +const { pipeUdpFramedServer } = require('../lib/udp-framed.js') +const { fakeSocket, spyLogger, frame, tick, withEchoUdpServer } = require('./helpers.js') + +test('pipeUdpFramedServer - forwards a framed message to the local service and frames the reply back', async function (t) { + await withEchoUdpServer(async (addr) => { + const stream = fakeSocket() + pipeUdpFramedServer(stream, null, { port: addr.port, host: '127.0.0.1' }) + + const payload = Buffer.from('ping-udp-server') + stream.push(frame(payload)) + await tick(40) + + t.is(stream.written.length, 1) + const reply = stream.written[0] + const len = reply.readUInt32BE(0) + t.is(len, payload.length) + t.alike(reply.subarray(4, 4 + len), payload) + + stream.emit('end') + await tick() + }) +}) + +test('pipeUdpFramedServer - drains a complete frame already present in leftover', async function (t) { + await withEchoUdpServer(async (addr) => { + const stream = fakeSocket() + const payload = Buffer.from('leftover-complete') + pipeUdpFramedServer(stream, frame(payload), { port: addr.port, host: '127.0.0.1' }) + await tick(40) + + const reply = stream.written[0] + const len = reply.readUInt32BE(0) + t.alike(reply.subarray(4, 4 + len), payload) + + stream.emit('end') + await tick() + }) +}) + +test('pipeUdpFramedServer - reassembles a frame split across leftover and later stream data', async function (t) { + await withEchoUdpServer(async (addr) => { + const stream = fakeSocket() + const full = frame(Buffer.from('split-across-chunks')) + const partial = full.subarray(0, 6) // 4-byte length prefix + 2 payload bytes + const rest = full.subarray(6) + + pipeUdpFramedServer(stream, partial, { port: addr.port, host: '127.0.0.1' }) + await tick(20) + t.is(stream.written.length, 0, 'incomplete frame is not forwarded yet') + + stream.push(rest) + await tick(40) + + const reply = stream.written[0] + const len = reply.readUInt32BE(0) + t.alike(reply.subarray(4, 4 + len), Buffer.from('split-across-chunks')) + + stream.emit('end') + await tick() + }) +}) + +test('pipeUdpFramedServer - handles back-to-back frames delivered in a single chunk', async function (t) { + await withEchoUdpServer(async (addr) => { + const stream = fakeSocket() + pipeUdpFramedServer(stream, null, { port: addr.port, host: '127.0.0.1' }) + + const combined = Buffer.concat([frame(Buffer.from('one')), frame(Buffer.from('two'))]) + stream.push(combined) + await tick(40) + + const payloads = stream.written.map((buf) => { + const len = buf.readUInt32BE(0) + return buf.subarray(4, 4 + len).toString() + }) + t.alike(payloads.sort(), ['one', 'two']) + + stream.emit('end') + await tick() + }) +}) + +test('pipeUdpFramedServer - a frame over maxFrameSize destroys the stream with a descriptive error', async function (t) { + const stream = fakeSocket() + let err + stream.on('error', (e) => { + err = e + }) + + pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1', maxFrameSize: 4 }) + stream.push(frame(Buffer.from('too-big-payload'))) + await tick(30) + + t.ok(stream.destroyed) + t.ok(err && /Frame too large: 15 > 4/.test(err.message)) +}) + +test('pipeUdpFramedServer - a remote stream error destroys the pipe', async function (t) { + const stream = fakeSocket() + const { logger, calls } = spyLogger() + pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1', logger }) + + const closed = new Promise((resolve) => stream.once('close', resolve)) + stream.emit('error', new Error('peer blew up')) + await closed + + t.ok(stream.destroyed) + t.ok(calls.info.some((args) => /Destroying pipe due to error: peer blew up/.test(args[0]))) +}) + +test('pipeUdpFramedServer - ends its writable side (no error) when the remote stream ends', async function (t) { + const stream = fakeSocket() + pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1' }) + + const finished = new Promise((resolve) => stream.once('finish', resolve)) + stream.emit('end') + await finished + + t.absent(stream.destroyed, 'a clean end() does not destroy the stream') +}) + +test('pipeUdpFramedServer - cleanup only runs once no matter how many teardown events fire', async function (t) { + const stream = fakeSocket() + const { logger, calls } = spyLogger() + pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1', logger }) + + stream.emit('close') + stream.emit('close') + stream.emit('end') + await tick() + + t.is(calls.info.filter((args) => /Destroying pipe/.test(args[0])).length, 1) +}) + +test('pipeUdpFramedServer - works with the default no-op logger', async function (t) { + const stream = fakeSocket() + pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1' }) + + const finished = new Promise((resolve) => stream.once('finish', resolve)) + stream.emit('end') + await finished + + t.absent(stream.destroyed) +}) From 69b5d6f1f5e2d24a3426b11717e897e61fc4e8c5 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Mon, 3 Aug 2026 11:39:21 +0530 Subject: [PATCH 08/10] name helpers properly --- test/connection-piper.js | 12 ++++++------ test/helpers.js | 4 ++-- test/udp-framed-proxy.js | 4 ++-- test/udp-framed-server.js | 6 +++--- 4 files changed, 13 insertions(+), 13 deletions(-) diff --git a/test/connection-piper.js b/test/connection-piper.js index c69db45..d63d47e 100644 --- a/test/connection-piper.js +++ b/test/connection-piper.js @@ -1,7 +1,7 @@ const test = require('brittle') const net = require('net') const { connPiper } = require('../lib/tcp-piper.js') -const { spyLogger, tick, listenTcp, tcpPair, closePair, readOnce } = require('./helpers.js') +const { createLogger, tick, listenTcp, tcpPair, closePair, readOnce } = require('./helpers.js') test('connPiper - pipes data from a to b', async function (t) { const A = await tcpPair() @@ -36,7 +36,7 @@ test('connPiper - pipes data from b to a', async function (t) { test('connPiper - bFactory throwing destroys a and logs the error', async function (t) { const A = await tcpPair() t.teardown(() => closePair(A)) - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() connPiper( A.local, @@ -55,7 +55,7 @@ test('connPiper - bFactory throwing destroys a and logs the error', async functi test('connPiper - bFactory returning null destroys a and logs a warning', async function (t) { const A = await tcpPair() t.teardown(() => closePair(A)) - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() connPiper(A.local, () => null, { logger }) await tick() @@ -135,7 +135,7 @@ test('connPiper - logs debug messages describing byte counts in both directions' closePair(A) closePair(B) }) - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() connPiper(A.local, () => B.local, { logger }) A.remote.write('abc') @@ -154,7 +154,7 @@ test('connPiper - logs Connected when b emits a connect event', async function ( t.teardown(() => bServer.close()) const addr = await listenTcp(bServer) - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() const b = net.connect(addr.port, '127.0.0.1') t.teardown(() => b.destroy()) // net.connect() hasn't fired 'connect' yet here — connPiper's own listener @@ -174,7 +174,7 @@ test('connPiper - a ending its readable side ends b in turn', async function (t) closePair(A) closePair(B) }) - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() connPiper(A.local, () => B.local, { logger }) const bFinished = new Promise((resolve) => B.local.once('finish', resolve)) diff --git a/test/helpers.js b/test/helpers.js index 39e7a8a..9047bdb 100644 --- a/test/helpers.js +++ b/test/helpers.js @@ -42,7 +42,7 @@ function splitEchoStream() { }) } -function spyLogger() { +function createLogger() { const calls = { debug: [], info: [], warn: [], error: [] } const logger = { debug: (...a) => calls.debug.push(a), @@ -120,7 +120,7 @@ module.exports = { fakeSocket, echoStream, splitEchoStream, - spyLogger, + createLogger, frame, tick, bindUdpClient, diff --git a/test/udp-framed-proxy.js b/test/udp-framed-proxy.js index 64ffe6f..abb1bb3 100644 --- a/test/udp-framed-proxy.js +++ b/test/udp-framed-proxy.js @@ -1,6 +1,6 @@ const test = require('brittle') const { createUdpFramedProxy } = require('../lib/udp-framed.js') -const { echoStream, splitEchoStream, spyLogger, tick, bindUdpClient } = require('./helpers.js') +const { echoStream, splitEchoStream, createLogger, tick, bindUdpClient } = require('./helpers.js') test('createUdpFramedProxy - round-trips a datagram through the tunnel factory', async function (t) { const proxy = createUdpFramedProxy( @@ -84,7 +84,7 @@ test('createUdpFramedProxy - gives distinct clients independent tunnels and repl test('createUdpFramedProxy - oversized reply frame destroys the tunnel and evicts the client', async function (t) { let calls = 0 - const { logger, calls: logs } = spyLogger() + const { logger, calls: logs } = createLogger() const proxy = createUdpFramedProxy( () => { calls++ diff --git a/test/udp-framed-server.js b/test/udp-framed-server.js index 3ed32fb..15dc8b4 100644 --- a/test/udp-framed-server.js +++ b/test/udp-framed-server.js @@ -1,6 +1,6 @@ const test = require('brittle') const { pipeUdpFramedServer } = require('../lib/udp-framed.js') -const { fakeSocket, spyLogger, frame, tick, withEchoUdpServer } = require('./helpers.js') +const { fakeSocket, createLogger, frame, tick, withEchoUdpServer } = require('./helpers.js') test('pipeUdpFramedServer - forwards a framed message to the local service and frames the reply back', async function (t) { await withEchoUdpServer(async (addr) => { @@ -98,7 +98,7 @@ test('pipeUdpFramedServer - a frame over maxFrameSize destroys the stream with a test('pipeUdpFramedServer - a remote stream error destroys the pipe', async function (t) { const stream = fakeSocket() - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1', logger }) const closed = new Promise((resolve) => stream.once('close', resolve)) @@ -122,7 +122,7 @@ test('pipeUdpFramedServer - ends its writable side (no error) when the remote st test('pipeUdpFramedServer - cleanup only runs once no matter how many teardown events fire', async function (t) { const stream = fakeSocket() - const { logger, calls } = spyLogger() + const { logger, calls } = createLogger() pipeUdpFramedServer(stream, null, { port: 1, host: '127.0.0.1', logger }) stream.emit('close') From 0485fd14de98a31e4515628e64981b2832e4028a Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Tue, 4 Aug 2026 11:31:34 +0530 Subject: [PATCH 09/10] fix leaking socket --- .gitignore | 2 +- README.md | 120 ++++++++++++++++++++++++++++++++++++--- test/connection-piper.js | 2 +- test/tcp-proxy.js | 50 +++++++++------- 4 files changed, 141 insertions(+), 33 deletions(-) diff --git a/.gitignore b/.gitignore index c3179ee..07ae377 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,4 @@ package-lock.json node_modules .idea -,DS_Store \ No newline at end of file +.DS_Store diff --git a/README.md b/README.md index a9823a2..267365a 100644 --- a/README.md +++ b/README.md @@ -1,19 +1,121 @@ -# Hyper Cmd Lib Net +# hyper-cmd-lib-net -Network library to pipe local and hypderdht udp/tcp streams. +Pipe local and remote (hyperdht) TCP/UDP streams together. -## Install +``` +npm i @holesail/hyper-cmd-lib-net +``` -`npm i @holesail/hyper-cmd-lib-net` +## Usage + +```js +const { createTcpProxy } = require('@holesail/hyper-cmd-lib-net') + +const proxy = createTcpProxy( + () => tunnel.createStream(), + { port: 8080, host: '127.0.0.1' }, + () => console.log('listening on', proxy.address()) +) +``` ## API -#### `connPiper` +#### `connPiper(a, bFactory, [opts])` + +Pipe two duplex streams together, and destroy both the moment either one closes or +errors. + +`bFactory()` is called once to make the second stream. Throw, or return `null`, to +reject the connection — `a` gets destroyed and `bFactory` is never called again. + +```js +connPiper(a, () => b, { + logger: null, // { debug, info, warn, error }, defaults to a no-op logger + onDestroy: (err) => {} // called once, when the pipe is torn down +}) +``` + +#### `const proxy = createTcpProxy(remoteTunnel, opts, onListen)` + +Start a TCP server. Every accepted connection is piped to a fresh `remoteTunnel()` +stream using `connPiper` above. + +```js +createTcpProxy( + () => tunnel.createStream(), + { + port: 0, + host: '127.0.0.1', + logger: null, + onDestroy: (err) => {} + }, + () => {} +) +``` + +Returns the underlying `net.Server`. `proxy.close()` also destroys any connections +still open, so it doesn't need `'close'` listeners on every socket to fully shut down. + +#### `pipeTcpServer(remoteStream, leftover, opts)` + +The other direction of `createTcpProxy` — connects out to a local TCP service and pipes +it to `remoteStream`. + +```js +pipeTcpServer(remoteStream, leftoverBuffer, { + port: 3000, + host: '127.0.0.1' +}) +``` + +`leftover` is a buffer of bytes already read off `remoteStream` (eg while sniffing a +protocol) that gets written to the local socket before piping starts. Pass `null` if +there isn't any. + +#### UDP framing + +UDP has no stream boundaries, so the two helpers below speak a small length-prefixed +framing over the tunnel instead: 4 bytes big-endian length, then that many bytes of +payload. They're meant to be run as a pair, one on each end of a tunnel. + +#### `const { proxySocket, clients } = createUdpFramedProxy(createTunnel, opts, onBind)` + +Bind a UDP socket. Every distinct sender (by `address:port`) gets its own tunnel +stream, made lazily via `createTunnel()` on its first packet. + +```js +createUdpFramedProxy( + () => tunnel.createStream(), + { + port: 0, + host: '127.0.0.1', + maxFrameSize: 65535, + logger: null + }, + () => {} +) +``` + +`clients` is a `Map` of `address:port` -> `{ stream, rinfo, buffer }`. A frame bigger +than `maxFrameSize` destroys that client's tunnel and evicts it — the next packet from +the same client just starts a new one. + +#### `pipeUdpFramedServer(stream, leftover, opts)` + +The other end of the pair. Unframes packets off `stream` and forwards them as plain +UDP to a local service, framing replies back onto `stream`. -#### `createTcpProxy` +```js +pipeUdpFramedServer(remoteStream, leftoverBuffer, { + port: 53, + host: '127.0.0.1', + maxFrameSize: 65535 +}) +``` -#### `pipeTcpServer` +A clean `end` on `stream` ends the local socket, no drama. An `error`, or a frame over +`maxFrameSize`, destroys it. -#### `createUdpFramedProxy` +## License -#### `pipeUdpFramedServer` +Apache-2.0 diff --git a/test/connection-piper.js b/test/connection-piper.js index d63d47e..2ca7cb3 100644 --- a/test/connection-piper.js +++ b/test/connection-piper.js @@ -150,7 +150,7 @@ test('connPiper - logs Connected when b emits a connect event', async function ( const A = await tcpPair() t.teardown(() => closePair(A)) - const bServer = net.createServer({ allowHalfOpen: true }) + const bServer = net.createServer() t.teardown(() => bServer.close()) const addr = await listenTcp(bServer) diff --git a/test/tcp-proxy.js b/test/tcp-proxy.js index 3ee994b..067828d 100644 --- a/test/tcp-proxy.js +++ b/test/tcp-proxy.js @@ -1,7 +1,7 @@ const test = require('brittle') const net = require('net') const { createTcpProxy } = require('../lib/tcp-piper.js') -const { echoStream, tick, listenTcp } = require('./helpers.js') +const { echoStream, tick } = require('./helpers.js') test('createTcpProxy - proxies a real TCP client round trip through the tunnel factory', async function (t) { const addr = await new Promise((resolve) => { @@ -40,17 +40,19 @@ test('createTcpProxy - calls onListen once the server is bound', async function test('createTcpProxy - invokes the tunnel factory once per incoming connection', async function (t) { let calls = 0 - const proxy = createTcpProxy( - () => { - calls++ - return echoStream() - }, - { port: 0, host: '127.0.0.1' }, - () => {} - ) + let proxy + const addr = await new Promise((resolve) => { + proxy = createTcpProxy( + () => { + calls++ + return echoStream() + }, + { port: 0, host: '127.0.0.1' }, + () => resolve(proxy.address()) + ) + }) t.teardown(() => proxy.close()) - const addr = await listenTcp(proxy) const c1 = net.connect(addr.port, '127.0.0.1') const c2 = net.connect(addr.port, '127.0.0.1') t.teardown(() => { @@ -67,14 +69,16 @@ test('createTcpProxy - invokes the tunnel factory once per incoming connection', }) test('createTcpProxy - two connections stay isolated from each other', async function (t) { - const proxy = createTcpProxy( - () => echoStream(), - { port: 0, host: '127.0.0.1' }, - () => {} - ) + let proxy + const addr = await new Promise((resolve) => { + proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => resolve(proxy.address()) + ) + }) t.teardown(() => proxy.close()) - const addr = await listenTcp(proxy) const c1 = net.connect(addr.port, '127.0.0.1') const c2 = net.connect(addr.port, '127.0.0.1') t.teardown(() => { @@ -96,13 +100,15 @@ test('createTcpProxy - two connections stay isolated from each other', async fun }) test('createTcpProxy - closing the proxy destroys any still-open connections', async function (t) { - const proxy = createTcpProxy( - () => echoStream(), - { port: 0, host: '127.0.0.1' }, - () => {} - ) + let proxy + const addr = await new Promise((resolve) => { + proxy = createTcpProxy( + () => echoStream(), + { port: 0, host: '127.0.0.1' }, + () => resolve(proxy.address()) + ) + }) - const addr = await listenTcp(proxy) const client = net.connect(addr.port, '127.0.0.1') client.on('error', () => {}) // the far end destroys mid-flight; ECONNRESET is expected here t.teardown(() => client.destroy()) From ce785621d3de5f8952e8cdd2492e066fc30990e7 Mon Sep 17 00:00:00 2001 From: supersuryaansh Date: Tue, 4 Aug 2026 18:20:40 +0530 Subject: [PATCH 10/10] bump bare net and tcp --- package.json | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/package.json b/package.json index 79eae22..9cfa6c3 100644 --- a/package.json +++ b/package.json @@ -30,9 +30,9 @@ "homepage": "https://github.com/holesail/hyper-cmd-lib-net/", "dependencies": { "bare-dgram": "^1.0.1", - "bare-net": "^2.3.2", - "bare-tcp": "^2.5.3", - "net": "npm:bare-net@^2.3.2" + "bare-net": "^2.3.3", + "bare-tcp": "^2.5.4", + "net": "npm:bare-net@^2.3.3" }, "devDependencies": { "brittle": "^4.0.0",