233 lines
14 KiB
JavaScript
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));
|
|
}
|