[serve] Push-based replica health: replica self-checks and heartbeats - #66421
johntaylor-cell wants to merge 11 commits into
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces a replica-pushed self-health mechanism to Ray Serve, allowing replicas to periodically run their own health checks and push the results to the controller, thereby reducing active pull-probing overhead. It also adds a mechanism to detect controller ingest lag and avoid incorrectly marking replicas as unhealthy due to probe timeouts when the controller is falling behind. Feedback on the changes highlights a potential ZeroDivisionError in expected_push_rate if the metrics interval or health check period is configured to zero, which should be guarded against.
There was a problem hiding this comment.
Code Review
This pull request introduces a push-based health check mechanism for Ray Serve replicas, allowing replicas to periodically run their own health checks and push the results to the controller. This mechanism includes a fallback to pull probes when pushed results go stale, and an ingest lag gate to suppress probe timeouts when the controller is lagging. Feedback on these changes highlights a potential memory leak in ReplicaHealthPushRegistry where dead replica entries are never pruned if the total count remains below _PRUNE_THRESHOLD (65,536), and suggests triggering pruning periodically regardless of the size threshold.
| if ( | ||
| len(self._state) > self._PRUNE_THRESHOLD | ||
| and now - self._last_prune_time > self._PRUNE_MIN_INTERVAL_S | ||
| ): |
There was a problem hiding this comment.
In clusters with high replica churn but a total lifetime replica count below _PRUNE_THRESHOLD (65,536), dead replica entries will never be pruned and will accumulate in self._state indefinitely, causing a slow memory leak. Since pruning is extremely cheap for smaller dictionary sizes, we should also trigger a prune periodically (e.g., every _PRUNE_MAX_AGE_S seconds) regardless of the size threshold, while still respecting the minimum interval to avoid thrashing.
| if ( | |
| len(self._state) > self._PRUNE_THRESHOLD | |
| and now - self._last_prune_time > self._PRUNE_MIN_INTERVAL_S | |
| ): | |
| if ( | |
| len(self._state) > self._PRUNE_THRESHOLD | |
| or now - self._last_prune_time > self._PRUNE_MAX_AGE_S | |
| ) and now - self._last_prune_time > self._PRUNE_MIN_INTERVAL_S: |
Adds the registry the controller records pushed replica self-health into, and the arbitration that lets a fresh push stand in for a pull probe: a probe is deferred while the newest observation is inside 1.5x the health-check period, and two watermarks keep a push and an in-flight probe from overwriting each other's verdict. An actor crash stays authoritative either way. Inert on its own. No replica pushes yet, so the registry stays empty, every replica is probed exactly as before, and the pull path is unchanged. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
…allers Makes record_pushed_health self-contained: it now keeps the newer of the stashed and incoming observations rather than relying on the registry's ordering guard. Unreachable through today's only caller, but the later PRs in this stack add producers upstream of that funnel. Also restores the rationale for clearing the failure flag while keeping the latency sample when a stale in-flight probe is discarded. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
…resh The deferral gate only consulted the last APPLIED push, but a push that a resolving probe beat to the tick stays stashed and unconsumed. With a probe outstanding longer than the health-check period -- a timeout, the case at scale -- the gate then armed a fresh probe against health already in hand. Take the newer of the applied and stashed arrival times instead. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
…failures Two push/probe arbitration tests read time.time() twice and relied on the second read being strictly greater, which is what the stale-probe drop is gated on. On Windows time.time() resolves to the ~15.6ms system timer tick, so the two reads tie, the drop never fires, and both tests fail there instead of pinning the arbitration they describe. Pin the two instants explicitly. Mirror the pushed consecutive-failure count with a max() against the count the controller already holds. The absolute assignment let a push arriving mid-failure-run lower a probe-derived count, moving a failing replica further from eviction than the pull-only path would have. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
test_apply_pushed_health_hands_off_to_wrapper asserted only that the wrapper was called, never with what. Since _apply_pushed_health forwards the registry entry as a bare splat, swapping the two floats was silent, and it would have fed replica-clock in as received_at, the exact skew the field exists to avoid. Assert the arguments, and cover the DeploymentReplica hop, which had no direct test. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
The age-based prune only runs when the registry holds more than _PRUNE_THRESHOLD entries, so a fleet that never crosses 65536 distinct replica ids keeps every dead replica's entry for the life of the controller. The terminal reap already evicts four sibling per-replica maps with the unique id in hand, so evict here too and leave the prune as the backstop for ids that stop pushing without a stop event. A review bot read this as an unbounded leak. It is not: the size gate is a level trigger, so once the registry crosses the threshold it re-prunes every _PRUNE_MIN_INTERVAL_S and the ceiling carries no uptime term (measured 12.4 MiB at the threshold). The real cost is retention time below the threshold, and that nothing reclaims once pushes stop, which a prune reachable only from record() cannot fix. MockReplicaActorWrapper gained record_pushed_health, without which no deployment-state-level test could exercise the push-apply path at all. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
9c0c77b to
0ea2058
Compare
…ence A crashed replica is invisible while a pushed result is still fresh, because ACTOR_CRASHED only ever comes from an outstanding probe ref and the window delays arming one. At 1.5x the health-check period that made crash detection slower than pull probing alone: ~15s against ~11s at the default period. The window is keyed on the push cadence, which is twice the probe cadence, so 0.75x is enough to never probe a replica whose pushes are arriving. Detection is then ~7.5s, better than before pushes existed. Pinned with absolute values rather than by calling the function under test, so changing the factor has to be deliberate. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
The controller probes every replica on a timer; at 8k the submit alone is ~465 us and GIL-bound, and the sweep accounts for most of deployment-state update time. Replicas now run their own check twice per health-check period and heartbeat the result, so the probe fires only when a pushed result goes stale (the fallback added in the previous PR). Behind RAY_SERVE_ENABLE_PUSH_HEALTH, default off. The flag exists so the rest of the stack can be reviewed without changing behavior, and the last PR in the stack deletes it. The in-flight guard is bounded: a push overdue past RAY_SERVE_METRICS_PUSH_STUCK_S is abandoned rather than waited on. An unbounded guard lets one stuck receiver silence a replica permanently, so the same bound now also covers the pre-existing metric push, which had the identical unguarded pattern. Probe timeouts are not charged to a replica while the controller's own ingest is behind: the registry compares arrivals against the rate the fleet should be publishing, and a timed-out probe then scores as no information rather than a failure. A crashed actor still reports ACTOR_CRASHED either way. With the feature off the expected rate stays 0, which keeps the gate inert. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
…cadence expected_push_rate() claimed a deployment owes one health-bearing report per metrics_interval_s whenever that is at most the health-check period. Nothing carries health on metric reports yet, so the only thing feeding the registry is the heartbeat, which fires twice per health-check period. With metrics_interval_s below half the period the expectation exceeded reality by more than 2x, so ingest_lagging() latched on permanently and probe timeouts were never charged. A replica whose health check hung would then never accumulate strikes and never be replaced. Measured at 100 replicas with a 10s period and a 2s interval: expected 50/s against an actual 20/s, and the gate reads lagging at anything under 25/s. The metric cadence becomes correct once reports carry health, and moves to that PR. health_check_period_s is a PositiveFloat, so the remaining division cannot be by zero. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
The ingest-lag gate could be held open by the silence it was judging. Replicas that hang stop heartbeating, which drives the observed rate below the expected one, which suppresses their probe timeouts, which keeps them RUNNING and still counted in the expectation. Suppression is now bounded per replica (RAY_SERVE_MAX_SUPPRESSED_HEALTH_TIMEOUTS, 3), so a transient controller backlog is still absorbed but a permanent silence resolves. Any real verdict, pushed or probed, refills the budget. The expectation also assumed two heartbeats per period, but the pusher sleeps its interval after the eval, so the real cadence is eval + period/2. A user check slower than half the period read as controller lag and pinned the gate on fleet-wide. Expect one per period instead, which the docstring's "erring low is safe" already called for. Dropped the reentrancy shortcut in _run_user_health_check. It let a waiter adopt the result of a check the controller had already timed out, so a check taking between one and two timeouts alternated timeout and success and never reached the failure threshold. This was reachable with the feature off. The lock stays, so concurrent callers still serialize. A cancelled check now confirms nothing: it neither refreshes the cached verdict nor flips _healthy, which backs /-/healthz and so the replica's place in the load balancer rotation. CancelledError is also caught in the pusher, since MetricsPusher only catches Exception and would otherwise retire the heartbeat for the replica's lifetime. A suppressed timeout no longer records a latency sample, which made the one episode where failures are ignored read as merely slow. Tests: the two wirings that joined the halves of the feature were untested and are now pinned, along with the suppression bound, the ratio from above, the cached verdict, the waiter, and cancellation. Reverting the bound or restoring the coalescing each fail a test. Also switched the flag to get_env_bool, un-orphaned a constant comment, and reverted the unused mock parameter. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
0ea2058 to
deb1a94
Compare
A replica counts an unhealthy self-check once per health-check period but heartbeats on every eval, so the same failure count reaches the controller twice. Each arrival was consumed as a fresh APP_FAILURE and advanced the controller's count, which reached the threshold in about half the configured time and replaced replicas with flaky checks sooner than intended. The controller now advances only when the replica's reported count actually increases, cancelling the chain's increment otherwise. A count that genuinely advances still mirrors as before, so a push stream starting mid-failure-run cannot lower what the controller already probed. Co-Authored-By: Claude <noreply@anthropic.com> Signed-off-by: john.taylor <john.taylor@anyscale.com>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.
Reviewed by Cursor Bugbot for commit 91d4d34. Configure here.
| self._consecutive_health_check_failures, | ||
| pushed_failures - 1, | ||
| ) | ||
| self._last_mirrored_push_failures = pushed_failures |
There was a problem hiding this comment.
Stale failure watermark skips new strikes
Medium Severity
_last_mirrored_push_failures is not cleared when a replica recovers. A later unhealthy push that restarts at count 1 matches the stale watermark, so the increment is cancelled and the strike is dropped. Intermittent fail-then-recover cycles can therefore never reach the unhealthy threshold while fresh heartbeats keep deferring pull probes.
Additional Locations (1)
Reviewed by Cursor Bugbot for commit 91d4d34. Configure here.


Stacked on #66243, which added the controller-side receive path.
The controller health-checks every replica by firing
check_health.remote()on atimer. At 8k replicas the submit alone is ~465 us and GIL-bound, and the sweep
accounts for essentially all of deployment-state update time.
Replicas now run their own check twice per health-check period and heartbeat the
result, so the pull probe fires only when a pushed result goes stale.
Behind
RAY_SERVE_ENABLE_PUSH_HEALTH, default off. The flag exists so the restof the stack can be reviewed without changing behavior; the last PR deletes it.
The in-flight guard is bounded. A push overdue past
RAY_SERVE_METRICS_PUSH_STUCK_S(20 s) is abandoned rather than waited on. Anunbounded guard lets one stuck receiver silence a replica permanently — measured
during development, one node's replicas dropped from ~420 to ~1.5 reports/s and
stayed there. The pre-existing metric push had the identical unguarded pattern, so
the same bound now covers it too.
Probe timeouts are not charged to a replica while the controller is behind. The
registry compares arrivals against the rate the fleet should be publishing; below
half of it, a timed-out probe scores as no information rather than a failure. It
deliberately declines to decide whose fault the silence is, because the only thing
it gates is whether to count a strike, and declining is safe either way. A returned
failure still counts, and a crashed actor still reports
ACTOR_CRASHED. With thefeature off the expected rate stays 0, so the gate is inert.
Testing. 24 new unit tests. Verified by mutation: each of ten independent
removals — the lag gate, the timeout flag, the lagging ratio, the rate window, the
expected-rate cadence, the stuck bound, once-per-period failure counting, the
threshold latch, the unhealthy bypass, and repeat-unhealthy suppression — is caught
by at least one test.
test_deployment_state.py294 passed andtest_metrics_utils.py108 passed, withtest_replica_backpressure,test_replica_quiesce,test_user_callable_wrapper,test_controller,test_application_stateandtest_deployment_rank_managerunchanged.