Skip to content
Open
Show file tree
Hide file tree
Changes from 2 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
128 changes: 128 additions & 0 deletions src/tests/test_k8s_pod_ip_service_discovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -125,3 +125,131 @@ def test_modified_not_ready_with_ip_removes_registered():
)

assert "pod-c" not in d.available_engines


# ---------------------------------------------------------------------------
# Reconciliation on watch reconnect
#
# The watch stream is the only source of removals, so a DELETED event lost
# during a disconnect leaves a ghost engine. Listing before each stream repairs
# it, and the watch starts at the list's resourceVersion so nothing slips
# between the two calls. Since the watch no longer replays ADDED events, the
# reconciliation is also the startup discovery path.
# ---------------------------------------------------------------------------


def _make_pod(name, ip="10.0.0.1", ready=True, terminating=False, labels=None):
pod = MagicMock()
pod.metadata.name = name
pod.metadata.labels = {} if labels is None else labels
pod.metadata.deletion_timestamp = None if not terminating else "2026-01-01T00:00:00Z"
pod.status.pod_ip = ip
pod.status.container_statuses = [MagicMock(ready=ready)]
return pod


def _make_reconciler(pods, engines=None):
d = _make_discovery()
d.label_selector = "environment=test"
d.available_engines = engines or {}
d.k8s_api = MagicMock()
d.k8s_api.list_namespaced_pod.return_value = MagicMock(
items=pods, metadata=MagicMock(resource_version="4242")
)
d._get_model_names = MagicMock(return_value=["m"])
d._get_model_label = MagicMock(side_effect=lambda pod: (pod.metadata.labels or {}).get("model"))
d._add_engine = MagicMock()
return d


def _registered(url="http://10.0.0.1:8000", model_label=None, sleep=False):
known = MagicMock(spec=EndpointInfo)
known.url = url
known.model_label = model_label
known.sleep = sleep
return known


def test_reconcile_drops_engine_whose_pod_is_gone():
"""Core regression: the DELETED event was missed, the list must repair it."""
d = _make_reconciler(
pods=[_make_pod("pod-alive")],
engines={"pod-gone": _registered(url="http://10.0.0.9:8000"), "pod-alive": _registered()},
)

d._reconcile_engines()

assert "pod-gone" not in d.available_engines
assert "pod-alive" in d.available_engines


def test_reconcile_discovers_pods_at_startup():
"""The watch no longer replays ADDED, so the list must populate the table."""
d = _make_reconciler(pods=[_make_pod("pod-a", ip="10.0.0.1"), _make_pod("pod-b", ip="10.0.0.2")])

d._reconcile_engines()

assert {c.args[0] for c in d._add_engine.call_args_list} == {"pod-a", "pod-b"}


def test_reconcile_skips_unchanged_pod():
"""A ready pod already registered at the same URL must not be re-queried."""
d = _make_reconciler(pods=[_make_pod("pod-a")], engines={"pod-a": _registered()})

d._reconcile_engines()

d._get_model_names.assert_not_called()
d._add_engine.assert_not_called()


def test_reconcile_drops_pod_that_became_not_ready():
"""A readiness change missed during the disconnect must still be applied."""
d = _make_reconciler(
pods=[_make_pod("pod-a", ready=False)], engines={"pod-a": _registered()}
)

d._reconcile_engines()

assert "pod-a" not in d.available_engines


def test_reconcile_drops_pod_with_cleared_ip():
"""An evicted pod keeps its object but loses its IP."""
d = _make_reconciler(pods=[_make_pod("pod-a", ip=None)], engines={"pod-a": _registered()})

d._reconcile_engines()

assert "pod-a" not in d.available_engines


def test_reconcile_rehandles_pod_whose_sleep_label_changed():
"""/sleep patches the pod; missing that event would keep routing to it."""
d = _make_reconciler(
pods=[_make_pod("pod-a", labels={"sleeping": "true"})],
engines={"pod-a": _registered(sleep=False)},
)

d._reconcile_engines()

d._add_engine.assert_called_once()


