rpc.rs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2025 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::{sync::Arc, thread};
  19. use num_bigint::BigUint;
  20. use randomx::RandomXVM;
  21. use smol::channel::{Receiver, Sender};
  22. use tracing::{debug, error, info};
  23. use url::Url;
  24. use darkfi::{
  25. blockchain::{Header, HeaderHash},
  26. rpc::{client::RpcClient, jsonrpc::JsonRequest, util::JsonValue},
  27. system::{sleep, ExecutorPtr, StoppableTask},
  28. util::encoding::base64,
  29. validator::pow::{generate_mining_vms, mine_block},
  30. Error, Result,
  31. };
  32. use darkfi_serial::deserialize_async;
  33. use crate::{MinerNode, MinerNodePtr};
  34. /// Structure to hold a JSON-RPC client and its config,
  35. /// so we can recreate it in case of an error.
  36. pub struct DarkfidRpcClient {
  37. endpoint: Url,
  38. ex: ExecutorPtr,
  39. client: Option<RpcClient>,
  40. }
  41. impl DarkfidRpcClient {
  42. pub async fn new(endpoint: Url, ex: ExecutorPtr) -> Self {
  43. let client = RpcClient::new(endpoint.clone(), ex.clone()).await.ok();
  44. Self { endpoint, ex, client }
  45. }
  46. /// Stop the client.
  47. pub async fn stop(&self) {
  48. if let Some(ref client) = self.client {
  49. client.stop().await
  50. }
  51. }
  52. }
  53. impl MinerNode {
  54. /// Auxiliary function to request configured darkfid daemon for
  55. /// its current and next RandomX keys.
  56. async fn randomx_keys(&self) -> Result<(HeaderHash, HeaderHash)> {
  57. loop {
  58. debug!(target: "minerd::rpc::randomx_keys", "Executing RandomX keys request to darkfid...");
  59. let params = match self
  60. .darkfid_daemon_request("miner.get_current_randomx_keys", &JsonValue::Array(vec![]))
  61. .await
  62. {
  63. Ok(params) => params,
  64. Err(e) => {
  65. error!(target: "minerd::rpc::randomx_keys", "darkfid request failed: {e}");
  66. self.sleep().await?;
  67. continue
  68. }
  69. };
  70. debug!(target: "minerd::rpc::randomx_keys", "Got reply: {params:?}");
  71. // Verify response parameters
  72. if !params.is_array() {
  73. error!(target: "minerd::rpc::randomx_keys", "darkfid responded with invalid params: {params:?}");
  74. self.sleep().await?;
  75. continue
  76. }
  77. let params = params.get::<Vec<JsonValue>>().unwrap();
  78. if params.is_empty() {
  79. debug!(target: "minerd::rpc::randomx_keys", "darkfid response is empty");
  80. self.sleep().await?;
  81. continue
  82. }
  83. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  84. error!(target: "minerd::rpc::randomx_keys", "darkfid responded with invalid params: {params:?}");
  85. self.sleep().await?;
  86. continue
  87. }
  88. // Parse parameters
  89. let Some(randomx_key_bytes) = base64::decode(params[0].get::<String>().unwrap()) else {
  90. error!(target: "minerd::rpc::randomx_keys", "Failed to parse RandomX key bytes");
  91. self.sleep().await?;
  92. continue
  93. };
  94. let Ok(randomx_key) = deserialize_async::<HeaderHash>(&randomx_key_bytes).await else {
  95. error!(target: "minerd::rpc::randomx_keys", "Failed to parse RandomX key");
  96. self.sleep().await?;
  97. continue
  98. };
  99. let Some(next_key_bytes) = base64::decode(params[1].get::<String>().unwrap()) else {
  100. error!(target: "minerd::rpc::randomx_keys", "Failed to parse next RandomX key bytes");
  101. self.sleep().await?;
  102. continue
  103. };
  104. let Ok(next_key) = deserialize_async::<HeaderHash>(&next_key_bytes).await else {
  105. error!(target: "minerd::rpc::randomx_keys", "Failed to parse next RandomX key");
  106. self.sleep().await?;
  107. continue
  108. };
  109. return Ok((randomx_key, next_key))
  110. }
  111. }
  112. /// Auxiliary function to poll configured darkfid daemon for a new
  113. /// mining job.
  114. async fn poll(&self, header: &str) -> Result<(HeaderHash, HeaderHash, BigUint, Header)> {
  115. loop {
  116. debug!(target: "minerd::rpc::poll", "Executing poll request to darkfid...");
  117. let mut request_params = self.config.wallet_config.clone();
  118. request_params.insert(String::from("header"), JsonValue::String(String::from(header)));
  119. let params = match self
  120. .darkfid_daemon_request("miner.get_header", &JsonValue::from(request_params))
  121. .await
  122. {
  123. Ok(params) => params,
  124. Err(e) => {
  125. error!(target: "minerd::rpc::poll", "darkfid poll failed: {e}");
  126. self.sleep().await?;
  127. continue
  128. }
  129. };
  130. debug!(target: "minerd::rpc::poll", "Got reply: {params:?}");
  131. // Verify response parameters
  132. if !params.is_array() {
  133. error!(target: "minerd::rpc::poll", "darkfid responded with invalid params: {params:?}");
  134. self.sleep().await?;
  135. continue
  136. }
  137. let params = params.get::<Vec<JsonValue>>().unwrap();
  138. if params.is_empty() {
  139. debug!(target: "minerd::rpc::poll", "darkfid response is empty");
  140. self.sleep().await?;
  141. continue
  142. }
  143. if params.len() != 4 ||
  144. !params[0].is_string() ||
  145. !params[1].is_string() ||
  146. !params[2].is_string() ||
  147. !params[3].is_string()
  148. {
  149. error!(target: "minerd::rpc::poll", "darkfid responded with invalid params: {params:?}");
  150. self.sleep().await?;
  151. continue
  152. }
  153. // Parse parameters
  154. let Some(randomx_key_bytes) = base64::decode(params[0].get::<String>().unwrap()) else {
  155. error!(target: "minerd::rpc::poll", "Failed to parse RandomX key bytes");
  156. self.sleep().await?;
  157. continue
  158. };
  159. let Ok(randomx_key) = deserialize_async::<HeaderHash>(&randomx_key_bytes).await else {
  160. error!(target: "minerd::rpc::poll", "Failed to parse RandomX key");
  161. self.sleep().await?;
  162. continue
  163. };
  164. let Some(next_key_bytes) = base64::decode(params[1].get::<String>().unwrap()) else {
  165. error!(target: "minerd::rpc::poll", "Failed to parse next RandomX key bytes");
  166. self.sleep().await?;
  167. continue
  168. };
  169. let Ok(next_key) = deserialize_async::<HeaderHash>(&next_key_bytes).await else {
  170. error!(target: "minerd::rpc::poll", "Failed to parse next RandomX key");
  171. self.sleep().await?;
  172. continue
  173. };
  174. let Some(target_bytes) = base64::decode(params[2].get::<String>().unwrap()) else {
  175. error!(target: "minerd::rpc::poll", "Failed to parse target bytes");
  176. self.sleep().await?;
  177. continue
  178. };
  179. let target = BigUint::from_bytes_le(&target_bytes);
  180. let Some(header_bytes) = base64::decode(params[3].get::<String>().unwrap()) else {
  181. error!(target: "minerd::rpc::poll", "Failed to parse header bytes");
  182. self.sleep().await?;
  183. continue
  184. };
  185. let Ok(header) = deserialize_async::<Header>(&header_bytes).await else {
  186. error!(target: "minerd::rpc::poll", "Failed to parse header");
  187. self.sleep().await?;
  188. continue
  189. };
  190. return Ok((randomx_key, next_key, target, header))
  191. }
  192. }
  193. /// Auxiliary function to submit a mining solution to configured
  194. /// darkfid daemon.
  195. async fn submit(&self, nonce: f64) -> String {
  196. debug!(target: "minerd::rpc::submit", "Executing submit request to darkfid...");
  197. let mut request_params = self.config.wallet_config.clone();
  198. request_params.insert(String::from("nonce"), JsonValue::Number(nonce));
  199. let result = match self
  200. .darkfid_daemon_request("miner.submit_solution", &JsonValue::from(request_params))
  201. .await
  202. {
  203. Ok(result) => result,
  204. Err(e) => return format!("darkfid submit failed: {e}"),
  205. };
  206. debug!(target: "minerd::rpc::submit", "Got reply: {result:?}");
  207. // Parse response
  208. match result.get::<String>() {
  209. Some(result) => result.clone(),
  210. None => format!("darkfid responded with invalid params: {result:?}"),
  211. }
  212. }
  213. /// Auxiliary function to execute a request towards the configured
  214. /// darkfid daemon JSON-RPC endpoint.
  215. async fn darkfid_daemon_request(&self, method: &str, params: &JsonValue) -> Result<JsonValue> {
  216. let mut lock = self.rpc_client.write().await;
  217. let req = JsonRequest::new(method, params.clone());
  218. // Check the client is initialized
  219. if let Some(ref client) = lock.client {
  220. // Execute request
  221. if let Ok(rep) = client.request(req.clone()).await {
  222. drop(lock);
  223. return Ok(rep);
  224. }
  225. }
  226. // Reset the rpc client in case of an error and try again
  227. let client = RpcClient::new(lock.endpoint.clone(), lock.ex.clone()).await?;
  228. let rep = client.request(req).await?;
  229. lock.client = Some(client);
  230. drop(lock);
  231. Ok(rep)
  232. }
  233. /// Auxiliary function to stop current JSON-RPC client, if its
  234. /// initialized.
  235. pub async fn stop_rpc_client(&self) {
  236. self.rpc_client.read().await.stop().await;
  237. }
  238. /// Auxiliary function to sleep for configured polling rate time.
  239. async fn sleep(&self) -> Result<()> {
  240. // Check if stop signal is received
  241. if self.mining_channel.1.is_full() {
  242. debug!(target: "minerd::rpc::sleep", "Stop signal received, exiting polling task");
  243. return Err(Error::DetachedTaskStopped);
  244. }
  245. debug!(target: "minerd::rpc::sleep", "Sleeping for {} until next poll...", self.config.polling_rate);
  246. sleep(self.config.polling_rate).await;
  247. Ok(())
  248. }
  249. }
  250. /// Async task to poll darkfid for new mining jobs. Once a new job is
  251. /// received, spawns a mining task in the background.
  252. pub async fn polling_task(miner: MinerNodePtr, ex: ExecutorPtr) -> Result<()> {
  253. // Cache current and next RandomX keys and current VMs
  254. let (mut current_randomx_key, mut next_randomx_key) = miner.randomx_keys().await?;
  255. info!(target: "minerd::rpc::mining_task", "Initializing {} mining VMs for key: {current_randomx_key}", miner.config.threads);
  256. let mut current_vms = Arc::new(generate_mining_vms(
  257. &current_randomx_key,
  258. miner.config.threads,
  259. &miner.mining_channel.1.clone(),
  260. )?);
  261. // Initialize the smol channel to send signal between the threads
  262. let (vms_sender, vms_receiver) = smol::channel::bounded(1);
  263. // Detach next RandomX VMs generation in the background if needed
  264. if current_randomx_key != next_randomx_key {
  265. let threads = miner.config.threads;
  266. let sender = vms_sender.clone();
  267. let stop_singal = miner.background_channel.1.clone();
  268. thread::spawn(move || {
  269. match vms_generation_task(next_randomx_key, threads, sender, stop_singal) {
  270. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  271. Err(e) => {
  272. error!(target: "minerd::rpc::polling_task", "RandomX VMs generation task failed: {e}")
  273. }
  274. }
  275. });
  276. }
  277. // Use the dummy Header on first poll
  278. let mut current_job = current_randomx_key.to_string();
  279. loop {
  280. // Poll darkfid for a mining job
  281. let (randomx_key, next_key, target, header) = miner.poll(&current_job).await?;
  282. let header_hash = header.hash().to_string();
  283. debug!(target: "minerd::rpc::polling_task", "Received job:");
  284. debug!(target: "minerd::rpc::polling_task", "\tRandomX key - {randomx_key}");
  285. debug!(target: "minerd::rpc::polling_task", "\tNext RandomX key - {next_key}");
  286. debug!(target: "minerd::rpc::polling_task", "\tTarget - {target}");
  287. debug!(target: "minerd::rpc::polling_task", "\tHeader - {header_hash}");
  288. // Check if we are already processing this job
  289. if header_hash == current_job {
  290. debug!(target: "minerd::rpc::polling_task", "Already received job, skipping...");
  291. miner.sleep().await?;
  292. continue
  293. }
  294. // Check if we reached the stop height
  295. if miner.config.stop_at_height > 0 && header.height > miner.config.stop_at_height {
  296. info!(target: "minerd::rpc::polling_task", "Reached requested mining height: {}", miner.config.stop_at_height);
  297. info!(target: "minerd::rpc::polling_task", "Daemon can be safely terminated now!");
  298. break
  299. }
  300. info!(target: "minerd::rpc::polling_task", "Received new job to mine block header {header_hash} with key {randomx_key} for target: 0x{target:064x}");
  301. // Abord pending mining job
  302. miner.abort_mining().await;
  303. // Check if the current RandomX key has changed
  304. if randomx_key != current_randomx_key {
  305. current_randomx_key = randomx_key;
  306. // Check if we should shift to next VMs
  307. if current_randomx_key == next_randomx_key {
  308. // Shift next generated VMs into current ones
  309. info!(target: "minerd::rpc::mining_task", "Grabing next mining VMs from channel for key: {randomx_key}");
  310. current_vms = Arc::new(vms_receiver.recv().await?);
  311. } else {
  312. // Generate the RandomX VMs for the key
  313. info!(target: "minerd::rpc::mining_task", "Initializing {} mining VMs for key: {randomx_key}", miner.config.threads);
  314. current_vms = Arc::new(generate_mining_vms(
  315. &randomx_key,
  316. miner.config.threads,
  317. &miner.mining_channel.1.clone(),
  318. )?);
  319. }
  320. }
  321. // Check if the next RandomX key has changed
  322. if next_key != next_randomx_key && next_key != randomx_key {
  323. // Abord pending VMs generation task
  324. miner.abort_background().await;
  325. // Consume VMs channel item so its empty
  326. if let Err(e) = vms_receiver.try_recv() {
  327. debug!(target: "minerd::rpc::mining_task", "Failed to cleanup VMs receiver: {e}");
  328. }
  329. // Detach next RandomX VMs generation in the background
  330. let threads = miner.config.threads;
  331. let sender = vms_sender.clone();
  332. let stop_singal = miner.background_channel.1.clone();
  333. thread::spawn(move || {
  334. match vms_generation_task(next_key, threads, sender, stop_singal) {
  335. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  336. Err(e) => {
  337. error!(target: "minerd::rpc::polling_task", "RandomX VMs generation task failed: {e}")
  338. }
  339. }
  340. });
  341. next_randomx_key = next_key;
  342. }
  343. // Detach mining task
  344. StoppableTask::new().start(
  345. mining_task(miner.clone(), current_vms.clone(), target, header),
  346. |res| async {
  347. match res {
  348. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  349. Err(e) => error!(target: "minerd::rpc::polling_task", "Failed starting mining task: {e}"),
  350. }
  351. },
  352. Error::DetachedTaskStopped,
  353. ex.clone(),
  354. );
  355. // Update current job
  356. current_job = header_hash;
  357. // Sleep until next poll
  358. miner.sleep().await?;
  359. }
  360. Ok(())
  361. }
  362. /// Async task to mine provided header and submit solution to darkfid.
  363. async fn mining_task(
  364. miner: MinerNodePtr,
  365. vms: Arc<Vec<Arc<RandomXVM>>>,
  366. target: BigUint,
  367. mut header: Header,
  368. ) -> Result<()> {
  369. // Mine provided block header
  370. let header_hash = header.hash().to_string();
  371. info!(target: "minerd::rpc::mining_task", "Mining block header {header_hash} for target: 0x{target:064x}");
  372. if let Err(e) = mine_block(&vms, &target, &mut header, &miner.mining_channel.1.clone()) {
  373. error!(target: "minerd::rpc::mining_task", "Failed mining block header {header_hash} with error: {e}");
  374. return Err(Error::DetachedTaskStopped)
  375. }
  376. info!(target: "minerd::rpc::mining_task", "Mined block header {header_hash} with nonce: {}", header.nonce);
  377. info!(target: "minerd::rpc::mining_task", "Mined block header hash: {}", header.hash());
  378. // Submit solution to darkfid
  379. info!(target: "minerd::rpc::submit", "Submitting solution to darkfid...");
  380. let result = miner.submit(header.nonce as f64).await;
  381. info!(target: "minerd::rpc::submit", "Submition result: {result}");
  382. Ok(())
  383. }
  384. /// Async task to generate RandomX VMs in the background and push them
  385. /// in provided channel.
  386. fn vms_generation_task(
  387. randomx_key: HeaderHash,
  388. threads: usize,
  389. sender: Sender<Vec<Arc<RandomXVM>>>,
  390. stop_signal: Receiver<()>,
  391. ) -> Result<()> {
  392. // Generate the RandomX VMs for the key
  393. info!(target: "minerd::rpc::vms_generation_task", "Initializing {threads} mining VMs for key: {randomx_key}");
  394. let vms = generate_mining_vms(&randomx_key, threads, &stop_signal)?;
  395. // Push them into the channel
  396. sender.send_blocking(vms)?;
  397. Ok(())
  398. }