From ab6c5c63d949995ad97152df5c0d53faf63f85b4 Mon Sep 17 00:00:00 2001 From: nilsmechtel Date: Sun, 27 Sep 2026 01:55:36 +0200 Subject: [PATCH 1/3] fix(apps): re-assert the registration record from the replica that is serving MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A successor proxy replica claims the record as unregistered at init. If it then never registers — a redeploy where the predecessor still holds the shared Hypha client_id, or a crash-loop — nothing put the record back: a registered replica takes the probe branch of the maintenance tick and never re-reported True. The worker withheld a working service address indefinitely. The reachability probe already establishes that our address resolves, so report True there. The record becomes self-correcting instead of write-once. Co-Authored-By: Claude Opus 5 (1M context) --- bioengine/apps/proxy_deployment.py | 6 ++ bioengine/cluster/proxy_actor.py | 13 ++- .../apps/test_service_id_registration_gate.py | 101 ++++++++++++++++++ 3 files changed, 115 insertions(+), 5 deletions(-) diff --git a/bioengine/apps/proxy_deployment.py b/bioengine/apps/proxy_deployment.py index e4c290de..6e0ac9a4 100644 --- a/bioengine/apps/proxy_deployment.py +++ b/bioengine/apps/proxy_deployment.py @@ -1481,6 +1481,12 @@ async def _maintenance_tick(self) -> None: # against a server that has dropped the registration. await self.server.get_service_info(self.websocket_service_id) self._probe_due_at = time.time() + _REACHABILITY_PROBE_INTERVAL_S + # Re-assert rather than assume the record still says what we last + # wrote: a successor's init-time claim replaces it with False, and + # if that successor never registers nothing else would ever put it + # back. Having just confirmed the address resolves, we are the + # replica entitled to say so. + self._report_service_registration(True) except Exception as e: # Not a health-check failure: flag for rebuild on the next tick. logger.warning( diff --git a/bioengine/cluster/proxy_actor.py b/bioengine/cluster/proxy_actor.py index ba119a85..ac77f125 100644 --- a/bioengine/cluster/proxy_actor.py +++ b/bioengine/cluster/proxy_actor.py @@ -923,11 +923,14 @@ def report_service_registration( belongs to. Across a replica generation the two reports race: the incoming replica claims the record, registers and reports True, and the outgoing one's ``__del__`` deregisters afterwards. Taking that last - write would pin a healthy app at False for good, because nothing - re-reports True until the next ``_register_services``, which the - running replica will not redo. A ``True``, and any report about an app - with no record, always takes over — those can only come from a replica - that is serving now. + write would leave a healthy app reading False until the serving + replica's next reachability probe re-asserts True. A ``True``, and any + report about an app with no record, always takes over — those can only + come from a replica that is serving now. + + That periodic re-assert, not this guard, is what makes the record + self-correcting: a claim from a successor that then never registers + would otherwise hold a serving app at False indefinitely. """ current = self.service_registrations.get(application_id) if ( diff --git a/tests/apps/test_service_id_registration_gate.py b/tests/apps/test_service_id_registration_gate.py index d1e2e4dc..df852651 100644 --- a/tests/apps/test_service_id_registration_gate.py +++ b/tests/apps/test_service_id_registration_gate.py @@ -19,6 +19,10 @@ actor existed — a newer worker, or an actor recreated after eviction — serves perfectly well and must keep its id, so the worker falls back to the replica-alive gate. +* The record is self-correcting, not write-once: a replica whose reachability + probe succeeds re-asserts ``True`` every probe interval. Otherwise a + successor that claims the record and then never registers withholds a + working address for as long as the predecessor keeps serving it. """ from __future__ import annotations @@ -460,6 +464,103 @@ def test_a_replica_that_died_without_deregistering_loses_the_record( assert actor.get_service_registration(APP_ID) is False +class _LiveServer: + """A Hypha connection whose reachability probe answers.""" + + def __init__(self) -> None: + self.probes = 0 + + async def get_service_info(self, _sid): + self.probes += 1 + + async def unregister_service(self, _sid): + return None + + +@pytest.mark.asyncio +async def test_a_successor_that_never_registers_cannot_pin_a_serving_app_at_false( + monkeypatch, +) -> None: + # Redeploy: Ray Serve builds the new proxy replica while the old one is + # still serving, and both want the same Hypha client_id, so the successor's + # registration is refused for as long as the predecessor holds it. Its + # init-time claim has already replaced the record with False. Nothing then + # re-reported True — a registered replica skips the registration branch of + # _maintenance_tick entirely — so the worker withheld a working address + # indefinitely, which is the mirror image of the bug the claim fixed. + actor = _bare_actor() + + old = _bare_proxy(_replica_id="old", _proxy_actor_handle=_ActorHandle(actor)) + server = _LiveServer() + + async def _register(): + old.server = server + old.websocket_service_id = "ws" + + old._register_services = _register + await old._maintenance_tick() + assert actor.get_service_registration(APP_ID) is True + + _construct_proxy(monkeypatch, _ActorHandle(actor), replica_tag="new") + # Asserted mid-way: without the claim actually landing there is no + # divergence left to heal, and the final assertion would pass vacuously. + assert actor.get_service_registration(APP_ID) is False + + # The predecessor is still serving, so its periodic probe is the only + # signal left that the app is reachable. + old._probe_due_at = 0.0 + await old._maintenance_tick() + + assert server.probes == 1 + assert actor.get_service_registration(APP_ID) is True + + +@pytest.mark.asyncio +async def test_the_probe_does_not_re_register_an_app_that_deregistered_itself() -> None: + # A sibling deployment going down closes the gate via _deregister_services. + # The self-healing re-assert must not undo that on the next tick: this + # replica reports False about itself and stays that way until it re-registers. + actor = _bare_actor() + inst = _bare_proxy( + _replica_id="r0", + _proxy_actor_handle=_ActorHandle(actor), + server=_LiveServer(), + websocket_service_id="ws", + ) + actor.report_service_registration(APP_ID, True, replica_id="r0") + + await inst._deregister_services() + inst._probe_due_at = 0.0 + await inst._maintenance_tick() + await inst._maintenance_tick() + + assert actor.get_service_registration(APP_ID) is False + + +@pytest.mark.asyncio +async def test_a_failing_probe_does_not_assert_the_app_as_registered() -> None: + # The probe is what establishes the address still resolves. When it fails + # the replica rebuilds its client instead, and must not claim reachability + # it has just failed to confirm. + class _AmnesiacServer: + async def get_service_info(self, sid): + raise KeyError(f"Service not found: {sid}@*") + + actor = _bare_actor() + inst = _bare_proxy( + _replica_id="r0", + _proxy_actor_handle=_ActorHandle(actor), + server=_AmnesiacServer(), + websocket_service_id="ws", + ) + actor.claim_service_registration(APP_ID, replica_id="newer") + + await inst._maintenance_tick() + + assert inst._connection_lost is True + assert actor.get_service_registration(APP_ID) is False + + def test_an_untagged_report_still_works() -> None: # A replica that could not read its Ray replica context reports None, which # must degrade to the plain last-writer behaviour rather than being dropped. From a8c6eae90491ad123a167b9d6364cb0e5bde970d Mon Sep 17 00:00:00 2001 From: nilsmechtel Date: Sun, 27 Sep 2026 02:47:40 +0200 Subject: [PATCH 2/3] fix(apps): do not re-assert registration for a replica that deregistered MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The probe suspends, so check_health's sibling-down path can deregister this replica while it is in flight. Reporting True on the answer to a question asked before that would leave a deregistered app advertised with no way back, since every later tick returns at the readiness gate. Also retires the stale rationale in the gate test that still claimed nothing re-reports True after a registration — the twin of the sentence already fixed in the proxy actor docstring. Co-Authored-By: Claude Opus 5 (1M context) --- bioengine/apps/proxy_deployment.py | 8 +++- .../apps/test_service_id_registration_gate.py | 42 +++++++++++++++++-- 2 files changed, 46 insertions(+), 4 deletions(-) diff --git a/bioengine/apps/proxy_deployment.py b/bioengine/apps/proxy_deployment.py index 6e0ac9a4..4ca44fec 100644 --- a/bioengine/apps/proxy_deployment.py +++ b/bioengine/apps/proxy_deployment.py @@ -1486,7 +1486,13 @@ async def _maintenance_tick(self) -> None: # if that successor never registers nothing else would ever put it # back. Having just confirmed the address resolves, we are the # replica entitled to say so. - self._report_service_registration(True) + # + # Re-read the gate: check_health can deregister us while the probe + # above is in flight, and answering a question asked before that + # would leave a deregistered app advertised with no path back — + # every later tick returns at the gate. + if self.entry_deployment_ready: + self._report_service_registration(True) except Exception as e: # Not a health-check failure: flag for rebuild on the next tick. logger.warning( diff --git a/tests/apps/test_service_id_registration_gate.py b/tests/apps/test_service_id_registration_gate.py index df852651..cfad3bc1 100644 --- a/tests/apps/test_service_id_registration_gate.py +++ b/tests/apps/test_service_id_registration_gate.py @@ -391,9 +391,10 @@ def test_reporting_survives_a_part_built_replica() -> None: def test_a_late_deregistration_from_the_old_replica_is_ignored() -> None: # Rolling update: the incoming replica claims the record, registers and # reports True; the outgoing replica's __del__ deregisters afterwards. - # Taking that last write pins a healthy app at False for good — nothing - # re-reports True until the next _register_services, which the running - # replica will not do. + # Taking that last write leaves a healthy app reading False until the + # serving replica's next reachability probe re-asserts True — a whole probe + # interval of a working address being withheld. The guard avoids the window + # rather than relying on the re-assert to close it. actor = _bare_actor() actor.report_service_registration(APP_ID, True, replica_id="old") actor.claim_service_registration(APP_ID, replica_id="new") @@ -561,6 +562,41 @@ async def get_service_info(self, sid): assert actor.get_service_registration(APP_ID) is False +@pytest.mark.asyncio +async def test_a_probe_answered_after_a_deregistration_does_not_re_register() -> None: + # check_health can deregister this replica while the probe is suspended. + # Reporting True on the answer to a question asked before that would leave + # a deregistered app advertised with no way back, since every later tick + # returns at the readiness gate. + entered = asyncio.Event() + release = asyncio.Event() + + class _SlowServer: + async def get_service_info(self, _sid): + entered.set() + await release.wait() + + async def unregister_service(self, _sid): + return None + + actor = _bare_actor() + inst = _bare_proxy( + _replica_id="r0", + _proxy_actor_handle=_ActorHandle(actor), + server=_SlowServer(), + websocket_service_id="ws", + ) + actor.report_service_registration(APP_ID, True, replica_id="r0") + + tick = asyncio.create_task(inst._maintenance_tick()) + await asyncio.wait_for(entered.wait(), timeout=5) + await inst._deregister_services() + release.set() + await asyncio.wait_for(tick, timeout=5) + + assert actor.get_service_registration(APP_ID) is False + + def test_an_untagged_report_still_works() -> None: # A replica that could not read its Ray replica context reports None, which # must degrade to the plain last-writer behaviour rather than being dropped. From eef74cfbd12b07ecd6db369ab79a8e3059c51083 Mon Sep 17 00:00:00 2001 From: Nils Mechtel <49943582+nilsmechtel@users.noreply.github.com> Date: Sun, 27 Sep 2026 03:05:17 +0200 Subject: [PATCH 3/3] chore(release): bump version to 0.16.26 --- bioengine/_version.py | 2 +- pyproject.toml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/bioengine/_version.py b/bioengine/_version.py index 0d261556..2c13dcad 100644 --- a/bioengine/_version.py +++ b/bioengine/_version.py @@ -13,4 +13,4 @@ Must stay in lock-step with ``pyproject.toml``'s ``version`` field. The ``version-check.yml`` CI workflow enforces the match. """ -__version__ = "0.16.25" +__version__ = "0.16.26" diff --git a/pyproject.toml b/pyproject.toml index d537a7dd..4170131e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "bioengine" -version = "0.16.25" +version = "0.16.26" description = "BioEngine — CLI and SDK for deploying and calling AI model services on BioEngine workers" requires-python = ">=3.11" authors = [