igneum/tools/observer/vendor-share.mjs

233 lines
14 KiB
JavaScript

// The vendor-share metric (Counter ASIC 3.0 item 7, docs/analysis/asic-resistance-history.md section 4.3 rank 7;
// docs/benchmarks/repro.md section 8). Hash rate by vendor (nvidia, amd, apple, intel, unknown), two ways, each
// named for what it is:
//
// (a) fleet_reported: the sum of "now=X MH/s wall" from the newest STATUS line of every miner worker that reported
// to the log intake (miner_logs, labels miner-<vendor>-<machine id8>-<card index + 1>) in the last 10 minutes.
// It covers ONLY the machines that upload to the intake: the project's own fleet and any app install that
// carries the intake key. It says nothing about the rest of the network.
// (b) chain_attributed: every miner id's (vote_key_hash) share of the BLUE blocks of the last 10 minutes
// (live_blocks.color = 'blue'; pending and red blocks are left out) times the node's network hash-rate
// estimate (live_state.hashes_per_second_estimate), attributed to a vendor where the vote key is one a fleet
// worker logged ("identity N '<label>-<k>' vote_key_hash=<hex>" in its upload; the vendor from the label, and
// for miner-other-* workers from the app's "cards:" line), and to "unknown" otherwise. A key that mined and
// belongs to no reporting machine is the honest hole in the metric.
//
// Exported pure functions (tested in vendor-share.test.mjs on fabricated rows); `vendorShareTick(sql, prefix, log)`
// is the observer's one hook: it reads the three queries below (no table is written but live_state.vendor_share,
// one UPDATE of that column every 60 s, the way live_state.proving is a jsonb column of the same row) and returns
// the object it wrote. `node tools/observer/vendor-share.mjs --dry` prints the object from the live tables and
// writes nothing (DATABASE_URL from ~/.config/igneum/env).
import { readFileSync } from 'node:fs';
import { homedir } from 'node:os';
export const VENDORS = ['nvidia', 'amd', 'apple', 'intel', 'unknown'];
export const WINDOW_S = 600;
/** The vendor of a card name as the app prints it in its "cards:" line. */
export function vendorOfCard(name) {
const n = String(name || '').toLowerCase();
if (/nvidia|geforce|rtx|gtx|quadro/.test(n)) return 'nvidia';
if (/amd|radeon|rx \d|gfx\d/.test(n)) return 'amd';
if (/apple m\d|apple gpu/.test(n)) return 'apple';
if (/intel|uhd|iris|arc/.test(n)) return 'intel';
return 'unknown';
}
/** The miner label's parts: miner-<vendor word>-<machine id8>-<n>; vendor word nvidia, amd, mac or other. */
export function parseMinerLabel(label) {
const m = /^miner-(nvidia|amd|mac|other)-([0-9a-f]{8})-(\d+)$/.exec(String(label || ''));
if (!m) return null;
return { word: m[1], machine: m[2], index: Number(m[3]) - 1 };
}
/** The vendor of a worker: from the label's word, or for "other" from the machine's card list (cards: lines). */
export function vendorOfWorker(label, cards) {
const p = parseMinerLabel(label);
if (!p) return 'unknown';
if (p.word === 'nvidia' || p.word === 'amd') return p.word;
if (p.word === 'mac') return 'apple';
const card = cardOfWorker(label, cards);
return card ? vendorOfCard(card) : 'unknown';
}
/** The card name of a worker from the machine's "cards:" line: the entry at the label's index when it exists. */
export function cardOfWorker(label, cards) {
const p = parseMinerLabel(label);
if (!p || !Array.isArray(cards) || !cards.length) return null;
if (p.index >= 0 && p.index < cards.length) return cards[p.index];
return cards[0];
}
/** "1791271300 cards: NVIDIA GeForce RTX 5090 [discrete, mining] | AMD Radeon(TM) Graphics [integrated, off]" -> names. */
export function parseCardsLine(line) {
const m = /cards: (.*)$/.exec(String(line || ''));
if (!m) return [];
return m[1].split('|').map(s => s.replace(/\s*\[[^\]]*\]\s*$/, '').trim()).filter(Boolean);
}
/** One STATUS line -> { t (unix s), label, nowMhs } ; null for any other line. */
export function parseStatus(line) {
const s = String(line || '');
const m = /^(\d+(?:\.\d+)?) STATUS '([^']+)' .*? now=(\d+(?:\.\d+)?) MH\/s wall/.exec(s);
if (!m) return null;
return { t: Number(m[1]), label: m[2], nowMhs: Number(m[3]) };
}
/** Every vote key hash a worker's upload carries ("identity N '...' vote_key_hash=<64 hex>"). */
export function parseVoteKeys(lines) {
const out = new Set();
for (const m of String(lines || '').matchAll(/vote_key_hash=([0-9a-f]{64})/g)) out.add(m[1]);
return [...out];
}
function emptyByVendor(fields) {
const o = {};
for (const v of VENDORS) { o[v] = {}; for (const f of fields) o[v][f] = 0; }
return o;
}
/**
* (a) fleet-reported. workers: [{ label, machine, received_at (ms or ISO), status (the newest STATUS line), keys [] }],
* machineCards: { id8 -> [card names] }. Only uploads received within the window count; a worker whose newest STATUS
* line is older than the window inside the upload is still counted (the upload is the freshness signal) but its
* status age is reported.
*/
export function fleetReported(workers, machineCards, nowMs = Date.now(), windowS = WINDOW_S) {
const by = emptyByVendor(['mhs', 'workers']);
const rows = [];
const machines = new Set();
for (const w of workers || []) {
const recv = typeof w.received_at === 'number' ? w.received_at : new Date(w.received_at).getTime();
if (!(recv > nowMs - windowS * 1000)) continue;
const p = parseMinerLabel(w.label);
if (!p) continue;
const cards = machineCards && machineCards[p.machine];
const vendor = vendorOfWorker(w.label, cards);
const st = parseStatus(w.status);
const mhs = st ? st.nowMhs : 0;
by[vendor].mhs += mhs; by[vendor].workers += 1;
machines.add(p.machine);
rows.push({ label: w.label, machine: p.machine, vendor, card: cardOfWorker(w.label, cards), mhs: Math.round(mhs * 100) / 100,
upload_age_s: Math.round((nowMs - recv) / 1000), status_age_s: st ? Math.max(0, Math.round(nowMs / 1000 - st.t)) : null, keys: (w.keys || []).length });
}
let total = 0; for (const v of VENDORS) { by[v].mhs = Math.round(by[v].mhs * 100) / 100; total += by[v].mhs; }
total = Math.round(total * 100) / 100;
for (const v of VENDORS) by[v].share = total > 0 ? Math.round(by[v].mhs / total * 10000) / 10000 : 0;
return { by_vendor: by, total_mhs: total, workers: rows.length, machines: machines.size, rows };
}
/** vote key -> vendor from the workers' uploads (the newest upload per worker; a key seen by two workers keeps the first). */
export function keyMapOf(workers, machineCards) {
const map = {};
for (const w of workers || []) {
const p = parseMinerLabel(w.label);
if (!p) continue;
const vendor = vendorOfWorker(w.label, machineCards && machineCards[p.machine]);
for (const k of w.keys || []) if (!(k in map)) map[k] = vendor;
}
return map;
}
/**
* (b) chain-attributed. blocks: [{ vote_key_hash, blue (count of blue blocks in the window) }], keyMap: vote key -> vendor,
* networkHps: the node's estimate in H/s (null when unknown: shares are still computed, MH/s is null).
*/
export function chainAttributed(blocks, keyMap, networkHps) {
const by = emptyByVendor(['blue_blocks', 'miners']);
let total = 0;
const unknownKeys = [];
for (const b of blocks || []) {
const n = Number(b.blue) || 0;
if (!b.vote_key_hash || n <= 0) continue;
const vendor = keyMap[b.vote_key_hash] || 'unknown';
if (vendor === 'unknown') unknownKeys.push(String(b.vote_key_hash).slice(0, 8));
by[vendor].blue_blocks += n; by[vendor].miners += 1; total += n;
}
const hps = networkHps === null || networkHps === undefined || !(Number(networkHps) > 0) ? null : Number(networkHps);
for (const v of VENDORS) {
by[v].share = total > 0 ? Math.round(by[v].blue_blocks / total * 10000) / 10000 : 0;
by[v].mhs = hps === null ? null : Math.round(by[v].share * hps / 1e6 * 100) / 100;
}
return { by_vendor: by, blue_blocks_total: total, miners: (blocks || []).filter(b => b.vote_key_hash && Number(b.blue) > 0).length,
mapped_keys: Object.keys(keyMap).length, unknown_miner_ids: unknownKeys, network_mhs: hps === null ? null : Math.round(hps / 1e6 * 100) / 100 };
}
/** The whole object from rows already read. */
export function compute({ workers, machineCards, blocks, networkHps, nowMs = Date.now(), hpsSource = null }) {
const keyMap = keyMapOf(workers, machineCards);
const fleet = fleetReported(workers, machineCards, nowMs);
const chain = chainAttributed(blocks, keyMap, networkHps);
const summary = {};
for (const v of VENDORS) summary[v] = { fleet_reported_mhs: fleet.by_vendor[v].mhs, fleet_share: fleet.by_vendor[v].share, chain_share: chain.by_vendor[v].share, chain_attributed_mhs: chain.by_vendor[v].mhs };
return {
at: new Date(nowMs).toISOString(), window_s: WINDOW_S,
network_hps: networkHps === null || networkHps === undefined ? null : Number(networkHps), network_hps_source: hpsSource,
fleet_reported: { what: 'sum of now=MH/s wall from the newest STATUS line of each worker that uploaded to the intake in the window; intake-reporting machines only', ...fleet },
chain_attributed: { what: 'blue blocks per miner id in the window as a share of all blue blocks, times the network hash-rate estimate; vendor from the fleet worker that logged the vote key, else unknown', ...chain },
by_vendor: summary,
// coverage: how much of the chain's attributed rate the fleet's own reports explain (1.0 = every blue block came from a reporting worker)
coverage: chain.blue_blocks_total > 0 ? Math.round((1 - chain.by_vendor.unknown.share) * 10000) / 10000 : null,
};
}
// ---- the queries (read-only; the only write is the one column in vendorShareTick) ----
// The newest upload per worker (its last STATUS line and run id) and the vote keys that upload still carries. A worker
// logs its identity lines once at start, and an upload is the last 256 KB of its log, so after a few hours they have
// scrolled out: the keys of a run come from its first uploads (KEYS_SQL), read once per (label, run_id) and cached.
const WORKERS_SQL = `
SELECT label, machine, run_id, received_at,
reverse(substring(reverse(lines) from '[^\n]* SUTATS [^\n]*')) AS status,
array(SELECT DISTINCT m[1] FROM regexp_matches(lines, 'vote_key_hash=([0-9a-f]{64})', 'g') m) AS keys
FROM (SELECT DISTINCT ON (label) label, machine, run_id, received_at, lines FROM miner_logs
WHERE label LIKE 'miner-%' AND received_at > now() - interval '1 day' ORDER BY label, received_at DESC) s`;
const KEYS_SQL = `
SELECT array(SELECT DISTINCT m[1] FROM regexp_matches(string_agg(lines, E'\n'), 'vote_key_hash=([0-9a-f]{64})', 'g') m) AS keys
FROM (SELECT lines FROM miner_logs WHERE label = $1 AND run_id = $2 ORDER BY received_at ASC LIMIT 3) s`;
const runKeys = new Map(); // `${label}|${run_id}` -> [vote key hashes], for this process's lifetime
const CARDS_SQL = `
SELECT label, reverse(substring(reverse(lines) from '[^\n]* :sdrac [^\n]*')) AS cards
FROM (SELECT DISTINCT ON (label) label, lines FROM miner_logs
WHERE (label LIKE 'win-%' OR label LIKE 'mac-%' OR label LIKE 'linux-%') AND received_at > now() - interval '1 day' ORDER BY label, received_at DESC) s`;
const BLOCKS_SQL = (T) => `
SELECT vote_key_hash, count(*) FILTER (WHERE color = 'blue')::int AS blue
FROM ${T}live_blocks WHERE received_at > now() - interval '${WINDOW_S} seconds' AND vote_key_hash IS NOT NULL GROUP BY 1`;
const STATE_SQL = (T) => `SELECT hashes_per_second_estimate FROM ${T}live_state WHERE id = 1`;
/** Reads the live tables (read-only) and returns the object; `sql(query, params)` returns rows. T = table prefix. */
export async function readVendorShare(sql, T = '') {
const [workers, cardRows, blocks, state] = await Promise.all([sql(WORKERS_SQL), sql(CARDS_SQL), sql(BLOCKS_SQL(T)), sql(STATE_SQL(T))]);
for (const w of workers) {
const k = `${w.label}|${w.run_id}`;
if (!runKeys.has(k)) { const r = await sql(KEYS_SQL, [w.label, w.run_id]); runKeys.set(k, (r[0] && r[0].keys) || []); }
w.keys = [...new Set([...(w.keys || []), ...runKeys.get(k)])];
}
const machineCards = {};
for (const r of cardRows) { const m = /^(?:win|mac|linux)-([0-9a-f]{8})$/.exec(r.label); if (m && r.cards) machineCards[m[1]] = parseCardsLine(r.cards); }
const hps = state.length && state[0].hashes_per_second_estimate !== null ? Number(state[0].hashes_per_second_estimate) : null;
return compute({ workers, machineCards, blocks, networkHps: hps, hpsSource: 'live_state.hashes_per_second_estimate (the node\'s estimateNetworkHashesPerSecond, else blue work per second over 10 min)' });
}
/** The observer's hook: read, write live_state.vendor_share, return the object. Errors are logged, never thrown. */
export async function vendorShareTick(sql, T = '', log = () => {}) {
try {
const v = await readVendorShare(sql, T);
await sql(`UPDATE ${T}live_state SET vendor_share = $1::jsonb WHERE id = 1`, [JSON.stringify(v)]);
return v;
} catch (e) { log('vendor share failed', e.message); return null; }
}
// ---- CLI: --dry prints the live reading and writes nothing ----
if (process.argv[1] && process.argv[1].endsWith('vendor-share.mjs') && process.argv.includes('--dry')) {
const env = readFileSync(`${homedir()}/.config/igneum/env`, 'utf8');
const m = /^DATABASE_URL=(.*)$/m.exec(env);
if (!m) { console.error('DATABASE_URL not found in ~/.config/igneum/env'); process.exit(1); }
const url = m[1].trim().replace(/^['"]|['"]$/g, '');
const host = new URL(url).hostname.replace('-pooler', '');
const sql = async (query, params = []) => {
const r = await fetch(`https://${host}/sql`, { method: 'POST', headers: { 'Neon-Connection-String': url, 'Content-Type': 'application/json' }, body: JSON.stringify({ query, params }) });
const j = await r.json(); if (!r.ok) throw new Error(j.message || JSON.stringify(j)); return j.rows || [];
};
const T = (process.env.LIVE_TABLE_PREFIX || '').replace(/[^a-z0-9_]/gi, '');
console.log(JSON.stringify(await readVendorShare(sql, T), null, 1));
}