TENT Failover#
TENT hides transfer failures from the application by recovering inside the data path. This document describes how the recovery works, which knobs control it, and how it is tested.
The design has two layers:
Cross-transport failover in
TransferEngineImpl. When a transport fails a task at the completion stage, the engine moves that task to the next available transport (for example RDMA → TCP). Submit-stage failures are not retried today; see Known Gaps.Intra-RDMA rail recovery in
RailMonitor. When a specific (local NIC, remote NIC) rail keeps failing, the monitor pauses it with exponential cooldown; a successful transfer, or a probe / Half-Open trial on the paused rail, brings it back early.
Application code submits a batch and polls getTransferStatus. It never sees a FAILED task as long as any healthy path remains and the failover budget is not exhausted.
Fault Model#
TENT focuses on three kinds of transient faults:
Fault |
Surface |
Recovery action |
|---|---|---|
Work request completion error (WC error) |
RDMA worker sees a bad completion |
Rail-level |
QP / endpoint failure |
|
Not retried today: task surfaces as |
Peer disconnect mid-transfer |
|
Cross-transport failover |
Permanent or application-visible errors (invalid arguments, out-of-memory, segment not found) are not retried; they are returned to the caller as-is.
Architecture#
+-------------------------+
submitTransfer | TransferEngineImpl |
---------------> classify by TransportType
| submitTransferTasks |----failure----+
+-------------------------+ |
| v
| resubmitTransferTask
| (bump priority, pick next
| transport, resubmit)
v
+---------------------------+
| RdmaTransport / workers |
| +---------------------+ |
| | RailMonitor | |
| | per-rail state | |
| | cooldown / recover | |
| +---------------------+ |
+---------------------------+
Each request is owned by one
TaskInfo.typenames the transport currently executing the task;xport_priorityis the index into the ranked fallback list;failover_countcaps how many times we may re-resolve the transport.The ranked fallback list comes from
getTransportType(req, priority). Priority 0 yields the best available transport; increasing priority walks down the list;UNSPECmeans no transport left.RDMA rail state lives in
RailMonitor. Its lifecycle is independent of the task-level state machine: a rail can be paused while tasks keep flowing on other rails.
State Machine#
Cross-transport failover#
resubmitTransferTask is the single entry point that promotes a failing task to the next transport:
++task.failover_count
if failover_count > max_failover_attempts -> return error (exhausted)
task.xport_priority++
type = resolveTransport(task.request, task.xport_priority)
if type == UNSPEC -> return error (no transport)
transport_list_[type]->submitTransferTasks(...)
It has two callers, one per recoverable failure surface:
Completion-stage failure.
getTransferStatus(batch_id, task_id, status)and the batch-form overload callresubmitTransferTaskonce perFAILEDcompletion. On success the task is re-markedPENDINGso the aggregated batch status does not latch toFAILEDbecause of a task that is actually retrying.Exhaustion. When the budget is hit,
resubmitTransferTasksets the returned status toInvalidEntry("Failover limit exceeded, all transports exhausted"). Callers leavetask.typeunchanged; the task then reportsFAILEDthrough the normal status flow.
Submit-stage failures (submitTransferTasks returning non-OK) are not retried today. They mark the task as UNSPEC, and getTransferStatus short-circuits to FAILED. See Known Gaps for why.
RDMA rail recovery#
Inside RdmaTransport, each completion drives the rail monitor:
Bad completion →
rail.markFailed(local_nic, remote_nic)Good completion →
rail.markRecovered(local_nic, remote_nic)
markFailed bumps error_count inside error_window_. Once the count hits error_threshold_ the rail is paused until now + cooldown_; the cooldown escalates only when a fresh pause arms (was_paused == false), never within one burst — a single outage of N error WQEs in one window is one trip, not N. On a repeat failure across recovery cycles it doubles up to kMaxCooldown (300 s).
markRecovered clears the error count, un-pauses the rail, and resets the exponential-backoff memory so the next failure cycle starts from the initial cooldown. A fast path returns without work when the rail is already healthy, which is the common case on the completion hot path.
The admission gate is split into a pure predicate and a mutating admit:
isAvailable(local, remote) constreturns true only for a Closed/healthy rail. It never mutates state and never callsupdateBestMapping, so it is safe to call fromupdateBestMapping(the mapping rebuild) without the recursion the old singleavailable()had.admit(local, remote)is the only path that arms a probe/trial. For a Closed rail it returns true with no mutation. While a rail is Open (cooldown running) it admits one exploratory probe everyprobe_interval_; a probe success closes the rail early (transient-fault recovery in seconds rather than the full cooldown), a probe failure is a no-op for escalation (the rail was already paused, so defect A holds). When the cooldown expires without a probe proving the rail healthy,admittransitions to Half-Open and admits exactly ONE trial — it does not fully reopen. Escalation happens only if that trial fails (markFaileddoubles the cooldown and re-arms); a trial success closes the rail. Elapsed time alone never proves the path is healthy, so clock expiry alone never escalates and never floods a still-dead peer with every slice.probe_in_flightgates a second probe/trial from racing the first while it is on the wire (single-threaded per-worker ownership makes a bool sufficient); a pre-wire failure that armed a probe but never reached the wire callscancelProbeto roll the admission back.
This produces two recovery signals — a probe/trial result, and a live success on a posted transfer — so a flaky rail returns to service at the first good completion instead of waiting for the full cooldown, and a still-dead peer is probed at one trial per probe_interval_ rather than slammed every cycle.
Configuration#
All knobs live in the top-level transfer-engine.json. Defaults are safe for production; tune only if you have evidence.
Key |
Default |
Meaning |
|---|---|---|
|
|
Controls whether |
|
|
Upper bound on |
|
|
Number of failures inside |
|
|
Sliding window for counting rail errors. A failure older than the window resets |
|
|
Initial cooldown after tripping. Escalates on a fresh pause across recovery cycles, capped at 300 s. |
|
|
While a rail is paused, |
The RDMA keys are read by RailMonitor::load. Example:
{
"enable_auto_failover_on_poll": true,
"max_failover_attempts": 3,
"transports": {
"rdma": {
"rail_error_threshold": 3,
"rail_error_window_secs": 10,
"rail_cooldown_secs": 30
}
}
}
Observability#
Metric#
tent_transport_failover_total is a counter incremented once per successful transport switch inside resubmitTransferTask. A non-zero rate means the engine is actively recovering; a sudden jump usually points at a single bad link or flaky peer.
The counter is only built when TENT is compiled with -DTENT_METRICS_ENABLED=ON (see metrics.md). Without that flag the macro is a no-op.
Log keywords#
Keyword |
Interpretation |
|---|---|
|
A task has successfully switched transports. |
|
Task exhausted its budget and will surface |
|
|
|
Cooldown elapsed; the next |
|
Live success (or a successful probe/trial) brought a previously paused rail back early. |
|
Cooldown expired and |
Testing#
Real hardware faults are hard to stage, so failover is tested by driving the real TransferEngineImpl with FakeTransport backends and a fault-injecting decorator. The engine is unmodified: a completion-stage FAILED looks like a WC error or a dropped peer, and resubmitTransferTask runs as it would in production.
The harness — why a fake Transport is enough, how fakes are swapped in, and what this can and cannot prove — is in TENT Testing. That page is the mechanism; this section only notes what failover uses it for:
Completion-stage
FAILEDon the primary must resubmit on the next available transport.Submit-stage failure: a synchronous non-OK from
submitTransferTasksmust fail the owner task over to the remaining candidates (bounded bymax_failover_attempts) instead of terminal-failing it.Derived (merged) alias tasks must follow their owner’s recovered route and must not cause a duplicate physical submission.
Exhausting
max_failover_attempts(including0and1) must surfaceFAILEDand must not touch a transport beyond the budget.failover_countis per-task: one failing request must not spend another request’s budget.With
enable_auto_failover_on_poll=false, status polling is observational;progressBatch/waitTransferCompletion/transferSyncstill recover.
Submit-stage recovery is exercised through the same FakeTransport harness with force_submit_fail backends (see SubmitStageFailureFailsOverToSecondary and friends in engine_failover_e2e_test.cpp).
Known Gaps#
Submit-stage failure recovery is implemented (previously a gap): a synchronous non-OK from
submitTransferTasksfails the owner tasks over to the remaining candidate transports (bounded bymax_failover_attempts), reusingresubmitTransferTask. The two hazards that originally made a naive retry unsafe are handled as follows:Merged requests. Only the owner task of each merged request is resubmitted; derived aliases mirror the owner’s recovered route (transport, sub-batch slot, status) instead of being resubmitted themselves, preserving the merge pass’s deduplication.
Partial enqueue.
NVLinkTransport::submitTransferTasksrolls back its half-appended task entries on synchronous failure (no I/O has started at that point). Transports that start per-request work inside the submit loop (e.g.ShmTransport, whose copies are synchronous and idempotent) may still cause a duplicate same-content transfer on the fallback; re-posting an identical request is a data-level no-op for KV-cache-style workloads. A failover target whose own synchronous submit also fails is retried within the same budget; a task that exhausts all candidates surfacesFAILEDattributed to the last attempted transport (neverunspec).
A Closed rail resets its backoff memory entirely on recovery (
markRecovered, or expiry reopen inisAvailable’s predecessor behavior). Escalation across cycles now happens via a Half-Open trial failure (markFaileddoubles the cooldown and re-arms), not via clock expiry — so a rail that flaps does back off harder when its trial probes keep failing, while a proven-healthy recovery starts fresh. If this ever proves too aggressive the fix is to decay rather than escalate on trial failure.Cross-transport failover is driven purely by return status; there is no latency-based “this transport is healthy but too slow, try another” signal. That belongs to the scheduler, not this document.
Runtime-layer failover is covered by FakeTransport tests in the
tent-cicuda-offlegs. DMA integrity, real WC errors, and staging under NVLink still need hardware runners; see TENT Testing.