Skip to content

client-rs: multi-rail RDMA reads — one Worker, many NICs - #28

Open
123123213weqw wants to merge 6 commits into
DaoCloud:mainfrom
123123213weqw:multirail-read
Open

123123213weqw wants to merge 6 commits into
DaoCloud:mainfrom
123123213weqw:multirail-read

Conversation

@123123213weqw

Copy link
Copy Markdown

Summary

Multi-rail RDMA read path: a single client Worker reads one striped object in parallel over several local RDMA devices ("rails"), on top of the existing stripe-subset descriptor GETs — no changes to object or on-disk stripe layout, and the existing single-NIC configuration and public APIs keep working unchanged.

The server already exposes one RDMA listener per NIC (CS_RDMA_DEVICES); this closes the client-side gap: per-rail resources, stripe scheduling across rails, integrity, transport/memory safety, backpressure, topology-aware selection, and per-rail observability — plus a Python ctypes entry point so the KVConnector Worker can use it directly.

Design

  • Rail management: a rail = one local HCA (verbs ctx / PD / CQ / QPs / cached MRs) with health + cooldown recovery. Connections are per (rail, endpoint); one dedicated worker thread per connection preserves the one-request-in-flight control-channel invariant.
  • Planning (kv-service/client-rs/src/multirail.rs): placement validation (missing / duplicated / out-of-bounds stripes are rejected before any I/O), byte-exact weighted stripe allocation (cross-multiplied least-loaded, ties toward the higher-weight rail), wave dispatch bounded by connection caps, optional intra-rail splitting (RailLimits::task_max_stripes) for per-rail concurrency.
  • Memory safety (late-write protection): a read returns only when no work request can still place data into the caller's buffer — success replies imply server-side WRITE completions (RC placement); failure/timeout/cancel paths drain outstanding replies, then quiesce: Stop → join → ibv_destroy_qp → MR dereg. Stray retransmissions hit a destroyed QP / invalid rkey and are dropped by the transport. The safe API (&mut [u8]) additionally evicts the buffer's cached registrations synchronously; the _raw variant offers a sticky-registration fast path for pinned pools with the contract documented.
  • Integrity & version consistency: every task carries the full descriptor (handle / generation / ETag / layout version); a mismatch surfaces as a typed StaleDescriptor (re-lookup required). Per-task byte counts and server-reported num_chunks are checked against the layout, and per-stripe xxh3-64 checksums are verified client-side when the placement carries them (same encoding as the server).
  • Backpressure: RailLimits caps total/per-rail connections (queue depth) and total/per-rail in-flight bytes; dispatch waits with a deadline for headroom.
  • Topology: sysfs NUMA/PCIe per rail; selection policies LeastLoaded (default), EndpointAffinity (GID /24 match + endpoint whitelist pinning via device[:port[:gid[:weight[:mtu]]]][@ep[;ep...]]), WeightedRoundRobin.
  • Cancellation: CancelToken (shared flag) checked at wave boundaries, backpressure waits and reply collection; cancelled reads keep the same quiesce guarantees.
  • Observability: per-rail throughput, request/error/timeout counts, latency (avg/max), in-flight requests/bytes, registered-memory bytes, connections created/quiesced — via RailSnapshot (Rust), cs_mr_rail_stats (C ABI) and RailStats (Python).
  • Python Worker: cs_mr_* C ABI in rdma-ffi + contextstore.storage.multirail_client.MultiRailReader (ctypes), so the vLLM/Dynamo KVConnector can drive multi-rail reads directly.
  • Config knobs added (all default-preserving): RC path MTU (with_path_mtu / rail spec field / server CS_RDMA_PATH_MTU; both peers must agree), GRH hop limit (with_hop_limit / CS_RDMA_HOP_LIMIT) for routed fabrics, control-plane io/connect timeouts.

