From 90ed9db7b353d47ee2ed0c768e663c1804d65a71 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Wed, 21 Jan 2026 14:11:27 +0530 Subject: [PATCH 01/28] Add support to prepare Fifo Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 120 +++++++++++++++--------------- src/plugins/uccl/uccl_backend.h | 5 +- 2 files changed, 62 insertions(+), 63 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index ba558d8472..bc18f086a6 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -177,7 +177,16 @@ nixlUcclEngine::getSupportedMems() const { nixl_status_t nixlUcclEngine::getPublicData(const nixlBackendMD *meta, std::string &str) const { nixlUcclBackendMD *priv = (nixlUcclBackendMD *)meta; - str = std::to_string(priv->mr_id); + + // Export fifo_item as hex string + str.clear(); + str.reserve(FIFO_ITEM_SIZE * 2); + for (int i = 0; i < FIFO_ITEM_SIZE; i++) { + char hex[3]; + snprintf(hex, sizeof(hex), "%02x", static_cast(priv->fifo_item[i])); + str += hex; + } + NIXL_DEBUG << "Exporting Meta Info =" << hex_str << sts::endl; return NIXL_SUCCESS; } @@ -292,6 +301,17 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, priv->length = mem.len; priv->ref_cnt = 1; priv->mr_id = reinterpret_cast(mr); // Store the memory region handle + + // Pre-compute fifo_item for one-sided RDMA operations + int result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, + mem.len, priv->fifo_item); + if (result != 0) { + NIXL_ERROR << "Failed to prepare fifo_item for memory region"; + uccl_engine_mr_destroy(mr); + delete priv; + return NIXL_ERR_BACKEND; + } + out = priv; mem_reg_info_[mem.addr] = priv; NIXL_DEBUG << "Registering memory: " << mem.addr << "Device: " << mem.devId @@ -349,7 +369,23 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, output_md->addr = (void *)input.addr; output_md->length = input.len; output_md->ref_cnt = 1; - output_md->mr_id = strtoul(input.metaInfo.c_str(), NULL, 10); + + // Decode fifo_item from hex string + const std::string &hex_str = input.metaInfo; + NIXL_DEBUG << "Meta Info =" << hex_str << sts::endl; + if (hex_str.length() == FIFO_ITEM_SIZE * 2) { + for (int i = 0; i < FIFO_ITEM_SIZE; i++) { + std::string byte_str = hex_str.substr(i * 2, 2); + output_md->fifo_item[i] = static_cast(strtoul(byte_str.c_str(), NULL, 16)); + } + NIXL_DEBUG << "Parsed fifo_item from remote metadata"; + } else { + NIXL_ERROR << "Invalid fifo_item hex string length: " << hex_str.length() + << " (expected " << FIFO_ITEM_SIZE * 2 << ")"; + delete output_md; + output = nullptr; + return NIXL_ERR_INVALID_PARAM; + } return NIXL_SUCCESS; } @@ -369,7 +405,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, const std::string &remote_agent, nixlBackendReqH *&handle, const nixl_opt_b_args_t *opt_args) const { - int result = 0; nixlUcclBackendMD *lmd; nixlUcclBackendMD *rmd; bool rcmode = false; @@ -411,9 +446,11 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_INVALID_PARAM; } } - // Collect all tx_data into vectors for batch sending - std::vector md_vector; - std::vector local_priv_vector; + handle = new nixlUcclReqH(conn); + nixlUcclReqH *uccl_handle = static_cast(handle); + + uccl_handle->fifo_items.clear(); + uccl_handle->fifo_items.resize(lcnt); std::lock_guard lock(mem_mutex_); for (size_t i = 0; i < lcnt; i++) { @@ -441,62 +478,21 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, NIXL_ERROR << "Local memory region not properly registered"; return NIXL_ERR_BACKEND; } - - // Prepare the memory region metadata for batch sending - md_t md; - tx_msg_t tx_data; - tx_data.data_ptr = remote_addr; - tx_data.data_size = rsize; - - // RC mode is supported for both READ/WRITE operations - // UC mode is supported only for WRITE operations - md.op = rcmode ? UCCL_RW_RC : UCCL_WRITE; - md.data.tx_data = tx_data; - - // Add to vectors for batch processing - md_vector.push_back(md); - local_priv_vector.push_back(local_priv); - } - - // Send all tx_data as a vector - result = uccl_engine_send_tx_md_vector(conn, md_vector.data(), md_vector.size()); - if (result < 0) { - NIXL_ERROR << "Failed to send transfer metadata vector"; - return NIXL_ERR_BACKEND; - } - - if (rcmode) { - if (!handle) { - handle = new nixlUcclReqH(conn); - } - nixlUcclReqH *uccl_handle = static_cast(handle); - - uccl_handle->fifo_items.clear(); - uccl_handle->fifo_items.resize(local_priv_vector.size()); - - for (size_t i = 0; i < local_priv_vector.size(); i++) { - char fifo_item[FIFO_ITEM_SIZE]; - int retry_count = 0; - const int max_retries = 50; - do { - result = uccl_engine_get_fifo_item(conn, i, &fifo_item); - if (result == 0) { - memcpy(uccl_handle->fifo_items[i].data(), fifo_item, FIFO_ITEM_SIZE); - break; - } - retry_count++; - if (retry_count < max_retries) { - NIXL_DEBUG << "Failed to get FIFO item, retry " << retry_count << "/" - << max_retries << " for item " << i; - std::this_thread::sleep_for(std::chrono::microseconds(10)); - } - } while (retry_count < max_retries); - - if (result != 0) { - NIXL_ERROR << "Failed to get FIFO item after " << max_retries - << " retries for item " << i; - return NIXL_ERR_BACKEND; - } + if (rcmode) { + memcpy(uccl_handle->fifo_items[i].data(), rmd->fifo_item, FIFO_ITEM_SIZE); + + // Update the address and size in the fifo_item for the actual transfer + // FifoItem layout: addr (8 bytes at offset 0), size (4 bytes at offset 8) + uint64_t *fifo_addr = reinterpret_cast(uccl_handle->fifo_items[i].data()); + uint32_t *fifo_size = reinterpret_cast(uccl_handle->fifo_items[i].data() + 8); + *fifo_addr = remote_addr; + *fifo_size = static_cast(rsize); + + NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex + << remote_addr << ", size=" << std::dec << rsize; + } else { + // TODO : Suppport UC one-sided + return NIXL_ERR_BACKEND; } } diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index 3348add3d2..4cb3340393 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -139,7 +139,9 @@ class nixlUcclEngine : public nixlBackendEngine { // UCCL Backend Memory Descriptor class nixlUcclBackendMD : public nixlBackendMD { public: - nixlUcclBackendMD(bool isPrivate) : nixlBackendMD(isPrivate) {} + nixlUcclBackendMD(bool isPrivate) : nixlBackendMD(isPrivate) { + memset(fifo_item, 0, FIFO_ITEM_SIZE); + } virtual ~nixlUcclBackendMD() {} @@ -147,6 +149,7 @@ class nixlUcclBackendMD : public nixlBackendMD { size_t length; int ref_cnt; uint64_t mr_id; // UCCL memory region id + char fifo_item[FIFO_ITEM_SIZE]; }; // UCCL Backend Request Handle From 8a354ea45dd11caf8100582ddbbff309f823fa9f Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Wed, 21 Jan 2026 14:20:22 +0530 Subject: [PATCH 02/28] Fix errors Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index bc18f086a6..05b2bf5f95 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -186,7 +186,7 @@ nixlUcclEngine::getPublicData(const nixlBackendMD *meta, std::string &str) const snprintf(hex, sizeof(hex), "%02x", static_cast(priv->fifo_item[i])); str += hex; } - NIXL_DEBUG << "Exporting Meta Info =" << hex_str << sts::endl; + NIXL_DEBUG << "Exporting Meta Info =" << str << std::endl; return NIXL_SUCCESS; } @@ -372,7 +372,7 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, // Decode fifo_item from hex string const std::string &hex_str = input.metaInfo; - NIXL_DEBUG << "Meta Info =" << hex_str << sts::endl; + NIXL_DEBUG << "Meta Info =" << hex_str << std::endl; if (hex_str.length() == FIFO_ITEM_SIZE * 2) { for (int i = 0; i < FIFO_ITEM_SIZE; i++) { std::string byte_str = hex_str.substr(i * 2, 2); From a9bcd2bafbe18e1b6a6ed74a0bd718fe265ffd67 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 22 Jan 2026 15:08:34 +0530 Subject: [PATCH 03/28] Add vector read/write support Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 156 ++++++++++++++++-------------- src/plugins/uccl/uccl_backend.h | 2 +- 2 files changed, 85 insertions(+), 73 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 05b2bf5f95..9aa943d970 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -107,10 +107,7 @@ nixlUcclEngine::~nixlUcclEngine() { std::lock_guard lock(mem_mutex_); for (auto &[addr, priv] : mem_reg_info_) { if (priv && priv->mr_id != 0) { - uccl_mr_t *mr = reinterpret_cast(priv->mr_id); - if (mr) { - uccl_engine_mr_destroy(mr); - } + uccl_engine_mr_destroy(engine_, priv->mr_id); } delete priv; } @@ -290,8 +287,10 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, } // Register memory with UCCL engine - uccl_mr_t *mr = uccl_engine_reg(engine_, mem.addr, mem.len); - if (!mr) { + uccl_mr_t mr; + // Register memory with UCCL engine + int result = uccl_engine_reg(engine_, mem.addr, mem.len, mr); + if (result != 0) { NIXL_ERROR << "Failed to register memory with UCCL engine"; return NIXL_ERR_BACKEND; } @@ -300,7 +299,7 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, priv->addr = (void *)mem.addr; priv->length = mem.len; priv->ref_cnt = 1; - priv->mr_id = reinterpret_cast(mr); // Store the memory region handle + priv->mr_id = mr; // Pre-compute fifo_item for one-sided RDMA operations int result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, @@ -328,14 +327,8 @@ nixlUcclEngine::deregisterMem(nixlBackendMD *meta) { if (priv->ref_cnt > 0) return NIXL_SUCCESS; // Deregister memory from UCCL engine - if (priv->mr_id != 0) { - uccl_mr_t *mr = reinterpret_cast(priv->mr_id); - if (mr) { - uccl_engine_mr_destroy(mr); - NIXL_DEBUG << "Deregistered memory: " << priv->addr << " mr_id: " << priv->mr_id; - } - priv->mr_id = 0; - } + uccl_engine_mr_destroy(engine_, priv->mr_id); + NIXL_DEBUG << "Deregistered memory: " << priv->addr << " mr_id: " << priv->mr_id; mem_reg_info_.erase((uint64_t)priv->addr); delete priv; @@ -481,15 +474,27 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, if (rcmode) { memcpy(uccl_handle->fifo_items[i].data(), rmd->fifo_item, FIFO_ITEM_SIZE); + // Log the original fifo_item address before overwriting + uint64_t original_fifo_addr = *reinterpret_cast(rmd->fifo_item); + uint32_t *rkey_ptr = reinterpret_cast(rmd->fifo_item + 32); // padding offset + NIXL_ERROR << "DEBUG: Original fifo_item addr=" << std::hex << original_fifo_addr + << ", rmd->addr=" << std::hex << (uintptr_t)rmd->addr + << ", remote_addr=" << std::hex << remote_addr + << ", first_rkey=" << rkey_ptr[0]; + // Update the address and size in the fifo_item for the actual transfer - // FifoItem layout: addr (8 bytes at offset 0), size (4 bytes at offset 8) - uint64_t *fifo_addr = reinterpret_cast(uccl_handle->fifo_items[i].data()); - uint32_t *fifo_size = reinterpret_cast(uccl_handle->fifo_items[i].data() + 8); - *fifo_addr = remote_addr; - *fifo_size = static_cast(rsize); - - NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex - << remote_addr << ", size=" << std::dec << rsize; + int result = uccl_engine_update_fifo( + uccl_handle->fifo_items[i].data(), + remote_addr, + static_cast(rsize) + ); + if (result != 0) { + NIXL_ERROR << "Failed to update FIFO item"; + return NIXL_ERR_BACKEND; + } + + NIXL_ERROR << "DEBUG: After update fifo_item addr=" << std::hex + << *fifo_addr << ", size=" << std::dec << *fifo_size; } else { // TODO : Suppport UC one-sided return NIXL_ERR_BACKEND; @@ -548,10 +553,11 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_INVALID_PARAM; } } - - // Process each descriptor pair - // TODO: Use a vector send async API to send all the transfers at once - std::lock_guard lock(mem_mutex_); // Lock once for the entire loop + std::vector mr_ids; + std::vector addr_v; + std::vector size_v; + + std::lock_guard lock(mem_mutex_); // Lock once for the entire operation for (size_t i = 0; i < lcnt; i++) { lmd = (nixlUcclBackendMD *)local[i].metadataP; rmd = (nixlUcclBackendMD *)remote[i].metadataP; @@ -579,61 +585,60 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, } auto local_priv = local_mem_iter->second; - if (local_priv->mr_id == 0) { - NIXL_ERROR << "Local memory region not properly registered"; - return NIXL_ERR_BACKEND; - } - uccl_mr_t *local_mr = reinterpret_cast(local_priv->mr_id); + mr_ids.push_back(local_priv->mr_id); + addr_v.push_back((void *)local_addr); + size_v.push_back(lsize); + } - int result = 0; - uint64_t transfer_id = 0; + // Perform the vector operation (single call for all transfers) + int result = 0; + uint64_t transfer_id = 0; + nixlUcclReqH *uccl_handle = static_cast(handle); - char *fifo_item_data = nullptr; - if (rcmode && handle) { - nixlUcclReqH *uccl_handle = static_cast(handle); - if (i < uccl_handle->fifo_items.size()) { - fifo_item_data = uccl_handle->fifo_items[i].data(); - } else { - NIXL_ERROR << "No FIFO item found for item: " << i - << ", fifo_items.size()=" << uccl_handle->fifo_items.size(); - return NIXL_ERR_BACKEND; - } - } + auto fifo_items = uccl_handle->fifo_items; - switch (operation) { - case NIXL_READ: { - result = uccl_engine_read( - conn, local_mr, (void *)local_addr, lsize, fifo_item_data, &transfer_id); - break; - } - case NIXL_WRITE: - if (rcmode) { - result = uccl_engine_write_rc( - conn, local_mr, (void *)local_addr, lsize, fifo_item_data, &transfer_id); - } else { - result = uccl_engine_write(conn, local_mr, (void *)local_addr, lsize, &transfer_id); - } - break; - default: - NIXL_ERROR << "Unsupported operation type: " << operation; + switch (operation) { + case NIXL_READ: { + if (!rcmode) { + NIXL_ERROR << "NIXL_READ operations require UCCL_RCMODE=1"; return NIXL_ERR_INVALID_PARAM; } - - if (result != 0) { - NIXL_ERROR << "UCCL operation failed with result: " << result; - return NIXL_ERR_BACKEND; + result = uccl_engine_read_vector( + conn, mr_ids, addr_v, size_v, fifo_items, lcnt, &transfer_id); + break; + } + case NIXL_WRITE: { + if (rcmode) { + // RC write_vector takes non-const void* + result = uccl_engine_write_vector_rc( + conn, mr_ids, addr_v, size_v, fifo_items, lcnt, &transfer_id); + } else { + // TODO: Non-RC write + // std::vector const_addr_v(addr_v.begin(), addr_v.end()); + // result = uccl_engine_write_vector( + // conn, mr_ids, const_addr_v, size_v, lcnt, &transfer_id); } + break; + } + default: + NIXL_ERROR << "Unsupported operation type: " << operation; + return NIXL_ERR_INVALID_PARAM; + } - if (!handle) { - handle = new nixlUcclReqH(conn); - } - uccl_handle = static_cast(handle); - uccl_handle->pending_transfer_ids.insert(transfer_id); + if (result != 0) { + NIXL_ERROR << "UCCL operation failed with result: " << result; + return NIXL_ERR_BACKEND; + } - NIXL_DEBUG << "Successfully posted " << (operation == NIXL_READ ? "READ" : "WRITE") - << " operation: " << lsize << " bytes with transfer_id: " << transfer_id; + if (!handle) { + handle = new nixlUcclReqH(conn); } + uccl_handle->pending_transfer_ids.insert(transfer_id); + + NIXL_DEBUG << "Successfully posted vector " << (operation == NIXL_READ ? "READ" : "WRITE") + << " operation with " << lcnt << " iovecs, transfer_id: " << transfer_id; + if (opt_args && opt_args->hasNotif) { uccl_handle->notif_msg = opt_args->notifMsg; } @@ -665,12 +670,16 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { uint64_t transfer_id = *it; int is_done = uccl_engine_xfer_status(conn, transfer_id); if (is_done) { + NIXL_ERROR << "DEBUG: Transfer " << transfer_id << " completed"; it = uccl_handle->pending_transfer_ids.erase(it); } else { ++it; } } bool all_done = uccl_handle->pending_transfer_ids.empty(); + if (all_done) { + NIXL_ERROR << "DEBUG: All transfers complete"; + } if (all_done && !uccl_handle->notif_msg.empty()) { nixlSerDes ser_des; ser_des.addStr("msg", uccl_handle->notif_msg); @@ -694,6 +703,7 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { } } + NIXL_ERROR << "DEBUG checkXfer: returning " << (all_done ? "SUCCESS" : "IN_PROG"); return (all_done) ? NIXL_SUCCESS : NIXL_IN_PROG; } @@ -715,7 +725,9 @@ nixl_status_t nixlUcclEngine::getNotifs(notif_list_t ¬if_list) { if (notif_list.size() != 0) return NIXL_ERR_INVALID_PARAM; + NIXL_ERROR << "DEBUG getNotifs: polling for notifications"; std::vector notify_msgs = uccl_engine_get_notifs(); + NIXL_ERROR << "DEBUG getNotifs: got " << notify_msgs.size() << " notifications"; for (size_t i = 0; i < notify_msgs.size(); i++) { size_t msg_len = sizeof(notify_msgs[i].msg); std::string serialized_str(notify_msgs[i].msg, msg_len); diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index 4cb3340393..57cccae16d 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -148,7 +148,7 @@ class nixlUcclBackendMD : public nixlBackendMD { void *addr; size_t length; int ref_cnt; - uint64_t mr_id; // UCCL memory region id + uccl_mr_t mr_id; // UCCL memory region id char fifo_item[FIFO_ITEM_SIZE]; }; From 8fd6a223c620f5fdea115d11e4931998cdf064e1 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 23 Jan 2026 13:07:34 +0530 Subject: [PATCH 04/28] Use prepare fifo APIs Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 57 +++++++------------------------ src/plugins/uccl/uccl_backend.h | 3 +- 2 files changed, 15 insertions(+), 45 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 9aa943d970..0e63006e85 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -288,7 +288,6 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, // Register memory with UCCL engine uccl_mr_t mr; - // Register memory with UCCL engine int result = uccl_engine_reg(engine_, mem.addr, mem.len, mr); if (result != 0) { NIXL_ERROR << "Failed to register memory with UCCL engine"; @@ -302,11 +301,11 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, priv->mr_id = mr; // Pre-compute fifo_item for one-sided RDMA operations - int result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, + result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, mem.len, priv->fifo_item); if (result != 0) { NIXL_ERROR << "Failed to prepare fifo_item for memory region"; - uccl_engine_mr_destroy(mr); + uccl_engine_mr_destroy(engine_, mr); delete priv; return NIXL_ERR_BACKEND; } @@ -442,7 +441,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, handle = new nixlUcclReqH(conn); nixlUcclReqH *uccl_handle = static_cast(handle); - uccl_handle->fifo_items.clear(); uccl_handle->fifo_items.resize(lcnt); std::lock_guard lock(mem_mutex_); @@ -466,35 +464,14 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_BACKEND; } - auto local_priv = local_mem_iter->second; - if (local_priv->mr_id == 0) { - NIXL_ERROR << "Local memory region not properly registered"; - return NIXL_ERR_BACKEND; - } if (rcmode) { - memcpy(uccl_handle->fifo_items[i].data(), rmd->fifo_item, FIFO_ITEM_SIZE); - - // Log the original fifo_item address before overwriting - uint64_t original_fifo_addr = *reinterpret_cast(rmd->fifo_item); - uint32_t *rkey_ptr = reinterpret_cast(rmd->fifo_item + 32); // padding offset - NIXL_ERROR << "DEBUG: Original fifo_item addr=" << std::hex << original_fifo_addr - << ", rmd->addr=" << std::hex << (uintptr_t)rmd->addr - << ", remote_addr=" << std::hex << remote_addr - << ", first_rkey=" << rkey_ptr[0]; - - // Update the address and size in the fifo_item for the actual transfer - int result = uccl_engine_update_fifo( - uccl_handle->fifo_items[i].data(), - remote_addr, - static_cast(rsize) - ); - if (result != 0) { - NIXL_ERROR << "Failed to update FIFO item"; - return NIXL_ERR_BACKEND; - } + // Deserialize fifo_item from char[] into FifoItem struct + deserialize_fifo_item(rmd->fifo_item, &uccl_handle->fifo_items[i]); - NIXL_ERROR << "DEBUG: After update fifo_item addr=" << std::hex - << *fifo_addr << ", size=" << std::dec << *fifo_size; + uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); + + NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex + << remote_addr << ", size=" << std::dec << rsize; } else { // TODO : Suppport UC one-sided return NIXL_ERR_BACKEND; @@ -594,9 +571,7 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, // Perform the vector operation (single call for all transfers) int result = 0; uint64_t transfer_id = 0; - nixlUcclReqH *uccl_handle = static_cast(handle); - - auto fifo_items = uccl_handle->fifo_items; + uccl_handle = static_cast(handle); switch (operation) { case NIXL_READ: { @@ -605,14 +580,13 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_INVALID_PARAM; } result = uccl_engine_read_vector( - conn, mr_ids, addr_v, size_v, fifo_items, lcnt, &transfer_id); + conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); break; } case NIXL_WRITE: { if (rcmode) { - // RC write_vector takes non-const void* - result = uccl_engine_write_vector_rc( - conn, mr_ids, addr_v, size_v, fifo_items, lcnt, &transfer_id); + result = uccl_engine_write_rc_vector( + conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); } else { // TODO: Non-RC write // std::vector const_addr_v(addr_v.begin(), addr_v.end()); @@ -677,9 +651,7 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { } } bool all_done = uccl_handle->pending_transfer_ids.empty(); - if (all_done) { - NIXL_ERROR << "DEBUG: All transfers complete"; - } + if (all_done && !uccl_handle->notif_msg.empty()) { nixlSerDes ser_des; ser_des.addStr("msg", uccl_handle->notif_msg); @@ -703,7 +675,6 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { } } - NIXL_ERROR << "DEBUG checkXfer: returning " << (all_done ? "SUCCESS" : "IN_PROG"); return (all_done) ? NIXL_SUCCESS : NIXL_IN_PROG; } @@ -725,9 +696,7 @@ nixl_status_t nixlUcclEngine::getNotifs(notif_list_t ¬if_list) { if (notif_list.size() != 0) return NIXL_ERR_INVALID_PARAM; - NIXL_ERROR << "DEBUG getNotifs: polling for notifications"; std::vector notify_msgs = uccl_engine_get_notifs(); - NIXL_ERROR << "DEBUG getNotifs: got " << notify_msgs.size() << " notifications"; for (size_t i = 0; i < notify_msgs.size(); i++) { size_t msg_len = sizeof(notify_msgs[i].msg); std::string serialized_str(notify_msgs[i].msg, msg_len); diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index 57cccae16d..94e3115c9d 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -35,6 +35,7 @@ #include "uccl_engine.h" #define FIFO_ITEM_SIZE 64 +// FifoItem and deserialize_fifo_item are now provided by uccl_engine.h class nixlUcclBackendMD; class nixlUcclReqH; @@ -162,7 +163,7 @@ class nixlUcclReqH : public nixlBackendReqH { uccl_conn_t *conn; std::unordered_set pending_transfer_ids; nixl_blob_t notif_msg; - std::vector> fifo_items; + std::vector fifo_items; }; #endif From 897f4374adf66c6cb06a6c1b25178c72ee0f5469 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 23 Jan 2026 14:06:49 +0530 Subject: [PATCH 05/28] Remove RCMODE Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 34 +-------------------------- test/gtest/plugins/uccl/uccl_test.cpp | 3 --- 2 files changed, 1 insertion(+), 36 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 0e63006e85..c2c9b11868 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -428,16 +428,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_INVALID_PARAM; } - const char *uccl_rcmode = std::getenv("UCCL_RCMODE"); - rcmode = (uccl_rcmode && std::strcmp(uccl_rcmode, "1") == 0); - - if (operation == NIXL_READ) { - if (!rcmode) { - NIXL_ERROR - << "UCCL_RCMODE environment variable must be set to 1 for NIXL_READ operations"; - return NIXL_ERR_INVALID_PARAM; - } - } handle = new nixlUcclReqH(conn); nixlUcclReqH *uccl_handle = static_cast(handle); @@ -520,16 +510,6 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_INVALID_PARAM; } - const char *uccl_rcmode = std::getenv("UCCL_RCMODE"); - rcmode = (uccl_rcmode && std::strcmp(uccl_rcmode, "1") == 0); - - if (operation == NIXL_READ) { - if (!rcmode) { - NIXL_ERROR - << "UCCL_RCMODE environment variable must be set to 1 for NIXL_READ operations"; - return NIXL_ERR_INVALID_PARAM; - } - } std::vector mr_ids; std::vector addr_v; std::vector size_v; @@ -575,24 +555,13 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, switch (operation) { case NIXL_READ: { - if (!rcmode) { - NIXL_ERROR << "NIXL_READ operations require UCCL_RCMODE=1"; - return NIXL_ERR_INVALID_PARAM; - } result = uccl_engine_read_vector( conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); break; } case NIXL_WRITE: { - if (rcmode) { - result = uccl_engine_write_rc_vector( + result = uccl_engine_write_vector( conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); - } else { - // TODO: Non-RC write - // std::vector const_addr_v(addr_v.begin(), addr_v.end()); - // result = uccl_engine_write_vector( - // conn, mr_ids, const_addr_v, size_v, lcnt, &transfer_id); - } break; } default: @@ -644,7 +613,6 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { uint64_t transfer_id = *it; int is_done = uccl_engine_xfer_status(conn, transfer_id); if (is_done) { - NIXL_ERROR << "DEBUG: Transfer " << transfer_id << " completed"; it = uccl_handle->pending_transfer_ids.erase(it); } else { ++it; diff --git a/test/gtest/plugins/uccl/uccl_test.cpp b/test/gtest/plugins/uccl/uccl_test.cpp index 50e9042063..30f96466c6 100644 --- a/test/gtest/plugins/uccl/uccl_test.cpp +++ b/test/gtest/plugins/uccl/uccl_test.cpp @@ -255,9 +255,6 @@ TestUcclBackend::TestUcclBackend() { template void TestUcclBackend::testXfer() { - if (op == NIXL_READ) { - m_env.addVar("UCCL_RCMODE", "1"); - } const std::string initiator_name = "initiator"; const std::string target_name = "target"; From 5921a683b41d26f12b773fb3c5ab6d7f94fd08ec Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 23 Jan 2026 14:20:58 +0530 Subject: [PATCH 06/28] Remove rcmode in prepXfer Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 16 +++++----------- 1 file changed, 5 insertions(+), 11 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index c2c9b11868..90d8370025 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -399,7 +399,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, const nixl_opt_b_args_t *opt_args) const { nixlUcclBackendMD *lmd; nixlUcclBackendMD *rmd; - bool rcmode = false; handle = nullptr; NIXL_DEBUG << "UCCL PrepXfer: " << operation << " remote_agent: " << remote_agent; @@ -454,18 +453,13 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, return NIXL_ERR_BACKEND; } - if (rcmode) { - // Deserialize fifo_item from char[] into FifoItem struct - deserialize_fifo_item(rmd->fifo_item, &uccl_handle->fifo_items[i]); + // Deserialize fifo_item from char[] into FifoItem struct + deserialize_fifo_item(rmd->fifo_item, &uccl_handle->fifo_items[i]); - uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); + uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); - NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex - << remote_addr << ", size=" << std::dec << rsize; - } else { - // TODO : Suppport UC one-sided - return NIXL_ERR_BACKEND; - } + NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex + << remote_addr << ", size=" << std::dec << rsize; } return NIXL_SUCCESS; From f03ee88f274fdaa43e017a50d0540a611b4511d0 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 23 Jan 2026 14:26:01 +0530 Subject: [PATCH 07/28] Remove rcmode from CI Signed-off-by: Pravein Govindan Kannan --- .gitlab/test_nixlbench.sh | 2 +- src/plugins/uccl/uccl_backend.cpp | 1 - 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/.gitlab/test_nixlbench.sh b/.gitlab/test_nixlbench.sh index 1d5aa2f376..b7a769825c 100755 --- a/.gitlab/test_nixlbench.sh +++ b/.gitlab/test_nixlbench.sh @@ -111,7 +111,7 @@ if $HAS_GPU ; then for op_type in READ WRITE; do for initiator in $seg_types; do for target in $seg_types; do - UCCL_RCMODE=1 run_nixlbench_two_workers --backend UCCL --op_type $op_type --initiator_seg_type $initiator --target_seg_type $target --check_consistency + run_nixlbench_two_workers --backend UCCL --op_type $op_type --initiator_seg_type $initiator --target_seg_type $target --check_consistency done done done diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 90d8370025..d6522a6043 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -475,7 +475,6 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, nixlUcclReqH *uccl_handle; nixlUcclBackendMD *lmd; nixlUcclBackendMD *rmd; - bool rcmode = false; NIXL_DEBUG << "UCCL PostXfer: " << operation << " remote_agent: " << remote_agent; From 76f8c81b0a8c6cc83151f62df14ac22727d35f90 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 23 Jan 2026 15:19:56 +0530 Subject: [PATCH 08/28] Fix cleanup Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 30 ++++++++++++++++++++++-------- src/plugins/uccl/uccl_backend.h | 1 + 2 files changed, 23 insertions(+), 8 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index d6522a6043..5963cbf737 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -89,7 +89,7 @@ getNixlParam(const nixl_b_params_t *custom_params, const std::string &key, int d } nixlUcclEngine::nixlUcclEngine(const nixlBackendInitParams *init_params) - : nixlBackendEngine(init_params) { + : nixlBackendEngine(init_params), stop_listener_(false) { local_agent_name_ = init_params->localAgent; nixl_b_params_t *custom_params = init_params->customParams; @@ -103,6 +103,8 @@ nixlUcclEngine::nixlUcclEngine(const nixlBackendInitParams *init_params) } nixlUcclEngine::~nixlUcclEngine() { + stop_listener_ = true; + { std::lock_guard lock(mem_mutex_); for (auto &[addr, priv] : mem_reg_info_) { @@ -128,27 +130,38 @@ nixlUcclEngine::~nixlUcclEngine() { connected_agents_.clear(); } - if (listener_thread_.joinable()) { - listener_thread_.detach(); - } - if (engine_) { - // Add a small delay to allow UCCL internal cleanup to complete - std::this_thread::sleep_for(std::chrono::milliseconds(1)); uccl_engine_destroy(engine_); engine_ = nullptr; } + + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + + if (listener_thread_.joinable()) { + listener_thread_.detach(); + } } void nixlUcclEngine::startListener() { // The listener waits for connections from remote agents NIXL_DEBUG << "UCCL accepting connections"; - while (true) { + while (!stop_listener_) { + // Check if engine is still valid before using it + if (!engine_) { + NIXL_DEBUG << "Engine destroyed, listener thread exiting"; + break; + } + char ip_buf[256]; int remote_gpu_idx; uccl_conn_t *conn = uccl_engine_accept(engine_, ip_buf, sizeof(ip_buf), &remote_gpu_idx); if (!conn) { + // Check if we should stop (engine destroyed or shutdown requested) + if (stop_listener_ || !engine_) { + NIXL_DEBUG << "Listener thread stopping"; + break; + } NIXL_ERROR << "Failed to accept connection from remote agent"; continue; } @@ -160,6 +173,7 @@ nixlUcclEngine::startListener() { connected_agents_[ip_buf] = reinterpret_cast(conn); } } + NIXL_DEBUG << "UCCL listener thread exiting"; } nixl_mem_list_t diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index 94e3115c9d..c6643137fe 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -135,6 +135,7 @@ class nixlUcclEngine : public nixlBackendEngine { std::unordered_map mem_reg_info_; std::unordered_map connected_agents_; // agent name -> conn_id std::thread listener_thread_; + std::atomic stop_listener_; }; // UCCL Backend Memory Descriptor From e80c85b56139451aef1613b4a3ff7698cb549993 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Wed, 28 Jan 2026 17:33:25 +0530 Subject: [PATCH 09/28] Maintain a single transfer ID Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 20 +++++--------------- src/plugins/uccl/uccl_backend.h | 2 +- 2 files changed, 6 insertions(+), 16 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 5963cbf737..cb1e7f4647 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -584,7 +584,7 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, if (!handle) { handle = new nixlUcclReqH(conn); } - uccl_handle->pending_transfer_ids.insert(transfer_id); + uccl_handle->transfer_id = transfer_id; NIXL_DEBUG << "Successfully posted vector " << (operation == NIXL_READ ? "READ" : "WRITE") << " operation with " << lcnt << " iovecs, transfer_id: " << transfer_id; @@ -615,19 +615,8 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { return NIXL_ERR_BACKEND; } - auto it = uccl_handle->pending_transfer_ids.begin(); - while (it != uccl_handle->pending_transfer_ids.end()) { - uint64_t transfer_id = *it; - int is_done = uccl_engine_xfer_status(conn, transfer_id); - if (is_done) { - it = uccl_handle->pending_transfer_ids.erase(it); - } else { - ++it; - } - } - bool all_done = uccl_handle->pending_transfer_ids.empty(); - - if (all_done && !uccl_handle->notif_msg.empty()) { + int is_done = uccl_engine_xfer_status(conn, uccl_handle->transfer_id); + if (is_done) { nixlSerDes ser_des; ser_des.addStr("msg", uccl_handle->notif_msg); std::string serialized = ser_des.exportStr(); @@ -648,9 +637,10 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { NIXL_DEBUG << "All transfers in handle completed, sent notification: " << uccl_handle->notif_msg; } + return NIXL_SUCCESS; } - return (all_done) ? NIXL_SUCCESS : NIXL_IN_PROG; + return NIXL_IN_PROG; } nixl_status_t diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index c6643137fe..ecfa609d59 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -162,7 +162,7 @@ class nixlUcclReqH : public nixlBackendReqH { virtual ~nixlUcclReqH() {} uccl_conn_t *conn; - std::unordered_set pending_transfer_ids; + uint64_t transfer_id; nixl_blob_t notif_msg; std::vector fifo_items; }; From 047fd3cb55fdf7b765c451703a60b9c1f341eede Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 29 Jan 2026 10:09:47 +0530 Subject: [PATCH 10/28] Fix formatting Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 37 +++++++++++++++---------------- 1 file changed, 18 insertions(+), 19 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index cb1e7f4647..215aa25fb3 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -89,7 +89,8 @@ getNixlParam(const nixl_b_params_t *custom_params, const std::string &key, int d } nixlUcclEngine::nixlUcclEngine(const nixlBackendInitParams *init_params) - : nixlBackendEngine(init_params), stop_listener_(false) { + : nixlBackendEngine(init_params), + stop_listener_(false) { local_agent_name_ = init_params->localAgent; nixl_b_params_t *custom_params = init_params->customParams; @@ -197,7 +198,6 @@ nixlUcclEngine::getPublicData(const nixlBackendMD *meta, std::string &str) const snprintf(hex, sizeof(hex), "%02x", static_cast(priv->fifo_item[i])); str += hex; } - NIXL_DEBUG << "Exporting Meta Info =" << str << std::endl; return NIXL_SUCCESS; } @@ -315,8 +315,7 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, priv->mr_id = mr; // Pre-compute fifo_item for one-sided RDMA operations - result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, - mem.len, priv->fifo_item); + result = uccl_engine_prepare_fifo(engine_, mr, (void *)mem.addr, mem.len, priv->fifo_item); if (result != 0) { NIXL_ERROR << "Failed to prepare fifo_item for memory region"; uccl_engine_mr_destroy(engine_, mr); @@ -326,7 +325,7 @@ nixlUcclEngine::registerMem(const nixlBlobDesc &mem, out = priv; mem_reg_info_[mem.addr] = priv; - NIXL_DEBUG << "Registering memory: " << mem.addr << "Device: " << mem.devId + NIXL_DEBUG << "Registering memory: " << std::hex << mem.addr << " Device: " << mem.devId << " ref_cnt: " << priv->ref_cnt << " mr_id: " << priv->mr_id; return NIXL_SUCCESS; @@ -341,7 +340,7 @@ nixlUcclEngine::deregisterMem(nixlBackendMD *meta) { // Deregister memory from UCCL engine uccl_engine_mr_destroy(engine_, priv->mr_id); - NIXL_DEBUG << "Deregistered memory: " << priv->addr << " mr_id: " << priv->mr_id; + NIXL_DEBUG << "Deregistered memory: " << std::hex << priv->addr << " mr_id: " << priv->mr_id; mem_reg_info_.erase((uint64_t)priv->addr); delete priv; @@ -351,7 +350,7 @@ nixlUcclEngine::deregisterMem(nixlBackendMD *meta) { nixl_status_t nixlUcclEngine::loadLocalMD(nixlBackendMD *input, nixlBackendMD *&output) { nixlUcclBackendMD *input_md = (nixlUcclBackendMD *)input; - NIXL_DEBUG << "UCCL Load Local MD: " << input_md->addr << "Meta Info:" << input_md->mr_id; + NIXL_DEBUG << "UCCL Load Local MD: " << std::hex << input_md->addr << "Meta Info:" << input_md->mr_id; nixlUcclBackendMD *output_md = (nixlUcclBackendMD *)output; output_md->addr = (void *)input_md->addr; @@ -367,7 +366,7 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, const nixl_mem_t &nixl_mem, const std::string &remote_agent, nixlBackendMD *&output) { - NIXL_DEBUG << "UCCL Load Remote MD: " << input.addr << "Meta Info:" << input.metaInfo + NIXL_DEBUG << "UCCL Load Remote MD: " << std::hex << input.addr << " Meta Info:" << input.metaInfo << " remote_agent: " << remote_agent; output = new nixlUcclBackendMD(true); @@ -378,7 +377,7 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, // Decode fifo_item from hex string const std::string &hex_str = input.metaInfo; - NIXL_DEBUG << "Meta Info =" << hex_str << std::endl; + if (hex_str.length() == FIFO_ITEM_SIZE * 2) { for (int i = 0; i < FIFO_ITEM_SIZE; i++) { std::string byte_str = hex_str.substr(i * 2, 2); @@ -386,8 +385,8 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, } NIXL_DEBUG << "Parsed fifo_item from remote metadata"; } else { - NIXL_ERROR << "Invalid fifo_item hex string length: " << hex_str.length() - << " (expected " << FIFO_ITEM_SIZE * 2 << ")"; + NIXL_ERROR << "Invalid fifo_item hex string length: " << hex_str.length() << " (expected " + << FIFO_ITEM_SIZE * 2 << ")"; delete output_md; output = nullptr; return NIXL_ERR_INVALID_PARAM; @@ -472,8 +471,8 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); - NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex - << remote_addr << ", size=" << std::dec << rsize; + NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex << remote_addr + << ", size=" << std::dec << rsize; } return NIXL_SUCCESS; @@ -518,9 +517,9 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, } std::vector mr_ids; - std::vector addr_v; + std::vector addr_v; std::vector size_v; - + std::lock_guard lock(mem_mutex_); // Lock once for the entire operation for (size_t i = 0; i < lcnt; i++) { lmd = (nixlUcclBackendMD *)local[i].metadataP; @@ -555,7 +554,7 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, size_v.push_back(lsize); } - // Perform the vector operation (single call for all transfers) + // Perform a vector operation int result = 0; uint64_t transfer_id = 0; uccl_handle = static_cast(handle); @@ -567,8 +566,8 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, break; } case NIXL_WRITE: { - result = uccl_engine_write_vector( - conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); + result = uccl_engine_write_vector( + conn, mr_ids, addr_v, size_v, uccl_handle->fifo_items, lcnt, &transfer_id); break; } default: @@ -588,7 +587,7 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, NIXL_DEBUG << "Successfully posted vector " << (operation == NIXL_READ ? "READ" : "WRITE") << " operation with " << lcnt << " iovecs, transfer_id: " << transfer_id; - + if (opt_args && opt_args->hasNotif) { uccl_handle->notif_msg = opt_args->notifMsg; } From 04384003d7740b475fbf4988ab4a553f62e1c1df Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 29 Jan 2026 10:10:08 +0530 Subject: [PATCH 11/28] Fix formatting after changes Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 215aa25fb3..da60a1dfde 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -350,7 +350,8 @@ nixlUcclEngine::deregisterMem(nixlBackendMD *meta) { nixl_status_t nixlUcclEngine::loadLocalMD(nixlBackendMD *input, nixlBackendMD *&output) { nixlUcclBackendMD *input_md = (nixlUcclBackendMD *)input; - NIXL_DEBUG << "UCCL Load Local MD: " << std::hex << input_md->addr << "Meta Info:" << input_md->mr_id; + NIXL_DEBUG << "UCCL Load Local MD: " << std::hex << input_md->addr + << "Meta Info:" << input_md->mr_id; nixlUcclBackendMD *output_md = (nixlUcclBackendMD *)output; output_md->addr = (void *)input_md->addr; @@ -366,8 +367,8 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, const nixl_mem_t &nixl_mem, const std::string &remote_agent, nixlBackendMD *&output) { - NIXL_DEBUG << "UCCL Load Remote MD: " << std::hex << input.addr << " Meta Info:" << input.metaInfo - << " remote_agent: " << remote_agent; + NIXL_DEBUG << "UCCL Load Remote MD: " << std::hex << input.addr + << " Meta Info:" << input.metaInfo << " remote_agent: " << remote_agent; output = new nixlUcclBackendMD(true); nixlUcclBackendMD *output_md = static_cast(output); From 83a8eeaa2c65595ca242abcc3e0339e710a4bd5d Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 14:07:02 +0530 Subject: [PATCH 12/28] Remove redundant FIFO_ITEM_SIZE Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 26 ++++++++------------------ src/plugins/uccl/uccl_backend.h | 7 ++----- 2 files changed, 10 insertions(+), 23 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index da60a1dfde..26bc2533a7 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -148,7 +148,6 @@ nixlUcclEngine::startListener() { // The listener waits for connections from remote agents NIXL_DEBUG << "UCCL accepting connections"; while (!stop_listener_) { - // Check if engine is still valid before using it if (!engine_) { NIXL_DEBUG << "Engine destroyed, listener thread exiting"; break; @@ -174,7 +173,6 @@ nixlUcclEngine::startListener() { connected_agents_[ip_buf] = reinterpret_cast(conn); } } - NIXL_DEBUG << "UCCL listener thread exiting"; } nixl_mem_list_t @@ -190,10 +188,11 @@ nixl_status_t nixlUcclEngine::getPublicData(const nixlBackendMD *meta, std::string &str) const { nixlUcclBackendMD *priv = (nixlUcclBackendMD *)meta; - // Export fifo_item as hex string + // Export fifo_item as hex string. + // The fifo_item is used to perform one-sided operation str.clear(); - str.reserve(FIFO_ITEM_SIZE * 2); - for (int i = 0; i < FIFO_ITEM_SIZE; i++) { + str.reserve(FIFO_SIZE * 2); + for (int i = 0; i < FIFO_SIZE; i++) { char hex[3]; snprintf(hex, sizeof(hex), "%02x", static_cast(priv->fifo_item[i])); str += hex; @@ -379,15 +378,14 @@ nixlUcclEngine::loadRemoteMD(const nixlBlobDesc &input, // Decode fifo_item from hex string const std::string &hex_str = input.metaInfo; - if (hex_str.length() == FIFO_ITEM_SIZE * 2) { - for (int i = 0; i < FIFO_ITEM_SIZE; i++) { + if (hex_str.length() == FIFO_SIZE * 2) { + for (int i = 0; i < FIFO_SIZE; i++) { std::string byte_str = hex_str.substr(i * 2, 2); output_md->fifo_item[i] = static_cast(strtoul(byte_str.c_str(), NULL, 16)); } - NIXL_DEBUG << "Parsed fifo_item from remote metadata"; } else { NIXL_ERROR << "Invalid fifo_item hex string length: " << hex_str.length() << " (expected " - << FIFO_ITEM_SIZE * 2 << ")"; + << FIFO_SIZE * 2 << ")"; delete output_md; output = nullptr; return NIXL_ERR_INVALID_PARAM; @@ -454,12 +452,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, uintptr_t local_addr = local[i].addr; uintptr_t remote_addr = remote[i].addr; - NIXL_DEBUG << "prepXfer iovec[" << i << "]: local[i].addr=" << std::hex << local_addr - << ", lmd->addr=" << std::hex << lmd->addr << ", lmd->mr_id=" << std::dec - << lmd->mr_id << ", remote[i].addr=" << std::hex << remote_addr - << ", rmd->addr=" << std::hex << rmd->addr << ", rmd->mr_id=" << std::dec - << rmd->mr_id; - // Validate the local address is registered auto local_mem_iter = mem_reg_info_.find((uint64_t)lmd->addr); if (local_mem_iter == mem_reg_info_.end()) { @@ -472,8 +464,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); - NIXL_DEBUG << "Using pre-shared fifo_item[" << i << "]: addr=" << std::hex << remote_addr - << ", size=" << std::dec << rsize; } return NIXL_SUCCESS; @@ -634,7 +624,7 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { NIXL_ERROR << "Failed to send notify message"; return NIXL_ERR_BACKEND; } - NIXL_DEBUG << "All transfers in handle completed, sent notification: " + NIXL_DEBUG << "Transfer complete, sent notification: " << uccl_handle->notif_msg; } return NIXL_SUCCESS; diff --git a/src/plugins/uccl/uccl_backend.h b/src/plugins/uccl/uccl_backend.h index ecfa609d59..97d1fd1416 100644 --- a/src/plugins/uccl/uccl_backend.h +++ b/src/plugins/uccl/uccl_backend.h @@ -34,9 +34,6 @@ #include "uccl_engine.h" -#define FIFO_ITEM_SIZE 64 -// FifoItem and deserialize_fifo_item are now provided by uccl_engine.h - class nixlUcclBackendMD; class nixlUcclReqH; @@ -142,7 +139,7 @@ class nixlUcclEngine : public nixlBackendEngine { class nixlUcclBackendMD : public nixlBackendMD { public: nixlUcclBackendMD(bool isPrivate) : nixlBackendMD(isPrivate) { - memset(fifo_item, 0, FIFO_ITEM_SIZE); + memset(fifo_item, 0, FIFO_SIZE); } virtual ~nixlUcclBackendMD() {} @@ -151,7 +148,7 @@ class nixlUcclBackendMD : public nixlBackendMD { size_t length; int ref_cnt; uccl_mr_t mr_id; // UCCL memory region id - char fifo_item[FIFO_ITEM_SIZE]; + char fifo_item[FIFO_SIZE]; }; // UCCL Backend Request Handle From d70747d316c9ead246180d86c046427d65618b98 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 14:11:22 +0530 Subject: [PATCH 13/28] Update commit SHA of UCCL Signed-off-by: Pravein Govindan Kannan --- .gitlab/build.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.gitlab/build.sh b/.gitlab/build.sh index 819209dcab..ae280bfe1c 100755 --- a/.gitlab/build.sh +++ b/.gitlab/build.sh @@ -33,7 +33,7 @@ LIBFABRIC_VERSION=${LIBFABRIC_VERSION:-v1.21.0} # LIBFABRIC_INSTALL_DIR can be set via environment variable, defaults to INSTALL_DIR LIBFABRIC_INSTALL_DIR=${LIBFABRIC_INSTALL_DIR:-$INSTALL_DIR} # UCCL_COMMIT_SHA is the commit SHA of UCCL. -UCCL_COMMIT_SHA="a962f611021afc2e3c9358f6da4ae96539cbca0f" +UCCL_COMMIT_SHA="3af0d38dd6ba5b1070a718ab631826999fef695f" if [ -z "$INSTALL_DIR" ]; then echo "Usage: $0 " From 92f0466c36f560725efdeea40e7670a904c44307 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 14:16:43 +0530 Subject: [PATCH 14/28] Remove unused variable Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 26bc2533a7..cf4c35034b 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -449,7 +449,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, lmd = (nixlUcclBackendMD *)local[i].metadataP; rmd = (nixlUcclBackendMD *)remote[i].metadataP; size_t rsize = remote[i].len; - uintptr_t local_addr = local[i].addr; uintptr_t remote_addr = remote[i].addr; // Validate the local address is registered @@ -521,11 +520,6 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, uintptr_t local_addr = local[i].addr; uintptr_t remote_addr = remote[i].addr; - NIXL_DEBUG << "postXfer iovec[" << i << "]: local[i].addr=" << std::hex << local_addr - << ", lsize=" << std::dec << lsize << ", remote[i].addr=" << std::hex - << remote_addr << ", rsize=" << std::dec << rsize << ", lmd->addr=" << std::hex - << lmd->addr << ", rmd->addr=" << std::hex << rmd->addr; - if (lsize != rsize) { NIXL_ERROR << "Local and remote sizes don't match: " << lsize << " != " << rsize; return NIXL_ERR_INVALID_PARAM; @@ -545,7 +539,7 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, size_v.push_back(lsize); } - // Perform a vector operation + // Perform a vector read/write operation int result = 0; uint64_t transfer_id = 0; uccl_handle = static_cast(handle); From d276d725d32c272c9325a45a29433136e178c296 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 14:19:58 +0530 Subject: [PATCH 15/28] Remove unused vars Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index cf4c35034b..a05f9b32cc 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -510,15 +510,13 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, std::vector addr_v; std::vector size_v; - std::lock_guard lock(mem_mutex_); // Lock once for the entire operation + std::lock_guard lock(mem_mutex_); for (size_t i = 0; i < lcnt; i++) { lmd = (nixlUcclBackendMD *)local[i].metadataP; - rmd = (nixlUcclBackendMD *)remote[i].metadataP; size_t lsize = local[i].len; size_t rsize = remote[i].len; // Use local[i].addr for the actual iovec address, not lmd->addr (which is base address) uintptr_t local_addr = local[i].addr; - uintptr_t remote_addr = remote[i].addr; if (lsize != rsize) { NIXL_ERROR << "Local and remote sizes don't match: " << lsize << " != " << rsize; From 800179bcf9961929a1ea1dbdad042f1b5e3299f5 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 14:20:59 +0530 Subject: [PATCH 16/28] Fix error Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 1 - 1 file changed, 1 deletion(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index a05f9b32cc..c13c3933e1 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -477,7 +477,6 @@ nixlUcclEngine::postXfer(const nixl_xfer_op_t &operation, const nixl_opt_b_args_t *opt_args) const { nixlUcclReqH *uccl_handle; nixlUcclBackendMD *lmd; - nixlUcclBackendMD *rmd; NIXL_DEBUG << "UCCL PostXfer: " << operation << " remote_agent: " << remote_agent; From be118d192f77e45db174cb111431d0f23d28823b Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 15:27:46 +0530 Subject: [PATCH 17/28] Add recent updates to README Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/README.md | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/src/plugins/uccl/README.md b/src/plugins/uccl/README.md index 7078a6b1f2..27db3584b8 100644 --- a/src/plugins/uccl/README.md +++ b/src/plugins/uccl/README.md @@ -34,20 +34,18 @@ UCCL engine would auto discover the right NIC to be used for the GPU based on th Refer to [README](https://github.com/uccl-project/uccl/tree/main/collective/rdma#environment-variables-in-uccl) for the complete list of environment variables that can be set to customize UCCL. -**Important**: For `NIXL_READ` operations set `UCCL_RCMODE=1`. By default, UCCL uses RDMA UC (Unreliable Connection). However, `READ` operations need to operate on RDMA RC (Reliable Connection). `WRITE` operations can work on both RDMA UC and RC. - ### Usage References 1) [NIXL Benchmark](https://github.com/uccl-project/uccl/blob/main/p2p/benchmarks/benchmark_nixl.py) in UCCL P2P: Refer to this [README](https://github.com/uccl-project/uccl/tree/main/p2p) on how to run the script. -2) [NIXL connector](https://github.com/vllm-project/vllm/commit/e731733d30d0aed3252dc60427927768bfc0ca73) in vLLM. vLLM's NIXL connector uses `NIXL_READ` operations, hence set env `UCCL_RCMODE` to 1. +2) [NIXL connector](https://github.com/vllm-project/vllm/commit/e731733d30d0aed3252dc60427927768bfc0ca73) in vLLM. ### Road Map -- [ ] Add Intra-node communication support +- ✅ Add asynchronous posting of reads over multiple workers to mitigate latency increase upon fragmentation -- [ ] Add Progress Thread support +- 🚧 Add Intra-node communication support -- [ ] Add asynchronous posting of reads over multiple workers to mitigate latency increase upon fragmentation +- 🚧 Add Progress Thread support -- [ ] Add support for other transport (TCP, TCP-X, etc.) \ No newline at end of file +- 🚧 Add support for other transport (TCP, TCP-X, etc.) \ No newline at end of file From dc3d822ede5172e19a401b4d1bef1ec690454178 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 15:51:44 +0530 Subject: [PATCH 18/28] Fix copyright and format Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 4 +--- test/gtest/plugins/uccl/uccl_test.cpp | 2 +- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index c13c3933e1..2a8650f06c 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -462,7 +462,6 @@ nixlUcclEngine::prepXfer(const nixl_xfer_op_t &operation, deserialize_fifo_item(rmd->fifo_item, &uccl_handle->fifo_items[i]); uccl_engine_update_fifo(uccl_handle->fifo_items[i], remote_addr, rsize); - } return NIXL_SUCCESS; @@ -615,8 +614,7 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { NIXL_ERROR << "Failed to send notify message"; return NIXL_ERR_BACKEND; } - NIXL_DEBUG << "Transfer complete, sent notification: " - << uccl_handle->notif_msg; + NIXL_DEBUG << "Transfer complete, sent notification: " << uccl_handle->notif_msg; } return NIXL_SUCCESS; } diff --git a/test/gtest/plugins/uccl/uccl_test.cpp b/test/gtest/plugins/uccl/uccl_test.cpp index 30f96466c6..5ac46e0f79 100644 --- a/test/gtest/plugins/uccl/uccl_test.cpp +++ b/test/gtest/plugins/uccl/uccl_test.cpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2025 NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. * SPDX-License-Identifier: Apache-2.0 * * Licensed under the Apache License, Version 2.0 (the "License"); From 90d27cb57804e058a3781d72d545102c0565de4b Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 3 Feb 2026 19:20:16 +0530 Subject: [PATCH 19/28] Increase CI image tag Signed-off-by: Pravein Govindan Kannan --- .ci/jenkins/lib/build-matrix.yaml | 2 +- .ci/jenkins/lib/test-matrix.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.ci/jenkins/lib/build-matrix.yaml b/.ci/jenkins/lib/build-matrix.yaml index a1e7636621..c24dac1b96 100644 --- a/.ci/jenkins/lib/build-matrix.yaml +++ b/.ci/jenkins/lib/build-matrix.yaml @@ -42,7 +42,7 @@ env: TEST_TIMEOUT: 30 UCX_TLS: "^shm" STORAGE_DRIVER: 'overlay' - CI_IMAGE_TAG: "20260129-1" + CI_IMAGE_TAG: "20260203-1" runs_on_dockers: diff --git a/.ci/jenkins/lib/test-matrix.yaml b/.ci/jenkins/lib/test-matrix.yaml index 79c516e58e..feeba48701 100644 --- a/.ci/jenkins/lib/test-matrix.yaml +++ b/.ci/jenkins/lib/test-matrix.yaml @@ -30,7 +30,7 @@ env: NIXL_INSTALL_DIR: /opt/nixl TEST_TIMEOUT: 30 NPROC: 32 - CI_IMAGE_TAG: "20260129-1" + CI_IMAGE_TAG: "20260203-1" docker_opt: "-e NPROC -e EXECUTOR_NUMBER --security-opt seccomp=unconfined --security-opt apparmor=unconfined --shm-size=1G --ulimit nofile=65535:65535 --ulimit stack=67108864 --ulimit memlock=-1:-1 --network=host --ipc=host --cap-add=SYS_PTRACE --cap-add=SYS_ADMIN --privileged --gpus all --device=/dev/infiniband --device=/dev/gdrdrv --entrypoint=''" From facdff67f6b6e56be8a9a661ce563eb81060a25d Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Wed, 4 Feb 2026 18:42:22 +0530 Subject: [PATCH 20/28] Use latest UCCL commit SHA Signed-off-by: Pravein Govindan Kannan --- .gitlab/build.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.gitlab/build.sh b/.gitlab/build.sh index 78342c7878..4e0a36e07e 100755 --- a/.gitlab/build.sh +++ b/.gitlab/build.sh @@ -35,7 +35,7 @@ LIBFABRIC_VERSION=${LIBFABRIC_VERSION:-v1.21.0} # LIBFABRIC_INSTALL_DIR can be set via environment variable, defaults to INSTALL_DIR LIBFABRIC_INSTALL_DIR=${LIBFABRIC_INSTALL_DIR:-$INSTALL_DIR} # UCCL_COMMIT_SHA is the commit SHA of UCCL. -UCCL_COMMIT_SHA="3af0d38dd6ba5b1070a718ab631826999fef695f" +UCCL_COMMIT_SHA="91f92325dced4e1313f6f35ac94c63d745f68df9" AZURITE_VER="3.35.0" TMPDIR=$(mktemp -d) From 99df3d5a55081083f701655f6863b25618878ff8 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 19 Feb 2026 12:31:35 +0530 Subject: [PATCH 21/28] Add recent UCCL Commit - Fixing build issue on ARM platform - Add TCP support Signed-off-by: Pravein Govindan Kannan --- .gitlab/build.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.gitlab/build.sh b/.gitlab/build.sh index 4e0a36e07e..80af1d4aae 100755 --- a/.gitlab/build.sh +++ b/.gitlab/build.sh @@ -35,7 +35,7 @@ LIBFABRIC_VERSION=${LIBFABRIC_VERSION:-v1.21.0} # LIBFABRIC_INSTALL_DIR can be set via environment variable, defaults to INSTALL_DIR LIBFABRIC_INSTALL_DIR=${LIBFABRIC_INSTALL_DIR:-$INSTALL_DIR} # UCCL_COMMIT_SHA is the commit SHA of UCCL. -UCCL_COMMIT_SHA="91f92325dced4e1313f6f35ac94c63d745f68df9" +UCCL_COMMIT_SHA="a550d613f9a61e4b3b08ddf577a232891fb75ec7" AZURITE_VER="3.35.0" TMPDIR=$(mktemp -d) From d1e653662d591d23ba37121bafebfffe618dbec8 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Tue, 24 Feb 2026 19:32:18 +0530 Subject: [PATCH 22/28] Remove sleep during cleanup Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 2 -- 1 file changed, 2 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 2a8650f06c..42faf662b0 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -136,8 +136,6 @@ nixlUcclEngine::~nixlUcclEngine() { engine_ = nullptr; } - std::this_thread::sleep_for(std::chrono::milliseconds(100)); - if (listener_thread_.joinable()) { listener_thread_.detach(); } From 3144cf975e3fa9efa497665642ababb4f7d18e0b Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Wed, 25 Feb 2026 19:51:46 +0530 Subject: [PATCH 23/28] Add stop_accept API for graceful cleanup Signed-off-by: Pravein Govindan Kannan --- .gitlab/build.sh | 2 +- src/plugins/uccl/uccl_backend.cpp | 12 ++++++++---- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/.gitlab/build.sh b/.gitlab/build.sh index 631bef6681..708c5f3070 100755 --- a/.gitlab/build.sh +++ b/.gitlab/build.sh @@ -35,7 +35,7 @@ LIBFABRIC_VERSION=${LIBFABRIC_VERSION:-v1.21.0} # LIBFABRIC_INSTALL_DIR can be set via environment variable, defaults to INSTALL_DIR LIBFABRIC_INSTALL_DIR=${LIBFABRIC_INSTALL_DIR:-$INSTALL_DIR} # UCCL_COMMIT_SHA is the commit SHA of UCCL. -UCCL_COMMIT_SHA="a550d613f9a61e4b3b08ddf577a232891fb75ec7" +UCCL_COMMIT_SHA="2de728f1a27ea3f3b66059baf838f940e243ebc6" AZURITE_VER="3.35.0" TMPDIR=$(mktemp -d) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 42faf662b0..466ab56d1b 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -106,6 +106,14 @@ nixlUcclEngine::nixlUcclEngine(const nixlBackendInitParams *init_params) nixlUcclEngine::~nixlUcclEngine() { stop_listener_ = true; + if (engine_) { + uccl_engine_stop_accept(engine_); + } + + if (listener_thread_.joinable()) { + listener_thread_.join(); + } + { std::lock_guard lock(mem_mutex_); for (auto &[addr, priv] : mem_reg_info_) { @@ -135,10 +143,6 @@ nixlUcclEngine::~nixlUcclEngine() { uccl_engine_destroy(engine_); engine_ = nullptr; } - - if (listener_thread_.joinable()) { - listener_thread_.detach(); - } } void From 1bc2d51ce77a7321f3b1a6872ed8588acfee4e5d Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 26 Feb 2026 10:37:05 +0530 Subject: [PATCH 24/28] Use CodeRabbit suggestions Signed-off-by: Pravein Govindan Kannan --- src/plugins/uccl/uccl_backend.cpp | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/src/plugins/uccl/uccl_backend.cpp b/src/plugins/uccl/uccl_backend.cpp index 466ab56d1b..7cc42ede77 100644 --- a/src/plugins/uccl/uccl_backend.cpp +++ b/src/plugins/uccl/uccl_backend.cpp @@ -150,17 +150,13 @@ nixlUcclEngine::startListener() { // The listener waits for connections from remote agents NIXL_DEBUG << "UCCL accepting connections"; while (!stop_listener_) { - if (!engine_) { - NIXL_DEBUG << "Engine destroyed, listener thread exiting"; - break; - } char ip_buf[256]; int remote_gpu_idx; uccl_conn_t *conn = uccl_engine_accept(engine_, ip_buf, sizeof(ip_buf), &remote_gpu_idx); if (!conn) { - // Check if we should stop (engine destroyed or shutdown requested) - if (stop_listener_ || !engine_) { + // Check if we should stop + if (stop_listener_) { NIXL_DEBUG << "Listener thread stopping"; break; } @@ -597,7 +593,7 @@ nixlUcclEngine::checkXfer(nixlBackendReqH *handle) const { return NIXL_ERR_BACKEND; } - int is_done = uccl_engine_xfer_status(conn, uccl_handle->transfer_id); + bool is_done = uccl_engine_xfer_status(conn, uccl_handle->transfer_id); if (is_done) { nixlSerDes ser_des; ser_des.addStr("msg", uccl_handle->notif_msg); From ca40dee35d7c624f80f080a163d03511af1abfb2 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 26 Feb 2026 19:10:00 +0530 Subject: [PATCH 25/28] Add cleanup change to address crash Signed-off-by: Pravein Govindan Kannan --- benchmark/nixlbench/src/worker/nixl/nixl_worker.cpp | 3 +++ 1 file changed, 3 insertions(+) diff --git a/benchmark/nixlbench/src/worker/nixl/nixl_worker.cpp b/benchmark/nixlbench/src/worker/nixl/nixl_worker.cpp index fbd95a544a..8cdaed4123 100644 --- a/benchmark/nixlbench/src/worker/nixl/nixl_worker.cpp +++ b/benchmark/nixlbench/src/worker/nixl/nixl_worker.cpp @@ -310,6 +310,9 @@ xferBenchNixlWorker::xferBenchNixlWorker(int *argc, char ***argv, std::vector Date: Thu, 26 Feb 2026 20:05:44 +0530 Subject: [PATCH 26/28] Bump up the IMAGE_TAG Signed-off-by: Pravein Govindan Kannan --- .ci/jenkins/lib/build-matrix.yaml | 2 +- .ci/jenkins/lib/test-matrix.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.ci/jenkins/lib/build-matrix.yaml b/.ci/jenkins/lib/build-matrix.yaml index 0177fbb7db..6b5585fbf9 100644 --- a/.ci/jenkins/lib/build-matrix.yaml +++ b/.ci/jenkins/lib/build-matrix.yaml @@ -43,7 +43,7 @@ env: TEST_TIMEOUT: 30 UCX_TLS: "^shm" STORAGE_DRIVER: 'overlay' - CI_IMAGE_TAG: "20260219-1" + CI_IMAGE_TAG: "20260226-1" runs_on_dockers: diff --git a/.ci/jenkins/lib/test-matrix.yaml b/.ci/jenkins/lib/test-matrix.yaml index 1fe42a0bdc..5f627c111f 100644 --- a/.ci/jenkins/lib/test-matrix.yaml +++ b/.ci/jenkins/lib/test-matrix.yaml @@ -49,7 +49,7 @@ env: SLURM_JOB_TIMEOUT: '02:20:00' SLURM_IMMEDIATE_TIMEOUT: "3600" STORAGE_DRIVER: overlay - CI_IMAGE_TAG: "20260219-1" + CI_IMAGE_TAG: "20260226-1" empty_volumes: - {mountPath: /var/lib/containers/storage, memory: false} From bb7b0a7fac86b681a1adb044fe369446f2720d73 Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Thu, 26 Feb 2026 21:58:48 +0530 Subject: [PATCH 27/28] Revert CI_IMAGE_TAG Signed-off-by: Pravein Govindan Kannan --- .ci/jenkins/lib/build-matrix.yaml | 2 +- .ci/jenkins/lib/test-matrix.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.ci/jenkins/lib/build-matrix.yaml b/.ci/jenkins/lib/build-matrix.yaml index 6b5585fbf9..0177fbb7db 100644 --- a/.ci/jenkins/lib/build-matrix.yaml +++ b/.ci/jenkins/lib/build-matrix.yaml @@ -43,7 +43,7 @@ env: TEST_TIMEOUT: 30 UCX_TLS: "^shm" STORAGE_DRIVER: 'overlay' - CI_IMAGE_TAG: "20260226-1" + CI_IMAGE_TAG: "20260219-1" runs_on_dockers: diff --git a/.ci/jenkins/lib/test-matrix.yaml b/.ci/jenkins/lib/test-matrix.yaml index 5f627c111f..1fe42a0bdc 100644 --- a/.ci/jenkins/lib/test-matrix.yaml +++ b/.ci/jenkins/lib/test-matrix.yaml @@ -49,7 +49,7 @@ env: SLURM_JOB_TIMEOUT: '02:20:00' SLURM_IMMEDIATE_TIMEOUT: "3600" STORAGE_DRIVER: overlay - CI_IMAGE_TAG: "20260226-1" + CI_IMAGE_TAG: "20260219-1" empty_volumes: - {mountPath: /var/lib/containers/storage, memory: false} From 872e33a4b3b80e4eb3cc5221c29a13399204420b Mon Sep 17 00:00:00 2001 From: Pravein Govindan Kannan Date: Fri, 27 Feb 2026 06:42:49 +0530 Subject: [PATCH 28/28] Update IMAGE_TAG Signed-off-by: Pravein Govindan Kannan --- .ci/jenkins/lib/build-matrix.yaml | 2 +- .ci/jenkins/lib/test-matrix.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.ci/jenkins/lib/build-matrix.yaml b/.ci/jenkins/lib/build-matrix.yaml index 0177fbb7db..6b5585fbf9 100644 --- a/.ci/jenkins/lib/build-matrix.yaml +++ b/.ci/jenkins/lib/build-matrix.yaml @@ -43,7 +43,7 @@ env: TEST_TIMEOUT: 30 UCX_TLS: "^shm" STORAGE_DRIVER: 'overlay' - CI_IMAGE_TAG: "20260219-1" + CI_IMAGE_TAG: "20260226-1" runs_on_dockers: diff --git a/.ci/jenkins/lib/test-matrix.yaml b/.ci/jenkins/lib/test-matrix.yaml index 1fe42a0bdc..5f627c111f 100644 --- a/.ci/jenkins/lib/test-matrix.yaml +++ b/.ci/jenkins/lib/test-matrix.yaml @@ -49,7 +49,7 @@ env: SLURM_JOB_TIMEOUT: '02:20:00' SLURM_IMMEDIATE_TIMEOUT: "3600" STORAGE_DRIVER: overlay - CI_IMAGE_TAG: "20260219-1" + CI_IMAGE_TAG: "20260226-1" empty_volumes: - {mountPath: /var/lib/containers/storage, memory: false}