//! 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, 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, pub seen: HashSet, pub daa_score: u64, pub blue_score: u64, } pub struct MemberInner { pub vardiff: Vardiff, pub jobs: VecDeque, pub last_parents: Option>, pub last_seeds: Option, pub accepted: u64, pub stale: u64, pub rejected: u64, pub blocks: u64, pub refused: u64, pub last_share: Option, 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, pub inner: Mutex, } impl Member { pub fn send(&self, line: String) { let _ = self.tx.send(line); } pub fn job(&self, job_id: u64) -> Option<(JobKey, Arc, 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, 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, pub state: Mutex, pub members: Mutex>>, pub next_member_id: AtomicU64, pub next_template_id: AtomicU64, pub next_job_id: AtomicU64, pub node: Arc, /// 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, pub extra_nodes: Vec>, /// 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, 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> { 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 } }