igneum/pool/src/state.rs

304 lines
10 KiB
Rust

//! The pool's ledger: shares, blocks, balances, payments and the per-miner counters. One `Mutex<State>` with short
//! critical sections; a JSON snapshot every `snapshot_interval_s` and at shutdown (v0 persistence, docs/plans/pool.md).
use crate::pplns::{Payee, Pplns, ShareRec};
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, VecDeque};
use std::path::Path;
pub const WEI_PER_IGN: u128 = 1_000_000_000_000_000_000;
/// 1 IGN = 10^8 sompi on the UTXO side = 10^18 wei on the EVM side (igneum/exec/src/config.rs WEI_PER_SOMPI)
pub const WEI_PER_SOMPI: u128 = 10_000_000_000;
pub fn unix_ms() -> u64 {
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map(|d| d.as_millis() as u64).unwrap_or(0)
}
pub fn ign(wei: u128) -> f64 {
wei as f64 / WEI_PER_IGN as f64
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct BlockRec {
pub hash: String,
pub daa_score: u64,
pub blue_score: u64,
pub found_ms: u64,
pub finder: String,
pub worker: String,
/// pending, confirmed, orphan
pub status: String,
pub reward_wei: u128,
pub fee_wei: u128,
pub effort: f64,
pub payees: Vec<Payee>,
pub confirmed_ms: Option<u64>,
pub nonce: String,
pub shift: u32,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct PaymentRec {
pub at_ms: u64,
pub address: String,
pub amount_wei: u128,
pub tx_hash: Option<String>,
/// dry-run, sent, confirmed, failed
pub status: String,
pub dry_run: bool,
pub nonce: Option<u64>,
pub note: String,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct WorkerRec {
pub accepted: u64,
pub stale: u64,
pub rejected: u64,
pub last_share_ms: u64,
pub last_seen_ms: u64,
pub blocks: u64,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct MinerRec {
pub accepted: u64,
pub stale: u64,
pub rejected: u64,
pub blocks: u64,
pub last_share_ms: u64,
pub first_seen_ms: u64,
pub workers: HashMap<String, WorkerRec>,
}
/// One accepted share for the hashrate windows: expected hashes at its target.
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Sample {
pub at_ms: u64,
pub address: String,
pub worker: String,
pub hashes: f64,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct State {
pub started_ms: u64,
pub pplns: Pplns,
pub blocks: Vec<BlockRec>,
pub balances: HashMap<String, u128>,
pub paid: HashMap<String, u128>,
pub payments: Vec<PaymentRec>,
pub miners: HashMap<String, MinerRec>,
pub accepted: u64,
pub stale: u64,
pub rejected: u64,
pub blocks_found: u64,
pub blocks_confirmed: u64,
pub blocks_orphaned: u64,
pub pool_fee_wei: u128,
pub submit_errors: u64,
#[serde(skip)]
pub samples: VecDeque<Sample>,
/// Verification cost of the last 20,000 shares, milliseconds
#[serde(skip)]
pub check_ms: VecDeque<f64>,
#[serde(skip)]
pub dirty: bool,
}
impl State {
pub fn new(window_blocks: f64) -> Self {
State {
started_ms: unix_ms(),
pplns: Pplns::new(window_blocks),
blocks: Vec::new(),
balances: HashMap::new(),
paid: HashMap::new(),
payments: Vec::new(),
miners: HashMap::new(),
accepted: 0,
stale: 0,
rejected: 0,
blocks_found: 0,
blocks_confirmed: 0,
blocks_orphaned: 0,
pool_fee_wei: 0,
submit_errors: 0,
samples: VecDeque::new(),
check_ms: VecDeque::new(),
dirty: false,
}
}
pub fn load(path: &Path, window_blocks: f64) -> Self {
match std::fs::read(path) {
Ok(bytes) => match serde_json::from_slice::<State>(&bytes) {
Ok(mut s) => {
s.pplns.window_blocks = window_blocks;
s.dirty = false;
eprintln!("state: loaded {} ({} blocks, {} payments, {} miners)", path.display(), s.blocks.len(), s.payments.len(), s.miners.len());
s
}
Err(e) => {
eprintln!("state: {} unreadable ({e}); starting empty", path.display());
State::new(window_blocks)
}
},
Err(_) => State::new(window_blocks),
}
}
pub fn save(&mut self, path: &Path) -> std::io::Result<()> {
if let Some(dir) = path.parent() {
std::fs::create_dir_all(dir)?;
}
let tmp = path.with_extension("json.tmp");
std::fs::write(&tmp, serde_json::to_vec(self)?)?;
std::fs::rename(tmp, path)?;
self.dirty = false;
Ok(())
}
fn miner(&mut self, address: &str, worker: &str, now: u64) -> &mut WorkerRec {
let m = self.miners.entry(address.to_string()).or_insert_with(|| MinerRec { first_seen_ms: now, ..Default::default() });
let w = m.workers.entry(worker.to_string()).or_default();
w.last_seen_ms = now;
w
}
pub fn accept_share(&mut self, address: &str, worker: &str, weight: f64, share_target64: u64, cost_ms: f64) {
let now = unix_ms();
self.accepted += 1;
self.pplns.push(ShareRec { address: address.into(), worker: worker.into(), weight, at_ms: now });
let hashes = if share_target64 == 0 { 0.0 } else { (u64::MAX as f64 + 1.0) / share_target64 as f64 };
self.samples.push_back(Sample { at_ms: now, address: address.into(), worker: worker.into(), hashes });
let cutoff = now.saturating_sub(15 * 60 * 1000);
while self.samples.front().is_some_and(|s| s.at_ms < cutoff) {
self.samples.pop_front();
}
self.check_ms.push_back(cost_ms);
while self.check_ms.len() > 20_000 {
self.check_ms.pop_front();
}
let w = self.miner(address, worker, now);
w.accepted += 1;
w.last_share_ms = now;
let m = self.miners.get_mut(address).unwrap();
m.accepted += 1;
m.last_share_ms = now;
self.dirty = true;
}
pub fn refuse_share(&mut self, address: &str, worker: &str, stale: bool, cost_ms: Option<f64>) {
let now = unix_ms();
if let Some(c) = cost_ms {
self.check_ms.push_back(c);
}
if stale {
self.stale += 1;
} else {
self.rejected += 1;
}
let w = self.miner(address, worker, now);
if stale {
w.stale += 1
} else {
w.rejected += 1
}
let m = self.miners.get_mut(address).unwrap();
if stale {
m.stale += 1
} else {
m.rejected += 1
}
self.dirty = true;
}
pub fn block_found(&mut self, mut rec: BlockRec) {
let now = unix_ms();
rec.payees = self.pplns.snapshot();
rec.effort = self.pplns.take_effort();
self.blocks_found += 1;
let w = self.miner(&rec.finder.clone(), &rec.worker.clone(), now);
w.blocks += 1;
self.miners.get_mut(&rec.finder).unwrap().blocks += 1;
self.blocks.push(rec);
while self.blocks.len() > 5000 {
self.blocks.remove(0);
}
self.dirty = true;
}
/// Credits a confirmed block to its payees (the PPLNS snapshot of the moment it was found).
pub fn confirm_block(&mut self, hash: &str, fee_percent: f64) -> Option<(u128, Vec<(String, u128)>)> {
let idx = self.blocks.iter().position(|b| b.hash == hash && b.status == "pending")?;
let reward = self.blocks[idx].reward_wei;
let (parts, kept) = crate::pplns::distribute(reward, fee_percent, &self.blocks[idx].payees);
for (a, wei) in &parts {
*self.balances.entry(a.clone()).or_insert(0) += wei;
}
self.pool_fee_wei += kept;
self.blocks_confirmed += 1;
let b = &mut self.blocks[idx];
b.status = "confirmed".into();
b.fee_wei = kept;
b.confirmed_ms = Some(unix_ms());
self.dirty = true;
Some((reward, parts))
}
pub fn orphan_block(&mut self, hash: &str) -> bool {
let Some(b) = self.blocks.iter_mut().find(|b| b.hash == hash && b.status == "pending") else { return false };
b.status = "orphan".into();
self.blocks_orphaned += 1;
self.dirty = true;
true
}
/// Hashes per second over the last `secs` seconds, pool-wide or for one address or worker.
pub fn hashrate(&self, secs: u64, address: Option<&str>, worker: Option<&str>) -> f64 {
let now = unix_ms();
let cutoff = now.saturating_sub(secs * 1000);
let span = (now.saturating_sub(self.started_ms.max(cutoff))).max(1000) as f64 / 1000.0;
let sum: f64 = self
.samples
.iter()
.filter(|s| s.at_ms >= cutoff && address.is_none_or(|a| s.address == a) && worker.is_none_or(|w| s.worker == w))
.map(|s| s.hashes)
.sum();
sum / span
}
pub fn blocks_in(&self, secs: u64) -> (u64, u64, u64) {
let cutoff = unix_ms().saturating_sub(secs * 1000);
let mut found = 0;
let mut confirmed = 0;
let mut orphan = 0;
for b in self.blocks.iter().filter(|b| b.found_ms >= cutoff) {
found += 1;
match b.status.as_str() {
"confirmed" => confirmed += 1,
"orphan" => orphan += 1,
_ => {}
}
}
(found, confirmed, orphan)
}
pub fn paid_in(&self, secs: u64) -> u128 {
let cutoff = unix_ms().saturating_sub(secs * 1000);
self.payments.iter().filter(|p| p.at_ms >= cutoff && !p.dry_run && p.status != "failed").map(|p| p.amount_wei).sum()
}
/// Share check cost: (count, mean ms, p50, p99, max)
pub fn check_cost(&self) -> (usize, f64, f64, f64, f64) {
if self.check_ms.is_empty() {
return (0, 0.0, 0.0, 0.0, 0.0);
}
let mut v: Vec<f64> = self.check_ms.iter().copied().collect();
v.sort_by(|a, b| a.partial_cmp(b).unwrap());
let n = v.len();
let mean = v.iter().sum::<f64>() / n as f64;
(n, mean, v[n / 2], v[(n * 99 / 100).min(n - 1)], v[n - 1])
}
}