igneum/pool/src/server.rs
igneum-labs c83ea0f5d4 Pool finished (mission item 11): pool-0's fee published beside the dev fee, TLS 1.3 on the member port with the authorize binding (spec 9.3, O-9.7 closed), the page rows Q68 to Q73, and the open pool: a share sidechain with no operator (spec 9.12)
Pool-0: welcome carries software_dev_fee_percent beside the pool fee, the page's Fees row and site/miner.html say pool-0 charges the same 1 percent as the solo dev fee; jobs and seeds carry the latency ladder's rung for 0.3.19.

TLS (pool/src/tls.rs): --tls-cert/--tls-key or --tls-self-signed (a P-256 pair under the data dir, the certificate pin printed at start); the binding is the member's BLS signature over this connection's exporter (label EXPORTER-igneum-pool-binding, context the chain id), refused on any other connection; testnet and mainnet refuse the clear without --allow-plain. HiveOS: pools:// and POOL_PIN.

Page rows: node state in words (Q68), one formatter set and luck in MiningPoolStats' convention (Q69), samples and check costs persisted (Q70), payments under the lookup (Q71), hourly history with a sparkline, --alert-webhook, /metrics, /health on the node (Q72), the network's finality state beside "payouts follow blue confirmation, not finality" (Q73).

