diff --git a/docs/api-reference/configuration.md b/docs/api-reference/configuration.md index 42abd7e8..bddb5fc6 100644 --- a/docs/api-reference/configuration.md +++ b/docs/api-reference/configuration.md @@ -60,6 +60,12 @@ These settings apply globally to both TX and RX: `memory_regions:` List of regions where packet buffers are stored. The number of regions and their `kind` determines the receive mode (CPU-only, header-data split, or batched GPU). +YAML describes region semantics but never contains process-local pointers. C++ and Python callers +can attach application-owned allocations by name at initialization; see +[Application-owned memory regions](cpp.md#application-owned-memory-regions). A runtime binding +replaces DAQIRI's allocation for that region. The legacy `owned: false` form requires a matching +runtime binding. + - **`name`**: Memory region name. Referenced by queue configurations. - type: `string` - **`kind`**: Memory type. diff --git a/docs/api-reference/cpp.md b/docs/api-reference/cpp.md index b0d8ec47..0128c97c 100644 --- a/docs/api-reference/cpp.md +++ b/docs/api-reference/cpp.md @@ -56,6 +56,44 @@ Only one engine may be active in a process. Calling `daqiri_init()` again before After `shutdown()` completes, a later `daqiri_init()` creates a fresh engine instance and allocates/registers its resources again. +### Application-owned memory regions + +Applications that need allocations in a specific CUDA context (including MPS applications) can +bind their own storage to configured memory-region names. Query the engine-adjusted requirements +first; the returned capacity includes packet headroom, slot alignment, and any DPDK buffer-count +safety adjustment. + +```cpp +daqiri::NetworkConfig config; +if (daqiri::parse_network_config("config.yaml", config) != daqiri::Status::SUCCESS) { + throw std::runtime_error("invalid DAQIRI config"); +} + +daqiri::MemoryRegionRequirements requirements; +daqiri::get_memory_region_requirements(config, requirements); +const auto& rx = requirements.at("RX_GPU"); + +void* rx_gpu = nullptr; +cudaMalloc(&rx_gpu, rx.capacity); // Uses the application's current context. +daqiri::MemoryRegionBindings bindings{{"RX_GPU", {rx_gpu, rx.capacity}}}; + +if (daqiri::daqiri_init(config, bindings) != daqiri::Status::SUCCESS) { + cudaFree(rx_gpu); + throw std::runtime_error("DAQIRI initialization failed"); +} + +// Return every outstanding burst before shutdown. +daqiri::shutdown(); +cudaFree(rx_gpu); +``` + +Bindings may cover only some regions; DAQIRI allocates the rest. Bound memory is borrowed, so its +allocation and CUDA context must remain alive until `shutdown()` completes. DAQIRI deregisters it +from the NIC but never frees, unpins, or unmaps it. Direct TCP/UDP sockets reject bindings because +their configured regions are not packet pools. DPDK also rejects externally bound `huge` regions: +EAL assumes ownership of pre-existing hugepage mappings during initialization and unmaps them at +cleanup. External `huge` regions remain supported by ibverbs/RDMA. + If GPU RX `reorder_configs` are configured for Raw Ethernet (`stream_type: "raw"`), set one CUDA stream per GPU reorder plan before pulling reordered bursts. CPU reorder configs do not use a CUDA stream. See the [Configuration YAML Reference](configuration.md#rx-reorder-configs) @@ -647,10 +685,12 @@ workflow sections above show the common call order and ownership rules. | `version_year()` / `version_month()` / `version_patch()` | Return the CalVer components. | | `abi_version()` | Return the DAQIRI shared-library ABI version. | | `daqiri_init(NetworkConfig &config)` | Initialize DAQIRI from an already-populated config object. | +| `daqiri_init(config, bindings)` | Initialize with non-owning external memory bindings. | | `daqiri_init(const std::string &yaml_string_or_path)` | Initialize from a YAML string or YAML file path. | | `daqiri_init_from_yaml_string(const std::string &yaml_string)` | Initialize from YAML content. | | `daqiri_init_from_yaml_file(const std::string &yaml_path)` | Initialize from a YAML file path. | | `parse_network_config(...)` | Parse YAML into `NetworkConfig` without starting the engine. | +| `get_memory_region_requirements(config, requirements)` | Return effective slot size, count, capacity, and alignment. | | `get_engine_type()` | Return the active engine type after initialization. | | `get_engine_type(config)` | Return the engine type selected by a config object. | | `shutdown()` | Stop DAQIRI and release engine-owned resources. | diff --git a/docs/api-reference/python.md b/docs/api-reference/python.md index 2a1455dc..6481d28e 100644 --- a/docs/api-reference/python.md +++ b/docs/api-reference/python.md @@ -144,6 +144,30 @@ config = { status = daqiri.daqiri_init(config) ``` +Application-owned CUDA buffers use the same requirements-and-bindings workflow as C++. Addresses +are integers, and the Python allocation owner must stay alive until after `shutdown()`: + +```python +status, parsed = daqiri.parse_network_config("config.yaml") +status, requirements = daqiri.get_memory_region_requirements(parsed) +rx_req = requirements["RX_GPU"] + +# For example, gpu_owner may be a CuPy allocation made in the desired context. +bindings = { + "RX_GPU": daqiri.ExternalMemoryRegion(int(gpu_owner.data.ptr), rx_req.capacity), +} +status = daqiri.daqiri_init(parsed, bindings) + +# Return outstanding bursts first. +daqiri.shutdown() +del gpu_owner +``` + +Bindings may cover only a subset of regions. They are supported by DPDK, raw ibverbs, and +RoCE/RDMA, but not direct TCP/UDP sockets. DPDK rejects external `huge` regions because EAL cannot +preserve ownership of hugepage mappings that predate initialization; use `host`, `host_pinned`, or +ibverbs/RDMA for caller-owned CPU hugepages. + Parse without starting the engine: ```python @@ -513,6 +537,8 @@ The workflow sections above show the common call order and ownership rules. | Function | Returns | | --- | --- | | `daqiri_init(config)` | `Status` | +| `daqiri_init(config, bindings)` | `Status` | +| `get_memory_region_requirements(config)` | `(Status, dict[str, MemoryRegionRequirement])` | | `daqiri_init_from_yaml_string(yaml_string)` | `Status` | | `daqiri_init_from_yaml_file(yaml_path)` | `Status` | | `parse_network_config(yaml_string_or_path)` | `(Status, NetworkConfig)` | @@ -709,6 +735,8 @@ names that mostly omit the trailing underscore from the C++ member name (e.g. | `RxQueueConfig` | RX queue wrapper with common queue fields, timeout, and `QueuePollMode`. | | `TxQueueConfig` | TX queue wrapper with common queue fields and `QueuePollMode`. | | `MemoryRegionConfig` | Memory region kind, affinity, access flags, sizes, counts, and ownership. | +| `ExternalMemoryRegion` | Non-owning integer address and byte capacity. | +| `MemoryRegionRequirement` | Effective kind, slot size, buffer count, capacity, and alignment. | | `VlanActionConfig` | VLAN push parameters: VLAN ID, priority, DEI, and ethertype. | | `TunnelConfig` | VXLAN, GRE, or NVGRE tunnel template fields for hardware encap/decap actions. | | `FlowAction` | Flow action type, scalar queue target `id`, queue-list target `ids`, optional VLAN config, and optional tunnel config. | diff --git a/docs/concepts.md b/docs/concepts.md index 99a1c1a5..0da9cc7b 100644 --- a/docs/concepts.md +++ b/docs/concepts.md @@ -402,6 +402,12 @@ buffer pools and can lead to `NO_FREE_BURST_BUFFERS`, When to call each `free_*` function is documented in the [C++ API Usage page](api-reference/cpp.md#rx-step-3-free-buffers). +The backing allocation may be DAQIRI-owned or application-owned. For application-owned regions, +the caller binds an address and capacity to a configured region name during initialization. +DAQIRI owns the temporary NIC registration and packet-pool bookkeeping; the caller retains the +allocation and its CUDA context until shutdown finishes. DAQIRI never frees externally bound +memory. + ## RX Packet Aggregation and Reorder DAQIRI can perform GPU- or CPU-side packet aggregation and reordering diff --git a/include/daqiri/common.h b/include/daqiri/common.h index fb0a2d91..28c3addc 100644 --- a/include/daqiri/common.h +++ b/include/daqiri/common.h @@ -107,9 +107,19 @@ inline int EnabledDirections(const std::string &dir) { * INTERNAL_ERROR: Internal error */ Status daqiri_init(NetworkConfig &config); +Status daqiri_init(NetworkConfig &config, const MemoryRegionBindings &bindings); Status daqiri_init(const std::string &yaml_string_or_path); +Status daqiri_init(const std::string &yaml_string_or_path, + const MemoryRegionBindings &bindings); Status daqiri_init_from_yaml_string(const std::string &yaml_string); +Status daqiri_init_from_yaml_string(const std::string &yaml_string, + const MemoryRegionBindings &bindings); Status daqiri_init_from_yaml_file(const std::string &yaml_path); +Status daqiri_init_from_yaml_file(const std::string &yaml_path, + const MemoryRegionBindings &bindings); + +Status get_memory_region_requirements(const NetworkConfig &config, + MemoryRegionRequirements &requirements); Status parse_network_config(const std::string &yaml_string_or_path, NetworkConfig &config); diff --git a/include/daqiri/types.h b/include/daqiri/types.h index 9680f9da..262aeb24 100644 --- a/include/daqiri/types.h +++ b/include/daqiri/types.h @@ -166,6 +166,23 @@ struct UDPIPV4Pkt { enum class MemoryKind { HOST, HOST_PINNED, HUGE, DEVICE, INVALID }; +struct ExternalMemoryRegion { + void* data = nullptr; + size_t capacity = 0; +}; + +using MemoryRegionBindings = std::unordered_map; + +struct MemoryRegionRequirement { + MemoryKind kind = MemoryKind::INVALID; + size_t slot_size = 0; + size_t num_bufs = 0; + size_t capacity = 0; + size_t alignment = 0; +}; + +using MemoryRegionRequirements = std::unordered_map; + enum MemoryAccess { MEM_ACCESS_LOCAL = 1U, MEM_ACCESS_RDMA_WRITE = 1U << 1, diff --git a/python/daqiri_common_pybind.cpp b/python/daqiri_common_pybind.cpp index c5274232..5165f60a 100644 --- a/python/daqiri_common_pybind.cpp +++ b/python/daqiri_common_pybind.cpp @@ -230,13 +230,13 @@ py::tuple segment_packet_bytes_impl(BurstParams *burst, int seg, int idx, return py::make_tuple(status, py::bytes(out)); } -Status daqiri_init_from_python(py::object config_obj) { +Status daqiri_init_from_python(py::object config_obj, const MemoryRegionBindings& bindings) { try { if (py::isinstance(config_obj) || py::isinstance(config_obj)) { const auto yaml_string_or_path = py::cast(config_obj); py::gil_scoped_release release; - return daqiri_init(yaml_string_or_path); + return daqiri_init(yaml_string_or_path, bindings); } if (py::isinstance(config_obj)) { @@ -245,13 +245,13 @@ Status daqiri_init_from_python(py::object config_obj) { config_obj, py::arg("default_flow_style") = false); const auto yaml_str = py::cast(py_yaml_str); py::gil_scoped_release release; - return daqiri_init_from_yaml_string(yaml_str); + return daqiri_init_from_yaml_string(yaml_str, bindings); } if (py::isinstance(config_obj)) { auto &config = config_obj.cast(); py::gil_scoped_release release; - return daqiri_init(config); + return daqiri_init(config, bindings); } if (py::hasattr(config_obj, "value")) { @@ -259,7 +259,7 @@ Status daqiri_init_from_python(py::object config_obj) { const auto yaml_string_or_path = py::cast(yaml_string_or_path_obj); py::gil_scoped_release release; - return daqiri_init(yaml_string_or_path); + return daqiri_init(yaml_string_or_path, bindings); } if (py::hasattr(config_obj, "as_dict")) { @@ -268,14 +268,14 @@ Status daqiri_init_from_python(py::object config_obj) { config_obj.attr("as_dict")(), py::arg("default_flow_style") = false); const auto yaml_str = py::cast(py_yaml_str); py::gil_scoped_release release; - return daqiri_init_from_yaml_string(yaml_str); + return daqiri_init_from_yaml_string(yaml_str, bindings); } const auto yaml_string_or_path_obj = py::str(config_obj); const auto yaml_string_or_path = py::cast(yaml_string_or_path_obj); py::gil_scoped_release release; - return daqiri_init(yaml_string_or_path); + return daqiri_init(yaml_string_or_path, bindings); } catch (const py::error_already_set &e) { DAQIRI_LOG_ERROR("Python config conversion failed: {}", e.what()); return Status::INTERNAL_ERROR; @@ -743,6 +743,29 @@ void bind_config_types(py::module_ &m) { .def_readwrite("num_bufs", &MemoryRegionConfig::num_bufs_) .def_readwrite("owned", &MemoryRegionConfig::owned_); + py::class_(m, "ExternalMemoryRegion") + .def(py::init([](uintptr_t address, size_t capacity) { + return ExternalMemoryRegion{reinterpret_cast(address), capacity}; + }), + "address"_a, "capacity"_a) + .def_property( + "address", + [](const ExternalMemoryRegion& region) { + return reinterpret_cast(region.data); + }, + [](ExternalMemoryRegion& region, uintptr_t address) { + region.data = reinterpret_cast(address); + }) + .def_readwrite("capacity", &ExternalMemoryRegion::capacity); + + py::class_(m, "MemoryRegionRequirement") + .def(py::init<>()) + .def_readonly("kind", &MemoryRegionRequirement::kind) + .def_readonly("slot_size", &MemoryRegionRequirement::slot_size) + .def_readonly("num_bufs", &MemoryRegionRequirement::num_bufs) + .def_readonly("capacity", &MemoryRegionRequirement::capacity) + .def_readonly("alignment", &MemoryRegionRequirement::alignment); + py::class_(m, "RxQueueConfig") .def(py::init<>()) .def_readwrite("common", &RxQueueConfig::common_) @@ -985,13 +1008,31 @@ PYBIND11_MODULE(_daqiri, m) { bind_enums(m); bind_config_types(m); - m.def("daqiri_init", &daqiri_init_from_python, "config"_a, + m.def("daqiri_init", &daqiri_init_from_python, "config"_a, "bindings"_a = MemoryRegionBindings{}, "Initialize DAQIRI from a YAML path, YAML string, dict, or config-like " "object"); - m.def("daqiri_init_from_yaml_string", &daqiri_init_from_yaml_string, - "yaml_string"_a, py::call_guard()); - m.def("daqiri_init_from_yaml_file", &daqiri_init_from_yaml_file, - "yaml_path"_a, py::call_guard()); + m.def( + "daqiri_init_from_yaml_string", + [](const std::string &yaml, const MemoryRegionBindings &bindings) { + return daqiri_init_from_yaml_string(yaml, bindings); + }, + "yaml_string"_a, "bindings"_a = MemoryRegionBindings{}, + py::call_guard()); + m.def( + "daqiri_init_from_yaml_file", + [](const std::string& path, const MemoryRegionBindings& bindings) { + return daqiri_init_from_yaml_file(path, bindings); + }, + "yaml_path"_a, "bindings"_a = MemoryRegionBindings{}, + py::call_guard()); + m.def( + "get_memory_region_requirements", + [](const NetworkConfig &config) { + MemoryRegionRequirements requirements; + const Status status = get_memory_region_requirements(config, requirements); + return py::make_tuple(status, requirements); + }, + "config"_a); m.def( "parse_network_config", [](const std::string &yaml_string_or_path) { diff --git a/src/common.cpp b/src/common.cpp index 17c96c56..0e90a4ae 100644 --- a/src/common.cpp +++ b/src/common.cpp @@ -608,6 +608,11 @@ void print_stats() { } Status daqiri_init(NetworkConfig& config) { + static const MemoryRegionBindings no_bindings; + return daqiri_init(config, no_bindings); +} + +Status daqiri_init(NetworkConfig& config, const MemoryRegionBindings& bindings) { if (g_daqiri_engine != nullptr) { DAQIRI_LOG_ERROR("DAQIRI is already initialized; call shutdown() before daqiri_init()"); return Status::INTERNAL_ERROR; @@ -647,6 +652,10 @@ Status daqiri_init(NetworkConfig& config) { if (config.common_.stream_type == StreamType::SOCKET && (config.common_.protocol == SocketProtocol::TCP || config.common_.protocol == SocketProtocol::UDP)) { + if (!bindings.empty()) { + DAQIRI_LOG_ERROR("External memory bindings are not supported for direct TCP/UDP sockets"); + return Status::NOT_SUPPORTED; + } std::unordered_set gpu_mrs; for (const auto& intf : config.ifs_) { for (const auto& q : intf.rx_.queues_) { @@ -682,10 +691,30 @@ Status daqiri_init(NetworkConfig& config) { } } + if (config.common_.engine_type == EngineType::DPDK) { + for (const auto& [name, binding] : bindings) { + (void)binding; + const auto mr = config.mrs_.find(name); + if (mr != config.mrs_.end() && mr->second.kind_ == MemoryKind::HUGE) { + DAQIRI_LOG_ERROR( + "DPDK cannot safely retain caller-owned hugepage mapping '{}' across EAL cleanup; " + "use kind=host, kind=host_pinned, or an ibverbs/RDMA engine", + name); + return Status::NOT_SUPPORTED; + } + } + } + EngineFactory::set_engine_type(config.common_.engine_type); auto engine = &(EngineFactory::get_active_engine()); + if (engine->set_external_memory_regions(config, bindings) != Status::SUCCESS) { + reset_active_engine(); + metrics::shutdown(); + return Status::INVALID_PARAMETER; + } + if (!engine->set_config_and_initialize(config)) { reset_active_engine(); metrics::shutdown(); @@ -783,25 +812,134 @@ Status parse_network_config(const std::string& yaml_string_or_path, NetworkConfi return parse_network_config_from_yaml_string(yaml_string_or_path, config); } +namespace { + +size_t ceil_power_of_two(size_t value) { + if (value <= 1) { + return 1; + } + --value; + for (size_t shift = 1; shift < sizeof(size_t) * 8; shift <<= 1) { + value |= value >> shift; + } + return value + 1; +} + +size_t ceil_to(size_t value, size_t alignment) { + return (value + alignment - 1) & ~(alignment - 1); +} + +EngineType effective_engine_type(const NetworkConfig& config) { + if (is_explicit_engine_type(config.common_.engine_type)) { + return config.common_.engine_type; + } + if (is_explicit_engine_type(config.common_.engine)) { + return config.common_.engine; + } + return engine_type_from_stream_type(config.common_.stream_type, config.common_.protocol); +} + +} // namespace + +Status get_memory_region_requirements(const NetworkConfig& config, + MemoryRegionRequirements& requirements) { + requirements.clear(); + const EngineType engine = effective_engine_type(config); + if (engine == EngineType::SOCKET && config.common_.protocol != SocketProtocol::ROCE) { + DAQIRI_LOG_ERROR("Direct TCP/UDP socket streams do not use configured memory-region pools"); + return Status::NOT_SUPPORTED; + } + + std::unordered_set queue_mrs; + for (const auto& intf : config.ifs_) { + for (const auto& queue : intf.rx_.queues_) { + queue_mrs.insert(queue.common_.mrs_.begin(), queue.common_.mrs_.end()); + } + for (const auto& queue : intf.tx_.queues_) { + queue_mrs.insert(queue.common_.mrs_.begin(), queue.common_.mrs_.end()); + } + } + + constexpr size_t gpu_page_size = 1UL << 16; + for (const auto& [name, mr] : config.mrs_) { + MemoryRegionRequirement req; + req.kind = mr.kind_; + req.num_bufs = mr.num_bufs_; + switch (engine) { + case EngineType::DPDK: +#if DAQIRI_ENGINE_DPDK + req.slot_size = mr.buf_size_ + RTE_PKTMBUF_HEADROOM; +#else + return Status::NOT_SUPPORTED; +#endif + if (queue_mrs.count(name) != 0 && req.num_bufs < 8192UL * 3 / 2) { + req.num_bufs = 8192UL * 3; + } + break; + case EngineType::IBVERBS: + req.slot_size = ceil_power_of_two(mr.buf_size_); + break; + case EngineType::RDMA: + req.slot_size = ceil_to(mr.buf_size_, gpu_page_size); + break; + default: + DAQIRI_LOG_ERROR("Cannot compute memory requirements for engine {}", + static_cast(engine)); + return Status::NOT_SUPPORTED; + } + if (engine == EngineType::DPDK) { + req.alignment = gpu_page_size; + } else if (engine == EngineType::IBVERBS && mr.kind_ == MemoryKind::DEVICE) { + req.alignment = 4096; + } else { + req.alignment = mr.kind_ == MemoryKind::DEVICE ? 256 : 128; + } + if (req.slot_size != 0 && req.num_bufs > std::numeric_limits::max() / req.slot_size) { + DAQIRI_LOG_ERROR("Memory region '{}' size overflows", name); + return Status::INVALID_PARAMETER; + } + req.capacity = ceil_to(req.slot_size * req.num_bufs, gpu_page_size); + requirements.emplace(name, req); + } + return Status::SUCCESS; +} + Status daqiri_init_from_yaml_string(const std::string& yaml_string) { + static const MemoryRegionBindings no_bindings; + return daqiri_init_from_yaml_string(yaml_string, no_bindings); +} + +Status daqiri_init_from_yaml_string(const std::string& yaml_string, + const MemoryRegionBindings& bindings) { NetworkConfig config; const Status parse_status = parse_network_config_from_yaml_string(yaml_string, config); if (parse_status != Status::SUCCESS) { return parse_status; } - return daqiri_init(config); + return daqiri_init(config, bindings); } Status daqiri_init_from_yaml_file(const std::string& yaml_path) { + static const MemoryRegionBindings no_bindings; + return daqiri_init_from_yaml_file(yaml_path, no_bindings); +} + +Status daqiri_init_from_yaml_file(const std::string& yaml_path, + const MemoryRegionBindings& bindings) { NetworkConfig config; const Status parse_status = parse_network_config_from_yaml_file(yaml_path, config); if (parse_status != Status::SUCCESS) { return parse_status; } - return daqiri_init(config); + return daqiri_init(config, bindings); } Status daqiri_init(const std::string& yaml_string_or_path) { + static const MemoryRegionBindings no_bindings; + return daqiri_init(yaml_string_or_path, no_bindings); +} + +Status daqiri_init(const std::string& yaml_string_or_path, const MemoryRegionBindings& bindings) { NetworkConfig config; const Status parse_status = parse_network_config(yaml_string_or_path, config); if (parse_status != Status::SUCCESS) { return parse_status; } - return daqiri_init(config); + return daqiri_init(config, bindings); } // Generic socket functions diff --git a/src/engine.cpp b/src/engine.cpp index afe91a02..8958bbfe 100644 --- a/src/engine.cpp +++ b/src/engine.cpp @@ -36,6 +36,7 @@ #include #include #include +#include #include #include #include @@ -360,6 +361,117 @@ size_t Engine::get_alignment(MemoryKind kind) { } } +Status Engine::set_external_memory_regions(const NetworkConfig& cfg, + const MemoryRegionBindings& bindings) { + external_mrs_.clear(); + std::vector> ranges; + ranges.reserve(bindings.size()); + + for (const auto& [name, binding] : bindings) { + const auto mr_it = cfg.mrs_.find(name); + if (mr_it == cfg.mrs_.end()) { + DAQIRI_LOG_ERROR("External memory binding '{}' does not name a configured memory region", + name); + return Status::INVALID_PARAMETER; + } + if (binding.data == nullptr || binding.capacity == 0) { + DAQIRI_LOG_ERROR("External memory binding '{}' has a null pointer or zero capacity", name); + return Status::INVALID_PARAMETER; + } + const uintptr_t begin = reinterpret_cast(binding.data); + if (binding.capacity > std::numeric_limits::max() - begin) { + DAQIRI_LOG_ERROR("External memory binding '{}' address range overflows", name); + return Status::INVALID_PARAMETER; + } + const uintptr_t end = begin + binding.capacity; + for (const auto& range : ranges) { + if (begin < range.second && range.first < end) { + DAQIRI_LOG_ERROR("External memory binding '{}' overlaps another bound region", name); + return Status::INVALID_PARAMETER; + } + } + ranges.emplace_back(begin, end); + + ResolvedExternalMemoryRegion resolved{binding.data, binding.capacity, nullptr, -1}; + const MemoryKind kind = mr_it->second.kind_; + if (kind == MemoryKind::DEVICE || kind == MemoryKind::HOST_PINNED) { + CUmemorytype memory_type{}; + CUcontext context = nullptr; + const CUdeviceptr ptr = reinterpret_cast(binding.data); + if (cuPointerGetAttribute(&context, CU_POINTER_ATTRIBUTE_CONTEXT, ptr) != CUDA_SUCCESS || + cuPointerGetAttribute(&memory_type, CU_POINTER_ATTRIBUTE_MEMORY_TYPE, ptr) != + CUDA_SUCCESS || + context == nullptr) { + DAQIRI_LOG_ERROR("External memory binding '{}' is not CUDA-registered memory", name); + return Status::INVALID_PARAMETER; + } + const CUmemorytype expected = + kind == MemoryKind::DEVICE ? CU_MEMORYTYPE_DEVICE : CU_MEMORYTYPE_HOST; + if (memory_type != expected) { + DAQIRI_LOG_ERROR("External memory binding '{}' does not match configured memory kind", + name); + return Status::INVALID_PARAMETER; + } + resolved.cuda_context = context; + { + CudaContextGuard guard(context); + CUdevice context_device = -1; + if (!guard.valid() || cuCtxGetDevice(&context_device) != CUDA_SUCCESS) { + DAQIRI_LOG_ERROR("Could not determine the CUDA context device for '{}'", name); + return Status::INVALID_PARAMETER; + } + resolved.cuda_device = context_device; + } + if (kind == MemoryKind::DEVICE) { + CUdeviceptr range_start = 0; + size_t range_size = 0; + if (cuPointerGetAttribute(&range_start, CU_POINTER_ATTRIBUTE_RANGE_START_ADDR, ptr) != + CUDA_SUCCESS || + cuPointerGetAttribute(&range_size, CU_POINTER_ATTRIBUTE_RANGE_SIZE, ptr) != + CUDA_SUCCESS || + ptr < range_start || (ptr - range_start) > range_size || + binding.capacity > range_size - (ptr - range_start)) { + DAQIRI_LOG_ERROR( + "External device binding '{}' capacity exceeds its CUDA allocation range", name); + return Status::INVALID_PARAMETER; + } + int ordinal = -1; + if (cuPointerGetAttribute(&ordinal, CU_POINTER_ATTRIBUTE_DEVICE_ORDINAL, ptr) != + CUDA_SUCCESS || + ordinal != mr_it->second.affinity_) { + DAQIRI_LOG_ERROR( + "External device binding '{}' is on CUDA device {}, configured affinity is {}", name, + ordinal, mr_it->second.affinity_); + return Status::INVALID_PARAMETER; + } + int capable = 0; + if (cuPointerGetAttribute(&capable, CU_POINTER_ATTRIBUTE_IS_GPU_DIRECT_RDMA_CAPABLE, ptr) == + CUDA_SUCCESS && + capable == 0) { + DAQIRI_LOG_ERROR("External device binding '{}' is not GPUDirect RDMA capable", name); + return Status::INVALID_PARAMETER; + } + CudaContextGuard guard(context); + unsigned int sync_memops = 1; + if (!guard.valid() || cuPointerSetAttribute(&sync_memops, CU_POINTER_ATTRIBUTE_SYNC_MEMOPS, + ptr) != CUDA_SUCCESS) { + DAQIRI_LOG_ERROR("Could not enable synchronous memory operations for '{}'", name); + return Status::INVALID_PARAMETER; + } + } + } + external_mrs_.emplace(name, resolved); + } + for (const auto& [name, mr] : cfg.mrs_) { + if (!mr.owned_ && external_mrs_.find(name) == external_mrs_.end()) { + DAQIRI_LOG_ERROR("Memory region '{}' has owned=false but no external binding", name); + external_mrs_.clear(); + return Status::INVALID_PARAMETER; + } + } + return Status::SUCCESS; +} + Status Engine::populate_pool(daqiri::Ring* ring, const std::string& mr_name) { auto mr = cfg_.mrs_[mr_name]; auto base = reinterpret_cast(ar_[mr_name].ptr_); @@ -438,7 +550,37 @@ Status Engine::allocate_memory_regions() { mr.second.ttl_size_ = align_ceil(mr.second.adj_size_ * mr.second.num_bufs_, GPU_PAGE_SIZE); ar.size_ = mr.second.ttl_size_; - if (mr.second.owned_) { + const auto external = external_mrs_.find(mr.first); + if (external != external_mrs_.end()) { + if (external->second.capacity < mr.second.ttl_size_) { + DAQIRI_LOG_ERROR( + "External memory region '{}' supplies {} bytes but requires {} bytes " + "({} buffers at a {} byte stride)", + mr.first, external->second.capacity, mr.second.ttl_size_, mr.second.num_bufs_, + mr.second.adj_size_); + return Status::INVALID_PARAMETER; + } + size_t alignment = get_alignment(mr.second.kind_); + if (cfg_.common_.engine_type == EngineType::DPDK) { + alignment = GPU_PAGE_SIZE; + } else if (cfg_.common_.engine_type == EngineType::IBVERBS && + mr.second.kind_ == MemoryKind::DEVICE) { + alignment = static_cast(sysconf(_SC_PAGESIZE)); + } + if ((reinterpret_cast(external->second.data) % alignment) != 0) { + DAQIRI_LOG_ERROR("External memory region '{}' pointer is not {}-byte aligned", mr.first, + alignment); + return Status::INVALID_PARAMETER; + } + ptr = external->second.data; + ar.external_ = true; + ar.cuda_context_ = external->second.cuda_context; + ar.cuda_device_ = external->second.cuda_device; + ar.deallocator_ = AllocRegion::Deallocator::NONE; + } else if (!mr.second.owned_) { + DAQIRI_LOG_ERROR("Memory region '{}' has owned=false but no external binding", mr.first); + return Status::INVALID_PARAMETER; + } else { switch (mr.second.kind_) { case MemoryKind::HOST: if (posix_memalign(&ptr, GPU_PAGE_SIZE, mr.second.ttl_size_) != 0) { @@ -453,6 +595,12 @@ Status Engine::allocate_memory_regions() { DAQIRI_LOG_CRITICAL("Failed to allocate CUDA pinned host memory!"); return Status::NULL_PTR; } + if (cuCtxGetCurrent(&ar.cuda_context_) != CUDA_SUCCESS || ar.cuda_context_ == nullptr) { + DAQIRI_LOG_CRITICAL("Failed to capture CUDA context for pinned host memory!"); + cudaFreeHost(ptr); + return Status::NULL_PTR; + } + ar.cuda_device_ = mr.second.affinity_; ar.deallocator_ = AllocRegion::Deallocator::CUDA_HOST; break; case MemoryKind::HUGE: @@ -464,6 +612,14 @@ Status Engine::allocate_memory_regions() { CUdeviceptr cuptr; CUcontext current = nullptr; + const auto driver_init_res = cuInit(0); + if (driver_init_res != CUDA_SUCCESS) { + const char* err_str = nullptr; + cuGetErrorString(driver_init_res, &err_str); + DAQIRI_LOG_CRITICAL("Could not initialize the CUDA driver: {}", + err_str != nullptr ? err_str : "unknown error"); + return Status::NULL_PTR; + } const auto current_res = cuCtxGetCurrent(¤t); if (current_res != CUDA_SUCCESS) { DAQIRI_LOG_CRITICAL("Could not query the current CUDA context"); @@ -529,6 +685,7 @@ Status Engine::allocate_memory_regions() { restore_previous(); return Status::NULL_PTR; } + ar.cuda_device_ = mr.second.affinity_; ar.deallocator_ = AllocRegion::Deallocator::CUDA_DEVICE; if (!restore_previous()) { if (cuCtxSetCurrent(ar.cuda_context_) == CUDA_SUCCESS) { diff --git a/src/engine.h b/src/engine.h index 4a055a9b..38e00cdd 100644 --- a/src/engine.h +++ b/src/engine.h @@ -18,6 +18,7 @@ #pragma once #include "src/daqiri_ring.h" +#include #include #include #include @@ -48,7 +49,39 @@ struct AllocRegion { size_t size_ = 0; int affinity_ = -1; CUcontext cuda_context_ = nullptr; + int cuda_device_ = -1; Deallocator deallocator_ = Deallocator::NONE; + bool external_ = false; +}; + +struct ResolvedExternalMemoryRegion { + void* data = nullptr; + size_t capacity = 0; + CUcontext cuda_context = nullptr; + int cuda_device = -1; +}; + +class CudaContextGuard { + public: + explicit CudaContextGuard(CUcontext context) : active_(context != nullptr) { + if (active_ && cuCtxPushCurrent(context) != CUDA_SUCCESS) { + active_ = false; + } + } + ~CudaContextGuard() { + if (active_) { + CUcontext popped = nullptr; + (void)cuCtxPopCurrent(&popped); + } + } + CudaContextGuard(const CudaContextGuard&) = delete; + CudaContextGuard& operator=(const CudaContextGuard&) = delete; + bool valid() const { + return active_; + } + + private: + bool active_ = false; }; /** @@ -62,6 +95,8 @@ class Engine { virtual void initialize() = 0; virtual bool is_initialized() const { return initialized_; } virtual bool set_config_and_initialize(const NetworkConfig& cfg) = 0; + Status set_external_memory_regions(const NetworkConfig& cfg, + const MemoryRegionBindings& bindings); virtual void run() = 0; // Common free functions to override @@ -174,6 +209,7 @@ class Engine { bool initialized_ = false; NetworkConfig cfg_; std::unordered_map ar_; + std::unordered_map external_mrs_; // shared_ptr to an incomplete type -- only populated by the DPDK engine // (engine_dpdk.cpp). Layout is identical in every build. std::unordered_map> ext_pktmbufs_; diff --git a/src/engine_dpdk.cpp b/src/engine_dpdk.cpp index 3196c555..f1be6e0d 100644 --- a/src/engine_dpdk.cpp +++ b/src/engine_dpdk.cpp @@ -272,8 +272,10 @@ Status Engine::register_memory_regions() { for (const auto& ar : ar_) { const auto& mr = cfg_.mrs_[ar.second.mr_name_]; - // Hugepages use the normal rte functions that don't require extmem - if (mr.kind_ == MemoryKind::HUGE) { + // DAQIRI-owned hugepages use the normal EAL-backed pool. Caller-owned + // hugepages must be registered as external memory so the pool uses the + // supplied address. + if (mr.kind_ == MemoryKind::HUGE && !ar.second.external_) { continue; } @@ -285,6 +287,13 @@ Status Engine::register_memory_regions() { int ret = 0; if (mr.kind_ == MemoryKind::DEVICE) { + const auto resolved = external_mrs_.find(mr.name_); + CudaContextGuard context_guard( + resolved == external_mrs_.end() ? nullptr : resolved->second.cuda_context); + if (resolved != external_mrs_.end() && !context_guard.valid()) { + DAQIRI_LOG_CRITICAL("Could not make the CUDA context for MR {} current", mr.name_); + return Status::INVALID_PARAMETER; + } int flag = 0; CUresult s = cuDeviceGetAttribute(&flag, CU_DEVICE_ATTRIBUTE_DMA_BUF_SUPPORTED, mr.affinity_); if (s != CUDA_SUCCESS) { @@ -374,7 +383,9 @@ struct rte_mempool* Engine::create_pktmbuf_pool(const std::string& name, struct rte_mempool* pool; #pragma GCC diagnostic push #pragma GCC diagnostic ignored "-Wdeprecated-declarations" - if (mr.kind_ == MemoryKind::HUGE) { + const auto ar_it = ar_.find(mr.name_); + const bool external = ar_it != ar_.end() && ar_it->second.external_; + if (mr.kind_ == MemoryKind::HUGE && !external) { pool = rte_pktmbuf_pool_create(name.c_str(), mr.num_bufs_, 0, 0, mr.adj_size_, numa_from_mem(mr)); } else { diff --git a/src/engines/dpdk/daqiri_dpdk_engine.cpp b/src/engines/dpdk/daqiri_dpdk_engine.cpp index 7a6f7961..8f7e2efe 100644 --- a/src/engines/dpdk/daqiri_dpdk_engine.cpp +++ b/src/engines/dpdk/daqiri_dpdk_engine.cpp @@ -729,6 +729,7 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ } int cuda_device_id = 0; + CUcontext cuda_context = nullptr; if (use_gpu_backend) { if (copy_src_mr->kind_ == MemoryKind::DEVICE) { cuda_device_id = copy_src_mr->affinity_; @@ -737,6 +738,29 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ } else { cuda_device_id = copy_src_mr->affinity_; } + const auto select_region_context = [&](const std::string& mr_name, const char* role) { + const auto& region = ar_.at(mr_name); + if (region.cuda_device_ != cuda_device_id) { + DAQIRI_LOG_ERROR( + "Reorder '{}' {} MR '{}' uses CUDA device {}, but the plan uses device {}", + reorder_cfg.name_, + role, + mr_name, + region.cuda_device_, + cuda_device_id); + return false; + } + if (cuda_context != nullptr && region.cuda_context_ != nullptr && + cuda_context != region.cuda_context_) { + DAQIRI_LOG_ERROR("Reorder '{}' source and output use different CUDA contexts", + reorder_cfg.name_); + return false; + } + cuda_context = region.cuda_context_; + return true; + }; + if (!select_region_context(source_mr_name, "source") || + !select_region_context(out_mr.name_, "output")) { return false; } } auto pool_it = reorder_output_pools_.find(out_mr.name_); @@ -767,9 +791,15 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ buffer.event_complete = true; }; const auto enable_cuda_events = [&destroy_cuda_events](ReorderOutputPool& pool, - int device_id) -> bool { + int device_id, + CUcontext context) -> bool { if (pool.cuda_events_enabled) { return true; } - cudaSetDevice(device_id); + CudaContextGuard context_guard(context); + if (context == nullptr) { + cudaSetDevice(device_id); + } else if (!context_guard.valid()) { + return false; + } for (size_t i = 0; i < pool.buffers.size(); ++i) { auto& buffer = pool.buffers[i]; if (cudaEventCreateWithFlags(&buffer.event, cudaEventDisableTiming) != cudaSuccess) { @@ -814,7 +844,8 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ auto* buf_ptr = base + (i * out_mr.adj_size_); output_pool->buffers[i].ptr = buf_ptr; } - if (use_gpu_backend && !enable_cuda_events(*output_pool, cuda_device_id)) { + output_pool->cuda_context = cuda_context; + if (use_gpu_backend && !enable_cuda_events(*output_pool, cuda_device_id, cuda_context)) { DAQIRI_LOG_ERROR("Failed to create CUDA events for reorder output MR '{}'", out_mr.name_); return false; } @@ -831,7 +862,11 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ cuda_device_id); return false; } - if (use_gpu_backend && !enable_cuda_events(*output_pool, cuda_device_id)) { + if (use_gpu_backend && output_pool->cuda_context != cuda_context) { + DAQIRI_LOG_ERROR("Reorder output MR '{}' is shared across CUDA contexts", out_mr.name_); + return false; + } + if (use_gpu_backend && !enable_cuda_events(*output_pool, cuda_device_id, cuda_context)) { DAQIRI_LOG_ERROR("Failed to create CUDA events for reorder output MR '{}'", out_mr.name_); return false; } @@ -863,13 +898,19 @@ bool DpdkEngine::init_reorder_queue_state(const InterfaceConfig& intf, const RxQ plan.h_input_ptrs.resize(plan.cuda_staging_capacity); plan.h_source_mbufs.resize(plan.cuda_staging_capacity); plan.cuda_device_id = cuda_device_id; + plan.cuda_context = cuda_context; plan.timeout_cycles = (qcfg.timeout_us_ == 0) ? 0 : (static_cast(qcfg.timeout_us_) * rte_get_timer_hz() / 1000000ULL); if (use_gpu_backend) { - cudaSetDevice(plan.cuda_device_id); + CudaContextGuard context_guard(plan.cuda_context); + if (plan.cuda_context == nullptr) { + cudaSetDevice(plan.cuda_device_id); + } else if (!context_guard.valid()) { + return false; + } if (cudaMalloc(reinterpret_cast(&plan.d_input_ptrs), sizeof(void*) * plan.cuda_staging_capacity) != cudaSuccess) { DAQIRI_LOG_ERROR("Failed to allocate CUDA staging buffers for reorder config '{}'", @@ -934,6 +975,7 @@ void DpdkEngine::cleanup_reorder_state() { for (auto& [qkey, qstate] : reorder_queue_states_) { (void)qkey; for (auto& plan : qstate.plans) { + CudaContextGuard context_guard(plan.cuda_context); #if DAQIRI_REORDER_GPU_PROFILE if (plan.gpu_profile.gpu_kernel_samples != 0) { const double avg_kernel_us = @@ -992,7 +1034,10 @@ void DpdkEngine::cleanup_reorder_state() { for (auto& [mr_name, output_pool] : reorder_output_pools_) { (void)mr_name; if (output_pool == nullptr) { continue; } - if (output_pool->cuda_events_enabled) { cudaSetDevice(output_pool->cuda_device_id); } + CudaContextGuard context_guard(output_pool->cuda_context); + if (output_pool->cuda_context == nullptr && output_pool->cuda_events_enabled) { + cudaSetDevice(output_pool->cuda_device_id); + } for (auto& buffer : output_pool->buffers) { if (buffer.event != nullptr) { cudaEventDestroy(buffer.event); @@ -1063,6 +1108,10 @@ void DpdkEngine::release_reorder_output_buffer(std::shared_ptrh_batch_id != nullptr) { + CudaContextGuard context_guard(ctx->output_pool != nullptr ? ctx->output_pool->cuda_context + : nullptr); if (burst->event != nullptr) { const cudaError_t event_status = cudaEventQuery(burst->event); if (event_status == cudaErrorNotReady) { return Status::NOT_READY; } @@ -1732,7 +1787,21 @@ Status DpdkEngine::set_reorder_cuda_stream(const std::string& interface_name, } const auto mr_it = cfg_.mrs_.find(plan.memory_region_name); if (mr_it == cfg_.mrs_.end()) { return Status::INVALID_PARAMETER; } - cudaSetDevice(plan.cuda_device_id); + CudaContextGuard context_guard(plan.cuda_context); + if (plan.cuda_context == nullptr) { + cudaSetDevice(plan.cuda_device_id); + } else if (!context_guard.valid()) { + return Status::INTERNAL_ERROR; + } + if (stream != nullptr && plan.cuda_context != nullptr) { + CUcontext stream_context = nullptr; + if (cuStreamGetCtx(reinterpret_cast(stream), &stream_context) != CUDA_SUCCESS || + stream_context != plan.cuda_context) { + DAQIRI_LOG_ERROR("CUDA stream for reorder '{}' belongs to a different context", + reorder_name); + return Status::INVALID_PARAMETER; + } + } plan.stream = stream; found = true; } diff --git a/src/engines/dpdk/daqiri_dpdk_engine.h b/src/engines/dpdk/daqiri_dpdk_engine.h index 147c71b7..4ec22341 100644 --- a/src/engines/dpdk/daqiri_dpdk_engine.h +++ b/src/engines/dpdk/daqiri_dpdk_engine.h @@ -382,6 +382,7 @@ class DpdkEngine : public Engine { std::vector buffers; size_t next_buffer = 0; int cuda_device_id = 0; + CUcontext cuda_context = nullptr; bool cuda_events_enabled = false; }; @@ -415,6 +416,7 @@ class DpdkEngine : public Engine { uint64_t timeout_cycles = 0; bool use_gpu_backend = false; int cuda_device_id = 0; + CUcontext cuda_context = nullptr; cudaStream_t stream = nullptr; uint32_t cuda_staging_capacity = 0; diff --git a/src/engines/ibverbs/daqiri_ibverbs_engine.cpp b/src/engines/ibverbs/daqiri_ibverbs_engine.cpp index 4f76c032..96cf31b2 100644 --- a/src/engines/ibverbs/daqiri_ibverbs_engine.cpp +++ b/src/engines/ibverbs/daqiri_ibverbs_engine.cpp @@ -605,7 +605,6 @@ Status IbverbsEngine::register_mr(struct ibv_pd* pd, const std::string& mr_name, DAQIRI_LOG_CRITICAL("Could not activate the owning CUDA context for MR {}", mr_name); return Status::INTERNAL_ERROR; } - const size_t page = sysconf(_SC_PAGESIZE); const auto va = reinterpret_cast(base); const uintptr_t aligned = va & ~(static_cast(page) - 1); @@ -3195,10 +3194,53 @@ Status IbverbsEngine::init_reorder(IbvRxQueue& q, const InterfaceConfig& intf, plan.copy_source_offset = rc.payload_byte_offset_; plan.slot_stride = static_cast(src_mr.buf_size_ - rc.payload_byte_offset_); plan.data_type_conversion = rdr_uses_conversion(rc); - plan.cuda_device_id = src_mr.affinity_; + if (src_mr.kind_ == MemoryKind::DEVICE && out_mr.kind_ == MemoryKind::DEVICE && + src_mr.affinity_ != out_mr.affinity_) { + DAQIRI_LOG_CRITICAL( + "Reorder '{}' requires input/output memory on the same GPU (src affinity {} " + "!= reorder affinity {})", + rc.name_, src_mr.affinity_, out_mr.affinity_); + return Status::INVALID_PARAMETER; + } + if (src_mr.kind_ == MemoryKind::DEVICE) { + plan.cuda_device_id = src_mr.affinity_; + } else if (out_mr.kind_ == MemoryKind::DEVICE) { + plan.cuda_device_id = out_mr.affinity_; + } else { + plan.cuda_device_id = src_mr.affinity_; + } + const auto select_region_context = [&](const std::string& mr_name, const char* role) { + const auto& region = ar_.at(mr_name); + if (region.cuda_device_ != plan.cuda_device_id) { + DAQIRI_LOG_CRITICAL( + "Reorder '{}' {} MR '{}' uses CUDA device {}, but the plan uses device {}", + rc.name_, + role, + mr_name, + region.cuda_device_, + plan.cuda_device_id); + return Status::INVALID_PARAMETER; + } + if (plan.cuda_context != nullptr && region.cuda_context_ != nullptr && + plan.cuda_context != region.cuda_context_) { + DAQIRI_LOG_CRITICAL("Reorder '{}' source and output use different CUDA contexts", rc.name_); + return Status::INVALID_PARAMETER; + } + plan.cuda_context = region.cuda_context_; + return Status::SUCCESS; + }; + if (select_region_context(q.mr_name, "source") != Status::SUCCESS || + select_region_context(rc.memory_region_, "output") != Status::SUCCESS) { + return Status::INVALID_PARAMETER; + } plan.acc_ptrs.reserve(plan.packets_per_batch); - cudaSetDevice(plan.cuda_device_id); + CudaContextGuard context_guard(plan.cuda_context); + if (plan.cuda_context == nullptr) { + cudaSetDevice(plan.cuda_device_id); + } else if (!context_guard.valid()) { + return Status::INTERNAL_ERROR; + } if (cudaMalloc(reinterpret_cast(&plan.d_input_ptrs), sizeof(void*) * plan.packets_per_batch) != cudaSuccess) { DAQIRI_LOG_CRITICAL("Reorder '{}' cudaMalloc(d_input_ptrs) failed", rc.name_); @@ -3299,7 +3341,12 @@ Status IbverbsEngine::reorder_flush_batch(IbvRxQueue& q, IbvReorderPlan& plan, B const uint32_t out_payload = plan.acc_output_payload_len; const uint32_t aggregate_len = plan.packets_per_batch * out_payload; - cudaSetDevice(plan.cuda_device_id); + CudaContextGuard context_guard(plan.cuda_context); + if (plan.cuda_context == nullptr) { + cudaSetDevice(plan.cuda_device_id); + } else if (!context_guard.valid()) { + return Status::INTERNAL_ERROR; + } cudaMemcpyAsync(plan.d_input_ptrs, plan.acc_ptrs.data(), sizeof(void*) * num_pkts, cudaMemcpyHostToDevice, plan.stream); packet_reorder_copy_payload_by_sequence( @@ -3438,6 +3485,7 @@ void IbverbsEngine::reorder_cleanup(IbvRxQueue& q) { return; } for (auto& plan : q.reorder->plans) { + CudaContextGuard context_guard(plan.cuda_context); for (auto& ob : plan.out_bufs) { if (ob.event) { cudaEventDestroy(ob.event); @@ -3466,6 +3514,15 @@ Status IbverbsEngine::set_reorder_cuda_stream(const std::string& interface_name, } for (auto& plan : q->reorder->plans) { if (plan.cfg.name_ == reorder_name) { + if (stream != nullptr && plan.cuda_context != nullptr) { + CUcontext stream_context = nullptr; + if (cuStreamGetCtx(reinterpret_cast(stream), &stream_context) != CUDA_SUCCESS || + stream_context != plan.cuda_context) { + DAQIRI_LOG_ERROR("CUDA stream for reorder '{}' belongs to a different context", + reorder_name); + return Status::INVALID_PARAMETER; + } + } plan.stream = stream; DAQIRI_LOG_INFO("Reorder '{}' stream set on port {} q{}", reorder_name, port, q->queue_id); return Status::SUCCESS; @@ -3490,6 +3547,7 @@ Status IbverbsEngine::get_reorder_burst_info(BurstParams* burst, ReorderBurstInf return Status::INVALID_PARAMETER; } if (ctx->h_batch_id != nullptr) { + CudaContextGuard context_guard(ctx->plan != nullptr ? ctx->plan->cuda_context : nullptr); if (burst->event != nullptr) { const cudaError_t s = cudaEventQuery(burst->event); if (s == cudaErrorNotReady) { diff --git a/src/engines/ibverbs/daqiri_ibverbs_engine.h b/src/engines/ibverbs/daqiri_ibverbs_engine.h index 4d71c4b4..c3daf715 100644 --- a/src/engines/ibverbs/daqiri_ibverbs_engine.h +++ b/src/engines/ibverbs/daqiri_ibverbs_engine.h @@ -71,6 +71,7 @@ struct IbvReorderPlan { uint32_t slot_stride = 0; bool data_type_conversion = false; int cuda_device_id = 0; + CUcontext cuda_context = nullptr; cudaStream_t stream = nullptr; void** d_input_ptrs = nullptr; std::vector out_bufs;