Pool daemon: refuses to start half-alive (listeners bound first, exit 2; a node that answers no template, exit 3; a STATUS line every 30 s)
The fleet agent's row from the 6 October night: the daemon ran 20 minutes as a process serving nothing while its node answered no template and its member port had been held. Now main binds the member and API listeners before anything else (server::bind_listener; a failure is "POOL NOT STARTED: cannot bind the members listener on <addr>: <os error>", exit 2), probes the node for a template with pow_epoch, retried for --node-wait-secs (default 60) with a line per attempt, else exit 3, and prints a STATUS line every 30 s (members, jobs issued, shares, blocks, template failures, the node's state: ok, NODE NOT ANSWERING, NO TEMPLATE YET). fetch_job counts template failures and the last success. Test server::bind_tests: a held port is refused with the address in the message, the freed port binds. Suite 23 of 23 on igneum-build-1. Protocol and API unchanged. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
d110bd16fb
commit
1797ebadfa
8 changed files with 132 additions and 17 deletions
|
|
@ -236,3 +236,17 @@ a wrong-size cache and mismatched silently); a pool on a class v4 network needs
|
|||
Building the pool crate on igneum-build-1: the pool reads the fork through the `vendor/igneum-node` symlink; lib.sh syncs
|
||||
the fork worktree it points at (`vendor/igneum-node-pr`) as a whole repository, and the symlink itself is created once on
|
||||
the box by hand (`ln -s igneum-node-pr /srv/builds/<worktree>/vendor/igneum-node`), the one step the scripts do not do.
|
||||
|
||||
### 9.1 7 October 2026: the daemon refuses to start half-alive
|
||||
|
||||
The fleet agent's night (6 October, 23:01Z to 23:32Z): the rebased pair never ran as a test. The daemon was started while
|
||||
the node still held its member port, then four members (the wind-down's weight rule freed four pods, not ten) sat
|
||||
connected with no job while every template fetch timed out; nothing in the daemon's own log said so except one line per
|
||||
member. The row from it, now code (pool `269364a`+): the member and API listeners are bound in `main` before anything
|
||||
else and a port that cannot be bound ends the start with exit code 2 and `POOL NOT STARTED: cannot bind the members
|
||||
listener on <addr>: <os error>`; the node must answer a template with its `pow_epoch` before the pool serves anyone,
|
||||
retried for `--node-wait-secs` (default 60, each attempt logged), else exit code 3; a `STATUS` line every 30 s
|
||||
(`members= jobs_issued= shares_accepted= shares_rejected= blocks= template_failures=; node ok | NODE NOT ANSWERING |
|
||||
NO TEMPLATE YET`). Test: `server::bind_tests` (a held port refused with the address in the message, the freed port binds).
|
||||
The 10-member run is still owed: a clean 30-minute window, the daemon bound and its `node ... answers templates` line
|
||||
printed before the first member, on the next pods the weight rule frees (the table must read over 75 percent signed first).
|
||||
|
|
|
|||
|
|
@ -170,3 +170,7 @@ Every Linux build and suite runs on the box through `tools/build-remote.sh` from
|
|||
crate reads the fork through the `vendor/igneum-node` symlink; the build library syncs the fork worktree it points at as a
|
||||
whole repository, and the symlink is made once on the box by hand: `ln -s <fork worktree name> /srv/builds/<worktree>/vendor/igneum-node`.
|
||||
Artefacts land in `pool/target-remote/release/igneum-pool` (x86_64 Linux); the Mac builds no Linux binary.
|
||||
|
||||
Start-up contract (7 October 2026): the daemon exits 2 when `--listen` or `--http` cannot be bound, exits 3 when the node
|
||||
answers no template with `pow_epoch` within `--node-wait-secs` (default 60), and prints `node <url> answers templates: ...`
|
||||
before it accepts a member; a `STATUS` line every 30 s names the member count, jobs issued, shares and the node's state.
|
||||
|
|
|
|||
|
|
@ -199,11 +199,7 @@ pub fn route(pool: &Pool, path: &str) -> Vec<u8> {
|
|||
}
|
||||
}
|
||||
|
||||
pub async fn serve(pool: Arc<Pool>) {
|
||||
let listener = tokio::net::TcpListener::bind(&pool.cfg.http).await.unwrap_or_else(|e| {
|
||||
eprintln!("cannot listen on {}: {e}", pool.cfg.http);
|
||||
std::process::exit(1)
|
||||
});
|
||||
pub async fn serve(pool: Arc<Pool>, listener: tokio::net::TcpListener) {
|
||||
println!("{} pool: stats API and page on http://{}/", crate::state::unix_ms(), pool.cfg.http);
|
||||
loop {
|
||||
let Ok((sock, _)) = listener.accept().await else { continue };
|
||||
|
|
|
|||
|
|
@ -33,12 +33,14 @@ pub struct Config {
|
|||
/// Shares sampled at shifts above this are still verified in v0 (spec sample_shift 8 is an allowance, not a duty)
|
||||
pub verify_threads: usize,
|
||||
pub snapshot_interval_s: u64,
|
||||
/// How long the start-up waits for the node to answer a template before the daemon refuses to start (exit 3)
|
||||
pub node_wait_secs: u64,
|
||||
}
|
||||
|
||||
impl Config {
|
||||
pub fn usage() -> ! {
|
||||
eprintln!(
|
||||
"usage: igneum-pool [--node grpc://127.0.0.1:26610]... [--evm-rpc http://127.0.0.1:26790] [--listen 0.0.0.0:4463]\n\
|
||||
"usage: igneum-pool [--node grpc://127.0.0.1:26610]... [--evm-rpc http://127.0.0.1:26790] [--listen 0.0.0.0:4463] [--node-wait-secs 60]\n\
|
||||
\x20 [--http 127.0.0.1:4480] [--data-dir ./data] [--payout-key <file>] [--fee-percent 1] [--min-payout 1.0]\n\
|
||||
\x20 [--pplns-window 2.0] [--payout-interval-s 60] [--dry-run] [--network devnet|testnet|mainnet]\n\
|
||||
\x20 [--share-interval-s 10] [--min-shift 0] [--max-shift 60] [--stale-grace-ms 2000] [--orphan-after-daa 120]\n\
|
||||
|
|
@ -70,6 +72,7 @@ impl Config {
|
|||
public_url: String::new(),
|
||||
verify_threads: 2,
|
||||
snapshot_interval_s: 15,
|
||||
node_wait_secs: 60,
|
||||
};
|
||||
let mut i = 0;
|
||||
while i < args.len() {
|
||||
|
|
@ -100,6 +103,7 @@ impl Config {
|
|||
"--public-url" => c.public_url = val(),
|
||||
"--verify-threads" => c.verify_threads = val().parse::<usize>().unwrap_or_else(|_| Config::usage()).max(1),
|
||||
"--snapshot-interval-s" => c.snapshot_interval_s = val().parse().unwrap_or_else(|_| Config::usage()),
|
||||
"--node-wait-secs" => c.node_wait_secs = val().parse().unwrap_or_else(|_| Config::usage()),
|
||||
"-h" | "--help" => Config::usage(),
|
||||
_ => Config::usage(),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ use config::Config;
|
|||
use kaspa_addresses::{Address, Prefix, Version};
|
||||
use kaspa_grpc_client::GrpcClient;
|
||||
use kaspa_pow::igneum::IgneumEngine;
|
||||
use kaspa_rpc_core::api::rpc::RpcApi;
|
||||
use pool::{NetInfo, Pool};
|
||||
use state::{State, WEI_PER_IGN};
|
||||
use std::collections::HashMap;
|
||||
|
|
@ -38,6 +39,17 @@ use std::time::{Duration, Instant};
|
|||
async fn main() {
|
||||
let args: Vec<String> = std::env::args().skip(1).collect();
|
||||
let cfg = Config::from_args(&args);
|
||||
// The two listeners first (7 October 2026, the fleet agent's row from the 6 October night: a daemon whose member
|
||||
// port another process still held ran 20 minutes as a process with no socket): a port that cannot be bound ends
|
||||
// the start here, exit code 2, before a node is contacted or a member is let down.
|
||||
let member_listener = server::bind_listener("members", &cfg.listen).await.unwrap_or_else(|e| {
|
||||
eprintln!("POOL NOT STARTED: {e}");
|
||||
std::process::exit(2)
|
||||
});
|
||||
let api_listener = server::bind_listener("stats API", &cfg.http).await.unwrap_or_else(|e| {
|
||||
eprintln!("POOL NOT STARTED: {e}");
|
||||
std::process::exit(2)
|
||||
});
|
||||
std::fs::create_dir_all(&cfg.data_dir).unwrap_or_else(|e| {
|
||||
eprintln!("data dir {}: {e}", cfg.data_dir.display());
|
||||
std::process::exit(1)
|
||||
|
|
@ -68,6 +80,33 @@ async fn main() {
|
|||
}
|
||||
}
|
||||
let node = clients.remove(0);
|
||||
// The node must answer a template with its pow_epoch before the pool serves anyone: a pool whose node answers no
|
||||
// template issues no job (the 6 October night: four members connected, "template timed out" per member, no job
|
||||
// for 20 minutes). Retried for --node-wait-secs, saying so; then exit code 3.
|
||||
{
|
||||
let probe_address = Address::new(Prefix::Devnet, Version::PubKey, &[0u8; 32]);
|
||||
let deadline = Instant::now() + Duration::from_secs(cfg.node_wait_secs);
|
||||
let mut attempt = 0u32;
|
||||
loop {
|
||||
attempt += 1;
|
||||
let why = match tokio::time::timeout(Duration::from_secs(10), node.get_block_template(probe_address.clone(), Vec::new())).await {
|
||||
Ok(Ok(t)) if t.pow_epoch.is_some() => {
|
||||
let i = t.pow_epoch.unwrap();
|
||||
println!("{} node {} answers templates: epoch {} class {} era {} daa {}", state::unix_ms(), cfg.nodes[0], i.epoch_index, i.class().name(), i.era_seed.map(|h| h.to_string()).unwrap_or_else(|| "-".into()), i.virtual_daa_score);
|
||||
break;
|
||||
}
|
||||
Ok(Ok(_)) => "the template carries no pow_epoch (a node before the 0.3.x line); the pool cannot name a program".to_string(),
|
||||
Ok(Err(e)) => e.to_string(),
|
||||
Err(_) => "template timed out after 10 s".to_string(),
|
||||
};
|
||||
if Instant::now() >= deadline {
|
||||
eprintln!("POOL NOT STARTED: node {} answered no usable template in {} s ({attempt} attempts; last: {why})", cfg.nodes[0], cfg.node_wait_secs);
|
||||
std::process::exit(3)
|
||||
}
|
||||
eprintln!("{} waiting for node {} to answer a template (attempt {attempt}: {why}); giving up at {} s", state::unix_ms(), cfg.nodes[0], cfg.node_wait_secs);
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
}
|
||||
}
|
||||
// the UTXO-side coinbase address: derived from the pool's EVM address; the execution layer pays the EVM side
|
||||
// (igneum/exec/src/executor.rs), the UTXO output is not spent by v0 (docs/plans/pool.md)
|
||||
let prefix = match cfg.network.as_str() {
|
||||
|
|
@ -99,6 +138,8 @@ async fn main() {
|
|||
verify_permits: tokio::sync::Semaphore::new(cfg.verify_threads),
|
||||
started: Instant::now(),
|
||||
cfg: cfg.clone(),
|
||||
template_failures: AtomicU64::new(0),
|
||||
last_template_ok_ms: AtomicU64::new(0),
|
||||
});
|
||||
let chain = evm.call("eth_chainId", serde_json::json!([])).await.ok().and_then(|v| v.as_str().map(|s| s.to_string()));
|
||||
println!(
|
||||
|
|
@ -119,8 +160,9 @@ async fn main() {
|
|||
tokio::spawn(node::confirm_loop(pool.clone()));
|
||||
tokio::spawn(node::net_loop(pool.clone()));
|
||||
tokio::spawn(node::vardiff_loop(pool.clone()));
|
||||
tokio::spawn(server::listen(pool.clone()));
|
||||
tokio::spawn(api::serve(pool.clone()));
|
||||
tokio::spawn(node::status_loop(pool.clone()));
|
||||
tokio::spawn(server::listen(pool.clone(), member_listener));
|
||||
tokio::spawn(api::serve(pool.clone(), api_listener));
|
||||
{
|
||||
let (pool, payer) = (pool.clone(), payer.clone());
|
||||
tokio::spawn(async move {
|
||||
|
|
|
|||
|
|
@ -39,10 +39,18 @@ pub fn extra_data(member: &Member, pool_address: &[u8; 20]) -> Vec<u8> {
|
|||
/// cache build), so it runs under `spawn_blocking`.
|
||||
pub async fn fetch_job(pool: &Arc<Pool>, member: &Arc<Member>) -> Result<(), String> {
|
||||
let t0 = Instant::now();
|
||||
let tmpl = tokio::time::timeout(Duration::from_secs(5), pool.node.get_block_template(pool.pay_address.clone(), extra_data(member, &pool.pool_address)))
|
||||
.await
|
||||
.map_err(|_| "template timed out".to_string())?
|
||||
.map_err(|e| e.to_string())?;
|
||||
let tmpl = match tokio::time::timeout(Duration::from_secs(5), pool.node.get_block_template(pool.pay_address.clone(), extra_data(member, &pool.pool_address))).await {
|
||||
Ok(Ok(t)) => t,
|
||||
Ok(Err(e)) => {
|
||||
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
||||
return Err(e.to_string());
|
||||
}
|
||||
Err(_) => {
|
||||
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
||||
return Err("template timed out".to_string());
|
||||
}
|
||||
};
|
||||
pool.last_template_ok_ms.store(unix_ms(), Ordering::Relaxed);
|
||||
let mut raw = tmpl.block;
|
||||
raw.header.vote_key_hash = member.key_hash;
|
||||
let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?;
|
||||
|
|
@ -344,6 +352,29 @@ pub async fn confirm_loop(pool: Arc<Pool>) {
|
|||
}
|
||||
}
|
||||
|
||||
/// The daemon's own status line every 30 s (7 October 2026: a night where four members sat connected with no job for
|
||||
/// 20 minutes was visible only as per-member "template timed out" lines).
|
||||
pub async fn status_loop(pool: Arc<Pool>) {
|
||||
loop {
|
||||
tokio::time::sleep(Duration::from_secs(30)).await;
|
||||
let members = pool.members_snapshot().len();
|
||||
let failures = pool.template_failures.load(Ordering::Relaxed);
|
||||
let last_ok = pool.last_template_ok_ms.load(Ordering::Relaxed);
|
||||
let jobs = pool.next_job_id.load(Ordering::Relaxed);
|
||||
let (accepted, rejected, blocks) = {
|
||||
let s = pool.state.lock().unwrap();
|
||||
(s.accepted, s.rejected, s.blocks_found)
|
||||
};
|
||||
let node_state = if last_ok == 0 {
|
||||
"NO TEMPLATE YET from the node: no job has been issued".to_string()
|
||||
} else {
|
||||
let ago = unix_ms().saturating_sub(last_ok) / 1000;
|
||||
if ago > 30 { format!("NODE NOT ANSWERING: last template {ago} s ago") } else { format!("node ok, last template {ago} s ago") }
|
||||
};
|
||||
println!("{} STATUS members={members} jobs_issued={jobs} shares_accepted={accepted} shares_rejected={rejected} blocks={blocks} template_failures={failures}; {node_state}", now());
|
||||
}
|
||||
}
|
||||
|
||||
/// Network numbers every 5 s.
|
||||
pub async fn net_loop(pool: Arc<Pool>) {
|
||||
loop {
|
||||
|
|
|
|||
|
|
@ -109,6 +109,10 @@ pub struct Pool {
|
|||
pub want_templates: tokio::sync::Notify,
|
||||
pub verify_permits: tokio::sync::Semaphore,
|
||||
pub started: Instant,
|
||||
/// Template fetches that failed or timed out since start, and the unix ms of the last one that succeeded (0 =
|
||||
/// never): a daemon whose node answers no template issues no job, and the STATUS line in the log says so
|
||||
pub template_failures: AtomicU64,
|
||||
pub last_template_ok_ms: AtomicU64,
|
||||
}
|
||||
|
||||
impl Pool {
|
||||
|
|
|
|||
|
|
@ -19,11 +19,14 @@ fn now() -> String {
|
|||
format!("{}.{:03}", t.as_secs(), t.subsec_millis())
|
||||
}
|
||||
|
||||
pub async fn listen(pool: Arc<Pool>) {
|
||||
let listener = tokio::net::TcpListener::bind(&pool.cfg.listen).await.unwrap_or_else(|e| {
|
||||
eprintln!("cannot listen on {}: {e}", pool.cfg.listen);
|
||||
std::process::exit(1)
|
||||
});
|
||||
/// Binds a listener or says exactly why not. Called from `main` BEFORE anything else starts (7 October 2026, the fleet
|
||||
/// agent's row from the 6 October night: a daemon whose member port another process still held ran for 20 minutes as a
|
||||
/// process with no socket; now a bind that fails ends the process with exit code 2 before a member can be let down).
|
||||
pub async fn bind_listener(what: &str, addr: &str) -> Result<tokio::net::TcpListener, String> {
|
||||
tokio::net::TcpListener::bind(addr).await.map_err(|e| format!("cannot bind the {what} listener on {addr}: {e} (another process holds the port, or the address is not this machine's)"))
|
||||
}
|
||||
|
||||
pub async fn listen(pool: Arc<Pool>, listener: tokio::net::TcpListener) {
|
||||
println!("{} pool: members on {} (chain id {}, {})", now(), pool.cfg.listen, pool.cfg.chain_id(), pool.cfg.network);
|
||||
loop {
|
||||
match listener.accept().await {
|
||||
|
|
@ -323,3 +326,20 @@ async fn handle_share(pool: &Arc<Pool>, m: &Arc<Member>, job_id: u64, nonce: u64
|
|||
tokio::spawn(submit_block(pool.clone(), m.clone(), raw, nonce, shift, daa, blue));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod bind_tests {
|
||||
use super::bind_listener;
|
||||
|
||||
/// Known good: a free port binds. Known failed: the same port, held, is refused with the address in the message
|
||||
/// (the shape the fleet saw on 6 October: a node still holding 4463).
|
||||
#[tokio::test]
|
||||
async fn a_held_port_is_refused_loudly_and_a_free_one_binds() {
|
||||
let held = bind_listener("members", "127.0.0.1:0").await.expect("a free port binds");
|
||||
let addr = held.local_addr().unwrap().to_string();
|
||||
let err = bind_listener("members", &addr).await.err().expect("the held port is refused");
|
||||
assert!(err.contains(&addr) && err.contains("members"), "{err}");
|
||||
drop(held);
|
||||
bind_listener("members", &addr).await.expect("released, it binds again");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue