//! Member connections: newline JSON over TCP (spec 9.3 wants TLS 1.3; v0 listens in the clear and the README puts a //! TLS terminator in front, docs/plans/pool.md). One reader task and one writer task per connection. 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_pop, KeyReveal}; use std::collections::VecDeque; use std::sync::atomic::Ordering; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::TcpStream; fn now() -> String { let t = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap(); format!("{}.{:03}", t.as_secs(), t.subsec_millis()) } pub async fn listen(pool: Arc) { 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) }); println!("{} pool: members on {} (chain id {}, {})", now(), pool.cfg.listen, pool.cfg.chain_id(), pool.cfg.network); loop { match listener.accept().await { Ok((sock, addr)) => { let pool = pool.clone(); tokio::spawn(async move { let _ = sock.set_nodelay(true); connection(pool, sock, addr.to_string()).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: 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 }, 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(), } } async fn connection(pool: Arc, sock: TcpStream, remote: String) { let (rd, mut wr) = sock.into_split(); let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); 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> = 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, .. } => { 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; } 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, }), }); 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::(&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, m: &Arc, 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)) } } } }; let (key, epoch, shift, raw, daa, blue) = 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()); 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)); } }