stratum.rs (22026B)
1 use std::{ 2 collections::BTreeMap, 3 net::{IpAddr, Ipv6Addr, SocketAddr}, 4 sync::{ 5 Arc, Mutex as StdMutex, 6 atomic::{AtomicU64, Ordering}, 7 }, 8 time::Duration, 9 }; 10 11 use anyhow::{Context, Result, bail}; 12 use serde_json::{Value, json}; 13 use tokio::{ 14 io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader}, 15 net::{TcpListener, TcpStream, tcp::OwnedWriteHalf}, 16 sync::{Mutex, OwnedSemaphorePermit, Semaphore}, 17 time::timeout, 18 }; 19 20 use crate::{ 21 adapters::p2p::GossipNetwork, 22 app::{ExternalMineJob, SharedNode, debug_logging_enabled}, 23 domain::{STRATUM_EXTRANONCE1_HEX, STRATUM_EXTRANONCE2_SIZE, StratumMineShare}, 24 }; 25 26 const STRATUM_MAX_LINE_BYTES: usize = 16 * 1024; 27 const STRATUM_MAX_JOBS_PER_SESSION: usize = 128; 28 const STRATUM_MAX_SESSIONS: usize = 64; 29 const STRATUM_MAX_SESSIONS_PER_SOURCE: usize = 4; 30 const STRATUM_IDLE_TIMEOUT: Duration = Duration::from_secs(120); 31 const STRATUM_WRITE_TIMEOUT: Duration = Duration::from_secs(10); 32 33 #[cfg(feature = "fuzzing")] 34 pub fn fuzz_parse_stratum_request(line: &str) -> Result<Value> { 35 parse_stratum_request(line) 36 } 37 38 fn parse_stratum_request(line: &str) -> Result<Value> { 39 if line.len() > STRATUM_MAX_LINE_BYTES { 40 bail!("Stratum request exceeds {STRATUM_MAX_LINE_BYTES} byte limit"); 41 } 42 serde_json::from_str(line).context("invalid Stratum JSON") 43 } 44 45 #[derive(Clone)] 46 pub struct StratumServer { 47 node: SharedNode, 48 gossip: GossipNetwork, 49 listen_addr: SocketAddr, 50 next_job_salt: Arc<AtomicU64>, 51 session_limiter: StratumSessionLimiter, 52 } 53 54 #[derive(Clone, Debug)] 55 struct StratumJob { 56 mine: ExternalMineJob, 57 } 58 59 #[derive(Clone)] 60 struct StratumSessionLimiter { 61 permits: Arc<Semaphore>, 62 active_by_source: Arc<StdMutex<BTreeMap<IpAddr, usize>>>, 63 } 64 65 struct StratumSessionPermit { 66 _global: OwnedSemaphorePermit, 67 source: IpAddr, 68 active_by_source: Arc<StdMutex<BTreeMap<IpAddr, usize>>>, 69 } 70 71 impl StratumSessionLimiter { 72 fn new(max_sessions: usize) -> Self { 73 Self { 74 permits: Arc::new(Semaphore::new(max_sessions)), 75 active_by_source: Arc::new(StdMutex::new(BTreeMap::new())), 76 } 77 } 78 79 fn try_acquire(&self, remote_ip: IpAddr) -> Option<StratumSessionPermit> { 80 let global = self.permits.clone().try_acquire_owned().ok()?; 81 let source = stratum_source_key(remote_ip); 82 let mut active_by_source = self.active_by_source.lock().ok()?; 83 let active = active_by_source.entry(source).or_default(); 84 if *active >= STRATUM_MAX_SESSIONS_PER_SOURCE { 85 return None; 86 } 87 *active += 1; 88 drop(active_by_source); 89 Some(StratumSessionPermit { 90 _global: global, 91 source, 92 active_by_source: Arc::clone(&self.active_by_source), 93 }) 94 } 95 } 96 97 impl Drop for StratumSessionPermit { 98 fn drop(&mut self) { 99 let Ok(mut active_by_source) = self.active_by_source.lock() else { 100 return; 101 }; 102 let Some(active) = active_by_source.get_mut(&self.source) else { 103 return; 104 }; 105 *active = active.saturating_sub(1); 106 if *active == 0 { 107 active_by_source.remove(&self.source); 108 } 109 } 110 } 111 112 fn stratum_source_key(ip: IpAddr) -> IpAddr { 113 match ip { 114 IpAddr::V4(_) => ip, 115 IpAddr::V6(ip) => { 116 let segments = ip.segments(); 117 IpAddr::V6(Ipv6Addr::new( 118 segments[0], 119 segments[1], 120 segments[2], 121 segments[3], 122 0, 123 0, 124 0, 125 0, 126 )) 127 } 128 } 129 } 130 131 impl StratumServer { 132 pub async fn start( 133 node: SharedNode, 134 gossip: GossipNetwork, 135 listen_addr: SocketAddr, 136 ) -> Result<Self> { 137 let listener = TcpListener::bind(listen_addr) 138 .await 139 .with_context(|| format!("failed to bind Stratum listener on {listen_addr}"))?; 140 let local_addr = listener.local_addr()?; 141 let server = Self { 142 node, 143 gossip, 144 listen_addr: local_addr, 145 next_job_salt: Arc::new(AtomicU64::new(1)), 146 session_limiter: StratumSessionLimiter::new(STRATUM_MAX_SESSIONS), 147 }; 148 tokio::spawn(run_listener(server.clone(), listener)); 149 Ok(server) 150 } 151 152 pub fn listen_addr(&self) -> SocketAddr { 153 self.listen_addr 154 } 155 } 156 157 async fn run_listener(server: StratumServer, listener: TcpListener) { 158 loop { 159 match listener.accept().await { 160 Ok((stream, remote)) => { 161 let Some(permit) = server.session_limiter.try_acquire(remote.ip()) else { 162 if debug_logging_enabled() { 163 eprintln!("stratum session with {remote} rejected: session limit reached"); 164 } 165 continue; 166 }; 167 let server = server.clone(); 168 tokio::spawn(async move { 169 let _permit = permit; 170 if let Err(error) = handle_connection(server, stream).await { 171 if debug_logging_enabled() { 172 eprintln!("stratum session with {remote} failed: {error:#}"); 173 } 174 } 175 }); 176 } 177 Err(error) if debug_logging_enabled() => { 178 eprintln!("stratum accept failed: {error:#}"); 179 } 180 Err(_) => {} 181 } 182 } 183 } 184 185 async fn handle_connection(server: StratumServer, stream: TcpStream) -> Result<()> { 186 let (read, write) = stream.into_split(); 187 let mut session = StratumSession { 188 server, 189 writer: Arc::new(Mutex::new(write)), 190 authorized_worker: None, 191 jobs: BTreeMap::new(), 192 next_job_id: 1, 193 }; 194 let mut lines = StratumLineReader::new(read); 195 loop { 196 let Some(line) = timeout(STRATUM_IDLE_TIMEOUT, lines.read_line()) 197 .await 198 .context("Stratum session idle timeout")?? 199 else { 200 break; 201 }; 202 if line.trim().is_empty() { 203 continue; 204 } 205 let request = parse_stratum_request(&line)?; 206 session.handle_request(request).await?; 207 } 208 Ok(()) 209 } 210 211 struct StratumLineReader<R> { 212 reader: BufReader<R>, 213 pending: Vec<u8>, 214 } 215 216 impl<R: AsyncRead + Unpin> StratumLineReader<R> { 217 fn new(reader: R) -> Self { 218 Self { 219 reader: BufReader::new(reader), 220 pending: Vec::new(), 221 } 222 } 223 224 async fn read_line(&mut self) -> Result<Option<String>> { 225 loop { 226 let available = self.reader.fill_buf().await?; 227 if available.is_empty() { 228 if self.pending.is_empty() { 229 return Ok(None); 230 } 231 bail!("client closed before completing a Stratum request"); 232 } 233 234 if let Some(newline) = available.iter().position(|byte| *byte == b'\n') { 235 if self.pending.len() + newline > STRATUM_MAX_LINE_BYTES { 236 bail!("Stratum request exceeds {STRATUM_MAX_LINE_BYTES} byte limit"); 237 } 238 self.pending.extend_from_slice(&available[..newline]); 239 self.reader.consume(newline + 1); 240 if self.pending.ends_with(b"\r") { 241 self.pending.pop(); 242 } 243 let bytes = std::mem::take(&mut self.pending); 244 return String::from_utf8(bytes) 245 .context("Stratum request is not valid UTF-8") 246 .map(Some); 247 } 248 249 if self.pending.len() + available.len() > STRATUM_MAX_LINE_BYTES { 250 bail!("Stratum request exceeds {STRATUM_MAX_LINE_BYTES} byte limit"); 251 } 252 let consumed = available.len(); 253 self.pending.extend_from_slice(available); 254 self.reader.consume(consumed); 255 } 256 } 257 } 258 259 struct StratumSession { 260 server: StratumServer, 261 writer: Arc<Mutex<OwnedWriteHalf>>, 262 authorized_worker: Option<String>, 263 jobs: BTreeMap<String, StratumJob>, 264 next_job_id: u64, 265 } 266 267 impl StratumSession { 268 async fn handle_request(&mut self, request: Value) -> Result<()> { 269 let id = request.get("id").cloned().unwrap_or(Value::Null); 270 let method = request 271 .get("method") 272 .and_then(Value::as_str) 273 .context("Stratum request is missing method")?; 274 match method { 275 "mining.subscribe" => { 276 self.send_response( 277 id, 278 json!([ 279 [["mining.set_difficulty", "iuna"], ["mining.notify", "iuna"]], 280 STRATUM_EXTRANONCE1_HEX, 281 STRATUM_EXTRANONCE2_SIZE 282 ]), 283 ) 284 .await?; 285 } 286 "mining.authorize" => { 287 let worker = request 288 .get("params") 289 .and_then(Value::as_array) 290 .and_then(|params| params.first()) 291 .and_then(Value::as_str) 292 .context("mining.authorize requires worker address")? 293 .to_string(); 294 self.normalized_worker_recipient(&worker).await?; 295 self.authorized_worker = Some(worker.clone()); 296 self.send_response(id, json!(true)).await?; 297 self.send_job(&worker, true).await?; 298 } 299 "mining.submit" => { 300 let accepted = self.handle_submit(&request).await; 301 match accepted { 302 Ok(true) => self.send_response(id, json!(true)).await?, 303 Ok(false) => { 304 self.send_error(id, 23, "duplicate share or transaction") 305 .await?; 306 } 307 Err(error) => self.send_error(id, 23, &format!("{error:#}")).await?, 308 } 309 } 310 "mining.configure" => { 311 self.send_response(id, json!({})).await?; 312 } 313 "mining.extranonce.subscribe" => { 314 self.send_response(id, json!(true)).await?; 315 } 316 _ => { 317 self.send_error(id, 20, &format!("unsupported method {method}")) 318 .await?; 319 } 320 } 321 Ok(()) 322 } 323 324 async fn send_job(&mut self, worker: &str, clean_jobs: bool) -> Result<()> { 325 let job_id = self.next_job_id.to_string(); 326 self.next_job_id = self.next_job_id.saturating_add(1); 327 let salt = self.server.next_job_salt.fetch_add(1, Ordering::Relaxed); 328 let recipient = self.normalized_worker_recipient(worker).await?; 329 let mine = self 330 .server 331 .node 332 .lock() 333 .await 334 .external_mine_job(recipient, salt)?; 335 let difficulty = stratum_difficulty_for_bits(mine.template.difficulty_bits); 336 self.send_notification("mining.set_difficulty", json!([difficulty])) 337 .await?; 338 self.send_notification( 339 "mining.notify", 340 json!([ 341 job_id, 342 mine.template.prev_hash_hex, 343 mine.template.coinb1_hex(), 344 "", 345 [], 346 mine.template.version_hex, 347 mine.template.nbits_hex, 348 mine.template.ntime_hex, 349 clean_jobs 350 ]), 351 ) 352 .await?; 353 insert_bounded_job(&mut self.jobs, job_id, StratumJob { mine }); 354 Ok(()) 355 } 356 357 async fn handle_submit(&mut self, request: &Value) -> Result<bool> { 358 let params = request 359 .get("params") 360 .and_then(Value::as_array) 361 .context("mining.submit requires params")?; 362 let worker = str_param(params, 0, "worker")?; 363 let job_id = str_param(params, 1, "job id")?; 364 let extranonce2 = hex_array_4(str_param(params, 2, "extranonce2")?)?; 365 let ntime = str_param(params, 3, "ntime")?; 366 let header_nonce = hex_array_4(str_param(params, 4, "nonce")?)?; 367 let authorized = self 368 .authorized_worker 369 .as_deref() 370 .context("worker is not authorized")?; 371 if worker != authorized { 372 bail!("submitted worker does not match authorized worker"); 373 } 374 let job = self 375 .jobs 376 .get(job_id) 377 .cloned() 378 .context("unknown Stratum job")?; 379 if ntime != job.mine.template.ntime_hex { 380 bail!("submitted ntime does not match job"); 381 } 382 383 let (result, outbox) = { 384 let mut node = self.server.node.lock().await; 385 let recipient = node.normalize_user_address(recipient_from_worker(worker))?; 386 let result = node.submit_external_mine( 387 recipient, 388 job.mine.template.clone(), 389 StratumMineShare { 390 extranonce2, 391 header_nonce, 392 }, 393 ); 394 let outbox = node.drain_outbox(); 395 (result, outbox) 396 }; 397 match result { 398 Ok(_) => { 399 self.server.gossip.broadcast(outbox).await?; 400 let worker = worker.to_string(); 401 self.send_job(&worker, false).await?; 402 Ok(true) 403 } 404 Err(error) => Err(error), 405 } 406 } 407 408 async fn normalized_worker_recipient(&self, worker: &str) -> Result<String> { 409 self.server 410 .node 411 .lock() 412 .await 413 .normalize_user_address(recipient_from_worker(worker)) 414 .context("invalid Stratum worker address") 415 } 416 417 async fn send_response(&self, id: Value, result: Value) -> Result<()> { 418 self.send(json!({ "id": id, "result": result, "error": null })) 419 .await 420 } 421 422 async fn send_error(&self, id: Value, code: i64, message: &str) -> Result<()> { 423 self.send(json!({ "id": id, "result": null, "error": [code, message, null] })) 424 .await 425 } 426 427 async fn send_notification(&self, method: &str, params: Value) -> Result<()> { 428 self.send(json!({ "id": null, "method": method, "params": params })) 429 .await 430 } 431 432 async fn send(&self, value: Value) -> Result<()> { 433 let mut payload = serde_json::to_vec(&value)?; 434 payload.push(b'\n'); 435 let mut writer = self.writer.lock().await; 436 timeout(STRATUM_WRITE_TIMEOUT, writer.write_all(&payload)) 437 .await 438 .context("Stratum response write timeout")??; 439 Ok(()) 440 } 441 } 442 443 fn insert_bounded_job(jobs: &mut BTreeMap<String, StratumJob>, job_id: String, job: StratumJob) { 444 jobs.insert(job_id, job); 445 while jobs.len() > STRATUM_MAX_JOBS_PER_SESSION { 446 let Some(oldest) = oldest_job_id(jobs) else { 447 break; 448 }; 449 jobs.remove(&oldest); 450 } 451 } 452 453 fn oldest_job_id(jobs: &BTreeMap<String, StratumJob>) -> Option<String> { 454 jobs.keys() 455 .min_by(|left, right| { 456 stratum_job_id_sort_key(left) 457 .cmp(&stratum_job_id_sort_key(right)) 458 .then_with(|| left.cmp(right)) 459 }) 460 .cloned() 461 } 462 463 fn stratum_job_id_sort_key(job_id: &str) -> u64 { 464 job_id.parse().unwrap_or(u64::MAX) 465 } 466 467 fn str_param<'a>(params: &'a [Value], index: usize, name: &str) -> Result<&'a str> { 468 params 469 .get(index) 470 .and_then(Value::as_str) 471 .with_context(|| format!("mining.submit requires {name}")) 472 } 473 474 fn recipient_from_worker(worker: &str) -> &str { 475 worker 476 .split_once('.') 477 .map_or(worker, |(recipient, _)| recipient) 478 } 479 480 fn hex_array_4(input: &str) -> Result<[u8; 4]> { 481 let bytes = decode_hex(input)?; 482 let len = bytes.len(); 483 bytes 484 .try_into() 485 .map_err(|_| anyhow::anyhow!("expected 4 hex bytes, got {len}")) 486 } 487 488 fn decode_hex(input: &str) -> Result<Vec<u8>> { 489 if input.len() % 2 != 0 { 490 bail!("hex string has odd length"); 491 } 492 let mut bytes = Vec::with_capacity(input.len() / 2); 493 for pair in input.as_bytes().chunks_exact(2) { 494 bytes.push((hex_value(pair[0])? << 4) | hex_value(pair[1])?); 495 } 496 Ok(bytes) 497 } 498 499 fn hex_value(byte: u8) -> Result<u8> { 500 match byte { 501 b'0'..=b'9' => Ok(byte - b'0'), 502 b'a'..=b'f' => Ok(byte - b'a' + 10), 503 b'A'..=b'F' => Ok(byte - b'A' + 10), 504 _ => bail!("invalid hex character"), 505 } 506 } 507 508 fn stratum_difficulty_for_bits(bits: u32) -> f64 { 509 2_f64.powi(bits as i32 - 16).max(0.000001) 510 } 511 512 #[cfg(test)] 513 mod tests { 514 use tokio::io::AsyncWriteExt; 515 516 use crate::{app::ExternalMineJob, domain::StratumMineTemplate}; 517 518 use super::{ 519 STRATUM_MAX_JOBS_PER_SESSION, STRATUM_MAX_LINE_BYTES, STRATUM_MAX_SESSIONS, 520 STRATUM_MAX_SESSIONS_PER_SOURCE, StratumJob, StratumLineReader, StratumSessionLimiter, 521 insert_bounded_job, parse_stratum_request, 522 }; 523 524 fn dummy_job() -> StratumJob { 525 StratumJob { 526 mine: ExternalMineJob { 527 template: StratumMineTemplate { 528 recipient: "0".repeat(64), 529 anchor: "0".repeat(64), 530 salt: 0, 531 difficulty_bits: 0, 532 coinbase_prefix: Vec::new(), 533 version_hex: "00000000".to_string(), 534 prev_hash_hex: "0".repeat(64), 535 nbits_hex: "00000000".to_string(), 536 ntime_hex: "00000000".to_string(), 537 }, 538 }, 539 } 540 } 541 542 #[tokio::test] 543 async fn stratum_line_reader_accepts_max_sized_line() { 544 let (mut client, server) = tokio::io::duplex(STRATUM_MAX_LINE_BYTES + 1); 545 let mut reader = StratumLineReader::new(server); 546 let line = vec![b'a'; STRATUM_MAX_LINE_BYTES]; 547 client.write_all(&line).await.unwrap(); 548 client.write_all(b"\n").await.unwrap(); 549 550 let read = reader.read_line().await.unwrap().unwrap(); 551 552 assert_eq!(read.len(), STRATUM_MAX_LINE_BYTES); 553 } 554 555 #[tokio::test] 556 async fn stratum_line_reader_rejects_oversized_line() { 557 let (mut client, server) = tokio::io::duplex(STRATUM_MAX_LINE_BYTES + 2); 558 let mut reader = StratumLineReader::new(server); 559 let line = vec![b'a'; STRATUM_MAX_LINE_BYTES + 1]; 560 client.write_all(&line).await.unwrap(); 561 client.write_all(b"\n").await.unwrap(); 562 563 let error = reader.read_line().await.unwrap_err(); 564 565 assert!(error.to_string().contains("Stratum request exceeds")); 566 } 567 568 #[test] 569 fn stratum_job_cache_prunes_oldest_jobs() { 570 let mut jobs = std::collections::BTreeMap::new(); 571 for id in 1..=STRATUM_MAX_JOBS_PER_SESSION + 2 { 572 insert_bounded_job(&mut jobs, id.to_string(), dummy_job()); 573 } 574 575 assert_eq!(jobs.len(), STRATUM_MAX_JOBS_PER_SESSION); 576 assert!(!jobs.contains_key("1")); 577 assert!(!jobs.contains_key("2")); 578 assert!(jobs.contains_key("3")); 579 } 580 581 #[test] 582 fn stratum_session_limiter_enforces_global_cap() { 583 let limiter = StratumSessionLimiter::new(STRATUM_MAX_SESSIONS); 584 let permits = (0..STRATUM_MAX_SESSIONS) 585 .map(|index| { 586 limiter 587 .try_acquire(format!("192.0.2.{}", index + 1).parse().unwrap()) 588 .expect("permit should be available") 589 }) 590 .collect::<Vec<_>>(); 591 592 assert!( 593 limiter 594 .try_acquire("198.51.100.1".parse().unwrap()) 595 .is_none() 596 ); 597 drop(permits); 598 assert!( 599 limiter 600 .try_acquire("198.51.100.1".parse().unwrap()) 601 .is_some() 602 ); 603 } 604 605 #[test] 606 fn stratum_session_limiter_prevents_one_source_from_exhausting_global_slots() { 607 let limiter = StratumSessionLimiter::new(STRATUM_MAX_SESSIONS); 608 let source = "192.0.2.10".parse().unwrap(); 609 let permits = (0..STRATUM_MAX_SESSIONS_PER_SOURCE) 610 .map(|_| limiter.try_acquire(source).unwrap()) 611 .collect::<Vec<_>>(); 612 613 assert!(limiter.try_acquire(source).is_none()); 614 assert!(limiter.try_acquire("192.0.2.11".parse().unwrap()).is_some()); 615 drop(permits); 616 assert!(limiter.try_acquire(source).is_some()); 617 } 618 619 #[test] 620 fn stratum_session_limiter_groups_ipv6_clients_by_prefix() { 621 let limiter = StratumSessionLimiter::new(STRATUM_MAX_SESSIONS); 622 let permits = (1..=STRATUM_MAX_SESSIONS_PER_SOURCE) 623 .map(|index| { 624 limiter 625 .try_acquire(format!("2001:db8:1:2::{index}").parse().unwrap()) 626 .unwrap() 627 }) 628 .collect::<Vec<_>>(); 629 630 assert!( 631 limiter 632 .try_acquire("2001:db8:1:2::ffff".parse().unwrap()) 633 .is_none() 634 ); 635 assert!( 636 limiter 637 .try_acquire("2001:db8:1:3::1".parse().unwrap()) 638 .is_some() 639 ); 640 drop(permits); 641 } 642 643 #[test] 644 fn performance_budget_stratum_request_parser_enforces_line_size() { 645 let prefix = r#"{"id":1,"method":"mining.configure","params":[""#; 646 let suffix = r#""]}"#; 647 let fill_len = STRATUM_MAX_LINE_BYTES - prefix.len() - suffix.len(); 648 let at_budget = format!("{prefix}{}{suffix}", "a".repeat(fill_len)); 649 650 let parsed = 651 parse_stratum_request(&at_budget).expect("request at line budget should parse"); 652 assert_eq!( 653 parsed.get("method").and_then(serde_json::Value::as_str), 654 Some("mining.configure") 655 ); 656 657 let over_budget = format!("{prefix}{}{suffix}", "a".repeat(fill_len + 1)); 658 let error = parse_stratum_request(&over_budget).unwrap_err(); 659 assert!(error.to_string().contains("Stratum request exceeds")); 660 } 661 }