igneum/pool/src/pool.rs
igneum-labs 9bfde8bf8a Pool daemon: the confirmation walk on its own gRPC connection, never restarted from the pruning point, bounded per tick
7 October 2026, 06:01Z and 06:12:55Z on pool-1: confirm_loop restarted its chain walk from the pruning point on any
failed getVirtualChainFromBlock and then fetched every chain block since (about 120,000) over the one connection the
templates used; one timed-out request under four parallel template fetches started it, every template request after it
timed out, no job was issued, 105 blocks on stale templates were orphans. Now: a second GrpcClient (pool.walker) for the
walk and the network numbers; a failed chain call keeps its cursor (the sink only when the node no longer knows it;
cursor_after_failure, node::walk_tests); at most 600 chain blocks per tick with a line saying so; a start with pending
blocks walks from the sink and says that older ones resolve by the orphan rule. docs/plans/pool.md section 9.2: the
re-run's record (the class question closed on pm-1 and pm-2, 71,240 shares, 0 mismatches), the cause, what stays open.
Suite 24 of 24 on igneum-build-1. Protocol and API unchanged.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-07 06:18:53 +00:00

140 lines
4.8 KiB
Rust

//! The live pool: configuration, the share verifier engine, the ledger, the connected members and their jobs.
use crate::config::Config;
use crate::state::State;
use crate::vardiff::Vardiff;
use crate::verify::JobKey;
use kaspa_consensus_core::finality::KeyReveal;
use kaspa_grpc_client::GrpcClient;
use kaspa_hashes::Hash;
use kaspa_pow::igneum::{EpochRef, EpochSeeds, IgneumEngine};
use kaspa_rpc_core::RpcRawBlock;
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Instant;
/// A job issued to one member: the template it was built on and what its shares are checked against.
pub struct JobRec {
pub job_id: u64,
pub template_id: u64,
pub raw: RpcRawBlock,
pub key: JobKey,
pub epoch: Arc<EpochRef>,
pub shift: u32,
pub issued: Instant,
/// Set when a job on a newer tip was issued: shares after the grace are stale
pub superseded_at: Option<Instant>,
pub seen: HashSet<u64>,
pub daa_score: u64,
pub blue_score: u64,
}
pub struct MemberInner {
pub vardiff: Vardiff,
pub jobs: VecDeque<JobRec>,
pub last_parents: Option<Vec<Hash>>,
pub last_seeds: Option<EpochSeeds>,
pub accepted: u64,
pub stale: u64,
pub rejected: u64,
pub blocks: u64,
pub refused: u64,
pub last_share: Option<Instant>,
pub revealed: bool,
/// The member's own report (spec `stats`): hashrate, workers, and the proving flag (consequence C6)
pub reported_hashrate: f64,
pub reported_workers: u32,
pub proving: bool,
pub current_target64: u64,
pub in_flight_checks: u32,
}
pub struct Member {
pub id: u64,
pub address: String,
pub worker: String,
pub key_hash: Hash,
pub reveal: KeyReveal,
pub remote: String,
pub connected_at: Instant,
pub tx: tokio::sync::mpsc::UnboundedSender<String>,
pub inner: Mutex<MemberInner>,
}
impl Member {
pub fn send(&self, line: String) {
let _ = self.tx.send(line);
}
pub fn job(&self, job_id: u64) -> Option<(JobKey, Arc<EpochRef>, u32, bool, u64, u64)> {
let g = self.inner.lock().unwrap();
g.jobs.iter().find(|j| j.job_id == job_id).map(|j| (j.key.clone(), j.epoch.clone(), j.shift, j.superseded_at.is_some(), j.daa_score, j.blue_score))
}
}
/// The network numbers the API reports beside the pool's own, refreshed every 5 s from the primary node.
#[derive(Clone, Debug, Default, serde::Serialize)]
pub struct NetInfo {
pub network: String,
pub chain_id: u64,
pub difficulty: f64,
pub hashrate: Option<f64>,
pub daa_score: u64,
pub block_count: u64,
pub blue_score: u64,
pub synced: bool,
pub node_version: String,
pub block_reward_ign: f64,
pub miner_reward_ign: f64,
pub updated_ms: u64,
pub epoch_seed: String,
pub epoch_index: u64,
}
pub struct Pool {
pub cfg: Config,
pub engine: Arc<IgneumEngine>,
pub state: Mutex<State>,
pub members: Mutex<HashMap<u64, Arc<Member>>>,
pub next_member_id: AtomicU64,
pub next_template_id: AtomicU64,
pub next_job_id: AtomicU64,
pub node: Arc<GrpcClient>,
/// A second connection to the same node for the confirmation walk and the network numbers (7 October 2026,
/// 06:01Z and 06:12:55Z: the walk's `get_block` calls shared the template connection and starved every template
/// request; the gRPC client is one request stream per connection)
pub walker: Arc<GrpcClient>,
pub extra_nodes: Vec<Arc<GrpcClient>>,
/// The pool's EVM coinbase address, lowercase 0x hex, named in every template's `IGNA` field
pub pool_address: [u8; 20],
pub pool_address_hex: String,
pub pay_address: kaspa_addresses::Address,
pub net: Mutex<NetInfo>,
pub want_templates: tokio::sync::Notify,
pub verify_permits: tokio::sync::Semaphore,
pub started: Instant,
/// Template fetches that failed or timed out since start, and the unix ms of the last one that succeeded (0 =
/// never): a daemon whose node answers no template issues no job, and the STATUS line in the log says so
pub template_failures: AtomicU64,
pub last_template_ok_ms: AtomicU64,
}
impl Pool {
pub fn members_snapshot(&self) -> Vec<Arc<Member>> {
self.members.lock().unwrap().values().cloned().collect()
}
pub fn online(&self) -> (usize, usize) {
let m = self.members.lock().unwrap();
let addresses: HashSet<&str> = m.values().map(|x| x.address.as_str()).collect();
(addresses.len(), m.len())
}
pub fn now_s(&self) -> f64 {
self.started.elapsed().as_secs_f64()
}
pub fn template_id(&self) -> u64 {
self.next_template_id.fetch_add(1, Ordering::Relaxed) + 1
}
pub fn job_id(&self) -> u64 {
self.next_job_id.fetch_add(1, Ordering::Relaxed) + 1
}
}