def test_reconcile_returns_the_list_resource_version():
"""The watch must start where the list stopped, leaving no gap."""
d = _make_reconciler(pods=[_make_pod("pod-a")], engines={"pod-a": _registered()})

assert d._reconcile_engines() == "4242"


def test_reconcile_survives_one_broken_pod():
"""One unreachable pod must not abort discovery for the others."""
d = _make_reconciler(pods=[_make_pod("pod-bad", ip="10.0.0.1"), _make_pod("pod-ok", ip="10.0.0.2")])
d._get_model_names = MagicMock(
side_effect=lambda ip: (_ for _ in ()).throw(RuntimeError("unreachable"))
if ip == "10.0.0.1"
else ["m"]
)

d._reconcile_engines()

assert [c.args[0] for c in d._add_engine.call_args_list] == ["pod-ok"]
97 changes: 97 additions & 0 deletions src/tests/test_k8s_service_name_service_discovery.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
"""Unit tests for K8sServiceNameServiceDiscovery reconciliation.

Same failure mode as the pod-IP class: the watch stream is the only source of
removals, so a DELETED event lost during a disconnect leaves a ghost engine.
Listing before each stream repairs it and provides the resourceVersion the
watch starts from.

The real __init__ loads a kubeconfig and starts a watcher thread, so we build
the instance with __new__ and inject only what the tested code touches.
"""

import threading
from unittest.mock import MagicMock

from vllm_router.service_discovery import EndpointInfo, K8sServiceNameServiceDiscovery


def _make_service(name):
svc = MagicMock()
svc.metadata.name = name
svc.metadata.labels = {}
return svc


def _make_reconciler(services, engines=None, ready=True):
d = K8sServiceNameServiceDiscovery.__new__(K8sServiceNameServiceDiscovery)
d.available_engines = engines or {}
d.available_engines_lock = threading.Lock()
d.known_models = set()
d.known_models_lock = threading.Lock()
d.namespace = "test-ns"
d.port = "8000"
d.label_selector = "environment=test"
d.k8s_api = MagicMock()
d.k8s_api.list_namespaced_service.return_value = MagicMock(
items=services, metadata=MagicMock(resource_version="99")
)
d._check_service_ready = MagicMock(return_value=ready)
d._get_model_names = MagicMock(return_value=["m"])
d._get_model_label = MagicMock(return_value=None)
d._add_engine = MagicMock()
return d


def _registered():
known = MagicMock(spec=EndpointInfo)
known.url = "http://svc-a:8000"
return known


def test_reconcile_drops_engine_whose_service_is_gone():
d = _make_reconciler(
services=[_make_service("svc-a")],
engines={"svc-gone": _registered(), "svc-a": _registered()},
)

d._reconcile_engines()

assert "svc-gone" not in d.available_engines
assert "svc-a" in d.available_engines


def test_reconcile_discovers_services_at_startup():
"""The watch no longer replays ADDED, so the list must populate the table."""
d = _make_reconciler(services=[_make_service("svc-a"), _make_service("svc-b")])

d._reconcile_engines()

assert {c.args[0] for c in d._add_engine.call_args_list} == {"svc-a", "svc-b"}


def test_reconcile_skips_ready_registered_service():
d = _make_reconciler(services=[_make_service("svc-a")], engines={"svc-a": _registered()})

d._reconcile_engines()

d._get_model_names.assert_not_called()


def test_reconcile_survives_one_broken_service():
"""A Service with no Endpoints object must not abort the whole pass."""
d = _make_reconciler(services=[_make_service("svc-bad"), _make_service("svc-ok")])
d._check_service_ready = MagicMock(
side_effect=lambda name, ns: (_ for _ in ()).throw(RuntimeError("404"))
if name == "svc-bad"
else True
)

d._reconcile_engines()

assert [c.args[0] for c in d._add_engine.call_args_list] == ["svc-ok"]


def test_reconcile_returns_the_list_resource_version():
d = _make_reconciler(services=[], engines={})

assert d._reconcile_engines() == "99"
Loading
Loading