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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/source/design/transfer-engine/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -487,6 +487,7 @@ For advanced users, TransferEngine provides the following advanced runtime optio
- `MC_ENABLE_DEST_DEVICE_AFFINITY` Enable device affinity for RDMA performance optimization. When enabled, Transfer Engine will prioritize communication with remote NICs that have the same name as local NICs to reduce QP count and improve network performance in rail-optimized topologies. The default value is false
- `MC_TRACK_RDMA_POSTED_SLICES` Enable RDMA posted-slice tracking for timeout diagnostics. When enabled, CQ timeout logs include stuck transfer groups by peer NIC path, slice count, bytes, oldest post age, and sample addresses. This adds synchronization on the RDMA post and poll hot paths, so it is disabled by default and should be enabled only while diagnosing stuck completions.
- `MC_ENABLE_PARALLEL_REG_MR` Control parallel memory region registration across multiple RDMA NICs. Valid values: -1 (auto, default), 0 (disabled), 1 (enabled). When set to -1, parallel registration is automatically enabled when multiple RNICs exist and memory has been pre-touched. Note: If memory hasn't been touched before registration, parallel registration can be slower than sequential registration
- `MC_MAX_CONCURRENT_REG_MR` Cap on how many buffers `registerLocalMemoryBatch` registers concurrently (EFA transport). The default 0 means unbounded — one thread per buffer, the historical behavior. Note the cap is **per process**, so a framework running one `TransferEngine` per TP rank multiplies it by the rank count. Capping can cut registration time substantially when a batch holds many large GPU buffers. Registration is CPU-bound, so a reasonable value is `cores / processes-per-node` — on a 192-core node running 8 ranks, around 16. Oversubscribing costs more than undersubscribing, and the result also depends on the order the caller passes buffers in, so a poorly chosen cap can be slower than unbounded — hence opt-in.
- `MC_FORCE_HCA` Force to use RDMA as the active transport, return error if no HCA has been found.
- `MC_FORCE_MNNVL` Force to use Multi-Node NVLink as the active transport regardless whether RDMA devices are installed.
- `MC_INTRA_NVLINK` Enable intra-node NVLINK transport, and cannot be used together with MC_FORCE_MNNVL.
Expand Down
5 changes: 5 additions & 0 deletions mooncake-transfer-engine/include/config.h
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,11 @@ struct GlobalConfig {
bool log_rdma_slice_affinity = false;
bool track_rdma_posted_slices = false;
int parallel_reg_mr = -1;
// Cap on concurrent buffer registrations in registerLocalMemoryBatch().
// 0 (default) = unbounded, one thread per buffer. Set via
// MC_MAX_CONCURRENT_REG_MR; the best value is platform-specific, see the
// measured tables in efa_transport.cpp before choosing one.
size_t max_concurrent_reg_mr = 0;
size_t eic_max_block_size = 64UL * 1024 * 1024;
EndpointStoreType endpoint_store_type = EndpointStoreType::SIEVE;
int ib_traffic_class = -1;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,12 @@ class EfaTransport : public Transport {
int unregisterLocalMemoryBatch(
const std::vector<void*>& addr_list) override;

// Indices into `buffer_list` ordered by ascending length -- the order
// registerLocalMemoryBatch() hands the batch to its workers. Exposed for
// testing; see the definition for why ascending is optimal.
static std::vector<size_t> registrationOrder(
const std::vector<BufferEntry>& buffer_list);

// Eagerly populate the address vector with every (local_ctx, peer_nic)
// handshake for `segment_name`.
//
Expand Down
19 changes: 19 additions & 0 deletions mooncake-transfer-engine/src/config.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -538,6 +538,25 @@ void loadGlobalConfig(GlobalConfig& config) {
}
}

const char* max_concurrent_reg_mr = std::getenv("MC_MAX_CONCURRENT_REG_MR");
if (max_concurrent_reg_mr) {
// Robust parse (not atol): a non-numeric typo must keep the default
// rather than silently resolve to 0, which here means "no cap" and so
// would read as a deliberate request for the old unbounded behavior.
// 0 is a valid explicit way to ask for no cap; negative and garbage are
// rejected.
size_t val = 0;
const char* end = max_concurrent_reg_mr + strlen(max_concurrent_reg_mr);
auto [ptr, ec] = std::from_chars(max_concurrent_reg_mr, end, val);
if (ec == std::errc() && ptr == end) {
config.max_concurrent_reg_mr = val;
} else {
LOG(WARNING) << "Invalid MC_MAX_CONCURRENT_REG_MR environment "
"value: "
<< max_concurrent_reg_mr << ", keeping default";
}
}

const char* endpoint_store_type_env = std::getenv("MC_ENDPOINT_STORE_TYPE");
if (endpoint_store_type_env) {
if (strcmp(endpoint_store_type_env, "FIFO") == 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,15 @@
#include <unistd.h>

#include <algorithm>
#include <atomic>
#include <cassert>
#include <chrono>
#include <cstddef>
#include <cstdlib>
#include <fstream>
#include <functional>
#include <future>
#include <limits>
#include <set>
#include <thread>

Expand Down Expand Up @@ -471,20 +474,22 @@ int EfaTransport::registerLocalMemoryInternal(void* addr, size_t length,
reg_start)
.count();

// Per-chunk detail is trace-only: a batch registration emits one of
// these per chunk, which is ~1450 lines at Kimi-K3 startup. Note that
// reg_duration_ms is the wall time this thread spent in the call, so
// when several buffers register concurrently it includes time waiting
// on the provider's locks, not just this chunk's own work -- compare it
// against the batch total logged by registerLocalMemoryBatch().
if (globalConfig().trace) {
LOG(INFO) << "EFA registerMemoryRegion: chunk " << ci
<< ", addr=" << chunk_addr << ", length=" << chunk_len
LOG(INFO) << "EFA registerMemoryRegion: chunk " << ci << "/"
<< chunks.size() << ", addr=" << chunk_addr
<< ", length=" << chunk_len
<< ", nics=" << assigned_nics.size() << "/"
<< context_list_.size()
<< ", parallel=" << (use_parallel_reg ? "true" : "false")
<< ", duration=" << reg_duration_ms << "ms";
}

LOG(WARNING) << "Chunk " << ci << "/" << chunks.size()
<< " registered on " << assigned_nics.size() << " NICs"
<< ", addr=" << chunk_addr << ", length=" << chunk_len
<< ", duration=" << reg_duration_ms << "ms";

// Collect keys: assigned NICs have valid keys, others get 0
BufferDesc buffer_desc;
for (auto& context : context_list_) {
Expand Down Expand Up @@ -624,52 +629,162 @@ int EfaTransport::allocateLocalSegmentID() {
return 0;
}

// Optional concurrency cap for the batch register/unregister fan-out below, set
// via MC_MAX_CONCURRENT_REG_MR. Unset (0) means unbounded -- one thread per
// buffer, which is what these entry points have always done.
//
// The cap is PER PROCESS, and registration is CPU-bound page-pinning, so what
// matters is cap x processes against the core count. Measured on p5.48xlarge
// (192 cores) replaying Kimi-K3's KV registration as SGLang issues it -- one
// TransferEngine per TP rank, 182 GPU buffers of 2.5 KB to 391 MB each, slowest
// rank -- the optimum tracks the cores and not the cap:
//
// 8 ranks, 192 cores -> cap 16 (99.8 s); cap 64 is 1.8x slower
// 4 ranks, 192 cores -> cap 32 (89.2 s); same 128 threads as above
// 8 ranks, 64 cores -> cap 8 (121.3 s); 64 threads, tracks the budget
//
// So a good value is roughly cores/processes. Oversubscribing costs more than
// undersubscribing. This cannot be a built-in default because a single engine
// does not know how many peer processes share the node; picking one from the
// core count alone would oversubscribe by exactly the rank count.
//
// Input order matters as much as the cap: at cap 16 the same batch takes 43 s
// in SGLang's pool order but 95-98 s largest-first. Unbounded ignores order
// (99-128 s) since nothing queues. Largest-first being worst is backwards from
// longest-processing-time scheduling and is unexplained, so no sort is applied
// here yet -- another reason the cap stays opt-in rather than a default.
static size_t maxConcurrentRegMr() {
size_t configured = globalConfig().max_concurrent_reg_mr;
// 0 (unset) means unbounded, i.e. one thread per buffer as before.
return configured > 0 ? configured : std::numeric_limits<size_t>::max();
}

// Run `fn(i)` for i in [0, count) on at most maxConcurrentRegMr() threads,
// returning the first non-zero result (all items are still attempted).
//
// With no cap set this spawns count-1 threads and runs the caller as a worker,
// which reproduces what both batch entry points did before: one
// std::async(std::launch::async) per buffer, which libstdc++ honours literally
// as one fresh thread each. Kimi-K3 registers ~1450 KV buffers at once, so a
// 192-core node peaks at ~1400 runnable threads inside fi_mr_regattr, and the
// per-buffer duration this logs inflates with the queueing delay of the ones
// ahead of it -- on 2x p6-b300 the median reached 106 s while the whole batch
// took 138 s. That inflation is not by itself a reason to cap: a badly chosen
// cap is slower still, see maxConcurrentRegMr().
//
// A plain thread pool rather than a semaphore over std::async, because if a cap
// is set the point is to avoid the thread *creation*, not just to gate entry
// into the provider -- admission control after the thread already exists would
// leave that cost in place.
static int runBoundedParallel(size_t count,
const std::function<int(size_t)>& fn) {
if (count == 0) return 0;

size_t workers = std::min(count, maxConcurrentRegMr());
std::atomic<size_t> next{0};
std::atomic<int> first_error{0};

auto worker = [&]() {
for (size_t i = next.fetch_add(1); i < count; i = next.fetch_add(1)) {
int ret = fn(i);
if (ret) {
int expected = 0;
first_error.compare_exchange_strong(expected, ret);
}
}
};

std::vector<std::thread> threads;
threads.reserve(workers - 1);
for (size_t w = 1; w < workers; ++w) threads.emplace_back(worker);
worker(); // the caller is a worker too
for (auto& t : threads) t.join();

return first_error.load();
}

// Smallest buffer first, because registering DEVICE memory on EFA costs roughly
// k x (bytes of device memory already registered on this domain): measured at
// ~260 ms/GiB on p5.48xlarge, flat from 0 to 19 GiB, and nearly independent of
// the buffer's own size. The same 2.79 MB CUDA buffer takes 70 ms on an empty
// domain but 2505 ms once 9.17 GiB is in.
//
// So a buffer's bytes are charged once per buffer registered after it, and the
// total is minimized by registering the large ones last. Ascending is a
// well-defined optimum, not a heuristic: registering the same 48 GPU buffers
// (9.23 GiB) serially takes 35.5 s ascending against 93.3 s descending.
//
// Host memory does not accumulate -- the same sweep on mmap'd host buffers is
// flat (a 391 MB buffer costs ~500-600 ms whether it is 1st or 48th) and order
// makes no difference (14-16 s either way, ordering within run-to-run spread).
// Sorting is therefore a no-op there rather than a regression, so it is applied
// unconditionally instead of only for VRAM.
//
// Only matters when MC_MAX_CONCURRENT_REG_MR is set. Unbounded, every buffer
// gets its own thread and nothing queues, so this order does not reach the
// provider (measured 112 s vs 120 s, i.e. noise). With a cap it is worth 3.1x
// on a descending batch.
//
// Ties keep the caller's relative order so the dispatch order stays a stable
// function of the input.
std::vector<size_t> EfaTransport::registrationOrder(
const std::vector<EfaTransport::BufferEntry>& buffer_list) {
std::vector<size_t> order(buffer_list.size());
for (size_t i = 0; i < order.size(); ++i) order[i] = i;
std::stable_sort(order.begin(), order.end(), [&](size_t a, size_t b) {
return buffer_list[a].length < buffer_list[b].length;
});
return order;
}

int EfaTransport::registerLocalMemoryBatch(
const std::vector<EfaTransport::BufferEntry>& buffer_list,
const std::string& location) {
std::vector<std::future<int>> results;
for (auto& buffer : buffer_list) {
results.emplace_back(
std::async(std::launch::async, [this, buffer, location]() -> int {
return registerLocalMemoryInternal(buffer.addr, buffer.length,
location, true, false, true);
}));
}

int first_error = 0;
for (size_t i = 0; i < buffer_list.size(); ++i) {
int ret = results[i].get();
auto start = std::chrono::steady_clock::now();

// Dispatch order only, nothing observable changes. A buffer's NIC set is
// derived from its own length and the NIC count, so no buffer lands
// anywhere else; and segment_desc->buffers was already in whatever order
// the workers happened to finish in, since they append concurrently.
const std::vector<size_t> order = registrationOrder(buffer_list);

int first_error = runBoundedParallel(buffer_list.size(), [&](size_t k) {
const size_t i = order[k];
int ret = registerLocalMemoryInternal(buffer_list[i].addr,
buffer_list[i].length, location,
true, false, true);
if (ret) {
LOG(WARNING) << "EfaTransport: Failed to register memory: addr "
<< buffer_list[i].addr << " length "
<< buffer_list[i].length;
if (!first_error) first_error = ret;
}
}
return ret;
});

auto elapsed = std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now() - start)
.count();
LOG(INFO) << "EfaTransport: registered " << buffer_list.size()
<< " buffers on "
<< std::min(buffer_list.size(), maxConcurrentRegMr())
<< " threads in " << elapsed << "ms";

if (first_error) return first_error;

return metadata_->updateLocalSegmentDesc();
}

int EfaTransport::unregisterLocalMemoryBatch(
const std::vector<void*>& addr_list) {
std::vector<std::future<int>> results;
for (auto& addr : addr_list) {
results.emplace_back(
std::async(std::launch::async, [this, addr]() -> int {
return unregisterLocalMemoryInternal(addr, false, true);
}));
}

int first_error = 0;
for (size_t i = 0; i < addr_list.size(); ++i) {
int ret = results[i].get();
int first_error = runBoundedParallel(addr_list.size(), [&](size_t i) {
int ret = unregisterLocalMemoryInternal(addr_list[i], false, true);
if (ret) {
LOG(WARNING) << "EfaTransport: Failed to unregister memory: addr "
<< addr_list[i];
if (!first_error) first_error = ret;
}
}
return ret;
});

int metadata_ret = metadata_->updateLocalSegmentDesc();
return first_error ? first_error : metadata_ret;
}
Expand Down
71 changes: 71 additions & 0 deletions mooncake-transfer-engine/tests/config_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -277,5 +277,76 @@ TEST_F(ConnPauseTtlEnvTest, EmptyStringKeepsDefault) {
EXPECT_EQ(config.conn_pause_ttl_ms, 17);
}

// MC_MAX_CONCURRENT_REG_MR caps how many buffers registerLocalMemoryBatch()
// registers at once; 0 (the default) means unbounded. 0 is therefore also what
// a silent atol() fallback would produce on a typo, which would read as "the
// knob was honored and asked for no cap" -- the opposite of what the operator
// wanted. So a typo must be rejected loudly and leave the field untouched.
class MaxConcurrentRegMrEnvTest : public ::testing::Test {
protected:
void TearDown() override { ::unsetenv("MC_MAX_CONCURRENT_REG_MR"); }
};

TEST_F(MaxConcurrentRegMrEnvTest, UnboundedWhenUnset) {
::unsetenv("MC_MAX_CONCURRENT_REG_MR");
GlobalConfig config;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 0u);
}

TEST_F(MaxConcurrentRegMrEnvTest, ValidOverrideIsApplied) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "8", 1), 0);
GlobalConfig config;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 8u);
}

