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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -130,7 +130,7 @@ clang-format -style=file -i -fallback-style=none <files>

`EngineType` (`include/daqiri/types.h`) is resolved from `(stream_type, engine)`: `raw` defaults to `EngineType::IBVERBS` when that engine is built (falling back to `EngineType::DPDK` in DPDK-only builds); `raw` + `engine: "dpdk"` explicitly selects `EngineType::DPDK`; `socket` + a `roce://` endpoint (or `engine: "ibverbs"`) selects `EngineType::RDMA`. The stream-aware `config_engine_from_string(str, stream_type)` overload encodes the `ibverbs`→`{IBVERBS for raw, RDMA for socket}` split. `EngineFactory` (in `engine.h`) is a singleton that instantiates the active engine. `daqiri_init(...)` resolves which engine to use from the `NetworkConfig`, runs the shared hardware-independent semantic validation, and only then delegates initialization through the `Engine` vtable. The standalone `daqiri_config_validate` tool calls the same parser and shared checks without creating an engine. There is only ever **one** active `Engine` per process.

The always-built socket engine implements Linux UDP/TCP streams directly. Applications that need kernel socket tuning call `socket_setsockopt(conn_id, level, optname, optval, optlen)` after resolving a TCP/UDP connection ID; DAQIRI passes the numeric Linux constants through without maintaining a symbolic option map. `socket_setsockopt` is not supported for `roce://` connections, which delegate to the RDMA/ibverbs path.
The always-built socket engine implements Linux UDP/TCP streams directly. Applications that need kernel socket tuning call `socket_setsockopt(conn_id, level, optname, optval, optlen)` after resolving a TCP/UDP connection ID; DAQIRI passes the numeric Linux constants through without maintaining a symbolic option map. TCP RX queues bound their internal burst backlog by the smallest `num_bufs` among the queue's referenced memory regions. Receive threads wait for readable data without reserving shared queue capacity, then reserve a slot before consuming; a full queue stops socket consumption so TCP backpressure reaches the sender, while an idle peer cannot monopolize capacity shared with active peers. UDP retains its existing receive behavior. `socket_setsockopt` is not supported for `roce://` connections, which delegate to the RDMA/ibverbs path.

The opt-in `dpdk` raw engine (`src/engines/dpdk/`, `DpdkEngine`) programs RX steering, send-to-kernel fallbacks (`flow_isolation: true`), and `tx_eth_src` TX offloads via DPDK RTE Flow during `daqiri_init()`. Standard UDP/IP (group 3), flex-item (group 1) and eCPRI-over-Ethernet (group 2, EtherType `0xAEFE`, via `RTE_FLOW_ITEM_TYPE_ECPRI` matching message type and pc_id/rtc_id) RX flows use separate flow groups; `validate_config()` rejects mixing these flow classes per interface, duplicate `flex_item_id` values, and unknown queue targets / `flex_item_id` values before NIC programming. Queue actions with two or more IDs use mlx5 Toeplitz IPv4/UDP RSS. The mlx5 PMD honors the native eCPRI flow item only under **firmware steering** (`dv_flow_en=1`); under HW steering (`dv_flow_en=2`, the default) the rule installs but silently never matches on ConnectX-class NICs. So `initialize()` auto-switches any interface carrying eCPRI RX flows to `dv_flow_en=1` (with a `WARN`), which means the async/template dynamic-RX-flow path is unavailable on that port. Flex-item parser handles are created per `(port, flex_item_id)` (scoped per interface). All programmed `rte_flow` rules and flex-item handles are tracked and destroyed in order on shutdown, init failure, and engine teardown (programmed flows → flex items → group-0 ETH jump rules). See `docs/benchmarks/raw_benchmarking.md` (Flow programming smoke test) for manual verification steps. It also reports 802.3x pause state at init and per-run pause counter deltas from `src/net_pause.{h,cpp}` when a kernel netdev is available.

