605 lines
30 KiB
Rust
605 lines
30 KiB
Rust
//! 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<u8> {
|
|
let mut v = member.reveal.to_extra_data();
|
|
v.extend_from_slice(&kaspa_consensus_core::evm::miner_address_extra_data(pool_address));
|
|
v
|
|
}
|
|
|
|
/// Rewrites a template's coinbase extra data in place and re-derives the header's merkle root exactly as the node's
|
|
/// body check does (`calc_block_hash_merkle_root` over the UTXO transactions and the EVM transactions, design D7),
|
|
/// so the block the member hashes is the block the node accepts. The open pool restamps every template with the
|
|
/// share chain's parent and the window's split (`open.rs`); the node is asked for a template once a second, not
|
|
/// once per share-chain tip.
|
|
pub fn restamp_extra_data(raw: &mut RpcRawBlock, extra: &[u8]) -> Result<(), String> {
|
|
let cb = raw.transactions.first_mut().ok_or("template without a coinbase")?;
|
|
let keep = {
|
|
let p = &cb.payload;
|
|
if p.len() < 19 {
|
|
return Err("coinbase payload too short".into());
|
|
}
|
|
19 + p[18] as usize
|
|
};
|
|
if cb.payload.len() < keep {
|
|
return Err("coinbase payload shorter than its script".into());
|
|
}
|
|
cb.payload.truncate(keep);
|
|
cb.payload.extend_from_slice(extra);
|
|
let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?;
|
|
raw.header.hash_merkle_root = kaspa_consensus_core::merkle::calc_block_hash_merkle_root(block.transactions.iter(), block.evm_transactions.iter());
|
|
Ok(())
|
|
}
|
|
|
|
/// One template for one member from the node, turned into a job (`issue_job`). The template is kept per member so
|
|
/// the open pool can re-stamp it on a new share-chain parent without asking the node again (`reissue_job`).
|
|
pub async fn fetch_job(pool: &Arc<Pool>, member: &Arc<Member>) -> Result<(), String> {
|
|
let t0 = Instant::now();
|
|
// the open pool: the member's own address in the coinbase (the window's split is stamped below), never the pool's
|
|
let base_extra = match &pool.open {
|
|
Some(o) => o.base_extra_data(member),
|
|
None => extra_data(member, &pool.pool_address),
|
|
};
|
|
// at most `--template-parallel` fetches in flight: the gRPC client is one request stream per connection and a
|
|
// node under load builds a template in over a second, so a burst of one fetch per member queued past the client's
|
|
// request timeout (pool-1, 7 October 2026, 12:54Z to 13:30Z: 1,891 "RPC request timeout" lines, no job issued)
|
|
let _permit = pool.template_permits.acquire().await.map_err(|e| e.to_string())?;
|
|
let tmpl = match tokio::time::timeout(Duration::from_secs(pool.cfg.template_timeout_s), pool.node.get_block_template(pool.pay_address.clone(), base_extra)).await {
|
|
Ok(Ok(t)) => t,
|
|
Ok(Err(e)) => {
|
|
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
|
return Err(e.to_string());
|
|
}
|
|
Err(_) => {
|
|
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
|
return Err(format!("template timed out after {} s (--template-timeout-s)", pool.cfg.template_timeout_s));
|
|
}
|
|
};
|
|
pool.last_template_ok_ms.store(unix_ms(), Ordering::Relaxed);
|
|
let info = *tmpl.pow_epoch.as_ref().ok_or("the node reports no pow_epoch; a devnet-v4 line node is needed")?;
|
|
member.inner.lock().unwrap().last_template = Some((tmpl.block.clone(), info));
|
|
let r = issue_job(pool, member, tmpl.block, &info).await;
|
|
let _ = t0;
|
|
r
|
|
}
|
|
|
|
/// The open pool: the member's last template re-stamped on the current share-chain parent (no node call).
|
|
pub async fn reissue_job(pool: &Arc<Pool>, member: &Arc<Member>) -> Result<(), String> {
|
|
let (raw, info) = member.inner.lock().unwrap().last_template.clone().ok_or("no template yet")?;
|
|
issue_job(pool, member, raw, &info).await
|
|
}
|
|
|
|
/// A job on a template: the member's key in the header, the open pool's stamp, the seeds, the share target, the
|
|
/// `seeds`, `template` and `job` lines. Blocks the thread on the first use of a seed pair (the cache build), under
|
|
/// `spawn_blocking`.
|
|
async fn issue_job(pool: &Arc<Pool>, member: &Arc<Member>, mut raw: RpcRawBlock, info: &kaspa_rpc_core::RpcPowEpochInfo) -> Result<(), String> {
|
|
raw.header.vote_key_hash = member.key_hash;
|
|
// the open pool: the share chain's parent and the window's split ride in the coinbase; the share target is the
|
|
// chain's (the member's vardiff target never goes below it)
|
|
let (open_parent, open_share_target64) = match &pool.open {
|
|
Some(o) => {
|
|
let (parent, target, extra) = o.stamp_for(member, info, &raw)?;
|
|
restamp_extra_data(&mut raw, &extra)?;
|
|
(Some(parent), target)
|
|
}
|
|
None => (None, 0),
|
|
};
|
|
let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?;
|
|
let header: Header = block.header.as_ref().clone();
|
|
// 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), shadow_reps: info.latency_ladder_reps as u16 };
|
|
if let Some(o) = &pool.open {
|
|
o.note_epoch(info, &seeds);
|
|
}
|
|
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<Hash> = 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;
|
|
// the open pool: a new share-chain parent supersedes the earlier jobs as a new tip does
|
|
let clean = g.last_parents.as_ref() != Some(&parents) || g.jobs.back().is_some_and(|j| j.open_parent != open_parent);
|
|
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);
|
|
// the open pool: the member never sends shares above the chain's target (a chain share is what pays);
|
|
// a member whose vardiff target is easier than the chain's still sends its own shares for the hashrate
|
|
// figure and the chain shares among them
|
|
let share_t = share_target(t64, shift).max(open_share_target64);
|
|
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,
|
|
open_parent,
|
|
open_share_target64,
|
|
});
|
|
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,
|
|
shadow_reps: info.latency_ladder_reps,
|
|
next_shadow_reps: info.next_latency_ladder_reps,
|
|
}
|
|
.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()),
|
|
shadow_reps: info.latency_ladder_reps,
|
|
clean,
|
|
}
|
|
.line(),
|
|
);
|
|
Ok(())
|
|
}
|
|
|
|
async fn subscribe(url: &str) -> Result<(GrpcClient, async_channel::Receiver<Notification>), 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<Pool>) {
|
|
let url = pool.cfg.nodes[0].clone();
|
|
let mut sub: Option<(GrpcClient, async_channel::Receiver<Notification>)> = 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;
|
|
// the open pool: a new share-chain tip re-stamps every member's last template without a node call
|
|
let mut restamp_only = false;
|
|
let tip_moved = async {
|
|
match &pool.open {
|
|
Some(o) => o.tip_changed.notified().await,
|
|
None => std::future::pending::<()>().await,
|
|
}
|
|
};
|
|
match &sub {
|
|
Some((_, rx)) => {
|
|
tokio::select! {
|
|
n = rx.recv() => { if n.is_err() { lost = true; } },
|
|
_ = pool.want_templates.notified() => {},
|
|
_ = tip_moved => { restamp_only = true; },
|
|
_ = tokio::time::sleep(Duration::from_secs(1)) => {},
|
|
}
|
|
if let Some((_, rx)) = &sub {
|
|
while rx.try_recv().is_ok() {}
|
|
}
|
|
}
|
|
None => {
|
|
tokio::select! {
|
|
_ = tokio::time::timeout(Duration::from_secs(1), pool.want_templates.notified()) => {},
|
|
_ = tip_moved => { restamp_only = true; },
|
|
}
|
|
}
|
|
}
|
|
if lost {
|
|
sub = None;
|
|
next_sub = Instant::now() + Duration::from_secs(2);
|
|
}
|
|
let members = pool.members_snapshot();
|
|
if restamp_only {
|
|
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) = reissue_job(&pool, &m).await
|
|
&& e != "no template yet"
|
|
{
|
|
eprintln!("{} re-stamp for member {} ({}): {e}", now(), m.id, m.worker);
|
|
}
|
|
}));
|
|
}
|
|
for h in handles {
|
|
let _ = h.await;
|
|
}
|
|
continue;
|
|
}
|
|
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<Pool>, member: Arc<Member>, 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();
|
|
if let Some(o) = &pool.open {
|
|
o.note_block(&raw, &hash, daa_score, blue_score, reward_wei, &member.address, &mut s);
|
|
} else {
|
|
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();
|
|
}
|
|
|
|
/// Where the confirmation walk continues from after a failed `getVirtualChainFromBlock`: the same cursor, retried
|
|
/// next tick. The first version restarted from the pruning point, and on a 120,000-block devnet chain that is a
|
|
/// `get_block` per chain block through the connection the templates shared: one RPC timeout under load (7 October 2026,
|
|
/// 06:01Z) started a walk that starved every template request for the rest of the window. A cursor the node no
|
|
/// longer knows (pruned, or from a state file of another chain) moves to the sink, and the pending blocks older than
|
|
/// it resolve by the orphan rule.
|
|
pub fn cursor_after_failure(cursor: Hash, sink: Hash, cursor_known: bool) -> Hash {
|
|
if cursor_known { cursor } else { sink }
|
|
}
|
|
|
|
/// At most this many chain blocks (and `get_block` calls) per 5 s tick; the rest continue next tick.
|
|
pub const WALK_MAX_PER_TICK: usize = 600;
|
|
|
|
/// 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. On its own connection (`pool.walker`).
|
|
pub async fn confirm_loop(pool: Arc<Pool>) {
|
|
let mut last_chain: Option<Hash> = 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.walker.get_block_dag_info().await {
|
|
Ok(i) => i,
|
|
Err(e) => {
|
|
eprintln!("{} confirm: getBlockDagInfo failed ({e})", now());
|
|
continue;
|
|
}
|
|
};
|
|
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;
|
|
}
|
|
// a first tick with pending blocks (a state file kept across a restart): from the sink, never the pruning
|
|
// point; blocks older than the sink that were blue are a payout lost to the restart, said once
|
|
let low = match last_chain {
|
|
Some(h) => h,
|
|
None => {
|
|
println!("{} confirm: {} pending block(s) at start; the walk begins at the sink, older ones resolve by the orphan rule", now(), pending.len());
|
|
last_chain = Some(info.sink);
|
|
info.sink
|
|
}
|
|
};
|
|
let t0 = Instant::now();
|
|
let chain = match pool.walker.get_virtual_chain_from_block(low, false, None).await {
|
|
Ok(c) => c,
|
|
Err(e) => {
|
|
let known = pool.walker.get_block(low, false).await.is_ok();
|
|
let next = cursor_after_failure(low, info.sink, known);
|
|
eprintln!("{} confirm: getVirtualChainFromBlock from {} failed ({e}); cursor {}", now(), low, if next == low { "kept, retried next tick".to_string() } else { format!("unknown to the node, moved to the sink {next}") });
|
|
last_chain = Some(next);
|
|
continue;
|
|
}
|
|
};
|
|
let added = &chain.added_chain_block_hashes;
|
|
let take = added.len().min(WALK_MAX_PER_TICK);
|
|
let mut blues: HashSet<String> = HashSet::new();
|
|
let mut fetched = 0usize;
|
|
for h in &added[..take] {
|
|
blues.insert(h.to_string());
|
|
match pool.walker.get_block(*h, false).await {
|
|
Ok(b) => {
|
|
fetched += 1;
|
|
if let Some(v) = b.verbose_data {
|
|
for m in v.merge_set_blues_hashes {
|
|
blues.insert(m.to_string());
|
|
}
|
|
}
|
|
}
|
|
Err(e) => {
|
|
eprintln!("{} confirm: getBlock {h} failed ({e}); the walk stops here and continues next tick", now());
|
|
break;
|
|
}
|
|
}
|
|
last_chain = Some(*h);
|
|
}
|
|
if added.len() > WALK_MAX_PER_TICK || t0.elapsed() > Duration::from_secs(2) {
|
|
println!("{} confirm: walked {fetched} of {} chain blocks in {:.0} ms ({} pending)", now(), added.len(), t0.elapsed().as_secs_f64() * 1e3, pending.len());
|
|
}
|
|
for (hash, daa) in pending {
|
|
if blues.contains(&hash) {
|
|
// the open pool: the coinbase paid the split; the ledger records it, no balance moves (open.rs)
|
|
let fee = if pool.open.is_some() { 0.0 } else { pool.cfg.fee_percent };
|
|
let r = pool.state.lock().unwrap().confirm_block(&hash, fee, pool.open.is_some());
|
|
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::<Vec<_>>().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);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// The daemon's own status line every 30 s (7 October 2026: a night where four members sat connected with no job for
|
|
/// 20 minutes was visible only as per-member "template timed out" lines).
|
|
pub async fn status_loop(pool: Arc<Pool>) {
|
|
loop {
|
|
tokio::time::sleep(Duration::from_secs(30)).await;
|
|
let members = pool.members_snapshot().len();
|
|
let failures = pool.template_failures.load(Ordering::Relaxed);
|
|
let last_ok = pool.last_template_ok_ms.load(Ordering::Relaxed);
|
|
let jobs = pool.next_job_id.load(Ordering::Relaxed);
|
|
let (accepted, rejected, blocks) = {
|
|
let s = pool.state.lock().unwrap();
|
|
(s.accepted, s.rejected, s.blocks_found)
|
|
};
|
|
let node_state = if last_ok == 0 {
|
|
"NO TEMPLATE YET from the node: no job has been issued".to_string()
|
|
} else {
|
|
let ago = unix_ms().saturating_sub(last_ok) / 1000;
|
|
if ago > 30 { format!("NODE NOT ANSWERING: last template {ago} s ago") } else { format!("node ok, last template {ago} s ago") }
|
|
};
|
|
println!("{} STATUS members={members} jobs_issued={jobs} shares_accepted={accepted} shares_rejected={rejected} blocks={blocks} template_failures={failures}; {node_state}", now());
|
|
}
|
|
}
|
|
|
|
/// Network numbers every 5 s.
|
|
pub async fn net_loop(pool: Arc<Pool>) {
|
|
loop {
|
|
let mut n = NetInfo { network: pool.cfg.network.clone(), chain_id: pool.cfg.chain_id(), finality: "unknown".into(), ..Default::default() };
|
|
if let Ok(i) = pool.walker.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.walker.estimate_network_hashes_per_second(1000, None).await.ok().map(|h| h as f64);
|
|
if let Ok(b) = pool.walker.get_sink_blue_score().await {
|
|
n.blue_score = b;
|
|
}
|
|
if let Ok(i) = pool.walker.get_info().await {
|
|
n.synced = i.is_synced;
|
|
n.node_version = i.server_version;
|
|
}
|
|
// Q73: the network's finality state, for the page; payouts follow blue confirmation, not finality
|
|
match pool.walker.get_finality_checkpoints(1).await {
|
|
Ok(f) => {
|
|
n.finality = if f.finality_active { "active".to_string() } else { f.finality_reason.clone() };
|
|
n.finality_locked_index = f.latest_locked_index;
|
|
n.finality_locked_age_s = f.checkpoints.iter().find(|c| c.index == f.latest_locked_index).map(|c| n.daa_score.saturating_sub(c.locked_at_daa)).unwrap_or(0);
|
|
}
|
|
Err(_) => n.finality = "unknown".into(),
|
|
}
|
|
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<Pool>) {
|
|
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);
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod walk_tests {
|
|
use super::*;
|
|
|
|
/// The cursor rule after a failed chain call: known cursor kept (the 6 October shape restarted from the pruning
|
|
/// point and walked 120,000 blocks), unknown cursor moved to the sink.
|
|
#[test]
|
|
fn a_failed_chain_call_keeps_a_known_cursor_and_never_goes_to_the_pruning_point() {
|
|
let cursor = Hash::from_bytes([1u8; 32]);
|
|
let sink = Hash::from_bytes([2u8; 32]);
|
|
let pruning = Hash::from_bytes([3u8; 32]);
|
|
assert_eq!(cursor_after_failure(cursor, sink, true), cursor);
|
|
assert_eq!(cursor_after_failure(cursor, sink, false), sink);
|
|
assert_ne!(cursor_after_failure(cursor, sink, false), pruning);
|
|
assert!(WALK_MAX_PER_TICK <= 1000, "a tick's walk stays bounded");
|
|
}
|
|
}
|