TEST_F(MaxConcurrentRegMrEnvTest, ExplicitZeroSelectsUnbounded) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "0", 1), 0);
GlobalConfig config;
config.max_concurrent_reg_mr = 99; // sentinel must be overwritten by 0
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 0u);
}

TEST_F(MaxConcurrentRegMrEnvTest, OneIsAcceptedAndSerializes) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "1", 1), 0);
GlobalConfig config;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 1u);
}

TEST_F(MaxConcurrentRegMrEnvTest, NegativeIsIgnored) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "-1", 1), 0);
GlobalConfig config;
config.max_concurrent_reg_mr = 11;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 11u);
}

TEST_F(MaxConcurrentRegMrEnvTest, NonNumericKeepsDefault) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "abc", 1), 0);
GlobalConfig config;
config.max_concurrent_reg_mr = 13;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 13u);
}

TEST_F(MaxConcurrentRegMrEnvTest, NumericSuffixKeepsDefault) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "8x", 1), 0);
GlobalConfig config;
config.max_concurrent_reg_mr = 15;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 15u);
}

TEST_F(MaxConcurrentRegMrEnvTest, EmptyStringKeepsDefault) {
ASSERT_EQ(::setenv("MC_MAX_CONCURRENT_REG_MR", "", 1), 0);
GlobalConfig config;
config.max_concurrent_reg_mr = 17;
loadGlobalConfig(config);
EXPECT_EQ(config.max_concurrent_reg_mr, 17u);
}

} // namespace
} // namespace mooncake
Loading
Loading