Expand Down
5 changes: 4 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,10 @@ DAQIRI provides direct NIC hardware access in userspace, bypassing the Linux ker
- **RDMA** — RDMA verbs (READ, WRITE, SEND) over RoCE on Ethernet NICs or InfiniBand.
- **Linux socket control** — TCP/UDP socket streams expose connection IDs and
`socket_setsockopt()` for native Linux `setsockopt` tuning without YAML option
name mappings.
name mappings. TCP RX queues bound their internal backlog using the smallest
referenced memory-region `num_bufs`, so a slow consumer propagates TCP
backpressure instead of growing an unbounded process queue. See
[TCP receive backpressure and queue sizing](https://nvidia.github.io/daqiri/benchmarks/socket_benchmarking/#tcp-receive-backpressure-and-queue-sizing).
- **Flow-control telemetry** — Raw Ethernet streams warn at `daqiri_init()` when 802.3x
pause is enabled on a port and report the pause frames exchanged during the run with the
shutdown stats. A paused link throttles the sender instead of dropping, so it caps
Expand Down
7 changes: 6 additions & 1 deletion docs/api-reference/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,12 @@ runtime binding.
- values: `local`, `rdma_read`, `rdma_write`
- **`num_bufs`**: Number of buffers in this region. Higher values give more processing
headroom but consume more memory (GPU BAR1 for `device`). Too low risks dropped packets
on RX or a TX stall. Raw DPDK queue regions use a floor of
on RX or a TX stall. For socket TCP RX queues, the smallest `num_bufs` among the queue's
memory regions bounds the number of bursts DAQIRI holds internally. When that queue fills,
DAQIRI stops calling `recv()`, allowing TCP flow control to backpressure the sender. UDP RX
does not wait for this internal queue capacity. See
[TCP receive backpressure and queue sizing](../benchmarks/socket_benchmarking.md#tcp-receive-backpressure-and-queue-sizing)
for sizing guidance and ownership details. Raw DPDK queue regions use a floor of
`max(1.5 * ring, ring + 2 * batch_size)`. DAQIRI bumps values below that floor to
`max(3 * ring, ring + 4 * batch_size)` and warns with the exact `num_bufs` to configure.
With the default 8192-descriptor ring and `batch_size: 10240`, the floor is 28672 and the
Expand Down
56 changes: 56 additions & 0 deletions docs/benchmarks/socket_benchmarking.md
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,62 @@ The shipped configs run both endpoints on `127.0.0.1` and are useful for a smoke
--seconds 10 --mode both
```

### TCP receive backpressure and queue sizing

DAQIRI bounds each TCP RX queue by the smallest `num_bufs` value among the
memory regions referenced by that queue. Each successful TCP `recv()` becomes
one DAQIRI burst and consumes one slot until the application calls
`get_rx_burst()` to remove it from the internal queue. For example, this queue
allows at most 64 received bursts to wait inside DAQIRI:

```yaml
memory_regions:
- name: "TCP_RX"
kind: "host"
affinity: 0
num_bufs: 64
buf_size: 1052672

interfaces:
- name: tcp_server
address: 10.250.0.2
socket_config:
mode: server
local_addr: "tcp://10.250.0.2:6001"
rx:
queues:
- name: "TCP_RX_Queue"
id: 0
cpu_core: 8
batch_size: 1
memory_regions: ["TCP_RX"]
```

When all slots are occupied, the receive threads stop consuming bytes from the
kernel sockets. The kernel TCP receive window then contracts and eventually
slows or blocks the sender. This is normal TCP flow control: a sustained slow
GPU/application consumer should reduce achieved throughput instead of causing
DAQIRI's internal queue to grow until the process runs out of memory. An idle
connection waits for readable data without reserving a slot, so it does not
prevent an active connection sharing the endpoint queue from receiving.

Choose `num_bufs` according to the amount of receive jitter to absorb:

- Increase it to tolerate longer temporary compute stalls, at the cost of more
queued payload memory and later backpressure.
- Decrease it to apply backpressure sooner and reduce the maximum internal
backlog. A value of `1` permits one waiting burst.
- Size from observed `recv()` chunks, not application message count. TCP is a
byte stream, so one application write is not guaranteed to equal one DAQIRI
burst.

The bound covers bursts waiting inside DAQIRI. It does not include bytes still
in the kernel socket receive buffer or bursts already returned to and retained
by the application. The application must continue to free every received burst
after processing it. TCP `rx.queues[].batch_size` does not change this bound;
the socket TCP path surfaces one packet per `recv()` chunk. UDP retains its
existing datagram receive behavior and does not wait for this TCP queue capacity.

For an on-wire namespace test, use separate server and client YAML files. The important fields are the endpoint URI scheme, namespace IPs, server port, `max_payload_size`, memory-region `buf_size`, and benchmark `message_size`.

For UDP, `rx.queues[].cpu_core` pins the DAQIRI socket I/O thread that drains
Expand Down
5 changes: 4 additions & 1 deletion docs/concepts.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,10 @@ URI schemes:
- **`udp://`** / **`tcp://`**: Linux kernel UDP and TCP sockets. No NIC
privileges required, no special hardware. Useful
as a comparison baseline against the kernel-bypass paths and as a way
to get first results on a system without an NVIDIA NIC.
to get first results on a system without an NVIDIA NIC. TCP RX queues use
the smallest referenced memory-region `num_bufs` as their internal burst
limit; when full, DAQIRI stops consuming socket data so kernel TCP flow
control reaches the sender. See [TCP receive backpressure and queue sizing](benchmarks/socket_benchmarking.md#tcp-receive-backpressure-and-queue-sizing).
- **`roce://` endpoints**: RDMA over Converged Ethernet, using the
open-source [`rdma-core`](https://github.com/linux-rdma/rdma-core)
library. A server/client connection model, NIC-level reliable
Expand Down
6 changes: 6 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,12 @@ hide:

DAQIRI's baseline requirements depend on which [stream type](concepts.md#stream-types) you plan to use. The Linux Sockets path (`stream_type: "socket"` with `udp://` or `tcp://` endpoints) runs on any modern Linux box. The Raw Ethernet kernel-bypass path and GPUDirect impose additional hardware requirements, listed below.

The built-in TCP path bounds each internal RX backlog by the smallest
memory-region `num_bufs` referenced by that queue. A sustained slow consumer
therefore applies TCP backpressure rather than allowing DAQIRI's process queue
to grow without limit. See
[TCP receive backpressure and queue sizing](benchmarks/socket_benchmarking.md#tcp-receive-backpressure-and-queue-sizing).

| Component | Requirement |
|-----------|-------------|
| **OS** | Linux (kernel 5.15+), Ubuntu 22.04 recommended |
Expand Down
2 changes: 1 addition & 1 deletion docs/tutorials/configuration-walkthrough.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ hide:
DAQIRI exposes a single API on top of multiple packet I/O stacks, selected at runtime with `stream_type` and endpoint URI schemes such as `udp://`, `tcp://`, and `roce://`. Pick the row that matches your hardware and the role of the other endpoint:

- **Raw Ethernet**: `stream_type: "raw"`. Kernel-bypass with GPUDirect zero-copy. Highest performance. Requires an [NVIDIA ConnectX-class NIC](https://www.nvidia.com/en-us/networking/ethernet-adapters/). `tx_port` and `rx_port` can share one physical NIC for a single-host closed-loop bench, or be split across two hosts.
- **Socket (UDP / TCP)**: `stream_type: "socket"` with `udp://` or `tcp://` endpoints. Plain Linux kernel sockets. No NIC, no privileges, no special CMake flags. Useful as a comparison baseline and as a path to first results on a system without an NVIDIA NIC. Socket options are runtime API calls: resolve a connection ID, then pass native Linux constants to `socket_setsockopt()` rather than adding option names to YAML.
- **Socket (UDP / TCP)**: `stream_type: "socket"` with `udp://` or `tcp://` endpoints. Plain Linux kernel sockets. No NIC, no privileges, no special CMake flags. Useful as a comparison baseline and as a path to first results on a system without an NVIDIA NIC. Socket options are runtime API calls: resolve a connection ID, then pass native Linux constants to `socket_setsockopt()` rather than adding option names to YAML. For TCP RX, the smallest memory-region `num_bufs` referenced by the queue bounds DAQIRI's internal backlog and determines when [TCP backpressure](../benchmarks/socket_benchmarking.md#tcp-receive-backpressure-and-queue-sizing) begins.
- **Socket (RoCE / RDMA)**: `stream_type: "socket"` and `roce://` endpoints. RDMA verbs over Ethernet, with a server/client connection model and a NIC-level reliable transport. Primarily intended for setups where **one** endpoint is a third-party RoCE implementation (FPGA, instrument, customer black box). When both peers run DAQIRI, prefer an upper-layer library such as MPI / NCCL / UCX instead.

If you don't have any NIC at all, the `*_sw_loopback*` variants of the Raw Ethernet configs need no hardware, which is useful for first-time build verification.
Expand Down
87 changes: 83 additions & 4 deletions src/engines/socket/daqiri_socket_engine.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@

#include <algorithm>
#include <array>
#include <limits>
#include <stdexcept>
#include <vector>

Expand Down Expand Up @@ -100,6 +101,8 @@ bool pin_udp_rx_thread(int cpu_core, uint16_t port) {

} // namespace

SocketEngine::SocketEngine() = default;

SocketEngine::~SocketEngine() {
shutdown();
}
Expand Down Expand Up @@ -215,6 +218,13 @@ void SocketEngine::initialize() {
ep->tx_batch_size = select_batch_size(if_cfg.tx_.queues_);
ep->max_packet_size = static_cast<size_t>(std::max(1, select_max_packet_size(if_cfg)));
ep->rx_queue_state = get_or_create_rx_queue(ep->port, ep->rx_queue);
ep->rx_queue_state->max_bursts = select_rx_queue_capacity(if_cfg);
if (cfg_.common_.protocol == SocketProtocol::TCP) {
DAQIRI_LOG_INFO("TCP RX queue port={} queue={} capacity={} bursts",
ep->port,
ep->rx_queue,
ep->rx_queue_state->max_bursts);
}
ep->rx_metrics = metrics::get_or_create_queue("socket",
if_cfg.name_.empty() ? if_cfg.address_
: if_cfg.name_,
Expand Down Expand Up @@ -522,7 +532,15 @@ void SocketEngine::close_all_connections() {

for (auto& conn : conn_copy) {
if (conn == nullptr) { continue; }
conn->running.store(false);
if (conn->rx_queue != nullptr) {
{
std::lock_guard<std::mutex> lock(conn->rx_queue->mutex);
conn->running.store(false);
}
conn->rx_queue->capacity_cv.notify_all();
} else {
conn->running.store(false);
}
close_fd(conn->fd);
}

Expand Down Expand Up @@ -628,19 +646,49 @@ BurstParams* SocketEngine::create_tx_burst_params() {
Status SocketEngine::pop_rx_burst(const std::shared_ptr<RxQueueState>& qstate, BurstParams** burst) {
if (burst == nullptr || qstate == nullptr) { return Status::INVALID_PARAMETER; }

std::lock_guard<std::mutex> lock(qstate->mutex);
std::unique_lock<std::mutex> lock(qstate->mutex);
if (qstate->bursts.empty()) {
*burst = nullptr;
return Status::NULL_PTR;
}

*burst = qstate->bursts.front();
qstate->bursts.pop();
lock.unlock();
qstate->capacity_cv.notify_one();
return Status::SUCCESS;
}

void SocketEngine::push_rx_burst(const std::shared_ptr<RxQueueState>& qstate, BurstParams* burst) {
bool SocketEngine::reserve_rx_burst(const std::shared_ptr<RxQueueState>& qstate,
const std::atomic<bool>& connection_running) {
if (qstate == nullptr) { return false; }

std::unique_lock<std::mutex> lock(qstate->mutex);
qstate->capacity_cv.wait(lock, [&] {
if (!running_.load() || !connection_running.load()) { return true; }
return qstate->reserved_bursts < qstate->max_bursts &&
qstate->bursts.size() < qstate->max_bursts - qstate->reserved_bursts;
});
if (!running_.load() || !connection_running.load()) { return false; }

++qstate->reserved_bursts;
return true;
}

void SocketEngine::cancel_rx_burst_reservation(const std::shared_ptr<RxQueueState>& qstate) {
if (qstate == nullptr) { return; }

{
std::lock_guard<std::mutex> lock(qstate->mutex);
if (qstate->reserved_bursts > 0) { --qstate->reserved_bursts; }
}
qstate->capacity_cv.notify_one();
}

void SocketEngine::push_rx_burst(const std::shared_ptr<RxQueueState>& qstate, BurstParams* burst,
bool reserved) {
if (burst == nullptr) {
if (reserved) { cancel_rx_burst_reservation(qstate); }
return;
}
if (qstate == nullptr) {
Expand All @@ -650,6 +698,7 @@ void SocketEngine::push_rx_burst(const std::shared_ptr<RxQueueState>& qstate, Bu
return;
}
std::lock_guard<std::mutex> lock(qstate->mutex);
if (reserved && qstate->reserved_bursts > 0) { --qstate->reserved_bursts; }
qstate->bursts.push(burst);
}

Expand Down Expand Up @@ -1031,6 +1080,17 @@ int SocketEngine::select_max_packet_size(const InterfaceConfig& if_cfg) const {
return max_size;
}

size_t SocketEngine::select_rx_queue_capacity(const InterfaceConfig& if_cfg) const {
if (if_cfg.rx_.queues_.empty()) { return 1; }

size_t capacity = std::numeric_limits<size_t>::max();
for (const auto& mr_name : if_cfg.rx_.queues_.front().common_.mrs_) {
const auto mr_it = cfg_.mrs_.find(mr_name);
if (mr_it != cfg_.mrs_.end()) { capacity = std::min(capacity, mr_it->second.num_bufs_); }
}
return capacity == std::numeric_limits<size_t>::max() ? 1 : std::max<size_t>(1, capacity);
Comment thread
cliffburdick marked this conversation as resolved.
}

uint16_t SocketEngine::select_queue_id(const std::vector<RxQueueConfig>& queues) const {
if (queues.empty()) { return 0; }
return static_cast<uint16_t>(queues.front().common_.id_);
Expand Down Expand Up @@ -1364,11 +1424,30 @@ void SocketEngine::tcp_rx_loop(std::shared_ptr<ConnectionState> conn) {
std::vector<uint8_t> tmp(max_size);

while (running_.load() && conn->running.load()) {
// Wait for data without consuming it or reserving shared queue capacity.
// Otherwise an idle peer could reserve the last slot while blocked in
// recv(), preventing active peers on the same endpoint from making progress.
uint8_t marker = 0;
const ssize_t ready = ::recv(conn->fd, &marker, sizeof(marker), MSG_PEEK);
if (ready == 0) { break; }
if (ready < 0) {
if (errno == EINTR) { continue; }
if (!running_.load()) { break; }
DAQIRI_LOG_WARN("TCP recv peek failed on conn_id={}: {}", conn->conn_id, strerror(errno));
break;
}

// Once the application falls behind, stop draining the kernel socket so
// TCP flow control can propagate backpressure without growing DAQIRI's heap.
if (!reserve_rx_burst(conn->rx_queue, conn->running)) { break; }

const ssize_t rx = ::recv(conn->fd, tmp.data(), tmp.size(), 0);
Comment thread
cliffburdick marked this conversation as resolved.
if (rx == 0) {
cancel_rx_burst_reservation(conn->rx_queue);
break;
}
if (rx < 0) {
cancel_rx_burst_reservation(conn->rx_queue);
if (errno == EINTR) { continue; }
if (!running_.load()) { break; }
DAQIRI_LOG_WARN("TCP recv failed on conn_id={}: {}", conn->conn_id, strerror(errno));
Expand All @@ -1389,7 +1468,7 @@ void SocketEngine::tcp_rx_loop(std::shared_ptr<ConnectionState> conn) {
burst->pkt_lens[0][0] = static_cast<uint32_t>(rx);
set_connection_id(burst, conn->conn_id);

push_rx_burst(conn->rx_queue, burst);
push_rx_burst(conn->rx_queue, burst, true);
rx_pkts_.fetch_add(1);
rx_bytes_.fetch_add(static_cast<uint64_t>(rx));
metrics::add_rx(conn->rx_metrics, 1, static_cast<uint64_t>(rx));
Expand Down
Loading
Loading