The open pool (pool/src/sidechain.rs, open.rs, p2p.rs; igneum-pool --open): the member's own node and daemon, no payout key, no balance; every coinbase carries the member's own address, the share chain's parent (IGNS) and the window's split (IGNP); templates re-stamped on a new chain tip without a node call; shares gossiped and checked by every member (structure, the node's seeds, the PoW); heaviest work wins, a fork's loser is stale; the dev fee as one split entry; shares.log and verify-share for the drop proof; the gate harness pool/tools/open-gate.mjs. The node side (the executor's split from pool_split_activation_daa) is on the fork branch pool-finish-node.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
(cherry picked from commit 57954a43eb)
2026-10-07 14:03:05 +00:00

413 lines
20 KiB
Rust

//! Member connections: newline JSON over TLS 1.3 (spec 9.3, `tls.rs`: `--tls-cert`/`--tls-key` or `--tls-self-signed`)
//! or over plain TCP (the devnet default; testnet and mainnet refuse plain without `--allow-plain`). One reader task
//! and one writer task per connection. Over TLS the `authorize` binding (O-9.7) is checked against this connection's
//! exporter, so an `authorize` replayed on another connection is refused.
use crate::node::{fetch_job, submit_block};
use crate::pool::{Member, MemberInner, Pool};
use crate::protocol::{hex_u64, parse_hex_array, parse_hex_u64, Msg, ShareScheme, MAX_LINE, VERSION};
use crate::vardiff::{share_target, share_weight, Vardiff};
use crate::verify::{check, Code};
use kaspa_consensus_core::finality::{verify_binding, verify_pop, KeyReveal, BINDING_EXPORTER_LEN};
use std::collections::VecDeque;
use std::sync::atomic::Ordering;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, BufReader};
fn now() -> String {
let t = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap();
format!("{}.{:03}", t.as_secs(), t.subsec_millis())
}
/// 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,
match &pool.tls {
Some(t) => format!("TLS 1.3, certificate pin {}{}", t.pin, if t.self_signed { " (self-signed: members pass it as --pool-pin)" } else { "" }),
None => "PLAIN TCP: spec 9.3 wants TLS 1.3 (--tls-self-signed or --tls-cert/--tls-key)".to_string(),
}
);
loop {
match listener.accept().await {
Ok((sock, addr)) => {
let pool = pool.clone();
tokio::spawn(async move {
let _ = sock.set_nodelay(true);
match &pool.tls {
Some(t) => {
let acceptor = t.acceptor.clone();
match tokio::time::timeout(Duration::from_secs(10), acceptor.accept(sock)).await {
Ok(Ok(tls)) => {
let exporter = match crate::tls::exporter(tls.get_ref().1, pool.cfg.chain_id()) {
Ok(e) => e,
Err(e) => {
eprintln!("{} {addr}: {e}", now());
return;
}
};
connection(pool, tls, addr.to_string(), Some(exporter)).await;
}
Ok(Err(e)) => eprintln!("{} {addr}: TLS handshake failed: {e} (a member without --pool-tls or --pool-pin, or the wrong pin)", now()),
Err(_) => eprintln!("{} {addr}: TLS handshake timed out", now()),
}
}
None => connection(pool, sock, addr.to_string(), None).await,
}
});
}
Err(e) => {
eprintln!("{} accept: {e}", now());
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
}
fn welcome(pool: &Pool) -> Msg {
Msg::Welcome {
version: VERSION.into(),
chain_id: pool.cfg.chain_id(),
network: pool.cfg.network.clone(),
pool_address: pool.pool_address_hex.clone(),
modes: vec!["A".into()],
vote_mode: "member".into(),
share_scheme: match &pool.open {
Some(o) => ShareScheme {
scheme: "pplns-sidechain".into(),
fee_percent: 0.0,
pplns_window_blocks: 0.0,
min_payout_ign: 0.0,
software_dev_fee_percent: o.dev_fee_percent as f64,
window_shares: pool.cfg.window_shares as u64,
chain_share_s: pool.cfg.chain_share_s,
},
None => ShareScheme {
scheme: "pplns".into(),
fee_percent: pool.cfg.fee_percent,
pplns_window_blocks: pool.cfg.pplns_window_blocks,
min_payout_ign: pool.cfg.min_payout_ign,
software_dev_fee_percent: 0.0,
window_shares: 0,
chain_share_s: 0.0,
},
},
min_shift: pool.cfg.min_shift,
max_shift: pool.cfg.max_shift,
share_interval_s: pool.cfg.share_interval_s,
stale_grace_ms: pool.cfg.stale_grace_ms,
pool_name: pool.cfg.name.clone(),
pool_mode: if pool.open.is_some() { "open".into() } else { "operator".into() },
tls: pool.tls.is_some(),
}
}
async fn connection<S: AsyncRead + AsyncWrite + Send + 'static>(pool: Arc<Pool>, sock: S, remote: String, exporter: Option<[u8; BINDING_EXPORTER_LEN]>) {
let (rd, mut wr) = tokio::io::split(sock);
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<String>();
let writer = tokio::spawn(async move {
while let Some(line) = rx.recv().await {
if wr.write_all(line.as_bytes()).await.is_err() {
break;
}
}
});
let mut lines = BufReader::with_capacity(64 << 10, rd).lines();
let mut member: Option<Arc<Member>> = None;
let mut said_hello = false;
let mut buf_guard = 0usize;
loop {
let l = match tokio::time::timeout(Duration::from_secs(300), lines.next_line()).await {
Ok(Ok(Some(l))) => l,
Ok(Ok(None)) => break,
Ok(Err(e)) => {
eprintln!("{} {remote}: read error: {e}", now());
break;
}
Err(_) => {
let _ = tx.send(Msg::Bye { reason: "idle for 300 s".into() }.line());
break;
}
};
buf_guard += l.len();
if l.len() > MAX_LINE {
let _ = tx.send(Msg::Bye { reason: "line over 4 MiB".into() }.line());
break;
}
let msg = match Msg::parse(&l) {
Ok(m) => m,
Err(e) => {
let _ = tx.send(Msg::Error { code: "parse".into(), detail: e }.line());
continue;
}
};
match msg {
Msg::Hello { versions, chain_id, client, .. } => {
if !versions.iter().any(|v| v == VERSION) {
let _ = tx.send(Msg::Bye { reason: format!("protocol {VERSION} only") }.line());
break;
}
if chain_id != pool.cfg.chain_id() {
let _ = tx.send(Msg::Bye { reason: format!("chain id {} here, you said {chain_id}", pool.cfg.chain_id()) }.line());
break;
}
said_hello = true;
println!("{} {remote}: hello from {client}", now());
let _ = tx.send(welcome(&pool).line());
}
Msg::Authorize { pubkey, pop, label, payout, binding } => {
if !said_hello {
let _ = tx.send(Msg::Error { code: "order".into(), detail: "hello first".into() }.line());
continue;
}
let (Some(pk), Some(pp)) = (parse_hex_array::<48>(&pubkey), parse_hex_array::<96>(&pop)) else {
let _ = tx.send(Msg::Bye { reason: "pubkey is 48 bytes hex and pop 96".into() }.line());
break;
};
if !verify_pop(&pk, &pp) {
let _ = tx.send(Msg::Bye { reason: "proof of possession does not verify".into() }.line());
break;
}
// spec 9.3, O-9.7: over TLS the binding must be this member's signature over THIS connection's exporter
if let Some(e) = &exporter {
let ok = parse_hex_array::<96>(&binding).is_some_and(|sig| verify_binding(&pk, pool.cfg.chain_id(), e, &sig));
if !ok {
println!("{} {remote}: authorize REFUSED: the binding does not sign this connection's TLS exporter (a replay, or a member without the binding)", now());
let _ = tx.send(Msg::Bye { reason: "binding: the authorize must sign this connection's TLS exporter under your key (spec 9.3)".into() }.line());
break;
}
}
let Some(addr) = parse_hex_array::<20>(&payout) else {
let _ = tx.send(Msg::Bye { reason: "payout is 0x followed by 40 hex".into() }.line());
break;
};
let reveal = KeyReveal { pubkey: pk, pop: pp };
let key_hash = reveal.key_hash();
let worker: String = label.chars().filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_' || *c == '.').take(32).collect();
let worker = if worker.is_empty() { "worker".to_string() } else { worker };
let id = pool.next_member_id.fetch_add(1, Ordering::Relaxed) + 1;
let m = Arc::new(Member {
id,
address: format!("0x{}", hex::encode(addr)),
worker: worker.clone(),
key_hash,
reveal,
remote: remote.clone(),
connected_at: Instant::now(),
tx: tx.clone(),
inner: Mutex::new(MemberInner {
vardiff: Vardiff::new(pool.cfg.min_shift, pool.cfg.min_shift, pool.cfg.max_shift, pool.cfg.share_interval_s as f64, pool.now_s()),
jobs: VecDeque::new(),
last_parents: None,
last_seeds: None,
accepted: 0,
stale: 0,
rejected: 0,
blocks: 0,
refused: 0,
last_share: None,
revealed: false,
reported_hashrate: 0.0,
reported_workers: 1,
proving: false,
current_target64: 0,
in_flight_checks: 0,
last_template: None,
}),
});
pool.members.lock().unwrap().insert(id, m.clone());
member = Some(m.clone());
println!("{} {remote}: member {id} '{worker}' key {} pays {}", now(), &key_hash.to_string()[..16], m.address);
let _ = tx.send(Msg::Authorized { member_id: id, revealed: false, key_hash: key_hash.to_string() }.line());
let pool2 = pool.clone();
let m2 = m.clone();
tokio::spawn(async move {
if let Err(e) = fetch_job(&pool2, &m2).await {
eprintln!("{} first job for member {}: {e}", now(), m2.id);
}
});
}
Msg::Share { job_id, nonce, hash } => {
let Some(m) = member.clone() else {
let _ = tx.send(Msg::Error { code: "order".into(), detail: "authorize first".into() }.line());
continue;
};
let (Some(n), Some(h)) = (parse_hex_u64(&nonce), parse_hex_u64(&hash)) else {
let _ = tx.send(Msg::ShareResult { job_id, nonce, accepted: false, code: "parse".into(), weight: 0.0, block: false }.line());
continue;
};
handle_share(&pool, &m, job_id, n, h, false).await;
}
Msg::Solution { job_id, nonce, hash, .. } => {
// a solution follows the same nonce's share (the member sends both, spec 9.5): the share already
// credited and submitted it, so a seen nonce is acknowledged and not counted twice
let Some(m) = member.clone() else {
let _ = tx.send(Msg::Error { code: "order".into(), detail: "authorize first".into() }.line());
continue;
};
let (Some(n), Some(h)) = (parse_hex_u64(&nonce), parse_hex_u64(&hash)) else {
let _ = tx.send(Msg::ShareResult { job_id, nonce, accepted: false, code: "parse".into(), weight: 0.0, block: false }.line());
continue;
};
handle_share(&pool, &m, job_id, n, h, true).await;
}
Msg::Stats { hashrate, workers, .. } => {
if let Some(m) = &member {
let mut g = m.inner.lock().unwrap();
g.reported_hashrate = hashrate;
g.reported_workers = workers.max(1);
// consequence C6: a member that also proves reports it in `stats` (the `proving` field of the
// v0 stats line), so a 4% share drop reads as proving and not as a sick card
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&l) {
g.proving = v.get("proving").and_then(|p| p.as_bool()).unwrap_or(false);
}
}
}
Msg::JobRefused { job_id, code, detail } => {
if let Some(m) = &member {
m.inner.lock().unwrap().refused += 1;
println!("{} member {} refused job {job_id}: {code} {detail}", now(), m.id);
}
}
Msg::Ping { id } => {
let _ = tx.send(Msg::Pong { id }.line());
}
Msg::Pong { .. } => {}
Msg::Bye { reason } => {
println!("{} {remote}: bye ({reason})", now());
break;
}
other => {
let _ = tx.send(Msg::Error { code: "unexpected".into(), detail: format!("{:?}", std::mem::discriminant(&other)) }.line());
}
}
if buf_guard > (64 << 20) {
buf_guard = 0;
}
}
if let Some(m) = member {
pool.members.lock().unwrap().remove(&m.id);
let g = m.inner.lock().unwrap();
println!("{} {remote}: member {} '{}' left after {:.0} s: accepted={} stale={} rejected={} blocks={}", now(), m.id, m.worker, m.connected_at.elapsed().as_secs_f64(), g.accepted, g.stale, g.rejected, g.blocks);
}
drop(tx);
let _ = writer.await;
}
async fn handle_share(pool: &Arc<Pool>, m: &Arc<Member>, job_id: u64, nonce: u64, claimed: u64, is_solution: bool) {
// Book-keeping under the lock: the job, duplicate and stale checks
let found = {
let mut g = m.inner.lock().unwrap();
let grace = Duration::from_millis(pool.cfg.stale_grace_ms);
match g.jobs.iter_mut().find(|j| j.job_id == job_id) {
None => Err(Code::UnknownJob),
Some(j) => {
if !j.seen.insert(nonce) {
Err(Code::Duplicate)
} else if j.superseded_at.is_some_and(|t| t.elapsed() > grace) {
Err(Code::Stale)
} else {
Ok((j.key.clone(), j.epoch.clone(), j.shift, j.raw.clone(), j.daa_score, j.blue_score, j.open_parent, j.open_share_target64))
}
}
}
};
let (key, epoch, shift, raw, daa, blue, open_parent, open_target) = match found {
Ok(x) => x,
Err(Code::Duplicate) if is_solution => return,
Err(code) => {
{
let mut s = pool.state.lock().unwrap();
s.refuse_share(&m.address, &m.worker, code == Code::Stale, None);
}
{
let mut g = m.inner.lock().unwrap();
if code == Code::Stale {
g.stale += 1
} else {
g.rejected += 1
}
}
m.send(Msg::ShareResult { job_id, nonce: hex_u64(nonce), accepted: false, code: code.as_str().into(), weight: 0.0, block: false }.line());
return;
}
};
// The CPU evaluation, off the async threads, bounded by the verifier permits
let permit = pool.verify_permits.acquire().await;
let key2 = key.clone();
let v = tokio::task::spawn_blocking(move || check(&epoch, &key2, nonce, claimed)).await.expect("verifier task");
drop(permit);
let now_s = pool.now_s();
if v.code != Code::Ok {
{
let mut s = pool.state.lock().unwrap();
s.refuse_share(&m.address, &m.worker, false, Some(v.cost_ms));
}
m.inner.lock().unwrap().rejected += 1;
m.send(Msg::ShareResult { job_id, nonce: hex_u64(nonce), accepted: false, code: v.code.as_str().into(), weight: 0.0, block: false }.line());
if v.code == Code::WrongHash {
println!("{} member {} ({}) WRONG HASH nonce={:#x} claimed={:016x} pool={:016x}", now(), m.id, m.worker, nonce, claimed, v.hash);
}
return;
}
let weight = share_weight(shift);
{
let mut s = pool.state.lock().unwrap();
s.accept_share(&m.address, &m.worker, weight, key.share_target64, v.cost_ms);
}
let retarget = {
let mut g = m.inner.lock().unwrap();
g.accepted += 1;
g.last_share = Some(Instant::now());
g.vardiff.on_share(now_s);
let t64 = g.current_target64;
g.vardiff.retarget(now_s, t64).map(|s| (s, share_target(t64, s)))
};
m.send(Msg::ShareResult { job_id, nonce: hex_u64(nonce), accepted: true, code: "ok".into(), weight, block: v.block }.line());
// the open pool: a hash at or below the chain's target is a share of the share chain (open.rs)
if let (Some(o), Some(parent)) = (&pool.open, open_parent)
&& open_target > 0
&& v.hash <= open_target
{
o.local_share(pool, m, &raw, nonce, v.hash, &key, open_target, parent);
}
if let Some((s, st)) = retarget {
m.send(Msg::SetTarget { shift: s, share_target64: hex_u64(st) }.line());
println!("{} VARDIFF member {} ({}) shift -> {} (share target {:016x})", now(), m.id, m.worker, s, st);
let (pool2, m2) = (pool.clone(), m.clone());
tokio::spawn(async move {
let _ = fetch_job(&pool2, &m2).await;
});
}
if v.block {
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");
}
}