igneum/pool/src/main.rs

183 lines
7.3 KiB
Rust

//! igneum-pool: a public mining pool for Igneum, version 0 (docs/plans/pool.md, docs/spec/09-pool-protocol.md).
//!
//! What it does. Connects to one or more nodes, builds a template per member (the member's vote key in the header,
//! its key reveal and the pool's payout address in the coinbase), hands jobs to members over the spec 09 protocol
//! (newline JSON), sets a share target per member so each sends about one share per 10 s, verifies every share on the
//! CPU warp verifier, submits found blocks to every node, pays PPLNS from the pool's coinbase address on the EVM side
//! with a configurable fee (default 1%), and serves the stats API and a page.
//!
//! ```text
//! igneum-pool --node grpc://127.0.0.1:26610 --evm-rpc http://127.0.0.1:26790 --listen 0.0.0.0:4463 --http 127.0.0.1:4480 \
//! --data-dir /var/lib/igneum-pool --fee-percent 1 --min-payout 1 [--dry-run]
//! ```
mod api;
mod config;
mod node;
mod payout;
mod pool;
mod pplns;
mod protocol;
mod server;
mod state;
mod vardiff;
mod verify;
use config::Config;
use kaspa_addresses::{Address, Prefix, Version};
use kaspa_grpc_client::GrpcClient;
use kaspa_pow::igneum::IgneumEngine;
use pool::{NetInfo, Pool};
use state::{State, WEI_PER_IGN};
use std::collections::HashMap;
use std::sync::atomic::AtomicU64;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
#[tokio::main]
async fn main() {
let args: Vec<String> = std::env::args().skip(1).collect();
let cfg = Config::from_args(&args);
std::fs::create_dir_all(&cfg.data_dir).unwrap_or_else(|e| {
eprintln!("data dir {}: {e}", cfg.data_dir.display());
std::process::exit(1)
});
let signer = payout::load_or_create_key(&cfg.payout_key).unwrap_or_else(|e| {
eprintln!("payout key: {e}");
std::process::exit(1)
});
let pool_address: [u8; 20] = signer.address().into_array();
let pool_address_hex = format!("0x{}", hex::encode(pool_address));
let evm = payout::EvmRpc::new(&cfg.evm_rpc).unwrap_or_else(|e| {
eprintln!("--evm-rpc: {e}");
std::process::exit(1)
});
// nodes
let mut clients = Vec::new();
for url in &cfg.nodes {
match tokio::time::timeout(Duration::from_secs(10), GrpcClient::connect(url.clone())).await {
Ok(Ok(c)) => clients.push(Arc::new(c)),
Ok(Err(e)) => {
eprintln!("node {url}: {e}");
std::process::exit(1)
}
Err(_) => {
eprintln!("node {url}: connect timed out");
std::process::exit(1)
}
}
}
let node = clients.remove(0);
// the UTXO-side coinbase address: derived from the pool's EVM address; the execution layer pays the EVM side
// (igneum/exec/src/executor.rs), the UTXO output is not spent by v0 (docs/plans/pool.md)
let prefix = match cfg.network.as_str() {
"mainnet" => Prefix::Mainnet,
"testnet" => Prefix::Testnet,
_ => Prefix::Devnet,
};
let pay_address = {
let mut h = kaspa_hashes::BlockHash::new();
use kaspa_hashes::HasherBase;
h.update(b"igneum-pool-utxo-payout-v0").update(pool_address);
Address::new(prefix, Version::PubKey, &h.finalize().as_bytes())
};
let state_path = cfg.data_dir.join("state.json");
let pool = Arc::new(Pool {
engine: Arc::new(IgneumEngine::new()),
state: Mutex::new(State::load(&state_path, cfg.pplns_window_blocks)),
members: Mutex::new(HashMap::new()),
next_member_id: AtomicU64::new(0),
next_template_id: AtomicU64::new(0),
next_job_id: AtomicU64::new(0),
node,
extra_nodes: clients,
pool_address,
pool_address_hex: pool_address_hex.clone(),
pay_address,
net: Mutex::new(NetInfo::default()),
want_templates: tokio::sync::Notify::new(),
verify_permits: tokio::sync::Semaphore::new(cfg.verify_threads),
started: Instant::now(),
cfg: cfg.clone(),
});
let chain = evm.call("eth_chainId", serde_json::json!([])).await.ok().and_then(|v| v.as_str().map(|s| s.to_string()));
println!(
"igneum-pool v0: {} on {} (chain id {} per --network, node says {}), payout address {}, fee {}%, min payout {} IGN, PPLNS window {} blocks{}",
cfg.name,
cfg.network,
cfg.chain_id(),
chain.unwrap_or_else(|| "unreachable".into()),
pool_address_hex,
cfg.fee_percent,
cfg.min_payout_ign,
cfg.pplns_window_blocks,
if cfg.dry_run { ", DRY RUN (payouts computed, not sent)" } else { "" }
);
let payer = Arc::new(payout::Payer { signer, rpc: evm, chain_id: cfg.chain_id(), dry_run: cfg.dry_run, min_payout_wei: (cfg.min_payout_ign * WEI_PER_IGN as f64) as u128 });
tokio::spawn(node::template_feed(pool.clone()));
tokio::spawn(node::confirm_loop(pool.clone()));
tokio::spawn(node::net_loop(pool.clone()));
tokio::spawn(node::vardiff_loop(pool.clone()));
tokio::spawn(server::listen(pool.clone()));
tokio::spawn(api::serve(pool.clone()));
{
let (pool, payer) = (pool.clone(), payer.clone());
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(pool.cfg.payout_interval_s.max(5))).await;
for l in payer.round(&pool.state).await {
println!("{} {l}", state::unix_ms());
}
for l in payer.receipts(&pool.state).await {
println!("{} {l}", state::unix_ms());
}
}
});
}
{
let (pool, path) = (pool.clone(), state_path.clone());
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(pool.cfg.snapshot_interval_s.max(1))).await;
let mut s = pool.state.lock().unwrap();
if s.dirty
&& let Err(e) = s.save(&path)
{
eprintln!("state: save failed: {e}");
}
}
});
}
// a status line every 30 s
{
let pool = pool.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(30)).await;
let (miners, workers) = pool.online();
let s = pool.state.lock().unwrap();
let (n, mean, p50, p99, max) = s.check_cost();
println!(
"{} STATUS miners={miners} workers={workers} accepted={} stale={} rejected={} hashrate={:.0} H/s blocks={} confirmed={} orphans={} window={:.3} blocks/{} shares check_ms n={n} mean={mean:.3} p50={p50:.3} p99={p99:.3} max={max:.3}",
state::unix_ms(),
s.accepted,
s.stale,
s.rejected,
s.hashrate(600, None, None),
s.blocks_found,
s.blocks_confirmed,
s.blocks_orphaned,
s.pplns.total_weight,
s.pplns.shares.len()
);
}
});
}
let _ = tokio::signal::ctrl_c().await;
println!("{} pool: stopping, saving state", state::unix_ms());
let mut s = pool.state.lock().unwrap();
if let Err(e) = s.save(&state_path) {
eprintln!("state: save failed: {e}");
}
}