commit d4ec66a1fbc9b8fe92301bff27d4d7fa4d067dc7 Author: Alloynetworks Date: Tue Sep 29 15:57:21 2026 +0300 Initial commit diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..f321c1e --- /dev/null +++ b/.env.example @@ -0,0 +1,15 @@ +PORT=8080 +HOST=0.0.0.0 +# One or more comma-separated XAHAUD WebSocket endpoints. +XAHAUD_WS_URLS=wss://xahau.network,ws://127.0.0.1:6008 +DEDUPLICATION_TTL_MS=300000 +VALIDATOR_LIST_URL=https://vl.xahau.org/ +# Relative to the validations-proxy directory. Defaults to this value. +VERSION_STORE_PATH=data/validator-versions.json +# Optional bearer token required from downstream WebSocket clients. +# Clients can send it as Authorization: Bearer . +DOWNSTREAM_AUTH_TOKEN= +# Optional comma-separated browser origins allowed to connect. +ALLOWED_ORIGINS= +MAX_CLIENTS=1000 +CLIENT_BUFFER_LIMIT_BYTES=1048576 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5754b17 --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +node_modules/ +.env +npm-debug.log* +data/*.json diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..ff61572 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,16 @@ +FROM node:22-alpine AS dependencies +WORKDIR /app +COPY package*.json ./ +RUN npm ci --omit=dev + +FROM node:22-alpine +ENV NODE_ENV=production +WORKDIR /app +COPY --from=dependencies /app/node_modules ./node_modules +COPY package.json ./ +COPY src ./src +COPY public ./public +RUN mkdir -p /app/data && chown node:node /app/data +USER node +EXPOSE 8080 +CMD ["node", "src/server.js"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..0d267f0 --- /dev/null +++ b/README.md @@ -0,0 +1,82 @@ +# Validations Proxy + +An active-active WebSocket relay for the XAHAUD `validations` stream. It maintains connections to one or more XAHAUD nodes and broadcasts each unique `validationReceived` message to downstream clients. Validation versions are the only persisted data. + +Validation messages are forwarded byte-for-byte as text frames. The proxy does not add, remove, rename, or resolve fields, so consumers receive XAHAUD's original `master_key`, `server_version`, flag-ledger fields, and any fields added by future XAHAUD versions. + +The root HTTP page is a live validation dashboard. It shows validator master keys, manifest domains, activity, ledgers, the next 256-ledger flag boundary, and decoded XAHAUD versions. Validators published by `vl.xahau.org` appear first. The activity table defaults to base-domain ordering, so a host such as `validator.example.org` sorts under `example.org`; the full hostname is used as a tie-breaker. Browser data is memory-only and disappears when the page closes. Validator versions are persisted atomically to `data/validator-versions.json` inside the project directory; validations, domains, activity, and validator-list data are not persisted. Relative `VERSION_STORE_PATH` values are always resolved inside the project directory. The dashboard defaults to a restrained light theme with an optional non-persistent dark mode. + +The dashboard and WebSocket share the same root path: an ordinary HTTP request receives the page, while a WebSocket upgrade receives the live stream. TLS is normally terminated by a reverse proxy. + +## WebSocket streams + +- `ws://localhost:8080/` forwards XAHAUD `validationReceived` messages unchanged. Use the same root path on your deployment's public `wss://` address. There is no replay. +- `ws://localhost:8080/?dashboard=1` sends the same validation stream plus `validatorMetadata` messages used by the sample dashboard. Use the same query parameter on your deployment's public `wss://` address. A metadata frame identifies the validator with `master_key` and may contain its manifest `domain`, `listed` status from `vl.xahau.org`, explicit `offline` status, last known encoded `server_version`, latest `ledger_index`, relay-observed `last_seen` time, and `monitoring_since` timestamp. + +When a dashboard connection opens, the relay immediately sends every metadata record it currently knows. Persisted versions can therefore appear before a new validation from that validator. A previously unknown domain is sent after the validator is first observed and the relay resolves its manifest. Listing status is based on the relay's periodically refreshed `vl.xahau.org` list and is sent again when that status changes. Version metadata is sent when a validation announces `server_version`; XAHAUD normally includes this on flag-ledger validations rather than every validation. + +Every validator in the current VL is represented in dashboard metadata even when it has not sent a validation. The relay immediately sets `offline: true` for a listed validator that has never been observed. An observed validator is marked offline after falling two ledgers behind the leading ledger, with a 15-second timeout used until ledger activity is available. The relay broadcasts another `validatorMetadata` frame whenever this status changes; the browser only renders the server-provided status. Availability state is held in memory; only versions are persisted. + +The metadata is public enrichment rather than part of XAHAUD's validation message format. Consumers that require the original stream should use the root endpoint. Consumers must tolerate duplicate validations around relay restarts and fields added by future XAHAUD releases. + +When multiple upstreams deliver the same signed validation, the first original message is forwarded and subsequent copies are dropped. Deduplication uses the signed `data` field and a five-minute in-memory cache by default. This cache is empty after a restart, so consumers should still tolerate occasional duplicates. + +## Run + +Requires Node.js 20 or newer. + +```sh +npm install +cp .env.example .env +set -a; . ./.env; set +a +npm start +``` + +Set `XAHAUD_WS_URLS` to a comma-separated list of XAHAUD node endpoints. You can collect validations from the public `wss://xahau.network` endpoint, a local node, or multiple nodes for redundancy. For example: + +```sh +XAHAUD_WS_URLS=wss://xahau.network,ws://127.0.0.1:6008 +``` + +The relay remains ready while at least one node is connected. To prevent arbitrary clients from using the relay, set `DOWNSTREAM_AUTH_TOKEN` to a strong secret. Browser deployments should also set `ALLOWED_ORIGINS`. + +Connect a consumer to `ws://localhost:8080/` with an authorization header: + +```js +import WebSocket from 'ws'; + +const socket = new WebSocket('ws://localhost:8080/', { + headers: { Authorization: `Bearer ${process.env.VALIDATIONS_TOKEN}` }, +}); + +socket.on('message', (data) => console.log(JSON.parse(data))); +``` + +Web browsers cannot set an `Authorization` header in the native WebSocket API. For browser consumers, put the relay behind an authenticating reverse proxy that injects the header, or use a cookie-aware gateway. + +## Endpoints + +- `GET /` — live validation dashboard +- `GET /healthz` — process liveness +- `GET /readyz` — returns 200 while at least one XAHAUD node is connected +- `GET /status` — per-node connection state plus relay and deduplication counters +- `WS /` — live validation messages only; no replay +- `WS /?dashboard=1` — live validations plus public validator metadata + +The proxy automatically reconnects to XAHAUD with exponential backoff and jitter. Slow downstream clients are disconnected instead of allowing unbounded memory growth. + +## Test + +```sh +npm test +``` + +## Docker + +```sh +docker build -t validations-proxy . +docker run --rm -p 8080:8080 \ + -e XAHAUD_WS_URLS=ws://host.docker.internal:6008,ws://xahaud-standby:6008 \ + -e DOWNSTREAM_AUTH_TOKEN=replace-me \ + validations-proxy +``` diff --git a/package-lock.json b/package-lock.json new file mode 100644 index 0000000..a81728a --- /dev/null +++ b/package-lock.json @@ -0,0 +1,39 @@ +{ + "name": "validations-proxy", + "version": "1.0.0", + "lockfileVersion": 3, + "requires": true, + "packages": { + "": { + "name": "validations-proxy", + "version": "1.0.0", + "dependencies": { + "ws": "^8.18.3" + }, + "engines": { + "node": ">=20" + } + }, + "node_modules/ws": { + "version": "8.21.3", + "resolved": "https://registry.npmjs.org/ws/-/ws-8.21.3.tgz", + "integrity": "sha512-201TZ/kPWxoPr/OKWjquZR1SWKXcvxdH+e1xrx89b3YbmzLMFCLfnaG1HFIgWzJOEWZ7MvpK++odZufgYR50Rw==", + "license": "MIT", + "engines": { + "node": ">=10.0.0" + }, + "peerDependencies": { + "bufferutil": "^4.0.1", + "utf-8-validate": ">=5.0.2" + }, + "peerDependenciesMeta": { + "bufferutil": { + "optional": true + }, + "utf-8-validate": { + "optional": true + } + } + } + } +} diff --git a/package.json b/package.json new file mode 100644 index 0000000..006fbb8 --- /dev/null +++ b/package.json @@ -0,0 +1,18 @@ +{ + "name": "validations-proxy", + "version": "1.0.0", + "private": true, + "description": "Stateless XAHAUD validations WebSocket relay", + "type": "module", + "scripts": { + "start": "node src/server.js", + "dev": "node --watch src/server.js", + "test": "node --test" + }, + "engines": { + "node": ">=20" + }, + "dependencies": { + "ws": "^8.18.3" + } +} diff --git a/public/index.html b/public/index.html new file mode 100644 index 0000000..68626a8 --- /dev/null +++ b/public/index.html @@ -0,0 +1,595 @@ + + + + + + + Xahau Validation Pulse + + + +
+
+
+
Xahau network telemetry
+

