Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .kiro/specs/tcp-support/tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -309,7 +309,7 @@ Build production-credible TCP support for dpdk-stdlib-rust, providing drop-in re
- Implement `Drop` for DpdkTcpStream: decrement app_refcount, send Close on last handle
- _Requirements: 9.4, 9.5, 9.8, 9.11, 9.13, 9.14_

- [ ] 8.4 Implement TcpStream public API
- [x] 8.4 Implement TcpStream public API
- Implement `TcpStream` with enum `Inner { Dpdk(DpdkTcpStream), Std(std::net::TcpStream) }`
- Implement `connect<A: ToSocketAddrs>` with v4/v6 dispatch (v4 → DPDK, v6 → kernel fallback)
- Implement `shutdown(how: Shutdown)`, `peer_addr()`, `local_addr()`
Expand All @@ -320,7 +320,7 @@ Build production-credible TCP support for dpdk-stdlib-rust, providing drop-in re
- Implement `peek(buf)` — non-destructive ring read
- _Requirements: 9.1, 9.4, 9.5, 9.6, 9.7, 9.8, 9.10, 9.11, 9.12_

- [ ] 8.5 Implement TcpListener public API
- [x] 8.5 Implement TcpListener public API
- Implement `TcpListener` with enum `Inner { Dpdk(DpdkTcpListener), Std(std::net::TcpListener) }`
- Implement `bind<A: ToSocketAddrs>` with v4/v6 dispatch
- Implement `accept() -> io::Result<(TcpStream, SocketAddr)>` — via oneshot to engine
Expand Down
2 changes: 1 addition & 1 deletion ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,7 @@ Seven property-based tests covering the full engine: (1) state machine validity,
`TcpStream` enum `Inner { Dpdk(DpdkTcpStream), Std(std::net::TcpStream) }` with full `std::net::TcpStream` surface: `connect<A: ToSocketAddrs>` (v4 → DPDK, v6 → kernel fallback), `shutdown`, `peer_addr`, `local_addr`, `set_read_timeout`, `set_write_timeout`, `read_timeout`, `write_timeout`, `set_nodelay`, `nodelay`, `set_ttl`, `ttl`, `set_linger`, `linger`, `set_nonblocking`, `take_error`, `peek`, `try_clone` (Unsupported on DPDK arm). `impl Read for &TcpStream` / `impl Write for &TcpStream` (serialized via read_mutex/write_mutex). `TcpListener` enum with `bind`, `accept() -> (TcpStream, SocketAddr)`, `local_addr`, `set_ttl`, `incoming`. (~400 LOC)

- Spec: `.kiro/specs/tcp-support/` · tasks `8.4`, `8.5`
- [ ] Complete · PR: —
- [x] Complete · PR: #92

---

