Skip to content

comms/uniflow: instrument the receive path across the controller and the transport (#3905) - #3905

Closed
cppccppccppc wants to merge 6 commits into
meta-pytorch:mainfrom
cppccppccppc:export-D117632428
Closed

comms/uniflow: instrument the receive path across the controller and the transport (#3905)#3905
cppccppccppc wants to merge 6 commits into
meta-pytorch:mainfrom
cppccppccppc:export-D117632428

Conversation

@cppccppccppc

@cppccppccppc cppccppccppc commented Sep 1, 2026

Copy link
Copy Markdown

Summary:

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

--- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

--- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Reviewed By: yexiangd

Differential Revision: D117632428

@meta-cla meta-cla Bot added the CLA Signed This label is managed by the Meta Open Source bot. label Sep 1, 2026
@meta-codesync

meta-codesync Bot commented Sep 1, 2026

Copy link
Copy Markdown
Contributor

@cppccppccppc has exported this pull request. If you are a Meta employee, you can view the originating Diff in D117632428.

cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 1, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
@meta-codesync meta-codesync Bot changed the title comms/uniflow: instrument the receive path across the controller and the transport comms/uniflow: instrument the receive path across the controller and the transport (#3905) Sep 1, 2026
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 1, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 1, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 2, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 2, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 2, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 2, 2026
…the transport (meta-pytorch#3905)

Summary:
Pull Request resolved: meta-pytorch#3905

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

 --- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

 --- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
@cppccppccppc
cppccppccppc force-pushed the export-D117632428 branch 2 times, most recently from 1d95a78 to d28fc90 Compare September 2, 2026 23:47
cppccppccppc pushed a commit to cppccppccppc/torchcomms that referenced this pull request Sep 2, 2026
…the transport (meta-pytorch#3905)

Summary:

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

--- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

--- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Differential Revision: D117632428
Peng Chen added 6 commits September 2, 2026 17:58
…rrected numbers (meta-pytorch#3897)

Summary:

--tcp-sockbuf exposes the data connection's SO_SNDBUF/SO_RCVBUF so the value can be
swept instead of argued about, and the bench brackets each size's timed loop with
TcpTransport::logAndResetPhaseStats() so every bandwidth line is accompanied by where
the time went.

The script header carried get ~1.0 GB/s, which predated the pinned-staging work and
was stale by roughly 9x. It misled a full round of planning, so it is corrected here
to measured values along with the link baseline (200G, MTU 1500, RTT 0.046 ms, iperf3
15.4 GB/s single stream and 23.3 GB/s over 8) that makes the numbers interpretable.

Recipe 8 records a disproof rather than a hypothesis. The theory was that the 1 MiB
buffer pin caps a stream at window/RTT; at 0.046 ms RTT, 1 MiB permits ~22.8 GB/s, so
the window was never the constraint. A 6-point sweep with 3 repeats does show the
1 MiB pin is the worst non-degenerate setting (8.63 GB/s at 1 GiB against 9.67
unpinned, disjoint ranges), but the mechanism is autotuning being disabled while the
reader does a multi-MiB copy, not the bandwidth-delay product. The 64K arm is a
deliberate control: it drops throughput 75%, which is what makes the flat region
above 1 MiB trustworthy rather than merely consistent with a dead knob.

Also records that a --no-verify arm must never be compared against a verifying one --
that mistake inverted this exact comparison once.

Reviewed By: yexiangd

Differential Revision: D117132684
…ytorch#3913)

Summary:

Phase 3d makes slab-backed VRAM ReadReply copies asynchronous on the caller stream. The transport retains the receive slab and destination write reservation until CUDA event completion, polls all pending events from the EventBase so later copies can retire out of order, and drains or quarantines copies safely across query errors and shutdown.

Vector-backed fallback remains synchronous because its source storage is reused immediately. The benchmark exposes --no-tcp-async-h2d for controlled comparisons and reports the active mode.

# A failed completion probe is not a failed copy

Two paths conflated "I could not observe this copy finish" with "this copy did not finish", and reported a transfer error for data that had demonstrably landed.

In `stageAsyncH2d`, a failing `eventRecord` arrives after `memcpyAsync` has already succeeded, so the copy is in flight and only the tracking mechanism is gone. The code falls back to `waitForH2dCopy`, and if that wait succeeds the copy has completed and the destination holds the payload -- but the old code then returned the `eventRecord` error, failing the caller's get for a copy that worked. It now completes with `Ok()`.

`pollPendingH2d` had the same shape: a failing `eventQuery` left its error in `result`, and a successful `waitForH2dCopy` fell through without clearing it, retiring a completed copy as failed. The recovery now sets `result = Ok()` before retirement.

The sibling site at the end of `stageAsyncH2d` deliberately keeps propagating `status`: there the error is the transport stopping or the pending record failing to track, which is a real operation failure even though the synchronize proves the copy finished. Only the two probe-failure paths change.

Reviewed By: yexiangd

Differential Revision: D117155852
…eta-pytorch#3918)

Summary:

TcpAsyncAcceptTest derived the address family by comparing the param's
clientHost against the "127.0.0.1" literal. That defaults every other
spelling -- "localhost", "[::1]", a resolvable hostname -- to AF_INET6
without saying so, and the two socket-buffer tests feed that family to
kernelDefaultRcvBuf/rcvBufForRequest. A param added later would probe the
wrong family and surface as a confusing skip or a wrong expectation rather
than a clear failure.

Store the family on AddrFamily and state it at the two
INSTANTIATE_TEST_SUITE_P entries, so a new param has to declare which
family it is. A helper function would have deduplicated the comparison but
kept the default. The five sites this replaces are the two family
derivations in the socket-buffer tests, the socket()/sockaddr pair in
AsyncAcceptRejectsNonUniflowClient, and the test-name lambda -- which now
derives the name from the same field the bodies use, so the two cannot
disagree.

AddrFamily is defined separately in four test files; this changes only
TcpAsyncAcceptTest.cpp. The same pattern remains in TcpConnTest.cpp:48
(inverse polarity) and the other three name lambdas.

Reviewed By: yexiangd

Differential Revision: D117407765
… sockets (meta-pytorch#3904)

Summary:

D117132557 threaded TcpSocketConfig into the accept path but
configureAcceptedSocket consumed only socketBufSize and hardcoded the other
seven fields. That is a worse shape than the old `int acceptRetryCnt`

Reviewed By: yexiangd

Differential Revision: D117409938
Summary:

sendAllVec took (iovec*, int) and tracked its position with a separate idx
cursor, so the loop carried `iov + idx`, `iovCnt - idx` and `iov[idx]`
arithmetic. std::span carries the count with the pointer, which lets the
cursor go away entirely: the retire step becomes iov = iov.subspan(1) and the
loop condition becomes !iov.empty().

That is the real win rather than the shorter signature. The arithmetic being
deleted lives in the short-write branch, which the comment there notes is not
reachable on a blocking socket except via a signal mid-transfer and is not
covered by tests -- so it is the least safe place in the function to keep
hand-rolled index math.

It also drops the static_cast<size_t> on msg_iovlen, since size() is already
size_t, and matches the idiom the rest of the interface uses:
std::span<const uint8_t> on send, std::span<uint8_t> on recv. <span> was
already included and sendAllVec is private with one call site, so there is no
ABI consideration.

The call site's count becomes size_t and passes std::span{iov}.first(iovCnt),
so the cast is removed rather than relocated to the caller.

Kept the doc comment. A non-const std::span<iovec> conveys no more about
mutation than a non-const iovec* did; what the comment carries is that the
mutation is destructive bookkeeping -- the iov must not be reused after the
call -- and no signature expresses that.

Reviewed By: yexiangd

Differential Revision: D117414473
…the transport (meta-pytorch#3905)

Summary:

Folded from two adjacent diffs; each half is stated separately below so the two
arguments stay reviewable on their own terms.

--- controller: recv header-wait and payload-drain timings (was D117632427) ---

TcpConn::syncRecv() populates the existing RecvPhaseStats with the time spent
waiting for the length prefix versus draining the payload, plus frame and byte
counts. Splitting the two phases is what makes a slow receive interpretable: a
large headerWaitNs means we were waiting on the peer, while a large
payloadDrainNs means the socket itself was the limit.

All four counters are relaxed fetch_add on atomics already declared in
Controller.h, so this adds two steady_clock reads per frame and no
synchronisation.

--- transport: receive-slab hits, misses and vector receives (was D117632428) ---

Adds receiveSlabAttempts_, receiveSlabMisses_, and vectorReceiveCount_ to
TcpTransport and reports them on the existing tcp phases log line. A miss means
the reader could not get a pinned slab and fell back to a vector-backed
receive, which changes the H2D path for that frame, so distinguishing the two
is necessary before drawing conclusions from an aggregate drain number.

The counters are relaxed atomics incremented in readerLoop().

Reviewed By: yexiangd

Differential Revision: D117632428
@meta-codesync

meta-codesync Bot commented Sep 3, 2026

Copy link
Copy Markdown
Contributor

This pull request has been merged in e8c03e5.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

CLA Signed This label is managed by the Meta Open Source bot. Merged meta-exported

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant