diff --git a/lib/discovery.js b/lib/discovery.js new file mode 100644 index 0000000..221bc38 --- /dev/null +++ b/lib/discovery.js @@ -0,0 +1,265 @@ +'use strict'; + +/** + * Discovery — pluggable peer discovery for SYM mesh nodes. + * + * Three implementations: + * - BonjourDiscovery: LAN discovery via dns-sd (macOS) or bonjour-service (Linux) + * - NullDiscovery: no-op for relay-only nodes and tests + * + * SymNode accepts a discovery instance via opts.discovery. If not provided, + * it creates BonjourDiscovery (default) or NullDiscovery (if relayOnly). + * + * See MMP v0.2.0 Section 5 (Connection, Layer 2). + * + * Copyright (c) 2026 SYM.BOT Ltd. Apache 2.0 License. + */ + +const net = require('net'); +const { EventEmitter } = require('events'); +const { TcpTransport } = require('./transport'); + +// ── Interface ──────────────────────────────────────────────── + +/** + * Base discovery class. Subclasses implement start/stop. + * Emits: + * 'peer-found' (address, port, peerId, peerName) — outbound connection needed + * 'inbound-connection' (transport, peerId, peerName) — peer connected to us + * 'error' (err) — non-fatal discovery error + */ +class Discovery extends EventEmitter { + /** + * Start discovery: listen for inbound connections + advertise + browse for peers. + * @param {object} identity — { nodeId, name, publicKey, hostname } + * @param {function} log — logging function + * @returns {Promise} listening port (0 if no listener) + */ + async start(identity, log) { return 0; } + + /** + * Stop discovery: close listener, stop browsing, clean up. + * @returns {Promise} + */ + async stop() {} +} + +// ── Bonjour Discovery ──────────────────────────────────────── + +/** + * LAN discovery via TCP server + Bonjour/mDNS. + * Uses system dns-sd on macOS, falls back to bonjour-service on Linux. + */ +class BonjourDiscovery extends Discovery { + /** + * @param {object} [opts] + * @param {boolean} [opts.mdns=true] — enable mDNS advertisement/browsing (set false for server-only mode) + */ + constructor(opts = {}) { + super(); + this._mdnsEnabled = opts.mdns !== false; + this._server = null; + this._dnssdRegister = null; + this._dnssdBrowse = null; + this._bonjour = null; + this._browser = null; + this._port = 0; + this._identity = null; + this._log = () => {}; + } + + async start(identity, log) { + this._identity = identity; + this._log = log || (() => {}); + + // Start TCP server + await this._startServer(); + + // Start Bonjour advertisement + browsing (skip in server-only mode) + if (this._mdnsEnabled) { + this._startDnsSd(); + } + + return this._port; + } + + async stop() { + // Kill dns-sd processes + if (this._dnssdRegister) { + try { this._dnssdRegister.kill(); } catch {} + this._dnssdRegister = null; + } + if (this._dnssdBrowse) { + try { this._dnssdBrowse.kill(); } catch {} + this._dnssdBrowse = null; + } + + // Stop bonjour-service fallback + if (this._browser) { + try { this._browser.stop(); } catch {} + this._browser = null; + } + if (this._bonjour) { + try { this._bonjour.destroy(); } catch {} + this._bonjour = null; + } + + // Close TCP server + if (this._server) { + await new Promise((resolve) => { + this._server.close(() => resolve()); + setTimeout(resolve, 1000); + }); + this._server = null; + } + } + + _startServer() { + return new Promise((resolve, reject) => { + this._server = net.createServer((socket) => { + this._handleInboundConnection(socket); + }); + this._server.on('error', (err) => { + this._log(`Server error: ${err.message}`); + reject(err); + }); + this._server.listen(0, '0.0.0.0', () => { + this._port = this._server.address().port; + resolve(); + }); + }); + } + + _handleInboundConnection(socket) { + const transport = new TcpTransport(socket); + let identified = false; + const timeout = setTimeout(() => { if (!identified) transport.close(); }, 10000); + + transport.on('message', (msg) => { + if (identified) return; + if (msg.type !== 'handshake') { transport.close(); return; } + identified = true; + clearTimeout(timeout); + + transport.removeAllListeners('message'); + this.emit('inbound-connection', transport, msg.nodeId, msg.name, msg); + }); + + transport.on('error', () => clearTimeout(timeout)); + } + + _startDnsSd() { + const { spawn } = require('child_process'); + const identity = this._identity; + + const txtParts = [ + `node-id=${identity.nodeId}`, + `node-name=${identity.name}`, + `public-key=${identity.publicKey}`, + `hostname=${identity.hostname}`, + ]; + this._dnssdRegister = spawn('dns-sd', [ + '-R', identity.nodeId, '_sym._tcp', 'local.', + String(this._port), ...txtParts, + ], { stdio: 'ignore' }); + + this._dnssdRegister.on('error', (err) => { + this._log(`dns-sd not available, falling back to bonjour-service: ${err.message}`); + this._dnssdRegister = null; + if (this._dnssdBrowse) { + try { this._dnssdBrowse.kill(); } catch {} + this._dnssdBrowse = null; + } + this._startBonjourFallback(); + }); + + this._dnssdBrowse = spawn('dns-sd', ['-B', '_sym._tcp'], { stdio: ['ignore', 'pipe', 'ignore'] }); + this._dnssdBrowse.on('error', () => { this._dnssdBrowse = null; }); + + let browseBuffer = ''; + this._dnssdBrowse.stdout.on('data', (data) => { + browseBuffer += data.toString(); + let idx; + while ((idx = browseBuffer.indexOf('\n')) !== -1) { + const line = browseBuffer.slice(0, idx).trim(); + browseBuffer = browseBuffer.slice(idx + 1); + const match = line.match(/\s+Add\s+\d+\s+\d+\s+\S+\s+_sym\._tcp\.\s+(.+)$/); + if (match) { + const instanceName = match[1].trim(); + if (instanceName === identity.nodeId) continue; + this._resolvePeer(instanceName); + } + } + }); + } + + _resolvePeer(instanceName) { + const { spawn } = require('child_process'); + const identity = this._identity; + const resolve = spawn('dns-sd', ['-L', instanceName, '_sym._tcp', 'local.'], { stdio: ['ignore', 'pipe', 'ignore'] }); + let resolveBuffer = ''; + const timeout = setTimeout(() => resolve.kill(), 5000); + + resolve.stdout.on('data', (data) => { + resolveBuffer += data.toString(); + const match = resolveBuffer.match(/can be reached at (.+?):(\d+)/); + if (!match) return; + + clearTimeout(timeout); + const host = match[1]; + const port = parseInt(match[2]); + + const nodeIdMatch = resolveBuffer.match(/node-id=(\S+)/); + const nodeNameMatch = resolveBuffer.match(/node-name=(\S+)/); + const peerId = nodeIdMatch ? nodeIdMatch[1] : instanceName; + const peerName = nodeNameMatch ? nodeNameMatch[1] : 'unknown'; + + resolve.kill(); + + if (peerId === identity.nodeId) return; + if (identity.nodeId < peerId) { + this.emit('peer-found', host, port, peerId, peerName); + } + }); + } + + _startBonjourFallback() { + const { Bonjour } = require('bonjour-service'); + const identity = this._identity; + this._bonjour = new Bonjour(); + + this._bonjour.publish({ + name: identity.nodeId, + type: 'sym', + port: this._port, + txt: { 'node-id': identity.nodeId, 'node-name': identity.name, 'public-key': identity.publicKey, 'hostname': identity.hostname }, + }); + + this._browser = this._bonjour.find({ type: 'sym' }); + + this._browser.on('up', (service) => { + const peerId = service.txt?.['node-id']; + if (!peerId || peerId === identity.nodeId) return; + const peerName = service.txt?.['node-name'] || 'unknown'; + const address = service.referer?.address || service.addresses?.[0]; + const port = service.port; + if (!address || !port) return; + if (identity.nodeId < peerId) { + this.emit('peer-found', address, port, peerId, peerName); + } + }); + } +} + +// ── Null Discovery ─────────────────────────────────────────── + +/** + * No-op discovery for relay-only nodes and testing. + * No TCP server, no Bonjour, no child processes. + */ +class NullDiscovery extends Discovery { + async start() { return 0; } + async stop() {} +} + +module.exports = { Discovery, BonjourDiscovery, NullDiscovery }; diff --git a/lib/node.js b/lib/node.js index 8d9e03a..7455b1e 100644 --- a/lib/node.js +++ b/lib/node.js @@ -30,6 +30,7 @@ const { } = require('@sym-bot/core'); const { TcpTransport } = require('./transport'); const { RelayConnection } = require('./relay'); +const { BonjourDiscovery, NullDiscovery } = require('./discovery'); class SymNode extends EventEmitter { /** @@ -108,12 +109,12 @@ class SymNode extends EventEmitter { // Peer state this._peers = new Map(); - this._server = null; - this._bonjour = null; - this._browser = null; this._port = 0; this._running = false; + // Discovery — pluggable for testability. See MMP v0.2.0 Section 5. + this._discovery = opts.discovery || (opts.relayOnly ? new NullDiscovery() : new BonjourDiscovery()); + // Wake this._wakeChannel = opts.wakeChannel || null; this._peerWakeChannels = new Map(); @@ -337,10 +338,27 @@ class SymNode extends EventEmitter { if (this._running) return; this._running = true; - if (!this._relayOnly) { - await this._startServer(); - this._startDiscovery(); - } + // Wire discovery events to peer management + this._discovery.on('peer-found', (address, port, peerId, peerName) => { + if (!this._peers.has(peerId)) { + this._connectToPeer(address, port, peerId, peerName); + } + }); + this._discovery.on('inbound-connection', (transport, peerId, peerName, handshakeMsg) => { + if (this._peers.has(peerId)) { transport.close(); return; } + if (handshakeMsg.e2ePublicKey) { + this._deriveAndStoreSecret(peerId, handshakeMsg.e2ePublicKey); + } + transport.on('message', (m) => { + const peer = this._peers.get(peerId); + if (peer) peer.lastSeen = Date.now(); + this._frameHandler.handle(peerId, peerName, m); + }); + const peer = this._createPeer(transport, peerId, peerName, false, 'bonjour'); + this._addPeer(peer); + }); + + this._port = await this._discovery.start(this._identity, (msg) => this._log(msg)); if (this._relayUrl) { this._relay.connect(); @@ -375,31 +393,8 @@ class SymNode extends EventEmitter { } this._peers.clear(); - if (this._dnssdRegister) { - try { this._dnssdRegister.kill(); } catch {} - this._dnssdRegister = null; - } - if (this._dnssdBrowse) { - try { this._dnssdBrowse.kill(); } catch {} - this._dnssdBrowse = null; - } - if (this._browser) { - try { this._browser.stop(); } catch {} - this._browser = null; - } - if (this._bonjour) { - try { this._bonjour.destroy(); } catch {} - this._bonjour = null; - } - - if (this._server) { - await new Promise((resolve) => { - this._server.close(() => resolve()); - // Timeout in case server.close() hangs - setTimeout(resolve, 1000); - }); - this._server = null; - } + await this._discovery.stop(); + this._discovery.removeAllListeners(); this._log('Stopped'); } @@ -662,168 +657,6 @@ class SymNode extends EventEmitter { }; } - // ── TCP Server ───────────────────────────────────────────── - - _startServer() { - return new Promise((resolve, reject) => { - this._server = net.createServer((socket) => { - this._handleInboundConnection(socket); - }); - this._server.on('error', (err) => { - this._log(`Server error: ${err.message}`); - reject(err); - }); - this._server.listen(0, '0.0.0.0', () => { - this._port = this._server.address().port; - resolve(); - }); - }); - } - - _handleInboundConnection(socket) { - const transport = new TcpTransport(socket); - let identified = false; - const timeout = setTimeout(() => { if (!identified) transport.close(); }, 10000); - - transport.on('message', (msg) => { - if (identified) return; - if (msg.type !== 'handshake') { transport.close(); return; } - identified = true; - clearTimeout(timeout); - if (this._peers.has(msg.nodeId)) { transport.close(); return; } - - // Derive E2E shared secret if peer sent public key - if (msg.e2ePublicKey) { - this._deriveAndStoreSecret(msg.nodeId, msg.e2ePublicKey); - } - - transport.removeAllListeners('message'); - transport.on('message', (m) => { - const peer = this._peers.get(msg.nodeId); - if (peer) peer.lastSeen = Date.now(); - this._frameHandler.handle(msg.nodeId, msg.name, m); - }); - - const peer = this._createPeer(transport, msg.nodeId, msg.name, false, 'bonjour'); - this._addPeer(peer); - }); - - transport.on('error', () => clearTimeout(timeout)); - } - - // ── Bonjour Discovery ────────────────────────────────────── - - _startDiscovery() { - const { spawn } = require('child_process'); - - // Register via system dns-sd (Apple's native mDNS responder). - // This ensures NWConnection on iOS can resolve the service endpoint. - // The JavaScript bonjour-service library uses its own multicast DNS - // which Apple's Network framework cannot resolve. - const txtParts = [ - `node-id=${this._identity.nodeId}`, - `node-name=${this.name}`, - `public-key=${this._identity.publicKey}`, - `hostname=${this._identity.hostname}`, - ]; - this._dnssdRegister = spawn('dns-sd', [ - '-R', this._identity.nodeId, '_sym._tcp', 'local.', - String(this._port), ...txtParts, - ], { stdio: 'ignore' }); - - this._dnssdRegister.on('error', (err) => { - // dns-sd not available (Linux, containers) — fallback to JavaScript mDNS - this._log(`dns-sd not available, falling back to bonjour-service: ${err.message}`); - this._dnssdRegister = null; - // Kill browse process too — it won't work without dns-sd - if (this._dnssdBrowse) { - try { this._dnssdBrowse.kill(); } catch {} - this._dnssdBrowse = null; - } - this._startBonjourFallback(); - }); - - // Browse for peers via dns-sd (only useful if dns-sd is available) - this._dnssdBrowse = spawn('dns-sd', ['-B', '_sym._tcp'], { stdio: ['ignore', 'pipe', 'ignore'] }); - this._dnssdBrowse.on('error', () => { - // Handled by register error above — just prevent unhandled error crash - this._dnssdBrowse = null; - }); - let browseBuffer = ''; - this._dnssdBrowse.stdout.on('data', (data) => { - browseBuffer += data.toString(); - let idx; - while ((idx = browseBuffer.indexOf('\n')) !== -1) { - const line = browseBuffer.slice(0, idx).trim(); - browseBuffer = browseBuffer.slice(idx + 1); - // Parse dns-sd -B output: "Timestamp A/R Flags if Domain Service Type Instance Name" - const match = line.match(/\s+Add\s+\d+\s+\d+\s+\S+\s+_sym\._tcp\.\s+(.+)$/); - if (match) { - const instanceName = match[1].trim(); - if (instanceName === this._identity.nodeId) continue; // self - this._resolvePeer(instanceName); - } - } - }); - } - - _resolvePeer(instanceName) { - const { spawn } = require('child_process'); - const resolve = spawn('dns-sd', ['-L', instanceName, '_sym._tcp', 'local.'], { stdio: ['ignore', 'pipe', 'ignore'] }); - let resolveBuffer = ''; - const timeout = setTimeout(() => resolve.kill(), 5000); - - resolve.stdout.on('data', (data) => { - resolveBuffer += data.toString(); - // Parse: "instance._sym._tcp.local. can be reached at hostname:port (interface N)" - const match = resolveBuffer.match(/can be reached at (.+?):(\d+)/); - if (!match) return; - - clearTimeout(timeout); - const host = match[1]; - const port = parseInt(match[2]); - - // Extract TXT record values - const nodeIdMatch = resolveBuffer.match(/node-id=(\S+)/); - const nodeNameMatch = resolveBuffer.match(/node-name=(\S+)/); - const peerId = nodeIdMatch ? nodeIdMatch[1] : instanceName; - const peerName = nodeNameMatch ? nodeNameMatch[1] : 'unknown'; - - resolve.kill(); - - if (peerId === this._identity.nodeId) return; - if (this._identity.nodeId < peerId && !this._peers.has(peerId)) { - this._connectToPeer(host, port, peerId, peerName); - } - }); - } - - _startBonjourFallback() { - const { Bonjour } = require('bonjour-service'); - this._bonjour = new Bonjour(); - - this._bonjour.publish({ - name: this._identity.nodeId, - type: 'sym', - port: this._port, - txt: { 'node-id': this._identity.nodeId, 'node-name': this.name, 'public-key': this._identity.publicKey, 'hostname': this._identity.hostname }, - }); - - this._browser = this._bonjour.find({ type: 'sym' }); - - this._browser.on('up', (service) => { - const peerId = service.txt?.['node-id']; - if (!peerId || peerId === this._identity.nodeId) return; - const peerName = service.txt?.['node-name'] || 'unknown'; - const address = service.referer?.address || service.addresses?.[0]; - const port = service.port; - if (!address || !port) return; - if (this._identity.nodeId < peerId && !this._peers.has(peerId)) { - this._connectToPeer(address, port, peerId, peerName); - } - }); - } - _connectToPeer(address, port, peerId, peerName) { if (this._peers.has(peerId)) return; const socket = net.createConnection({ host: address, port }, () => { diff --git a/tests/discovery.test.js b/tests/discovery.test.js new file mode 100644 index 0000000..b63e1d3 --- /dev/null +++ b/tests/discovery.test.js @@ -0,0 +1,105 @@ +'use strict'; + +const { describe, it } = require('node:test'); +const assert = require('node:assert'); +const { Discovery, BonjourDiscovery, NullDiscovery } = require('../lib/discovery'); + +describe('NullDiscovery', () => { + it('should return port 0', async () => { + const d = new NullDiscovery(); + const port = await d.start({}, () => {}); + assert.strictEqual(port, 0); + }); + + it('should stop without error', async () => { + const d = new NullDiscovery(); + await d.start({}, () => {}); + await d.stop(); // should not throw + }); + + it('should be an EventEmitter', () => { + const d = new NullDiscovery(); + assert.ok(typeof d.on === 'function'); + assert.ok(typeof d.emit === 'function'); + }); +}); + +describe('BonjourDiscovery', () => { + it('should be constructable', () => { + const d = new BonjourDiscovery({ mdns: false }); + assert.ok(d instanceof Discovery); + }); + + it('should start a TCP server and return a port', async () => { + const d = new BonjourDiscovery({ mdns: false }); + const identity = { nodeId: 'test-id', name: 'test', publicKey: 'pk', hostname: 'host' }; + const port = await d.start(identity, () => {}); + assert.ok(port > 0, `should get a real port, got ${port}`); + await d.stop(); + }); + + it('should emit inbound-connection on valid handshake', async () => { + const d = new BonjourDiscovery({ mdns: false }); + const identity = { nodeId: 'test-id', name: 'test', publicKey: 'pk', hostname: 'host' }; + const port = await d.start(identity, () => {}); + + const connections = []; + d.on('inbound-connection', (transport, peerId, peerName) => { + connections.push({ peerId, peerName }); + }); + + // Connect and send handshake + const net = require('net'); + const { sendFrame } = require('../lib/frame-parser'); + const client = net.createConnection({ host: '127.0.0.1', port }, () => { + sendFrame(client, { type: 'handshake', nodeId: 'peer-abc', name: 'peer-node' }); + }); + + // Wait for the event + await new Promise(resolve => setTimeout(resolve, 100)); + + assert.strictEqual(connections.length, 1); + assert.strictEqual(connections[0].peerId, 'peer-abc'); + assert.strictEqual(connections[0].peerName, 'peer-node'); + + client.destroy(); + await d.stop(); + }); + + it('should reject non-handshake first frames', async () => { + const d = new BonjourDiscovery({ mdns: false }); + const identity = { nodeId: 'test-id', name: 'test', publicKey: 'pk', hostname: 'host' }; + const port = await d.start(identity, () => {}); + + const connections = []; + d.on('inbound-connection', () => connections.push(true)); + + const net = require('net'); + const { sendFrame } = require('../lib/frame-parser'); + const client = net.createConnection({ host: '127.0.0.1', port }, () => { + sendFrame(client, { type: 'ping' }); // not a handshake + }); + + await new Promise(resolve => setTimeout(resolve, 100)); + assert.strictEqual(connections.length, 0, 'should not accept non-handshake'); + + client.destroy(); + await d.stop(); + }); + + it('should stop cleanly', async () => { + const d = new BonjourDiscovery({ mdns: false }); + const identity = { nodeId: 'test-id', name: 'test', publicKey: 'pk', hostname: 'host' }; + await d.start(identity, () => {}); + await d.stop(); + await d.stop(); // double stop should not throw + }); +}); + +describe('Discovery base class', () => { + it('should have start and stop methods', () => { + const d = new Discovery(); + assert.ok(typeof d.start === 'function'); + assert.ok(typeof d.stop === 'function'); + }); +}); diff --git a/tests/node.test.js b/tests/node.test.js index a5170b2..fe8d823 100644 --- a/tests/node.test.js +++ b/tests/node.test.js @@ -4,6 +4,7 @@ const { describe, it, after } = require('node:test'); const assert = require('node:assert'); const fs = require('fs'); const { nodeDir } = require('../lib/config'); +const { NullDiscovery } = require('../lib/discovery'); // Use unique names to avoid state conflicts const nodeName = `test-node-${Date.now()}-${Math.random().toString(36).slice(2, 6)}`; @@ -11,7 +12,6 @@ const nodeName = `test-node-${Date.now()}-${Math.random().toString(36).slice(2, describe('SymNode', () => { let SymNode; - // Dynamic require to avoid top-level import issues it('should load SymNode', () => { SymNode = require('../lib/node').SymNode; assert.ok(SymNode, 'SymNode should be exported'); @@ -28,12 +28,12 @@ describe('SymNode', () => { assert.strictEqual(node.nodeId.length, 36, 'nodeId should be full UUID'); }); - // Note: lifecycle tests use relayOnly to avoid bonjour-service UDP socket leak - // (known issue: bonjour-service .destroy() doesn't close its multicast socket, - // preventing clean Node.js process exit). Bonjour path is tested in local dev (macOS). + // Lifecycle tests inject NullDiscovery — no TCP server, no Bonjour, no child processes. + // This tests the node's business logic in isolation from networking. + // Bonjour integration is validated in local dev and e2e tests. it('should return full nodeId in status()', async () => { - const node = new SymNode({ name: nodeName, silent: true, relayOnly: true }); + const node = new SymNode({ name: nodeName, silent: true, discovery: new NullDiscovery() }); await node.start(); const s = node.status(); assert.strictEqual(s.nodeId, node.nodeId); @@ -46,13 +46,19 @@ describe('SymNode', () => { it('should start and stop without error', async () => { const name = `test-lifecycle-${Date.now()}`; - const node = new SymNode({ name, silent: true, relayOnly: true }); + const node = new SymNode({ name, silent: true, discovery: new NullDiscovery() }); await node.start(); assert.strictEqual(node.status().running, true); await node.stop(); fs.rmSync(nodeDir(name), { recursive: true, force: true }); }); + it('should use NullDiscovery when relayOnly is true', () => { + const node = new SymNode({ name: nodeName, silent: true, relayOnly: true }); + // relayOnly creates NullDiscovery internally — no server, no discovery + assert.ok(node._discovery instanceof NullDiscovery); + }); + it('should return empty peers when no connections', () => { const node = new SymNode({ name: nodeName, silent: true }); const peers = node.peers(); @@ -68,7 +74,7 @@ describe('SymNode', () => { it('should remember and recall', async () => { const name = `test-memory-${Date.now()}`; - const node = new SymNode({ name, silent: true, relayOnly: true }); + const node = new SymNode({ name, silent: true, discovery: new NullDiscovery() }); await node.start(); const entry = node.remember({