Merge pull request #10 from fccview/feature/security-improvement

Ehm.. refactored a bunch.. sorry
This commit is contained in:
Kilian Tyler
2026-01-02 19:00:56 -05:00
committed by GitHub
26 changed files with 1257 additions and 694 deletions
+128
View File
@@ -0,0 +1,128 @@
name: Build
on:
workflow_call:
inputs:
push-to-registry:
description: "Whether to push images to the container registry"
required: false
default: false
type: boolean
env:
REGISTRY: ghcr.io
IMAGE_NAME: ${{ github.repository }}
jobs:
build:
runs-on: ${{ matrix.runner }}
strategy:
fail-fast: false
matrix:
include:
- platform: linux/amd64
arch: amd64
runner: ubuntu-latest
test: true
- platform: linux/arm64
arch: arm64
runner: ubuntu-24.04-arm
test: true
- platform: linux/arm/v7
arch: armv7
runner: ubuntu-24.04-arm
test: false
steps:
- name: Checkout repository
uses: actions/checkout@v4
- name: Set up QEMU
uses: docker/setup-qemu-action@v3
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
- name: Log in to the Container registry
if: inputs.push-to-registry
uses: docker/login-action@v3
with:
registry: ${{ env.REGISTRY }}
username: ${{ github.actor }}
password: ${{ secrets.GITHUB_TOKEN }}
- name: Extract metadata (tags, labels) for Docker
if: inputs.push-to-registry
id: meta
uses: docker/metadata-action@v5
with:
images: ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}
- name: Build Docker image
id: build
uses: docker/build-push-action@v6
with:
context: .
platforms: ${{ matrix.platform }}
labels: ${{ steps.meta.outputs.labels }}
load: ${{ matrix.test }}
tags: hypermind-test:${{ github.sha }}
cache-from: type=gha
cache-to: type=gha,mode=max
- name: Run container for testing
if: matrix.test
run: docker run -d --name hypermind-test -p 3000:3000 hypermind-test:${{ github.sha }}
- name: Wait for server to be ready
if: matrix.test
run: |
for i in {1..30}; do
if curl -sf http://localhost:3000/api/stats; then
echo "Server is ready"
exit 0
fi
echo "Waiting for server... ($i/30)"
sleep 1
done
echo "Server failed to start"
docker logs hypermind-test
exit 1
- name: Verify API response
if: matrix.test
run: |
response=$(curl -sf http://localhost:3000/api/stats)
echo "Response: $response"
echo "$response" | jq -e '.count != null and .id != null'
- name: Cleanup test container
if: always() && matrix.test
run: docker rm -f hypermind-test || true
- name: Push by digest
if: inputs.push-to-registry
id: push
uses: docker/build-push-action@v6
with:
context: .
platforms: ${{ matrix.platform }}
labels: ${{ steps.meta.outputs.labels }}
outputs: type=image,name=${{ env.REGISTRY }}/${{ env.IMAGE_NAME }},push-by-digest=true,name-canonical=true,push=true
cache-from: type=gha
- name: Export digest
if: inputs.push-to-registry
run: |
mkdir -p ${{ runner.temp }}/digests
digest="${{ steps.push.outputs.digest }}"
touch "${{ runner.temp }}/digests/${digest#sha256:}"
- name: Upload digest
if: inputs.push-to-registry
uses: actions/upload-artifact@v4
with:
name: digests-${{ matrix.arch }}
path: ${{ runner.temp }}/digests/*
if-no-files-found: error
retention-days: 1
+13
View File
@@ -0,0 +1,13 @@
name: CI
on:
pull_request:
branches: ["main"]
jobs:
build:
uses: ./.github/workflows/build.yml
permissions:
contents: read
with:
push-to-registry: false
+4 -63
View File
@@ -1,8 +1,6 @@
name: Publish Docker Image
on:
pull_request:
branches: ["main"]
push:
branches: ["main"]
release:
@@ -15,73 +13,16 @@ env:
jobs:
build:
runs-on: ${{ matrix.runner }}
uses: ./.github/workflows/build.yml
permissions:
contents: read
packages: write
strategy:
fail-fast: false
matrix:
include:
- platform: linux/amd64
runner: ubuntu-latest
- platform: linux/arm64
runner: ubuntu-24.04-arm
- platform: linux/arm/v7
runner: ubuntu-24.04-arm
steps:
- name: Checkout repository
uses: actions/checkout@v4
- name: Set up QEMU
uses: docker/setup-qemu-action@v3
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
- name: Log in to the Container registry
if: github.event_name != 'pull_request'
uses: docker/login-action@v3
with:
registry: ${{ env.REGISTRY }}
username: ${{ github.actor }}
password: ${{ secrets.GITHUB_TOKEN }}
- name: Extract metadata (tags, labels) for Docker
id: meta
uses: docker/metadata-action@v5
with:
images: ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}
- name: Build and push by digest
id: build
uses: docker/build-push-action@v6
with:
context: .
platforms: ${{ matrix.platform }}
labels: ${{ steps.meta.outputs.labels }}
outputs: ${{ github.event_name != 'pull_request' && format('type=image,name={0}/{1},push-by-digest=true,name-canonical=true,push=true', env.REGISTRY, env.IMAGE_NAME) || 'type=docker' }}
- name: Export digest
if: github.event_name != 'pull_request'
run: |
mkdir -p ${{ runner.temp }}/digests
digest="${{ steps.build.outputs.digest }}"
touch "${{ runner.temp }}/digests/${digest#sha256:}"
- name: Upload digest
if: github.event_name != 'pull_request'
uses: actions/upload-artifact@v4
with:
name: digests-${{ matrix.platform == 'linux/amd64' && 'amd64' || matrix.platform == 'linux/arm64' && 'arm64' || 'armv7' }}
path: ${{ runner.temp }}/digests/*
if-no-files-found: error
retention-days: 1
secrets: inherit
with:
push-to-registry: true
merge:
runs-on: ubuntu-latest
if: github.event_name != 'pull_request'
needs: build
permissions:
contents: read
+2
View File
@@ -6,7 +6,9 @@ COPY package*.json ./
RUN npm ci --omit=dev
COPY public/ ./public/
COPY server.js hypermind2.svg LICENSE ./
COPY src/ ./src/
ENV PORT=3000
ENV NODE_ENV=production
+7
View File
@@ -73,6 +73,13 @@ services:
```
## Environment Variables
| Variable | Default | Description |
|----------|---------|-------------|
| `PORT` | `3000` | The port the web dashboard listens on. Since `--network host` is used, this port opens directly on the host. |
| `MAX_PEERS` | `10000` | Maximum number of peers to track in the swarm. Unless you're expecting the entire internet to join, the default is probably fine. |
## Usage
Open your browser to: `http://localhost:3000`
+1
View File
@@ -6,3 +6,4 @@ services:
restart: unless-stopped
environment:
- PORT=3000
- MAX_PEERS=10000
+13
View File
@@ -9,6 +9,7 @@
"version": "1.0.0",
"license": "MIT",
"dependencies": {
"dotenv": "^17.2.3",
"express": "^5.2.1",
"hyperswarm": "^4.16.0"
}
@@ -321,6 +322,18 @@
"udx-native": "^1.5.3"
}
},
"node_modules/dotenv": {
"version": "17.2.3",
"resolved": "https://registry.npmjs.org/dotenv/-/dotenv-17.2.3.tgz",
"integrity": "sha512-JVUnt+DUIzu87TABbhPmNfVdBDt18BLOWjMUFJMSi/Qqg7NTYtabbvSNJGOJ7afbRuv9D/lngizHtP7QyLQ+9w==",
"license": "BSD-2-Clause",
"engines": {
"node": ">=12"
},
"funding": {
"url": "https://dotenvx.com"
}
},
"node_modules/dunder-proto": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/dunder-proto/-/dunder-proto-1.0.1.tgz",
+1
View File
@@ -26,6 +26,7 @@
"homepage": "https://github.com/lklynet/hypermind#readme",
"type": "commonjs",
"dependencies": {
"dotenv": "^17.2.3",
"express": "^5.2.1",
"hyperswarm": "^4.16.0"
}
+141
View File
@@ -0,0 +1,141 @@
const countEl = document.getElementById('count');
const directEl = document.getElementById('direct');
const canvas = document.getElementById('network');
const ctx = canvas.getContext('2d');
let particles = [];
function resize() {
canvas.width = window.innerWidth;
canvas.height = window.innerHeight;
}
window.addEventListener('resize', resize);
resize();
class Particle {
constructor() {
this.x = Math.random() * canvas.width;
this.y = Math.random() * canvas.height;
this.vx = (Math.random() - 0.5) * 1;
this.vy = (Math.random() - 0.5) * 1;
this.size = 3;
}
update() {
this.x += this.vx;
this.y += this.vy;
if (this.x < 0 || this.x > canvas.width) this.vx *= -1;
if (this.y < 0 || this.y > canvas.height) this.vy *= -1;
}
draw() {
ctx.beginPath();
ctx.arc(this.x, this.y, this.size, 0, Math.PI * 2);
ctx.fillStyle = '#4ade80';
ctx.fill();
}
}
const updateParticles = (count) => {
const VISUAL_LIMIT = 500;
const visualCount = Math.min(count, VISUAL_LIMIT);
const currentCount = particles.length;
if (visualCount > currentCount) {
for (let i = 0; i < visualCount - currentCount; i++) {
particles.push(new Particle());
}
} else if (visualCount < currentCount) {
particles.splice(visualCount, currentCount - visualCount);
}
}
const animate = () => {
ctx.clearRect(0, 0, canvas.width, canvas.height);
ctx.strokeStyle = 'rgba(74, 222, 128, 0.15)';
ctx.lineWidth = 1;
for (let i = 0; i < particles.length; i++) {
for (let j = i + 1; j < particles.length; j++) {
const dx = particles[i].x - particles[j].x;
const dy = particles[i].y - particles[j].y;
const distance = Math.sqrt(dx * dx + dy * dy);
if (distance < 150) {
ctx.beginPath();
ctx.moveTo(particles[i].x, particles[i].y);
ctx.lineTo(particles[j].x, particles[j].y);
ctx.stroke();
}
}
}
particles.forEach(p => {
p.update();
p.draw();
});
requestAnimationFrame(animate);
}
const openDiagnostics = () => {
document.getElementById('diagnosticsModal').classList.add('active');
}
const closeDiagnostics = () => {
document.getElementById('diagnosticsModal').classList.remove('active');
}
document.getElementById('diagnosticsModal').addEventListener('click', (e) => {
if (e.target.id === 'diagnosticsModal') {
closeDiagnostics();
}
});
document.addEventListener('keydown', (e) => {
if (e.key === 'Escape') {
closeDiagnostics();
}
});
const evtSource = new EventSource("/events");
evtSource.onmessage = (event) => {
const data = JSON.parse(event.data);
updateParticles(data.count);
if (countEl.innerText != data.count) {
countEl.innerText = data.count;
countEl.classList.remove('pulse');
void countEl.offsetWidth;
countEl.classList.add('pulse');
}
directEl.innerText = data.direct;
if (data.diagnostics) {
const d = data.diagnostics;
document.getElementById('diag-heartbeats-rx').innerText = d.heartbeatsReceived.toLocaleString();
document.getElementById('diag-heartbeats-tx').innerText = d.heartbeatsRelayed.toLocaleString();
document.getElementById('diag-new-peers').innerText = d.newPeersAdded.toLocaleString();
document.getElementById('diag-dup-seq').innerText = d.duplicateSeq.toLocaleString();
document.getElementById('diag-invalid-pow').innerText = d.invalidPoW.toLocaleString();
document.getElementById('diag-invalid-sig').innerText = d.invalidSig.toLocaleString();
document.getElementById('diag-bandwidth-in').innerText = (d.bytesReceived / 1024).toFixed(1) + ' KB';
document.getElementById('diag-bandwidth-out').innerText = (d.bytesRelayed / 1024).toFixed(1) + ' KB';
document.getElementById('diag-leave').innerText = d.leaveMessages.toLocaleString();
}
};
evtSource.onerror = (err) => {
// Removing console error here as it's extremely spammy in the browser console and it will reconnct automatically anyway, so pretty redundant.
};
const initialCount = parseInt(countEl.dataset.initialCount) || 0;
countEl.innerText = initialCount;
countEl.classList.add('loaded');
updateParticles(initialCount);
animate();
Binary file not shown.

After

Width:  |  Height:  |  Size: 15 KiB

+71
View File
@@ -0,0 +1,71 @@
<!DOCTYPE html>
<html>
<head>
<title>Hypermind</title>
<meta name="viewport" content="width=device-width, initial-scale=1">
<link rel="icon" href="/favicon.ico">
<link rel="stylesheet" href="/style.css">
</head>
<body>
<canvas id="network"></canvas>
<div class="container">
<div id="count" class="count" data-initial-count="{{COUNT}}">{{COUNT}}</div>
<div class="label">Active Nodes</div>
<div class="footer">
powered by <a href="https://github.com/lklynet/hypermind" target="_blank">hypermind</a>
</div>
<div class="debug">
ID: {{ID}}<br>
Direct Connections: <span id="direct">{{DIRECT}}</span><br>
<span class="debug-link" onclick="openDiagnostics()">diagnostics</span>
</div>
</div>
<div id="diagnosticsModal" class="modal">
<div class="modal-content">
<button class="close-btn" onclick="closeDiagnostics()">×</button>
<div class="modal-title">Diagnostics</div>
<div class="stat-row">
<span class="stat-label">Heartbeats Received</span>
<span class="stat-value" id="diag-heartbeats-rx">0</span>
</div>
<div class="stat-row">
<span class="stat-label">Heartbeats Relayed</span>
<span class="stat-value" id="diag-heartbeats-tx">0</span>
</div>
<div class="stat-row">
<span class="stat-label">New Peers Added</span>
<span class="stat-value" id="diag-new-peers">0</span>
</div>
<div class="stat-row">
<span class="stat-label">Duplicate/Old Seq</span>
<span class="stat-value" id="diag-dup-seq">0</span>
</div>
<div class="stat-row">
<span class="stat-label">Invalid PoW</span>
<span class="stat-value" id="diag-invalid-pow">0</span>
</div>
<div class="stat-row">
<span class="stat-label">Invalid Signatures</span>
<span class="stat-value" id="diag-invalid-sig">0</span>
</div>
<div class="stat-row">
<span class="stat-label">Bandwidth In</span>
<span class="stat-value" id="diag-bandwidth-in">0 KB</span>
</div>
<div class="stat-row">
<span class="stat-label">Bandwidth Out</span>
<span class="stat-value" id="diag-bandwidth-out">0 KB</span>
</div>
<div class="stat-row">
<span class="stat-label">LEAVE Messages</span>
<span class="stat-value" id="diag-leave">0</span>
</div>
<div class="update-time" id="last-update">last 10 seconds</div>
</div>
</div>
<script src="/app.js"></script>
</body>
</html>
+87
View File
@@ -0,0 +1,87 @@
* { margin: 0; padding: 0; box-sizing: border-box; }
body {
font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, Helvetica, Arial, sans-serif;
display: flex;
justify-content: center;
align-items: center;
height: 100vh;
background: #111;
color: #eee;
margin: 0;
}
.container { text-align: center; position: relative; z-index: 10; }
#network { position: fixed; top: 0; left: 0; width: 100%; height: 100%; z-index: 1; }
.count { font-size: 8rem; font-weight: bold; color: #4ade80; transition: color 0.2s; visibility: hidden; }
.count.loaded { visibility: visible; }
.label { font-size: 1.5rem; color: #9ca3af; margin-top: 1rem; }
.footer { margin-top: 2rem; font-size: 0.9rem; color: #4b5563; }
.debug { font-size: 0.8rem; color: #4b5563; margin-top: 1rem; }
.debug-link { color: #4b5563; border-bottom: 1px dotted #4b5563; cursor: pointer; }
.debug-link:hover { color: #9ca3af; }
a { color: #4b5563; text-decoration: none; border-bottom: 1px dotted #4b5563; }
.pulse { animation: pulse 0.5s ease-in-out; }
@keyframes pulse {
0% { transform: scale(1); }
50% { transform: scale(1.1); color: #fff; }
100% { transform: scale(1); }
}
.modal {
display: none;
position: fixed;
z-index: 1000;
left: 0;
top: 0;
width: 100%;
height: 100%;
background: rgba(0, 0, 0, 0.8);
}
.modal.active { display: flex; align-items: center; justify-content: center; }
.modal-content {
background: #111;
border: 1px solid #222;
padding: 2rem;
max-width: 500px;
width: 90%;
position: relative;
}
.modal-title {
font-size: 0.9rem;
color: #666;
margin-bottom: 1.5rem;
text-transform: uppercase;
letter-spacing: 1px;
}
.close-btn {
position: absolute;
top: 1.5rem;
right: 1.5rem;
background: none;
border: none;
color: #333;
font-size: 1.2rem;
cursor: pointer;
}
.close-btn:hover { color: #666; }
.stat-row {
display: flex;
justify-content: space-between;
padding: 0.5rem 0;
border-bottom: 1px solid #1a1a1a;
font-size: 0.85rem;
}
.stat-row:last-child { border-bottom: none; }
.stat-label { color: #4b5563; }
.stat-value {
color: #9ca3af;
font-variant-numeric: tabular-nums;
}
.update-time {
text-align: center;
font-size: 0.7rem;
color: #333;
margin-top: 1rem;
}
+54 -631
View File
@@ -1,646 +1,69 @@
const express = require("express");
const Hyperswarm = require("hyperswarm");
const crypto = require("crypto");
require('dotenv').config();
const app = express();
const PORT = process.env.PORT || 3000;
const { generateIdentity } = require("./src/core/identity");
const { PeerManager } = require("./src/state/peers");
const { DiagnosticsManager } = require("./src/state/diagnostics");
const { MessageHandler } = require("./src/p2p/messaging");
const { relayMessage } = require("./src/p2p/relay");
const { SwarmManager } = require("./src/p2p/swarm");
const { SSEManager } = require("./src/web/sse");
const { createServer, startServer } = require("./src/web/server");
const { DIAGNOSTICS_INTERVAL } = require("./src/config/constants");
// --- CONFIGURATION ---
const TOPIC_NAME = "hypermind-lklynet-v1";
const TOPIC = crypto.createHash("sha256").update(TOPIC_NAME).digest();
const main = async () => {
const identity = generateIdentity();
const peerManager = new PeerManager();
const diagnostics = new DiagnosticsManager();
const sseManager = new SSEManager();
// Gossip protocol tuning
const GOSSIP_FANOUT = 10; // Relay to max 10 random peers instead of all
const HEARTBEAT_INTERVAL_FAST = 1000; // 1 second during startup
const HEARTBEAT_INTERVAL_SLOW = 15000; // 15 seconds at steady state
const STARTUP_DURATION = 120000; // Stay in fast mode for 2 minutes
const PEER_STALE_TIMEOUT = 90000; // 90 seconds before peer considered stale
peerManager.addOrUpdatePeer(identity.id, peerManager.getSeq(), null);
// --- BLOOM FILTER FOR MESSAGE DEDUPLICATION ---
// Simple bloom filter to prevent re-relaying messages we've already seen
class BloomFilter {
constructor(size = 10000, hashCount = 3) {
this.size = size;
this.hashCount = hashCount;
this.bits = new Uint8Array(Math.ceil(size / 8));
}
const broadcastUpdate = () => {
sseManager.broadcastUpdate({
count: peerManager.size,
direct: swarmManager.getSwarm().connections.size,
id: identity.id,
diagnostics: diagnostics.getStats(),
});
};
_hash(str, seed) {
let h = seed;
for (let i = 0; i < str.length; i++) {
h = (h * 31 + str.charCodeAt(i)) >>> 0;
}
return h % this.size;
}
const messageHandler = new MessageHandler(
peerManager,
diagnostics,
(msg, sourceSocket) => relayMessage(msg, sourceSocket, swarmManager.getSwarm(), diagnostics),
broadcastUpdate
);
add(item) {
for (let i = 0; i < this.hashCount; i++) {
const idx = this._hash(item, i * 0x9e3779b9);
this.bits[idx >>> 3] |= (1 << (idx & 7));
}
}
const swarmManager = new SwarmManager(
identity,
peerManager,
diagnostics,
messageHandler,
(msg, sourceSocket) => relayMessage(msg, sourceSocket, swarmManager.getSwarm(), diagnostics),
broadcastUpdate
);
has(item) {
for (let i = 0; i < this.hashCount; i++) {
const idx = this._hash(item, i * 0x9e3779b9);
if ((this.bits[idx >>> 3] & (1 << (idx & 7))) === 0) {
return false;
}
}
return true;
}
await swarmManager.start();
clear() {
this.bits.fill(0);
}
}
diagnostics.startLogging(
() => peerManager.size,
() => swarmManager.getSwarm().connections.size
);
// Time-bucketed bloom filter - rotates every 30 seconds
let currentBloom = new BloomFilter();
let previousBloom = new BloomFilter();
function rotateBloomFilters() {
previousBloom = currentBloom;
currentBloom = new BloomFilter();
}
// Check if we've recently relayed this message
function hasRelayedMessage(id, seq) {
const key = `${id}:${seq}`;
return currentBloom.has(key) || previousBloom.has(key);
}
// Mark message as relayed
function markRelayed(id, seq) {
const key = `${id}:${seq}`;
currentBloom.add(key);
}
// Rotate bloom filters periodically
setInterval(rotateBloomFilters, 30000);
// --- HYPERLOGLOG FOR PEER COUNTING ---
// Approximate unique peer count with fixed ~1.5KB memory
// Accuracy: ~2% error rate, can count millions of peers
class HyperLogLog {
constructor(precision = 10) {
// 2^precision registers, precision=10 gives 1024 registers (~1KB)
this.precision = precision;
this.registerCount = 1 << precision;
this.registers = new Uint8Array(this.registerCount);
this.alphaMM = this._getAlpha() * this.registerCount * this.registerCount;
}
_getAlpha() {
// Bias correction constant
switch (this.precision) {
case 4: return 0.673;
case 5: return 0.697;
case 6: return 0.709;
default: return 0.7213 / (1 + 1.079 / this.registerCount);
}
}
_hash(str) {
// Simple 32-bit hash (good enough for HLL)
let h = 0x811c9dc5;
for (let i = 0; i < str.length; i++) {
h ^= str.charCodeAt(i);
h = (h * 0x01000193) >>> 0;
}
return h;
}
_countLeadingZeros(value, maxBits) {
if (value === 0) return maxBits;
let count = 0;
while ((value & (1 << (maxBits - 1 - count))) === 0 && count < maxBits) {
count++;
}
return count;
}
add(item) {
const hash = this._hash(item);
// Use first 'precision' bits for register index
const registerIndex = hash >>> (32 - this.precision);
// Use remaining bits to count leading zeros
const remainingBits = hash << this.precision;
const leadingZeros = this._countLeadingZeros(remainingBits, 32 - this.precision) + 1;
// Store maximum leading zeros seen for this register
if (leadingZeros > this.registers[registerIndex]) {
this.registers[registerIndex] = leadingZeros;
}
}
count() {
// Harmonic mean of 2^register values
let harmonicSum = 0;
let zeroRegisters = 0;
for (let i = 0; i < this.registerCount; i++) {
harmonicSum += Math.pow(2, -this.registers[i]);
if (this.registers[i] === 0) zeroRegisters++;
}
let estimate = this.alphaMM / harmonicSum;
// Small range correction (linear counting)
if (estimate <= 2.5 * this.registerCount && zeroRegisters > 0) {
estimate = this.registerCount * Math.log(this.registerCount / zeroRegisters);
}
return Math.round(estimate);
}
}
// Global peer counter - tracks all unique peers ever seen
const peerCounter = new HyperLogLog(10); // ~1KB, 2% error
// --- SECURITY ---
// We use Ed25519 for signatures and a PoW puzzle to prevent Sybil attacks.
// Difficulty: Hash(ID + nonce) must start with '0000'
const POW_PREFIX = "0000";
console.log("[Security] Generating Identity & Solving PoW...");
const { publicKey, privateKey } = crypto.generateKeyPairSync("ed25519");
const MY_ID = publicKey.export({ type: "spki", format: "der" }).toString("hex");
let MY_NONCE = 0;
while (true) {
const hash = crypto
.createHash("sha256")
.update(MY_ID + MY_NONCE)
.digest("hex");
if (hash.startsWith(POW_PREFIX)) break;
MY_NONCE++;
}
console.log(
`[Security] Identity ready. ID: ${MY_ID.slice(0, 8)}... Nonce: ${MY_NONCE}`
);
let mySeq = 0;
const seenPeers = new Map();
const sseClients = new Set();
seenPeers.set(MY_ID, { seq: mySeq, lastSeen: Date.now() });
peerCounter.add(MY_ID); // Count ourselves
// Throttle updates to once per second
let lastBroadcast = 0;
function broadcastUpdate() {
const now = Date.now();
if (now - lastBroadcast < 1000) return;
lastBroadcast = now;
const data = JSON.stringify({
count: peerCounter.count(), // Use HyperLogLog for total peer count
direct: swarm.connections.size,
id: MY_ID,
});
for (const client of sseClients) {
client.write(`data: ${data}\n\n`);
}
}
const swarm = new Hyperswarm();
swarm.on("connection", (socket) => {
// Start adaptive heartbeat on first connection
startHeartbeatIfNeeded();
const sig = crypto
.sign(null, Buffer.from(`seq:${mySeq}`), privateKey)
.toString("hex");
const hello = JSON.stringify({
type: "HEARTBEAT",
id: MY_ID,
seq: mySeq,
hops: 0,
nonce: MY_NONCE,
sig,
});
socket.write(hello);
broadcastUpdate();
socket.on("data", (data) => {
try {
const msgs = data
.toString()
.split("\n")
.filter((x) => x.trim());
for (const msgStr of msgs) {
const msg = JSON.parse(msgStr);
handleMessage(msg, socket);
}
} catch (e) {
// console.error('Invalid message', e);
}
});
socket.on("close", () => {
if (socket.peerId && seenPeers.has(socket.peerId)) {
seenPeers.delete(socket.peerId);
}
setInterval(() => {
broadcastUpdate();
});
}, DIAGNOSTICS_INTERVAL);
socket.on("error", () => {});
});
const app = createServer(identity, peerManager, swarmManager, sseManager, diagnostics);
startServer(app, identity);
const discovery = swarm.join(TOPIC);
discovery.flushed().then(() => {
console.log("[P2P] Joined topic:", TOPIC_NAME);
});
const handleShutdown = () => {
diagnostics.stopLogging();
swarmManager.shutdown();
};
function handleMessage(msg, sourceSocket) {
if (msg.type === "HEARTBEAT") {
const { id, seq, hops, nonce, sig } = msg;
// 1. Verify PoW
if (!nonce) return;
const powHash = crypto
.createHash("sha256")
.update(id + nonce)
.digest("hex");
if (!powHash.startsWith(POW_PREFIX)) return; // Invalid PoW
// 2. Check Sequence (Optimization: Drop duplicates before expensive verify)
const stored = seenPeers.get(id);
if (stored && seq <= stored.seq) return; // Ignore old/duplicate messages
// 3. Verify Signature
if (!sig) return;
try {
let key;
if (stored && stored.key) {
key = stored.key;
} else {
key = crypto.createPublicKey({
key: Buffer.from(id, "hex"),
format: "der",
type: "spki",
});
}
const verified = crypto.verify(
null,
Buffer.from(`seq:${seq}`),
key,
Buffer.from(sig, "hex")
);
if (!verified) return; // Invalid Signature
// Track unique peer in HyperLogLog counter (always, even if over MAX_PEERS)
const prevCount = peerCounter.count();
peerCounter.add(id);
const countChanged = peerCounter.count() !== prevCount;
// Update Peer
if (hops === 0) {
sourceSocket.peerId = id;
}
const now = Date.now();
const wasNew = !stored;
// Store in seenPeers only if we have room (memory limit)
// But we still count and relay even if we can't store
const canStore = stored || seenPeers.size < MAX_PEERS;
if (canStore) {
seenPeers.set(id, { seq, lastSeen: now, keyDer });
}
if ((wasNew && canStore) || countChanged) broadcastUpdate();
// Only relay if we haven't already relayed this message (bloom filter check)
if (hops < 3 && !hasRelayedMessage(id, seq)) {
markRelayed(id, seq);
relayMessage({ ...msg, hops: hops + 1 }, sourceSocket);
}
} catch (e) {
return;
}
} else if (msg.type === "LEAVE") {
const { id, hops } = msg;
if (seenPeers.has(id)) {
seenPeers.delete(id);
broadcastUpdate();
// Use id:leave as key for LEAVE messages
if (hops < 3 && !hasRelayedMessage(id, "leave")) {
markRelayed(id, "leave");
relayMessage({ ...msg, hops: hops + 1 }, sourceSocket);
}
}
}
process.on("SIGINT", handleShutdown);
process.on("SIGTERM", handleShutdown);
}
// Fisher-Yates shuffle for random peer selection
function shuffleArray(array) {
for (let i = array.length - 1; i > 0; i--) {
const j = Math.floor(Math.random() * (i + 1));
[array[i], array[j]] = [array[j], array[i]];
}
return array;
}
function relayMessage(msg, sourceSocket) {
const data = JSON.stringify(msg) + "\n";
// Get all eligible sockets (excluding source)
const eligibleSockets = [...swarm.connections].filter(s => s !== sourceSocket);
// Apply fanout limiting - only relay to GOSSIP_FANOUT random peers
const targetSockets = eligibleSockets.length <= GOSSIP_FANOUT
? eligibleSockets
: shuffleArray(eligibleSockets).slice(0, GOSSIP_FANOUT);
for (const socket of targetSockets) {
socket.write(data);
}
}
// Adaptive Heartbeat - fast at startup, slows down after STARTUP_DURATION
// Timer starts when first connection is established, not at process start
let heartbeatStartTime = null;
let heartbeatStarted = false;
function getHeartbeatInterval() {
if (!heartbeatStartTime) return HEARTBEAT_INTERVAL_FAST;
const elapsed = Date.now() - heartbeatStartTime;
return elapsed < STARTUP_DURATION ? HEARTBEAT_INTERVAL_FAST : HEARTBEAT_INTERVAL_SLOW;
}
function sendHeartbeat() {
mySeq++;
seenPeers.set(MY_ID, { seq: mySeq, lastSeen: Date.now() });
const sig = crypto
.sign(null, Buffer.from(`seq:${mySeq}`), privateKey)
.toString("hex");
const heartbeat =
JSON.stringify({
type: "HEARTBEAT",
id: MY_ID,
seq: mySeq,
hops: 0,
nonce: MY_NONCE,
sig,
}) + "\n";
for (const socket of swarm.connections) {
socket.write(heartbeat);
}
const now = Date.now();
let changed = false;
for (const [id, data] of seenPeers) {
if (now - data.lastSeen > PEER_STALE_TIMEOUT) {
seenPeers.delete(id);
changed = true;
}
}
if (changed) broadcastUpdate();
// Schedule next heartbeat with adaptive interval
setTimeout(sendHeartbeat, getHeartbeatInterval());
}
// Start heartbeat loop on first connection
function startHeartbeatIfNeeded() {
if (!heartbeatStarted) {
heartbeatStarted = true;
heartbeatStartTime = Date.now();
console.log("[P2P] First connection established, starting fast heartbeat...");
sendHeartbeat();
}
}
// Graceful Shutdown
function handleShutdown() {
console.log("[P2P] Shutting down, sending goodbye...");
const goodbye = JSON.stringify({ type: "LEAVE", id: MY_ID, hops: 0 }) + "\n";
for (const socket of swarm.connections) {
socket.write(goodbye);
}
setTimeout(() => {
process.exit(0);
}, 500);
}
process.on("SIGINT", handleShutdown);
process.on("SIGTERM", handleShutdown);
// --- WEB SERVER ---
app.get("/", (req, res) => {
const count = peerCounter.count(); // HyperLogLog approximate count
const directPeers = swarm.connections.size;
res.send(`
<!DOCTYPE html>
<html>
<head>
<title>Hypermind</title>
<link rel="icon" type="image/svg+xml" href="/favicon.svg">
<meta name="viewport" content="width=device-width, initial-scale=1">
<style>
body {
font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, Helvetica, Arial, sans-serif;
display: flex;
justify-content: center;
align-items: center;
height: 100vh;
background: #111;
color: #eee;
margin: 0;
}
.container { text-align: center; position: relative; z-index: 10; }
#network { position: fixed; top: 0; left: 0; width: 100%; height: 100%; z-index: 1; }
.count { font-size: 8rem; font-weight: bold; color: #4ade80; transition: color 0.2s; }
.label { font-size: 1.5rem; color: #9ca3af; margin-top: 1rem; }
.footer { margin-top: 2rem; font-size: 0.9rem; color: #4b5563; }
.debug { font-size: 0.8rem; color: #333; margin-top: 1rem; }
a { color: #4b5563; text-decoration: none; border-bottom: 1px dotted #4b5563; }
.pulse { animation: pulse 0.5s ease-in-out; }
@keyframes pulse {
0% { transform: scale(1); }
50% { transform: scale(1.1); color: #fff; }
100% { transform: scale(1); }
}
</style>
</head>
<body>
<canvas id="network"></canvas>
<div class="container">
<div id="count" class="count">${count}</div>
<div class="label">Active Nodes</div>
<div class="footer">
powered by <a href="https://github.com/lklynet/hypermind" target="_blank">hypermind</a>
</div>
<div class="debug">
ID: ${MY_ID.slice(0, 8)}...<br>
Direct Connections: <span id="direct">${directPeers}</span>
</div>
</div>
<script>
const countEl = document.getElementById('count');
const directEl = document.getElementById('direct');
// Particle System
const canvas = document.getElementById('network');
const ctx = canvas.getContext('2d');
let particles = [];
function resize() {
canvas.width = window.innerWidth;
canvas.height = window.innerHeight;
}
window.addEventListener('resize', resize);
resize();
class Particle {
constructor() {
this.x = Math.random() * canvas.width;
this.y = Math.random() * canvas.height;
this.vx = (Math.random() - 0.5) * 1;
this.vy = (Math.random() - 0.5) * 1;
this.size = 3;
}
update() {
this.x += this.vx;
this.y += this.vy;
if (this.x < 0 || this.x > canvas.width) this.vx *= -1;
if (this.y < 0 || this.y > canvas.height) this.vy *= -1;
}
draw() {
ctx.beginPath();
ctx.arc(this.x, this.y, this.size, 0, Math.PI * 2);
ctx.fillStyle = '#4ade80';
ctx.fill();
}
}
function updateParticles(count) {
// Limit visual particles to 500 to prevent browser crash
const VISUAL_LIMIT = 500;
const visualCount = Math.min(count, VISUAL_LIMIT);
const currentCount = particles.length;
if (visualCount > currentCount) {
for (let i = 0; i < visualCount - currentCount; i++) {
particles.push(new Particle());
}
} else if (visualCount < currentCount) {
particles.splice(visualCount, currentCount - visualCount);
}
}
// Initialize with server-rendered count
updateParticles(${count});
function animate() {
ctx.clearRect(0, 0, canvas.width, canvas.height);
// Draw connections
ctx.strokeStyle = 'rgba(74, 222, 128, 0.15)';
ctx.lineWidth = 1;
for (let i = 0; i < particles.length; i++) {
for (let j = i + 1; j < particles.length; j++) {
const dx = particles[i].x - particles[j].x;
const dy = particles[i].y - particles[j].y;
const distance = Math.sqrt(dx * dx + dy * dy);
if (distance < 150) {
ctx.beginPath();
ctx.moveTo(particles[i].x, particles[i].y);
ctx.lineTo(particles[j].x, particles[j].y);
ctx.stroke();
}
}
}
particles.forEach(p => {
p.update();
p.draw();
});
requestAnimationFrame(animate);
}
animate();
// Use Server-Sent Events for realtime updates
const evtSource = new EventSource("/events");
evtSource.onmessage = (event) => {
const data = JSON.parse(event.data);
updateParticles(data.count);
// Only update and animate if changed
if (countEl.innerText != data.count) {
countEl.innerText = data.count;
countEl.classList.remove('pulse');
void countEl.offsetWidth; // trigger reflow
countEl.classList.add('pulse');
}
directEl.innerText = data.direct;
};
evtSource.onerror = (err) => {
console.error("EventSource failed:", err);
};
</script>
</body>
</html>
`);
});
// SSE Endpoint
app.get("/favicon.svg", (req, res) => {
res.sendFile(__dirname + "/hypermind2.svg");
});
app.get("/events", (req, res) => {
res.setHeader("Content-Type", "text/event-stream");
res.setHeader("Cache-Control", "no-cache");
res.setHeader("Connection", "keep-alive");
res.flushHeaders();
sseClients.add(res);
const data = JSON.stringify({
count: peerCounter.count(),
direct: swarm.connections.size,
id: MY_ID,
});
res.write(`data: ${data}\n\n`);
req.on("close", () => {
sseClients.delete(res);
});
});
app.get("/api/stats", (req, res) => {
res.json({
count: peerCounter.count(),
direct: swarm.connections.size,
id: MY_ID,
});
});
app.listen(PORT, () => {
console.log(`Hypermind Node running on port ${PORT}`);
console.log(`ID: ${MY_ID}`);
});
main().catch(console.error);
+38
View File
@@ -0,0 +1,38 @@
const crypto = require("crypto");
const TOPIC_NAME = "hypermind-lklynet-v1";
const TOPIC = crypto.createHash("sha256").update(TOPIC_NAME).digest();
/**
* fccview here, frankly I don't think I can make this more secure, we can change it to `00000` but
* that means until everyone upgrade there'll be a divide between nodes.
*
* I ran it that way and I was fairly isolated, with hundreds of failed POW, shame.
* adding an extra 0 makes it very expensive on attacker to make it worth the fun for them, so maybe consider it.
*
*/
const POW_PREFIX = "0000";
const MAX_PEERS = parseInt(process.env.MAX_PEERS) || 10000;
const MAX_MESSAGE_SIZE = 2048;
const MAX_RELAY_HOPS = 3;
const HEARTBEAT_INTERVAL = 5000;
const PEER_TIMEOUT = 15000;
const BROADCAST_THROTTLE = 1000;
const DIAGNOSTICS_INTERVAL = 10000;
const PORT = process.env.PORT || 3000;
module.exports = {
TOPIC_NAME,
TOPIC,
POW_PREFIX,
MAX_PEERS,
MAX_MESSAGE_SIZE,
MAX_RELAY_HOPS,
HEARTBEAT_INTERVAL,
PEER_TIMEOUT,
BROADCAST_THROTTLE,
DIAGNOSTICS_INTERVAL,
PORT,
};
+21
View File
@@ -0,0 +1,21 @@
const crypto = require("crypto");
const { POW_PREFIX } = require("../config/constants");
const generateIdentity = () => {
const { publicKey, privateKey } = crypto.generateKeyPairSync("ed25519");
const id = publicKey.export({ type: "spki", format: "der" }).toString("hex");
let nonce = 0;
while (true) {
const hash = crypto
.createHash("sha256")
.update(id + nonce)
.digest("hex");
if (hash.startsWith(POW_PREFIX)) break;
nonce++;
}
return { publicKey, privateKey, id, nonce };
}
module.exports = { generateIdentity };
+43
View File
@@ -0,0 +1,43 @@
const crypto = require("crypto");
const { POW_PREFIX } = require("../config/constants");
const verifyPoW = (id, nonce) => {
if (!nonce) return false;
const powHash = crypto
.createHash("sha256")
.update(id + nonce)
.digest("hex");
return powHash.startsWith(POW_PREFIX);
}
const signMessage = (message, privateKey) => {
return crypto.sign(null, Buffer.from(message), privateKey).toString("hex");
}
const verifySignature = (message, signature, publicKey) => {
try {
return crypto.verify(
null,
Buffer.from(message),
publicKey,
Buffer.from(signature, "hex")
);
} catch (e) {
return false;
}
}
const createPublicKey = (id) => {
return crypto.createPublicKey({
key: Buffer.from(id, "hex"),
format: "der",
type: "spki",
});
}
module.exports = {
verifyPoW,
signMessage,
verifySignature,
createPublicKey,
};
+132
View File
@@ -0,0 +1,132 @@
const { verifyPoW, verifySignature, createPublicKey } = require("../core/security");
const { MAX_RELAY_HOPS } = require("../config/constants");
const { BloomFilterManager } = require("../state/bloom");
class MessageHandler {
constructor(peerManager, diagnostics, relayCallback, broadcastCallback) {
this.peerManager = peerManager;
this.diagnostics = diagnostics;
this.relayCallback = relayCallback;
this.broadcastCallback = broadcastCallback;
this.bloomFilter = new BloomFilterManager();
this.bloomFilter.start();
}
handleMessage(msg, sourceSocket) {
if (!validateMessage(msg)) {
return;
}
if (msg.type === "HEARTBEAT") {
this.handleHeartbeat(msg, sourceSocket);
} else if (msg.type === "LEAVE") {
this.handleLeave(msg, sourceSocket);
}
}
handleHeartbeat(msg, sourceSocket) {
this.diagnostics.increment("heartbeatsReceived");
const { id, seq, hops, nonce, sig } = msg;
if (!verifyPoW(id, nonce)) {
this.diagnostics.increment("invalidPoW");
return;
}
const stored = this.peerManager.getPeer(id);
if (stored && seq <= stored.seq) {
this.diagnostics.increment("duplicateSeq");
return;
}
if (!sig) return;
try {
let key;
if (stored && stored.key) {
key = stored.key;
} else {
if (!this.peerManager.canAcceptPeer(id)) return;
key = createPublicKey(id);
}
if (!verifySignature(`seq:${seq}`, sig, key)) {
this.diagnostics.increment("invalidSig");
return;
}
if (hops === 0) {
sourceSocket.peerId = id;
}
const wasNew = this.peerManager.addOrUpdatePeer(id, seq, key);
if (wasNew) {
this.diagnostics.increment("newPeersAdded");
this.broadcastCallback();
}
// Only relay if we haven't already relayed this message (bloom filter check)
if (hops < MAX_RELAY_HOPS && !this.bloomFilter.hasRelayed(id, seq)) {
this.bloomFilter.markRelayed(id, seq);
this.diagnostics.increment("heartbeatsRelayed");
this.relayCallback({ ...msg, hops: hops + 1 }, sourceSocket);
}
} catch (e) {
return;
}
}
handleLeave(msg, sourceSocket) {
this.diagnostics.increment("leaveMessages");
const { id, hops, sig } = msg;
if (!sig) return;
const stored = this.peerManager.getPeer(id);
if (!stored || !stored.key) return;
if (!verifySignature(`type:LEAVE:${id}`, sig, stored.key)) {
this.diagnostics.increment("invalidSig");
return;
}
if (this.peerManager.hasPeer(id)) {
this.peerManager.removePeer(id);
this.broadcastCallback();
// Use id:leave as key for LEAVE messages
if (hops < MAX_RELAY_HOPS && !this.bloomFilter.hasRelayed(id, "leave")) {
this.bloomFilter.markRelayed(id, "leave");
this.relayCallback({ ...msg, hops: hops + 1 }, sourceSocket);
}
}
}
}
const validateMessage = (msg) => {
if (!msg || typeof msg !== 'object') return false;
if (!msg.type) return false;
const msgSize = JSON.stringify(msg).length;
if (msgSize > require("../config/constants").MAX_MESSAGE_SIZE) return false;
if (msg.type === "HEARTBEAT") {
const allowedFields = ['type', 'id', 'seq', 'hops', 'nonce', 'sig'];
const fields = Object.keys(msg);
return fields.every(f => allowedFields.includes(f)) &&
msg.id && typeof msg.seq === 'number' &&
typeof msg.hops === 'number' && msg.nonce && msg.sig;
}
if (msg.type === "LEAVE") {
const allowedFields = ['type', 'id', 'hops', 'sig'];
const fields = Object.keys(msg);
return fields.every(f => allowedFields.includes(f)) &&
msg.id && typeof msg.hops === 'number' && msg.sig;
}
return false;
}
module.exports = { MessageHandler, validateMessage };
+16
View File
@@ -0,0 +1,16 @@
const relayMessage = (msg, sourceSocket, swarm, diagnostics) => {
const data = JSON.stringify(msg) + "\n";
const relayCount = swarm.connections.size - 1;
if (diagnostics) {
diagnostics.increment("bytesRelayed", data.length * relayCount);
}
for (const socket of swarm.connections) {
if (socket !== sourceSocket) {
socket.write(data);
}
}
}
module.exports = { relayMessage };
+118
View File
@@ -0,0 +1,118 @@
const Hyperswarm = require("hyperswarm");
const { signMessage } = require("../core/security");
const { TOPIC, TOPIC_NAME, HEARTBEAT_INTERVAL } = require("../config/constants");
class SwarmManager {
constructor(identity, peerManager, diagnostics, messageHandler, relayFn, broadcastFn) {
this.identity = identity;
this.peerManager = peerManager;
this.diagnostics = diagnostics;
this.messageHandler = messageHandler;
this.relayFn = relayFn;
this.broadcastFn = broadcastFn;
this.swarm = new Hyperswarm();
this.heartbeatInterval = null;
}
async start() {
this.swarm.on("connection", (socket) => this.handleConnection(socket));
const discovery = this.swarm.join(TOPIC);
await discovery.flushed();
this.startHeartbeat();
}
handleConnection(socket) {
const sig = signMessage(`seq:${this.peerManager.getSeq()}`, this.identity.privateKey);
const hello = JSON.stringify({
type: "HEARTBEAT",
id: this.identity.id,
seq: this.peerManager.getSeq(),
hops: 0,
nonce: this.identity.nonce,
sig,
});
socket.write(hello);
this.broadcastFn();
socket.on("data", (data) => {
this.diagnostics.increment("bytesReceived", data.length);
try {
const msgs = data
.toString()
.split("\n")
.filter((x) => x.trim());
for (const msgStr of msgs) {
const msg = JSON.parse(msgStr);
this.messageHandler.handleMessage(msg, socket);
}
} catch (e) {
}
});
socket.on("close", () => {
if (socket.peerId && this.peerManager.hasPeer(socket.peerId)) {
this.peerManager.removePeer(socket.peerId);
}
this.broadcastFn();
});
socket.on("error", () => { });
}
startHeartbeat() {
this.heartbeatInterval = setInterval(() => {
const seq = this.peerManager.incrementSeq();
this.peerManager.addOrUpdatePeer(this.identity.id, seq, null);
const sig = signMessage(`seq:${seq}`, this.identity.privateKey);
const heartbeat = JSON.stringify({
type: "HEARTBEAT",
id: this.identity.id,
seq,
hops: 0,
nonce: this.identity.nonce,
sig,
}) + "\n";
for (const socket of this.swarm.connections) {
socket.write(heartbeat);
}
const removed = this.peerManager.cleanupStalePeers();
if (removed > 0) {
this.broadcastFn();
}
}, HEARTBEAT_INTERVAL);
}
shutdown() {
const sig = signMessage(`type:LEAVE:${this.identity.id}`, this.identity.privateKey);
const goodbye = JSON.stringify({
type: "LEAVE",
id: this.identity.id,
hops: 0,
sig,
}) + "\n";
for (const socket of this.swarm.connections) {
socket.write(goodbye);
}
if (this.heartbeatInterval) {
clearInterval(this.heartbeatInterval);
}
setTimeout(() => {
process.exit(0);
}, 500);
}
getSwarm() {
return this.swarm;
}
}
module.exports = { SwarmManager };
+78
View File
@@ -0,0 +1,78 @@
/**
* Simple Bloom filter for message deduplication
* Prevents re-relaying messages we've already seen
*/
class BloomFilter {
constructor(size = 10000, hashCount = 3) {
this.size = size;
this.hashCount = hashCount;
this.bits = new Uint8Array(Math.ceil(size / 8));
}
_hash(str, seed) {
let h = seed;
for (let i = 0; i < str.length; i++) {
h = (h * 31 + str.charCodeAt(i)) >>> 0;
}
return h % this.size;
}
add(item) {
for (let i = 0; i < this.hashCount; i++) {
const idx = this._hash(item, i * 0x9e3779b9);
this.bits[idx >>> 3] |= (1 << (idx & 7));
}
}
has(item) {
for (let i = 0; i < this.hashCount; i++) {
const idx = this._hash(item, i * 0x9e3779b9);
if ((this.bits[idx >>> 3] & (1 << (idx & 7))) === 0) {
return false;
}
}
return true;
}
clear() {
this.bits.fill(0);
}
}
/**
* Time-bucketed bloom filter manager
* Rotates every 30 seconds to prevent unbounded growth
*/
class BloomFilterManager {
constructor() {
this.currentBloom = new BloomFilter();
this.previousBloom = new BloomFilter();
this.rotationInterval = null;
}
start() {
this.rotationInterval = setInterval(() => {
this.previousBloom = this.currentBloom;
this.currentBloom = new BloomFilter();
}, 30000);
}
stop() {
if (this.rotationInterval) {
clearInterval(this.rotationInterval);
this.rotationInterval = null;
}
}
hasRelayed(id, seq) {
const key = `${id}:${seq}`;
return this.currentBloom.has(key) || this.previousBloom.has(key);
}
markRelayed(id, seq) {
const key = `${id}:${seq}`;
this.currentBloom.add(key);
}
}
module.exports = { BloomFilter, BloomFilterManager };
+48
View File
@@ -0,0 +1,48 @@
const { DIAGNOSTICS_INTERVAL } = require("../config/constants");
class DiagnosticsManager {
constructor() {
this.stats = {
heartbeatsReceived: 0,
heartbeatsRelayed: 0,
invalidPoW: 0,
duplicateSeq: 0,
invalidSig: 0,
newPeersAdded: 0,
bytesReceived: 0,
bytesRelayed: 0,
leaveMessages: 0,
};
this.interval = null;
}
increment(key, amount = 1) {
if (this.stats.hasOwnProperty(key)) {
this.stats[key] += amount;
}
}
getStats() {
return { ...this.stats };
}
reset() {
Object.keys(this.stats).forEach(k => this.stats[k] = 0);
}
startLogging(getPeerCount, getConnectionCount) {
this.interval = setInterval(() => {
this.reset();
}, DIAGNOSTICS_INTERVAL);
}
stopLogging() {
if (this.interval) {
clearInterval(this.interval);
this.interval = null;
}
}
}
module.exports = { DiagnosticsManager };
+66
View File
@@ -0,0 +1,66 @@
class HyperLogLog {
constructor(precision = 10) {
this.precision = precision;
this.registerCount = 1 << precision;
this.registers = new Uint8Array(this.registerCount);
this.alphaMM = this._getAlpha() * this.registerCount * this.registerCount;
}
_getAlpha() {
switch (this.precision) {
case 4: return 0.673;
case 5: return 0.697;
case 6: return 0.709;
default: return 0.7213 / (1 + 1.079 / this.registerCount);
}
}
_hash(str) {
let h = 0x811c9dc5;
for (let i = 0; i < str.length; i++) {
h ^= str.charCodeAt(i);
h = (h * 0x01000193) >>> 0;
}
return h;
}
_countLeadingZeros(value, maxBits) {
if (value === 0) return maxBits;
let count = 0;
while ((value & (1 << (maxBits - 1 - count))) === 0 && count < maxBits) {
count++;
}
return count;
}
add(item) {
const hash = this._hash(item);
const registerIndex = hash >>> (32 - this.precision);
const remainingBits = hash << this.precision;
const leadingZeros = this._countLeadingZeros(remainingBits, 32 - this.precision) + 1;
if (leadingZeros > this.registers[registerIndex]) {
this.registers[registerIndex] = leadingZeros;
}
}
count() {
let harmonicSum = 0;
let zeroRegisters = 0;
for (let i = 0; i < this.registerCount; i++) {
harmonicSum += Math.pow(2, -this.registers[i]);
if (this.registers[i] === 0) zeroRegisters++;
}
let estimate = this.alphaMM / harmonicSum;
if (estimate <= 2.5 * this.registerCount && zeroRegisters > 0) {
estimate = this.registerCount * Math.log(this.registerCount / zeroRegisters);
}
return Math.round(estimate);
}
}
module.exports = { HyperLogLog };
+66
View File
@@ -0,0 +1,66 @@
const { MAX_PEERS, PEER_TIMEOUT } = require("../config/constants");
class PeerManager {
constructor() {
this.seenPeers = new Map();
this.mySeq = 0;
}
addOrUpdatePeer(id, seq, key) {
const stored = this.seenPeers.get(id);
const wasNew = !stored;
this.seenPeers.set(id, {
seq,
lastSeen: Date.now(),
key,
});
return wasNew;
}
canAcceptPeer(id) {
if (this.seenPeers.has(id)) return true;
return this.seenPeers.size < MAX_PEERS;
}
getPeer(id) {
return this.seenPeers.get(id);
}
removePeer(id) {
return this.seenPeers.delete(id);
}
hasPeer(id) {
return this.seenPeers.has(id);
}
cleanupStalePeers() {
const now = Date.now();
let removed = 0;
for (const [id, data] of this.seenPeers) {
if (now - data.lastSeen > PEER_TIMEOUT) {
this.seenPeers.delete(id);
removed++;
}
}
return removed;
}
get size() {
return this.seenPeers.size;
}
incrementSeq() {
return ++this.mySeq;
}
getSeq() {
return this.mySeq;
}
}
module.exports = { PeerManager };
+56
View File
@@ -0,0 +1,56 @@
const express = require("express");
const fs = require("fs");
const path = require("path");
const HTML_TEMPLATE = fs.readFileSync(
path.join(__dirname, "../../public/index.html"),
"utf-8"
);
const setupRoutes = (app, identity, peerManager, swarm, sseManager, diagnostics) => {
app.get("/", (req, res) => {
const count = peerManager.size;
const directPeers = swarm.getSwarm().connections.size;
const html = HTML_TEMPLATE
.replace(/\{\{COUNT\}\}/g, count)
.replace(/\{\{ID\}\}/g, identity.id.slice(0, 8) + "...")
.replace(/\{\{DIRECT\}\}/g, directPeers);
res.send(html);
});
app.get("/events", (req, res) => {
res.setHeader("Content-Type", "text/event-stream");
res.setHeader("Cache-Control", "no-cache");
res.setHeader("Connection", "keep-alive");
res.flushHeaders();
sseManager.addClient(res);
const data = JSON.stringify({
count: peerManager.size,
direct: swarm.getSwarm().connections.size,
id: identity.id,
diagnostics: diagnostics.getStats(),
});
res.write(`data: ${data}\n\n`);
req.on("close", () => {
sseManager.removeClient(res);
});
});
app.get("/api/stats", (req, res) => {
res.json({
count: peerManager.size,
direct: swarm.getSwarm().connections.size,
id: identity.id,
diagnostics: diagnostics.getStats(),
});
});
app.use(express.static(path.join(__dirname, "../../public")));
}
module.exports = { setupRoutes };
+20
View File
@@ -0,0 +1,20 @@
const express = require("express");
const { PORT } = require("../config/constants");
const { setupRoutes } = require("./routes");
const createServer = (identity, peerManager, swarm, sseManager, diagnostics) => {
const app = express();
setupRoutes(app, identity, peerManager, swarm, sseManager, diagnostics);
return app;
}
const startServer = (app, identity) => {
app.listen(PORT, () => {
console.log(`Hypermind Node running on port ${PORT}`);
console.log(`ID: ${identity.id}`);
});
}
module.exports = { createServer, startServer };
+33
View File
@@ -0,0 +1,33 @@
const { BROADCAST_THROTTLE } = require("../config/constants");
class SSEManager {
constructor() {
this.clients = new Set();
this.lastBroadcast = 0;
}
addClient(res) {
this.clients.add(res);
}
removeClient(res) {
this.clients.delete(res);
}
broadcastUpdate(data) {
const now = Date.now();
if (now - this.lastBroadcast < BROADCAST_THROTTLE) return;
this.lastBroadcast = now;
const message = JSON.stringify(data);
for (const client of this.clients) {
client.write(`data: ${message}\n\n`);
}
}
get size() {
return this.clients.size;
}
}
module.exports = { SSEManager };