304 lines
10 KiB
Rust
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])
|
|
}
|
|
}
|