//! The node side: templates per member (the member's vote key in the header, its key reveal and the pool's `IGNA` //! address in the coinbase extra data, spec 9.4.1 items 1 and 2), jobs on them, found blocks submitted to every //! node, the blue-set walk that confirms or orphans them, and the network numbers for the API. use crate::pool::{JobRec, Member, NetInfo, Pool}; use crate::protocol::{hex_bytes, hex_u64, Msg}; use crate::state::{unix_ms, BlockRec, WEI_PER_SOMPI}; use crate::vardiff::{share_target, Vardiff}; use crate::verify::JobKey; use kaspa_consensus_core::block::Block; use kaspa_consensus_core::header::Header; use kaspa_consensus_core::igneum::{block_subsidy, install_pow_genesis, install_pow_schedule, install_program_class_v3_activation, pow_genesis_dataset_log2, pow_genesis_day_index, pow_schedule, producer_share, program_class_v3_activation_daa, PowSchedule}; use kaspa_grpc_client::GrpcClient; use kaspa_hashes::Hash; use kaspa_notify::{listener::ListenerId, scope::NewBlockTemplateScope}; use kaspa_pow::igneum::{day_index, header_prehash, target64, EpochSeeds, IgneumEngine}; use kaspa_rpc_core::api::rpc::RpcApi; use kaspa_rpc_core::{Notification, RpcRawBlock}; use rand::Rng; use std::collections::HashSet; use std::sync::atomic::Ordering; use std::sync::Arc; use std::time::{Duration, Instant}; fn now() -> String { let t = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap(); format!("{}.{:03}", t.as_secs(), t.subsec_millis()) } /// The extra data of a member's template: its key reveal (so its first block reveals the key) and the pool's /// payout address (`IGNA`, the address the execution layer pays). pub fn extra_data(member: &Member, pool_address: &[u8; 20]) -> Vec { let mut v = member.reveal.to_extra_data(); v.extend_from_slice(&kaspa_consensus_core::evm::miner_address_extra_data(pool_address)); v } /// One template for one member, turned into a job. Blocks the thread on the first use of a seed pair (the 256 MiB /// cache build), so it runs under `spawn_blocking`. pub async fn fetch_job(pool: &Arc, member: &Arc) -> Result<(), String> { let t0 = Instant::now(); let tmpl = tokio::time::timeout(Duration::from_secs(5), pool.node.get_block_template(pool.pay_address.clone(), extra_data(member, &pool.pool_address))) .await .map_err(|_| "template timed out".to_string())? .map_err(|e| e.to_string())?; let mut raw = tmpl.block; raw.header.vote_key_hash = member.key_hash; let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?; let header: Header = block.header.as_ref().clone(); let info = tmpl.pow_epoch.as_ref().ok_or("the node reports no pow_epoch; a devnet-v4 line node is needed")?; // The node's PoW schedule, genesis day and dataset size, and class v3 activation are the pool's (what the solo // miner's `template()` installs): the share verifier's day cache and program must be the node's exactly let wanted = PowSchedule::clamped(info.epoch_blocks, info.epoch_lead, info.day_ms); if wanted != pow_schedule() { install_pow_schedule(wanted); eprintln!("{} node PoW schedule: {} DAA per epoch, lead {}, day {} ms", now(), wanted.epoch_blocks, wanted.epoch_lead, wanted.day_ms); } if (info.genesis_day_index, info.genesis_dataset_log2) != (pow_genesis_day_index(), pow_genesis_dataset_log2()) { install_pow_genesis(info.genesis_day_index, info.genesis_dataset_log2); eprintln!("{} node genesis day index {} and genesis dataset 2^{} words: the day cache follows them", now(), info.genesis_day_index, info.genesis_dataset_log2); } if info.program_class_v3_activation_daa != program_class_v3_activation_daa() { install_program_class_v3_activation(info.program_class_v3_activation_daa); eprintln!("{} node program class v3 activation: {} (epoch {} is class {}, the next epoch class {})", now(), info.program_class_v3_activation_daa, info.epoch_index, info.class().name(), info.next_class().name()); } // Counter ASIC 2.0: the class and the era seed of the epoch ride with the template; the verifier hashes the // program they name, the same one the members' workers compile (6 October 2026: a fixed class here or on the // member refused every GPU share of the fleet run) let seeds = EpochSeeds { epoch: info.epoch_seed, day: day_index(header.timestamp), class: info.class(), era: info.era_seed.unwrap_or(kaspa_hashes::ZERO_HASH) }; let engine = pool.engine.clone(); let epoch = tokio::task::spawn_blocking(move || engine.epoch_for(&seeds)).await.map_err(|e| e.to_string())?; let prehash = header_prehash(&header); let t64 = target64(header.bits); let parents: Vec = header.direct_parents().to_vec(); let template_id = pool.template_id(); let job_id = pool.job_id(); let hi: u32 = rand::thread_rng().r#gen(); let nonce_start = (hi as u64) << 32; let (_shift, clean, seeds_changed, share_t) = { let mut g = member.inner.lock().unwrap(); if g.vardiff.changes == 0 && g.current_target64 == 0 { // first job: size the shift from the target let s = Vardiff::initial_shift(t64, pool.cfg.min_shift, pool.cfg.max_shift); g.vardiff.shift = s; } g.current_target64 = t64; let cap = Vardiff::cap_for(t64).min(pool.cfg.max_shift); if g.vardiff.shift > cap { g.vardiff.shift = cap; } let shift = g.vardiff.shift; let clean = g.last_parents.as_ref() != Some(&parents); if clean { for j in g.jobs.iter_mut() { if j.superseded_at.is_none() { j.superseded_at = Some(Instant::now()); } } } g.last_parents = Some(parents); let seeds_changed = g.last_seeds != Some(seeds); g.last_seeds = Some(seeds); let share_t = share_target(t64, shift); g.jobs.push_back(JobRec { job_id, template_id, raw: raw.clone(), key: JobKey { prehash, target64: t64, share_target64: share_t, seeds }, epoch, shift, issued: Instant::now(), superseded_at: None, seen: HashSet::new(), daa_score: header.daa_score, blue_score: header.blue_score, }); while g.jobs.len() > 12 { g.jobs.pop_front(); } (shift, clean, seeds_changed, share_t) }; if seeds_changed { let (_, day_bytes) = IgneumEngine::seed_bytes(&seeds); member.send( Msg::Seeds { epoch_seed: seeds.epoch.to_string(), day: seeds.day, day_seed: hex_bytes(&day_bytes), seed_source: "node pow_epoch".into(), next_epoch_seed: info.next_epoch_seed.map(|h| h.to_string()), next_at_daa: info.boundary_daa_score, epoch_blocks: info.epoch_blocks, epoch_lead: info.epoch_lead, day_ms: info.day_ms, program_class: info.program_class, next_program_class: info.next_program_class, era_seed: info.era_seed.map(|h| h.to_string()), era_index: info.era_index, genesis_day_index: info.genesis_day_index, genesis_dataset_log2: info.genesis_dataset_log2, program_class_v3_activation_daa: info.program_class_v3_activation_daa, } .line(), ); println!("{} SEEDS for member {} ({}): epoch {} day {} class {} era {}", now(), member.id, member.worker, seeds.epoch, seeds.day, seeds.class.name(), seeds.era); } member.send(Msg::Template { template_id, mode: "A".into(), block: Box::new(raw), daa_score: header.daa_score, bits: header.bits, votes_carried: Vec::new() }.line()); member.send( Msg::Job { job_id, template_id, prehash: hex_bytes(&prehash), target64: hex_u64(t64), share_target64: hex_u64(share_t), nonce_start: hex_u64(nonce_start), nonce_count: hex_u64(1 << 32), epoch_seed: seeds.epoch.to_string(), day: seeds.day, program_class: seeds.class.generator_version(), era_seed: (seeds.class != kaspa_consensus_core::igneum::ProgramClass::V2).then(|| seeds.era.to_string()), clean, } .line(), ); let _ = t0; Ok(()) } async fn subscribe(url: &str) -> Result<(GrpcClient, async_channel::Receiver), String> { let client = tokio::time::timeout(Duration::from_secs(5), GrpcClient::connect(url.to_string())).await.map_err(|_| "connect timed out".to_string())?.map_err(|e| e.to_string())?; tokio::time::timeout(Duration::from_secs(5), client.start_notify(ListenerId::default(), NewBlockTemplateScope {}.into())) .await .map_err(|_| "start_notify timed out".to_string())? .map_err(|e| e.to_string())?; let rx = client.notification_channel_receiver(); Ok((client, rx)) } /// The template feed: on every NewBlockTemplate notification, every second, or on demand, a fresh template and job /// for every authorised member (fetched in parallel). pub async fn template_feed(pool: Arc) { let url = pool.cfg.nodes[0].clone(); let mut sub: Option<(GrpcClient, async_channel::Receiver)> = None; let mut next_sub = Instant::now(); loop { if sub.is_none() && Instant::now() >= next_sub { match subscribe(&url).await { Ok(s) => { eprintln!("{} templates: subscribed to NewBlockTemplate at {url}", now()); sub = Some(s); } Err(e) => { eprintln!("{} templates: subscription failed ({e}); polling once a second", now()); next_sub = Instant::now() + Duration::from_secs(5); } } } let mut lost = false; match &sub { Some((_, rx)) => { tokio::select! { n = rx.recv() => { if n.is_err() { lost = true; } }, _ = pool.want_templates.notified() => {}, _ = tokio::time::sleep(Duration::from_secs(1)) => {}, } if let Some((_, rx)) = &sub { while rx.try_recv().is_ok() {} } } None => { let _ = tokio::time::timeout(Duration::from_secs(1), pool.want_templates.notified()).await; } } if lost { sub = None; next_sub = Instant::now() + Duration::from_secs(2); } let members = pool.members_snapshot(); let mut handles = Vec::with_capacity(members.len()); for m in members { let pool = pool.clone(); handles.push(tokio::spawn(async move { if let Err(e) = fetch_job(&pool, &m).await { eprintln!("{} template for member {} ({}): {e}", now(), m.id, m.worker); } })); } for h in handles { let _ = h.await; } } } /// Submits a found block to every node; records it as pending with the PPLNS snapshot. pub async fn submit_block(pool: Arc, member: Arc, mut raw: RpcRawBlock, nonce: u64, shift: u32, daa_score: u64, blue_score: u64) { raw.header.nonce = nonce; let hash = Block::try_from(raw.clone()).map(|b| b.hash().to_string()).unwrap_or_default(); let mut accepted = false; let mut detail = String::new(); for (i, c) in std::iter::once(&pool.node).chain(pool.extra_nodes.iter()).enumerate() { match tokio::time::timeout(Duration::from_secs(5), c.submit_block(raw.clone(), false)).await { Ok(Ok(r)) if r.report.is_success() => { accepted = true; } Ok(Ok(r)) => detail = format!("node {}: {:?}", i + 1, r.report), Ok(Err(e)) => detail = format!("node {}: {e}", i + 1), Err(_) => detail = format!("node {}: timed out", i + 1), } } let subsidy = block_subsidy(daa_score, 1); let reward_wei = producer_share(subsidy) as u128 * WEI_PER_SOMPI; if accepted { let mut s = pool.state.lock().unwrap(); s.block_found(BlockRec { hash: hash.clone(), daa_score, blue_score, found_ms: unix_ms(), finder: member.address.clone(), worker: member.worker.clone(), status: "pending".into(), reward_wei, fee_wei: 0, effort: 0.0, payees: Vec::new(), confirmed_ms: None, nonce: hex_u64(nonce), shift, }); member.inner.lock().unwrap().blocks += 1; println!("{} BLOCK {} daa={} by {} ({}) nonce={:#x} reward={} IGN expected, pending", now(), &hash[..16.min(hash.len())], daa_score, member.address, member.worker, nonce, crate::state::ign(reward_wei)); } else { pool.state.lock().unwrap().submit_errors += 1; eprintln!("{} block {} REJECTED by every node: {detail}", now(), &hash[..16.min(hash.len())]); } pool.want_templates.notify_one(); } /// Confirms pending blocks by walking the selected chain: a block is paid when it is a chain block or in the /// mergeset blues of one (the execution layer pays exactly those, executor.rs `b.is_blue`); a pending block more /// than `orphan_after_daa` behind the virtual with no blue merge is an orphan. pub async fn confirm_loop(pool: Arc) { let mut last_chain: Option = None; loop { tokio::time::sleep(Duration::from_secs(5)).await; let pending: Vec<(String, u64)> = { let s = pool.state.lock().unwrap(); s.blocks.iter().filter(|b| b.status == "pending").map(|b| (b.hash.clone(), b.daa_score)).collect() }; let info = match pool.node.get_block_dag_info().await { Ok(i) => i, Err(_) => continue, }; if last_chain.is_none() { last_chain = Some(info.pruning_point_hash); } if pending.is_empty() { // keep the cursor near the tip so the first pending block costs one short walk last_chain = Some(info.sink); continue; } let pending_set: HashSet = pending.iter().map(|p| p.0.clone()).collect(); let low = last_chain.unwrap(); let chain = match pool.node.get_virtual_chain_from_block(low, false, None).await { Ok(c) => c, Err(e) => { eprintln!("{} confirm: getVirtualChainFromBlock failed ({e}); restarting from the pruning point", now()); last_chain = Some(info.pruning_point_hash); continue; } }; let mut blues: HashSet = HashSet::new(); for h in &chain.added_chain_block_hashes { blues.insert(h.to_string()); if let Ok(b) = pool.node.get_block(*h, false).await && let Some(v) = b.verbose_data { for m in v.merge_set_blues_hashes { blues.insert(m.to_string()); } } } if let Some(h) = chain.added_chain_block_hashes.last() { last_chain = Some(*h); } for (hash, daa) in pending { if blues.contains(&hash) { let r = pool.state.lock().unwrap().confirm_block(&hash, pool.cfg.fee_percent); if let Some((reward, parts)) = r { println!( "{} CONFIRMED {} reward {} IGN paid to {} addresses: {}", now(), &hash[..16], crate::state::ign(reward), parts.len(), parts.iter().map(|(a, w)| format!("{}={:.6}", &a[..10], crate::state::ign(*w))).collect::>().join(" ") ); } } else if info.virtual_daa_score.saturating_sub(daa) > pool.cfg.orphan_after_daa && pool.state.lock().unwrap().orphan_block(&hash) { println!("{} ORPHAN {} (daa {}, virtual {}): not blue within {} DAA", now(), &hash[..16], daa, info.virtual_daa_score, pool.cfg.orphan_after_daa); } } let _ = pending_set; } } /// Network numbers every 5 s. pub async fn net_loop(pool: Arc) { loop { let mut n = NetInfo { network: pool.cfg.network.clone(), chain_id: pool.cfg.chain_id(), ..Default::default() }; if let Ok(i) = pool.node.get_block_dag_info().await { n.difficulty = i.difficulty; n.daa_score = i.virtual_daa_score; n.block_count = i.block_count; n.network = i.network.to_string(); } n.hashrate = pool.node.estimate_network_hashes_per_second(1000, None).await.ok().map(|h| h as f64); if let Ok(b) = pool.node.get_sink_blue_score().await { n.blue_score = b; } if let Ok(i) = pool.node.get_info().await { n.synced = i.is_synced; n.node_version = i.server_version; } let subsidy = block_subsidy(n.daa_score, 1); n.block_reward_ign = subsidy as f64 / 1e8; n.miner_reward_ign = producer_share(subsidy) as f64 / 1e8; n.updated_ms = unix_ms(); { let members = pool.members_snapshot(); let seeds = members.first().and_then(|m| m.inner.lock().unwrap().last_seeds); if let Some(seeds) = seeds { n.epoch_seed = seeds.epoch.to_string(); n.epoch_index = n.daa_score / kaspa_consensus_core::igneum::pow_epoch_blocks(); } } *pool.net.lock().unwrap() = n; tokio::time::sleep(Duration::from_secs(5)).await; } } /// Vardiff ticks for every member once a second (shares also retarget on arrival). pub async fn vardiff_loop(pool: Arc) { loop { tokio::time::sleep(Duration::from_secs(1)).await; let now_s = pool.now_s(); for m in pool.members_snapshot() { let changed = { let mut g = m.inner.lock().unwrap(); let t64 = g.current_target64; if t64 == 0 { None } else { g.vardiff.retarget(now_s, t64).map(|s| (s, share_target(t64, s))) } }; if let Some((s, st)) = changed { m.send(Msg::SetTarget { shift: s, share_target64: hex_u64(st) }.line()); println!("{} VARDIFF member {} ({}) shift -> {} (share target {:016x})", now(), m.id, m.worker, s, st); let pool = pool.clone(); let m = m.clone(); tokio::spawn(async move { let _ = fetch_job(&pool, &m).await; }); } } let _ = pool.next_member_id.load(Ordering::Relaxed); } }