Pool, the devnet-4 pair's class (8 October 2026): no job on an unsynced node's state. state_provider.rs reads igneum_getExecStatus beside the class v5 leaves and refuses a stream while the node reads synced false, blocked or re-executing (ExecStatus; a missing field is not synced, never a job on a guess), so the engine builds no epoch and the template feed issues no job on catch-up leaves (the pair: identical seeds, era and the freeze's generator on node, pool and member, every share wrong_hash because the node in class v5 catch-up answered igneum_getPowStateLeaves from its still-settling state); the STATUS line says 'node not synced (class v5 catch-up): no jobs' with the exec RPC named while it holds. Known-failed first: an_unsynced_node_hands_the_pool_no_state_and_no_job read red against the c24d5080 provider on build-6 (it served the leaves), green with the guard; the mock RPC answers igneum_getExecStatus and igneum_getPowStateLeaves. Gate from this tree against node 7cfa422a: 52 of 52 on build-6
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
(cherry picked from commit b5531f46b5)
This commit is contained in:
parent
2323d091d0
commit
7ffd29af4d
3 changed files with 138 additions and 2 deletions
|
|
@ -187,7 +187,11 @@ async fn main() {
|
|||
};
|
||||
let state_path = cfg.data_dir.join("state.json");
|
||||
// class v5: the engine's day state comes from the node's exec RPC (the pool's own provider), installed before the first epoch
|
||||
kaspa_pow::igneum::install_day_state_provider(Arc::new(crate::state_provider::RpcStateProvider::new(&cfg.exec_rpc)));
|
||||
// class v5 day state from the node's exec RPC; the provider refuses a node whose execution is not settled
|
||||
// (igneum_getExecStatus synced false, blocked or re-executing), so no job is built on catch-up leaves (the
|
||||
// devnet-4 pair, 8 October 2026: every member share wrong_hash on an unsynced node); the STATUS line says so
|
||||
let day_state = Arc::new(crate::state_provider::RpcStateProvider::new(&cfg.exec_rpc));
|
||||
kaspa_pow::igneum::install_day_state_provider(day_state.clone());
|
||||
eprintln!("class v5 day state: {} ({})", cfg.exec_rpc, "igneum_getPowStateLeaves");
|
||||
let engine = Arc::new(IgneumEngine::new());
|
||||
let walker_for_open = walker.clone();
|
||||
|
|
@ -298,9 +302,13 @@ async fn main() {
|
|||
// a status line every 30 s
|
||||
{
|
||||
let pool = pool.clone();
|
||||
let day_state = day_state.clone();
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
tokio::time::sleep(Duration::from_secs(30)).await;
|
||||
if let Some(why) = day_state.last_refusal() {
|
||||
println!("{} STATUS {why} (the node's exec RPC {}; members are held without work until it reads synced)", state::unix_ms(), pool.cfg.exec_rpc);
|
||||
}
|
||||
let (miners, workers) = pool.online();
|
||||
let s = pool.state.lock().unwrap();
|
||||
let (n, mean, p50, p99, max) = s.check_cost();
|
||||
|
|
|
|||
|
|
@ -51,6 +51,11 @@ pub struct Chain {
|
|||
/// Counters the tests read
|
||||
pub send_calls: u64,
|
||||
pub calls: Vec<String>,
|
||||
/// The node's execution state (`igneum_getExecStatus`) and the class v5 leaves it serves by block hash
|
||||
pub exec_synced: bool,
|
||||
pub exec_blocked: Option<String>,
|
||||
pub exec_reexecuting: Option<String>,
|
||||
pub state_leaves: std::collections::HashMap<String, Vec<u8>>,
|
||||
}
|
||||
|
||||
impl Chain {
|
||||
|
|
@ -70,6 +75,10 @@ impl Chain {
|
|||
paused: false,
|
||||
send_calls: 0,
|
||||
calls: Vec::new(),
|
||||
exec_synced: true,
|
||||
exec_blocked: None,
|
||||
exec_reexecuting: None,
|
||||
state_leaves: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -205,6 +214,14 @@ impl Chain {
|
|||
Ok(json!({"finalityActive": self.finality_active, "latestLockKind": if self.finality_active { json!(self.lock_kind) } else { Value::Null }, "latestLockedIndex": "0x10", "pausedSinceMs": if self.paused { json!("0x1") } else { Value::Null }, "checkpoints": []})),
|
||||
false,
|
||||
),
|
||||
"igneum_getExecStatus" => (Ok(json!({"synced": self.exec_synced, "blocked": self.exec_blocked, "reexecuting": self.exec_reexecuting, "executedTip": q(self.height)})), false),
|
||||
"igneum_getPowStateLeaves" => {
|
||||
let h = params.get(0).and_then(|v| v.as_str()).unwrap_or("");
|
||||
match self.state_leaves.get(h) {
|
||||
Some(b) => (Ok(json!({"block": h, "streamHex": format!("0x{}", hex::encode(b))})), false),
|
||||
None => (Err(format!("no state stream for block {h}")), false),
|
||||
}
|
||||
}
|
||||
other => (Err(format!("the mock has no {other}")), false),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,11 +11,61 @@ use std::time::Duration;
|
|||
|
||||
pub struct RpcStateProvider {
|
||||
url: String,
|
||||
/// The last refusal of the node's state (`node not synced ...`), read by the STATUS line; cleared on a stream served
|
||||
last_refusal: std::sync::Mutex<Option<String>>,
|
||||
}
|
||||
|
||||
/// The node's own word on its execution state (`igneum_getExecStatus`): leaves are served only from a node that reads
|
||||
/// `synced` true with no `blocked` and no `reexecuting`. The devnet-4 pair of 8 October 2026: a node in class v5
|
||||
/// catch-up answered `igneum_getPowStateLeaves` from its still-settling state, the pool built jobs on it, and every
|
||||
/// member share read wrong_hash; the generator, the seeds and the era were identical on all three binaries.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub struct ExecStatus {
|
||||
pub synced: bool,
|
||||
pub blocked: Option<String>,
|
||||
pub reexecuting: Option<String>,
|
||||
}
|
||||
|
||||
impl ExecStatus {
|
||||
/// Parses the node's answer; a field the node does not send reads as not synced (never a job on a guess).
|
||||
pub fn parse(body: &str) -> Result<ExecStatus, String> {
|
||||
let v: serde_json::Value = serde_json::from_str(body).map_err(|e| format!("igneum_getExecStatus: bad JSON: {e}"))?;
|
||||
if let Some(err) = v.get("error") {
|
||||
return Err(format!("igneum_getExecStatus: {}", err.get("message").and_then(|m| m.as_str()).unwrap_or("error")));
|
||||
}
|
||||
let r = v.get("result").ok_or("igneum_getExecStatus: no result")?;
|
||||
let text = |k: &str| r.get(k).and_then(|x| if x.is_null() { None } else { Some(x.as_str().map(|s| s.to_string()).unwrap_or_else(|| x.to_string())) });
|
||||
Ok(ExecStatus { synced: r.get("synced").and_then(|b| b.as_bool()).unwrap_or(false), blocked: text("blocked"), reexecuting: text("reexecuting") })
|
||||
}
|
||||
/// Why the node's state is not settled, or none when it is.
|
||||
pub fn refusal(&self) -> Option<String> {
|
||||
if let Some(b) = &self.blocked {
|
||||
return Some(format!("node not synced (execution blocked: {b}): no jobs"));
|
||||
}
|
||||
if let Some(r) = &self.reexecuting {
|
||||
return Some(format!("node not synced (re-executing after a reorg: {r}): no jobs"));
|
||||
}
|
||||
if !self.synced {
|
||||
return Some("node not synced (class v5 catch-up): no jobs".into());
|
||||
}
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
impl RpcStateProvider {
|
||||
pub fn new(url: &str) -> Self {
|
||||
Self { url: url.to_string() }
|
||||
Self { url: url.to_string(), last_refusal: std::sync::Mutex::new(None) }
|
||||
}
|
||||
|
||||
/// The last reason the node's state was refused, for the STATUS line; none while streams are served.
|
||||
pub fn last_refusal(&self) -> Option<String> {
|
||||
self.last_refusal.lock().unwrap().clone()
|
||||
}
|
||||
|
||||
/// The node's execution status, by the same socket the leaves come over.
|
||||
pub fn exec_status(&self) -> Result<ExecStatus, String> {
|
||||
let resp = self.post("{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"igneum_getExecStatus\",\"params\":[]}")?;
|
||||
ExecStatus::parse(&resp)
|
||||
}
|
||||
|
||||
fn post(&self, body: &str) -> Result<String, String> {
|
||||
|
|
@ -35,6 +85,13 @@ impl RpcStateProvider {
|
|||
|
||||
impl kaspa_pow::igneum::DayStateProvider for RpcStateProvider {
|
||||
fn state_stream(&self, block: Hash) -> Result<Vec<u8>, String> {
|
||||
// the node's state is settled, or no stream and no job on it
|
||||
let status = self.exec_status()?;
|
||||
if let Some(why) = status.refusal() {
|
||||
*self.last_refusal.lock().unwrap() = Some(why.clone());
|
||||
return Err(format!("{why} (igneum_getExecStatus on {}: synced {}, blocked {:?}, reexecuting {:?})", self.url, status.synced, status.blocked, status.reexecuting));
|
||||
}
|
||||
*self.last_refusal.lock().unwrap() = None;
|
||||
let body = format!("{{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"igneum_getPowStateLeaves\",\"params\":[\"{block}\"]}}");
|
||||
let resp = self.post(&body)?;
|
||||
const KEY: &str = "\"streamHex\":\"0x";
|
||||
|
|
@ -50,3 +107,57 @@ impl kaspa_pow::igneum::DayStateProvider for RpcStateProvider {
|
|||
"the node's exec RPC (the pool's provider)"
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::mock_rpc;
|
||||
use kaspa_pow::igneum::DayStateProvider;
|
||||
|
||||
/// The devnet-4 pair's class, 8 October 2026: a node still in class v5 catch-up (`synced` false) must hand the
|
||||
/// pool no state stream, so no job is built on unsettled leaves; the refusal names the cause for the STATUS line
|
||||
/// and `igneum_getPowStateLeaves` is never asked. Known-failed first against the provider of 8ff7a0f4, which
|
||||
/// asked for the leaves of any node that answered. Known-pass: the node synced, the stream served, the refusal
|
||||
/// cleared; blocked and re-executing refuse the same way.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn an_unsynced_node_hands_the_pool_no_state_and_no_job() {
|
||||
let (chain, url) = mock_rpc::start().await;
|
||||
let block = Hash::from_bytes([0x5au8; 32]);
|
||||
{
|
||||
let mut c = chain.lock().unwrap();
|
||||
c.exec_synced = false;
|
||||
c.state_leaves.insert(format!("{block}"), vec![1, 2, 3, 4]);
|
||||
}
|
||||
let p = RpcStateProvider::new(&url);
|
||||
let p2 = std::sync::Arc::new(p);
|
||||
let (pp, b) = (p2.clone(), block);
|
||||
let r = tokio::task::spawn_blocking(move || pp.state_stream(b)).await.unwrap();
|
||||
let e = r.err().expect("refused while the node is not synced");
|
||||
assert!(e.contains("node not synced (class v5 catch-up): no jobs"), "{e}");
|
||||
assert_eq!(p2.last_refusal().as_deref(), Some("node not synced (class v5 catch-up): no jobs"));
|
||||
assert!(!chain.lock().unwrap().calls.iter().any(|m| m == "igneum_getPowStateLeaves"), "the leaves were never asked for: {:?}", chain.lock().unwrap().calls);
|
||||
// blocked, and re-executing, refuse with their cause
|
||||
chain.lock().unwrap().exec_synced = true;
|
||||
chain.lock().unwrap().exec_blocked = Some("snapshot below the tip".into());
|
||||
let (pp, b) = (p2.clone(), block);
|
||||
let e = tokio::task::spawn_blocking(move || pp.state_stream(b)).await.unwrap().err().unwrap();
|
||||
assert!(e.contains("execution blocked: snapshot below the tip"), "{e}");
|
||||
chain.lock().unwrap().exec_blocked = None;
|
||||
chain.lock().unwrap().exec_reexecuting = Some("0x1234".into());
|
||||
let (pp, b) = (p2.clone(), block);
|
||||
let e = tokio::task::spawn_blocking(move || pp.state_stream(b)).await.unwrap().err().unwrap();
|
||||
assert!(e.contains("re-executing after a reorg"), "{e}");
|
||||
// synced: the stream is served and the refusal is cleared
|
||||
chain.lock().unwrap().exec_reexecuting = None;
|
||||
let (pp, b) = (p2.clone(), block);
|
||||
let bytes = tokio::task::spawn_blocking(move || pp.state_stream(b)).await.unwrap().expect("served once synced");
|
||||
assert_eq!(bytes, vec![1, 2, 3, 4]);
|
||||
assert_eq!(p2.last_refusal(), None);
|
||||
assert!(chain.lock().unwrap().calls.iter().any(|m| m == "igneum_getPowStateLeaves"));
|
||||
// a node that answers no status field at all is not synced (never a job on a guess)
|
||||
let s = ExecStatus::parse(r#"{"jsonrpc":"2.0","id":1,"result":{}}"#).unwrap();
|
||||
assert_eq!(s, ExecStatus { synced: false, blocked: None, reexecuting: None });
|
||||
assert!(s.refusal().is_some());
|
||||
assert!(ExecStatus::parse(r#"{"jsonrpc":"2.0","id":1,"error":{"message":"no exec"}}"#).is_err());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue