//! 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 kaspa_rpc_core::api::rpc::RpcApi; 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 = std::env::args().skip(1).collect(); let cfg = Config::from_args(&args); // The two listeners first (7 October 2026, the fleet agent's row from the 6 October night: a daemon whose member // port another process still held ran 20 minutes as a process with no socket): a port that cannot be bound ends // the start here, exit code 2, before a node is contacted or a member is let down. let member_listener = server::bind_listener("members", &cfg.listen).await.unwrap_or_else(|e| { eprintln!("POOL NOT STARTED: {e}"); std::process::exit(2) }); let api_listener = server::bind_listener("stats API", &cfg.http).await.unwrap_or_else(|e| { eprintln!("POOL NOT STARTED: {e}"); std::process::exit(2) }); 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 node must answer a template with its pow_epoch before the pool serves anyone: a pool whose node answers no // template issues no job (the 6 October night: four members connected, "template timed out" per member, no job // for 20 minutes). Retried for --node-wait-secs, saying so; then exit code 3. { let probe_address = Address::new(Prefix::Devnet, Version::PubKey, &[0u8; 32]); let deadline = Instant::now() + Duration::from_secs(cfg.node_wait_secs); let mut attempt = 0u32; loop { attempt += 1; let why = match tokio::time::timeout(Duration::from_secs(10), node.get_block_template(probe_address.clone(), Vec::new())).await { Ok(Ok(t)) if t.pow_epoch.is_some() => { let i = t.pow_epoch.unwrap(); println!("{} node {} answers templates: epoch {} class {} era {} daa {}", state::unix_ms(), cfg.nodes[0], i.epoch_index, i.class().name(), i.era_seed.map(|h| h.to_string()).unwrap_or_else(|| "-".into()), i.virtual_daa_score); break; } Ok(Ok(_)) => "the template carries no pow_epoch (a node before the 0.3.x line); the pool cannot name a program".to_string(), Ok(Err(e)) => e.to_string(), Err(_) => "template timed out after 10 s".to_string(), }; if Instant::now() >= deadline { eprintln!("POOL NOT STARTED: node {} answered no usable template in {} s ({attempt} attempts; last: {why})", cfg.nodes[0], cfg.node_wait_secs); std::process::exit(3) } eprintln!("{} waiting for node {} to answer a template (attempt {attempt}: {why}); giving up at {} s", state::unix_ms(), cfg.nodes[0], cfg.node_wait_secs); tokio::time::sleep(Duration::from_secs(5)).await; } } // 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(), template_failures: AtomicU64::new(0), last_template_ok_ms: AtomicU64::new(0), }); 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(node::status_loop(pool.clone())); tokio::spawn(server::listen(pool.clone(), member_listener)); tokio::spawn(api::serve(pool.clone(), api_listener)); { 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}"); } }