Browse Source

minerd: cache current randomx key vms and pregenerate next ones in background

skoupidi 8 months ago
parent
commit
5155425c19

+ 1 - 0
Cargo.lock

@@ -4318,6 +4318,7 @@ dependencies = [
  "darkfi-serial",
  "easy-parallel",
  "num-bigint",
+ "randomx",
  "serde",
  "signal-hook",
  "signal-hook-async-std",

+ 1 - 10
bin/darkfid/src/lib.rs

@@ -33,9 +33,7 @@ use darkfi::{
         settings::RpcSettings,
     },
     system::{ExecutorPtr, StoppableTask, StoppableTaskPtr},
-    validator::{
-        consensus::Fork, utils::best_fork_index, Validator, ValidatorConfig, ValidatorPtr,
-    },
+    validator::{Validator, ValidatorConfig, ValidatorPtr},
     zk::{empty_witnesses, ProvingKey, ZkCircuit},
     zkas::ZkBinary,
     Error, Result,
@@ -112,13 +110,6 @@ impl DarkfiNode {
             powrewardv1_zk,
         }))
     }
-
-    /// Auxiliary function to grab best current fork.
-    pub async fn best_current_fork(&self) -> Result<Fork> {
-        let forks = self.validator.consensus.forks.read().await;
-        let index = best_fork_index(&forks)?;
-        forks[index].full_clone()
-    }
 }
 
 /// ZK data used to generate the "coinbase" transaction in a block

+ 1 - 0
bin/darkfid/src/rpc.rs