Expand Down
78 changes: 78 additions & 0 deletions docs/perf-test-log.md
Original file line number Diff line number Diff line change
Expand Up @@ -4956,3 +4956,81 @@ The IPv6 UDP checksum validation adds an IPv6 parse fallback path to `process_fr
**tokio-dpdk at 350K PPS, 64B**: 319,238 RX (8.8% drop) — slightly improved from Run #30's 304,177 (13.1%). Async overhead pattern unchanged.

**Conclusion**: TCP sync socket implementation has zero impact on UDP datapath performance, as expected (separate crate, no shared hot-path code).

---

## Run #32: TCP Sync Socket — TcpStream and TcpListener Public API

| Field | Value |
|-------|-------|
| **Date** | 2026-06-18 |
| **Git Hash** | `42a1497` |
| **Branch** | `agent/tcp-public-api` |
| **PR** | [#92](https://github.com/gspivey/dpdk-stdlib-rust/pull/92) |
| **GH Actions Run** | [27773053334](https://github.com/gspivey/dpdk-stdlib-rust/actions/runs/27773053334) |
| **Instance Type** | c6in.xlarge (4 vCPU, 6.25 Gbps baseline / 30 Gbps burst) |
| **Traffic Generator** | TRex |

### Changes Since Run #31

1. **`42a1497` — TCP sync socket: TcpStream and TcpListener public API (tasks 8.4, 8.5).** Adds `TcpStream` enum wrapper with v4→DPDK / v6→kernel dispatch and full `std::net::TcpStream` surface (connect, connect_timeout, shutdown, peer_addr, local_addr, timeouts, nodelay, ttl, linger, set_nonblocking, take_error, peek, try_clone). Adds `TcpListener` with bind, accept, local_addr, set_ttl, incoming. Adds `peek()` to `SpscByteRing`. Adds `TcpContext` + `init_tcp_context()` for process-wide engine bootstrapping.

### Results: Hardware (TRex)

#### 64-byte packets

| Target PPS | rust-dpdk RX | Drop | Kernel RX | Drop | native-dpdk RX | Drop |
|-----------|-------------|------|----------|------|---------------|------|
| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% |
| 140,000 | 139,000 | 0.7% | 138,966 | 0.7% | 140,000 | 0.0% |
| 350,000 | 348,997 | 0.3% | 348,961 | 0.3% | 349,963 | 0.0% |
| 700,000 | 698,338 | 0.2% | 612,429 | 12.5% | 699,621 | 0.1% |

#### 512-byte packets

| Target PPS | rust-dpdk RX | Drop | Kernel RX | Drop | native-dpdk RX | Drop |
|-----------|-------------|------|----------|------|---------------|------|
| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% |
| 140,000 | 139,000 | 0.7% | 138,995 | 0.7% | 140,000 | 0.0% |
| 350,000 | 348,721 | 0.4% | 348,887 | 0.3% | 350,000 | 0.0% |
| 700,000 | 697,424 | 0.4% | 530,897 | 24.2% | 699,773 | 0.0% |

#### 1400-byte packets (near MTU)

| Target PPS | rust-dpdk RX | Drop | Kernel RX | Drop | native-dpdk RX | Drop |
|-----------|-------------|------|----------|------|---------------|------|
| 70,000 | 69,000 | 1.4% | 69,000 | 1.4% | 70,000 | 0.0% |
| 140,000 | 138,965 | 0.7% | 138,995 | 0.7% | 140,000 | 0.0% |
| 350,000 | 348,999 | 0.3% | 348,676 | 0.4% | 349,900 | 0.0% |
| 700,000 | 475,762 | 0.2% | 435,708 | 8.6% | 476,507 | 0.0% |

#### 8500-byte packets (jumbo)

| Target PPS | rust-dpdk RX | Drop | Kernel RX | Drop | native-dpdk RX | Drop |
|-----------|-------------|------|----------|------|---------------|------|
| 70,000 | 68,995 | 1.4% | 36,429 | 48.0% | 69,994 | 0.0% |
| 140,000 | 77,727 | 0.8% | 76,576 | 2.3% | 77,163 | 1.5% |
| 350,000 | 70,632 | 9.8% | 71,249 | 9.1% | 75,230 | 4.0% |

#### tokio-dpdk (async compat layer)

| Target PPS | tokio-dpdk RX | Drop |
|-----------|--------------|------|
| 70,000 | 69,000 | 1.4% |
| 140,000 | 138,997 | 0.7% |
| 350,000 | 310,526 | 11.3% |
| 700,000 | 307,918 | 56.0% |

### Analysis

**No performance regression from TcpStream/TcpListener public API changes.** This PR adds public API wrapper types to `dpdk-stdlib-tcp` — a separate crate from the UDP datapath with no shared hot-path code.

**rust-dpdk at 700K PPS, 64B**: 698,338 RX (0.2% drop) — consistent with Run #31's 698,965 (0.1%). Within normal ENA variance.

**rust-dpdk at 700K PPS, 512B**: 697,424 RX (0.4% drop) — consistent with Run #31's 698,847 (0.2%). Near-zero drop at line rate.

**rust-dpdk at 700K PPS, 1400B**: 475,762 RX (0.2% drop at TX-capped ~476K) — ENA bandwidth ceiling reached, matching native-dpdk's 476,507.

**tokio-dpdk at 350K PPS, 64B**: 310,526 RX (11.3% drop) — consistent with Run #31's 319,238 (8.8%). Async overhead pattern unchanged.

**Conclusion**: TcpStream/TcpListener public API implementation has zero impact on UDP datapath performance, as expected (separate crate, no shared hot-path code).
6 changes: 6 additions & 0 deletions dpdk-stdlib-tcp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,19 @@ pub mod seq;
pub mod state;
pub mod stream;
pub mod tcb;
pub mod tcp_listener;
pub mod tcp_stream;
pub mod timer;

// Re-export codec public API at crate root for convenience.
pub use codec::{
build_tcp_frame, build_tcp_packet, compute_mss, parse_tcp_packet, tcp_checksum,
};

// Re-export public socket API types.
pub use tcp_stream::{TcpStream, TcpContext, init_tcp_context};
pub use tcp_listener::{TcpListener, Incoming};

// --- Constants ---

/// Maximum TCP payload for IPv4 (MTU 1500 - 20 IPv4 - 20 TCP).
Expand Down
32 changes: 32 additions & 0 deletions dpdk-stdlib-tcp/src/ring.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,38 @@ impl SpscByteRing {
pub fn is_empty(&self) -> bool {
self.available_read() == 0
}

/// Peek at available bytes without advancing the read pointer.
/// Returns the number of bytes copied into `buf`.
pub fn peek(&self, buf: &mut [u8]) -> usize {
let tail = self.tail.load(Ordering::Relaxed);
let head = self.head.load(Ordering::Acquire);
let available = head.wrapping_sub(tail);
let n = buf.len().min(available);
if n == 0 {
return 0;
}

let mask = self.capacity - 1;
let start = tail & mask;
let first_chunk = n.min(self.capacity - start);

unsafe {
std::ptr::copy_nonoverlapping(
self.buf.as_ptr().add(start),
buf.as_mut_ptr(),
first_chunk,
);
if first_chunk < n {
std::ptr::copy_nonoverlapping(
self.buf.as_ptr(),
buf.as_mut_ptr().add(first_chunk),
n - first_chunk,
);
}
}
n
}
}

// Safety: SpscByteRing is Send+Sync because atomic operations guard head/tail,
Expand Down
175 changes: 175 additions & 0 deletions dpdk-stdlib-tcp/src/tcp_listener.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
//! Public `TcpListener` API — drop-in replacement for `std::net::TcpListener`.
//!
//! Dispatches IPv4 to the DPDK engine path and IPv6 to kernel fallback.

use std::io;
use std::net::{self, SocketAddr, ToSocketAddrs};

use crate::contract::{
CommandSender, EngineCommand, oneshot_channel,
};
use crate::stream::DpdkTcpStream;
use crate::tcp_stream::{get_tcp_context, resolve_addr, TcpStream};

/// A TCP socket server, either backed by DPDK (IPv4) or the kernel (IPv6 fallback).
///
/// Provides the full `std::net::TcpListener` API surface.
pub struct TcpListener {
inner: ListenerInner,
}

enum ListenerInner {
Dpdk(DpdkTcpListener),
Std(net::TcpListener),
}

struct DpdkTcpListener {
addr: SocketAddr,
cmd_tx: CommandSender,
}

impl TcpListener {
/// Creates a new `TcpListener` bound to the specified address.
///
/// IPv4 addresses use the DPDK path; IPv6 falls back to the kernel.
pub fn bind<A: ToSocketAddrs>(addr: A) -> io::Result<TcpListener> {
let addr = resolve_addr(addr)?;
match addr {
SocketAddr::V4(_) => {
let ctx = get_tcp_context()?;
let (resp_tx, resp_rx) = oneshot_channel();
ctx.cmd_tx
.send(EngineCommand::Listen {
addr,
backlog: 128,
response: resp_tx,
})
.map_err(|_| {
io::Error::new(io::ErrorKind::BrokenPipe, "engine channel closed")
})?;

let result = resp_rx.recv();
match result {
Ok(()) => Ok(TcpListener {
inner: ListenerInner::Dpdk(DpdkTcpListener {
addr,
cmd_tx: ctx.cmd_tx.clone(),
}),
}),
Err(e) => Err(e.into()),
}
}
SocketAddr::V6(_) => {
let listener = net::TcpListener::bind(addr)?;
Ok(TcpListener {
inner: ListenerInner::Std(listener),
})
}
}
}

/// Accept a new incoming connection.
///
/// Blocks until a connection is available.
pub fn accept(&self) -> io::Result<(TcpStream, SocketAddr)> {
match &self.inner {
ListenerInner::Dpdk(listener) => {
let (resp_tx, resp_rx) = oneshot_channel();
listener
.cmd_tx
.send(EngineCommand::Accept {
listen_addr: listener.addr,
response: resp_tx,
})
.map_err(|_| {
io::Error::new(io::ErrorKind::BrokenPipe, "engine channel closed")
})?;

let result = resp_rx.recv();
match result {
Ok((key, handle)) => {
let remote = key.remote;
let stream = DpdkTcpStream::new(handle, key);
Ok((TcpStream::from_dpdk(stream), remote))
}
Err(e) => Err(e.into()),
}
}
ListenerInner::Std(listener) => {
let (stream, addr) = listener.accept()?;
Ok((TcpStream::from_std(stream), addr))
}
}
}

/// Returns the local socket address.
pub fn local_addr(&self) -> io::Result<SocketAddr> {
match &self.inner {
ListenerInner::Dpdk(listener) => Ok(listener.addr),
ListenerInner::Std(listener) => listener.local_addr(),
}
}

/// Sets the TTL value for this listener's socket.
pub fn set_ttl(&self, ttl: u32) -> io::Result<()> {
match &self.inner {
ListenerInner::Dpdk(_) => {
// TTL on a listener is a no-op for DPDK (applies to accepted streams).
Ok(())
}
ListenerInner::Std(listener) => listener.set_ttl(ttl),
}
}

/// Gets the TTL value.
pub fn ttl(&self) -> io::Result<u32> {
match &self.inner {
ListenerInner::Dpdk(_) => Ok(64),
ListenerInner::Std(listener) => listener.ttl(),
}
}

/// Returns an iterator over incoming connections.
pub fn incoming(&self) -> Incoming<'_> {
Incoming { listener: self }
}
}

/// An iterator over incoming TCP connections on a `TcpListener`.
pub struct Incoming<'a> {
listener: &'a TcpListener,
}

impl<'a> Iterator for Incoming<'a> {
type Item = io::Result<TcpStream>;

fn next(&mut self) -> Option<Self::Item> {
Some(self.listener.accept().map(|(stream, _)| stream))
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn tcp_listener_bind_v6_falls_back_to_kernel() {
// Binding to [::1]:0 should use the kernel path (Std variant).
// This may fail if IPv6 is disabled, which is fine — we're testing dispatch.
let result = TcpListener::bind("[::1]:0");
// Either succeeds (uses kernel) or fails with a kernel error — both are correct.
if let Ok(listener) = result {
let addr = listener.local_addr().unwrap();
assert!(addr.is_ipv6());
}
}

#[test]
fn tcp_listener_bind_v4_without_context_returns_error() {
// Without TCP context initialized, V4 bind should fail.
// (Context may or may not be initialized from other tests, so
// we just verify the function doesn't panic.)
let _result = TcpListener::bind("10.0.0.1:0");
// Result depends on whether TCP_CONTEXT is initialized.
}
}
Loading
Loading