igneum/pool/src/api.rs

353 lines
21 KiB
Rust

//! The stats API in the shape WhatToMine and MiningPoolStats read, and the pool's page. A small HTTP/1.1 server:
//! GET only, JSON, `Connection: close`, CORS open. The field lists are the contract the fixture test asserts.
use crate::pool::Pool;
use crate::state::{ign, unix_ms};
use serde_json::{json, Value};
use std::sync::Arc;
use std::time::Duration;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
pub const STATS_FIELDS: &[&str] = &["ok", "now", "pool", "network", "source"];
pub const POOL_FIELDS: &[&str] = &[
"name", "url", "address", "algorithm", "scheme", "fee_percent", "min_payout_ign", "pplns_window_blocks", "hashrate", "hashrate_unit", "hashrate_1h",
"miners", "workers", "shares", "blocks_24h", "blocks_confirmed_24h", "blocks_orphaned_24h", "blocks_total", "blocks_confirmed_total", "blocks_orphaned_total",
"last_block", "effort_current", "luck_24h", "paid_24h_ign", "paid_total_ign", "pending_balance_ign", "pool_fee_total_ign", "uptime_s", "share_check_ms", "dry_run",
"mode", "tls", "tls_pin", "software_dev_fee_percent", "node_state",
];
pub const NETWORK_FIELDS: &[&str] = &[
"name", "chain_id", "difficulty", "hashrate", "hashrate_unit", "daa_score", "block_count", "blue_score", "block_reward_ign", "miner_reward_ign", "block_time_target_s", "synced", "node_version", "epoch_seed", "epoch_index", "updated_ms",
"finality", "finality_locked_index", "finality_locked_age_s", "payout_rule",
];
pub const BLOCK_FIELDS: &[&str] = &["hash", "daa_score", "blue_score", "time", "ts_ms", "finder", "worker", "status", "reward_ign", "fee_ign", "effort", "payees", "confirmed_ms", "nonce", "shift"];
pub const MINER_FIELDS: &[&str] = &["ok", "address", "hashrate", "hashrate_1h", "hashrate_reported", "shares", "workers", "blocks", "balance_ign", "paid_ign", "payments", "last_share_ms", "first_seen_ms", "online", "proving", "history"];
pub const PAYMENT_FIELDS: &[&str] = &["time", "ts_ms", "address", "amount_ign", "tx_hash", "status", "dry_run", "nonce"];
pub const HIVE_FIELDS: &[&str] = &["hashrate", "miners", "workers", "blocks", "lastBlock", "fee", "minPayout", "scheme", "symbol", "algo", "difficulty", "networkHashrate", "height"];
fn block_json(b: &crate::state::BlockRec) -> Value {
json!({
"hash": b.hash, "daa_score": b.daa_score, "blue_score": b.blue_score,
"time": iso(b.found_ms), "ts_ms": b.found_ms, "finder": b.finder, "worker": b.worker, "status": b.status,
"reward_ign": ign(b.reward_wei), "fee_ign": ign(b.fee_wei), "effort": round(b.effort, 4),
"payees": b.payees.iter().map(|p| json!({"address": p.address, "fraction": round(p.fraction, 6)})).collect::<Vec<_>>(),
"confirmed_ms": b.confirmed_ms, "nonce": b.nonce, "shift": b.shift,
})
}
fn payment_json(p: &crate::state::PaymentRec) -> Value {
json!({ "time": iso(p.at_ms), "ts_ms": p.at_ms, "address": p.address, "amount_ign": ign(p.amount_wei), "tx_hash": p.tx_hash, "status": p.status, "dry_run": p.dry_run, "nonce": p.nonce })
}
fn iso(ms: u64) -> String {
// UTC, seconds; no chrono dependency
let s = ms / 1000;
let (days, rem) = (s / 86400, s % 86400);
let (h, m, sec) = (rem / 3600, (rem % 3600) / 60, rem % 60);
// civil from days (Howard Hinnant's algorithm)
let z = days as i64 + 719468;
let era = z.div_euclid(146097);
let doe = z.rem_euclid(146097);
let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365;
let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = doy - (153 * mp + 2) / 5 + 1;
let mo = if mp < 10 { mp + 3 } else { mp - 9 };
let y = if mo <= 2 { y + 1 } else { y };
format!("{y:04}-{mo:02}-{d:02}T{h:02}:{m:02}:{sec:02}Z")
}
fn round(v: f64, places: i32) -> f64 {
let f = 10f64.powi(places);
(v * f).round() / f
}
pub fn stats(pool: &Pool) -> Value {
let net = pool.net.lock().unwrap().clone();
let (miners, workers) = pool.online();
let s = pool.state.lock().unwrap();
let (f24, c24, o24) = s.blocks_in(86400);
let last = s.blocks.iter().rev().find(|b| b.status != "orphan").or(s.blocks.last());
let (n, mean, p50, p99, max) = s.check_cost();
let pending: u128 = s.balances.values().sum();
let paid_total: u128 = s.paid.values().sum();
// luck over 24 h in MiningPoolStats' convention (Q69, Q72): expected work over actual, so over 100% is lucky.
// Expected: one block of work per block found; actual: the share weight the blocks took (effort) plus the
// current round. Earlier versions reported the inverse (actual over expected).
let shares_24h_weight: f64 = s.blocks.iter().filter(|b| b.found_ms >= unix_ms().saturating_sub(86_400_000)).map(|b| b.effort).sum::<f64>() + s.pplns.since_block;
let luck = if f24 > 0 && shares_24h_weight > 0.0 { f24 as f64 / shares_24h_weight } else { 0.0 };
let template_age_s = { let t = pool.last_template_ok_ms.load(std::sync::atomic::Ordering::Relaxed); if t == 0 { u64::MAX } else { unix_ms().saturating_sub(t) / 1000 } };
let node_state = if !net.synced && net.updated_ms == 0 { "unreachable" } else if template_age_s > 30 { "no template" } else if !net.synced { "syncing" } else { "ok" };
json!({
"ok": true,
"now": iso(unix_ms()),
"pool": {
"name": pool.cfg.name, "url": pool.cfg.public_url, "address": pool.pool_address_hex,
"algorithm": "igneum",
"scheme": "PPLNS", "fee_percent": pool.cfg.fee_percent, "min_payout_ign": pool.cfg.min_payout_ign, "pplns_window_blocks": pool.cfg.pplns_window_blocks,
"hashrate": round(s.hashrate(600, None, None), 0), "hashrate_unit": "H/s", "hashrate_1h": round(s.hashrate(3600, None, None), 0),
"miners": miners, "workers": workers,
"shares": { "accepted": s.accepted, "stale": s.stale, "rejected": s.rejected, "window_weight_blocks": round(s.pplns.total_weight, 4), "window_shares": s.pplns.shares.len(), "rejected_by_code": s.rejected_by_code, "rejected_by_code_this_run": s.rejected_by_code_run, "counters_since": "ledger (rejected, stale, accepted, rejected_by_code are the ledger's life across restarts; rejected_by_code_this_run is this process)" },
"blocks_24h": f24, "blocks_confirmed_24h": c24, "blocks_orphaned_24h": o24,
"blocks_total": s.blocks_found, "blocks_confirmed_total": s.blocks_confirmed, "blocks_orphaned_total": s.blocks_orphaned,
"last_block": last.map(block_json),
"effort_current": round(s.pplns.since_block, 4), "luck_24h": round(luck, 4),
"paid_24h_ign": ign(s.paid_in(86400)), "paid_total_ign": ign(paid_total), "pending_balance_ign": ign(pending), "pool_fee_total_ign": ign(s.pool_fee_wei),
"uptime_s": pool.started.elapsed().as_secs(),
"share_check_ms": { "count": n, "mean": round(mean, 3), "p50": round(p50, 3), "p99": round(p99, 3), "max": round(max, 3) },
"dry_run": pool.cfg.dry_run,
"mode": if pool.open.is_some() { "open" } else { "operator" },
"tls": pool.tls.is_some(), "tls_pin": pool.tls.as_ref().map(|t| t.pin.clone()),
"software_dev_fee_percent": pool.open.as_ref().map(|o| o.dev_fee_percent as f64).unwrap_or(0.0),
"node_state": node_state,
},
"network": {
"name": net.network, "chain_id": net.chain_id, "difficulty": net.difficulty, "hashrate": net.hashrate, "hashrate_unit": "H/s",
"daa_score": net.daa_score, "block_count": net.block_count, "blue_score": net.blue_score,
"block_reward_ign": net.block_reward_ign, "miner_reward_ign": net.miner_reward_ign, "block_time_target_s": 1,
"synced": net.synced, "node_version": net.node_version, "epoch_seed": net.epoch_seed, "epoch_index": net.epoch_index, "updated_ms": net.updated_ms,
"finality": net.finality, "finality_locked_index": net.finality_locked_index, "finality_locked_age_s": net.finality_locked_age_s,
"payout_rule": "payouts follow blue confirmation, not finality",
},
"source": "igneum-pool v0: shares verified on the CPU warp verifier; network numbers from the pool's node (getBlockDagInfo, estimateNetworkHashesPerSecond); reward by spec 2.5 at the node's DAA score",
})
}
pub fn blocks(pool: &Pool, limit: usize) -> Value {
let s = pool.state.lock().unwrap();
json!({ "ok": true, "blocks": s.blocks.iter().rev().take(limit).map(block_json).collect::<Vec<_>>() })
}
pub fn payments(pool: &Pool, limit: usize) -> Value {
let s = pool.state.lock().unwrap();
json!({ "ok": true, "payments": s.payments.iter().rev().take(limit).map(payment_json).collect::<Vec<_>>() })
}
pub fn miner(pool: &Pool, address: &str) -> Value {
let address = address.to_lowercase();
if address.len() != 42 || !address.starts_with("0x") || !address[2..].chars().all(|c| c.is_ascii_hexdigit()) {
return json!({ "ok": false, "error": "address is 0x followed by 40 hex" });
}
let live: Vec<(String, f64, u32, bool, u32)> = pool
.members_snapshot()
.iter()
.filter(|m| m.address == address)
.map(|m| {
let g = m.inner.lock().unwrap();
(m.worker.clone(), g.reported_hashrate, g.vardiff.shift, g.proving, g.reported_workers)
})
.collect();
let s = pool.state.lock().unwrap();
let rec = s.miners.get(&address).cloned().unwrap_or_default();
let workers: Vec<Value> = rec
.workers
.iter()
.map(|(name, w)| {
let l = live.iter().find(|x| &x.0 == name);
json!({
"name": name, "hashrate": round(s.hashrate(600, Some(&address), Some(name)), 0), "accepted": w.accepted, "stale": w.stale, "rejected": w.rejected, "blocks": w.blocks,
"last_share_ms": w.last_share_ms, "online": l.is_some(), "shift": l.map(|x| x.2), "proving": l.map(|x| x.3).unwrap_or(false), "hashrate_reported": l.map(|x| x.1),
})
})
.collect();
json!({
"ok": true, "address": address,
"hashrate": round(s.hashrate(600, Some(&address), None), 0), "hashrate_1h": round(s.hashrate(3600, Some(&address), None), 0),
"hashrate_reported": round(live.iter().map(|x| x.1).sum::<f64>(), 0),
"shares": { "accepted": rec.accepted, "stale": rec.stale, "rejected": rec.rejected },
"workers": workers, "blocks": rec.blocks,
"balance_ign": ign(s.balances.get(&address).copied().unwrap_or(0)), "paid_ign": ign(s.paid.get(&address).copied().unwrap_or(0)),
"payments": s.payments.iter().rev().filter(|p| p.address == address).take(50).map(payment_json).collect::<Vec<_>>(),
"last_share_ms": rec.last_share_ms, "first_seen_ms": rec.first_seen_ms, "online": !live.is_empty(),
"proving": live.iter().any(|x| x.3),
"history": s.history_of(&address).iter().map(|(t, hs, n)| json!({"hour_ms": t, "hashrate": round(*hs, 0), "shares": n})).collect::<Vec<_>>(),
})
}
/// Q72: a Prometheus text exposition of the pool's numbers.
pub fn metrics(pool: &Pool) -> String {
let v = stats(pool);
let p = &v["pool"];
let n = &v["network"];
let g = |k: &str, val: &Value, help: &str| format!("# HELP igneum_pool_{k} {help}\n# TYPE igneum_pool_{k} gauge\nigneum_pool_{k} {}\n", val.as_f64().unwrap_or(0.0));
let mut out = String::new();
out += &g("hashrate_hs", &p["hashrate"], "pool hash rate over 10 minutes, H/s");
out += &g("miners", &p["miners"], "addresses online");
out += &g("workers", &p["workers"], "workers online");
out += &g("shares_accepted_total", &p["shares"]["accepted"], "accepted shares");
out += &g("shares_stale_total", &p["shares"]["stale"], "stale shares");
out += &g("shares_rejected_total", &p["shares"]["rejected"], "rejected shares");
out += &g("blocks_total", &p["blocks_total"], "blocks found");
out += &g("blocks_confirmed_total", &p["blocks_confirmed_total"], "blocks confirmed blue");
out += &g("blocks_orphaned_total", &p["blocks_orphaned_total"], "blocks orphaned");
out += &g("paid_total_ign", &p["paid_total_ign"], "IGN paid");
out += &g("pending_balance_ign", &p["pending_balance_ign"], "IGN owed");
out += &g("share_check_ms_p99", &p["share_check_ms"]["p99"], "share check p99, ms");
out += &g("network_hashrate_hs", &n["hashrate"], "network hash rate, H/s");
out += &g("network_daa_score", &n["daa_score"], "network DAA score");
out += &g("network_finality_locked_index", &n["finality_locked_index"], "latest locked checkpoint");
out += &format!("# HELP igneum_pool_node_ok 1 when the node is synced and answered a template in the last 30 s\n# TYPE igneum_pool_node_ok gauge\nigneum_pool_node_ok {}\n", if p["node_state"] == "ok" { 1 } else { 0 });
out += &format!("# HELP igneum_pool_finality_active 1 when the network's finality is active\n# TYPE igneum_pool_finality_active gauge\nigneum_pool_finality_active {}\n", if n["finality"] == "active" { 1 } else { 0 });
if let Some(o) = &pool.open {
let s = o.status();
out += &g("open_height", &s["height"], "share chain height");
out += &g("open_stale_total", &s["stale"], "stale shares");
out += &g("open_reorgs_total", &s["reorgs"], "share chain reorgs");
out += &g("open_rejected_total", &s["rejected"], "shares refused");
}
out
}
/// Hive-style pool stats: the flat camelCase object pool dashboards poll (approximate: Hive fixes no schema for a
/// custom pool, so this mirrors the common shape of MiningPoolStats pool JSON).
pub fn hive(pool: &Pool) -> Value {
let v = stats(pool);
let p = &v["pool"];
let n = &v["network"];
json!({
"hashrate": p["hashrate"], "miners": p["miners"], "workers": p["workers"], "blocks": p["blocks_total"],
"lastBlock": p["last_block"].get("ts_ms").cloned().unwrap_or(Value::Null), "fee": p["fee_percent"], "minPayout": p["min_payout_ign"],
"scheme": "PPLNS", "symbol": "IGN", "algo": "igneum", "difficulty": n["difficulty"], "networkHashrate": n["hashrate"], "height": n["daa_score"],
})
}
fn page(pool: &Pool) -> String {
let (fee, min_payout, mode) = match &pool.open {
Some(o) => (format!("{} (software dev fee, a split entry)", o.dev_fee_percent), "none: the coinbase pays".to_string(), "open"),
None => (format!("{}", pool.cfg.fee_percent), format!("{}", pool.cfg.min_payout_ign), "operator"),
};
include_str!("../web/index.html")
.replace("__POOL_NAME__", &pool.cfg.name)
.replace("__POOL_URL__", &pool.cfg.public_url)
.replace("__POOL_ADDRESS__", &pool.pool_address_hex)
.replace("__NETWORK__", &pool.cfg.network)
.replace("__FEE__", &fee)
.replace("__MIN_PAYOUT__", &min_payout)
.replace("__MODE__", mode)
.replace("__TLS__", &match &pool.tls {
Some(t) if t.self_signed => format!("--pool-pin {}", t.pin),
Some(_) => "--pool-tls".to_string(),
None => String::new(),
})
}
fn respond(status: &str, ctype: &str, body: &[u8]) -> Vec<u8> {
let mut v = format!(
"HTTP/1.1 {status}\r\nContent-Type: {ctype}\r\nContent-Length: {}\r\nAccess-Control-Allow-Origin: *\r\nCache-Control: public, max-age=5\r\nX-Content-Type-Options: nosniff\r\nConnection: close\r\n\r\n",
body.len()
)
.into_bytes();
v.extend_from_slice(body);
v
}
pub fn route(pool: &Pool, path: &str) -> Vec<u8> {
let (path, query) = path.split_once('?').unwrap_or((path, ""));
let limit = query.split('&').find_map(|kv| kv.strip_prefix("limit=")).and_then(|v| v.parse::<usize>().ok()).unwrap_or(100).min(1000);
let json = |v: Value| respond("200 OK", "application/json; charset=utf-8", v.to_string().as_bytes());
let from = query.split('&').find_map(|kv| kv.strip_prefix("from=")).and_then(|v| v.parse::<u64>().ok()).unwrap_or(0);
match path {
"/" | "/index.html" => respond("200 OK", "text/html; charset=utf-8", page(pool).as_bytes()),
"/api/stats" => json(stats(pool)),
"/api/blocks" => json(blocks(pool, limit)),
"/api/payments" => json(payments(pool, limit)),
"/api/pool-stats" | "/api/hive" => json(hive(pool)),
"/metrics" => respond("200 OK", "text/plain; version=0.0.4; charset=utf-8", metrics(pool).as_bytes()),
// Q72: health reflects the node (synced, a template in the last 30 s), not the process alone
"/health" => {
let v = stats(pool);
let ok = v["pool"]["node_state"] == "ok";
let body = json!({"ok": ok, "node_state": v["pool"]["node_state"], "synced": v["network"]["synced"], "finality": v["network"]["finality"]});
respond(if ok { "200 OK" } else { "503 Service Unavailable" }, "application/json; charset=utf-8", body.to_string().as_bytes())
}
"/api/open" => match &pool.open {
Some(o) => json(json!({"ok": true, "open": o.status()})),
None => json(json!({"ok": false, "error": "not an open pool"})),
},
"/api/open/shares" => match &pool.open {
Some(o) => json(o.shares_json(from, limit)),
None => json(json!({"ok": false, "error": "not an open pool"})),
},
p if p.starts_with("/api/open/share/") => match pool.open.as_ref().and_then(|o| o.share_json(&p["/api/open/share/".len()..])) {
Some(v) => json(json!({"ok": true, "share": v})),
None => respond("404 Not Found", "application/json; charset=utf-8", br#"{"ok":false,"error":"unknown share"}"#),
},
p if p.starts_with("/api/miners/") => json(miner(pool, &p["/api/miners/".len()..])),
_ => respond("404 Not Found", "application/json; charset=utf-8", br#"{"ok":false,"error":"not found"}"#),
}
}
pub async fn serve(pool: Arc<Pool>, listener: tokio::net::TcpListener) {
println!("{} pool: stats API and page on http://{}/", crate::state::unix_ms(), pool.cfg.http);
loop {
let Ok((sock, _)) = listener.accept().await else { continue };
let pool = pool.clone();
tokio::spawn(async move {
let (rd, mut wr) = sock.into_split();
let mut lines = BufReader::new(rd).lines();
let first = match tokio::time::timeout(Duration::from_secs(10), lines.next_line()).await {
Ok(Ok(Some(l))) => l,
_ => return,
};
// drain the headers
while let Ok(Ok(Some(l))) = tokio::time::timeout(Duration::from_secs(2), lines.next_line()).await {
if l.is_empty() {
break;
}
}
let mut parts = first.split_whitespace();
let method = parts.next().unwrap_or("");
let path = parts.next().unwrap_or("/");
let out = if method != "GET" && method != "HEAD" {
respond("405 Method Not Allowed", "application/json; charset=utf-8", br#"{"ok":false,"error":"GET only"}"#)
} else {
route(&pool, path)
};
let _ = wr.write_all(&out).await;
let _ = wr.shutdown().await;
});
}
}
#[cfg(test)]
mod tests {
use super::*;
fn keys(v: &Value) -> Vec<String> {
let mut k: Vec<String> = v.as_object().expect("object").keys().cloned().collect();
k.sort();
k
}
fn sorted(f: &[&str]) -> Vec<String> {
let mut k: Vec<String> = f.iter().map(|s| s.to_string()).collect();
k.sort();
k
}
/// The fixture (`tests/fixtures/stats.json`, a response captured from the private network run) has exactly the
/// fields the contract names, and the contract matches what the handler builds.
#[test]
fn fixture_matches_the_field_contract() {
let f: Value = serde_json::from_str(include_str!("../tests/fixtures/stats.json")).unwrap();
assert_eq!(keys(&f), sorted(STATS_FIELDS));
assert_eq!(keys(&f["pool"]), sorted(POOL_FIELDS));
assert_eq!(keys(&f["network"]), sorted(NETWORK_FIELDS));
assert_eq!(keys(&f["pool"]["last_block"]), sorted(BLOCK_FIELDS));
let m: Value = serde_json::from_str(include_str!("../tests/fixtures/miner.json")).unwrap();
assert_eq!(keys(&m), sorted(MINER_FIELDS));
let p: Value = serde_json::from_str(include_str!("../tests/fixtures/payments.json")).unwrap();
assert_eq!(keys(&p["payments"][0]), sorted(PAYMENT_FIELDS));
let h: Value = serde_json::from_str(include_str!("../tests/fixtures/pool-stats.json")).unwrap();
assert_eq!(keys(&h), sorted(HIVE_FIELDS));
// the numbers WhatToMine's form asks for are present and typed
assert!(f["pool"]["hashrate"].is_number() && f["pool"]["fee_percent"].is_number() && f["pool"]["min_payout_ign"].is_number());
assert!(f["pool"]["miners"].is_number() && f["pool"]["workers"].is_number() && f["pool"]["blocks_24h"].is_number());
assert!(f["network"]["difficulty"].is_number() && f["network"]["block_reward_ign"].is_number());
}
#[test]
fn iso_dates_are_utc() {
assert_eq!(iso(0), "1970-01-01T00:00:00Z");
assert_eq!(iso(1_790_985_600_000), "2026-10-03T00:00:00Z");
}
}