@@ -84,6 +84,7 @@ impl RequestHandler<DefaultRpcHandler> for DarkfiNode {
             // =============
             // Miner methods
             // =============
+            "miner.get_current_randomx_keys" => self.miner_get_current_randomx_keys(req.id, req.params).await,
             "miner.get_header" => self.miner_get_header(req.id, req.params).await,
             "miner.submit_solution" => self.miner_submit_solution(req.id, req.params).await,
 

+ 58 - 2
bin/darkfid/src/rpc_miner.rs

@@ -66,6 +66,8 @@ pub struct BlockTemplate {
     pub block: BlockInfo,
     /// The base64 encoded RandomX key used
     randomx_key: String,
+    /// The base64 encoded next RandomX key used
+    next_randomx_key: String,
     /// The base64 encoded mining target used
     target: String,
     /// The signing secret for this block
@@ -73,6 +75,51 @@ pub struct BlockTemplate {
 }
 
 impl DarkfiNode {
+    // RPCAPI:
+    // Queries the validator for the current and next RandomX keys.
+    // If no forks exist, retrieves the canonical ones.
+    // Returns the current and next RandomX keys, both encoded as
+    // base64 strings.
+    //
+    // **Params:**
+    // * `None`
+    //
+    // **Returns:**
+    // * `String`: Current RandomX key (base64 encoded)
+    // * `String`: Current next RandomX key (base64 encoded)
+    //
+    // --> {"jsonrpc": "2.0", "method": "miner.get_current_randomx_keys", "params": [], "id": 1}
+    // <-- {"jsonrpc": "2.0", "result": ["randomx_key", "next_randomx_key"], "id": 1}
+    pub async fn miner_get_current_randomx_keys(&self, id: u16, params: JsonValue) -> JsonResult {
+        // Verify request params
+        let Some(params) = params.get::<Vec<JsonValue>>() else {
+            return JsonError::new(InvalidParams, None, id).into()
+        };
+        if !params.is_empty() {
+            return JsonError::new(InvalidParams, None, id).into()
+        }
+
+        // Grab current RandomX keys
+        let (randomx_key, next_randomx_key) = match self.validator.current_randomx_keys().await {
+            Ok(keys) => keys,
+            Err(e) => {
+                error!(
+                    target: "darkfid::rpc::miner_get_current_randomx_keys",
+                    "[RPC] Retrieving current RandomX keys failed: {e}",
+                );
+                return JsonError::new(ErrorCode::InternalError, None, id).into()
+            }
+        };
+
+        // Encode them and build response
+        let response = JsonValue::Array(vec![
+            JsonValue::String(base64::encode(&serialize_async(&randomx_key).await)),
+            JsonValue::String(base64::encode(&serialize_async(&next_randomx_key).await)),
+        ]);
+
+        JsonResponse::new(response, id).into()
+    }
+
     // RPCAPI:
     // Queries the validator for the current best fork next header to
     // mine.
@@ -87,11 +134,12 @@ impl DarkfiNode {
     //
     // **Returns:**
     // * `String`: Current best fork RandomX key (base64 encoded)
+    // * `String`: Current best fork next RandomX key (base64 encoded)
     // * `String`: Current best fork mining target (base64 encoded)
     // * `String`: Current best fork next block header (base64 encoded)
     //
     // --> {"jsonrpc": "2.0", "method": "miner.get_header", "params": {"header": "hash", "recipient": "address"}, "id": 1}
-    // <-- {"jsonrpc": "2.0", "result": ["randomx_key", "target", "header"], "id": 1}
+    // <-- {"jsonrpc": "2.0", "result": ["randomx_key", "next_randomx_key", "target", "header"], "id": 1}
     pub async fn miner_get_header(&self, id: u16, params: JsonValue) -> JsonResult {
         // Check if node is synced before responding to miner
         if !*self.validator.synced.read().await {
@@ -176,7 +224,7 @@ impl DarkfiNode {
         // released when this function exits.
         let address_bytes = serialize_async(&(recipient, spend_hook, user_data)).await;
         let mut blocktemplates = self.blocktemplates.lock().await;
-        let mut extended_fork = match self.best_current_fork().await {
+        let mut extended_fork = match self.validator.best_current_fork().await {
             Ok(f) => f,
             Err(e) => {
                 error!(
@@ -202,6 +250,7 @@ impl DarkfiNode {
                     JsonResponse::new(
                         JsonValue::Array(vec![
                             JsonValue::String(blocktemplate.randomx_key.clone()),
+                            JsonValue::String(blocktemplate.next_randomx_key.clone()),
                             JsonValue::String(blocktemplate.target.clone()),
                             JsonValue::String(base64::encode(
                                 &serialize_async(&blocktemplate.block.header).await,
@@ -263,6 +312,11 @@ impl DarkfiNode {
             base64::encode(&serialize_async(&extended_fork.module.darkfi_rx_keys.0).await)
         };
 
+        // Grab the next RandomX key to use so miner can pregenerate
+        // mining VMs.
+        let next_randomx_key =
+            base64::encode(&serialize_async(&extended_fork.module.darkfi_rx_keys.1).await);
+
         // Convert the target
         let target = base64::encode(&target.to_bytes_le());
 
@@ -270,6 +324,7 @@ impl DarkfiNode {
         let blocktemplate = BlockTemplate {
             block,
             randomx_key: randomx_key.clone(),
+            next_randomx_key: next_randomx_key.clone(),
             target: target.clone(),
             secret,
         };
@@ -286,6 +341,7 @@ impl DarkfiNode {
 
         let response = JsonValue::Array(vec![
             JsonValue::String(randomx_key),
+            JsonValue::String(next_randomx_key),
             JsonValue::String(target),
             JsonValue::String(header),
         ]);

+ 1 - 1
bin/darkfid/src/rpc_xmr.rs

@@ -231,7 +231,7 @@ impl DarkfiNode {
         // multiple times and potentially missing a job. The lock is
         // released when this function exits.
         let mut mm_blocktemplates = self.mm_blocktemplates.lock().await;
-        let mut extended_fork = match self.best_current_fork().await {
+        let mut extended_fork = match self.validator.best_current_fork().await {
             Ok(f) => f,
             Err(e) => {
                 error!(

+ 1 - 0
bin/minerd/Cargo.toml

@@ -16,6 +16,7 @@ darkfi-serial = {version = "0.5.0", features = ["async"]}
 
 # Misc
 bs58 = "0.5.1"
+randomx = {git = "https://codeberg.org/darkrenaissance/RandomX"}
 tracing = "0.1.41"
 num-bigint = "0.4.6"
 

+ 44 - 24
bin/minerd/src/lib.rs

@@ -80,10 +80,10 @@ pub type MinerNodePtr = Arc<MinerNode>;
 pub struct MinerNode {
     /// Node configuration
     config: MinerNodeConfig,
-    /// Sender to stop miner threads
-    sender: Sender<()>,
-    /// Receiver to stop miner threads
-    stop_signal: Receiver<()>,
+    /// Sender and receiver to stop mining threads
+    mining_channel: (Sender<()>, Receiver<()>),
+    /// Sender and receiver to stop background threads
+    background_channel: (Sender<()>, Receiver<()>),
     /// JSON-RPC client to execute requests to darkfid daemon
     rpc_client: RwLock<DarkfidRpcClient>,
 }
@@ -91,40 +91,60 @@ pub struct MinerNode {
 impl MinerNode {
     pub async fn new(config: MinerNodeConfig, endpoint: Url, ex: &ExecutorPtr) -> MinerNodePtr {
         // Initialize the smol channels to send signal between the threads
-        let (sender, stop_signal) = smol::channel::bounded(1);
+        let mining_channel = smol::channel::bounded(1);
+        let background_channel = smol::channel::bounded(1);
 
         // Initialize JSON-RPC client
         let rpc_client = RwLock::new(DarkfidRpcClient::new(endpoint, ex.clone()).await);
 
-        Arc::new(Self { config, sender, stop_signal, rpc_client })
+        Arc::new(Self { config, mining_channel, background_channel, rpc_client })
     }
 
-    /// Auxiliary function to abort pending job.
-    pub async fn abort_pending(&self) {
-        // Check if a pending request is being processed
-        debug!(target: "minerd::abort_pending", "Checking if a pending job is being processed...");
-        if self.stop_signal.receiver_count() <= 1 {
-            debug!(target: "minerd::rpc", "No pending job!");
+    /// Auxiliary function to abort all pending tasks.
+    pub async fn abort(&self) {
+        self.abort_mining().await;
+        self.abort_background().await;
+    }
+
+    /// Auxiliary function to abort pending mining task.
+    pub async fn abort_mining(&self) {
+        Self::abort_task(&self.mining_channel.0, &self.mining_channel.1, "mining").await;
+    }
+
+    /// Auxiliary function to abort pending background Randomx VMs
+    /// generation task.
+    pub async fn abort_background(&self) {
+        Self::abort_task(&self.background_channel.0, &self.background_channel.1, "VMs generation")
+            .await;
+    }
+
+    /// Auxiliary function to abort pending task by signaling provided
+    /// channels.
+    async fn abort_task(sender: &Sender<()>, stop_signal: &Receiver<()>, task: &str) {
+        // Check if a pending task is being processed
+        debug!(target: "minerd::abort_task", "Checking if a pending {task} task is being processed...");
+        if stop_signal.receiver_count() <= 1 {
+            debug!(target: "minerd::abort_task", "No pending {task} task!");
             return
         }
 
-        info!(target: "minerd::abort_pending", "Pending job is in progress, sending stop signal...");
+        info!(target: "minerd::abort_task", "Pending {task} is in progress, sending stop signal...");
         // Send stop signal to worker
-        if let Err(e) = self.sender.try_send(()) {
-            error!(target: "minerd::abort_pending", "Failed to stop pending job: {e}");
+        if let Err(e) = sender.try_send(()) {
+            error!(target: "minerd::abort_task", "Failed to stop pending {task} task: {e}");
             return
         }
 
         // Wait for worker to terminate
-        info!(target: "minerd::abort_pending", "Waiting for job to terminate...");
-        while self.stop_signal.receiver_count() > 1 {
+        info!(target: "minerd::abort_task", "Waiting for {task} task to terminate...");
+        while stop_signal.receiver_count() > 1 {
             sleep(1).await;
         }
-        info!(target: "minerd::abort_pending", "Pending job terminated!");
+        info!(target: "minerd::abort_task", "Pending {task} task terminated!");
 
         // Consume channel item so its empty again
-        if let Err(e) = self.stop_signal.try_recv() {
-            error!(target: "minerd::abort_pending", "Failed to cleanup stop signal channel: {e}");
+        if let Err(e) = stop_signal.try_recv() {
+            error!(target: "minerd::abort_task", "Failed to cleanup stop signal channel: {e}");
         }
     }
 }
@@ -185,14 +205,14 @@ impl Minerd {
     pub async fn stop(&self) {
         info!(target: "minerd::Minerd::stop", "Terminating mining daemon...");
 
+        // Stop the mining node
+        info!(target: "minerd::Minerd::stop", "Stopping miner background tasks...");
+        self.node.abort().await;
+
         // Stop the polling task
         info!(target: "minerd::Minerd::stop", "Stopping polling task...");
         self.polling_task.stop().await;
 
-        // Stop the mining node
-        info!(target: "minerd::Minerd::stop", "Stopping miner threads...");
-        self.node.abort_pending().await;
-
         // Close the JSON-RPC client
         info!(target: "minerd::Minerd::stop", "Stopping JSON-RPC client...");
         self.node.stop_rpc_client().await;

+ 183 - 23
bin/minerd/src/rpc.rs

@@ -16,7 +16,11 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
+use std::sync::Arc;
+
 use num_bigint::BigUint;
+use randomx::RandomXVM;
+use smol::channel::{Receiver, Sender};
 use tracing::{debug, error, info};
 use url::Url;
 
@@ -25,7 +29,7 @@ use darkfi::{
     rpc::{client::RpcClient, jsonrpc::JsonRequest, util::JsonValue},
     system::{sleep, ExecutorPtr, StoppableTask},
     util::encoding::base64,
-    validator::pow::mine_block,
+    validator::pow::{generate_mining_vms, mine_block},
     Error, Result,
 };
 use darkfi_serial::deserialize_async;
@@ -55,9 +59,71 @@ impl DarkfidRpcClient {
 }
 
 impl MinerNode {
+    /// Auxiliary function to request configured darkfid daemon for
+    /// its current and next RandomX keys.
+    async fn randomx_keys(&self) -> Result<(HeaderHash, HeaderHash)> {
+        loop {
+            debug!(target: "minerd::rpc::randomx_keys", "Executing RandomX keys request to darkfid...");
+            let params = match self
+                .darkfid_daemon_request("miner.get_current_randomx_keys", &JsonValue::Array(vec![]))
+                .await
+            {
+                Ok(params) => params,
+                Err(e) => {
+                    error!(target: "minerd::rpc::randomx_keys", "darkfid request failed: {e}");
+                    self.sleep().await?;
+                    continue
+                }
+            };
+            debug!(target: "minerd::rpc::randomx_keys", "Got reply: {params:?}");
+
+            // Verify response parameters
+            if !params.is_array() {
+                error!(target: "minerd::rpc::randomx_keys", "darkfid responded with invalid params: {params:?}");
+                self.sleep().await?;
+                continue
+            }
+            let params = params.get::<Vec<JsonValue>>().unwrap();
+            if params.is_empty() {
+                debug!(target: "minerd::rpc::randomx_keys", "darkfid response is empty");
+                self.sleep().await?;
+                continue
+            }
+            if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
+                error!(target: "minerd::rpc::randomx_keys", "darkfid responded with invalid params: {params:?}");
+                self.sleep().await?;
+                continue
+            }
+
+            // Parse parameters
+            let Some(randomx_key_bytes) = base64::decode(params[0].get::<String>().unwrap()) else {
+                error!(target: "minerd::rpc::randomx_keys", "Failed to parse RandomX key bytes");
+                self.sleep().await?;
+                continue
+            };
+            let Ok(randomx_key) = deserialize_async::<HeaderHash>(&randomx_key_bytes).await else {
+                error!(target: "minerd::rpc::randomx_keys", "Failed to parse RandomX key");
+                self.sleep().await?;
+                continue
+            };
+            let Some(next_key_bytes) = base64::decode(params[1].get::<String>().unwrap()) else {
+                error!(target: "minerd::rpc::randomx_keys", "Failed to parse next RandomX key bytes");
+                self.sleep().await?;
+                continue
+            };
+            let Ok(next_key) = deserialize_async::<HeaderHash>(&next_key_bytes).await else {
+                error!(target: "minerd::rpc::randomx_keys", "Failed to parse next RandomX key");
+                self.sleep().await?;
+                continue
+            };
+
+            return Ok((randomx_key, next_key))
+        }
+    }
+
     /// Auxiliary function to poll configured darkfid daemon for a new
     /// mining job.
-    async fn poll(&self, header: &str) -> Result<(HeaderHash, BigUint, Header)> {
+    async fn poll(&self, header: &str) -> Result<(HeaderHash, HeaderHash, BigUint, Header)> {
         loop {
             debug!(target: "minerd::rpc::poll", "Executing poll request to darkfid...");
             let mut request_params = self.config.wallet_config.clone();
@@ -87,10 +153,11 @@ impl MinerNode {
                 self.sleep().await?;
                 continue
             }
-            if params.len() != 3 ||
+            if params.len() != 4 ||
                 !params[0].is_string() ||
                 !params[1].is_string() ||
-                !params[2].is_string()
+                !params[2].is_string() ||
+                !params[3].is_string()
             {
                 error!(target: "minerd::rpc::poll", "darkfid responded with invalid params: {params:?}");
                 self.sleep().await?;
@@ -108,13 +175,23 @@ impl MinerNode {
                 self.sleep().await?;
                 continue
             };
-            let Some(target_bytes) = base64::decode(params[1].get::<String>().unwrap()) else {
+            let Some(next_key_bytes) = base64::decode(params[1].get::<String>().unwrap()) else {
+                error!(target: "minerd::rpc::poll", "Failed to parse next RandomX key bytes");
+                self.sleep().await?;
+                continue
+            };
+            let Ok(next_key) = deserialize_async::<HeaderHash>(&next_key_bytes).await else {
+                error!(target: "minerd::rpc::poll", "Failed to parse next RandomX key");
+                self.sleep().await?;
+                continue
+            };
+            let Some(target_bytes) = base64::decode(params[2].get::<String>().unwrap()) else {
                 error!(target: "minerd::rpc::poll", "Failed to parse target bytes");
                 self.sleep().await?;
                 continue
             };
             let target = BigUint::from_bytes_le(&target_bytes);
-            let Some(header_bytes) = base64::decode(params[2].get::<String>().unwrap()) else {
+            let Some(header_bytes) = base64::decode(params[3].get::<String>().unwrap()) else {
                 error!(target: "minerd::rpc::poll", "Failed to parse header bytes");
                 self.sleep().await?;
                 continue
@@ -125,7 +202,7 @@ impl MinerNode {
                 continue
             };
 
-            return Ok((randomx_key, target, header))
+            return Ok((randomx_key, next_key, target, header))
         }
     }
 
@@ -183,7 +260,7 @@ impl MinerNode {
     /// Auxiliary function to sleep for configured polling rate time.
     async fn sleep(&self) -> Result<()> {
         // Check if stop signal is received
-        if self.stop_signal.is_full() {
+        if self.mining_channel.1.is_full() {
             debug!(target: "minerd::rpc::sleep", "Stop signal received, exiting polling task");
             return Err(Error::DetachedTaskStopped);
         }
@@ -196,14 +273,42 @@ impl MinerNode {
 /// Async task to poll darkfid for new mining jobs. Once a new job is
 /// received, spawns a mining task in the background.
 pub async fn polling_task(miner: MinerNodePtr, ex: ExecutorPtr) -> Result<()> {
-    // Initialize a dummy Header to use on first poll
-    let mut current_job = Header::default().hash().to_string();
+    // Cache current and next RandomX keys and current VMs
+    let (mut current_randomx_key, mut next_randomx_key) = miner.randomx_keys().await?;
+    info!(target: "minerd::rpc::mining_task", "Initializing {} mining VMs for key: {current_randomx_key}", miner.config.threads);
+    let mut current_vms = Arc::new(generate_mining_vms(
+        &current_randomx_key,
+        miner.config.threads,
+        &miner.mining_channel.1.clone(),
+    )?);
+
+    // Initialize the smol channel to send signal between the threads
+    let (vms_sender, vms_receiver) = smol::channel::bounded(1);
+
+    // Detach next RandomX VMs generation in the background if needed
+    if current_randomx_key != next_randomx_key {
+        StoppableTask::new().start(
+            vms_generation_task(next_randomx_key, miner.config.threads, vms_sender.clone(), miner.background_channel.1.clone()),
+            |res| async {
+                match res {
+                    Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                    Err(e) => error!(target: "minerd::rpc::polling_task", "Failed starting next RandomX VMs generation task: {e}"),
+                }
+            },
+            Error::DetachedTaskStopped,
+            ex.clone(),
+        );
+    }
+
+    // Use the dummy Header on first poll
+    let mut current_job = current_randomx_key.to_string();
     loop {
         // Poll darkfid for a mining job
-        let (randomx_key, target, header) = miner.poll(&current_job).await?;
+        let (randomx_key, next_key, target, header) = miner.poll(&current_job).await?;
         let header_hash = header.hash().to_string();
         debug!(target: "minerd::rpc::polling_task", "Received job:");
         debug!(target: "minerd::rpc::polling_task", "\tRandomX key - {randomx_key}");
+        debug!(target: "minerd::rpc::polling_task", "\tNext RandomX key - {next_key}");
         debug!(target: "minerd::rpc::polling_task", "\tTarget - {target}");
         debug!(target: "minerd::rpc::polling_task", "\tHeader - {header_hash}");
 
@@ -223,12 +328,57 @@ pub async fn polling_task(miner: MinerNodePtr, ex: ExecutorPtr) -> Result<()> {
 
         info!(target: "minerd::rpc::polling_task", "Received new job to mine block header {header_hash} with key {randomx_key} for target: 0x{target:064x}");
 
-        // Abord pending job
-        miner.abort_pending().await;
+        // Abord pending mining job
+        miner.abort_mining().await;
+
+        // Check if the current RandomX key has changed
+        if randomx_key != current_randomx_key {
+            current_randomx_key = randomx_key;
+
+            // Check if we should shift to next VMs
+            if current_randomx_key == next_randomx_key {
+                // Shift next generated VMs into current ones
+                info!(target: "minerd::rpc::mining_task", "Grabing next mining VMs from channel for key: {randomx_key}");
+                current_vms = Arc::new(vms_receiver.recv().await?);
+            } else {
+                // Generate the RandomX VMs for the key
+                info!(target: "minerd::rpc::mining_task", "Initializing {} mining VMs for key: {randomx_key}", miner.config.threads);
+                current_vms = Arc::new(generate_mining_vms(
+                    &randomx_key,
+                    miner.config.threads,
+                    &miner.mining_channel.1.clone(),
+                )?);
+            }
+        }
+
+        // Check if the next RandomX key has changed
+        if next_key != next_randomx_key && next_key != randomx_key {
+            // Abord pending VMs generation task
+            miner.abort_background().await;
+
+            // Consume VMs channel item so its empty
+            if let Err(e) = vms_receiver.try_recv() {
+                debug!(target: "minerd::rpc::mining_task", "Failed to cleanup VMs receiver: {e}");
+            }
+
+            // Detach next RandomX VMs generation in the background
+            StoppableTask::new().start(
+                vms_generation_task(next_key, miner.config.threads, vms_sender.clone(), miner.background_channel.1.clone()),
+                |res| async {
+                    match res {
+                        Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                        Err(e) => error!(target: "minerd::rpc::polling_task", "Failed starting next RandomX VMs generation task: {e}"),
+                    }
+                },
+                Error::DetachedTaskStopped,
+                ex.clone(),
+            );
+            next_randomx_key = next_key;
+        }
 
         // Detach mining task
         StoppableTask::new().start(
-            mining_task(miner.clone(), randomx_key, target, header),
+            mining_task(miner.clone(), current_vms.clone(), target, header),
             |res| async {
                 match res {
                     Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
@@ -252,20 +402,14 @@ pub async fn polling_task(miner: MinerNodePtr, ex: ExecutorPtr) -> Result<()> {
 /// Async task to mine provided header and submit solution to darkfid.
 async fn mining_task(
     miner: MinerNodePtr,
-    randomx_key: HeaderHash,
+    vms: Arc<Vec<Arc<RandomXVM>>>,
     target: BigUint,
     mut header: Header,
 ) -> Result<()> {
     // Mine provided block header
     let header_hash = header.hash().to_string();
-    info!(target: "minerd::rpc::mining_task", "Mining block header {header_hash} with key {randomx_key} for target: 0x{target:064x}");
-    if let Err(e) = mine_block(
-        &randomx_key,
-        &target,
-        &mut header,
-        miner.config.threads,
-        &miner.stop_signal.clone(),
-    ) {
+    info!(target: "minerd::rpc::mining_task", "Mining block header {header_hash} for target: 0x{target:064x}");
+    if let Err(e) = mine_block(&vms, &target, &mut header, &miner.mining_channel.1.clone()) {
         error!(target: "minerd::rpc::mining_task", "Failed mining block header {header_hash} with error: {e}");
         return Err(Error::DetachedTaskStopped)
     }
@@ -279,3 +423,19 @@ async fn mining_task(
 
     Ok(())
 }
+
+/// Async task to generate RandomX VMs in the background and push them
+/// in provided channel.
+async fn vms_generation_task(
+    randomx_key: HeaderHash,
+    threads: usize,
+    sender: Sender<Vec<Arc<RandomXVM>>>,
+    stop_signal: Receiver<()>,
+) -> Result<()> {
+    // Generate the RandomX VMs for the key
+    info!(target: "minerd::rpc::vms_generation_task", "Initializing {threads} mining VMs for key: {randomx_key}");
+    let vms = generate_mining_vms(&randomx_key, threads, &stop_signal)?;
+    // Push them into the channel
+    sender.send(vms).await?;
+    Ok(())
+}

+ 22 - 0
src/validator/consensus.rs

@@ -448,6 +448,28 @@ impl Consensus {
         Ok(proposals)
     }
 
+    /// Auxiliary function to grab current RandomX keys.
+    /// If no forks exist, returns canonical keys.
+    pub async fn current_randomx_keys(&self) -> Result<(HeaderHash, HeaderHash)> {
+        // Grab a lock over current forks
+        let forks = self.forks.read().await;
+
+        // If no forks exist, return canonical keys
+        if forks.is_empty() {
+            return Ok(self.module.read().await.darkfi_rx_keys)
+        }
+
+        // Grab best fork keys
+        Ok(forks[best_fork_index(&forks)?].module.darkfi_rx_keys)
+    }
+
+    /// Auxiliary function to grab best current fork full clone.
+    pub async fn best_current_fork(&self) -> Result<Fork> {
+        let forks = self.forks.read().await;
+        let index = best_fork_index(&forks)?;
+        forks[index].full_clone()
+    }
+
     /// Auxiliary function to retrieve current best fork last header.
     /// If no forks exist, grab the last header from canonical.
     pub async fn best_fork_last_header(&self) -> Result<(u32, HeaderHash)> {

+ 11 - 0
src/validator/mod.rs

@@ -825,6 +825,17 @@ impl Validator {
         Ok(())
     }
 
+    /// Auxiliary function to grab current RandomX keys.
+    /// If no forks exist, returns canonical keys.
+    pub async fn current_randomx_keys(&self) -> Result<(HeaderHash, HeaderHash)> {
+        self.consensus.current_randomx_keys().await
+    }
+
+    /// Auxiliary function to grab best current fork full clone.
+    pub async fn best_current_fork(&self) -> Result<Fork> {
+        self.consensus.best_current_fork().await
+    }
+
     /// Auxiliary function to retrieve current best fork next block height.
     pub async fn best_fork_next_block_height(&self) -> Result<u32> {
         let forks = self.consensus.forks.read().await;

+ 95 - 94
src/validator/pow.rs

@@ -17,7 +17,6 @@
  */
 
 use std::{
-    collections::VecDeque,
     sync::{
         atomic::{AtomicBool, AtomicU64, Ordering},
         Arc,
@@ -407,10 +406,14 @@ impl PoWModule {
             &self.darkfi_rx_keys.0
         };
 
+        // Generate the RandomX VMs for the key
+        let vms = generate_mining_vms(randomx_key, threads, stop_signal)?;
+
         // Grab the next mine target
         let target = self.next_mine_target()?;
 
-        mine_block(randomx_key, &target, header, threads, stop_signal)
+        // Mine the block
+        mine_block(&vms, &target, header, stop_signal)
     }
 }
 
@@ -433,146 +436,145 @@ fn get_mining_flags() -> RandomXFlags {
     RandomXFlags::get_recommended_flags() | RandomXFlags::FULLMEM
 }
 
-/// Auxiliary function to mine provided header using a single thread.
-fn single_thread_mine(
+/// Auxiliary function to generate mining VMs for provided RandomX key.
+pub fn generate_mining_vms(
     input: &HeaderHash,
-    target: &BigUint,
-    header: &mut Header,
+    threads: usize,
     stop_signal: &Receiver<()>,
-) -> Result<()> {
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Initializing RandomX cache and dataset...");
+) -> Result<Vec<Arc<RandomXVM>>> {
+    debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initializing RandomX cache and dataset...");
+    debug!(target: "validator::pow::generate_mining_vms", "[MINER] PoW input: {input}");
     let setup_start = Instant::now();
+    let ds_start = Instant::now();
     let flags = get_mining_flags();
     let cache = RandomXCache::new(flags, &input.inner()[..])?;
     let dataset_item_count = RandomXDataset::count()?;
-    let dataset = RandomXDataset::new_init(flags, cache, 0, dataset_item_count)?;
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Setup time: {:?}", setup_start.elapsed());
 
-    // Check if stop signal is received
-    if stop_signal.is_full() {
-        debug!(target: "validator::pow::single_thread_mine", "[MINER] Stop signal received, exiting");
-        return Err(Error::MinerTaskStopped);
+    // Single thread mining VM
+    if threads == 1 {
+        let dataset = RandomXDataset::new_init(flags, cache, 0, dataset_item_count)?;
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initialized RandomX cache and dataset: {:?}", ds_start.elapsed());
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initializing RandomX VM...");
+        let vm_start = Instant::now();
+        let vm = Arc::new(RandomXVM::new(flags, None, Some(dataset))?);
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initialized RandomX VM in {:?}", vm_start.elapsed());
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Setup time: {:?}", setup_start.elapsed());
+        return Ok(vec![vm])
     }
 
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Initializing RandomX VM...");
-    let vm_start = Instant::now();
-    let vm = RandomXVM::new(flags, None, Some(dataset))?;
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Initialized RandomX VM in {:?}", vm_start.elapsed());
+    // Multi thread mining VMs
+    let dataset = RandomXDataset::new(flags, cache, dataset_item_count)?;
+    debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initialized RandomX cache and dataset: {:?}", ds_start.elapsed());
+    let mut vms = Vec::with_capacity(threads);
+    let threads_u32 = threads as u32;
+    for t in 0..threads_u32 {
+        // Check if stop signal is received
+        if stop_signal.is_full() {
+            debug!(target: "validator::pow::generate_mining_vms", "[MINER] Stop signal received, exiting");
+            return Err(Error::MinerTaskStopped);
+        }
+
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initializing RandomX dataset for thread #{t}...");
+        let ds_start = Instant::now();
+        let a = (dataset_item_count * t) / threads_u32;
+        let b = (dataset_item_count * (t + 1)) / threads_u32;
+        let dataset = dataset.subset_init(a, b - a);
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initialized RandomX dataset for thread #{t} in {:?}",
+            ds_start.elapsed()
+        );
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initializing RandomX VM #{t}...");
+        let vm_start = Instant::now();
+        vms.push(Arc::new(RandomXVM::new(flags, None, Some(dataset))?));
+        debug!(target: "validator::pow::generate_mining_vms", "[MINER] Initialized RandomX VM #{t} in {:?}", vm_start.elapsed());
+    }
+    debug!(target: "validator::pow::generate_mining_vms", "[MINER] Setup time: {:?}", setup_start.elapsed());
+    Ok(vms)
+}
 
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Mining started!");
+/// Auxiliary function to mine provided header using provided RandomX
+/// VM (single thread).
+fn randomx_vm_mine(
+    vm: &RandomXVM,
+    target: &BigUint,
+    header: &mut Header,
+    stop_signal: &Receiver<()>,
+) -> Result<()> {
+    debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Mining started!");
     let mining_start = Instant::now();
     loop {
         // Check if stop signal is received
         if stop_signal.is_full() {
-            debug!(target: "validator::pow::single_thread_mine", "[MINER] Stop signal received, exiting");
+            debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Stop signal received, exiting");
             return Err(Error::MinerTaskStopped);
         }
 
         let out_hash = vm.calculate_hash(header.hash().inner())?;
         let out_hash = BigUint::from_bytes_le(&out_hash);
         if &out_hash <= target {
-            debug!(target: "validator::pow::single_thread_mine", "[MINER] Found block header using nonce {}", header.nonce);
-            debug!(target: "validator::pow::single_thread_mine", "[MINER] Block header hash {}", header.hash());
-            debug!(target: "validator::pow::single_thread_mine", "[MINER] RandomX output: 0x{out_hash:064x}");
+            debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Found block header using nonce {}", header.nonce);
+            debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Block header hash {}", header.hash());
+            debug!(target: "validator::pow::randomx_vm_mine", "[MINER] RandomX output: 0x{out_hash:064x}");
             break;
         }
 
         header.nonce += 1;
     }
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Completed mining in {:?}", mining_start.elapsed());
-    debug!(target: "validator::pow::single_thread_mine", "[MINER] Mined header: {header:?}");
+    debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Completed mining in {:?}", mining_start.elapsed());
+    debug!(target: "validator::pow::randomx_vm_mine", "[MINER] Mined header: {header:?}");
     Ok(())
 }
 
-/// Auxiliary function to mine provided header using a multiple threads.
-fn multi_thread_mine(
-    input: &HeaderHash,
+/// Auxiliary function to mine provided header using provided RandomX
+/// VMs corresponding to a multiple threads setup.
+fn randomx_vms_mine(
+    vms: &[Arc<RandomXVM>],
     target: &BigUint,
     header: &mut Header,
-    threads: usize,
     stop_signal: &Receiver<()>,
 ) -> Result<()> {
-    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initializing RandomX cache and dataset...");
-    let setup_start = Instant::now();
-    let flags = get_mining_flags();
-    let cache = RandomXCache::new(flags, &input.inner()[..])?;
-    let dataset_item_count = RandomXDataset::count()?;
-    let dataset = RandomXDataset::new(flags, cache, dataset_item_count)?;
-    let mut subsets = VecDeque::with_capacity(threads);
-
-    // Multithreaded dataset init
-    let threads_u32 = threads as u32;
-    for t in 0..threads_u32 {
-        // Check if stop signal is received
-        if stop_signal.is_full() {
-            debug!(target: "validator::pow::multi_thread_mine", "[MINER] Stop signal received, exiting");
-            return Err(Error::MinerTaskStopped);
-        }
-
-        debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initializing RandomX dataset for thread #{t}...");
-        let ds_start = Instant::now();
-        let a = (dataset_item_count * t) / threads_u32;
-        let b = (dataset_item_count * (t + 1)) / threads_u32;
-        subsets.push_back(dataset.subset_init(a, b - a));
-        debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initialized RandomX dataset for thread #{t} in {:?}",
-            ds_start.elapsed()
-        );
-    }
-    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Setup time: {:?}", setup_start.elapsed());
-
-    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initializing mining threads...");
-    let mut handles = Vec::with_capacity(threads);
+    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Initializing mining threads...");
+    let mut handles = Vec::with_capacity(vms.len());
     let found_header = Arc::new(AtomicBool::new(false));
     let found_nonce = Arc::new(AtomicU64::new(0));
-    let threads_u64 = threads as u64;
+    let threads = vms.len() as u64;
     let mining_start = Instant::now();
-    for t in 0..threads_u64 {
+    for t in 0..threads {
         // Check if stop signal is received
         if stop_signal.is_full() {
-            debug!(target: "validator::pow::multi_thread_mine", "[MINER] Stop signal received, threads creation loop exiting");
+            debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Stop signal received, threads creation loop exiting");
             break
         }
 
         if found_header.load(Ordering::SeqCst) {
-            debug!(target: "validator::pow::multi_thread_mine", "[MINER] Block header found, threads creation loop exiting");
+            debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Block header found, threads creation loop exiting");
             break
         }
 
+        let vm = vms[t as usize].clone();
         let target = target.clone();
         let mut thread_header = header.clone();
         thread_header.nonce = t;
-        let found_header = Arc::clone(&found_header);
-        let found_nonce = Arc::clone(&found_nonce);
-        let dataset = subsets.pop_front().unwrap();
+        let found_header = found_header.clone();
+        let found_nonce = found_nonce.clone();
         let stop_signal = stop_signal.clone();
 
         handles.push(thread::spawn(move || {
-            debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initializing RandomX VM #{t}...");
-            let vm_start = Instant::now();
-            let vm = match RandomXVM::new(flags, None, Some(dataset)) {
-                Ok(vm) => vm,
-                Err(e) => {
-                    error!(target: "validator::pow::multi_thread_mine", "[MINER] Initialized RandomX VM #{t} failed: {e}");
-                    return
-                }
-            };
-            debug!(target: "validator::pow::multi_thread_mine", "[MINER] Initialized RandomX VM #{t} in {:?}", vm_start.elapsed());
             loop {
                 // Check if stop signal was received
                 if stop_signal.is_full() {
-                    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Stop signal received, thread #{t} exiting");
+                    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Stop signal received, thread #{t} exiting");
                     break
                 }
 
                 if found_header.load(Ordering::SeqCst) {
-                    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Block header found, thread #{t} exiting");
+                    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Block header found, thread #{t} exiting");
                     break;
                 }
 
                 let out_hash = match vm.calculate_hash(thread_header.hash().inner()) {
                     Ok(hash) => hash,
                     Err(e) => {
-                        error!(target: "validator::pow::multi_thread_mine", "[MINER] Calculating hash in thread #{t} failed: {e}");
+                        error!(target: "validator::pow::randomx_vms_mine", "[MINER] Calculating hash in thread #{t} failed: {e}");
                         break
                     }
                 };
@@ -580,17 +582,17 @@ fn multi_thread_mine(
                 if out_hash <= target {
                     found_header.store(true, Ordering::SeqCst);
                     found_nonce.store(thread_header.nonce, Ordering::SeqCst);
-                    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Thread #{t} found block header using nonce {}",
+                    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Thread #{t} found block header using nonce {}",
                         thread_header.nonce
                     );
-                    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Block header hash {}", thread_header.hash());
-                    debug!(target: "validator::pow::multi_thread_mine", "[MINER] RandomX output: 0x{out_hash:064x}");
+                    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Block header hash {}", thread_header.hash());
+                    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] RandomX output: 0x{out_hash:064x}");
                     break;
                 }
 
                 // This means thread 0 will use nonces, 0, 4, 8, ...
                 // and thread 1 will use nonces, 1, 5, 9, ...
-                thread_header.nonce += threads_u64;
+                thread_header.nonce += threads;
             }
         }));
     }
@@ -602,26 +604,25 @@ fn multi_thread_mine(
 
     // Check if stop signal is received
     if stop_signal.is_full() {
-        debug!(target: "validator::pow::multi_thread_mine", "[MINER] Stop signal received, exiting");
+        debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Stop signal received, exiting");
         return Err(Error::MinerTaskStopped);
     }
 
-    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Completed mining in {:?}", mining_start.elapsed());
+    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Completed mining in {:?}", mining_start.elapsed());
     header.nonce = found_nonce.load(Ordering::SeqCst);
-    debug!(target: "validator::pow::multi_thread_mine", "[MINER] Mined header: {header:?}");
+    debug!(target: "validator::pow::randomx_vms_mine", "[MINER] Mined header: {header:?}");
     Ok(())
 }
 
-/// Mine provided header, based on provided PoW module next mine target.
+/// Mine provided header, based on provided PoW module next mine target,
+/// using provided RandomX VMs setup.
 pub fn mine_block(
-    input: &HeaderHash,
+    vms: &[Arc<RandomXVM>],
     target: &BigUint,
     header: &mut Header,
-    threads: usize,
     stop_signal: &Receiver<()>,
 ) -> Result<()> {
     debug!(target: "validator::pow::mine_block", "[MINER] Mine target: 0x{target:064x}");
-    debug!(target: "validator::pow::mine_block", "[MINER] PoW input: {input}");
 
     // Check if stop signal is received
     if stop_signal.is_full() {
@@ -629,13 +630,13 @@ pub fn mine_block(
         return Err(Error::MinerTaskStopped);
     }
 
-    match threads {
+    match vms.len() {
         0 => {
-            error!(target: "validator::pow::mine_block", "[MINER] Can't use 0 threads!");
+            error!(target: "validator::pow::mine_block", "[MINER] No VMs were provided!");
             Err(Error::MinerTaskStopped)
         }
-        1 => single_thread_mine(input, target, header, stop_signal),
-        _ => multi_thread_mine(input, target, header, threads, stop_signal),
+        1 => randomx_vm_mine(&vms[0], target, header, stop_signal),
+        _ => randomx_vms_mine(vms, target, header, stop_signal),
     }
 }