Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
44 changes: 44 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,47 @@ def test_modified_not_ready_with_ip_removes_registered():
)

assert "pod-c" not in d.available_engines


def test_reconcile_drops_engine_whose_pod_is_gone():
"""A DELETED event missed while the watch was down leaves a ghost.

The watch stream is the only source of removals, so an engine whose pod
disappeared during a disconnect is never dropped and keeps taking
traffic. Listing on reconnection must remove it.
"""
d = _make_discovery()
d.label_selector = "environment=test"
_register(d, "pod-gone")
_register(d, "pod-alive")

still_running = MagicMock()
still_running.metadata.name = "pod-alive"
d.k8s_api = MagicMock()
d.k8s_api.list_namespaced_pod.return_value = MagicMock(items=[still_running])

d._reconcile_engines()

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


def test_reconcile_keeps_every_live_engine():
"""The list must never drop an engine whose pod is still there."""
d = _make_discovery()
d.label_selector = "environment=test"
_register(d, "pod-a")
_register(d, "pod-b")

pods = []
for name in ("pod-a", "pod-b"):
pod = MagicMock()
pod.metadata.name = name
pods.append(pod)

d.k8s_api = MagicMock()
d.k8s_api.list_namespaced_pod.return_value = MagicMock(items=pods)

d._reconcile_engines()

assert set(d.available_engines) == {"pod-a", "pod-b"}
26 changes: 26 additions & 0 deletions src/vllm_router/service_discovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -668,9 +668,35 @@ def _get_model_label(self, pod) -> Optional[str]:
return None
return pod.metadata.labels.get("model")

def _reconcile_engines(self):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

There is a subtle race condition between _reconcile_engines() and self.k8s_watcher.stream(). Because they are two separate API calls (list_namespaced_pod in _reconcile_engines and another list_namespaced_pod internally inside self.k8s_watcher.stream), if a pod is deleted in the short window between these two calls:\n\n1. _reconcile_engines will see the pod as alive (since it was alive during the first list call) and won't remove it.\n2. self.k8s_watcher.stream will start its watch after the pod is already deleted, so it won't see the pod in its initial list and won't receive a DELETED event for it.\n\nAs a result, the deleted pod will remain in self.available_engines as a ghost pod indefinitely until the next reconnect.\n\nTo completely eliminate this race condition and also avoid redundant HTTP requests to all pods on every reconnect (since the watch stream without resource_version lists all pods and triggers ADDED events for all of them), you can:\n1. Perform the list call once in _reconcile_engines and return both the live pods and the resource_version (pods.metadata.resource_version).\n2. Populate/reconcile self.available_engines using that list.\n3. Start the watch stream using that resource_version so it only streams subsequent events.

"""
Drop engines whose pod is gone from the cluster.

The watch stream is the only source of removals, so a DELETED event
that lands while the connection is down is lost for good and the
engine keeps receiving traffic. Listing on every (re)connection
closes that window: the watch keeps handling the steady state, the
list repairs whatever it missed.
"""
pods = self.k8s_api.list_namespaced_pod(
namespace=self.namespace,
label_selector=self.label_selector,
)
live_pods = {pod.metadata.name for pod in pods.items}

with self.available_engines_lock:
stale = set(self.available_engines) - live_pods
for engine_name in stale:
logger.warning(
f"Serving engine {engine_name} no longer exists but was "
f"still registered: dropping it"
)
del self.available_engines[engine_name]

def _watch_engines(self):
while self.running:
try:
self._reconcile_engines()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

The K8sServiceNameServiceDiscovery class implements a very similar watch loop in its _watch_engines method but does not have a reconciliation mechanism. It is vulnerable to the exact same watch-reconnect bug where a service deleted while the watch is down becomes a permanent ghost service.\n\nPlease implement a similar reconciliation method (e.g., _reconcile_services) for K8sServiceNameServiceDiscovery to ensure consistency and prevent the same bug there.

for event in self.k8s_watcher.stream(
self.k8s_api.list_namespaced_pod,
namespace=self.namespace,
Expand Down