Correctness & performance notes

  • Soft-RoCE testbed (dual rxe, single node, two listeners): byte-exact balance across rails (e.g. 1280/1280 MiB), zero errors; aggregate scaling 1.81x with two client processes (rxe's kernel CPU path is the in-process ceiling — quantified in the notes below the patch series).
  • --task-max-stripes splitting raises in-process dual-rail throughput up to 1.29x on the same testbed.
  • The hardware-gated e2e suite covers: dual-rail read equality vs single-rail, failure injection with buffer reuse (proves no late RDMA WRITE corrupts a subsequent read), and stale-descriptor reporting.

Testing

  • 21 hardware-independent unit tests (planning, validation, weights, whitelists, waves, checksum, topology, cooldown, cancellation plumbing)
  • Hardware-gated e2e (--test multirail_e2e -- --ignored): 3/3 on a dual-rxe node
  • Python: Worker end-to-end via ctypes (128 MiB, byte-verified, per-rail stats) + upstream integration suite (pytest 5/5, full repo suite 116 passed)
  • Upstream regression: make build, make e2e, make lint clean; cargo fmt applied to touched crates
  • Two latent upstream issues fixed along the way: the --features rdma unit test did not compile (build_put_request arity), and newer clippy flags tonic-generated stubs (result_large_err allowed crate-wide with a comment)

Files

  • kv-service/client-rs/src/multirail.rs — core (rails, planning, safety, metrics) + unit tests
  • kv-service/client-rs/src/rdma.rs — additive: timeouts, device listing, GID query, detailed GET outcome, MR eviction, path MTU / hop limit
  • kv-service/rdma-ffi/src/multirail.rs — cs_mr_* C ABI
  • src/contextstore/storage/multirail_client.py — ctypes binding
  • kv-service/client-rs/src/bin/multirail_bench.rs, kv-service/client-rs/tests/multirail_e2e.rs — benchmark driver + hardware-gated tests

wzu added 5 commits September 26, 2026 08:56
Add MultiRailClient: N local RDMA devices (rails) read one striped
object in parallel over existing stripe-subset descriptor GETs.

- Rail management: per-rail verbs ctx/PD/CQ/QP/cached-MR, health +
  cooldown recovery, endpoint whitelists (rail-optimized pinning),
  sysfs topology (NUMA/PCIe) reporting.
- Planning: placement validation (missing/dup/out-of-bounds stripes
  before any I/O), byte-exact weighted stripe allocation (cross-
  multiplied least-loaded), wave dispatch under connection caps.
- Safety: control-plane io/connect timeouts; failed reads quiesce
  connections (Stop -> join -> QP destroy -> MR dereg) before
  returning, so late server WRITEs cannot touch freed/reused memory;
  safe API evicts cached registrations synchronously.
- Integrity: per-task byte/chunk-count checks against the layout,
  StaleDescriptor typed error on generation/etag mismatch, optional
  per-stripe xxh3-64 verification matching the server encoding.
- Backpressure: per-rail/total connection and in-flight-byte limits.
- Observability: RailSnapshot per-rail stats.
- Tools/tests: cs-multirail-bench, 21 hardware-independent unit
  tests, hardware-gated e2e (dual-rail correctness, failure
  injection + buffer-reuse safety, stale descriptor).
- rdma.rs: additive io/connect timeouts, device listing, GID query,
  num_chunks GET outcome, synchronous MR eviction; fix upstream
  rdma-bench test compile + clippy nits; allow result_large_err on
  tonic-generated stubs (newer clippy).
…hardening

- RC path MTU configurable end to end (RdmaClientConfig with_path_mtu,
  RailConfig 5th spec field, server CS_RDMA_PATH_MTU); default stays
  1024. Both peers must agree - the HELLO exchange does not carry it.
- RailLimits task_max_stripes splits per-rail stripes over several
  connections (intra-rail concurrency / queue depth); exercised by the
  e2e suite via CS_MR_TASK_MAX_STRIPES.
- Bench: --sticky (cached registrations), --task-max-stripes.
- Testbed moved to a conflict-free 192.168.250.0/24 veth pair with an
  up/down setup script; full netns isolation is blocked by an rxe
  in-netns transport bug (vanilla upstream PUT flushes with
  WR_FLUSH_ERR), documented in the notes.
…oint pin syntax

- rdma-ffi exposes MultiRailClient over a C ABI: cs_mr_new/read/
  rail_stats/free with CsMrDescriptor/CsMrChunk mirroring the gRPC
  fields; Python ctypes binding MultiRailReader in
  contextstore.storage.multirail_client closes the KVConnector Worker
  loop (verified end to end: 128MiB over 2 rails, byte-exact balance).
- CancelToken: shared-flag cooperative cancellation checked at wave
  boundaries, backpressure waits and reply collection; cancelled reads
  still drain in-flight replies and quiesce connections before
  returning (memory-safety contract unchanged).
- RailConfig::parse gains an @host[:port][;...] endpoint whitelist
  suffix for rail-optimized pinning from a plain spec string.
- Python integration suite (pytest tests/) passes 5/5.
RailStats/RailSnapshot/FFI/Python now expose per-request latency
(avg/max over successful reads) and the bytes currently pinned by
cached registrations (accounted at register/evict time; the safe API
returns to zero after each read, sticky pooling holds the buffer
size). Closes the observability items of the requirements list.
Self-review findings fixed before opening the upstream PR:
- aborted reads now drain outstanding task replies (bounded by the read
  deadline) before quiescing, releasing inflight_requests/inflight_bytes
  — previously an aborted wave could leak counters and eventually starve
  the backpressure wait;
- a connection worker subtracts its still-cached registration bytes from
  the rail registered_bytes counter on teardown, keeping the metric
  honest after failures;
- cargo fmt applied to the touched crates (reverted on untouched
  upstream files to keep the diff noise-free).
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant