Motivation
In KVService, large objects are striped across multiple NVMe devices / nodes, so
aggregate disk bandwidth scales with stripe count. A single client worker,
however, reads through ONE RDMA path (one NIC / QP). Once
stripe_count × disk_bw > single-NIC bandwidth, the network becomes the
bottleneck while the disks still have headroom — adding disks or concurrency on
the storage side no longer shortens read time.
Multi-rail lets one client worker transfer disjoint stripe subsets of the
same object over several independent RDMA paths concurrently, aggregating NIC
bandwidth.
What the codebase already gives us (why this is a generalization, not a rewrite)
RdmaClient.get_descriptor_stripes_into / get_descriptor_stripes_sge
(wire tag 15, server RDMA-WRITEs straight into the caller's scattered
destination buffers) already supports per-node concurrency: "callers issue
one call per owning node … disjoint target regions let all nodes transfer
concurrently".
- Read completion is signalled over the TCP control channel (
GET_RESP), so
aggregating N rails = spawn one task per rail and join — no server-side
polling changes needed.
- Consistency fields are already on the wire:
object_generation in
ObjectDescriptor, chunk_checksums in PlacementChunk / StripingInfo.
So multi-rail = generalize "concurrent per node" to "concurrent per rail
(independent NIC/QP)", plus three missing pieces: path lifecycle management,
cross-rail completion aggregation, and memory-consistency guards.
Proposed design
Worker (vLLM / Dynamo / NIXL)
│ upper interface unchanged: LookupObject / ReadByDescriptor
▼
RailManager ── N independent rails (NIC + QP + CQ); single NIC = 1 rail
▼
Stripe→Rail scheduler ── locality affinity + round-robin (least-loaded later)
▼
per rail, concurrently: get_descriptor_stripes_sge(descriptor, stripe subset, SGE list)
│ each rail writes disjoint regions of the same destination buffer
▼
Completion aggregation + integrity layer
(object_generation check · per-stripe checksum · epoch buffer guard)
▼
Bottleneck attribution / observability (per-rail timing, classify
NIC / PCIe / memory / disk / software)
Invariants (hard constraints):
- On-disk stripe layout (
StripingInfo) is never touched.
- Upper read interface (
LookupObject / ReadByDescriptor) is never touched.
- Single-NIC deployments keep working unchanged: with one rail the scheduler
is a pass-through.
- v1 failure semantics: if any required rail fails, the whole request fails
safely — no partial data returned, no transparent in-request retry (later
versions may add retry).
Memory safety (late RDMA WRITE hazard): if a cancelled request's buffer is
freed and reused, a late-arriving RDMA WRITE can corrupt the new owner's data.
Three defences:
- pin-until-completion: the read holds the buffer and every rail's
registration until join; cancel only flips a logical flag, MRs are not
deregistered early.
- epoch buffer guard: late writes only land while
live == epoch; bumping
live on cancel makes the guard skip them.
- two-phase commit + per-stripe checksum: stripes are verified against
chunk_checksums after aggregation; a cancelled request never reaches
commit.
What is already working (Draft PR: #34)
-
multi_rail.rs — RailReader trait, RailManager, MultiRailReader
(stripe planning → scoped-thread concurrent reads → aggregation → checksum
verify), RailReadStats with a bottleneck classifier.
-
mock_rail.rs — Mock transport with configurable per-rail bandwidth and
fault injection (late_write_after_cancel + epoch guard), so the full
scheduling / state-machine / buffer-lifecycle path is testable on any
machine (no RDMA hardware).
-
rdma.rs — RailReader implemented for RdmaClient (real verbs path,
works for Soft-RoCE and physical NICs).
-
Tests (11 cases, default features, hardware-free): multi-rail aggregation
with byte-exact reconstruction, single-rail pass-through, A/B demonstration
of the late-write hazard (without cancel: corruption reproduced; with
cancel: guard blocks it), stripe-planning round-robin, per-stripe checksum
failure detection, single-rail disconnect safe-fail, partial-completion /
under-delivery detection, read timeout with the caller's buffer untouched,
inflight-budget backpressure rejection, and object version-change
(generation) rejection.
-
cs-mock-bench — capacity model sweeping rail count (per-rail bandwidth
configurable), measured on WSL2: 64 MiB object, 4 MiB chunks, 16 stripes,
125 MB/s per rail:
rails agg_MB/s bal_MB/s theo_MB/s speedup bottleneck
1 123.5 125.0 125.0 0.99x network: slowest rail caps the object
2 244.9 250.0 250.0 1.96x software: post-read stripe verification (client CPU)
3 327.4 333.3 375.0 2.62x software: post-read stripe verification
4 484.1 500.0 500.0 3.87x software: post-read stripe verification
Near-linear aggregate scaling (3.87x at 4 rails). bal accounts for stripe
quantization (the slowest rail carries ceil(stripes/rails) stripes — hence
the 3-rail dip). Notably, once the NIC bottleneck is removed the classifier
already flags the next ceiling — client-side checksum CPU — which is
exactly the per-layer overhead attribution this work instruments, and
motivates aligning on the server's xxh3 checksums (question 3 below).
Environment honesty
All functional validation so far uses the Mock transport (the 11
hardware-free tests above). A Soft-RoCE (RXE) dual-path setup on WSL2
(two rxe devices, server listening via CS_RDMA_DEVICES=rxe0:...,rxe1:...)
is being wired up for end-to-end verbs validation — the
examples/softroce_dual_rail.rs harness is included in the PR. Even once
complete, Mock/Soft-RoCE prove functionality, failure semantics, scheduling
and resource governance — they do NOT prove hardware aggregate bandwidth,
HCA offload, PCIe/NUMA effects or zero-copy behaviour on physical NICs.
All numbers are labelled accordingly. No GPU is involved.
Questions for maintainers
- Where should rail configuration live (client config vs. endpoint list
syntax)? CS_RDMA_DEVICES already accepts multiple server-side NICs —
should the client mirror that syntax?
- The server can already listen on multiple RDMA devices
(CS_RDMA_DEVICES=dev0:...,dev1:...). Any objection to using that as the
multi-rail server shape, with per-stripe-subset GETs pinned to the rail
that reaches the owning endpoint?
- Client-side stripe checksum is currently a lightweight local hash; the
server's chunk_checksums are xxh3-based. Align the client on xxh3-64 so
multi-rail verification checks the server's checksums?
- Preferred surface for per-rail metrics (throughput / latency / errors /
in-flight / registered memory)?
Non-goals (v1)
- No change to disk stripe layout or placement policy.
- No transparent retry within a single request.
- No GDS / GPU path changes.
- No claim of physical multi-NIC aggregate bandwidth from Mock/Soft-RoCE data.
Happy to adjust scope or split this into smaller PRs — feedback welcome.
Motivation
In KVService, large objects are striped across multiple NVMe devices / nodes, so
aggregate disk bandwidth scales with stripe count. A single client worker,
however, reads through ONE RDMA path (one NIC / QP). Once
stripe_count × disk_bw > single-NIC bandwidth, the network becomes thebottleneck while the disks still have headroom — adding disks or concurrency on
the storage side no longer shortens read time.
Multi-rail lets one client worker transfer disjoint stripe subsets of the
same object over several independent RDMA paths concurrently, aggregating NIC
bandwidth.
What the codebase already gives us (why this is a generalization, not a rewrite)
RdmaClient.get_descriptor_stripes_into/get_descriptor_stripes_sge(wire tag 15, server RDMA-WRITEs straight into the caller's scattered
destination buffers) already supports per-node concurrency: "callers issue
one call per owning node … disjoint target regions let all nodes transfer
concurrently".
GET_RESP), soaggregating N rails = spawn one task per rail and join — no server-side
polling changes needed.
object_generationinObjectDescriptor,chunk_checksumsinPlacementChunk/StripingInfo.So multi-rail = generalize "concurrent per node" to "concurrent per rail
(independent NIC/QP)", plus three missing pieces: path lifecycle management,
cross-rail completion aggregation, and memory-consistency guards.
Proposed design
Invariants (hard constraints):
StripingInfo) is never touched.LookupObject/ReadByDescriptor) is never touched.is a pass-through.
safely — no partial data returned, no transparent in-request retry (later
versions may add retry).
Memory safety (late RDMA WRITE hazard): if a cancelled request's buffer is
freed and reused, a late-arriving RDMA WRITE can corrupt the new owner's data.
Three defences:
registration until join; cancel only flips a logical flag, MRs are not
deregistered early.
live == epoch; bumpingliveon cancel makes the guard skip them.chunk_checksumsafter aggregation; a cancelled request never reachescommit.
What is already working (Draft PR: #34)
multi_rail.rs—RailReadertrait,RailManager,MultiRailReader(stripe planning → scoped-thread concurrent reads → aggregation → checksum
verify),
RailReadStatswith a bottleneck classifier.mock_rail.rs— Mock transport with configurable per-rail bandwidth andfault injection (
late_write_after_cancel+ epoch guard), so the fullscheduling / state-machine / buffer-lifecycle path is testable on any
machine (no RDMA hardware).
rdma.rs—RailReaderimplemented forRdmaClient(real verbs path,works for Soft-RoCE and physical NICs).
Tests (11 cases, default features, hardware-free): multi-rail aggregation
with byte-exact reconstruction, single-rail pass-through, A/B demonstration
of the late-write hazard (without cancel: corruption reproduced; with
cancel: guard blocks it), stripe-planning round-robin, per-stripe checksum
failure detection, single-rail disconnect safe-fail, partial-completion /
under-delivery detection, read timeout with the caller's buffer untouched,
inflight-budget backpressure rejection, and object version-change
(generation) rejection.
cs-mock-bench— capacity model sweeping rail count (per-rail bandwidthconfigurable), measured on WSL2: 64 MiB object, 4 MiB chunks, 16 stripes,
125 MB/s per rail:
Near-linear aggregate scaling (3.87x at 4 rails).
balaccounts for stripequantization (the slowest rail carries ceil(stripes/rails) stripes — hence
the 3-rail dip). Notably, once the NIC bottleneck is removed the classifier
already flags the next ceiling — client-side checksum CPU — which is
exactly the per-layer overhead attribution this work instruments, and
motivates aligning on the server's xxh3 checksums (question 3 below).
Environment honesty
All functional validation so far uses the Mock transport (the 11
hardware-free tests above). A Soft-RoCE (RXE) dual-path setup on WSL2
(two rxe devices, server listening via
CS_RDMA_DEVICES=rxe0:...,rxe1:...)is being wired up for end-to-end verbs validation — the
examples/softroce_dual_rail.rsharness is included in the PR. Even oncecomplete, Mock/Soft-RoCE prove functionality, failure semantics, scheduling
and resource governance — they do NOT prove hardware aggregate bandwidth,
HCA offload, PCIe/NUMA effects or zero-copy behaviour on physical NICs.
All numbers are labelled accordingly. No GPU is involved.
Questions for maintainers
syntax)?
CS_RDMA_DEVICESalready accepts multiple server-side NICs —should the client mirror that syntax?
(
CS_RDMA_DEVICES=dev0:...,dev1:...). Any objection to using that as themulti-rail server shape, with per-stripe-subset GETs pinned to the rail
that reaches the owning endpoint?
server's
chunk_checksumsare xxh3-based. Align the client on xxh3-64 somulti-rail verification checks the server's checksums?
in-flight / registered memory)?
Non-goals (v1)
Happy to adjust scope or split this into smaller PRs — feedback welcome.