Initial commit

This commit is contained in:
2026-09-29 15:57:21 +03:00
commit d4ec66a1fb
10 changed files with 1635 additions and 0 deletions

15
.env.example Normal file
View File

@@ -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 <token>.
DOWNSTREAM_AUTH_TOKEN=
# Optional comma-separated browser origins allowed to connect.
ALLOWED_ORIGINS=
MAX_CLIENTS=1000
CLIENT_BUFFER_LIMIT_BYTES=1048576

4
.gitignore vendored Normal file
View File

@@ -0,0 +1,4 @@
node_modules/
.env
npm-debug.log*
data/*.json

16
Dockerfile Normal file
View File

@@ -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"]

82
README.md Normal file
View File

@@ -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
```

39
package-lock.json generated Normal file
View File

@@ -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
}
}
}
}
}

18
package.json Normal file
View File

@@ -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"
}
}

595
public/index.html Normal file
View File

@@ -0,0 +1,595 @@
<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<meta name="color-scheme" content="light dark">
<title>Xahau Validation Pulse</title>
<style>
:root {
color-scheme: light;
--bg: #f7f8f7;
--line: #d9dfdc;
--text: #17201d;
--muted: #62706a;
--green: #167a55;
--cyan: #286c73;
--amber: #946b16;
--red: #a13f3f;
--mono: ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas, monospace;
--sans: Inter, ui-sans-serif, system-ui, -apple-system, BlinkMacSystemFont, "Segoe UI", sans-serif;
}
:root[data-theme="dark"] {
color-scheme: dark;
--bg: #111614;
--line: #2d3935;
--text: #edf2f0;
--muted: #9aa9a3;
--green: #58c99a;
--cyan: #73bdc3;
--amber: #d3ad61;
--red: #e07b7b;
}
* { box-sizing: border-box; }
html { background: var(--bg); }
body {
margin: 0;
min-width: 320px;
min-height: 100vh;
color: var(--text);
background: var(--bg);
font-family: var(--sans);
}
.shell { width: min(1320px, calc(100% - 40px)); margin: 0 auto; padding: 30px 0 52px; }
header { display: flex; justify-content: space-between; align-items: flex-start; gap: 24px; margin-bottom: 34px; }
.eyebrow { color: var(--green); font: 600 12px/1.4 var(--mono); letter-spacing: .16em; text-transform: uppercase; }
h1 { margin: 8px 0 6px; font-size: clamp(28px, 4vw, 42px); font-weight: 520; letter-spacing: -.035em; line-height: 1.1; }
.subtitle { margin: 0; color: var(--muted); font-size: 15px; }
.header-actions { display: flex; align-items: center; gap: 20px; }
.theme-toggle { border: 1px solid var(--line); background: transparent; color: var(--text); border-radius: 4px; padding: 7px 10px; font: 12px/1 var(--sans); cursor: pointer; }
.theme-toggle:hover { background: color-mix(in srgb, var(--text) 5%, transparent); }
.connection { display: flex; align-items: center; gap: 9px; color: var(--muted); font-size: 13px; white-space: nowrap; }
.dot { width: 9px; height: 9px; border-radius: 50%; background: var(--amber); box-shadow: 0 0 0 5px rgba(241, 197, 109, .08); }
.connection.live .dot { background: var(--green); box-shadow: none; }
.connection.offline .dot { background: var(--red); box-shadow: 0 0 0 5px rgba(255, 125, 125, .08); }
.stream-callout { display: flex; justify-content: space-between; align-items: center; gap: 24px; margin-bottom: 34px; padding: 18px 20px; border: 1px solid var(--line); background: color-mix(in srgb, var(--green) 4%, var(--bg)); }
.stream-callout p { margin: 3px 0 0; color: var(--muted); font-size: 13px; }
.stream-callout .metadata-note { margin-top: 8px; font-size: 12px; }
.stream-label { color: var(--green); font-size: 13px; font-weight: 600; }
.stream-endpoints { display: grid; gap: 7px; }
.stream-endpoint { display: flex; justify-content: space-between; gap: 18px; padding: 8px 10px; color: var(--text); background: var(--bg); border: 1px solid var(--line); font: 500 12px/1 var(--mono); white-space: nowrap; }
.stream-endpoint small { color: var(--muted); font: 11px/1 var(--sans); }
.consensus-section { margin: 16px 0 42px; }
.section-head { display: flex; justify-content: space-between; align-items: baseline; gap: 16px; margin-bottom: 15px; }
h2 { margin: 0; font-size: 15px; font-weight: 550; letter-spacing: .01em; }
.section-note { color: var(--muted); font: 12px/1.4 var(--mono); }
.consensus-graphic { border-top: 1px solid var(--line); border-bottom: 1px solid var(--line); }
.consensus-svg { display: block; width: min(100%, 600px); height: auto; margin: 0 auto; }
.ring { fill: none; stroke: var(--line); stroke-width: 1; }
.spoke { stroke: var(--line); stroke-width: 1; }
.spoke.agree { stroke: color-mix(in srgb, var(--green) 45%, var(--line)); }
.spoke.differ { stroke: color-mix(in srgb, var(--red) 55%, var(--line)); }
.validator-node { stroke: var(--bg); stroke-width: 2; }
.arrival { fill: none; stroke: var(--green); stroke-width: 1.5; animation: arrival 1.4s ease-out forwards; pointer-events: none; }
.ring-label { fill: var(--muted); font: 11px/1 var(--sans); text-anchor: middle; }
.ring-label { letter-spacing: .08em; text-transform: uppercase; }
.consensus-meta { display: flex; justify-content: space-between; align-items: center; gap: 20px; padding: 13px 0; border-top: 1px solid var(--line); color: var(--muted); font-size: 12px; }
.ring-stats { display: flex; flex-wrap: wrap; gap: 10px 24px; }
.ring-stat strong { margin-left: 6px; color: var(--text); font: 500 13px/1 var(--mono); }
.legend { display: flex; flex-wrap: wrap; justify-content: flex-end; gap: 14px; white-space: nowrap; }
.legend-item { display: inline-flex; align-items: center; gap: 6px; }
.legend-dot { width: 8px; height: 8px; border-radius: 50%; background: var(--line); }
.legend-dot.agree { background: var(--green); }
.legend-dot.differ { background: var(--red); }
@keyframes arrival { from { r: 8; opacity: .8; } to { r: 16; opacity: 0; } }
.offline-summary { margin: 0 0 24px; padding: 17px 20px; border: 1px solid color-mix(in srgb, var(--red) 45%, var(--line)); background: color-mix(in srgb, var(--red) 5%, var(--bg)); }
.offline-heading { display: flex; justify-content: space-between; align-items: baseline; gap: 16px; }
.offline-count { color: var(--red); font: 500 13px/1 var(--mono); }
.offline-description { margin: 6px 0 0; color: var(--muted); font-size: 12px; }
.offline-list { display: grid; grid-template-columns: repeat(auto-fit, minmax(310px, 1fr)); gap: 8px 24px; margin: 15px 0 0; padding: 0; list-style: none; }
.offline-list li { display: flex; justify-content: space-between; gap: 12px; padding-top: 9px; border-top: 1px solid var(--line); min-width: 0; }
.offline-name { overflow: hidden; color: var(--text); font: 12px/1.4 var(--mono); text-overflow: ellipsis; white-space: nowrap; }
.offline-age { flex: none; color: var(--red); font-size: 12px; }
.table-wrap { overflow-x: auto; border-top: 1px solid var(--line); }
table { width: 100%; border-collapse: collapse; min-width: 850px; }
th { padding: 13px 14px; color: var(--muted); font-size: 11px; font-weight: 520; letter-spacing: .1em; text-align: left; text-transform: uppercase; border-bottom: 1px solid var(--line); }
.sort-button { display: inline-flex; align-items: center; gap: 6px; padding: 0; border: 0; color: inherit; background: none; font: inherit; letter-spacing: inherit; text-transform: inherit; cursor: pointer; }
.sort-button::after { content: '↕'; color: var(--line); font-size: 10px; }
th[aria-sort="ascending"] .sort-button::after { content: '↑'; color: var(--text); }
th[aria-sort="descending"] .sort-button::after { content: '↓'; color: var(--text); }
th:last-child .sort-button { margin-left: auto; }
td { padding: 15px 14px; border-bottom: 1px solid var(--line); vertical-align: middle; }
th:first-child, td:first-child { padding-left: 0; }
th:last-child, td:last-child { padding-right: 0; text-align: right; }
tbody tr { transition: background .25s ease; }
tbody tr.fresh { background: color-mix(in srgb, var(--green) 6%, transparent); }
tbody tr.offline { background: color-mix(in srgb, var(--red) 4%, transparent); }
.key { font: 12px/1.5 var(--mono); color: var(--text); white-space: nowrap; }
.key-short { display: none; }
.vl-mark { display: inline-block; margin-left: 9px; padding: 3px 6px; border: 1px solid color-mix(in srgb, var(--green) 50%, var(--line)); color: var(--green); font: 700 11px/1 var(--sans); letter-spacing: .07em; }
.version { font: 12px/1.4 var(--mono); color: var(--cyan); white-space: nowrap; }
.version.pending { color: var(--muted); font-family: var(--sans); }
.age { color: var(--muted); font-size: 12px; white-space: nowrap; }
.empty { padding: 44px 0; color: var(--muted); text-align: center; font-size: 14px; }
footer { display: flex; justify-content: space-between; gap: 20px; margin-top: 24px; color: var(--muted); font-size: 11px; }
footer code { font-family: var(--mono); color: #b9cec7; }
@media (max-width: 720px) {
.shell { width: min(100% - 24px, 1440px); padding-top: 22px; }
header { display: block; margin-bottom: 25px; }
.connection { margin-top: 16px; }
.stream-callout { display: block; margin-bottom: 25px; padding: 15px; }
.stream-endpoints { margin-top: 13px; overflow-x: auto; }
.consensus-meta { display: block; }
.legend { justify-content: flex-start; margin-top: 10px; }
.key-full { display: none; }
.key-short { display: inline; }
footer { display: block; line-height: 1.7; }
}
@media (prefers-reduced-motion: reduce) {
*, *::before, *::after { transition: none !important; }
.arrival { animation: none; display: none; }
}
</style>
</head>
<body>
<main class="shell">
<header>
<div>
<div class="eyebrow">Xahau network telemetry</div>
<h1>Validation pulse</h1>
<p class="subtitle">Live consensus validations, observed as they arrive.</p>
</div>
<div class="header-actions">
<button id="theme-toggle" class="theme-toggle" type="button" aria-pressed="false">Dark mode</button>
<div id="connection" class="connection" role="status" aria-live="polite">
<span class="dot" aria-hidden="true"></span>
<span id="connection-text">Connecting</span>
</div>
</div>
</header>
<section class="offline-summary" aria-labelledby="offline-title">
<div class="offline-heading">
<h2 id="offline-title">Offline VL validators</h2>
<span id="offline-count" class="offline-count">Assessing…</span>
</div>
<p class="offline-description">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.</p>
<ul id="offline-list" class="offline-list">
<li><span class="offline-name">Waiting for validator-list data…</span></li>
</ul>
</section>
<aside class="stream-callout" aria-label="Public validation WebSocket">
<div>
<div class="stream-label">Use the live validation stream</div>
<p>This dashboard is one example. Connect your own application directly to the public WebSocket.</p>
<p class="metadata-note"><strong>Public metadata</strong> means the validator's master key, manifest domain, <code>vl.xahau.org</code> listing and offline status, last known server version, latest ledger and last-seen time. With <code>?dashboard=1</code>, known metadata is sent when you connect and updated as it changes; versions may arrive immediately from the relay's version store.</p>
</div>
<div class="stream-endpoints">
<code class="stream-endpoint"><span>wss://xahauvalidations.alloy.ee/</span><small>unchanged validations</small></code>
<code class="stream-endpoint"><span>wss://xahauvalidations.alloy.ee/?dashboard=1</span><small>+ public metadata</small></code>
</div>
</aside>
<section class="consensus-section" aria-labelledby="consensus-title">
<div class="section-head">
<h2 id="consensus-title">Live consensus</h2>
<span class="section-note">current ledger validations</span>
</div>
<div class="consensus-graphic">
<svg id="consensus-ring" class="consensus-svg" viewBox="0 0 600 360" role="img" aria-labelledby="consensus-title consensus-description">
<desc id="consensus-description">Validators arranged in fixed positions around the latest ledger. Inner nodes are listed by vl.xahau.org.</desc>
</svg>
<div class="consensus-meta">
<div class="ring-stats" aria-label="Live validation metrics">
<span class="ring-stat">Current ledger <strong id="ledger">—</strong></span>
<span class="ring-stat">Next flag ledger <strong id="next-flag-ledger">—</strong> <span id="flag-ledgers-remaining"></span></span>
<span class="ring-stat">Validators <strong id="validators">0</strong></span>
</div>
<div class="legend" aria-label="Consensus legend">
<span class="legend-item"><span class="legend-dot agree"></span>Leading hash</span>
<span class="legend-item"><span class="legend-dot differ"></span>Different hash</span>
<span class="legend-item"><span class="legend-dot"></span>Previous ledger</span>
</div>
</div>
</div>
</section>
<section aria-labelledby="validators-title">
<div class="section-head">
<h2 id="validators-title">Validator activity</h2>
<span id="updated" class="section-note">Waiting for data</span>
</div>
<div class="table-wrap">
<table>
<thead>
<tr>
<th><button class="sort-button" type="button" data-sort="masterKey">Master key</button></th>
<th aria-sort="ascending"><button class="sort-button" type="button" data-sort="domain">Domain</button></th>
<th><button class="sort-button" type="button" data-sort="version">Version</button></th>
<th><button class="sort-button" type="button" data-sort="lastSeen">Last seen</button></th>
</tr>
</thead>
<tbody id="validator-rows">
<tr><td class="empty" colspan="4">Waiting for the first validation…</td></tr>
</tbody>
</table>
</div>
</section>
<footer>
<span>Versions are announced on flag-ledger validations. No data is persisted in this browser.</span>
<span>WebSocket <code>wss://xahauvalidations.alloy.ee/</code></span>
</footer>
</main>
<script>
(() => {
const MAX_ROWS = 100;
const validators = new Map();
const metadata = new Map();
let latestLedger = 0;
let reconnectDelay = 1000;
let sortKey = 'domain';
let sortDirection = 1;
let socket;
const elements = {
connection: document.getElementById('connection'),
connectionText: document.getElementById('connection-text'),
validators: document.getElementById('validators'),
ledger: document.getElementById('ledger'),
nextFlagLedger: document.getElementById('next-flag-ledger'),
flagLedgersRemaining: document.getElementById('flag-ledgers-remaining'),
consensusRing: document.getElementById('consensus-ring'),
offlineCount: document.getElementById('offline-count'),
offlineList: document.getElementById('offline-list'),
rows: document.getElementById('validator-rows'),
updated: document.getElementById('updated'),
themeToggle: document.getElementById('theme-toggle'),
sortButtons: Array.from(document.querySelectorAll('.sort-button')),
};
function setConnection(mode, label) {
elements.connection.className = `connection ${mode}`;
elements.connectionText.textContent = label;
}
function decodeVersion(raw) {
if (raw === undefined || raw === null || raw === '') return null;
try {
const encoded = BigInt(String(raw));
const implementation = Number((encoded >> 48n) & 0xffffn);
if (implementation !== 0x183b) return `Unknown · ${raw}`;
const year = Number((encoded >> 40n) & 0xffn) + 2023;
const minor = Number((encoded >> 32n) & 0xffn);
const patch = Number((encoded >> 24n) & 0xffn);
const stage = Number((encoded >> 16n) & 0xffn);
const build = Number(encoded & 0xffffn);
const kind = stage & 0xc0;
const number = stage & 0x3f;
const prerelease = kind === 0x80 ? `-rc${number}` : kind === 0x40 ? `-b${number}` : '';
return `${year}.${minor}.${patch}${prerelease}+${build}`;
} catch {
return `Unknown · ${raw}`;
}
}
const SVG_NS = 'http://www.w3.org/2000/svg';
function svgElement(name, attributes = {}) {
const element = document.createElementNS(SVG_NS, name);
for (const [key, value] of Object.entries(attributes)) element.setAttribute(key, value);
return element;
}
function stableAngle(masterKey) {
let hash = 2166136261;
for (let index = 0; index < masterKey.length; index += 1) {
hash ^= masterKey.charCodeAt(index);
hash = Math.imul(hash, 16777619);
}
return ((hash >>> 0) / 4294967296) * Math.PI * 2 - Math.PI / 2;
}
function renderConsensus() {
const svg = elements.consensusRing;
svg.replaceChildren();
const description = svgElement('desc', { id: 'consensus-description' });
description.textContent = 'Validators arranged in fixed positions around the latest ledger. Inner nodes are listed by vl.xahau.org.';
svg.appendChild(description);
const centerX = 300;
const centerY = 180;
const current = Array.from(validators.values()).filter((validator) => Number(validator.ledger) === latestLedger);
const hashCounts = new Map();
for (const validator of current) {
if (validator.hash) hashCounts.set(validator.hash, (hashCounts.get(validator.hash) || 0) + 1);
}
const leadingHash = Array.from(hashCounts.entries()).sort((a, b) => b[1] - a[1])[0]?.[0];
const outerRing = svgElement('circle', { cx: centerX, cy: centerY, r: 142, class: 'ring' });
const innerRing = svgElement('circle', { cx: centerX, cy: centerY, r: 96, class: 'ring' });
svg.append(outerRing, innerRing);
const innerLabel = svgElement('text', { x: centerX, y: 70, class: 'ring-label' });
innerLabel.textContent = 'VL validators';
svg.appendChild(innerLabel);
for (const validator of Array.from(validators.values()).sort((a, b) => a.masterKey.localeCompare(b.masterKey))) {
const angle = stableAngle(validator.masterKey);
const radius = validator.listed ? 96 : 142;
const x = centerX + Math.cos(angle) * radius;
const y = centerY + Math.sin(angle) * radius;
const onCurrentLedger = Number(validator.ledger) === latestLedger;
const agrees = onCurrentLedger && (!leadingHash || validator.hash === leadingHash);
const differs = onCurrentLedger && leadingHash && validator.hash && validator.hash !== leadingHash;
const color = agrees ? 'var(--green)' : differs ? 'var(--red)' : 'var(--line)';
const spoke = svgElement('line', {
x1: centerX,
y1: centerY,
x2: x,
y2: y,
class: `spoke${agrees ? ' agree' : differs ? ' differ' : ''}`,
});
svg.appendChild(spoke);
if (Date.now() - validator.lastSeen < 1500 && onCurrentLedger) {
svg.appendChild(svgElement('circle', { cx: x, cy: y, r: 8, class: 'arrival' }));
}
const node = svgElement('circle', {
cx: x,
cy: y,
r: validator.listed ? 7 : 6,
fill: color,
class: 'validator-node',
});
svg.appendChild(node);
}
}
function relativeTime(timestamp) {
if (!Number.isFinite(timestamp)) return 'Never observed';
const seconds = Math.max(0, Math.floor((Date.now() - timestamp) / 1000));
if (seconds < 5) return 'now';
if (seconds < 60) return `${seconds}s ago`;
return `${Math.floor(seconds / 60)}m ago`;
}
function domainSortValue(domain) {
if (!domain) return '\uffff';
const normalized = domain.toLowerCase().replace(/\.$/, '');
const labels = normalized.split('.');
const baseDomain = labels.length > 1 ? labels.slice(-2).join('.') : normalized;
return `${baseDomain}\u0000${normalized}`;
}
function compareValidators(a, b) {
const listedOrder = Number(Boolean(b.listed)) - Number(Boolean(a.listed));
if (listedOrder) return listedOrder;
let comparison = 0;
if (sortKey === 'lastSeen') {
comparison = (a.lastSeen || 0) - (b.lastSeen || 0);
} else if (sortKey === 'domain') {
comparison = domainSortValue(a.domain).localeCompare(domainSortValue(b.domain));
} else {
comparison = String(a[sortKey] || '').localeCompare(String(b[sortKey] || ''));
}
return comparison * sortDirection || a.masterKey.localeCompare(b.masterKey);
}
function updateSortIndicators() {
for (const button of elements.sortButtons) {
const active = button.dataset.sort === sortKey;
const header = button.closest('th');
if (active) header.setAttribute('aria-sort', sortDirection === 1 ? 'ascending' : 'descending');
else header.removeAttribute('aria-sort');
}
}
function isOffline(validator) {
return validator.listed === true && validator.offline === true;
}
function renderOfflineSummary() {
const listed = Array.from(validators.values()).filter((validator) => validator.listed);
const offline = listed.filter(isOffline)
.sort((a, b) => domainSortValue(a.domain).localeCompare(domainSortValue(b.domain)) || a.masterKey.localeCompare(b.masterKey));
elements.offlineCount.textContent = listed.length === 0
? 'VL unavailable'
: `${offline.length} of ${listed.length} offline`;
elements.offlineList.replaceChildren();
if (listed.length === 0) {
const item = document.createElement('li');
item.innerHTML = '<span class="offline-name">Waiting for validator-list data…</span>';
elements.offlineList.appendChild(item);
return;
}
if (offline.length === 0) {
const item = document.createElement('li');
item.innerHTML = '<span class="offline-name">No VL validators are currently offline.</span>';
elements.offlineList.appendChild(item);
return;
}
for (const validator of offline) {
const item = document.createElement('li');
const name = document.createElement('span');
name.className = 'offline-name';
name.textContent = validator.domain ? `${validator.domain} · ${validator.masterKey}` : validator.masterKey;
name.title = name.textContent;
const age = document.createElement('span');
age.className = 'offline-age';
age.textContent = relativeTime(validator.lastSeen);
item.append(name, age);
elements.offlineList.appendChild(item);
}
}
function renderMetrics() {
const observed = Array.from(validators.values()).filter((validator) => Number.isFinite(validator.lastSeen));
elements.validators.textContent = observed.length.toLocaleString();
elements.ledger.textContent = latestLedger ? latestLedger.toLocaleString() : '—';
if (latestLedger) {
const nextFlagLedger = latestLedger + (256 - (latestLedger % 256));
const remaining = nextFlagLedger - latestLedger;
elements.nextFlagLedger.textContent = nextFlagLedger.toLocaleString();
elements.flagLedgersRemaining.textContent = `(${remaining} ledger${remaining === 1 ? '' : 's'})`;
} else {
elements.nextFlagLedger.textContent = '—';
elements.flagLedgersRemaining.textContent = '';
}
}
function appendCell(row, className, value) {
const cell = document.createElement('td');
cell.className = className;
cell.textContent = value;
row.appendChild(cell);
return cell;
}
function renderRows() {
const ordered = Array.from(validators.values())
.sort(compareValidators)
.slice(0, MAX_ROWS);
elements.rows.replaceChildren();
if (ordered.length === 0) {
const row = document.createElement('tr');
appendCell(row, 'empty', 'Waiting for the first validation…').colSpan = 4;
elements.rows.appendChild(row);
return;
}
for (const validator of ordered) {
const row = document.createElement('tr');
if (isOffline(validator)) row.className = 'offline';
else if (Date.now() - validator.lastSeen < 1200) row.className = 'fresh';
const keyCell = document.createElement('td');
keyCell.className = 'key';
const full = document.createElement('span');
full.className = 'key-full';
full.textContent = validator.masterKey;
const short = document.createElement('span');
short.className = 'key-short';
short.textContent = `${validator.masterKey.slice(0, 12)}…${validator.masterKey.slice(-8)}`;
keyCell.append(full, short);
if (validator.listed) {
const mark = document.createElement('span');
mark.className = 'vl-mark';
mark.textContent = 'VL';
mark.setAttribute('aria-label', 'Listed by vl.xahau.org');
keyCell.appendChild(mark);
}
row.appendChild(keyCell);
appendCell(row, validator.domain ? 'version' : 'version pending', validator.domain || '—');
appendCell(row, validator.version ? 'version' : 'version pending', validator.version || 'Awaiting flag ledger');
appendCell(row, 'age', relativeTime(validator.lastSeen));
elements.rows.appendChild(row);
}
}
function handleValidation(message) {
if (message.type === 'validatorMetadata') {
const details = metadata.get(message.master_key) || {};
if (Object.prototype.hasOwnProperty.call(message, 'domain')) details.domain = message.domain;
if (Object.prototype.hasOwnProperty.call(message, 'listed')) details.listed = message.listed;
if (Object.prototype.hasOwnProperty.call(message, 'offline')) details.offline = message.offline;
if (message.server_version !== undefined) details.version = decodeVersion(message.server_version);
if (message.ledger_index !== undefined) details.ledger = message.ledger_index;
if (message.last_seen !== undefined) details.lastSeen = Date.parse(message.last_seen);
if (message.monitoring_since !== undefined) details.monitoringSince = Date.parse(message.monitoring_since);
metadata.set(message.master_key, details);
const current = validators.get(message.master_key) || (details.listed ? { masterKey: message.master_key } : null);
if (current) {
Object.assign(current, details);
if (!current.listed && !Number.isFinite(current.lastSeen)) validators.delete(message.master_key);
else validators.set(message.master_key, current);
renderMetrics();
renderConsensus();
renderOfflineSummary();
renderRows();
}
return;
}
if (message.type !== 'validationReceived') return;
const masterKey = message.master_key;
if (!masterKey) return;
const now = Date.now();
const ledger = Number.parseInt(message.ledger_index, 10) || 0;
latestLedger = Math.max(latestLedger, ledger);
const decoded = decodeVersion(message.server_version);
const current = validators.get(masterKey) || { masterKey, ...(metadata.get(masterKey) || {}) };
current.lastSeen = now;
current.ledger = message.ledger_index || current.ledger;
current.hash = message.ledger_hash || current.hash;
if (decoded) current.version = decoded;
validators.set(masterKey, current);
renderMetrics();
elements.updated.textContent = `Updated ${new Date(now).toLocaleTimeString()}`;
renderConsensus();
renderOfflineSummary();
renderRows();
}
function connect() {
setConnection('', 'Connecting');
const protocol = location.protocol === 'https:' ? 'wss:' : 'ws:';
socket = new WebSocket(`${protocol}//${location.host}/?dashboard=1`);
socket.addEventListener('open', () => {
reconnectDelay = 1000;
setConnection('live', 'Live');
});
socket.addEventListener('message', (event) => {
try { handleValidation(JSON.parse(event.data)); } catch {}
});
socket.addEventListener('close', () => {
setConnection('offline', `Reconnecting in ${Math.ceil(reconnectDelay / 1000)}s`);
window.setTimeout(connect, reconnectDelay);
reconnectDelay = Math.min(reconnectDelay * 2, 30000);
});
socket.addEventListener('error', () => socket.close());
}
window.setInterval(() => {
renderConsensus();
renderOfflineSummary();
renderRows();
}, 1000);
elements.themeToggle.addEventListener('click', () => {
const dark = document.documentElement.dataset.theme !== 'dark';
document.documentElement.dataset.theme = dark ? 'dark' : '';
elements.themeToggle.textContent = dark ? 'Light mode' : 'Dark mode';
elements.themeToggle.setAttribute('aria-pressed', String(dark));
});
for (const button of elements.sortButtons) {
button.addEventListener('click', () => {
const nextKey = button.dataset.sort;
if (sortKey === nextKey) sortDirection *= -1;
else {
sortKey = nextKey;
sortDirection = 1;
}
updateSortIndicators();
renderRows();
});
}
updateSortIndicators();
connect();
})();
</script>
</body>
</html>

518
src/relay.js Normal file
View File

@@ -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 };
}

41
src/server.js Normal file
View File

@@ -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'));

307
test/relay.test.js Normal file
View File

@@ -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();
});