Validation pulse

+

Live consensus validations, observed as they arrive.

+
+
+ +
+ + Connecting +
+
+
+ +
+
+

Offline VL validators

+ Assessing… +
+

A VL validator is shown offline immediately if it has never been observed, or after it misses two ledgers. A 15-second timeout is used until ledger activity is available.

+
    +
  • Waiting for validator-list data…
  • +
+
+ + + +
+
+

Live consensus

+ current ledger validations +
+
+ + Validators arranged in fixed positions around the latest ledger. Inner nodes are listed by vl.xahau.org. + +
+
+ Current ledger — + Next flag ledger — + Validators 0 +
+
+ Leading hash + Different hash + Previous ledger +
+
+
+
+ +
+
+

Validator activity

+ Waiting for data +
+
+ + + + + + + + + + + + +
Waiting for the first validation…
+
+
+ +
+ Versions are announced on flag-ledger validations. No data is persisted in this browser. + WebSocket wss://xahauvalidations.alloy.ee/ +
+
+ + + + diff --git a/src/relay.js b/src/relay.js new file mode 100644 index 0000000..4c11e2c --- /dev/null +++ b/src/relay.js @@ -0,0 +1,518 @@ +import http from 'node:http'; +import { mkdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'; +import { createHash } from 'node:crypto'; +import { dirname, resolve, sep } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { WebSocket, WebSocketServer } from 'ws'; + +const OPEN = WebSocket.OPEN; +const BASE_PATH = resolve(fileURLToPath(new URL('..', import.meta.url))); +const dashboardHtml = readFileSync(new URL('../public/index.html', import.meta.url)); +const RIPPLE_BASE58_ALPHABET = 'rpshnaf39wBUDNEGHJKLM4PQRST7VWXYZ2bcdeCg65jkm8oFqi1tuvAxyz'; + +function resolveProjectPath(path) { + const resolved = resolve(BASE_PATH, path); + if (resolved !== BASE_PATH && !resolved.startsWith(`${BASE_PATH}${sep}`)) { + throw new Error('VERSION_STORE_PATH must stay inside the validations-proxy directory'); + } + return resolved; +} + +export function nodePublicKeyFromHex(hex) { + const payload = Buffer.concat([Buffer.from([0x1c]), Buffer.from(hex, 'hex')]); + const firstHash = createHash('sha256').update(payload).digest(); + const checksum = createHash('sha256').update(firstHash).digest().subarray(0, 4); + const bytes = Buffer.concat([payload, checksum]); + let value = BigInt(`0x${bytes.toString('hex')}`); + let encoded = ''; + while (value > 0n) { + encoded = RIPPLE_BASE58_ALPHABET[Number(value % 58n)] + encoded; + value /= 58n; + } + for (const byte of bytes) { + if (byte !== 0) break; + encoded = RIPPLE_BASE58_ALPHABET[0] + encoded; + } + return encoded; +} + +function json(res, status, body) { + const data = JSON.stringify(body); + res.writeHead(status, { + 'content-type': 'application/json; charset=utf-8', + 'content-length': Buffer.byteLength(data), + 'cache-control': 'no-store', + }); + res.end(data); +} + +function bearerToken(req) { + const header = req.headers.authorization; + return header?.startsWith('Bearer ') ? header.slice(7) : undefined; +} + +export function createRelay(options = {}) { + const config = { + upstreamUrls: options.upstreamUrls ?? [options.upstreamUrl ?? 'ws://127.0.0.1:6008'], + authToken: options.authToken || undefined, + allowedOrigins: options.allowedOrigins ?? [], + maxClients: options.maxClients ?? 1000, + clientBufferLimit: options.clientBufferLimit ?? 1024 * 1024, + reconnectMinMs: options.reconnectMinMs ?? 500, + reconnectMaxMs: options.reconnectMaxMs ?? 30_000, + deduplicationTtlMs: options.deduplicationTtlMs ?? 5 * 60_000, + validatorListUrl: options.validatorListUrl === undefined ? 'https://vl.xahau.org/' : options.validatorListUrl, + validatorListRefreshMs: options.validatorListRefreshMs ?? 15 * 60_000, + versionStorePath: options.versionStorePath === undefined + ? resolveProjectPath('data/validator-versions.json') + : options.versionStorePath && resolveProjectPath(options.versionStorePath), + versionPersistDelayMs: options.versionPersistDelayMs ?? 500, + logger: options.logger ?? console, + }; + + if (config.upstreamUrls.length === 0) throw new Error('At least one XAHAUD upstream URL is required'); + + const state = { + startedAt: new Date().toISOString(), + upstreamConnected: false, + upstreamConnectedAt: null, + lastValidationAt: null, + validationsRelayed: 0, + duplicatesDropped: 0, + reconnects: 0, + upstreams: config.upstreamUrls.map((url) => ({ + url, + connected: false, + connectedAt: null, + lastValidationAt: null, + validationsReceived: 0, + reconnects: 0, + })), + }; + + let stopped = false; + let manifestRequestId = 0; + let validatorListTimer; + let availabilityTimer; + let versionPersistTimer; + let versionsDirty = false; + let latestLedgerIndex = 0; + + const clients = new Set(); + const dashboardClients = new Set(); + const seenValidations = new Map(); + const validatorMetadata = new Map(); + const manifestLookups = new Set(); + const manifestRequests = new Map(); + let listedValidators = new Set(); + + function loadVersions() { + if (!config.versionStorePath) return; + try { + const stored = JSON.parse(readFileSync(config.versionStorePath, 'utf8')); + for (const [masterKey, version] of Object.entries(stored.versions ?? {})) { + if (typeof version === 'string') validatorMetadata.set(masterKey, { server_version: version }); + } + } catch (error) { + if (error.code !== 'ENOENT') config.logger.warn?.(`Could not load version cache: ${error.message}`); + } + } + + function persistVersions() { + if (!config.versionStorePath || !versionsDirty) return; + const versions = {}; + for (const [masterKey, metadata] of validatorMetadata) { + if (metadata.server_version !== undefined) versions[masterKey] = metadata.server_version; + } + const temporaryPath = `${config.versionStorePath}.${process.pid}.tmp`; + try { + mkdirSync(dirname(config.versionStorePath), { recursive: true }); + writeFileSync(temporaryPath, `${JSON.stringify({ versions }, null, 2)}\n`, { mode: 0o600 }); + renameSync(temporaryPath, config.versionStorePath); + versionsDirty = false; + } catch (error) { + config.logger.error?.(`Could not persist version cache: ${error.message}`); + } + } + + function scheduleVersionPersist() { + versionsDirty = true; + clearTimeout(versionPersistTimer); + versionPersistTimer = setTimeout(persistVersions, config.versionPersistDelayMs); + versionPersistTimer.unref?.(); + } + + loadVersions(); + const upstreams = state.upstreams.map((status) => ({ + status, + socket: undefined, + reconnectTimer: undefined, + reconnectDelay: config.reconnectMinMs, + })); + const downstream = new WebSocketServer({ noServer: true }); + + function updateAggregateState() { + const connected = state.upstreams.filter((item) => item.connected); + state.upstreamConnected = connected.length > 0; + state.upstreamConnectedAt = connected + .map((item) => item.connectedAt) + .filter(Boolean) + .sort()[0] ?? null; + } + + const server = http.createServer((req, res) => { + const url = new URL(req.url, 'http://localhost'); + + if ((req.method === 'GET' || req.method === 'HEAD') && url.pathname === '/') { + res.writeHead(200, { + 'content-type': 'text/html; charset=utf-8', + 'content-length': dashboardHtml.length, + 'cache-control': 'no-cache', + 'x-content-type-options': 'nosniff', + 'content-security-policy': "default-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; connect-src 'self' ws: wss:; img-src 'self' data:; base-uri 'none'; frame-ancestors 'none'", + }); + return res.end(req.method === 'HEAD' ? undefined : dashboardHtml); + } + + if (req.method === 'GET' && url.pathname === '/healthz') { + return json(res, 200, { status: 'ok' }); + } + + if (req.method === 'GET' && url.pathname === '/readyz') { + return json(res, state.upstreamConnected ? 200 : 503, { + status: state.upstreamConnected ? 'ready' : 'not_ready', + upstreamConnected: state.upstreamConnected, + }); + } + + if (req.method === 'GET' && url.pathname === '/status') { + return json(res, 200, { ...state, clients: clients.size }); + } + + json(res, 404, { error: 'not_found' }); + }); + + server.on('upgrade', (req, socket, head) => { + const url = new URL(req.url, 'http://localhost'); + const origin = req.headers.origin; + const originAllowed = config.allowedOrigins.length === 0 || + (origin && config.allowedOrigins.includes(origin)); + const authorized = !config.authToken || bearerToken(req) === config.authToken; + + if (url.pathname !== '/') { + socket.write('HTTP/1.1 404 Not Found\r\nConnection: close\r\n\r\n'); + return socket.destroy(); + } + if (!authorized) { + socket.write('HTTP/1.1 401 Unauthorized\r\nConnection: close\r\n\r\n'); + return socket.destroy(); + } + if (!originAllowed) { + socket.write('HTTP/1.1 403 Forbidden\r\nConnection: close\r\n\r\n'); + return socket.destroy(); + } + if (clients.size >= config.maxClients) { + socket.write('HTTP/1.1 503 Service Unavailable\r\nConnection: close\r\n\r\n'); + return socket.destroy(); + } + + downstream.handleUpgrade(req, socket, head, (ws) => { + downstream.emit('connection', ws, req); + }); + }); + + downstream.on('connection', (ws, req) => { + clients.add(ws); + const url = new URL(req.url, 'http://localhost'); + if (url.searchParams.get('dashboard') === '1') { + dashboardClients.add(ws); + for (const [masterKey, metadata] of validatorMetadata) { + ws.send(JSON.stringify({ type: 'validatorMetadata', master_key: masterKey, ...metadata })); + } + } + const remove = () => { + clients.delete(ws); + dashboardClients.delete(ws); + }; + ws.on('close', remove); + ws.on('error', remove); + }); + + function broadcast(data) { + for (const client of clients) { + if (client.readyState !== OPEN) continue; + if (client.bufferedAmount > config.clientBufferLimit) { + client.close(1013, 'Client is too slow'); + continue; + } + client.send(data); + } + } + + function broadcastDashboardMetadata(masterKey) { + const metadata = validatorMetadata.get(masterKey); + if (!metadata) return; + const data = JSON.stringify({ type: 'validatorMetadata', master_key: masterKey, ...metadata }); + for (const client of dashboardClients) { + if (client.readyState === OPEN) client.send(data); + } + } + + function validatorIsOffline(metadata, now = Date.now()) { + if (metadata.listed !== true) return false; + if (!metadata.last_seen) return true; + const validatorLedger = Number.parseInt(metadata.ledger_index, 10); + if (latestLedgerIndex && Number.isFinite(validatorLedger)) { + return latestLedgerIndex - validatorLedger >= 2; + } + const lastSeen = Date.parse(metadata.last_seen); + return !Number.isFinite(lastSeen) || now - lastSeen >= 15_000; + } + + function refreshAvailability() { + const now = Date.now(); + for (const [masterKey, metadata] of validatorMetadata) { + if (metadata.listed !== true) continue; + const offline = validatorIsOffline(metadata, now); + if (metadata.offline !== offline) { + validatorMetadata.set(masterKey, { ...metadata, offline }); + broadcastDashboardMetadata(masterKey); + } + } + } + + async function refreshValidatorList() { + if (!config.validatorListUrl || stopped) return; + try { + const response = await fetch(config.validatorListUrl, { signal: AbortSignal.timeout(10_000) }); + if (!response.ok) throw new Error(`HTTP ${response.status}`); + const published = await response.json(); + const blob = JSON.parse(Buffer.from(published.blob, 'base64').toString('utf8')); + const refreshedAt = new Date().toISOString(); + const refreshedValidators = new Set(blob.validators.map((validator) => + nodePublicKeyFromHex(validator.validation_public_key))); + listedValidators = refreshedValidators; + + for (const masterKey of refreshedValidators) { + const metadata = validatorMetadata.get(masterKey) ?? {}; + const updated = { + ...metadata, + listed: true, + monitoring_since: metadata.listed === true && metadata.monitoring_since + ? metadata.monitoring_since + : refreshedAt, + }; + updated.offline = validatorIsOffline(updated); + validatorMetadata.set(masterKey, updated); + if (metadata.listed !== true) broadcastDashboardMetadata(masterKey); + } + + const availableUpstream = upstreams.find((upstream) => upstream.socket?.readyState === OPEN); + if (availableUpstream) { + for (const masterKey of refreshedValidators) requestManifest(availableUpstream, masterKey); + } + + for (const [masterKey, metadata] of validatorMetadata) { + const listed = refreshedValidators.has(masterKey); + if (metadata.listed !== listed) { + validatorMetadata.set(masterKey, { ...metadata, listed, offline: listed && validatorIsOffline({ ...metadata, listed }) }); + broadcastDashboardMetadata(masterKey); + } + } + } catch (error) { + config.logger.warn?.(`Could not refresh ${config.validatorListUrl}: ${error.message}`); + } finally { + if (!stopped) { + validatorListTimer = setTimeout(refreshValidatorList, config.validatorListRefreshMs); + validatorListTimer.unref?.(); + } + } + } + + function requestManifest(upstream, masterKey) { + if (!masterKey || manifestLookups.has(masterKey) || upstream.socket?.readyState !== OPEN) return; + manifestLookups.add(masterKey); + const id = `dashboard-manifest:${++manifestRequestId}`; + manifestRequests.set(id, { masterKey, socket: upstream.socket }); + upstream.socket.send(JSON.stringify({ + id, + command: 'manifest', + public_key: masterKey, + api_version: 1, + })); + } + + function handleManifestResponse(message) { + const request = manifestRequests.get(message?.id); + if (!request) return false; + manifestRequests.delete(message.id); + const domain = message?.result?.details?.domain; + const existing = validatorMetadata.get(request.masterKey) ?? {}; + validatorMetadata.set(request.masterKey, { + ...existing, + domain: typeof domain === 'string' && domain.length > 0 ? domain : null, + }); + broadcastDashboardMetadata(request.masterKey); + return true; + } + + function validationKey(message, text) { + if (typeof message.data === 'string' && message.data.length > 0) return `data:${message.data}`; + if (message.validation_public_key && message.signature && message.ledger_hash) { + return `fields:${message.validation_public_key}:${message.signature}:${message.ledger_hash}`; + } + return `text:${text}`; + } + + function isDuplicate(message, text, now = Date.now()) { + const key = validationKey(message, text); + const seenAt = seenValidations.get(key); + seenValidations.set(key, now); + + // Opportunistic cleanup keeps memory bounded without a background timer. + if (seenValidations.size % 1000 === 0) { + const cutoff = now - config.deduplicationTtlMs; + for (const [storedKey, timestamp] of seenValidations) { + if (timestamp < cutoff) seenValidations.delete(storedKey); + } + } + + return seenAt !== undefined && now - seenAt <= config.deduplicationTtlMs; + } + + function scheduleReconnect(upstream) { + if (stopped || upstream.reconnectTimer) return; + const jitteredDelay = Math.round(upstream.reconnectDelay * (0.8 + Math.random() * 0.4)); + upstream.reconnectTimer = setTimeout(() => { + upstream.reconnectTimer = undefined; + state.reconnects += 1; + upstream.status.reconnects += 1; + connectUpstream(upstream); + }, jitteredDelay); + upstream.reconnectTimer.unref?.(); + upstream.reconnectDelay = Math.min(upstream.reconnectDelay * 2, config.reconnectMaxMs); + } + + function connectUpstream(upstream) { + if (stopped) return; + const socket = new WebSocket(upstream.status.url); + upstream.socket = socket; + + socket.on('open', () => { + upstream.status.connected = true; + upstream.status.connectedAt = new Date().toISOString(); + upstream.reconnectDelay = config.reconnectMinMs; + updateAggregateState(); + socket.send(JSON.stringify({ + id: 'validations-proxy', + command: 'subscribe', + streams: ['validations'], + api_version: 1, + })); + for (const masterKey of listedValidators) requestManifest(upstream, masterKey); + config.logger.info?.(`Connected to XAHAUD at ${upstream.status.url}`); + }); + + socket.on('message', (data, isBinary) => { + if (isBinary) return; + const text = data.toString(); + let message; + try { + message = JSON.parse(text); + } catch { + return; + } + + if (handleManifestResponse(message)) return; + if (message?.type !== 'validationReceived') return; + const receivedAt = new Date().toISOString(); + latestLedgerIndex = Math.max(latestLedgerIndex, Number.parseInt(message.ledger_index, 10) || 0); + upstream.status.lastValidationAt = receivedAt; + upstream.status.validationsReceived += 1; + if (message.master_key) { + const existing = validatorMetadata.get(message.master_key) ?? {}; + const versionChanged = message.server_version !== undefined && + String(message.server_version) !== existing.server_version; + const version = message.server_version === undefined + ? existing.server_version + : String(message.server_version); + const updated = { + ...existing, + server_version: version, + ledger_index: message.ledger_index, + last_seen: receivedAt, + listed: listedValidators.has(message.master_key), + }; + updated.offline = validatorIsOffline(updated); + validatorMetadata.set(message.master_key, updated); + if (versionChanged) scheduleVersionPersist(); + if (message.server_version !== undefined || (updated.listed && existing.offline !== updated.offline)) { + broadcastDashboardMetadata(message.master_key); + } + requestManifest(upstream, message.master_key); + } + refreshAvailability(); + if (isDuplicate(message, text)) { + state.duplicatesDropped += 1; + return; + } + state.lastValidationAt = receivedAt; + state.validationsRelayed += 1; + broadcast(text); + }); + + socket.on('close', () => { + if (upstream.socket !== socket) return; + for (const [id, request] of manifestRequests) { + if (request.socket === socket) { + manifestRequests.delete(id); + manifestLookups.delete(request.masterKey); + } + } + upstream.status.connected = false; + upstream.status.connectedAt = null; + updateAggregateState(); + scheduleReconnect(upstream); + }); + + socket.on('error', (error) => { + config.logger.error?.(`XAHAUD WebSocket error (${upstream.status.url}): ${error.message}`); + }); + } + + async function start({ port = 8080, host = '0.0.0.0' } = {}) { + stopped = false; + refreshValidatorList(); + availabilityTimer = setInterval(refreshAvailability, 1000); + availabilityTimer.unref?.(); + for (const upstream of upstreams) connectUpstream(upstream); + await new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(port, host, resolve); + }); + return server.address(); + } + + async function stop() { + stopped = true; + clearTimeout(validatorListTimer); + clearInterval(availabilityTimer); + clearTimeout(versionPersistTimer); + persistVersions(); + for (const upstream of upstreams) { + clearTimeout(upstream.reconnectTimer); + upstream.socket?.close(); + } + seenValidations.clear(); + validatorMetadata.clear(); + manifestLookups.clear(); + manifestRequests.clear(); + for (const client of clients) client.close(1001, 'Server shutting down'); + await new Promise((resolve) => downstream.close(resolve)); + if (server.listening) { + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + } + } + + return { server, downstream, state, start, stop }; +} diff --git a/src/server.js b/src/server.js new file mode 100644 index 0000000..5a611cc --- /dev/null +++ b/src/server.js @@ -0,0 +1,41 @@ +import { createRelay } from './relay.js'; + +function integer(name, fallback) { + const value = process.env[name]; + if (!value) return fallback; + const parsed = Number.parseInt(value, 10); + if (!Number.isFinite(parsed) || parsed <= 0) throw new Error(`${name} must be a positive integer`); + return parsed; +} + +const relay = createRelay({ + upstreamUrls: (process.env.XAHAUD_WS_URLS || process.env.XAHAUD_WS_URL || 'ws://127.0.0.1:6008') + .split(',') + .map((value) => value.trim()) + .filter(Boolean), + authToken: process.env.DOWNSTREAM_AUTH_TOKEN, + allowedOrigins: process.env.ALLOWED_ORIGINS + ? process.env.ALLOWED_ORIGINS.split(',').map((value) => value.trim()).filter(Boolean) + : [], + maxClients: integer('MAX_CLIENTS', 1000), + clientBufferLimit: integer('CLIENT_BUFFER_LIMIT_BYTES', 1024 * 1024), + deduplicationTtlMs: integer('DEDUPLICATION_TTL_MS', 5 * 60_000), + validatorListUrl: process.env.VALIDATOR_LIST_URL || 'https://vl.xahau.org/', + versionStorePath: process.env.VERSION_STORE_PATH || undefined, +}); + +const address = await relay.start({ + port: integer('PORT', 8080), + host: process.env.HOST || '0.0.0.0', +}); + +console.log(`Validations proxy listening on ${address.address}:${address.port}`); + +async function shutdown(signal) { + console.log(`Received ${signal}; shutting down`); + await relay.stop(); + process.exit(0); +} + +process.on('SIGINT', () => shutdown('SIGINT')); +process.on('SIGTERM', () => shutdown('SIGTERM')); diff --git a/test/relay.test.js b/test/relay.test.js new file mode 100644 index 0000000..efd77cb --- /dev/null +++ b/test/relay.test.js @@ -0,0 +1,307 @@ +import assert from 'node:assert/strict'; +import { once } from 'node:events'; +import { mkdirSync, mkdtempSync, readFileSync, rmSync } from 'node:fs'; +import { createServer } from 'node:http'; +import { join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import test from 'node:test'; +import { WebSocket, WebSocketServer } from 'ws'; +import { createRelay, nodePublicKeyFromHex } from '../src/relay.js'; + +async function upstreamServer() { + const server = new WebSocketServer({ port: 0, host: '127.0.0.1' }); + await once(server, 'listening'); + const { port } = server.address(); + return { server, url: `ws://127.0.0.1:${port}` }; +} + +async function stopWsServer(server) { + for (const client of server.clients) client.terminate(); + await new Promise((resolve) => server.close(resolve)); +} + +async function waitFor(predicate, timeoutMs = 1000) { + const deadline = Date.now() + timeoutMs; + while (!predicate()) { + if (Date.now() >= deadline) throw new Error('Timed out waiting for condition'); + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +test('serves the live dashboard at the root URL', async (t) => { + const upstream = await upstreamServer(); + const relay = createRelay({ upstreamUrl: upstream.url, validatorListUrl: null, versionStorePath: null, logger: {} }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + t.after(async () => { + await relay.stop(); + await stopWsServer(upstream.server); + }); + + const response = await fetch(`http://127.0.0.1:${address.port}/`); + assert.equal(response.status, 200); + assert.match(response.headers.get('content-type'), /text\/html/); + const body = await response.text(); + assert.match(body, /Validation pulse/); + assert.match(body, /Live consensus/); + assert.match(body, /Next flag ledger/); + assert.match(body, /Offline VL validators/); + assert.match(body, /Master key/); + assert.match(body, /Domain/); + assert.match(body, /Version/); +}); + +test('sends manifest domains only as dashboard metadata', async (t) => { + const upstream = await upstreamServer(); + const relay = createRelay({ upstreamUrl: upstream.url, validatorListUrl: null, versionStorePath: null, logger: {} }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + t.after(async () => { + await relay.stop(); + await stopWsServer(upstream.server); + }); + + const upstreamConnection = await once(upstream.server, 'connection').then(([ws]) => ws); + await once(upstreamConnection, 'message'); + const dashboard = new WebSocket(`ws://127.0.0.1:${address.port}/?dashboard=1`); + await once(dashboard, 'open'); + + const validation = { + type: 'validationReceived', + ledger_index: '789', + data: 'signed-validation-domain', + master_key: 'nMasterWithDomain', + }; + const manifestMessage = once(upstreamConnection, 'message'); + upstreamConnection.send(JSON.stringify(validation)); + const [rawMessage] = await once(dashboard, 'message'); + assert.equal(rawMessage.toString(), JSON.stringify(validation)); + + const [manifestData] = await manifestMessage; + const manifestRequest = JSON.parse(manifestData); + assert.equal(manifestRequest.command, 'manifest'); + assert.equal(manifestRequest.public_key, 'nMasterWithDomain'); + + upstreamConnection.send(JSON.stringify({ + id: manifestRequest.id, + type: 'response', + result: { details: { domain: 'validator.example' } }, + })); + const [metadataData] = await once(dashboard, 'message'); + const metadata = JSON.parse(metadataData); + assert.equal(metadata.type, 'validatorMetadata'); + assert.equal(metadata.master_key, 'nMasterWithDomain'); + assert.equal(metadata.domain, 'validator.example'); +}); + +test('subscribes upstream and relays only validation messages', async (t) => { + const upstream = await upstreamServer(); + const relay = createRelay({ upstreamUrl: upstream.url, validatorListUrl: null, versionStorePath: null, logger: {} }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + t.after(async () => { + await relay.stop(); + await stopWsServer(upstream.server); + }); + + const upstreamConnection = await once(upstream.server, 'connection').then(([ws]) => ws); + const subscription = JSON.parse((await once(upstreamConnection, 'message'))[0]); + assert.deepEqual(subscription.streams, ['validations']); + + const client = new WebSocket(`ws://127.0.0.1:${address.port}/`); + await once(client, 'open'); + + upstreamConnection.send(JSON.stringify({ type: 'response', status: 'success' })); + const validation = { + type: 'validationReceived', + ledger_index: '123', + signature: 'abc', + validation_public_key: 'nSigning', + master_key: 'nMaster', + }; + upstreamConnection.send(JSON.stringify(validation)); + + const [message] = await once(client, 'message'); + assert.equal(message.toString(), JSON.stringify(validation)); + assert.equal(relay.state.validationsRelayed, 1); +}); + +test('collates two upstreams, removes duplicates, and survives one disconnect', async (t) => { + const first = await upstreamServer(); + const second = await upstreamServer(); + const firstConnected = once(first.server, 'connection'); + const secondConnected = once(second.server, 'connection'); + const relay = createRelay({ + upstreamUrls: [first.url, second.url], + validatorListUrl: null, + versionStorePath: null, + reconnectMinMs: 10_000, + logger: {}, + }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + t.after(async () => { + await relay.stop(); + await stopWsServer(first.server); + await stopWsServer(second.server); + }); + + const firstConnection = (await firstConnected)[0]; + const secondConnection = (await secondConnected)[0]; + await Promise.all([once(firstConnection, 'message'), once(secondConnection, 'message')]); + + const client = new WebSocket(`ws://127.0.0.1:${address.port}/`); + await once(client, 'open'); + + const duplicate = { + type: 'validationReceived', + ledger_index: '200', + data: 'signed-validation-one', + }; + firstConnection.send(JSON.stringify(duplicate)); + const [firstMessage] = await once(client, 'message'); + assert.equal(firstMessage.toString(), JSON.stringify(duplicate)); + secondConnection.send(JSON.stringify({ ...duplicate, master_key: 'nMaster' })); + + await waitFor(() => relay.state.duplicatesDropped === 1); + assert.equal(relay.state.validationsRelayed, 1); + assert.equal(relay.state.duplicatesDropped, 1); + + firstConnection.terminate(); + await once(firstConnection, 'close'); + const unique = { + type: 'validationReceived', + ledger_index: '201', + data: 'signed-validation-two', + }; + secondConnection.send(JSON.stringify(unique)); + const [secondMessage] = await once(client, 'message'); + assert.equal(secondMessage.toString(), JSON.stringify(unique)); + assert.equal(relay.state.upstreamConnected, true); + assert.equal(relay.state.validationsRelayed, 2); +}); + +test('rejects downstream clients without the configured bearer token', async (t) => { + const upstream = await upstreamServer(); + const relay = createRelay({ upstreamUrl: upstream.url, validatorListUrl: null, versionStorePath: null, authToken: 'secret', logger: {} }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + t.after(async () => { + await relay.stop(); + await stopWsServer(upstream.server); + }); + + const client = new WebSocket(`ws://127.0.0.1:${address.port}/`); + const [error] = await once(client, 'error'); + assert.match(error.message, /401/); +}); + +test('converts validator-list hex keys to Xahau master keys', () => { + assert.equal( + nodePublicKeyFromHex('ED02E3102D348B688CCDAF2D40FA9549E7C25EE926A15880E40D5BB0E2168A14DE'), + 'nHB45nBNgjKMssrRqaNVr2tpCq3t55J5APRRDD6ov1U41JfVFjr6', + ); +}); + +test('publishes listed validators before they send a validation', async (t) => { + const publicKey = 'ED02E3102D348B688CCDAF2D40FA9549E7C25EE926A15880E40D5BB0E2168A14DE'; + const masterKey = nodePublicKeyFromHex(publicKey); + const validatorList = createServer((req, res) => { + const blob = Buffer.from(JSON.stringify({ validators: [{ validation_public_key: publicKey }] })).toString('base64'); + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(JSON.stringify({ blob })); + }); + validatorList.listen(0, '127.0.0.1'); + await once(validatorList, 'listening'); + const listPort = validatorList.address().port; + const upstream = await upstreamServer(); + const upstreamConnected = once(upstream.server, 'connection'); + const relay = createRelay({ + upstreamUrl: upstream.url, + validatorListUrl: `http://127.0.0.1:${listPort}/`, + versionStorePath: null, + logger: {}, + }); + const address = await relay.start({ port: 0, host: '127.0.0.1' }); + const upstreamConnection = (await upstreamConnected)[0]; + await once(upstreamConnection, 'message'); + t.after(async () => { + await relay.stop(); + await stopWsServer(upstream.server); + await new Promise((resolve) => validatorList.close(resolve)); + }); + + const dashboard = new WebSocket(`ws://127.0.0.1:${address.port}/?dashboard=1`); + const metadataPromise = once(dashboard, 'message'); + await once(dashboard, 'open'); + const metadata = JSON.parse((await metadataPromise)[0]); + assert.equal(metadata.type, 'validatorMetadata'); + assert.equal(metadata.master_key, masterKey); + assert.equal(metadata.listed, true); + assert.equal(metadata.offline, true); + assert.equal(metadata.last_seen, undefined); + assert.ok(Date.parse(metadata.monitoring_since)); + + const onlineMetadata = new Promise((resolve) => { + const listener = (data) => { + const message = JSON.parse(data); + if (message.type === 'validatorMetadata' && message.master_key === masterKey && message.offline === false) { + dashboard.off('message', listener); + resolve(message); + } + }; + dashboard.on('message', listener); + }); + upstreamConnection.send(JSON.stringify({ + type: 'validationReceived', + master_key: masterKey, + ledger_index: '901', + data: 'listed-validator-is-online', + })); + const online = await onlineMetadata; + assert.equal(online.last_seen.length > 0, true); +}); + +test('persists and reloads validator versions', async (t) => { + const dataDirectory = fileURLToPath(new URL('../data', import.meta.url)); + mkdirSync(dataDirectory, { recursive: true }); + const directory = mkdtempSync(join(dataDirectory, 'test-')); + const storePath = join(directory, 'versions.json'); + const upstream = await upstreamServer(); + t.after(async () => { + await stopWsServer(upstream.server); + rmSync(directory, { recursive: true, force: true }); + }); + + const firstRelay = createRelay({ + upstreamUrl: upstream.url, + validatorListUrl: null, + versionStorePath: storePath, + logger: {}, + }); + await firstRelay.start({ port: 0, host: '127.0.0.1' }); + const firstUpstream = await once(upstream.server, 'connection').then(([ws]) => ws); + await once(firstUpstream, 'message'); + firstUpstream.send(JSON.stringify({ + type: 'validationReceived', + master_key: 'nPersistedMaster', + ledger_index: '900', + server_version: '1745991531128425009', + })); + await waitFor(() => firstRelay.state.validationsRelayed === 1); + await firstRelay.stop(); + + const stored = JSON.parse(readFileSync(storePath, 'utf8')); + assert.equal(stored.versions.nPersistedMaster, '1745991531128425009'); + + const secondRelay = createRelay({ + upstreamUrl: upstream.url, + validatorListUrl: null, + versionStorePath: storePath, + logger: {}, + }); + const address = await secondRelay.start({ port: 0, host: '127.0.0.1' }); + const dashboard = new WebSocket(`ws://127.0.0.1:${address.port}/?dashboard=1`); + const metadataPromise = once(dashboard, 'message'); + await once(dashboard, 'open'); + const [metadataData] = await metadataPromise; + const metadata = JSON.parse(metadataData); + assert.equal(metadata.master_key, 'nPersistedMaster'); + assert.equal(metadata.server_version, '1745991531128425009'); + await secondRelay.stop(); +});