Răsfoiți Sursa

darkfid/registry: try to refresh jobs when a submittion fails

skoupidi 7 luni în urmă
părinte
comite
86f6e8712b

+ 25 - 10
bin/darkfid/src/registry/mod.rs

@@ -342,15 +342,15 @@ impl DarkfiMinersRegistry {
         Ok(())
         Ok(())
     }
     }
 
 
-    /// Refresh outdated jobs in the registry based on provided
-    /// validator state.
-    pub async fn refresh(&self, validator: &ValidatorPtr) -> Result<()> {
-        // Grab registry locks
-        let submit_lock = self.submit_lock.write().await;
-        let mut jobs = self.jobs.write().await;
-        let mut mm_jobs = self.mm_jobs.write().await;
-        let mut block_templates = self.block_templates.write().await;
-
+    /// Refresh outdated jobs in the provided registry maps based on
+    /// provided validator state.
+    pub async fn refresh_jobs(
+        &self,
+        block_templates: &mut HashMap<String, BlockTemplate>,
+        jobs: &mut HashMap<String, MinerClient>,
+        mm_jobs: &mut HashMap<String, String>,
+        validator: &ValidatorPtr,
+    ) -> Result<()> {
         // Find inactive native jobs and drop them
         // Find inactive native jobs and drop them
         let mut dropped_jobs = vec![];
         let mut dropped_jobs = vec![];
         let mut active_templates = HashSet::new();
         let mut active_templates = HashSet::new();
@@ -456,10 +456,25 @@ impl DarkfiMinersRegistry {
             client.publisher.notify(notification).await;
             client.publisher.notify(notification).await;
         }
         }
 
 
+        Ok(())
+    }
+
+    /// Refresh outdated jobs in the registry based on provided
+    /// validator state.
+    pub async fn refresh(&self, validator: &ValidatorPtr) -> Result<()> {
+        // Grab registry locks
+        let submit_lock = self.submit_lock.write().await;
+        let mut block_templates = self.block_templates.write().await;
+        let mut jobs = self.jobs.write().await;
+        let mut mm_jobs = self.mm_jobs.write().await;
+
+        // Refresh jobs
+        self.refresh_jobs(&mut block_templates, &mut jobs, &mut mm_jobs, validator).await?;
+
         // Release registry locks
         // Release registry locks
         drop(block_templates);
         drop(block_templates);
-        drop(mm_jobs);
         drop(jobs);
         drop(jobs);
+        drop(mm_jobs);
         drop(submit_lock);
         drop(submit_lock);
 
 
         Ok(())
         Ok(())

+ 22 - 2
bin/darkfid/src/rpc/rpc_stratum.rs

@@ -221,7 +221,7 @@ impl DarkfiNode {
         };
         };
 
 
         // If we don't know about this client, we can just abort here
         // If we don't know about this client, we can just abort here
-        let jobs = self.registry.jobs.read().await;
+        let mut jobs = self.registry.jobs.write().await;
         let Some(client) = jobs.get(client_id) else {
         let Some(client) = jobs.get(client_id) else {
             return server_error(RpcError::MinerUnknownClient, id, None)
             return server_error(RpcError::MinerUnknownClient, id, None)
         };
         };
@@ -298,9 +298,29 @@ impl DarkfiNode {
             self.registry.submit(&self.validator, &self.subscribers, &self.p2p_handler, block).await
             self.registry.submit(&self.validator, &self.subscribers, &self.p2p_handler, block).await
         {
         {
             error!(
             error!(
-                target: "darkfid::rpc::rpc_xmr::xmr_merge_submit_solution",
+                target: "darkfid::rpc::rpc_stratum::stratum_submit",
                 "[RPC-STRATUM] Error submitting new block: {e}",
                 "[RPC-STRATUM] Error submitting new block: {e}",
             );
             );
+
+            // Try to refresh the jobs before returning error
+            let mut mm_jobs = self.registry.mm_jobs.write().await;
+            if let Err(e) = self
+                .registry
+                .refresh_jobs(&mut block_templates, &mut jobs, &mut mm_jobs, &self.validator)
+                .await
+            {
+                error!(
+                    target: "darkfid::rpc::rpc_stratum::stratum_submit",
+                    "[RPC-STRATUM] Error refreshing registry jobs: {e}",
+                );
+            }
+
+            // Release all locks
+            drop(block_templates);
+            drop(jobs);
+            drop(mm_jobs);
+            drop(submit_lock);
+
             return JsonResponse::new(
             return JsonResponse::new(
                 JsonValue::from(HashMap::from([(
                 JsonValue::from(HashMap::from([(
                     "status".to_string(),
                     "status".to_string(),

+ 24 - 4
bin/darkfid/src/rpc/rpc_xmr.rs

@@ -275,8 +275,8 @@ impl DarkfiNode {
         }
         }
 
 
         // If we don't know about this mm job, we can just abort here
         // If we don't know about this mm job, we can just abort here
-        let jobs = self.registry.mm_jobs.read().await;
-        let Some(wallet) = jobs.get(aux_hash) else {
+        let mut mm_jobs = self.registry.mm_jobs.write().await;
+        let Some(wallet) = mm_jobs.get(aux_hash) else {
             return server_error(RpcError::MinerUnknownJob, id, None)
             return server_error(RpcError::MinerUnknownJob, id, None)
         };
         };
 
 
@@ -399,9 +399,29 @@ impl DarkfiNode {
             self.registry.submit(&self.validator, &self.subscribers, &self.p2p_handler, block).await
             self.registry.submit(&self.validator, &self.subscribers, &self.p2p_handler, block).await
         {
         {
             error!(
             error!(
-                target: "darkfid::rpc::rpc_xmr::xmr_merge_submit_solution",
+                target: "darkfid::rpc::rpc_xmr::xmr_merge_mining_submit_solution",
                 "[RPC-XMR] Error submitting new block: {e}",
                 "[RPC-XMR] Error submitting new block: {e}",
             );
             );
+
+            // Try to refresh the jobs before returning error
+            let mut jobs = self.registry.jobs.write().await;
+            if let Err(e) = self
+                .registry
+                .refresh_jobs(&mut block_templates, &mut jobs, &mut mm_jobs, &self.validator)
+                .await
+            {
+                error!(
+                    target: "darkfid::rpc::rpc_xmr::xmr_merge_mining_submit_solution",
+                    "[RPC-XMR] Error refreshing registry jobs: {e}",
+                );
+            }
+
+            // Release all locks
+            drop(block_templates);
+            drop(jobs);
+            drop(mm_jobs);
+            drop(submit_lock);
+
             return JsonResponse::new(
             return JsonResponse::new(
                 JsonValue::from(HashMap::from([(
                 JsonValue::from(HashMap::from([(
                     "status".to_string(),
                     "status".to_string(),
@@ -417,7 +437,7 @@ impl DarkfiNode {
 
 
         // Release all locks
         // Release all locks
         drop(block_templates);
         drop(block_templates);
-        drop(jobs);
+        drop(mm_jobs);
         drop(submit_lock);
         drop(submit_lock);
 
 
         JsonResponse::new(
         JsonResponse::new(