nemo_rl.models.generation.fleet_health#
Liveness bookkeeping for the vLLM generation fleet.
Both SingleController rollout paths pick a generation shard by static round-robin with
no idea whether the shard is alive: the native path in
VllmGeneration._async_generate_base, and the NeMo-Gym path inside Gym’s own
_resolve_client. This module owns the missing half — which shards are eligible to
serve, and why.
Deliberately a pure state machine. Probing, restarting and pushing membership are I/O and belong to the caller, which keeps every transition here testable without Ray, a network, or a GPU — and keeps one description of “eligible” that both routing adapters read.
The transition that carries the most weight is the one that does not exist: a shard
cannot go from DEAD back to HEALTHY on its own. A restarted engine holds whatever
weights it loaded at init, so re-admitting it because it started answering probes again
would feed training rollouts generated from a checkpoint hundreds of steps stale –
invisible, and far worse than the outage that caused it. Recovery must pass through
STALE and a completed refit.
Module Contents#
Classes#
Lifecycle of one generation data-parallel shard. |
|
Everything known about one shard. |
|
Thresholds resolved from |
|
Tracks per-shard health and exposes the current serving set. |
|
Picks a serving shard, preferring the one with the fewest requests in flight. |
Data#
API#
- class nemo_rl.models.generation.fleet_health.ShardState#
Bases:
str,enum.EnumLifecycle of one generation data-parallel shard.
Initialization
Initialize self. See help(type(self)) for accurate signature.
- HEALTHY#
‘healthy’
- SUSPECT#
‘suspect’
- DEAD#
‘dead’
- RESTARTING#
‘restarting’
- STALE#
‘stale’
- RETIRED#
‘retired’
- nemo_rl.models.generation.fleet_health._SERVING_STATES#
‘frozenset(…)’
- nemo_rl.models.generation.fleet_health._ABSENT_STATES#
‘frozenset(…)’
- class nemo_rl.models.generation.fleet_health.ShardHealth#
Everything known about one shard.
- dp_shard_idx: int#
None
- base_url: Optional[str]#
None
- consecutive_probe_failures: int#
0
- consecutive_probe_successes: int#
0
- consecutive_reported_failures: int#
0
- restart_attempts: int#
0
- state_before_partial: Optional[nemo_rl.models.generation.fleet_health.ShardState]#
None
- weight_version: int#
0
- last_ok_at: float#
0.0
- last_error: str = <Multiline-String>#
- property is_serving: bool#
- exception nemo_rl.models.generation.fleet_health.GenerationFleetExhausted#
Bases:
RuntimeErrorToo few shards remain for the run to be worth continuing.
Initialization
Initialize self. See help(type(self)) for accurate signature.
- class nemo_rl.models.generation.fleet_health.FleetHealthPolicy#
Thresholds resolved from
async_rl.generation_fleet_health.Mirrors :class:
nemo_rl.experience.rollout_manager.RolloutTimeoutsin shape: an internal dataclass so the BaseModel stays the single home for user-facing defaults.- unhealthy_threshold: int#
3
- healthy_threshold: int#
2
- max_restart_attempts_per_shard: int#
5
- min_healthy_shards: int#
1
- __post_init__() None#
- class nemo_rl.models.generation.fleet_health.GenerationFleetHealth(
- *,
- shard_count: int,
- policy: nemo_rl.models.generation.fleet_health.FleetHealthPolicy,
- base_urls: Optional[list[Optional[str]]] = None,
- clock: Optional[Callable[[], float]] = None,
Tracks per-shard health and exposes the current serving set.
- Parameters:
shard_count – Number of generation data-parallel shards.
policy – Thresholds governing the transitions.
base_urls – Per-shard OpenAI base URLs, when the deployment has them.
clock – Monotonic time source; injectable so tests need not sleep.
Initialization
- property membership_epoch: int#
- property shard_count: int#
- snapshot() list[nemo_rl.models.generation.fleet_health.ShardHealth]#
Per-shard health, ordered by shard index.
- serving_shards() list[int]#
Shard indices currently eligible to be handed traffic.
- absent_shards() list[int]#
Shard indices whose process cannot take part in a collective.
The set the weight-refit path cares about, and deliberately not the complement of
- Meth:
serving_shards: a SUSPECT or STALE shard is withheld from traffic but its process is alive and joins the refit normally.
- suspected_shards() list[int]#
Shards the ledger already doubted when the refit went wrong.
SUSPECT, or STALE having been SUSPECT immediately before an aborted refit marked it partial. The second case is the common one at recovery time and would otherwise be invisible:
mark_weights_partialruns over every serving shard, and SUSPECT is a serving state, so the suspicion is overwritten by STALE microseconds before anyone asks.state_before_partialis what preserves it.Not for routing –
serving_shardsdecides that. This exists for the one caller that has independent evidence something in the collective stopped participating and needs to know which shard the ledger was already unhappy with.
- state_of(
- shard_idx: int,
- serving_base_urls() list[str]#
Base URLs of the serving shards, for pushing to the NeMo-Gym router.
- counts_by_state() dict[nemo_rl.models.generation.fleet_health.ShardState, int]#
- as_metrics() dict[str, float]#
Flatten into a metric dict for the SingleController logger.
- record_probe(shard_idx: int, *, ok: bool, error: str = '') None#
Fold one probe result into a shard’s state.
Probes never resurrect a shard that has left the serving set:
DEAD,RESTARTING,STALEandRETIREDall ignore them, because getting an answer says nothing about whether the weights are current.
- record_actor_death(shard_idx: int, error: str = '') None#
Record proof that a shard’s process is gone. DEAD at once, no counting.
The counters in :meth:
record_probeexist to tell a slow shard from a dead one, because a probe timeout cannot distinguish them. Some evidence carries no such ambiguity: Ray reporting its actor dead means the process is gone, full stop, and making that wait forunhealthy_thresholdmore rounds of the same answer only delays the conclusion.The delay was not academic. Detection took
probe_interval_s * unhealthy_threshold, which the refit deadline could expire inside – so a refit hung on a dead rank aborted while the monitor still had that rank SUSPECT, and the rebuild the abort exists to trigger had an empty absent set to work from. Job 5925668.Ignores shards that are already absent, so a repeat report is idempotent, and RETIRED is never disturbed.
- condemn_silent_participant(shard_idx: int, *, reason: str) None#
Quarantine a shard that is alive but did not take part in a collective.
Distinct from :meth:
record_actor_death, which means “Ray says the process is gone”. This one means “the process answers and still broke the refit”, which is the failure modeis_alivecannot see: it is answered by the Ray actor and never touches the engine, so a wedged engine reads as healthy forever.DEAD rather than STALE, because the point is to make it absent: STALE keeps it in the refit, which is precisely what just hung. The caller owes an independent reason to believe this shard is the culprit – see the single-suspect rule in
_recover_from_failed_refit– because this is a judgement, not an observation.
- report_failure(shard_idx: int, error: BaseException) None#
Record a failure observed by a routing adapter rather than by a probe.
The adapters are the only components actually issuing generation requests, so they see failures a liveness probe cannot – a shard that answers
/healthand still errors on every generation.Counted on its own streak rather than folded into :meth:
record_probe. That used to be the implementation, and it could not condemn the very shard it exists for: a wedged engine answersis_alive, so an ok-probe everyprobe_interval_sreset the shared counter, and reported failures only accumulated ifunhealthy_thresholdof them landed inside one probe window. Under load they do; under a trickle the shard oscillated HEALTHY<->SUSPECT forever. A streak that only a successful generation clears says what is actually meant: this shard has failed N requests in a row.
- report_success(shard_idx: int) None#
Record a generation that completed on this shard.
The reset half of :meth:
report_failure. Without it the reported streak is monotonic and every shard eventually reachesunhealthy_thresholdgiven a long enough run, however healthy it is.
- shard_for_base_url(url: str) Optional[int]#
Reverse the shard -> base URL mapping, for failures reported by URL.
The NeMo-Gym router knows its backends only as URLs; the ledger keys everything by shard index.
- mark_restarting(shard_idx: int) None#
A replacement engine is being brought up for a dead shard.
- mark_loaded(
- shard_idx: int,
- *,
- base_url: Optional[str] = None,
The replacement finished loading. It holds stale weights until refit.
- Parameters:
base_url – the replacement’s URL. A restarted engine binds a new port, so without this the monitor keeps publishing the dead one – and
serving_base_urls()feeds the NeMo-Gym router, which would send every rollout to a socket nobody is listening on.
- mark_weights_partial(shard_idx: int) None#
An aborted refit left this shard holding a mix of old and new weights.
STALE rather than DEAD: the process is alive and refits normally, it just must not serve until a refit completes. That is exactly what STALE already means, and because DEAD -> HEALTHY is unreachable and STALE ignores successful probes, the only route back into the serving set is report_refit – which is the property that makes partial weights safe to hold.
Absent shards are left alone; they are not serving and their problem is not weights.
- mark_restart_failed(shard_idx: int, *, error: str = '') None#
A restart attempt did not bring the engine up. Back to DEAD.
Needed as its own transition because
record_probedeliberately ignores non-serving states – a probe must never resurrect a shard – so a failed restart reported that way would leave the shard stuck in RESTARTING: never retried, because it is no longer DEAD, and never retired, because retirement is driven by restart attempts.errormatters more here than the two siblings that already take one.retirehas exactly one caller – attempts exhausted – solast_erroris the only structured field that can say what the reloads failed on. Without this it still holds the original death, which is wrong and plausible enough to be believed: a GPU that an orphaned EngineCore held for 370s and aseed=Noneconfig error want completely different responses and would otherwise be indistinguishable in the record.
- record_weight_version(shard_idx: int, *, weight_version: int) None#
This shard received this refit’s weights. Nothing else changes.
The stamp without the promotion, for the shards a refit reached that are not STALE – which on a run where nothing has died is all of them.
report_refitcannot serve that case: it is a state transition, and driving one on a HEALTHY shard after every refit would clear the reported-failure streak that is the only counter able to condemn an engine that still answersis_alive.Absent shards received nothing, so the caller filters them out before calling; the RETIRED guard here is belt-and-braces, matching
report_refit.
- report_refit(shard_idx: int, *, weight_version: int) None#
A completed refit is the only way back into the serving set.
- retire(shard_idx: int, *, reason: str) None#
Remove a shard permanently. Training continues on what is left.
- raise_if_exhausted() None#
Raise once too few shards remain for the run to be worth continuing.
- _transition(
- shard: nemo_rl.models.generation.fleet_health.ShardHealth,
- new_state: nemo_rl.models.generation.fleet_health.ShardState,
- _refresh_membership() None#
- class nemo_rl.models.generation.fleet_health.HealthyShardSelector#
Picks a serving shard, preferring the one with the fewest requests in flight.
Least-outstanding rather than round-robin because it is strictly better for LLM serving: it steers away from a shard that is merely slow or wedged without needing that to be diagnosed first.
- _inflight: dict[int, int]#
‘field(…)’
- next_shard() int#
Return the shard index to serve the next request.
- Raises:
NoHealthyShards – No shard is currently eligible.
- acquire(shard_idx: int) None#
- release(shard_idx: int) None#
- inflight(shard_idx: int) int#