From 89243387fc322f19d5968ae2cf3c3c8cf49029f4 Mon Sep 17 00:00:00 2001 From: Eli Gottlieb Date: Fri, 24 Jul 2026 12:18:34 -0700 Subject: [PATCH 1/6] fix(v1): detect env-server worker death instead of hanging MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An env server that died was invisible to its client. Three gaps, all of which turned a crash into an unbounded wait: - The broker reported its address as soon as the frontend socket was bound, before any worker had loaded the env. A spawner was told "ready" about a pool whose env was still loading — or already dead — so an env that fails to import surfaced much later as a rollout that never returns. - A dead worker was never reaped. Its `active` count never drained, so least-busy dispatch elected the corpse for every subsequent request, and its in-flight requests went unanswered forever (`EnvClient` awaits rollout replies untimed, by design). - `run_eval_server` read the address queue with a bare 600s timeout, so a server that died on startup was reported as a timeout rather than a crash. The broker now reports its address only once a worker answers `health`, polls worker exit codes on its serve loop, answers a dead worker's in-flight requests with an error the client raises on, and replaces the worker up to MAX_WORKER_RESTARTS before failing requests outright. `wait_for_address` is the shared queue reader that watches the process's exit code alongside the queue and quotes the tail of its log file — used here by `run_eval_server`, and by prime-rl's orchestrator for the servers it spawns. Co-Authored-By: Claude Opus 5 (1M context) --- verifiers/v1/cli/eval/runner.py | 10 +- verifiers/v1/serve/__init__.py | 10 +- verifiers/v1/serve/pool.py | 330 ++++++++++++++++++++++++-------- 3 files changed, 271 insertions(+), 79 deletions(-) diff --git a/verifiers/v1/cli/eval/runner.py b/verifiers/v1/cli/eval/runner.py index 8eb1fb4e9..3e925b54a 100644 --- a/verifiers/v1/cli/eval/runner.py +++ b/verifiers/v1/cli/eval/runner.py @@ -109,7 +109,12 @@ async def run_eval_server(config: EvalConfig) -> list[Episode]: from verifiers.v1.utils.logging import setup_logging from verifiers.v1.configs.cli.env import pool_serve_kwargs - from verifiers.v1.serve import EnvClient, env_config_data, serve_env + from verifiers.v1.serve import ( + EnvClient, + env_config_data, + serve_env, + wait_for_address, + ) legacy = config.is_legacy server_kwargs = ( @@ -155,9 +160,8 @@ async def run_eval_server(config: EvalConfig) -> list[Episode]: proc.start() child_conn.close() # the child holds its end; we keep parent_conn so our exit closes it try: - address = await asyncio.to_thread(address_queue.get, timeout=600) + address = await wait_for_address(proc, address_queue, log_file=log_file) client = EnvClient(address=address) - await client.wait_for_server_startup(timeout=600) # A v1 run dispatches — and resumes — tasks by content: the client owns them, # and `resume.task_key` is their identity. Only the legacy bridge is addressed # by dataset row (its dataset lives server-side, reported via `info`), and diff --git a/verifiers/v1/serve/__init__.py b/verifiers/v1/serve/__init__.py index fdfce74e8..ceb746891 100644 --- a/verifiers/v1/serve/__init__.py +++ b/verifiers/v1/serve/__init__.py @@ -1,5 +1,11 @@ from verifiers.v1.serve.client import EnvClient -from verifiers.v1.serve.pool import EnvServerPool, env_config_data, serve_env +from verifiers.v1.serve.pool import ( + ENV_SERVER_SPAWN_TIMEOUT, + EnvServerPool, + env_config_data, + serve_env, + wait_for_address, +) from verifiers.v1.serve.server import EnvServer from verifiers.v1.serve.types import ( HealthRequest, @@ -17,6 +23,8 @@ "EnvServerPool", "serve_env", "env_config_data", + "wait_for_address", + "ENV_SERVER_SPAWN_TIMEOUT", "EnvClient", "HealthRequest", "HealthResponse", diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index 3d0b2a255..12434bf0d 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -16,9 +16,15 @@ Scaling is elastic but upscale-only: a new worker is spawned when in-flight requests reach 90% of current capacity (`workers * multiplex`). Workers are spawned `spawn`-style (own env, own loop) and monitor a death pipe so an orphaned worker self-exits if the -broker dies. TODO: downscale idle workers, per-worker restart-on-death, stats/lag -monitors (v0 had them; omitted here — rollout errors are returned as data, not crashes, -so worker death is rare). +broker dies. TODO: downscale idle workers, stats/lag monitors (v0 had them). + +A worker that dies is reaped: `EnvClient` awaits rollout replies untimed, so an +unanswered request is an unbounded hang, and a dead worker's `active` count never drains +— least-busy dispatch would elect the corpse for every subsequent request. The broker +therefore polls its workers' exit codes, answers their in-flight requests with an error, +and replaces them (up to `MAX_WORKER_RESTARTS`). Readiness is a worker handshake, not a +socket bind: the broker reports its address only once a worker answers `health`, so a +spawner is never told "ready" about a pool whose env failed to load. """ import asyncio @@ -26,10 +32,14 @@ import logging import multiprocessing as mp import os +import queue import signal import threading +import time import uuid from collections.abc import Callable +from multiprocessing.process import BaseProcess +from pathlib import Path import msgpack import zmq @@ -37,12 +47,28 @@ from verifiers.v1.configs.env import EnvConfig from verifiers.v1.serve.server import EnvServer -from verifiers.v1.serve.types import HealthResponse, RunGroupRequest +from verifiers.v1.serve.types import BaseResponse, HealthResponse, RunGroupRequest logger = logging.getLogger(__name__) _HEALTH = msgpack.packb(HealthResponse().model_dump(mode="json"), use_bin_type=True) +MAX_WORKER_RESTARTS = 3 +"""Replacements the broker spawns over its lifetime: enough to ride out a one-off crash +(an OOM, a segfaulting sandbox), few enough that an env dying on every load fails loudly +instead of respawning forever.""" + +ENV_SERVER_SPAWN_TIMEOUT = 600.0 +"""Default deadline for a spawned server to report its address. Generous: it reports only +once its env is loaded, which for a legacy bridge includes pulling a dataset.""" + +# Broker loop wake-up, so exited workers are reaped even while no socket is readable. +_REAP_INTERVAL_S = 0.25 +_REAP_INTERVAL_MS = int(_REAP_INTERVAL_S * 1000) +# Startup handshake: a `health` request whose reply proves a worker loaded its env. +_STARTUP_PROBE = b"startup" +_EMPTY_PAYLOAD = msgpack.packb({}, use_bin_type=True) + def _arm_teardown(death_pipe=None) -> None: """Arm a spawned process (serve_env broker/single server, or pool worker) for clean @@ -74,10 +100,10 @@ class EnvServerPool: With `elastic=True` (default) it starts with a single worker and spawns another whenever in-flight requests reach 90% of current capacity (`workers * multiplex`), up - to `max_workers`. Upscale-only for now — workers are never reclaimed. `elastic=False` - pre-spawns all `max_workers` upfront (the old fixed-pool behavior). The broker forwards - opaque request frames, so workers can be `EnvServer` (v1) or `LegacyEnvServer` (v0) - without the broker caring.""" + to `max_workers`. Upscale-only for now — a live worker is never reclaimed, though a dead + one is (see the module docstring). `elastic=False` pre-spawns all `max_workers` upfront + (the old fixed-pool behavior). The broker forwards opaque request frames, so workers can + be `EnvServer` (v1) or `LegacyEnvServer` (v0) without the broker caring.""" def __init__( self, @@ -88,6 +114,7 @@ def __init__( log_setup: Callable[[], None] | None = None, multiplex: int = 128, elastic: bool = True, + address_queue=None, ) -> None: self.server_kwargs = server_kwargs self.max_workers = max_workers @@ -95,10 +122,17 @@ def __init__( self.elastic = elastic self.legacy = legacy self.log_setup = log_setup + self.address_queue = address_queue self.session = uuid.uuid4().hex[:12] self.workers: list[dict] = [] + # request_id -> {client_id, worker, rollout_slots} + self.pending: dict[bytes, dict] = {} + self.in_flight = 0 self._mpctx = mp.get_context("spawn") self._poller: zmq.asyncio.Poller | None = None + self._next_index = 0 + self._restarts = 0 + self._last_reap = 0.0 self.ctx = zmq.asyncio.Context() self.frontend = self.ctx.socket(zmq.ROUTER) @@ -113,7 +147,8 @@ def _worker_path(self, i: int) -> str: return f"/tmp/vf-pool-{self.session}-{i}" def _spawn_worker(self) -> None: - i = len(self.workers) # upscale-only, so the next index is the current count + i = self._next_index # monotonic: a retired worker's path is never reused + self._next_index += 1 address = f"ipc://{self._worker_path(i)}" parent_conn, child_conn = self._mpctx.Pipe() proc = self._mpctx.Process( @@ -145,7 +180,7 @@ def _spawn_worker(self) -> None: if self._poller is not None: self._poller.register(dealer, zmq.POLLIN) - def _maybe_scale_up(self, in_flight: int) -> None: + def _maybe_scale_up(self) -> None: """Spawn one more worker when in-flight rollout slots reach 90% of capacity. A new worker starts at `active=0`, so least-busy dispatch funnels the backlog to @@ -153,19 +188,155 @@ def _maybe_scale_up(self, in_flight: int) -> None: up once already saturated. `max_workers=None` scales without a cap.""" if self.max_workers is not None and len(self.workers) >= self.max_workers: return - if in_flight >= 0.9 * len(self.workers) * self.multiplex: + if self.in_flight >= 0.9 * len(self.workers) * self.multiplex: self._spawn_worker() logger.info( "EnvServerPool scaled up to %d/%s workers (in_flight=%d)", len(self.workers), self._cap_str, - in_flight, + self.in_flight, ) @property def _cap_str(self) -> str: return "inf" if self.max_workers is None else str(self.max_workers) + async def _wait_for_worker(self, w: dict) -> None: + """Block until a worker answers `health`. + + A worker binds its socket before loading the env, so its first reply — not its + bind — is what proves the env loaded. Untimed on purpose: a slow load (a legacy + dataset pull) is the spawner's deadline to enforce, not ours. We only insist that + the process is still alive.""" + dealer = w["dealer"] + await dealer.send_multipart([_STARTUP_PROBE, b"health", _EMPTY_PAYLOAD]) + while not await dealer.poll(timeout=_REAP_INTERVAL_MS): + if w["process"].exitcode is not None: + raise RuntimeError( + f"env server worker {w['index']} died while loading the env " + f"(exit code {w['process'].exitcode})" + ) + await dealer.recv_multipart() + + async def _reply_error( + self, client_id: bytes, request_id: bytes, error: str + ) -> None: + """Answer a request with a failure the client raises on. `EnvClient` awaits + rollout replies untimed, so leaving one unanswered is an unbounded hang.""" + data = msgpack.packb( + BaseResponse(success=False, error=error).model_dump(mode="json"), + use_bin_type=True, + ) + with contextlib.suppress(zmq.ZMQError): + await self.frontend.send_multipart([client_id, request_id, data]) + + def _retire(self, w: dict) -> None: + """Drop a worker from dispatch and release everything it held.""" + self.workers.remove(w) + if self._poller is not None: + self._poller.unregister(w["dealer"]) + with contextlib.suppress(Exception): + w["dealer"].close() + with contextlib.suppress(Exception): + w["pipe"].close() + with contextlib.suppress(OSError): + os.unlink(self._worker_path(w["index"])) + + def _replace_worker(self) -> None: + """Spawn a replacement for a dead worker, within the restart budget.""" + if self._restarts >= MAX_WORKER_RESTARTS: + logger.error( + "EnvServerPool exhausted its %d worker restart(s) — not replacing", + MAX_WORKER_RESTARTS, + ) + return + self._restarts += 1 + self._spawn_worker() + logger.warning( + "EnvServerPool replaced a dead worker (restart %d/%d)", + self._restarts, + MAX_WORKER_RESTARTS, + ) + + async def _reap_dead_workers(self) -> None: + """Retire exited workers, fail their in-flight requests, and replace them. + + Rate-limited: reading `exitcode` is a waitpid per worker, and the broker loop turns + once per request.""" + now = time.monotonic() + if now - self._last_reap < _REAP_INTERVAL_S: + return + self._last_reap = now + for w in [w for w in self.workers if w["process"].exitcode is not None]: + exitcode = w["process"].exitcode + self._retire(w) + orphaned = [ + rid for rid, entry in self.pending.items() if entry["worker"] is w + ] + logger.error( + "EnvServerPool worker %d died (exit code %s) — failing %d in-flight request(s)", + w["index"], + exitcode, + len(orphaned), + ) + for request_id in orphaned: + entry = self.pending.pop(request_id) + self.in_flight -= entry["rollout_slots"] + await self._reply_error( + entry["client_id"], + request_id, + f"env server worker died (exit code {exitcode})", + ) + self._replace_worker() + + async def _on_request( + self, client_id: bytes, request_id: bytes, method: bytes, payload: bytes + ) -> None: + """Answer `health` inline; dispatch everything else to the least-busy worker.""" + if method == b"health": + await self.frontend.send_multipart([client_id, request_id, _HEALTH]) + return + if not self.workers: + await self._reply_error( + client_id, request_id, "env server has no live workers" + ) + return + # Pool capacity is measured in rollouts; one group request carries n. + rollout_slots = 1 + if method == b"run_group": + with contextlib.suppress(Exception): + request = RunGroupRequest.model_validate( + msgpack.unpackb(payload, raw=False) + ) + rollout_slots = max(1, request.n) + worker = min(self.workers, key=lambda w: w["active"]) + worker["active"] += rollout_slots + self.pending[request_id] = { + "client_id": client_id, + "worker": worker, + "rollout_slots": rollout_slots, + } + self.in_flight += rollout_slots + # forward without client_id — the DEALER identity is the worker's `client_id`; + # we hold the real one in `pending`. + await worker["dealer"].send_multipart([request_id, method, payload]) + if self.elastic: + self._maybe_scale_up() + + async def _on_reply(self, w: dict) -> None: + """Relay one worker reply back to the client that asked for it.""" + request_id, data = await w["dealer"].recv_multipart(copy=False) + # Copy only the routing key; relay the response Frames unchanged. + entry = self.pending.pop(request_id.bytes, None) + if entry is None: + return + entry["worker"]["active"] -= entry["rollout_slots"] + self.in_flight -= entry["rollout_slots"] + with contextlib.suppress(zmq.ZMQError): + await self.frontend.send_multipart( + [entry["client_id"], request_id, data], copy=False + ) + async def run(self) -> None: self._poller = zmq.asyncio.Poller() self._poller.register(self.frontend, zmq.POLLIN) @@ -173,68 +344,32 @@ async def run(self) -> None: # (`max_workers` is a concrete count when elastic is off). for _ in range(1 if self.elastic else (self.max_workers or 1)): self._spawn_worker() - # request_id -> {client_id, worker, rollout_slots} - pending: dict[bytes, dict] = {} - logger.info( - "EnvServerPool up: address=%s workers=%d/%s multiplex=%d elastic=%s", - self.address, - len(self.workers), - self._cap_str, - self.multiplex, - self.elastic, - ) try: - in_flight = 0 + # Report the address only once a worker can actually serve. The frontend binds + # before any env is loaded, so a spawner told "ready" at bind time would connect + # happily to a pool whose workers are still loading — or already dead. + for w in list(self.workers): + await self._wait_for_worker(w) + if self.address_queue is not None: + self.address_queue.put(self.address) + logger.info( + "EnvServerPool up: address=%s workers=%d/%s multiplex=%d elastic=%s", + self.address, + len(self.workers), + self._cap_str, + self.multiplex, + self.elastic, + ) while True: - events = dict(await self._poller.poll()) - if self.frontend in events: - ( - client_id, - request_id, - method, - payload, - ) = await self.frontend.recv_multipart() - if method == b"health": - await self.frontend.send_multipart( - [client_id, request_id, _HEALTH] - ) - else: - # Pool capacity is measured in rollouts; one group request carries n. - rollout_slots = 1 - if method == b"run_group": - with contextlib.suppress(Exception): - request = RunGroupRequest.model_validate( - msgpack.unpackb(payload, raw=False) - ) - rollout_slots = max(1, request.n) - worker = min(self.workers, key=lambda w: w["active"]) - worker["active"] += rollout_slots - pending[request_id] = { - "client_id": client_id, - "worker": worker, - "rollout_slots": rollout_slots, - } - in_flight += rollout_slots - # forward without client_id — the DEALER identity is the worker's - # `client_id`; we hold the real one in `pending`. - await worker["dealer"].send_multipart( - [request_id, method, payload] - ) - if self.elastic: - self._maybe_scale_up(in_flight) - for w in self.workers: + events = dict(await self._poller.poll(timeout=_REAP_INTERVAL_MS)) + # Drain replies before reaping, so a worker that answered and then died + # still has its reply relayed instead of failed. + for w in list(self.workers): if w["dealer"] in events: - request_id, data = await w["dealer"].recv_multipart(copy=False) - # Copy only the routing key; relay the response Frames unchanged. - entry = pending.pop(request_id.bytes, None) - if entry is None: - continue - entry["worker"]["active"] -= entry["rollout_slots"] - in_flight -= entry["rollout_slots"] - with contextlib.suppress(zmq.ZMQError): - await self.frontend.send_multipart( - [entry["client_id"], request_id, data], copy=False - ) + await self._on_reply(w) + await self._reap_dead_workers() + if self.frontend in events: + await self._on_request(*await self.frontend.recv_multipart()) except (asyncio.CancelledError, KeyboardInterrupt): pass finally: @@ -283,8 +418,9 @@ def serve_env( """Serve one env over ZMQ: a single in-process `EnvServer` when `max_workers <= 1`, else an `EnvServerPool` broker over up to `max_workers` worker processes (`None` = unbounded). The frontend speaks the same protocol either way, so the client is - identical. Reports the bound address on `address_queue` (for a spawner that passed an - OS-assigned `:0`). + identical. Reports the served address on `address_queue` (for a spawner that passed an + OS-assigned `:0`) once the env is loaded and serving — read it with `wait_for_address`, + which fails fast if this process dies first. `elastic` (default True) starts the pool at one worker and scales up to `max_workers` as load grows; `multiplex` is the per-worker capacity for the scale-up trigger (spawn @@ -323,9 +459,10 @@ def serve_env( log_setup, multiplex, elastic, + # The pool reports its own address, once a worker answers `health` — see + # EnvServerPool.run. + address_queue=address_queue, ) - if address_queue is not None: - address_queue.put(pool.address) asyncio.run(pool.run()) else: from verifiers.v1.legacy import LegacyEnvServer @@ -347,3 +484,46 @@ def serve_env( except Exception: logger.exception("Env server failed") raise + + +def _log_tail(log_file: str | Path | None, lines: int = 20) -> str: + """The tail of a spawned server's log, to quote in a failure message. A spawner that + redirects the child's output has the real traceback there, so this is what turns "the + server died" into the actual cause.""" + if log_file is None: + return "" + with contextlib.suppress(OSError): + tail = Path(log_file).read_text(errors="replace").splitlines()[-lines:] + if tail: + return "\n".join(["", f"--- tail of {log_file} ---", *tail]) + return "" + + +async def wait_for_address( + process: BaseProcess, + address_queue, + *, + name: str = "env server", + timeout: float = ENV_SERVER_SPAWN_TIMEOUT, + log_file: str | Path | None = None, +) -> str: + """Wait for a spawned `serve_env` process to report its address, failing fast if it + dies first. + + The address arrives only once the env is loaded and serving, so a server that dies on + startup (a bad import, missing credentials, an OOM) never reports one. Waiting on the + queue alone would then block for the full `timeout` and blame a timeout for what was a + crash — so poll the exit code alongside it and raise as soon as the process is gone.""" + deadline = time.monotonic() + timeout + while True: + with contextlib.suppress(queue.Empty): + return await asyncio.to_thread(address_queue.get, timeout=1.0) + if process.exitcode is not None: + raise RuntimeError( + f"{name} died on startup with exit code {process.exitcode}" + f"{_log_tail(log_file)}" + ) + if time.monotonic() >= deadline: + raise TimeoutError( + f"{name} did not report its address within {timeout}s{_log_tail(log_file)}" + ) From 25bd48ddd0b365ee28cf863a9941fe0d95542507 Mon Sep 17 00:00:00 2001 From: Eli Gottlieb Date: Fri, 24 Jul 2026 13:59:39 -0700 Subject: [PATCH 2/6] fix(v1): keep the broker alive on spawn failure, drain replies before reaping MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two problems in the death handling itself, from review: - `_replace_worker` called `_spawn_worker` unguarded. A worker death caused by memory pressure is exactly when spawning its replacement can fail, and that exception escaped the serve loop — killing the broker and stranding every other client's in-flight requests. Mid-serve spawns now go through `_try_spawn_worker`, which logs and leaves the pool one worker short. Scale-up had the same latent bug and takes the same path. Startup spawns still propagate: a pool with no workers has nothing to serve. - Reaping closed a dead worker's socket after relaying only one reply, so a worker that answered several requests and then exited had the rest discarded and reported to clients as `worker died`. The reaper now drains the socket before retiring it. The serve loop still takes one reply per pass, so a chatty worker can't starve the frontend. Co-Authored-By: Claude Opus 5 (1M context) --- verifiers/v1/serve/pool.py | 43 ++++++++++++++++++++++++++++++-------- 1 file changed, 34 insertions(+), 9 deletions(-) diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index 12434bf0d..087197260 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -180,6 +180,18 @@ def _spawn_worker(self) -> None: if self._poller is not None: self._poller.register(dealer, zmq.POLLIN) + def _try_spawn_worker(self) -> bool: + """Spawn a worker mid-serve, where failing to spawn must not be fatal: this process + holds every client's in-flight requests, so dying here strands all of them. A pool + one worker short still serves. (At startup the opposite holds — a pool with no + workers has nothing to serve — so `run` lets that failure propagate.)""" + try: + self._spawn_worker() + return True + except Exception: + logger.exception("EnvServerPool failed to spawn a worker") + return False + def _maybe_scale_up(self) -> None: """Spawn one more worker when in-flight rollout slots reach 90% of capacity. @@ -188,8 +200,10 @@ def _maybe_scale_up(self) -> None: up once already saturated. `max_workers=None` scales without a cap.""" if self.max_workers is not None and len(self.workers) >= self.max_workers: return - if self.in_flight >= 0.9 * len(self.workers) * self.multiplex: - self._spawn_worker() + if ( + self.in_flight >= 0.9 * len(self.workers) * self.multiplex + and self._try_spawn_worker() + ): logger.info( "EnvServerPool scaled up to %d/%s workers (in_flight=%d)", len(self.workers), @@ -250,13 +264,13 @@ def _replace_worker(self) -> None: MAX_WORKER_RESTARTS, ) return - self._restarts += 1 - self._spawn_worker() - logger.warning( - "EnvServerPool replaced a dead worker (restart %d/%d)", - self._restarts, - MAX_WORKER_RESTARTS, - ) + self._restarts += 1 # counts attempts, so a failing spawn can't loop hot + if self._try_spawn_worker(): + logger.warning( + "EnvServerPool replaced a dead worker (restart %d/%d)", + self._restarts, + MAX_WORKER_RESTARTS, + ) async def _reap_dead_workers(self) -> None: """Retire exited workers, fail their in-flight requests, and replace them. @@ -269,6 +283,10 @@ async def _reap_dead_workers(self) -> None: self._last_reap = now for w in [w for w in self.workers if w["process"].exitcode is not None]: exitcode = w["process"].exitcode + # Relay what it already answered before closing its socket: a worker can queue + # several replies and then exit, and retiring it first would discard them — + # reporting finished rollouts as deaths. + await self._drain_replies(w) self._retire(w) orphaned = [ rid for rid, entry in self.pending.items() if entry["worker"] is w @@ -323,6 +341,13 @@ async def _on_request( if self.elastic: self._maybe_scale_up() + async def _drain_replies(self, w: dict) -> None: + """Relay every reply already queued from this worker, then return. The serve loop + takes one reply per pass (so a chatty worker can't starve the frontend); a worker + about to be retired instead has to be emptied, since its socket is closing.""" + while await w["dealer"].poll(timeout=0): + await self._on_reply(w) + async def _on_reply(self, w: dict) -> None: """Relay one worker reply back to the client that asked for it.""" request_id, data = await w["dealer"].recv_multipart(copy=False) From f0c87f2e997213f2e82c08bac79093ed95d5b04f Mon Sep 17 00:00:00 2001 From: Eli Gottlieb Date: Fri, 24 Jul 2026 14:21:34 -0700 Subject: [PATCH 3/6] fix(v1): retry the replacement spawn when the pool is left empty MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A replacement that failed to spawn left the pool with no workers, and an empty pool has no dead worker left to drive another attempt — so `_reap_dead_workers` never tried again and every later request was refused with `no live workers`, with the restart budget still unspent. One transient failure to fork permanently disabled the pool. The reaper now retries while budget remains, bounded by its own interval and by MAX_WORKER_RESTARTS. Once the budget is spent the pool stays empty and fails requests fast, which is the honest end state: a host that cannot fork is not something the broker should paper over. Co-Authored-By: Claude Opus 5 (1M context) --- verifiers/v1/serve/pool.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index 087197260..d839d7975 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -306,6 +306,14 @@ async def _reap_dead_workers(self) -> None: f"env server worker died (exit code {exitcode})", ) self._replace_worker() + # A replacement that failed to spawn leaves the pool empty, and an empty pool has no + # dead worker left to drive another attempt — so retry here while budget remains, + # rather than refusing every request forever over one transient failure. Bounded by + # this method's own interval and by MAX_WORKER_RESTARTS; once that is spent the pool + # stays empty and fails requests fast, which is the honest end state (a host that + # cannot fork is not something the broker should paper over). + if not self.workers: + self._replace_worker() async def _on_request( self, client_id: bytes, request_id: bytes, method: bytes, payload: bytes From 6b4ad69ad62725669f372ede3d0520636590a6d8 Mon Sep 17 00:00:00 2001 From: eligotts <78387377+eligotts@users.noreply.github.com> Date: Tue, 28 Jul 2026 23:21:11 +0000 Subject: [PATCH 4/6] fix(v1): gate replacement workers on readiness --- verifiers/v1/serve/pool.py | 52 ++++++++++++++++++++++++++++++++------ 1 file changed, 44 insertions(+), 8 deletions(-) diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index ad63980b8..ee339234a 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -54,9 +54,9 @@ _HEALTH = msgpack.packb(HealthResponse().model_dump(mode="json"), use_bin_type=True) MAX_WORKER_RESTARTS = 3 -"""Replacements the broker spawns over its lifetime: enough to ride out a one-off crash -(an OOM, a segfaulting sandbox), few enough that an env dying on every load fails loudly -instead of respawning forever.""" +"""Consecutive recovery attempts before the broker stops spawning workers. The budget +resets when a worker becomes ready: transient crashes do not permanently reduce capacity, +while an env that repeatedly dies during startup still fails loudly.""" ENV_SERVER_SPAWN_TIMEOUT = 600.0 """Default deadline for a spawned server to report its address. Generous: it reports only @@ -175,6 +175,8 @@ def _spawn_worker(self) -> None: "pipe": parent_conn, "active": 0, "index": i, + "ready": False, + "probe_sent": False, } ) if self._poller is not None: @@ -195,9 +197,8 @@ def _try_spawn_worker(self) -> bool: def _maybe_scale_up(self) -> None: """Spawn one more worker when in-flight rollout slots reach 90% of capacity. - A new worker starts at `active=0`, so least-busy dispatch funnels the backlog to - it as it comes online (a few seconds to load the env) — fine, since we only scale - up once already saturated. `max_workers=None` scales without a cap.""" + A new worker receives traffic only after its health probe succeeds. + `max_workers=None` scales without a cap.""" if self._restarts >= MAX_WORKER_RESTARTS: return if self.max_workers is not None and len(self.workers) >= self.max_workers: @@ -225,7 +226,7 @@ async def _wait_for_worker(self, w: dict) -> None: dataset pull) is the spawner's deadline to enforce, not ours. We only insist that the process is still alive.""" dealer = w["dealer"] - await dealer.send_multipart([_STARTUP_PROBE, b"health", _EMPTY_PAYLOAD]) + await self._probe_worker(w) while not await dealer.poll(timeout=_REAP_INTERVAL_MS): if w["process"].exitcode is not None: raise RuntimeError( @@ -233,6 +234,30 @@ async def _wait_for_worker(self, w: dict) -> None: f"(exit code {w['process'].exitcode})" ) await dealer.recv_multipart() + if w["process"].exitcode is not None: + raise RuntimeError( + f"env server worker {w['index']} died while loading the env " + f"(exit code {w['process'].exitcode})" + ) + self._mark_worker_ready(w) + + async def _probe_worker(self, w: dict) -> None: + if w["probe_sent"]: + return + await w["dealer"].send_multipart([_STARTUP_PROBE, b"health", _EMPTY_PAYLOAD]) + w["probe_sent"] = True + + async def _probe_new_workers(self) -> None: + for w in self.workers: + if not w["ready"]: + await self._probe_worker(w) + + def _mark_worker_ready(self, w: dict) -> None: + if w["ready"]: + return + w["ready"] = True + self._restarts = 0 + logger.info("EnvServerPool worker %d ready", w["index"]) async def _reply_error( self, client_id: bytes, request_id: bytes, error: str @@ -331,6 +356,12 @@ async def _on_request( client_id, request_id, "env server has no live workers" ) return + ready_workers = [w for w in self.workers if w["ready"]] + if not ready_workers: + await self._reply_error( + client_id, request_id, "env server has no ready workers" + ) + return # Pool capacity is measured in rollouts; one group request carries n. rollout_slots = 1 if method == b"run_group": @@ -339,7 +370,7 @@ async def _on_request( msgpack.unpackb(payload, raw=False) ) rollout_slots = max(1, request.n) - worker = min(self.workers, key=lambda w: w["active"]) + worker = min(ready_workers, key=lambda w: w["active"]) worker["active"] += rollout_slots self.pending[request_id] = { "client_id": client_id, @@ -363,6 +394,10 @@ async def _drain_replies(self, w: dict) -> None: async def _on_reply(self, w: dict) -> None: """Relay one worker reply back to the client that asked for it.""" request_id, data = await w["dealer"].recv_multipart(copy=False) + if request_id.bytes == _STARTUP_PROBE: + if w["process"].exitcode is None: + self._mark_worker_ready(w) + return # Copy only the routing key; relay the response Frames unchanged. entry = self.pending.pop(request_id.bytes, None) if entry is None: @@ -398,6 +433,7 @@ async def run(self) -> None: self.elastic, ) while True: + await self._probe_new_workers() events = dict(await self._poller.poll(timeout=_REAP_INTERVAL_MS)) # Drain replies before reaping, so a worker that answered and then died # still has its reply relayed instead of failed. From ca9fbbf76e62ebc434098768187b8aa295379b69 Mon Sep 17 00:00:00 2001 From: eligotts <78387377+eligotts@users.noreply.github.com> Date: Tue, 28 Jul 2026 23:28:33 +0000 Subject: [PATCH 5/6] fix(v1): bound failed pool dispatch --- verifiers/v1/serve/pool.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index ee339234a..e4ba117e9 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -203,10 +203,10 @@ def _maybe_scale_up(self) -> None: return if self.max_workers is not None and len(self.workers) >= self.max_workers: return - if ( - self.in_flight >= 0.9 * len(self.workers) * self.multiplex - and self._try_spawn_worker() - ): + if self.in_flight >= 0.9 * len(self.workers) * self.multiplex: + if not self._try_spawn_worker(): + self._restarts += 1 + return logger.info( "EnvServerPool scaled up to %d/%s workers (in_flight=%d)", len(self.workers), @@ -356,7 +356,9 @@ async def _on_request( client_id, request_id, "env server has no live workers" ) return - ready_workers = [w for w in self.workers if w["ready"]] + ready_workers = [ + w for w in self.workers if w["ready"] and w["process"].exitcode is None + ] if not ready_workers: await self._reply_error( client_id, request_id, "env server has no ready workers" From 40154d6b4a28293f9247a02f1c1a1d384915c54e Mon Sep 17 00:00:00 2001 From: eligotts <78387377+eligotts@users.noreply.github.com> Date: Tue, 28 Jul 2026 23:36:36 +0000 Subject: [PATCH 6/6] fix(v1): isolate worker recovery budgets --- verifiers/v1/serve/pool.py | 96 ++++++++++++++++++++++---------------- 1 file changed, 56 insertions(+), 40 deletions(-) diff --git a/verifiers/v1/serve/pool.py b/verifiers/v1/serve/pool.py index e4ba117e9..032ec3d91 100644 --- a/verifiers/v1/serve/pool.py +++ b/verifiers/v1/serve/pool.py @@ -22,9 +22,9 @@ unanswered request is an unbounded hang, and a dead worker's `active` count never drains — least-busy dispatch would elect the corpse for every subsequent request. The broker therefore polls its workers' exit codes, answers their in-flight requests with an error, -and replaces them (up to `MAX_WORKER_RESTARTS`). Readiness is a worker handshake, not a -socket bind: the broker reports its address only once a worker answers `health`, so a -spawner is never told "ready" about a pool whose env failed to load. +and replaces each failed worker up to `MAX_WORKER_RESTARTS` times. Readiness is a worker +handshake, not a socket bind: the broker reports its address only once a worker answers +`health`, so a spawner is never told "ready" about a pool whose env failed to load. """ import asyncio @@ -54,9 +54,8 @@ _HEALTH = msgpack.packb(HealthResponse().model_dump(mode="json"), use_bin_type=True) MAX_WORKER_RESTARTS = 3 -"""Consecutive recovery attempts before the broker stops spawning workers. The budget -resets when a worker becomes ready: transient crashes do not permanently reduce capacity, -while an env that repeatedly dies during startup still fails loudly.""" +"""Consecutive replacement attempts per worker slot. A ready worker resets its slot's +budget; repeated startup failures exhaust only that slot, not the rest of the pool.""" ENV_SERVER_SPAWN_TIMEOUT = 600.0 """Default deadline for a spawned server to report its address. Generous: it reports only @@ -65,6 +64,7 @@ # Broker loop wake-up, so exited workers are reaped even while no socket is readable. _REAP_INTERVAL_S = 0.25 _REAP_INTERVAL_MS = int(_REAP_INTERVAL_S * 1000) +_MAX_SCALE_UP_BACKOFF_S = 30.0 # Startup handshake: a `health` request whose reply proves a worker loaded its env. _STARTUP_PROBE = b"startup" _EMPTY_PAYLOAD = msgpack.packb({}, use_bin_type=True) @@ -131,7 +131,9 @@ def __init__( self._mpctx = mp.get_context("spawn") self._poller: zmq.asyncio.Poller | None = None self._next_index = 0 - self._restarts = 0 + self._replacement_queue: list[int] = [] + self._scale_up_failures = 0 + self._scale_up_retry_at = 0.0 self._last_reap = 0.0 self.ctx = zmq.asyncio.Context() @@ -146,7 +148,7 @@ def __init__( def _worker_path(self, i: int) -> str: return f"/tmp/vf-pool-{self.session}-{i}" - def _spawn_worker(self) -> None: + def _spawn_worker(self, restarts: int = 0) -> None: i = self._next_index # monotonic: a retired worker's path is never reused self._next_index += 1 address = f"ipc://{self._worker_path(i)}" @@ -177,18 +179,19 @@ def _spawn_worker(self) -> None: "index": i, "ready": False, "probe_sent": False, + "restarts": restarts, } ) if self._poller is not None: self._poller.register(dealer, zmq.POLLIN) - def _try_spawn_worker(self) -> bool: + def _try_spawn_worker(self, restarts: int = 0) -> bool: """Spawn a worker mid-serve, where failing to spawn must not be fatal: this process holds every client's in-flight requests, so dying here strands all of them. A pool one worker short still serves. (At startup the opposite holds — a pool with no workers has nothing to serve — so `run` lets that failure propagate.)""" try: - self._spawn_worker() + self._spawn_worker(restarts) return True except Exception: logger.exception("EnvServerPool failed to spawn a worker") @@ -199,14 +202,22 @@ def _maybe_scale_up(self) -> None: A new worker receives traffic only after its health probe succeeds. `max_workers=None` scales without a cap.""" - if self._restarts >= MAX_WORKER_RESTARTS: + now = time.monotonic() + if now < self._scale_up_retry_at: return - if self.max_workers is not None and len(self.workers) >= self.max_workers: + slots = len(self.workers) + len(self._replacement_queue) + if self.max_workers is not None and slots >= self.max_workers: return - if self.in_flight >= 0.9 * len(self.workers) * self.multiplex: + if self.in_flight >= 0.9 * slots * self.multiplex: if not self._try_spawn_worker(): - self._restarts += 1 + self._scale_up_failures += 1 + self._scale_up_retry_at = now + min( + 2 ** min(self._scale_up_failures - 1, 5), + _MAX_SCALE_UP_BACKOFF_S, + ) return + self._scale_up_failures = 0 + self._scale_up_retry_at = 0.0 logger.info( "EnvServerPool scaled up to %d/%s workers (in_flight=%d)", len(self.workers), @@ -256,7 +267,7 @@ def _mark_worker_ready(self, w: dict) -> None: if w["ready"]: return w["ready"] = True - self._restarts = 0 + w["restarts"] = 0 logger.info("EnvServerPool worker %d ready", w["index"]) async def _reply_error( @@ -283,21 +294,28 @@ def _retire(self, w: dict) -> None: with contextlib.suppress(OSError): os.unlink(self._worker_path(w["index"])) - def _replace_worker(self) -> None: - """Spawn a replacement for a dead worker, within the restart budget.""" - if self._restarts >= MAX_WORKER_RESTARTS: - logger.error( - "EnvServerPool exhausted its %d worker restart(s) — not replacing", - MAX_WORKER_RESTARTS, - ) - return - self._restarts += 1 # counts attempts, so a failing spawn can't loop hot - if self._try_spawn_worker(): - logger.warning( - "EnvServerPool replaced a dead worker (restart %d/%d)", - self._restarts, - MAX_WORKER_RESTARTS, - ) + def _retry_replacements(self) -> None: + """Try each missing worker slot once, preserving independent restart budgets.""" + missing = self._replacement_queue + self._replacement_queue = [] + for restarts in missing: + if restarts >= MAX_WORKER_RESTARTS: + self._replacement_queue.append(restarts) + continue + attempt = restarts + 1 + if self._try_spawn_worker(attempt): + logger.warning( + "EnvServerPool replaced a dead worker (restart %d/%d)", + attempt, + MAX_WORKER_RESTARTS, + ) + else: + self._replacement_queue.append(attempt) + if attempt == MAX_WORKER_RESTARTS: + logger.error( + "EnvServerPool exhausted a worker slot's %d restart(s)", + MAX_WORKER_RESTARTS, + ) async def _reap_dead_workers(self) -> None: """Retire exited workers, fail their in-flight requests, and replace them. @@ -333,16 +351,14 @@ async def _reap_dead_workers(self) -> None: request_id, f"env server worker died (exit code {exitcode})", ) - for _ in dead_workers: - self._replace_worker() - # A replacement that failed to spawn leaves the pool empty, and an empty pool has no - # dead worker left to drive another attempt — so retry here while budget remains, - # rather than refusing every request forever over one transient failure. Bounded by - # this method's own interval and by MAX_WORKER_RESTARTS; once that is spent the pool - # stays empty and fails requests fast, which is the honest end state (a host that - # cannot fork is not something the broker should paper over). - if not dead_workers and not self.workers: - self._replace_worker() + self._replacement_queue.append(w["restarts"]) + if w["restarts"] >= MAX_WORKER_RESTARTS: + logger.error( + "EnvServerPool exhausted worker %d's %d restart(s) — not replacing", + w["index"], + MAX_WORKER_RESTARTS, + ) + self._retry_replacements() async def _on_request( self, client_id: bytes, request_id: bytes, method: bytes, payload: bytes