Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion bioengine/_version.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
12 changes: 12 additions & 0 deletions bioengine/apps/proxy_deployment.py
Original file line number Diff line number Diff line change
Expand Up @@ -1481,6 +1481,18 @@ 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.
#
# 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(
Expand Down
13 changes: 8 additions & 5 deletions bioengine/cluster/proxy_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 = [
Expand Down
143 changes: 140 additions & 3 deletions tests/apps/test_service_id_registration_gate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -387,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")
Expand Down Expand Up @@ -460,6 +465,138 @@ 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


@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.
Expand Down
Loading