diff --git a/ROADMAP.md b/ROADMAP.md index 8d512af..51c75b4 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -140,7 +140,7 @@ Create `dpdk-stdlib-tcp/` crate (depends on `dpdk-stdlib-net` + `dpdk-stdlib`, C `build_tcp_frame(params) -> Vec` (Eth + IPv4 + TCP; SYN/SYN-ACK frames include MSS, WScale, SACK-Perm, Timestamps). `tcp_checksum` with parameterized pseudo-header. `compute_mss(mtu, ip_hdr_len)`. `parse_tcp_packet` (validate data-offset ≥5, parse all options). `build_tcp_packet(mbuf, params)` (zero-copy DPDK path, byte-identical to `build_tcp_frame`). 7 property tests: round-trip, Mbuf equivalence, invalid frame rejection, SYN required options, MSS bound, checksum flip fails, sequence transitivity. (~500 LOC) - Spec: `.kiro/specs/tcp-support/` · tasks `3.7`, `3.8`, `3.9`, `3.10`, `3.11` -- [ ] Complete · PR: — +- [x] Complete · PR: #79 --- @@ -149,7 +149,7 @@ Create `dpdk-stdlib-tcp/` crate (depends on `dpdk-stdlib-net` + `dpdk-stdlib`, C `ConnectionHandle` (rx_ring/tx_ring SpscByteRing, AtomicU8 state, AtomicBool eof, Mutex error, condvar + notify_lock, AtomicWaker × 2, app_refcount, cmd_tx, key, linger). `EngineCommand` enum. `SocketOption` enum. `CommandSender` (wraps mpsc::Sender + Arc; every send signals engine_wakeup). `OneshotSender/Receiver`. `TcpState` (11 states). `FourTuple`. `SystemClock` + `MockClock` (with `advance()`). `IsnGenerator` (RFC 6528: 128-bit per-boot secret via `getrandom`, SipHash-2-4 of FourTuple, M = elapsed µs / 4). (~400 LOC) - Spec: `.kiro/specs/tcp-support/` · tasks `5.1`, `5.2`, `5.3`, `5.4` -- [ ] Complete · PR: — +- [x] Complete · PR: #80 --- diff --git a/docs/perf-test-log.md b/docs/perf-test-log.md index 0f1be81..e0b2b1b 100644 --- a/docs/perf-test-log.md +++ b/docs/perf-test-log.md @@ -6,6 +6,77 @@ Each entry captures the git context, test configuration, results, and analysis. **Standard benchmarks** (include in every run entry): 1. **Hardware PPS** — TRex on c6in.xlarge (measures NIC + DPDK + application stack) +## Run #42: dpdk-stdlib-tcp Contract Types, TcpState, Clock, IsnGenerator — No Regression + +| Field | Value | +|-------|-------| +| **Date** | 2026-06-13 | +| **Git Hash** | `22c0cd5` | +| **Branch** | `agent/tcp-contract-types-state-clock-isn` | +| **PR** | [#80](https://github.com/gspivey/dpdk-stdlib-rust/pull/80) | +| **GH Actions Run** | [27476064107](https://github.com/gspivey/dpdk-stdlib-rust/actions/runs/27476064107) | +| **Instance Type** | c6in.xlarge (4 vCPU, 6.25 Gbps baseline / 30 Gbps burst) | +| **Traffic Generator** | TRex | + +### Changes Since Run #41 + +1. **`22c0cd5` — dpdk-stdlib-tcp contract types, TcpState, Clock, and IsnGenerator.** New modules: `state.rs` (TcpState enum + FourTuple), `clock.rs` (Clock trait + SystemClock + MockClock), `isn.rs` (IsnGenerator per RFC 6528), `contract.rs` (ConnectionHandle, EngineCommand, SocketOption, CommandSender, OneshotSender/Receiver, EngineWakeup, AtomicWaker). Zero changes to any existing data-path crate. + +### Results: Hardware (TRex) + +#### 64-byte packets + +| Target PPS | plain-rust RX | Drop | rust-dpdk RX | Drop | tokio-dpdk RX | Drop | native-dpdk RX | Drop | +|-----------|--------------|------|-------------|------|--------------|------|---------------|------| +| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% | +| 140,000 | 138,997 | 0.7% | 139,000 | 0.7% | 139,000 | 0.7% | 140,000 | 0.0% | +| 350,000 | 348,985 | 0.3% | 348,990 | 0.3% | 311,544 | 11.0% | 349,973 | 0.0% | +| 700,000 | 381,095 | 45.6% | 658,996 | 5.9% | 312,244 | 55.4% | 656,738 | 6.2% | + +#### 512-byte packets + +| Target PPS | plain-rust RX | Drop | rust-dpdk RX | Drop | tokio-dpdk RX | Drop | native-dpdk RX | Drop | +|-----------|--------------|------|-------------|------|--------------|------|---------------|------| +| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% | +| 140,000 | 138,978 | 0.7% | 139,000 | 0.7% | 139,000 | 0.7% | 140,000 | 0.0% | +| 350,000 | 348,815 | 0.3% | 348,990 | 0.3% | 231,845 | 33.8% | 349,978 | 0.0% | +| 700,000 | 398,868 | 43.0% | 620,084 | 11.4% | 231,659 | 66.9% | 598,345 | 14.5% | + +#### 1400-byte packets (near MTU) + +| Target PPS | plain-rust RX | Drop | rust-dpdk RX | Drop | tokio-dpdk RX | Drop | native-dpdk RX | Drop | +|-----------|--------------|------|-------------|------|--------------|------|---------------|------| +| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% | +| 140,000 | 138,990 | 0.7% | 139,000 | 0.7% | 139,000 | 0.7% | 140,000 | 0.0% | +| 350,000 | 348,770 | 0.4% | 349,000 | 0.3% | 150,532 | 57.0% | 349,998 | 0.0% | +| 700,000 | 447,288 | 6.0%* | 466,142 | 2.0%* | 160,320 | 66.3%* | 460,415 | 3.1%* | + +\* TX capped at ~476K pps (ENA line-rate limit at 1400B) + +#### 8500-byte packets (jumbo) + +| Target PPS | plain-rust RX | Drop | rust-dpdk RX | Drop | tokio-dpdk RX | Drop | native-dpdk RX | Drop | +|-----------|--------------|------|-------------|------|--------------|------|---------------|------| +| 70,000 | 38,549 | 44.9% | 69,000 | 1.4% | 55,327 | 21.0% | 70,000 | 0.0% | +| 140,000 | 77,780 | 0.7%* | 77,724 | 0.8%* | 58,906 | 24.8%* | 78,266 | 0.0%* | +| 350,000 | 77,905 | 0.5%* | 77,900 | 0.6%* | 59,005 | 24.7%* | 77,960 | 0.4%* | + +\* TX capped at ~78K pps (30 Gbps ENA burst limit at 8500B) + +### Analysis + +**No performance regression.** This PR adds new contract types, state machine enum, clock abstraction, and ISN generator to `dpdk-stdlib-tcp`. Zero changes to any existing data-path crate or networking code. + +**rust-dpdk at 700K PPS, 64B**: 658,996 RX (5.9% drop) — consistent with Run #41's 659,356 (5.8% drop). Within normal ENA variance. + +**native-dpdk at 700K PPS, 64B**: 656,738 RX (6.2% drop) — consistent with Run #41's 674,638 (3.6% drop). Normal saturation-point variance. + +**rust-dpdk vs native-dpdk gap**: Near parity at 350K and below across all packet sizes. At 700K saturation, rust-dpdk tracks native-dpdk closely. + +**Conclusion**: Adding TCP contract types is performance-neutral as expected. + +--- + ## Run #41: dpdk-stdlib-tcp Crate Skeleton and Codec Types — No Regression | Field | Value | diff --git a/dpdk-stdlib-tcp/Cargo.toml b/dpdk-stdlib-tcp/Cargo.toml index 9814cb3..850145b 100644 --- a/dpdk-stdlib-tcp/Cargo.toml +++ b/dpdk-stdlib-tcp/Cargo.toml @@ -12,6 +12,8 @@ description = "DPDK-accelerated TCP stack for dpdk-stdlib-rust" dpdk-stdlib-net = { version = "0.2.0", path = "../dpdk-stdlib-net" } dpdk = { version = "0.2.0", path = "../dpdk", package = "dpdk-stdlib" } thiserror = { workspace = true } +getrandom = "0.2" +siphasher = "1" [dev-dependencies] proptest = "1.4" diff --git a/dpdk-stdlib-tcp/src/clock.rs b/dpdk-stdlib-tcp/src/clock.rs new file mode 100644 index 0000000..abe12c7 --- /dev/null +++ b/dpdk-stdlib-tcp/src/clock.rs @@ -0,0 +1,100 @@ +//! Clock abstraction for deterministic testing of timer-driven behavior. + +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +/// Clock trait for injectable time source. +pub trait Clock: Send + Sync { + fn now(&self) -> Instant; +} + +/// System clock delegating to `std::time::Instant::now()`. +pub struct SystemClock; + +impl Clock for SystemClock { + #[inline] + fn now(&self) -> Instant { + Instant::now() + } +} + +/// Mock clock for deterministic testing. Advances only via explicit calls. +pub struct MockClock { + inner: Arc>, +} + +impl MockClock { + pub fn new() -> Self { + Self { + inner: Arc::new(Mutex::new(Instant::now())), + } + } + + /// Create with a specific starting instant. + pub fn with_instant(instant: Instant) -> Self { + Self { + inner: Arc::new(Mutex::new(instant)), + } + } + + /// Advance the clock by the given duration. + pub fn advance(&self, duration: Duration) { + let mut t = self.inner.lock().unwrap(); + *t += duration; + } + + /// Set the clock to an exact instant. + pub fn set(&self, instant: Instant) { + let mut t = self.inner.lock().unwrap(); + *t = instant; + } +} + +impl Clock for MockClock { + fn now(&self) -> Instant { + *self.inner.lock().unwrap() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn system_clock_monotonic() { + let clock = SystemClock; + let t1 = clock.now(); + let t2 = clock.now(); + assert!(t2 >= t1); + } + + #[test] + fn mock_clock_advance() { + let clock = MockClock::new(); + let t1 = clock.now(); + clock.advance(Duration::from_secs(5)); + let t2 = clock.now(); + assert_eq!(t2 - t1, Duration::from_secs(5)); + } + + #[test] + fn mock_clock_set() { + let clock = MockClock::new(); + let start = clock.now(); + clock.advance(Duration::from_secs(10)); + let after_advance = clock.now(); + assert_eq!(after_advance - start, Duration::from_secs(10)); + clock.set(start); + assert_eq!(clock.now(), start); + } + + #[test] + fn mock_clock_multiple_advances() { + let clock = MockClock::new(); + let start = clock.now(); + clock.advance(Duration::from_millis(100)); + clock.advance(Duration::from_millis(200)); + clock.advance(Duration::from_millis(300)); + assert_eq!(clock.now() - start, Duration::from_millis(600)); + } +} diff --git a/dpdk-stdlib-tcp/src/contract.rs b/dpdk-stdlib-tcp/src/contract.rs new file mode 100644 index 0000000..b0dd1bf --- /dev/null +++ b/dpdk-stdlib-tcp/src/contract.rs @@ -0,0 +1,467 @@ +//! App↔engine contract types. +//! +//! Pure data/enum types with no dependency back on TcpEngine. +//! These define the shared interface between app threads and the engine thread. + +use std::net::{Shutdown, SocketAddr}; +use std::sync::atomic::{AtomicBool, AtomicU8, AtomicUsize, Ordering}; +use std::sync::mpsc::{self, SendError}; +use std::sync::{Arc, Condvar, Mutex}; +use std::time::Duration; + +use crate::error::TcpError; +use crate::ring::SpscByteRing; +use crate::state::{FourTuple, TcpState}; + +// --- EngineWakeup --- + +/// Engine wakeup signal. Uses AtomicBool + Condvar (portable/test). +pub struct EngineWakeup { + flag: AtomicBool, + condvar: Condvar, + mutex: Mutex<()>, +} + +impl EngineWakeup { + pub fn new() -> Self { + Self { + flag: AtomicBool::new(false), + condvar: Condvar::new(), + mutex: Mutex::new(()), + } + } + + /// Signal the engine to wake up. + pub fn signal(&self) { + self.flag.store(true, Ordering::Release); + let _guard = self.mutex.lock().unwrap(); + self.condvar.notify_one(); + } + + /// Wait for a signal (with timeout). Returns true if signaled. + pub fn wait(&self, timeout: Duration) -> bool { + let guard = self.mutex.lock().unwrap(); + if self.flag.swap(false, Ordering::AcqRel) { + return true; + } + let (_guard, result) = self.condvar.wait_timeout(guard, timeout).unwrap(); + self.flag.swap(false, Ordering::AcqRel) || !result.timed_out() + } + + /// Check and clear the signal without blocking. + pub fn try_recv(&self) -> bool { + self.flag.swap(false, Ordering::AcqRel) + } +} + +// --- CommandSender --- + +/// Wrapper around mpsc::Sender that also signals engine_wakeup on every send. +#[derive(Clone)] +pub struct CommandSender { + inner: mpsc::Sender, + wakeup: Arc, +} + +impl CommandSender { + pub fn new(sender: mpsc::Sender, wakeup: Arc) -> Self { + Self { + inner: sender, + wakeup, + } + } + + pub fn send(&self, cmd: EngineCommand) -> Result<(), SendError> { + let result = self.inner.send(cmd); + self.wakeup.signal(); + result + } +} + +// --- Oneshot channel --- + +/// Oneshot sender (no tokio dependency). +pub struct OneshotSender { + inner: Arc<(Mutex>, Condvar)>, +} + +/// Oneshot receiver (no tokio dependency). +pub struct OneshotReceiver { + inner: Arc<(Mutex>, Condvar)>, +} + +/// Create a oneshot channel pair. +pub fn oneshot_channel() -> (OneshotSender, OneshotReceiver) { + let inner = Arc::new((Mutex::new(None), Condvar::new())); + ( + OneshotSender { + inner: inner.clone(), + }, + OneshotReceiver { inner }, + ) +} + +impl OneshotSender { + /// Send a value, waking the receiver. + pub fn send(self, value: T) { + let (lock, condvar) = &*self.inner; + let mut slot = lock.lock().unwrap(); + *slot = Some(value); + condvar.notify_one(); + } +} + +impl OneshotReceiver { + /// Block until a value is received. + pub fn recv(self) -> T { + let (lock, condvar) = &*self.inner; + let mut slot = lock.lock().unwrap(); + loop { + if let Some(val) = slot.take() { + return val; + } + slot = condvar.wait(slot).unwrap(); + } + } + + /// Block until a value is received, with timeout. + pub fn recv_timeout(self, timeout: Duration) -> Option { + let (lock, condvar) = &*self.inner; + let mut slot = lock.lock().unwrap(); + let deadline = std::time::Instant::now() + timeout; + loop { + if let Some(val) = slot.take() { + return Some(val); + } + let remaining = deadline.saturating_duration_since(std::time::Instant::now()); + if remaining.is_zero() { + return None; + } + let (new_slot, _) = condvar.wait_timeout(slot, remaining).unwrap(); + slot = new_slot; + } + } +} + +// --- SocketOption --- + +/// Keepalive configuration. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct KeepaliveConfig { + pub idle: Duration, + pub interval: Duration, + pub count: u32, +} + +/// Socket options routed through EngineCommand::SetOption. +#[derive(Debug, Clone)] +pub enum SocketOption { + Nodelay(bool), + Keepalive(Option), + Linger(Option), + RecvBufSize(usize), + SendBufSize(usize), + ReuseAddr(bool), + Ttl(u8), + ReadTimeout(Option), + WriteTimeout(Option), + Nonblocking(bool), +} + +// --- EngineCommand --- + +/// Commands sent from app threads to the engine thread via mpsc channel. +pub enum EngineCommand { + Connect { + local: SocketAddr, + remote: SocketAddr, + src_mac: [u8; 6], + dst_mac: [u8; 6], + handle: Arc, + response: OneshotSender>, + }, + Listen { + addr: SocketAddr, + backlog: usize, + response: OneshotSender>, + }, + Accept { + listen_addr: SocketAddr, + response: OneshotSender), TcpError>>, + }, + Shutdown { + key: FourTuple, + how: Shutdown, + }, + SetOption { + key: FourTuple, + option: SocketOption, + }, + Close { + key: FourTuple, + linger: Option, + }, +} + +// --- AtomicWaker (minimal, no tokio) --- + +/// Minimal atomic waker for async integration. +/// Stores a single waker that can be registered and woken. +pub struct AtomicWaker { + waker: Mutex>, +} + +impl AtomicWaker { + pub fn new() -> Self { + Self { + waker: Mutex::new(None), + } + } + + /// Register a waker (replaces any previous waker). + pub fn register(&self, waker: &std::task::Waker) { + let mut slot = self.waker.lock().unwrap(); + *slot = Some(waker.clone()); + } + + /// Wake the registered waker (if any). + pub fn wake(&self) { + if let Some(w) = self.waker.lock().unwrap().take() { + w.wake(); + } + } +} + +// --- ConnectionHandle --- + +/// Shared state between app threads and engine thread (via Arc). +pub struct ConnectionHandle { + /// Received data: engine writes, app reads. + pub rx_ring: SpscByteRing, + /// Send data: app writes, engine reads. + pub tx_ring: SpscByteRing, + + /// Current TCP state (engine updates with Release). + pub state: AtomicU8, + /// Explicit EOF flag — set after final rx bytes enqueued on FIN. + pub eof: AtomicBool, + /// Latched connection error — sticky (peek/clone, never take). + pub error: Mutex>, + + /// Condvar + notify_lock for blocking wake (recheck-under-lock). + pub condvar: Condvar, + pub notify_lock: Mutex<()>, + + /// Async wakers for read/write. + pub read_waker: AtomicWaker, + pub write_waker: AtomicWaker, + + /// Serialize concurrent (&stream).read() calls. + pub read_mutex: Mutex<()>, + /// Serialize concurrent (&stream).write() calls. + pub write_mutex: Mutex<()>, + + /// Number of live app handles (TcpStream + split halves). + pub app_refcount: AtomicUsize, + /// Command sender for Close on last-handle-drop. + pub cmd_tx: CommandSender, + /// Connection key. + pub key: FourTuple, + /// SO_LINGER setting. + pub linger: Mutex>, +} + +impl ConnectionHandle { + /// Create a new connection handle with default buffer sizes. + pub fn new( + rx_capacity: usize, + tx_capacity: usize, + cmd_tx: CommandSender, + key: FourTuple, + ) -> Self { + Self { + rx_ring: SpscByteRing::new(rx_capacity), + tx_ring: SpscByteRing::new(tx_capacity), + state: AtomicU8::new(TcpState::Closed as u8), + eof: AtomicBool::new(false), + error: Mutex::new(None), + condvar: Condvar::new(), + notify_lock: Mutex::new(()), + read_waker: AtomicWaker::new(), + write_waker: AtomicWaker::new(), + read_mutex: Mutex::new(()), + write_mutex: Mutex::new(()), + app_refcount: AtomicUsize::new(1), + cmd_tx, + key, + linger: Mutex::new(None), + } + } + + /// Get the current TCP state. + pub fn tcp_state(&self) -> TcpState { + TcpState::from_u8(self.state.load(Ordering::Acquire)).unwrap_or(TcpState::Closed) + } + + /// Set the TCP state (called by engine). + pub fn set_state(&self, state: TcpState) { + self.state.store(state as u8, Ordering::Release); + } + + /// Latch a connection error (sticky — never cleared). + pub fn latch_error(&self, err: TcpError) { + let mut slot = self.error.lock().unwrap(); + if slot.is_none() { + *slot = Some(err); + } + } + + /// Peek at the latched error (clone, never take). + pub fn peek_error(&self) -> Option { + self.error.lock().unwrap().clone() + } + + /// Notify blocked readers/writers and async wakers. + pub fn notify_all(&self) { + let _guard = self.notify_lock.lock().unwrap(); + self.condvar.notify_all(); + self.read_waker.wake(); + self.write_waker.wake(); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn make_test_handle() -> Arc { + let (tx, _rx) = mpsc::channel(); + let wakeup = Arc::new(EngineWakeup::new()); + let cmd_tx = CommandSender::new(tx, wakeup); + let key = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + Arc::new(ConnectionHandle::new(65536, 65536, cmd_tx, key)) + } + + #[test] + fn handle_initial_state() { + let h = make_test_handle(); + assert_eq!(h.tcp_state(), TcpState::Closed); + assert!(!h.eof.load(Ordering::Acquire)); + assert!(h.peek_error().is_none()); + assert_eq!(h.app_refcount.load(Ordering::Acquire), 1); + } + + #[test] + fn handle_state_transitions() { + let h = make_test_handle(); + h.set_state(TcpState::SynSent); + assert_eq!(h.tcp_state(), TcpState::SynSent); + h.set_state(TcpState::Established); + assert_eq!(h.tcp_state(), TcpState::Established); + } + + #[test] + fn handle_latch_error_sticky() { + let h = make_test_handle(); + h.latch_error(TcpError::ConnectionReset); + assert!(matches!(h.peek_error(), Some(TcpError::ConnectionReset))); + // Second latch doesn't overwrite + h.latch_error(TcpError::TimedOut); + assert!(matches!(h.peek_error(), Some(TcpError::ConnectionReset))); + } + + #[test] + fn oneshot_send_recv() { + let (tx, rx) = oneshot_channel(); + tx.send(42u32); + assert_eq!(rx.recv(), 42); + } + + #[test] + fn oneshot_recv_timeout_success() { + let (tx, rx) = oneshot_channel(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(10)); + tx.send(99u32); + }); + let val = rx.recv_timeout(Duration::from_secs(1)); + assert_eq!(val, Some(99)); + } + + #[test] + fn oneshot_recv_timeout_expires() { + let (_tx, rx) = oneshot_channel::(); + let val = rx.recv_timeout(Duration::from_millis(10)); + assert_eq!(val, None); + } + + #[test] + fn command_sender_signals_wakeup() { + let (tx, _rx) = mpsc::channel(); + let wakeup = Arc::new(EngineWakeup::new()); + let cmd_tx = CommandSender::new(tx, wakeup.clone()); + let key = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + cmd_tx + .send(EngineCommand::Close { + key, + linger: None, + }) + .unwrap(); + // Wakeup should have been signaled + assert!(wakeup.try_recv()); + } + + #[test] + fn engine_wakeup_signal_and_wait() { + let wakeup = Arc::new(EngineWakeup::new()); + let w2 = wakeup.clone(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(10)); + w2.signal(); + }); + assert!(wakeup.wait(Duration::from_secs(1))); + } + + #[test] + fn engine_wakeup_timeout() { + let wakeup = EngineWakeup::new(); + assert!(!wakeup.wait(Duration::from_millis(10))); + } + + #[test] + fn atomic_waker_register_and_wake() { + use std::sync::atomic::AtomicBool; + use std::task::{RawWaker, RawWakerVTable, Waker}; + + static WOKEN: AtomicBool = AtomicBool::new(false); + + fn clone_fn(ptr: *const ()) -> RawWaker { + RawWaker::new(ptr, &VTABLE) + } + fn wake_fn(_: *const ()) { + WOKEN.store(true, Ordering::Release); + } + fn wake_by_ref_fn(_: *const ()) { + WOKEN.store(true, Ordering::Release); + } + fn drop_fn(_: *const ()) {} + + static VTABLE: RawWakerVTable = + RawWakerVTable::new(clone_fn, wake_fn, wake_by_ref_fn, drop_fn); + + let raw = RawWaker::new(std::ptr::null(), &VTABLE); + let waker = unsafe { Waker::from_raw(raw) }; + + let aw = AtomicWaker::new(); + aw.register(&waker); + WOKEN.store(false, Ordering::Release); + aw.wake(); + assert!(WOKEN.load(Ordering::Acquire)); + } +} diff --git a/dpdk-stdlib-tcp/src/isn.rs b/dpdk-stdlib-tcp/src/isn.rs new file mode 100644 index 0000000..113c4d0 --- /dev/null +++ b/dpdk-stdlib-tcp/src/isn.rs @@ -0,0 +1,130 @@ +//! Initial Sequence Number generator per RFC 6528. +//! +//! Uses a per-boot 128-bit secret + SipHash-2-4 of the 4-tuple + M +//! (elapsed microseconds / 4 since boot) to produce unpredictable ISNs. + +use std::hash::Hasher; +use std::time::Instant; + +use siphasher::sip::SipHasher24; + +use crate::clock::Clock; +use crate::seq::SeqNum; +use crate::state::FourTuple; + +/// ISN generator with per-boot secret (RFC 6528). +pub struct IsnGenerator { + secret: [u8; 16], + boot_instant: Instant, +} + +impl IsnGenerator { + /// Create a new ISN generator with a random per-boot secret. + pub fn new(clock: &dyn Clock) -> Self { + let mut secret = [0u8; 16]; + getrandom::getrandom(&mut secret).expect("getrandom failed"); + Self { + secret, + boot_instant: clock.now(), + } + } + + /// Create with a known secret (for deterministic testing). + #[cfg(test)] + pub fn with_secret(secret: [u8; 16], boot_instant: Instant) -> Self { + Self { + secret, + boot_instant, + } + } + + /// Generate an ISN for the given 4-tuple. + /// M = elapsed µs since boot / 4 (wraps ~4.7 hours, fine for ISN). + pub fn generate(&self, four_tuple: &FourTuple, clock: &dyn Clock) -> SeqNum { + let elapsed = clock.now().duration_since(self.boot_instant); + let m = (elapsed.as_micros() / 4) as u32; + + let key_lo = u64::from_le_bytes(self.secret[0..8].try_into().unwrap()); + let key_hi = u64::from_le_bytes(self.secret[8..16].try_into().unwrap()); + let mut hasher = SipHasher24::new_with_keys(key_lo, key_hi); + hasher.write(&four_tuple.to_bytes()); + let hash = hasher.finish() as u32; + + SeqNum(m.wrapping_add(hash)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::clock::MockClock; + use std::time::Duration; + + #[test] + fn isn_deterministic_for_same_inputs() { + let clock = MockClock::new(); + let gen = IsnGenerator::with_secret([1u8; 16], clock.now()); + let ft = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let isn1 = gen.generate(&ft, &clock); + let isn2 = gen.generate(&ft, &clock); + assert_eq!(isn1, isn2); + } + + #[test] + fn isn_differs_for_different_tuples() { + let clock = MockClock::new(); + let gen = IsnGenerator::with_secret([2u8; 16], clock.now()); + let ft1 = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let ft2 = FourTuple { + local: "10.0.0.1:1235".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + assert_ne!(gen.generate(&ft1, &clock), gen.generate(&ft2, &clock)); + } + + #[test] + fn isn_advances_with_time() { + let clock = MockClock::new(); + let gen = IsnGenerator::with_secret([3u8; 16], clock.now()); + let ft = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let isn1 = gen.generate(&ft, &clock); + clock.advance(Duration::from_micros(400)); // M advances by 100 + let isn2 = gen.generate(&ft, &clock); + // ISN should have advanced (M component changed) + assert_ne!(isn1, isn2); + } + + #[test] + fn isn_unpredictable_different_secrets() { + let clock = MockClock::new(); + let gen1 = IsnGenerator::with_secret([4u8; 16], clock.now()); + let gen2 = IsnGenerator::with_secret([5u8; 16], clock.now()); + let ft = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + assert_ne!(gen1.generate(&ft, &clock), gen2.generate(&ft, &clock)); + } + + #[test] + fn isn_new_uses_random_secret() { + let clock = MockClock::new(); + let gen1 = IsnGenerator::new(&clock); + let gen2 = IsnGenerator::new(&clock); + let ft = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + // Overwhelmingly likely to differ (random secrets) + assert_ne!(gen1.generate(&ft, &clock), gen2.generate(&ft, &clock)); + } +} diff --git a/dpdk-stdlib-tcp/src/lib.rs b/dpdk-stdlib-tcp/src/lib.rs index d1884af..337863c 100644 --- a/dpdk-stdlib-tcp/src/lib.rs +++ b/dpdk-stdlib-tcp/src/lib.rs @@ -5,10 +5,14 @@ //! //! Depends on `dpdk-stdlib-net` for `PacketBackend` — does NOT depend on `dpdk-udp`. +pub mod clock; pub mod codec; +pub mod contract; pub mod error; +pub mod isn; pub mod ring; pub mod seq; +pub mod state; // Re-export codec public API at crate root for convenience. pub use codec::{ diff --git a/dpdk-stdlib-tcp/src/state.rs b/dpdk-stdlib-tcp/src/state.rs new file mode 100644 index 0000000..9d22eb0 --- /dev/null +++ b/dpdk-stdlib-tcp/src/state.rs @@ -0,0 +1,136 @@ +//! TCP state machine states and connection 4-tuple. + +use std::net::SocketAddr; + +/// The 11 states of the TCP state machine (RFC 9293). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[repr(u8)] +pub enum TcpState { + Closed = 0, + Listen = 1, + SynSent = 2, + SynReceived = 3, + Established = 4, + FinWait1 = 5, + FinWait2 = 6, + CloseWait = 7, + Closing = 8, + LastAck = 9, + TimeWait = 10, +} + +impl TcpState { + /// Convert from raw u8 (e.g. from AtomicU8). + pub fn from_u8(v: u8) -> Option { + match v { + 0 => Some(Self::Closed), + 1 => Some(Self::Listen), + 2 => Some(Self::SynSent), + 3 => Some(Self::SynReceived), + 4 => Some(Self::Established), + 5 => Some(Self::FinWait1), + 6 => Some(Self::FinWait2), + 7 => Some(Self::CloseWait), + 8 => Some(Self::Closing), + 9 => Some(Self::LastAck), + 10 => Some(Self::TimeWait), + _ => None, + } + } +} + +/// A TCP connection identified by its 4-tuple. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub struct FourTuple { + pub local: SocketAddr, + pub remote: SocketAddr, +} + +impl FourTuple { + /// Serialize to bytes for hashing (ISN generation). + pub fn to_bytes(&self) -> [u8; 12] { + let mut out = [0u8; 12]; + match (self.local, self.remote) { + (SocketAddr::V4(l), SocketAddr::V4(r)) => { + out[0..4].copy_from_slice(&l.ip().octets()); + out[4..6].copy_from_slice(&l.port().to_be_bytes()); + out[6..10].copy_from_slice(&r.ip().octets()); + out[10..12].copy_from_slice(&r.port().to_be_bytes()); + } + _ => {} // IPv6 deferred + } + out + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn tcp_state_roundtrip() { + for v in 0..=10u8 { + let state = TcpState::from_u8(v).unwrap(); + assert_eq!(state as u8, v); + } + assert!(TcpState::from_u8(11).is_none()); + assert!(TcpState::from_u8(255).is_none()); + } + + #[test] + fn tcp_state_all_variants() { + let states = [ + TcpState::Closed, + TcpState::Listen, + TcpState::SynSent, + TcpState::SynReceived, + TcpState::Established, + TcpState::FinWait1, + TcpState::FinWait2, + TcpState::CloseWait, + TcpState::Closing, + TcpState::LastAck, + TcpState::TimeWait, + ]; + assert_eq!(states.len(), 11); + } + + #[test] + fn four_tuple_eq_and_hash() { + use std::collections::HashSet; + let ft1 = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let ft2 = FourTuple { + local: "10.0.0.1:1234".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let ft3 = FourTuple { + local: "10.0.0.1:1235".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + assert_eq!(ft1, ft2); + assert_ne!(ft1, ft3); + let mut set = HashSet::new(); + set.insert(ft1); + assert!(set.contains(&ft2)); + assert!(!set.contains(&ft3)); + } + + #[test] + fn four_tuple_to_bytes_deterministic() { + let ft = FourTuple { + local: "192.168.1.1:5000".parse().unwrap(), + remote: "10.0.0.2:80".parse().unwrap(), + }; + let b1 = ft.to_bytes(); + let b2 = ft.to_bytes(); + assert_eq!(b1, b2); + // Check content + assert_eq!(&b1[0..4], &[192, 168, 1, 1]); + assert_eq!(&b1[4..6], &5000u16.to_be_bytes()); + assert_eq!(&b1[6..10], &[10, 0, 0, 2]); + assert_eq!(&b1[10..12], &80u16.to_be_bytes()); + } +}