igneum/pool/src/main.rs

285 lines
13 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 open;
mod p2p;
mod payout;
mod pool;
mod pplns;
mod protocol;
mod server;
mod sidechain;
mod state;
mod tls;
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<String> = std::env::args().skip(1).collect();
// `igneum-pool verify-share <json line> [--block <raw block json>]`: the drop proof from a member's own log
if args.first().map(|a| a.as_str()) == Some("verify-share") {
std::process::exit(sidechain::verify_share_cli(&args[1..]));
}
let cfg = Config::from_args(&args);
// spec 9.3: TLS 1.3 on the member port; the public networks refuse the clear unless told
let tls = if cfg.tls_on() {
match tls::load(cfg.tls_cert.as_deref(), cfg.tls_key.as_deref(), cfg.tls_self_signed.then_some(cfg.data_dir.as_path()), &cfg.name) {
Ok(t) => Some(t),
Err(e) => {
eprintln!("POOL NOT STARTED: TLS: {e}");
std::process::exit(2)
}
}
} else {
if cfg.network != "devnet" && !cfg.allow_plain {
eprintln!("POOL NOT STARTED: --network {} wants TLS 1.3 on the member port (spec 9.3): --tls-self-signed or --tls-cert/--tls-key, or --allow-plain to say otherwise", cfg.network);
std::process::exit(2)
}
None
};
// 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)
});
// the open pool holds no key: no payout key is made or read, the members' own addresses are in every coinbase
let signer = if cfg.open {
None
} else {
Some(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.as_ref().map(|s| s.address().into_array()).unwrap_or([0u8; 20]);
let pool_address_hex = if cfg.open { "none (open pool: every coinbase pays the window)".to_string() } else { 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);
let walker = match tokio::time::timeout(Duration::from_secs(10), GrpcClient::connect(cfg.nodes[0].clone())).await {
Ok(Ok(c)) => Arc::new(c),
Ok(Err(e)) => {
eprintln!("node {} (second connection): {e}", cfg.nodes[0]);
std::process::exit(1)
}
Err(_) => {
eprintln!("node {} (second connection): connect timed out", cfg.nodes[0]);
std::process::exit(1)
}
};
// 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 engine = Arc::new(IgneumEngine::new());
let walker_for_open = walker.clone();
let pool = Arc::new(Pool {
engine: engine.clone(),
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,
walker,
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),
tls,
open: cfg.open.then(|| Arc::new(open::Open::new(&cfg, &state_path, engine.clone(), walker_for_open))),
});
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 { "" }
);
if let Some(o) = &pool.open {
println!(
"{} OPEN POOL: share chain '{}' ({} s per share, window {} shares, software dev fee {}% as a split entry), p2p on {}, {} peer(s); no operator key, no balance: every coinbase pays the window from DAA {} (pool_split_activation_daa on the node)",
state::unix_ms(),
cfg.open_chain,
cfg.chain_share_s,
cfg.window_shares,
o.dev_fee_percent,
cfg.p2p_listen,
cfg.peers.len(),
o.activation_note()
);
tokio::spawn(p2p::run(pool.clone(), o.clone()));
}
let payer = signer.map(|signer| 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(state::alert_loop(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));
if let Some(payer) = payer {
let pool = pool.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}");
}
}