From 91dc8f2c83d0d1862e6ececd86c2b4c8d1cadddc Mon Sep 17 00:00:00 2001 From: Cliff Burdick Date: Wed, 30 Sep 2026 21:52:25 +0000 Subject: [PATCH 1/4] #323 - Apply TCP receive backpressure Signed-off-by: Cliff Burdick --- docs/api-reference/configuration.md | 5 +- src/engines/socket/daqiri_socket_engine.cpp | 73 +++++++++++- src/engines/socket/daqiri_socket_engine.h | 13 ++- tests/cpp/CMakeLists.txt | 22 ++++ tests/cpp/socket_rx_queue_test.cpp | 123 ++++++++++++++++++++ 5 files changed, 230 insertions(+), 6 deletions(-) create mode 100644 tests/cpp/socket_rx_queue_test.cpp diff --git a/docs/api-reference/configuration.md b/docs/api-reference/configuration.md index d604416d..6bdfd712 100644 --- a/docs/api-reference/configuration.md +++ b/docs/api-reference/configuration.md @@ -128,7 +128,10 @@ 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. 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 diff --git a/src/engines/socket/daqiri_socket_engine.cpp b/src/engines/socket/daqiri_socket_engine.cpp index a8f1e127..30c3fca0 100644 --- a/src/engines/socket/daqiri_socket_engine.cpp +++ b/src/engines/socket/daqiri_socket_engine.cpp @@ -36,6 +36,7 @@ #include #include +#include #include #include @@ -215,6 +216,13 @@ void SocketEngine::initialize() { ep->tx_batch_size = select_batch_size(if_cfg.tx_.queues_); ep->max_packet_size = static_cast(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_, @@ -522,7 +530,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 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); } @@ -628,7 +644,7 @@ BurstParams* SocketEngine::create_tx_burst_params() { Status SocketEngine::pop_rx_burst(const std::shared_ptr& qstate, BurstParams** burst) { if (burst == nullptr || qstate == nullptr) { return Status::INVALID_PARAMETER; } - std::lock_guard lock(qstate->mutex); + std::unique_lock lock(qstate->mutex); if (qstate->bursts.empty()) { *burst = nullptr; return Status::NULL_PTR; @@ -636,11 +652,41 @@ Status SocketEngine::pop_rx_burst(const std::shared_ptr& qstate, B *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& qstate, BurstParams* burst) { +bool SocketEngine::reserve_rx_burst(const std::shared_ptr& qstate, + const std::atomic& connection_running) { + if (qstate == nullptr) { return false; } + + std::unique_lock 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& qstate) { + if (qstate == nullptr) { return; } + + { + std::lock_guard 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& qstate, BurstParams* burst, + bool reserved) { if (burst == nullptr) { + if (reserved) { cancel_rx_burst_reservation(qstate); } return; } if (qstate == nullptr) { @@ -650,6 +696,7 @@ void SocketEngine::push_rx_burst(const std::shared_ptr& qstate, Bu return; } std::lock_guard lock(qstate->mutex); + if (reserved && qstate->reserved_bursts > 0) { --qstate->reserved_bursts; } qstate->bursts.push(burst); } @@ -1031,6 +1078,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::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::max() ? 1 : std::max(1, capacity); +} + uint16_t SocketEngine::select_queue_id(const std::vector& queues) const { if (queues.empty()) { return 0; } return static_cast(queues.front().common_.id_); @@ -1364,11 +1422,18 @@ void SocketEngine::tcp_rx_loop(std::shared_ptr conn) { std::vector tmp(max_size); while (running_.load() && conn->running.load()) { + // Reserve queue capacity before recv(). Once the application falls behind, + // this thread stops draining the kernel socket so TCP flow control can + // propagate backpressure to the sender 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); 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)); @@ -1389,7 +1454,7 @@ void SocketEngine::tcp_rx_loop(std::shared_ptr conn) { burst->pkt_lens[0][0] = static_cast(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(rx)); metrics::add_rx(conn->rx_metrics, 1, static_cast(rx)); diff --git a/src/engines/socket/daqiri_socket_engine.h b/src/engines/socket/daqiri_socket_engine.h index 35b03769..ee44b58c 100644 --- a/src/engines/socket/daqiri_socket_engine.h +++ b/src/engines/socket/daqiri_socket_engine.h @@ -35,6 +35,7 @@ namespace daqiri { class RdmaEngine; +class SocketEngineQueueTestPeer; class SocketEngine : public Engine { public: @@ -112,10 +113,15 @@ class SocketEngine : public Engine { RDMAOpCode rdma_get_opcode(BurstParams* burst) override; private: + friend class SocketEngineQueueTestPeer; + struct RxQueueState { uint16_t port = 0; uint16_t queue = 0; + size_t max_bursts = 1; + size_t reserved_bursts = 0; std::mutex mutex; + std::condition_variable capacity_cv; std::queue bursts; }; @@ -193,6 +199,7 @@ class SocketEngine : public Engine { const EndpointState* endpoint_for_port(uint16_t port) const; int select_max_packet_size(const InterfaceConfig& if_cfg) const; + size_t select_rx_queue_capacity(const InterfaceConfig& if_cfg) const; uint16_t select_queue_id(const std::vector& queues) const; uint16_t select_queue_id(const std::vector& queues) const; int select_cpu_core(const std::vector& queues) const; @@ -200,7 +207,11 @@ class SocketEngine : public Engine { uint32_t select_batch_size(const std::vector& queues) const; Status pop_rx_burst(const std::shared_ptr& qstate, BurstParams** burst); - void push_rx_burst(const std::shared_ptr& qstate, BurstParams* burst); + bool reserve_rx_burst(const std::shared_ptr& qstate, + const std::atomic& connection_running); + void cancel_rx_burst_reservation(const std::shared_ptr& qstate); + void push_rx_burst(const std::shared_ptr& qstate, BurstParams* burst, + bool reserved = false); void free_packet_arrays(BurstParams* burst); void close_fd(int& fd); diff --git a/tests/cpp/CMakeLists.txt b/tests/cpp/CMakeLists.txt index e43ba454..abd45ab0 100644 --- a/tests/cpp/CMakeLists.txt +++ b/tests/cpp/CMakeLists.txt @@ -23,3 +23,25 @@ else() endif() add_test(NAME daqiri_init_validation_test COMMAND daqiri_init_validation_test) + +add_executable(daqiri_socket_rx_queue_test socket_rx_queue_test.cpp) +target_compile_features(daqiri_socket_rx_queue_test PRIVATE cxx_std_17) +target_include_directories(daqiri_socket_rx_queue_test PRIVATE ${PROJECT_SOURCE_DIR}) +if(BUILD_SHARED_LIBS AND UNIX AND NOT APPLE) + # daqiri_common and the engine DSOs reference each other. Keep daqiri_common + # loaded while this internal test links directly to the socket engine. + target_link_options(daqiri_socket_rx_queue_test PRIVATE "LINKER:--no-as-needed") +endif() + +if(NOT BUILD_SHARED_LIBS AND UNIX AND NOT APPLE) + target_link_libraries(daqiri_socket_rx_queue_test PRIVATE + "-Wl,--start-group" + ${DAQIRI_TEST_STATIC_LIBS} + "-Wl,--end-group" + ${DAQIRI_YAML_TARGET}) +else() + target_link_libraries(daqiri_socket_rx_queue_test PRIVATE daqiri_socket daqiri::daqiri + ${DAQIRI_YAML_TARGET}) +endif() + +add_test(NAME daqiri_socket_rx_queue_test COMMAND daqiri_socket_rx_queue_test) diff --git a/tests/cpp/socket_rx_queue_test.cpp b/tests/cpp/socket_rx_queue_test.cpp new file mode 100644 index 00000000..13116ed1 --- /dev/null +++ b/tests/cpp/socket_rx_queue_test.cpp @@ -0,0 +1,123 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include +#include +#include +#include +#include + +#include "src/engines/socket/daqiri_socket_engine.h" + +namespace daqiri { + +class SocketEngineQueueTestPeer { + public: + static bool bounded_queue_blocks_until_pop() { + SocketEngine engine; + engine.running_.store(true); + + auto queue = std::make_shared(); + queue->max_bursts = 1; + auto* queued_burst = new BurstParams{}; + engine.push_rx_burst(queue, queued_burst); + + std::atomic connection_running{true}; + auto reservation = std::async( + std::launch::async, [&] { return engine.reserve_rx_burst(queue, connection_running); }); + + const bool initially_blocked = + reservation.wait_for(std::chrono::milliseconds(100)) == std::future_status::timeout; + + BurstParams* popped_burst = nullptr; + const bool popped = engine.pop_rx_burst(queue, &popped_burst) == Status::SUCCESS && + popped_burst == queued_burst; + if (!popped) { + { + std::lock_guard lock(queue->mutex); + connection_running.store(false); + } + queue->capacity_cv.notify_all(); + } + + bool reservation_ready = + reservation.wait_for(std::chrono::seconds(1)) == std::future_status::ready; + if (!reservation_ready) { + { + std::lock_guard lock(queue->mutex); + connection_running.store(false); + } + queue->capacity_cv.notify_all(); + reservation_ready = + reservation.wait_for(std::chrono::seconds(1)) == std::future_status::ready; + } + const bool reserved = reservation_ready && reservation.get(); + if (reserved) { + engine.cancel_rx_burst_reservation(queue); + } + engine.running_.store(false); + engine.free_rx_burst(popped_burst); + return initially_blocked && popped && reserved && queue->bursts.empty() && + queue->reserved_bursts == 0; + } + + static bool shutdown_wakes_blocked_receiver() { + SocketEngine engine; + engine.running_.store(true); + + auto queue = std::make_shared(); + queue->max_bursts = 1; + auto* queued_burst = new BurstParams{}; + engine.push_rx_burst(queue, queued_burst); + + std::atomic connection_running{true}; + auto reservation = std::async( + std::launch::async, [&] { return engine.reserve_rx_burst(queue, connection_running); }); + + const bool initially_blocked = + reservation.wait_for(std::chrono::milliseconds(100)) == std::future_status::timeout; + + { + std::lock_guard lock(queue->mutex); + connection_running.store(false); + } + queue->capacity_cv.notify_all(); + const bool stopped = + reservation.wait_for(std::chrono::seconds(1)) == std::future_status::ready && + !reservation.get(); + + BurstParams* popped_burst = nullptr; + engine.pop_rx_burst(queue, &popped_burst); + engine.running_.store(false); + engine.free_rx_burst(popped_burst); + return initially_blocked && stopped; + } +}; + +} // namespace daqiri + +int main() { + if (!daqiri::SocketEngineQueueTestPeer::bounded_queue_blocks_until_pop()) { + std::cerr << "bounded TCP RX queue did not unblock after a pop\n"; + return 1; + } + if (!daqiri::SocketEngineQueueTestPeer::shutdown_wakes_blocked_receiver()) { + std::cerr << "blocked TCP RX queue did not wake for shutdown\n"; + return 1; + } + return 0; +} From a7299f5752d0eccb36dab317390da905ecfbb4fc Mon Sep 17 00:00:00 2001 From: Cliff Burdick Date: Wed, 30 Sep 2026 22:06:34 +0000 Subject: [PATCH 2/4] #323 - Prevent idle peer queue starvation Signed-off-by: Cliff Burdick --- src/engines/socket/daqiri_socket_engine.cpp | 18 ++++- tests/cpp/socket_rx_queue_test.cpp | 78 +++++++++++++++++++++ 2 files changed, 93 insertions(+), 3 deletions(-) diff --git a/src/engines/socket/daqiri_socket_engine.cpp b/src/engines/socket/daqiri_socket_engine.cpp index 30c3fca0..4e1b32aa 100644 --- a/src/engines/socket/daqiri_socket_engine.cpp +++ b/src/engines/socket/daqiri_socket_engine.cpp @@ -1422,9 +1422,21 @@ void SocketEngine::tcp_rx_loop(std::shared_ptr conn) { std::vector tmp(max_size); while (running_.load() && conn->running.load()) { - // Reserve queue capacity before recv(). Once the application falls behind, - // this thread stops draining the kernel socket so TCP flow control can - // propagate backpressure to the sender without growing DAQIRI's heap. + // 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); diff --git a/tests/cpp/socket_rx_queue_test.cpp b/tests/cpp/socket_rx_queue_test.cpp index 13116ed1..79de92c8 100644 --- a/tests/cpp/socket_rx_queue_test.cpp +++ b/tests/cpp/socket_rx_queue_test.cpp @@ -20,6 +20,9 @@ #include #include #include +#include +#include +#include #include "src/engines/socket/daqiri_socket_engine.h" @@ -27,6 +30,77 @@ namespace daqiri { class SocketEngineQueueTestPeer { public: + static bool idle_connection_does_not_reserve_capacity() { + SocketEngine engine; + engine.running_.store(true); + + auto queue = std::make_shared(); + queue->max_bursts = 1; + + int idle_fds[2] = {-1, -1}; + int active_fds[2] = {-1, -1}; + if (::socketpair(AF_UNIX, SOCK_STREAM, 0, idle_fds) != 0 || + ::socketpair(AF_UNIX, SOCK_STREAM, 0, active_fds) != 0) { + if (idle_fds[0] >= 0) { + ::close(idle_fds[0]); + } + if (idle_fds[1] >= 0) { + ::close(idle_fds[1]); + } + return false; + } + + auto idle = std::make_shared(); + idle->fd = idle_fds[0]; + idle->conn_id = 1; + idle->rx_queue = queue; + idle->running.store(true); + std::thread idle_thread(&SocketEngine::tcp_rx_loop, &engine, idle); + + // Give the idle receiver time to block. The old implementation reserved + // the only queue slot before entering recv() and starved the active peer. + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + + auto active = std::make_shared(); + active->fd = active_fds[0]; + active->conn_id = 2; + active->rx_queue = queue; + active->running.store(true); + std::thread active_thread(&SocketEngine::tcp_rx_loop, &engine, active); + + constexpr uint8_t payload = 42; + const bool sent = ::send(active_fds[1], &payload, sizeof(payload), 0) == sizeof(payload); + + BurstParams* received = nullptr; + const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(1); + while (sent && std::chrono::steady_clock::now() < deadline && received == nullptr) { + if (engine.pop_rx_burst(queue, &received) != Status::SUCCESS) { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + } + + engine.running_.store(false); + { + std::lock_guard lock(queue->mutex); + idle->running.store(false); + active->running.store(false); + } + queue->capacity_cv.notify_all(); + ::shutdown(idle_fds[1], SHUT_RDWR); + ::shutdown(active_fds[1], SHUT_RDWR); + ::close(idle_fds[1]); + ::close(active_fds[1]); + idle_thread.join(); + active_thread.join(); + + const bool correct_payload = received != nullptr && received->hdr.hdr.num_pkts == 1 && + received->pkt_lens[0][0] == sizeof(payload) && + *reinterpret_cast(received->pkts[0][0]) == payload; + engine.free_all_packets(received); + engine.free_rx_burst(received); + return sent && correct_payload; + } + static bool bounded_queue_blocks_until_pop() { SocketEngine engine; engine.running_.store(true); @@ -111,6 +185,10 @@ class SocketEngineQueueTestPeer { } // namespace daqiri int main() { + if (!daqiri::SocketEngineQueueTestPeer::idle_connection_does_not_reserve_capacity()) { + std::cerr << "idle TCP connection blocked an active peer\n"; + return 1; + } if (!daqiri::SocketEngineQueueTestPeer::bounded_queue_blocks_until_pop()) { std::cerr << "bounded TCP RX queue did not unblock after a pop\n"; return 1; From 48a880c3318e924dbfcb05726b7cc7467ca807e5 Mon Sep 17 00:00:00 2001 From: Cliff Burdick Date: Wed, 30 Sep 2026 22:15:23 +0000 Subject: [PATCH 3/4] #323 - Document TCP receive backpressure Signed-off-by: Cliff Burdick --- docs/api-reference/configuration.md | 4 +- docs/benchmarks/socket_benchmarking.md | 56 ++++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 1 deletion(-) diff --git a/docs/api-reference/configuration.md b/docs/api-reference/configuration.md index 6bdfd712..0681285e 100644 --- a/docs/api-reference/configuration.md +++ b/docs/api-reference/configuration.md @@ -131,7 +131,9 @@ runtime binding. 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. Raw DPDK queue regions use a floor of + 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 diff --git a/docs/benchmarks/socket_benchmarking.md b/docs/benchmarks/socket_benchmarking.md index 0ef443c1..2c2916d7 100644 --- a/docs/benchmarks/socket_benchmarking.md +++ b/docs/benchmarks/socket_benchmarking.md @@ -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 From 61ff3c02b1a104e576ae31db55e1567864144a9e Mon Sep 17 00:00:00 2001 From: Cliff Burdick Date: Wed, 30 Sep 2026 22:29:18 +0000 Subject: [PATCH 4/4] #323 - Complete TCP backpressure validation Signed-off-by: Cliff Burdick --- AGENTS.md | 2 +- README.md | 5 ++++- docs/concepts.md | 5 ++++- docs/getting-started.md | 6 ++++++ docs/tutorials/configuration-walkthrough.md | 2 +- src/engines/socket/daqiri_socket_engine.cpp | 2 ++ src/engines/socket/daqiri_socket_engine.h | 2 +- 7 files changed, 19 insertions(+), 5 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index c98b60cd..6eb1ab9c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -130,7 +130,7 @@ clang-format -style=file -i -fallback-style=none `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. diff --git a/README.md b/README.md index 2117ef1f..c37d99d1 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/docs/concepts.md b/docs/concepts.md index 14e81646..beda3996 100644 --- a/docs/concepts.md +++ b/docs/concepts.md @@ -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 diff --git a/docs/getting-started.md b/docs/getting-started.md index 144b163a..6e9342ce 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -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 | diff --git a/docs/tutorials/configuration-walkthrough.md b/docs/tutorials/configuration-walkthrough.md index 90d1c959..e0557498 100644 --- a/docs/tutorials/configuration-walkthrough.md +++ b/docs/tutorials/configuration-walkthrough.md @@ -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. diff --git a/src/engines/socket/daqiri_socket_engine.cpp b/src/engines/socket/daqiri_socket_engine.cpp index 4e1b32aa..3300e39f 100644 --- a/src/engines/socket/daqiri_socket_engine.cpp +++ b/src/engines/socket/daqiri_socket_engine.cpp @@ -101,6 +101,8 @@ bool pin_udp_rx_thread(int cpu_core, uint16_t port) { } // namespace +SocketEngine::SocketEngine() = default; + SocketEngine::~SocketEngine() { shutdown(); } diff --git a/src/engines/socket/daqiri_socket_engine.h b/src/engines/socket/daqiri_socket_engine.h index ee44b58c..18f4985f 100644 --- a/src/engines/socket/daqiri_socket_engine.h +++ b/src/engines/socket/daqiri_socket_engine.h @@ -39,7 +39,7 @@ class SocketEngineQueueTestPeer; class SocketEngine : public Engine { public: - SocketEngine() = default; + SocketEngine(); ~SocketEngine() override; bool set_config_and_initialize(const NetworkConfig& cfg) override;