Repository navigation
refactor(multimodal): extract engine-neutral RDMA pixel transport into crates/mm_rdma - #1908
Conversation
|
Warning Review limit reachedYou’ve reached a temporary PR review limit under our Fair Usage Limits Policy. Next review available in: 18 minutes Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available. How can I continue?After more reviews become available, a review can be triggered using the To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews. How do review limits work?CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability. For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window. Please refer docs for additional details. Review details⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (3)
📝 WalkthroughWalkthroughA new ChangesMultimodal RDMA export
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant ProtoWrapper
participant RdmaExporter
participant SlotPool
participant NIXL
ProtoWrapper->>RdmaExporter: export(slot_key, pixels)
RdmaExporter->>SlotPool: lease and frame payload
RdmaExporter->>NIXL: expose registered arena descriptor
NIXL-->>ProtoWrapper: return remote descriptor
RdmaExporter->>SlotPool: reclaim notified or expired slot
Possibly related PRs
Suggested labels: Suggested reviewers: Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 864716c596
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| let gw = init_agent(&cfg)?; | ||
| let agent = Mutex::new(gw); | ||
| let arena = build_arena(&agent, &cfg)?; |
There was a problem hiding this comment.
Skip RDMA initialization without a listener IP
When the mm-rdma feature is enabled and the lane is selected (SMG_MM_TENSOR_TRANSPORT=rdma or SMG_MM_PIXEL_RDMA=true) but SMG_RDMA_LISTEN_IP is unset, export() later returns inline because listen_ip is empty; however this constructor has already created the NIXL agent and registered the default 2 GiB arena. The old path checked the empty listener IP before building the arena, so this common misconfiguration stayed a cheap inline fallback instead of consuming memory/binding a listener or failing initialization. Treat an empty cfg.listen_ip as unavailable before calling init_agent/build_arena.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Code Review
This pull request refactors the multimodal pixel RDMA (NIXL) transport by moving it from the model_gateway crate into a dedicated, engine-neutral workspace crate smg-mm-rdma. This decouples the NIXL mechanics and wire format from the gateway's policy, providing a real implementation under the nixl feature and a no-op stub otherwise. The review feedback highlights several critical improvement opportunities: switching from a LIFO stack to a FIFO queue for slot allocation to reduce recycling race conditions, using checked_duration_since to prevent panics from potential multi-socket CPU clock drift, and propagating thread-spawning errors from the background reaper to prevent silent pool exhaustion.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| base: usize, | ||
| slot_bytes: usize, | ||
| /// Available slot indices. | ||
| free: Mutex<Vec<u32>>, |
There was a problem hiding this comment.
Using a LIFO stack (Vec::pop and Vec::push) for slot allocation causes the same few slots to be repeatedly reused under low concurrency. This significantly increases the likelihood of a race condition where a slot is recycled while a worker's slow or delayed RDMA READ is still in-flight, triggering a generation mismatch and forcing a fallback to the inline path. Switching to a FIFO queue (VecDeque) provides round-robin reuse across all slots, maximizing the safety window before any slot is recycled.
free: Mutex<std::collections::VecDeque<u32>>,| let stale: Vec<i64> = self | ||
| .occupied | ||
| .iter() | ||
| .filter(|e| now.duration_since(e.value().at) >= ttl) | ||
| .map(|e| *e.key()) | ||
| .collect(); |
There was a problem hiding this comment.
Using Instant::duration_since can panic if now is earlier than e.value().at. While Instant is monotonic, in multi-threaded environments with multi-socket CPUs, minor TSC drift across cores can cause Instant::now() in the reaper thread to be slightly behind the Instant recorded during the lease on the hot path. To prevent potential panics, use checked_duration_since. Additionally, to avoid silently ignoring this failure, use unwrap_or_else to log a warning instead of unwrap_or_default.
let stale: Vec<i64> = self
.occupied
.iter()
.filter(|e| {
now.checked_duration_since(e.value().at)
.unwrap_or_else(|| {
log::warn!("Instant::now() is earlier than slot lease time");
std::time::Duration::ZERO
}) >= ttl
})
.map(|e| *e.key())
.collect();References
- Instead of silently ignoring potential failures, log them as warnings to aid in debugging. In Rust, prefer using
unwrap_or_elseto log an error overunwrap_or_defaultwhich would fail silently. - To avoid potential deadlocks in DashMap when iterating and removing entries, collect the keys into a separate collection before removing them.
| fn spawn_reaper(inner: Arc<ExporterInner>) { | ||
| std::thread::Builder::new() | ||
| .name("epd-rdma-reaper".into()) | ||
| .spawn(move || loop { | ||
| std::thread::sleep(REAPER_TICK); | ||
| if let Ok(mut notifs) = NotificationMap::new() { | ||
| { | ||
| let guard = inner.agent.lock(); | ||
| let _ = guard.agent.get_notifications(&mut notifs, None); | ||
| } | ||
| if let Ok(map) = notifs.take_notifs() { | ||
| for (_agent, tags) in map { | ||
| for tag in tags { | ||
| if let Ok(key) = tag.parse::<i64>() { | ||
| inner.arena.pool.free_slot_key(key); | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
| // TTL sweep: reclaim slots whose READ-notif never arrived. `slot_ttl` is | ||
| // derived (by the gateway) to exceed the worker's max hold, so this never | ||
| // races a live READ. | ||
| let _ = inner | ||
| .arena | ||
| .pool | ||
| .reap_stale(Instant::now(), inner.cfg.slot_ttl); | ||
| }) | ||
| .map_err(|e| error!(error = ?e, "EPD RDMA: reaper thread spawn failed")) | ||
| .ok(); | ||
| } |
There was a problem hiding this comment.
The background reaper thread is critical for reclaiming leased slots. If spawn_reaper fails to spawn the thread (e.g., due to OS thread limits or OOM), the failure is logged but ignored, and RdmaExporter::new still returns Ok. Without the reaper, slots will never be reclaimed, leading to immediate pool exhaustion. Consider propagating the spawn error up to RdmaExporter::new so the gateway can safely fall back to the inline path.
fn spawn_reaper(inner: Arc<ExporterInner>) -> Result<(), std::io::Error> {
std::thread::Builder::new()
.name("epd-rdma-reaper".into())
.spawn(move || loop {
std::thread::sleep(REAPER_TICK);
if let Ok(mut notifs) = NotificationMap::new() {
{
let guard = inner.agent.lock();
let _ = guard.agent.get_notifications(&mut notifs, None);
}
if let Ok(map) = notifs.take_notifs() {
for (_agent, tags) in map {
for tag in tags {
if let Ok(key) = tag.parse::<i64>() {
inner.arena.pool.free_slot_key(key);
}
}
}
}
}
// TTL sweep: reclaim slots whose READ-notif never arrived. `slot_ttl` is
// derived (by the gateway) to exceed the worker's max hold, so this never
// races a live READ.
let _ = inner
.arena
.pool
.reap_stale(Instant::now(), inner.cfg.slot_ttl);
})?;
Ok(())
}There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/mm_rdma/src/lib.rs`:
- Around line 43-46: Update build_rdma_config_from_env() to clamp or reject
SMG_RDMA_SLOT_TTL_S values below worker_max_hold() + RDMA_SLOT_TTL_SLACK,
ensuring slot_ttl cannot violate the worker hold-time requirement. Add a gateway
test covering an environment override below this minimum and its expected
handling.
In `@crates/mm_rdma/src/nixl.rs`:
- Around line 198-218: Update build_arena so the raw arena allocation is
reclaimed whenever OptArgs::new, add_backend, or register_memory fails. Retain
ownership of the boxed buffer until registration succeeds, then transfer it with
Box::into_raw; ensure successful SlotArena construction still keeps the
allocation alive for the process lifetime.
In `@model_gateway/src/routers/grpc/multimodal/transport.rs`:
- Around line 213-224: Bound the values assigned to pool_slots and slot_bytes in
build_rdma_config_from_env before they reach
cfg.pool_slots.saturating_mul(cfg.slot_bytes) and the arena allocator. Reuse or
extend the existing rdma_env_positive parsing flow to clamp or reject
configurations whose combined arena size exceeds a safe maximum, preserving
normal valid configurations and allowing oversized settings to fall back safely
instead of attempting an unbounded allocation.
- Around line 238-250: Validate the parsed SMG_RDMA_SLOT_TTL_S override in
derive_rdma_slot_ttl so it cannot be less than or equal to worker_max_hold();
only accept values that preserve the documented TTL invariant, otherwise fall
back to the derived worker_max_hold() + RDMA_SLOT_TTL_SLACK value.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: d05ca5ba-de0b-4553-9492-44268b9bd415
📒 Files selected for processing (14)
Cargo.tomlcrates/mm_rdma/Cargo.tomlcrates/mm_rdma/src/lib.rscrates/mm_rdma/src/nixl.rscrates/mm_rdma/src/slot_pool.rscrates/mm_rdma/src/stub.rsmodel_gateway/Cargo.tomlmodel_gateway/src/routers/grpc/mm_rdma/mod.rsmodel_gateway/src/routers/grpc/mm_rdma/nixl.rsmodel_gateway/src/routers/grpc/mm_rdma/stub.rsmodel_gateway/src/routers/grpc/mod.rsmodel_gateway/src/routers/grpc/multimodal/mod.rsmodel_gateway/src/routers/grpc/multimodal/transport.rsmodel_gateway/src/routers/grpc/proto_wrapper.rs
💤 Files with no reviewable changes (4)
- model_gateway/src/routers/grpc/mm_rdma/mod.rs
- model_gateway/src/routers/grpc/mm_rdma/nixl.rs
- model_gateway/src/routers/grpc/mm_rdma/stub.rs
- model_gateway/src/routers/grpc/mod.rs
| /// How long a leased slot may live without a free-notif before the reaper | ||
| /// force-reclaims it. MUST exceed the worker's max hold (the gateway derives it | ||
| /// so); the per-lease gen framing makes correctness independent of the value. | ||
| pub slot_ttl: Duration, |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | ⚡ Quick win
Enforce the minimum TTL for environment overrides.
build_rdma_config_from_env() accepts SMG_RDMA_SLOT_TTL_S directly. A value below the worker wait-plus-READ window lets the reaper reclaim and reuse a slot while its descriptor remains valid. Clamp or reject overrides below worker_max_hold() + RDMA_SLOT_TTL_SLACK, and cover that override path with a gateway test.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/mm_rdma/src/lib.rs` around lines 43 - 46, Update
build_rdma_config_from_env() to clamp or reject SMG_RDMA_SLOT_TTL_S values below
worker_max_hold() + RDMA_SLOT_TTL_SLACK, ensuring slot_ttl cannot violate the
worker hold-time requirement. Add a gateway test covering an environment
override below this minimum and its expected handling.
…o crates/mm_rdma Move the NIXL multimodal pixel-transport mechanics out of the gateway (model_gateway/src/routers/grpc/mm_rdma) into a standalone, engine-neutral crate `smg-mm-rdma`. Behavior-preserving; the descriptor wire format is unchanged. - The crate owns only the NIXL mechanics + wire format and reads no env and no router globals. The gateway injects all policy via `RdmaConfig` and constructs one exporter (`RdmaExporter::new` / `.export`). - The optional native dependency (nixl-sys, bindgen-heavy) is isolated behind the crate's `nixl` feature; the default build compiles a no-op stub, so ordinary gateway builds need no clang/bindgen/libnixl. The gateway's `mm-rdma` feature now reads as `["smg-mm-rdma/nixl"]`. - Gateway owns the exporter singleton + env-derived config (multimodal/transport.rs): `mm_rdma_exporter`, `rdma_lane_enabled`, `build_rdma_config_from_env`, and the slot-TTL derivation moved here from the old module. - proto_wrapper `try_export_*` call `mm_rdma_exporter().export(slot_key, ..)`; `room` renamed to `slot_key` (Rust-internal; wire layout unchanged). - The pure `SlotPool` (free-list + lease/reclaim + gen framing) is split out for hardware-free unit tests (6 tests); the TTL-derivation policy test moves to the gateway (transport.rs). Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
864716c to
e9aec76
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: e9aec76eaf
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| pub fn try_export_nixl_remote(self, room: i64) -> Self { | ||
| if !mm_rdma::rdma_enabled() { | ||
| pub fn try_export_nixl_remote(self, slot_key: i64) -> Self { | ||
| let Some(exporter) = mm_rdma_exporter() else { |
There was a problem hiding this comment.
Defer RDMA exporter lookup until after payload checks
When the RDMA lane is enabled, mm_rdma_exporter() constructs the NIXL agent and registers the default 2 GiB arena, but this call now happens before checking whether the tensor is actually inline and non-empty. In configurations such as legacy SMG_MM_PIXEL_RDMA=true with SHM/auto-local transport, or for already-SHM/remote/empty encoder inputs, the function immediately returns the original payload below while still pinning the RDMA arena; the old rdma_enabled() guard was only a boolean and export_pixel_buffer was reached only after these payload checks. Move the exporter lookup to just before exporter.export and avoid eager preflight checks so fallback transports do not allocate RDMA resources.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/routers/grpc/multimodal/transport.rs`:
- Around line 185-195: Update the exporter initialization closure around
RdmaExporter::new to build the RDMA configuration first and return None when its
listen_ip is empty or unset. Only construct RdmaExporter after validating a
non-empty listener IP, while preserving the existing rdma_lane_enabled check and
initialization error fallback.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 30234463-128b-49f0-8f0e-a1b6748c4ec1
📒 Files selected for processing (14)
Cargo.tomlcrates/mm_rdma/Cargo.tomlcrates/mm_rdma/src/lib.rscrates/mm_rdma/src/nixl.rscrates/mm_rdma/src/slot_pool.rscrates/mm_rdma/src/stub.rsmodel_gateway/Cargo.tomlmodel_gateway/src/routers/grpc/mm_rdma/mod.rsmodel_gateway/src/routers/grpc/mm_rdma/nixl.rsmodel_gateway/src/routers/grpc/mm_rdma/stub.rsmodel_gateway/src/routers/grpc/mod.rsmodel_gateway/src/routers/grpc/multimodal/mod.rsmodel_gateway/src/routers/grpc/multimodal/transport.rsmodel_gateway/src/routers/grpc/proto_wrapper.rs
💤 Files with no reviewable changes (4)
- model_gateway/src/routers/grpc/mm_rdma/mod.rs
- model_gateway/src/routers/grpc/mod.rs
- model_gateway/src/routers/grpc/mm_rdma/stub.rs
- model_gateway/src/routers/grpc/mm_rdma/nixl.rs
| /// its own room-matched export path because the room must also be injected | ||
| /// into the encode->prefill handshake. |
There was a problem hiding this comment.
🟡 Nit: Stale doc comment — the PR renamed room → slot_key throughout but this comment still says "room-matched export path because the room must also be injected". Should be updated to match the new terminology (e.g. "slot_key-matched").
Address bot review comments on the crate extraction (all behavior-preserving hardening; wire format unchanged): - Skip building the NIXL agent + arena when SMG_RDMA_LISTEN_IP is unset: the lane can't do the cross-node exchange without it, so exports would fall back to inline anyway. Restores the old "empty IP => cheap inline" behavior instead of eagerly allocating the 2 GiB arena. (Codex) - SlotPool free-list is now a FIFO VecDeque, not a LIFO Vec: round-robin reuse maximizes the window before any freed slot is re-leased under a slow worker READ. (Gemini) - reap_stale uses saturating_duration_since to avoid a potential panic if the reaper ever observes `now` before a lease stamp (cross-core clock skew). (Gemini) - Propagate the reaper thread spawn error out of RdmaExporter::new (new RdmaError::Reaper variant) so the gateway falls back to inline rather than running with a dead reaper that would leak every leased slot. (Gemini) - build_arena frees the arena buffer (leaked via Box::into_raw) on the NIXL registration error path instead of burning up to pool_slots*slot_bytes for the process lifetime. (CodeRabbit) - Reject an SMG_RDMA_SLOT_TTL_S override that does not exceed worker_max_hold (would reintroduce the recycled-under-READ cross-wire); fall back to the derived hold+slack with a warning. Extracted resolve_slot_ttl as a pure fn with a unit test. (CodeRabbit, Gemini) - Clamp the arena size (pool_slots * slot_bytes) to MAX_RDMA_ARENA_BYTES so a fat-fingered env pair can't turn into an unbounded startup allocation that aborts the process; unit-tested. (CodeRabbit) Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 1d96253f08
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| kv-index.workspace = true | ||
| smg-data-connector.workspace = true | ||
| llm-multimodal = { workspace = true, default-features = false } | ||
| smg-mm-rdma.workspace = true |
There was a problem hiding this comment.
Add smg-mm-rdma to the crates release tiers
This adds a published smg dependency on smg-mm-rdma, but the crates release workflow still publishes only the existing tier1/tier2 crates before smg (.github/workflows/release-crates.yml:17-84 has no smg-mm-rdma entry; I also checked repo-wide with rg smg-mm-rdma). When model_gateway is published, Cargo will strip the path dependency and require smg-mm-rdma = 0.1.0 from crates.io, so the tier3 cargo publish for smg will fail unless the new crate is published first (and ideally added to the release-version registry too).
Useful? React with 👍 / 👎.
Description
Problem
RDMA/NIXL is the cross-node transport for large multimodal tensors (image/video
pixel_values), but unlike the engine-neutralinlineand/dev/shmtransports it is coupled to the gateway internals: the NIXL mechanics live insidemodel_gateway/src/routers/grpc/mm_rdma, read env vars and router globals directly, and drag the bindgen-heavynixl-sysnative dependency into thesmgcrate. That coupling blocks making RDMA a first-class, engine-neutral transport (and adding vLLM support).Solution
Extract the NIXL pixel-transport mechanics into a standalone, engine-neutral crate
smg-mm-rdma, with all policy injected by the gateway via a config struct. Behavior-preserving — theSMGRDMA1descriptor wire format is byte-for-byte unchanged. This is PR1 of a 3-PR stack (move → generalize → wire vLLM).Changes
crates/mm_rdma(smg-mm-rdma) owns only the NIXL mechanics + wire format. It reads no env and no router globals — the gateway injects all policy viaRdmaConfigand constructs oneRdmaExporter(RdmaExporter::new(cfg)/.export(slot_key, bytes)).nixl-sysmoves behind the crate'snixlfeature; the default build compiles a no-op stub, so ordinary gateway builds need no clang/bindgen/libnixl. The gateway'smm-rdmafeature now reads as["smg-mm-rdma/nixl"](matching theopencv-video = ["llm-multimodal/opencv-video"]passthrough style).multimodal/transport.rs):mm_rdma_exporter(),rdma_lane_enabled()(folds in the oldmm_default_transport_is_rdmaplus the legacySMG_MM_PIXEL_RDMAfallback),build_rdma_config_from_env(), and the slot-TTL derivation — all moved out of the deleted module.proto_wrapper.rstry_export_*now callmm_rdma_exporter().export(slot_key, ..).room→slot_keyis a Rust-internal rename only; the descriptor byte layout is untouched.SlotPool(free-list + lease/reclaim +[gen][payload][gen]framing) is split into its own module so the lease/TTL-race logic is unit-tested with no NIXL agent or hardware (6 tests). The TTL-derivation policy test moves to the gateway (transport.rs).Follow-ups (not in this PR): PR2 generalizes the payload path (
MmTensorPayload::Remote) and hoists the Python puller; PR3 wires vLLM end-to-end (emit + pull + bf16 + capability gate).Test Plan
Behavior-preserving refactor with no wire-format change, verified in both feature configurations:
cargo check -p smg(default) andcargo check -p smg --features mm-rdma— both compile.cargo test -p smg-mm-rdma --features nixl— 6SlotPooltests pass (gen framing, recycle detection, oversize rejection, TTL hold).cargo test -p smg --lib ...::multimodal::transport— gateway transport tests pass, incl.derived_rdma_slot_ttl_exceeds_worker_max_hold.cargo clippyclean:-p smgdefault and--features mm-rdma -- -D warnings;-p smg-mm-rdmastub and--features nixl -- -D warnings.cargo +nightly fmt --all --checkclean;codespellclean.Live NIXL loopback / vLLM validation on the GB300 box lands with PR3, when the transport is actually wired end-to-end.
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit