Completion Semantics of NCCL Failure Recovery

After the sender sees completion, the receiver may still not have consumed the message. Should failure recovery trust the sender's completion, or the application effect the receiver has already produced?

This note is about one concrete question: when an NCCL transfer fails halfway, how do we decide whether a message must be sent again?

Suppose the sender has emitted a piece of KV cache or an expert output, and the CUDA event for that send has already completed. The receiver process is then killed. Recovery has two choices, resend or do not resend. Not resending is correct only if sender completion already proves that the receiver finished the work. Resending is correct only if the receiver has not yet used the message. Both premises can fail at once.

The conjecture we started from was simple. The completion NCCL exposes to the sender, together with the communicator’s error state, might already be enough to define the retry boundary. If that were true, the layer above would not need per-message state. Recovery would resend only the operations whose completion the sender had not yet observed.

What has to be tested is therefore not whether NCCL can report an error. It is whether this equivalence holds:

sender completion ≈ the receiver has already produced the application effect

The rest of the experiment is organized around that relation.

Why this starts to matter

When large-scale training hit a communication fault, the usual response was to fail the whole group and restart from a checkpoint. That is expensive, and the semantics are simple: once the step is discarded, nobody has to know how far any in-flight message got.

Current systems are less willing to pay that cost. torchft can drop a failed replica at a step boundary. NCCL’s ncclCommShrink and ncclCommRevoke let a communicator be rebuilt after a fault. Inference also keeps long-lived state. Expert-parallel token dispatch, KV-cache transfer between prefill and decode, and long RL rollouts all want to keep work that already finished, instead of wiping the step and starting over (Qin et al., FAST 2025).

flowchart LR
A["Then: rank fails"] --> B["Drop the step"] --> C["Restart from checkpoint"]
D["Now: rank fails"] --> E["Revoke / rebuild"] --> F["Keep finished work"]
F --> G["Which messages took effect?"]

That picture is where the problem comes from. Message-level completion does not matter while the whole step is thrown away. The moment the system tries to keep partial progress, message-level semantics become the basis of correct recovery. A KV block that a decode worker has already attached to later token generation can create duplicate state if it is replayed. The opposite miss is just as bad. Data that only reached an intermediate buffer, and was never consumed, is dropped for good if sender completion is treated as a reason not to resend.

Completion, delivery, and consumption are not the same event

Three events have to stay separate. The first is local completion. The sender has handed the send to the NCCL and CUDA pipeline, and the send buffer may now be reused. ncclSend itself is asynchronous. It only enqueues work on a CUDA stream. The caller usually records a CUDA event afterward, so the event means “the send-side work has finished,” not “the peer application has processed the bytes.” The second is delivery. The bytes have landed in receiver memory. The third is consumption. A receiver kernel or the application logic has read those bytes and turned them into an effect that cannot be ignored, such as attaching a KV-cache block to later decode, or adding an expert output into the final result.

A parcel is a useful picture. Local completion is the drop box accepting the package, so the sender can leave. Delivery is the package arriving at the address. Consumption is the recipient opening it and using what was inside. In a healthy run the three stay in order, and it is easy to call all of them “done.” Recovery is hard because a fault can cut exactly between them.

flowchart LR
S["1. ncclSend"] --> C["2. Sender completion"] --> D["3. Delivery"] --> U["4. Consumption"]

The object of this note is the layer at which failure recovery places its commit point. A commit at sender completion can believe too early that a message is done. A graph-level or batch-level completion can confirm too late, and replay a message the receiver already consumed.

On an RDMA reliable connection, a completed write means the remote NIC has acknowledged that the data was written into remote memory. That is delivery. It does not include application consumption (InfiniBand specification). After a queue pair enters the error state, outstanding requests return as flush errors, but some of those writes may already have happened on the peer. A failed request does not imply that nothing happened remotely.

This experiment has no GPUDirect RDMA, so the path is longer. The NIC cannot write remote data straight into GPU HBM. The bytes land first in a CPU-visible host bounce buffer, NCCL’s proxy and copy path moves them into the GPU, and only then do the receive kernel and the later consumer kernel use them. That is an extra handoff between “the parcel reached the building” and “someone opened it in the room.” Each extra stage is another window a fault can enter. The gap is a property of this data path, not an abstract argument about words.

The recovery primitives we already have are group-scoped. ncclCommAbort drops every in-flight operation. ncclCommRevoke stops in-flight operations and marks the communicator in progress, after which it can be destroyed or shrunk once it is quiescent. ncclCommGetAsyncError returns one error code for the whole communicator (NCCL communicator API). PyTorch’s ProcessGroupNCCL adds a watchdog that fails every outstanding Work on the group. MPI’s fault-tolerance extension ULFM provides revoke, shrink, and agree. Agree lets surviving processes reach a fault-tolerant consensus on an integer (Bland et al., 2013), but a failed request itself only reports process failure or revocation. UCX’s peer error-handling mode guarantees that every send request ends in success or error, and an error ending does not say whether the peer received the data (UCX).

The window is large in absolute time. In this setup, abort or revoke itself takes about 0.5 s, mostly from the NCCL proxy thread’s 500 ms poll. A 16 B send/recv round trip is only tens of microseconds. The control interval from “one side has decided to give up” to “the other side is actually quiescent” is about four orders of magnitude longer than one message lifetime. The send window allows at most 32 operations in flight, and those 32 can keep crossing the boundaries of sent, delivered, and consumed during that half second. The network did not suddenly become slow. The recovery control plane is coarse and slow, while the data plane is still moving forward.

TL;DR

A communication library telling the sender “I am done” does not mean the receiver application has used the message. On two V100s, EDR InfiniBand, and NCCL 2.32.3 without GPUDirect RDMA, we injected 634 faults. When resend decisions used only sender completion, eager and revoke recovery missed about 1–4 messages per fault that were “send complete” but not yet consumed. A 16-op CUDA Graph, because its completion is that coarse, duplicated about 4–9. After recovery followed the receiver’s actual consumption sequence, duplicates and losses both fell to 0, and the extra cost stayed inside measurement noise. The stable commit point is the application effect, not transport or runtime completion.

What we actually conjectured

The simplest conjecture is that sender completion can stand in for a message’s recovery commit point. In-flight messages then fall into two classes. Those whose completion the sender has seen are not resent. Those still incomplete are sent again after recovery.

The conjecture looks reasonable because, on a healthy run, completion and consumption sit next to each other. NCCL also hides a lot of RDMA, proxy, buffer, and stream detail from the application. A CUDA event is the most natural signal the upper layer can get that “this communication is over.”

It has two openings. First, completion can be too early. The sender believes the send is finished, but the receiver has only received the bytes and has not run the consumer kernel. Second, completion can be too coarse. A CUDA Graph folds many sends into one graph-level event. The receiver may already have consumed part of the batch, while the sender can still only see that the whole graph is not done.

We therefore also ran a stronger control. Every message carries a sequence number. The receiver writes that number into a log only after it has actually consumed the message. On recovery the sender does not consult its own completion. It asks the receiver which sequence numbers have produced an application effect, and it resends only the missing ones. That is the end-to-end argument: the communication layer can report how much was transferred, and only the endpoint knows whether the operation happened in the application’s semantics (Saltzer, Reed, and Clark, 1984).

The experiment compares two commit points:

flowchart LR
A["Fault"] --> B{"Trust?"}
B --> C["Sender completion"]
B --> D["Receiver log"]
C --> C1["done: no resend"]
C --> C2["open: resend"]
D --> D1["consumed: no resend"]
D --> D2["missing: resend"]
C1 --> E["may under-send"]
C2 --> F["batch may over-send"]
D1 --> G["application effect"]
D2 --> G

If sender completion and receiver consumption agree, the two paths should produce the same result. If they systematically diverge inside the fault window, the error shows up directly.

How the experiment separates the three events

Message lifetime from sender fill and ncclSend, through host bounce, to receiver consumption and the durable log. The shaded window is completed at the sender and not yet consumed.
Figure 1. When the sender sees completion, the message may still be in the host bounce buffer, inside the receive kernel, or before the log. On the graph path, completion is one event for a batch of 16.

Figure 1 is one message lifetime. A kernel on the sender GPU fills the message and ncclSend emits it. The CUDA event after that send is the completion the sender sees. The data is written through an InfiniBand reliable connection into the receiver’s host buffer, then copied to the GPU. After ncclRecv completes, the receiver’s consumer kernel checks and applies the message, then appends the sequence number to a durable log. We define “consumed” by that log. The shaded region is the window where the sender has reported completion and the receiver has not consumed. A fault there, combined with “resend only what is not complete,” loses the message. The note under the figure is the graph path: 16 sends are captured in one CUDA graph, and the sender sees a single completion event for the whole graph.

The hardware is two nodes, one V100-PCIe-32GB each, on EDR InfiniBand, without GPUDirect RDMA, so the data path uses the host bounce buffer. The software is NCCL 2.32.3 and CUDA 12.4. The first 16 B of each message are a 64-bit sequence number and a checksum. Each payload word is determined by the sequence number, and the receiver checks every message. Message size alternates between 16 B and 64 KB across trials. The send window is 32. On the graph path each graph holds 16 messages, with two graphs in flight. A separate TCP control connection carries the communicator unique id, the fault plan, the abort notice, and the log query after the fact.

There are three recovery paths. Eager is per-message send/recv, then ncclCommAbort and a new communicator. The graph path also uses abort. Revoke uses a nonblocking communicator, then ncclCommRevoke, destroy, and a new communicator. There are three faults. F1 sends SIGKILL to the receiver in the middle of the stream, and a supervisor restarts it and brings it back. F2 is the sender aborting itself mid-stream. F3 is SIGSTOP on the receiver for 100 to 800 ms, then SIGCONT. The pause includes NCCL’s proxy thread. Injection time is uniform over the length of a fault-free stream. Detection uses only signals that already exist: control-socket disconnect, ncclCommGetAsyncError, a 500 ms stall, and F2’s self-abort.

We compare two recovery rules. The first uses only the completion the sender can see. It resends every message whose completion the sender has not observed, and the receiver does not deduplicate. The second uses endpoint consumption. Messages carry sequence numbers. The receiver keeps a dedup bitmap and writes the log only after real consumption. After a restart it restores that log. Once communication is back, the sender queries the log and resends only missing sequence numbers. A duplicate is the same sequence number consumed more than once. A loss is a sequence number still absent from the log after recovery. Stale is the count, at the fault, of messages whose sender completion and receiver consumption disagree. Each cell is 210 injections. In F3, a pause shorter than 500 ms trips no detector and only adds latency, so the trigger count is below 210. Fault-free controls, 40 runs each, have all four counters at 0, which rules out harness false positives. A fourth fault, link jitter, would require changing the NIC port state and was not allowed on the shared nodes, so it was not measured.

Where sender completion and receiver consumption diverge

Bar chart of duplicate consumes and losses per tripped fault for eager, graph, revoke, and sequence-log recovery under receiver kill, sender abort, and receiver pause.
Figure 2. Under fixed retry, eager and revoke lose messages the sender already completed, while the graph path consumes duplicates. The end-to-end protocol is zero on all three faults.

Figure 2 has one panel per fault. The horizontal axis is the recovery path and policy. The vertical axis is the average number of duplicate messages (solid) or lost messages (hatched) per fault that actually triggered recovery. The absolute height matters less than the direction of the error. Eager and revoke lean toward loss. The graph path leans toward duplication. Those two directions match two completion granularities. The first decides too early that a single send is done. The second waits until a whole batch is done.

On the eager path with fixed retry, killing the receiver (F1) loses 4.0 messages per trigger (675 messages across 169 triggers). A pause (F3) loses about 1.0. A sender abort (F2) loses 28 messages and duplicates 1 across 171 triggers. Against a window of at most 32 in-flight operations, 4 is not the whole window being wrong. It is the last few messages caught in the gap where the sender has already released the operation and the receiver has not yet consumed it. That magnitude matches the data path. Most messages are either not yet complete or already consumed. Only a short tail of the pipeline sits in the ambiguous region.

The lost messages share one property. The sender has already seen completion, so recovery will not resend them. The receiver may only have the bytes in a host or GPU buffer, with the consumer kernel not yet run, when the process is killed. Or abort invalidates data that was delivered and not yet consumed. The eager error is not that retry is too aggressive. It is that send completion is too early to use as the acknowledgement point. Revoke has the same shape. F1 loses 2.1 messages per trigger. F2 and F3 accumulate 27 and 91 losses, with no duplicates. Revoke shortens the kill path from eager’s roughly 880 ms to 385–520 ms. That changes how long it takes to stop. It does not change which event counts as the message having happened.

The graph path runs the other way, and the dominant error is duplication. Under F2, 210 triggers produce 1901 duplicate consumptions, about 9 per fault. F1 and F3 are about 3.9 and 6.2. The direct picture is that a CUDA Graph gives the sender one receipt for 16 messages. At the fault, the receiver may already have consumed the first 9, while the graph completion event is still dark. Recovery can only treat all 16 as unfinished and send them again, so the first 9 are applied twice.

Nine is not noise. It is the natural size of a fault in the middle of a 16-op batch. The closer the fault sits to the end of the batch, the more messages have been consumed without a batch completion. Once the graph coarsens completion, the error flips from sending too little to sending too much. If consumption is not idempotent, for example accumulating a gradient, merging an expert output, or advancing a state machine, that duplicate is not just wasted work. It can become a silent numerical error.

The end-to-end column is 0 on all three faults. On the eager path, F1, F2, and F3 trigger 178, 177, and 68 times (F2 has 4 extra injections, 214 in total). Duplicates and losses are both 0. The protocol does not use a more elaborate communication primitive. The receiver records which sequence numbers it actually consumed, and recovery fills the gaps in that ledger. Sender completion is a record of send-side progress. The receiver log is a record of application effect. Moving the recovery basis from the first to the second removes both of the opposite errors above.

The cost is small. On a fault-free stream, turning dedup on changes 16 B throughput by +1.7% and 64 KB throughput by -0.8%, both inside run-to-run noise. Sequence comparisons and bitmap lookups are light per-message metadata next to the data movement and GPU kernel launch. On the graph path we only ran fixed retry online. Offline, adding receiver dedup removes the duplicates. The remaining gaps, 121 on F1 and 6 on F3, are exactly sequence numbers absent from the receiver log, which is the set a log-based resend would identify. That online graph recovery path was not measured separately.

Three failure modes the runs exposed directly

ncclCommGetAsyncError returned success on every triggered fault, including after the peer process had been killed with SIGKILL. The moment a peer is dead and the moment the communicator has a queryable asynchronous error are not the same moment. Detection in the experiment came from the control socket closing, from a 500 ms stall, or from the sender aborting itself. An upper layer that treats this API as an immediate reset, in the sense of a TCP connection reset, waits forever in this environment.

Abort also does not retract kernels already queued on the GPU stream. The receiver still schedules a large number of consumer kernels. Their recv slots were never filled with new data, and each cell shows about 2000–6400 of these empty reads. The communication control plane has stopped, and the GPU execution queue does not rewind itself. We clear the 16 B header before every recv, so those kernels see an empty header and refuse to consume. Another 14 messages had a header that looked new and a payload that was half new and half old. The checksum caught the tear. Without an explicit empty marker and an integrity check, both cases look like legal input.

Cleanup order decides whether recovery can finish. While a CUDA graph exec still holds the communicator, ncclCommAbort does not return for more than 35 s. Destroying the graph exec first, then aborting, returns normally. The lifetime dependency is inverted. The obvious order is to stop communication and then clean up the graph, but the graph still holds the communication resource, so abort waits on an object that will not be released first. For a graph-based inference engine this is more dangerous than a slow path. The wrong cleanup order turns a recoverable fault into a hang inside recovery itself.

How endpoint consumption removes the ambiguity

The rightmost group in Figure 2 draws a sharp boundary. The completion and communicator error that NCCL exposes are not enough to decide per-message recovery. The receiver’s consumption state is enough. Across 634 injections, sequence numbers, dedup, and resend-from-the-receiver-log brought both duplicates and losses to 0. On fault-free streams, 16 B throughput moved by +1.7% and 64 KB by -0.8%, both inside run-to-run noise.

The information that does the work is not a finer NIC completion, and it is not a faster communicator teardown. It is “which sequence number the application has consumed.” The two sides keep different books. Sender completion is an execution log: which operations I handed off. The receiver consumption log is an effect log: which operations actually changed application state. Recovery needs the second book. Reading only the first under-sends on the eager path and over-sends on the graph path. Filling the difference from the second book removes both errors. Under partial failure, the correct recovery boundary crosses the communication library and lands on the end-to-end application effect.

The result also marks what is still open. This is two-rank point-to-point communication. Partial completion of a collective has no per-message counterpart here. There is no GPUDirect RDMA, and the host bounce buffer widens the gap between local completion and consumption. The consumer kernel’s dedup check and log append sit in the same kernel. We did not inject a fault between them, so consumption that has an external side effect still needs its own atomicity story. Link jitter and MSCCL++ were not measured. MSCCL++ 0.9.0 has no abort, revoke, or per-operation status. The source shows a coarser recovery boundary, and that reading is not a substitute for fault-injection data.

What these data support directly

On NCCL 2.32.3 point-to-point send/recv, a completed sender CUDA event means only that the local send has finished. When the receiver dies or is aborted, each fault leaves 1 to 4 messages that were reported complete and never consumed. That is two V100s, EDR InfiniBand without GPUDirect RDMA, a window of 32, and 16 B or 64 KB messages. Deciding what to resend from sender completion drops those messages.

Communication captured in a CUDA Graph reports completion for the whole graph. At the fault the receiver may already have consumed part of the batch. Resending every operation that has not completed duplicates about 4 to 9 messages per fault. One graph here is 16 operations. A coarser grain means more duplicates. Graph communication needs a sequence number on each message. It cannot rely on completion of the graph.

When the peer is killed or paused, or when the local rank aborts itself, ncclCommGetAsyncError keeps returning success. Fault detection needs a separate control channel, or a timeout in the application.

After abort, receive kernels already queued still run on buffers that were never filled, and a payload can be torn. Clear the header before every recv, and put a checksum on the payload. Destroy the graph exec before ncclCommAbort.

The state that makes per-message exactly-once recovery work sits at the endpoint, not in the communication library. Sequence numbers, dedup, and a resend against the peer’s consumption log already cover the completion ambiguity seen here, at very low cost. What is still open is the boundary this protocol does not cover: whether consumption and the log have to be atomic, how to handle consumption that has an external side effect, and partial completion of a collective that cannot be written as one message.

Recovery should commit on what the receiver has used

The opening question was whether a given message should be resent when NCCL recovery runs. The experiment’s answer is not “read the error code,” and it is not “read sender completion.” It is to name the semantics you actually need to protect.

If the semantics are only “the send buffer may be reused now,” sender completion is enough. If they are “the bytes reached the peer,” you need delivery. If they are “this message has already changed receiver application state,” the commit point has to be consumption. Calling all three events completion is harmless on a healthy run. Under partial failure, the gap between them becomes a loss or a duplicate.

The result worth keeping is not the quirk of one NCCL API. It is a more general systems judgment. Failure recovery has to define commit around the application effect. It cannot treat transport or runtime completion as application completion. On these two V100s with EDR InfiniBand and no GPUDirect RDMA, that difference is already enough to produce a stable 1–4 missed resends or 4–9 duplicates. Once recovery queries the receiver’s consumption facts, both errors disappear.

The line that matters is who holds the last information about whether the operation really happened. In this experiment, that is the receiver application.