|
| 1 | +// De-risk gossip mode on a FULL mesh: 3 peers all connected; the source sends a |
| 2 | +// block to ONE neighbor only; gossip forwarding carries it to the third peer |
| 3 | +// (the one the source never sent to). Proves the ?gossip=1 browser demo. |
| 4 | +import { readFile } from 'node:fs/promises'; |
| 5 | +import { WebSocket } from 'ws'; |
| 6 | +import nodeDataChannel from 'node-datachannel'; |
| 7 | +import { Codec } from './engine/codec/codec.js'; |
| 8 | +import { P2pEngine } from './engine/codec/p2p.js'; |
| 9 | +import { startSignaling } from './signaling-stub.mjs'; |
| 10 | + |
| 11 | +const D = new URL('./', import.meta.url); |
| 12 | +const jl = async (n) => JSON.parse(await readFile(new URL(`engine/schema/${n}.jsonld`, D), 'utf8')); |
| 13 | +const PORT = 9091, ROOM = 'b17c0192abad1deacafe', SIGNAL = `ws://localhost:${PORT}/.webrtc`, ICE = ['stun:stun.l.google.com:19302']; |
| 14 | + |
| 15 | +const codec = new Codec(await jl('core'), await jl('proof'), await jl('p2p')); |
| 16 | +const p2p = P2pEngine.fromSchemas(codec, await jl('p2p'), await jl('chain'), 'btc:testnet4'); |
| 17 | +const blockHex = (await readFile(new URL('data/block-26000.hex', D), 'utf8')).trim(); |
| 18 | +const expectHash = codec.blockHash(codec.decode('Block', blockHex).header); |
| 19 | +const log = (m) => console.log(m); |
| 20 | +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); |
| 21 | +const gather = (pc) => new Promise((res) => { if (pc.gatheringState() === 'complete') return res(pc.localDescription()); pc.onGatheringStateChange((s) => { if (s === 'complete') res(pc.localDescription()); }); }); |
| 22 | + |
| 23 | +// Full-mesh node peer WITH gossip forwarding (mirror of peer-mesh.js + onMessage). |
| 24 | +function makeMeshNode(name, { onBlock } = {}) { |
| 25 | + const ws = new WebSocket(SIGNAL); |
| 26 | + const offerPCs = new Map(), openDCs = [], seen = new Set(); |
| 27 | + let joined = false; |
| 28 | + const bufs = new WeakMap(); |
| 29 | + const rx = (dc) => (m) => { let b = bufs.get(dc) || Buffer.alloc(0); b = Buffer.concat([b, Buffer.isBuffer(m) ? m : Buffer.from(m)]); const r = p2p.decodeStream(new Uint8Array(b)); bufs.set(dc, b.subarray(r.consumed)); for (const msg of r.messages) { if (msg.command !== 'block') continue; const h = codec.blockHash(msg.payload.header); if (seen.has(h)) continue; seen.add(h); onBlock?.(name, h); for (const o of openDCs) if (o !== dc && o.isOpen()) { try { o.sendMessageBinary(Buffer.from(p2p.encodeMessage('block', msg.payload))); } catch {} } } }; |
| 30 | + const setupDC = (dc) => { bufs.set(dc, Buffer.alloc(0)); dc.onOpen(() => openDCs.push(dc)); dc.onMessage(rx(dc)); }; |
| 31 | + ws.on('open', () => ws.send(JSON.stringify({ type: 'announce', resource: ROOM, offers: [] }))); |
| 32 | + ws.on('message', async (data) => { let m; try { m = JSON.parse(data.toString()); } catch { return; } |
| 33 | + if (m.type === 'resource-peers' && !joined) { joined = true; const offers = []; for (let i = 0; i < m.count; i++) { const pc = new nodeDataChannel.PeerConnection(`${name}o${i}`, { iceServers: ICE }); setupDC(pc.createDataChannel('mesh')); const d = await gather(pc); const id = `${name}-${i}`; offerPCs.set(id, pc); offers.push({ sdp: d.sdp, offer_id: id }); } if (offers.length) ws.send(JSON.stringify({ type: 'announce', resource: ROOM, offers })); } |
| 34 | + else if (m.type === 'offer' && m.resource === ROOM) { const pc = new nodeDataChannel.PeerConnection(`${name}a`, { iceServers: ICE }); pc.onDataChannel(setupDC); pc.setRemoteDescription(m.sdp, 'offer'); const d = await gather(pc); ws.send(JSON.stringify({ type: 'answer', resource: ROOM, to: m.from, offer_id: m.offer_id, sdp: d.sdp })); } |
| 35 | + else if (m.type === 'answer' && m.resource === ROOM) offerPCs.get(m.offer_id)?.setRemoteDescription(m.sdp, 'answer'); }); |
| 36 | + return { name, openDCs, sendToOne: (cmd, payload) => openDCs[0]?.sendMessageBinary(Buffer.from(p2p.encodeMessage(cmd, payload))) }; |
| 37 | +} |
| 38 | + |
| 39 | +const sig = startSignaling(PORT, { log }); |
| 40 | +let done = false; const end = (c) => { if (done) return; done = true; try { sig.close(); } catch {} try { nodeDataChannel.cleanup(); } catch {} setImmediate(() => process.exit(c)); }; |
| 41 | +const to = setTimeout(() => { console.log('\n❌ timeout'); end(1); }, 40000); |
| 42 | + |
| 43 | +const got = new Set(); |
| 44 | +const onBlock = (who, h) => { if (h !== expectHash) return; got.add(who); log(` ${who}: received + validated the block`); if (got.has('B') && got.has('C')) { clearTimeout(to); log('\n✅ gossip mode: source A sent to ONE neighbor; B and C both got it (one via gossip) on a full mesh'); end(0); } }; |
| 45 | + |
| 46 | +log(`\nfull 3-mesh; A will send to ONE neighbor only…`); |
| 47 | +const A = makeMeshNode('A'); |
| 48 | +await sleep(1200); |
| 49 | +const B = makeMeshNode('B', { onBlock }); |
| 50 | +await sleep(1200); |
| 51 | +const C = makeMeshNode('C', { onBlock }); |
| 52 | +for (let i = 0; i < 40 && !(A.openDCs.length >= 2 && B.openDCs.length >= 2 && C.openDCs.length >= 2); i++) await sleep(400); |
| 53 | +log(` mesh up — A:${A.openDCs.length} B:${B.openDCs.length} C:${C.openDCs.length}; A sending block to ONE neighbor`); |
| 54 | +if (A.openDCs.length < 2) { console.log('\n❌ mesh did not fully form'); end(1); } |
| 55 | +else A.sendToOne('block', blockHex); |
0 commit comments