commit ee236f5f8a08d52f3ded388158e1448081691c9b
parent 36cfbfc9a31ae66a025b32a6d6a83de09ec491c5
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Wed, 19 Aug 2026 14:14:45 +0200
Add Stratum DoS limits
Diffstat:
2 files changed, 198 insertions(+), 6 deletions(-)
diff --git a/ROADMAP.md b/ROADMAP.md
@@ -40,7 +40,7 @@ These items are not protocol rules. They are the attack and reliability checks t
- [x] Block validation with burn attestations remains independent of local mempool contents, including empty and conflicting mempools.
- [x] Post-genesis transactions cannot spend with `genesis` input signatures.
- [x] P2P envelope item limits reject batches only above their configured boundaries.
-- [ ] Stratum endpoint has explicit DoS limits: maximum line size, maximum jobs per session, idle timeout, and connection/session caps.
+- [x] Stratum endpoint has explicit DoS limits: maximum line size, maximum jobs per session, idle timeout, and connection/session caps.
- [ ] Fork and snapshot adversarial tests cover same-height leader-quality choice, taller valid forks inside finality, invalid late snapshot blocks, and pending transaction carry-forward after reorg.
- [ ] Compact snapshot decoder has malformed-input tests for huge lengths, oversized varints, trailing bytes, truncated payloads, invalid tags, and random byte inputs without panics or excessive allocation.
- [ ] Supply invariant tests cover mixed burns, fees, PoW mine actions, reorgs, no replay, and no double spend.
diff --git a/src/adapters/stratum.rs b/src/adapters/stratum.rs
@@ -5,14 +5,16 @@ use std::{
Arc,
atomic::{AtomicU64, Ordering},
},
+ time::Duration,
};
use anyhow::{Context, Result, bail};
use serde_json::{Value, json};
use tokio::{
- io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
+ io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader},
net::{TcpListener, TcpStream, tcp::OwnedWriteHalf},
- sync::Mutex,
+ sync::{Mutex, OwnedSemaphorePermit, Semaphore},
+ time::timeout,
};
use crate::{
@@ -21,12 +23,18 @@ use crate::{
domain::{STRATUM_EXTRANONCE1_HEX, STRATUM_EXTRANONCE2_SIZE, StratumMineShare},
};
+const STRATUM_MAX_LINE_BYTES: usize = 16 * 1024;
+const STRATUM_MAX_JOBS_PER_SESSION: usize = 128;
+const STRATUM_MAX_SESSIONS: usize = 64;
+const STRATUM_IDLE_TIMEOUT: Duration = Duration::from_secs(120);
+
#[derive(Clone)]
pub struct StratumServer {
node: SharedNode,
gossip: GossipNetwork,
listen_addr: SocketAddr,
next_job_salt: Arc<AtomicU64>,
+ session_limiter: StratumSessionLimiter,
}
#[derive(Clone, Debug)]
@@ -34,6 +42,23 @@ struct StratumJob {
mine: ExternalMineJob,
}
+#[derive(Clone)]
+struct StratumSessionLimiter {
+ permits: Arc<Semaphore>,
+}
+
+impl StratumSessionLimiter {
+ fn new(max_sessions: usize) -> Self {
+ Self {
+ permits: Arc::new(Semaphore::new(max_sessions)),
+ }
+ }
+
+ fn try_acquire(&self) -> Option<OwnedSemaphorePermit> {
+ self.permits.clone().try_acquire_owned().ok()
+ }
+}
+
impl StratumServer {
pub async fn start(
node: SharedNode,
@@ -49,6 +74,7 @@ impl StratumServer {
gossip,
listen_addr: local_addr,
next_job_salt: Arc::new(AtomicU64::new(1)),
+ session_limiter: StratumSessionLimiter::new(STRATUM_MAX_SESSIONS),
};
tokio::spawn(run_listener(server.clone(), listener));
Ok(server)
@@ -63,8 +89,15 @@ async fn run_listener(server: StratumServer, listener: TcpListener) {
loop {
match listener.accept().await {
Ok((stream, remote)) => {
+ let Some(permit) = server.session_limiter.try_acquire() else {
+ if debug_logging_enabled() {
+ eprintln!("stratum session with {remote} rejected: session limit reached");
+ }
+ continue;
+ };
let server = server.clone();
tokio::spawn(async move {
+ let _permit = permit;
if let Err(error) = handle_connection(server, stream).await {
if debug_logging_enabled() {
eprintln!("stratum session with {remote} failed: {error:#}");
@@ -89,8 +122,14 @@ async fn handle_connection(server: StratumServer, stream: TcpStream) -> Result<(
jobs: BTreeMap::new(),
next_job_id: 1,
};
- let mut lines = BufReader::new(read).lines();
- while let Some(line) = lines.next_line().await? {
+ let mut lines = StratumLineReader::new(read);
+ loop {
+ let Some(line) = timeout(STRATUM_IDLE_TIMEOUT, lines.read_line())
+ .await
+ .context("Stratum session idle timeout")??
+ else {
+ break;
+ };
if line.trim().is_empty() {
continue;
}
@@ -100,6 +139,54 @@ async fn handle_connection(server: StratumServer, stream: TcpStream) -> Result<(
Ok(())
}
+struct StratumLineReader<R> {
+ reader: BufReader<R>,
+ pending: Vec<u8>,
+}
+
+impl<R: AsyncRead + Unpin> StratumLineReader<R> {
+ fn new(reader: R) -> Self {
+ Self {
+ reader: BufReader::new(reader),
+ pending: Vec::new(),
+ }
+ }
+
+ async fn read_line(&mut self) -> Result<Option<String>> {
+ loop {
+ let available = self.reader.fill_buf().await?;
+ if available.is_empty() {
+ if self.pending.is_empty() {
+ return Ok(None);
+ }
+ bail!("client closed before completing a Stratum request");
+ }
+
+ if let Some(newline) = available.iter().position(|byte| *byte == b'\n') {
+ if self.pending.len() + newline > STRATUM_MAX_LINE_BYTES {
+ bail!("Stratum request exceeds {STRATUM_MAX_LINE_BYTES} byte limit");
+ }
+ self.pending.extend_from_slice(&available[..newline]);
+ self.reader.consume(newline + 1);
+ if self.pending.ends_with(b"\r") {
+ self.pending.pop();
+ }
+ let bytes = std::mem::take(&mut self.pending);
+ return String::from_utf8(bytes)
+ .context("Stratum request is not valid UTF-8")
+ .map(Some);
+ }
+
+ if self.pending.len() + available.len() > STRATUM_MAX_LINE_BYTES {
+ bail!("Stratum request exceeds {STRATUM_MAX_LINE_BYTES} byte limit");
+ }
+ let consumed = available.len();
+ self.pending.extend_from_slice(available);
+ self.reader.consume(consumed);
+ }
+ }
+}
+
struct StratumSession {
server: StratumServer,
writer: Arc<Mutex<OwnedWriteHalf>>,
@@ -192,7 +279,7 @@ impl StratumSession {
]),
)
.await?;
- self.jobs.insert(job_id, StratumJob { mine });
+ insert_bounded_job(&mut self.jobs, job_id, StratumJob { mine });
Ok(())
}
@@ -271,6 +358,30 @@ impl StratumSession {
}
}
+fn insert_bounded_job(jobs: &mut BTreeMap<String, StratumJob>, job_id: String, job: StratumJob) {
+ jobs.insert(job_id, job);
+ while jobs.len() > STRATUM_MAX_JOBS_PER_SESSION {
+ let Some(oldest) = oldest_job_id(jobs) else {
+ break;
+ };
+ jobs.remove(&oldest);
+ }
+}
+
+fn oldest_job_id(jobs: &BTreeMap<String, StratumJob>) -> Option<String> {
+ jobs.keys()
+ .min_by(|left, right| {
+ stratum_job_id_sort_key(left)
+ .cmp(&stratum_job_id_sort_key(right))
+ .then_with(|| left.cmp(right))
+ })
+ .cloned()
+}
+
+fn stratum_job_id_sort_key(job_id: &str) -> u64 {
+ job_id.parse().unwrap_or(u64::MAX)
+}
+
fn str_param<'a>(params: &'a [Value], index: usize, name: &str) -> Result<&'a str> {
params
.get(index)
@@ -315,3 +426,84 @@ fn hex_value(byte: u8) -> Result<u8> {
fn stratum_difficulty_for_bits(bits: u32) -> f64 {
2_f64.powi(bits as i32 - 16).max(0.000001)
}
+
+#[cfg(test)]
+mod tests {
+ use tokio::io::AsyncWriteExt;
+
+ use crate::{app::ExternalMineJob, domain::StratumMineTemplate};
+
+ use super::{
+ STRATUM_MAX_JOBS_PER_SESSION, STRATUM_MAX_LINE_BYTES, STRATUM_MAX_SESSIONS, StratumJob,
+ StratumLineReader, StratumSessionLimiter, insert_bounded_job,
+ };
+
+ fn dummy_job() -> StratumJob {
+ StratumJob {
+ mine: ExternalMineJob {
+ template: StratumMineTemplate {
+ recipient: "0".repeat(64),
+ anchor: "0".repeat(64),
+ salt: 0,
+ difficulty_bits: 0,
+ coinbase_prefix: Vec::new(),
+ version_hex: "00000000".to_string(),
+ prev_hash_hex: "0".repeat(64),
+ nbits_hex: "00000000".to_string(),
+ ntime_hex: "00000000".to_string(),
+ },
+ },
+ }
+ }
+
+ #[tokio::test]
+ async fn stratum_line_reader_accepts_max_sized_line() {
+ let (mut client, server) = tokio::io::duplex(STRATUM_MAX_LINE_BYTES + 1);
+ let mut reader = StratumLineReader::new(server);
+ let line = vec![b'a'; STRATUM_MAX_LINE_BYTES];
+ client.write_all(&line).await.unwrap();
+ client.write_all(b"\n").await.unwrap();
+
+ let read = reader.read_line().await.unwrap().unwrap();
+
+ assert_eq!(read.len(), STRATUM_MAX_LINE_BYTES);
+ }
+
+ #[tokio::test]
+ async fn stratum_line_reader_rejects_oversized_line() {
+ let (mut client, server) = tokio::io::duplex(STRATUM_MAX_LINE_BYTES + 2);
+ let mut reader = StratumLineReader::new(server);
+ let line = vec![b'a'; STRATUM_MAX_LINE_BYTES + 1];
+ client.write_all(&line).await.unwrap();
+ client.write_all(b"\n").await.unwrap();
+
+ let error = reader.read_line().await.unwrap_err();
+
+ assert!(error.to_string().contains("Stratum request exceeds"));
+ }
+
+ #[test]
+ fn stratum_job_cache_prunes_oldest_jobs() {
+ let mut jobs = std::collections::BTreeMap::new();
+ for id in 1..=STRATUM_MAX_JOBS_PER_SESSION + 2 {
+ insert_bounded_job(&mut jobs, id.to_string(), dummy_job());
+ }
+
+ assert_eq!(jobs.len(), STRATUM_MAX_JOBS_PER_SESSION);
+ assert!(!jobs.contains_key("1"));
+ assert!(!jobs.contains_key("2"));
+ assert!(jobs.contains_key("3"));
+ }
+
+ #[test]
+ fn stratum_session_limiter_enforces_global_cap() {
+ let limiter = StratumSessionLimiter::new(STRATUM_MAX_SESSIONS);
+ let permits = (0..STRATUM_MAX_SESSIONS)
+ .map(|_| limiter.try_acquire().expect("permit should be available"))
+ .collect::<Vec<_>>();
+
+ assert!(limiter.try_acquire().is_none());
+ drop(permits);
+ assert!(limiter.try_acquire().is_some());
+ }
+}