diff --git a/miner/src/core/args.rs b/miner/src/core/args.rs new file mode 100644 index 0000000..33303d3 --- /dev/null +++ b/miner/src/core/args.rs @@ -0,0 +1,68 @@ +use clap::Parser; +use nyks_protocol::consensus::network::Network; +use nyks_rpc_client::RpcApi; +use nyks_rpc_client::http::HttpClient; +use nyks_standards::wallet::keys::address::Address; +use nyks_standards::wallet::keys::address::Recipient; + +#[derive(Parser)] +#[command(name = "nyks-miner")] +#[command(about = "A nyks CPU miner")] +pub struct Args { + /// RPC URL to use (JSON/HTTP) + #[arg(long, default_value = "http://localhost:9797")] + pub rpc_url: String, + + /// Address to mine for (coinbase reward receiver) + #[arg(long)] + pub address: String, + + /// Network we are going to mine on. + #[arg(long, default_value_t = Network::Main)] + pub network: Network, + + /// Minimum guesser reward as a percentage of total block reward (0-100), default 10. + #[arg(long, default_value_t = 10.0)] + pub min_percentage_reward: f64, +} + +impl Args { + /// Validates the percentage is in range and converts it to a fraction + /// (e.g. 10.0 -> 0.10) for internal use. + pub fn min_reward_fraction(&self) -> f64 { + assert!( + (0.0..=100.0).contains(&self.min_percentage_reward), + "min-percentage-reward must be between 0 and 100, got {}", + self.min_percentage_reward + ); + + self.min_percentage_reward / 100.0 + } + + /// Creates an RPC client and verifies that the connected node is on the + /// expected network. + pub async fn rpc_client(&self) -> HttpClient { + let client = HttpClient::new(self.rpc_url.clone()); + + let remote_network = client + .network() + .await + .expect("Failed to connect to RPC node") + .network + .parse::() + .unwrap(); + + assert_eq!( + self.network, remote_network, + "Network mismatch: expected {:?}, connected node is on {:?}", + self.network, remote_network, + ); + + client + } + + /// Validates that `address` is a well-formed address for the selected network. + pub fn validate_address(&self) { + Address::from_bech32m(&self.address, self.network).expect("Invalid address"); + } +} diff --git a/miner/src/core/mod.rs b/miner/src/core/mod.rs new file mode 100644 index 0000000..6e10f4a --- /dev/null +++ b/miner/src/core/mod.rs @@ -0,0 +1 @@ +pub mod args; diff --git a/miner/src/flow.rs b/miner/src/flow.rs deleted file mode 100644 index f1c66aa..0000000 --- a/miner/src/flow.rs +++ /dev/null @@ -1,70 +0,0 @@ -use std::sync::Arc; -use std::time::Duration; - -use num_traits::Zero; -use nyks_protocol::consensus::type_scripts::native_currency_amount::NativeCurrencyAmount; -use nyks_rpc_client::RpcApi; -use nyks_rpc_client::http::HttpClient; -use tokio::sync::RwLock; -use tracing::info; - -use crate::guesser::Guesser; - -#[derive(Clone, Debug)] -pub struct Miner { - client: HttpClient, - address: String, - guesser_reward: Arc>, - guesser: Guesser, -} - -impl Miner { - pub fn new(client: HttpClient, address: String) -> Self { - Miner { - client: client.clone(), - address, - guesser_reward: Arc::new(RwLock::new(NativeCurrencyAmount::zero())), - guesser: Guesser::new(client), - } - } - - pub async fn main_loop(&self) { - let mut interval = tokio::time::interval(Duration::from_secs(5)); - - loop { - interval.tick().await; - self.scan_templates().await; - } - } - - pub async fn scan_templates(&self) { - let template = self - .client - .get_block_template(Some(self.address.clone())) - .await - .unwrap() - .template; - - if let Some(template) = template { - let new_guesser_reward = template.metadata.total_guesser_reward.0; - let mut current_guesser_reward = self.guesser_reward.write().await; - - if new_guesser_reward > *current_guesser_reward { - info!( - "Switching to mining of new template with {} NYKS reward.", - new_guesser_reward - ); - *current_guesser_reward = new_guesser_reward; - - self.guesser - .override_task(template.metadata.prev_block, template) - .await; - } - } else if self.guesser.is_running().await { - info!("New tip is found, waiting for a composed template..."); - - *self.guesser_reward.write().await = NativeCurrencyAmount::zero(); - self.guesser.stop().await; - } - } -} diff --git a/miner/src/main.rs b/miner/src/main.rs index 0b5084a..3dc1e8f 100644 --- a/miner/src/main.rs +++ b/miner/src/main.rs @@ -2,26 +2,14 @@ use std::panic; use anyhow::Result; use clap::Parser; -use nyks_rpc_client::http::HttpClient; use tracing::info; use tracing_subscriber::EnvFilter; -use crate::flow::Miner; - -pub mod flow; -pub mod guesser; - -#[derive(Parser)] -#[command(name = "nyks-miner")] -#[command(about = "A nyks CPU miner")] -struct Args { - /// Address to mine for (coinbase reward receiver) - #[arg(long)] - address: String, - /// RPC URL to use (JSON/HTTP) - #[arg(long)] - rpc_url: String, -} +use crate::core::args::Args; +use crate::miner::flow::Miner; + +pub mod core; +pub mod miner; #[tokio::main] async fn main() -> Result<()> { @@ -37,8 +25,10 @@ async fn main() -> Result<()> { info!("Initializing nyks-miner, the operator of nyks blocks..."); - let client = HttpClient::new(args.rpc_url); - let miner = Miner::new(client, args.address); + args.validate_address(); + + let client = args.rpc_client().await; + let miner = Miner::new(client, args.address.clone(), args.min_reward_fraction()); miner.main_loop().await; diff --git a/miner/src/miner/flow.rs b/miner/src/miner/flow.rs new file mode 100644 index 0000000..619e858 --- /dev/null +++ b/miner/src/miner/flow.rs @@ -0,0 +1,110 @@ +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use num_traits::Zero; +use nyks_protocol::consensus::block::Block; +use nyks_protocol::consensus::type_scripts::native_currency_amount::NativeCurrencyAmount; +use nyks_rpc_client::RpcApi; +use nyks_rpc_client::http::HttpClient; +use tokio::sync::RwLock; +use tracing::debug; +use tracing::info; + +use crate::miner::guesser::Guesser; + +#[derive(Clone, Debug)] +pub struct Miner { + client: HttpClient, + address: String, + min_reward_fraction: f64, + guesser_reward: Arc>, + guesser: Guesser, + composing_since: Arc>>, +} + +impl Miner { + pub fn new(client: HttpClient, address: String, min_reward_fraction: f64) -> Self { + Miner { + client: client.clone(), + address, + min_reward_fraction, + guesser_reward: Arc::new(RwLock::new(NativeCurrencyAmount::zero())), + guesser: Guesser::new(client), + composing_since: Arc::new(RwLock::new(None)), + } + } + + pub async fn main_loop(&self) { + let mut interval = tokio::time::interval(Duration::from_secs(5)); + loop { + interval.tick().await; + self.scan_templates().await; + } + } + + pub async fn scan_templates(&self) { + let template = self + .client + .get_block_template(Some(self.address.clone())) + .await + .unwrap() + .template; + + if let Some(template) = template { + // If we were waiting on a template, capture how long composing took + let composing_time = { + let mut composing_since = self.composing_since.write().await; + composing_since + .take() + .map(|start| start.elapsed().as_secs_f64()) + }; + + let new_guesser_reward = template.metadata.total_guesser_reward.0; + let total_reward = Block::block_subsidy(template.block.kernel.header.height); + let guesser_share = new_guesser_reward.to_nau() as f64 / total_reward.to_nau() as f64; + + if guesser_share < self.min_reward_fraction { + debug!( + "Skipping template: guesser share {:.2}% below minimum {:.2}%.", + guesser_share * 100.0, + self.min_reward_fraction * 100.0 + ); + return; + } + + let mut current_guesser_reward = self.guesser_reward.write().await; + + if new_guesser_reward > *current_guesser_reward { + match composing_time { + Some(secs) => info!( + "Switching to mining of new template with {} NYKS reward (composed in {:.2}s).", + new_guesser_reward, secs + ), + None => info!( + "Switching to mining of new template with {} NYKS reward.", + new_guesser_reward + ), + } + *current_guesser_reward = new_guesser_reward; + + self.guesser + .override_task(template.metadata.prev_block, template) + .await; + } + } else { + { + let mut composing_since = self.composing_since.write().await; + if composing_since.is_none() { + info!("Waiting for a template..."); + *composing_since = Some(Instant::now()); + } + } + + if self.guesser.is_running().await { + info!("New tip is found, waiting for a composed template..."); + *self.guesser_reward.write().await = NativeCurrencyAmount::zero(); + self.guesser.stop().await; + } + } + } +} diff --git a/miner/src/guesser.rs b/miner/src/miner/guesser.rs similarity index 63% rename from miner/src/guesser.rs rename to miner/src/miner/guesser.rs index 34e3380..645d6d4 100644 --- a/miner/src/guesser.rs +++ b/miner/src/miner/guesser.rs @@ -1,6 +1,9 @@ use std::sync::Arc; use std::sync::atomic::AtomicBool; +use std::sync::atomic::AtomicU64; use std::sync::atomic::Ordering; +use std::time::Duration; +use std::time::Instant; use nyks_protocol::BFieldElement; use nyks_protocol::consensus::block::block_header::BlockPow; @@ -29,7 +32,7 @@ struct MinerBuffer { // Holds the cancel flag so stop() can signal the blocking thread directly. #[derive(Debug)] struct MinerTask { - template: RpcBlockTemplate, + template: Arc, cancel: Arc, handle: JoinHandle<()>, } @@ -58,15 +61,18 @@ impl Guesser { self.stop().await; let guesser_buffer = self.get_or_recompute_buffer(prev_block_digest).await; + let template = Arc::new(template); + info!( "Switching to mining template {}...", template.block.kernel.mast_hash().to_hex() ); + let cancel = Arc::new(AtomicBool::new(false)); + let client = self.client.clone(); let task_template = template.clone(); - let cancel = Arc::new(AtomicBool::new(false)); - let task_cancel = Arc::clone(&cancel); + let task_cancel = cancel.clone(); let handle = tokio::spawn(async move { Self::run_mining_task(client, task_template, guesser_buffer, task_cancel).await; @@ -100,11 +106,11 @@ impl Guesser { *guard = Some(MinerBuffer { digest: prev_block_digest, - buffer: Arc::clone(&new_buffer), + buffer: new_buffer.clone(), }); new_buffer } else { - Arc::clone(&guard.as_ref().unwrap().buffer) + guard.as_ref().unwrap().buffer.clone() } } @@ -123,43 +129,79 @@ impl Guesser { async fn run_mining_task( client: HttpClient, - template: RpcBlockTemplate, + template: Arc, guesser_buffer: Arc>, cancel: Arc, ) { - let mining_template = template.clone(); - let mine_result = - tokio::task::spawn_blocking(move || mine(&mining_template, &guesser_buffer, &cancel)) - .await; - let pow = match mine_result { - Ok(Some(pow)) => pow, - - Ok(None) => { - info!("Mining stopped (cancelled or nonce space exhausted)."); - return; - } + let hashes_done = Arc::new(AtomicU64::new(0)); + let logger_handle = Self::spawn_hashrate_logger(hashes_done.clone(), cancel.clone()); - Err(e) => { - info!("Mining task panicked: {e}"); - return; + loop { + if cancel.load(Ordering::Relaxed) { + break; } - }; - info!("Found the solution! Submitting to node..."); + let mining_template = template.clone(); + let guesser_buffer_for_mine = guesser_buffer.clone(); + let cancel_for_mine = cancel.clone(); + let hashes_done_for_mine = hashes_done.clone(); + + let mine_result = tokio::task::spawn_blocking(move || { + mine( + &mining_template, + &guesser_buffer_for_mine, + &cancel_for_mine, + &hashes_done_for_mine, + ) + }) + .await; + + let pow = match mine_result { + Ok(Some(pow)) => pow, + Ok(None) => { + info!("Mining stopped (cancelled or nonce space exhausted)."); + break; + } + Err(e) => { + info!("Mining task panicked: {e}"); + break; + } + }; + + info!("Found the solution! Submitting to node..."); - match client.submit_block(template.block, pow.into()).await { - Ok(response) => { - if response.success { + match client + .submit_block(template.block.clone(), pow.into()) + .await + { + Ok(response) if response.success => { info!("Block is accepted by node."); - } else { - warn!("Block is rejected by node, channel error."); + break; } - } - Err(e) => { - // TODO: RETRY - warn!("Block is rejected by node, reason: {}", e); + Ok(_) => warn!("Block is rejected by node, channel error. Retrying..."), + Err(e) => warn!("Block is rejected by node, reason: {}. Retrying...", e), } } + + logger_handle.abort(); + } + + fn spawn_hashrate_logger( + hashes_done: Arc, + cancel: Arc, + ) -> JoinHandle<()> { + tokio::spawn(async move { + let mut interval = tokio::time::interval(Duration::from_secs(10)); + let start = Instant::now(); + interval.tick().await; + + while !cancel.load(Ordering::Relaxed) { + interval.tick().await; + let total = hashes_done.load(Ordering::Relaxed); + let rate = total as f64 / start.elapsed().as_secs_f64().max(0.001); + info!("Hashrate: {:.2} MH/s ({total} total)", rate / 1e6); + } + }) } } @@ -168,6 +210,7 @@ fn mine( template: &RpcBlockTemplate, guesser_buffer: &GuesserBuffer, cancel: &AtomicBool, + hashes_done: &AtomicU64, ) -> Option { // Check for cancellation every ~500k nonces rather than every iteration. const CHECKPOINT_DISTANCE: u64 = 1 << 19; @@ -190,8 +233,11 @@ fn mine( (0u64..u64::MAX) .into_par_iter() .find_map_any(|i| { - if i % CHECKPOINT_DISTANCE == 0 && cancel.load(Ordering::Relaxed) { - return Some(None); + if i % CHECKPOINT_DISTANCE == 0 { + hashes_done.fetch_add(CHECKPOINT_DISTANCE, Ordering::Relaxed); + if cancel.load(Ordering::Relaxed) { + return Some(None); + } } let nonce = Digest(bfe_array![n0, n1, n2, n3, i]); diff --git a/miner/src/miner/mod.rs b/miner/src/miner/mod.rs new file mode 100644 index 0000000..0bea828 --- /dev/null +++ b/miner/src/miner/mod.rs @@ -0,0 +1,2 @@ +pub mod flow; +pub mod guesser;