Conversation
The watch stream is the only source of engine removals in K8sPodIPServiceDiscovery, so a DELETED event that lands while the connection is down is lost for good: the engine stays in available_engines and keeps taking traffic. Requests hang until they fail with ClientConnectorError, while /v1/models and /health stay green because both are served from the router's own state. This is distinct from vllm-project#1012, which handles the eviction case (MODIFIED with podIP=None): here the pod is deleted normally and the event is simply never delivered. List the pods before opening each stream and drop the engines that no longer exist. The loop already restarts on every disconnect, so the list runs exactly when events may have been missed, and only once at startup in steady state. Reproduced on a scale-to-zero GPU cluster by blocking the router's egress to the API server, deleting the worker while blind, then restoring: without this change the router served the dead pod indefinitely; with it the engine is dropped on reconnect. Fixes vllm-project#1091 Signed-off-by: Julien Renaud <julien.renaud@sancare.fr>
There was a problem hiding this comment.
Code Review
This pull request introduces a reconciliation mechanism (_reconcile_engines) in K8sPodIPServiceDiscovery to clean up stale engines whose Kubernetes pods were deleted while the watch stream was disconnected, along with corresponding unit tests. The review feedback highlights two important points: first, a similar watch-reconnect vulnerability exists in K8sServiceNameServiceDiscovery which should also implement a reconciliation mechanism; second, a potential race condition exists between the reconciliation list call and the start of the watch stream, which can be resolved by using the resource_version from the list call to initialize the watch stream.
| def _watch_engines(self): | ||
| while self.running: | ||
| try: | ||
| self._reconcile_engines() |
There was a problem hiding this comment.
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.
| return None | ||
| return pod.metadata.labels.get("model") | ||
|
|
||
| def _reconcile_engines(self): |
There was a problem hiding this comment.
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.
…rsion Addresses the review on the first revision. The list now also populates available_engines and returns its resourceVersion, which the watch starts from. Without it there was still a window: a pod deleted between the list and the stream was seen alive by the list and absent from the stream's own initial list, so no DELETED event ever arrived. Starting the watch at the list's resourceVersion closes it, and drops the redundant /v1/models request to every pod on each reconnect. Since the stream no longer replays ADDED events, the reconciliation is also the startup discovery path. A ready pod already registered at the same URL, model label and sleep label is skipped, so the common case costs one list. K8sServiceNameServiceDiscovery gets the same treatment: it had the identical watch-reconnect hole. Each object is reconciled in isolation. One unreachable pod, or a Service with no Endpoints object, would otherwise abort the whole pass and leave the router with no engines at all. Tests: startup discovery, skip of an unchanged pod, readiness and cleared-IP transitions missed during a disconnect, sleep-label change, resourceVersion propagation, per-object error isolation, and the service-name equivalents. Signed-off-by: Julien Renaud <julien.renaud@sancare.fr>
|
Thanks — both points were right, and acting on them found a third problem. Pushed a second commit. Race between the list and the watch. Fixed as suggested: One consequence worth stating explicitly: with a
Third problem, found while doing the above. Because the reconciliation is now the discovery path, a single failing object was enough to break everything: a Known limitation, not addressed here: the reconciliation is sequential, and an unknown pod that is k8s-ready but HTTP-unreachable costs Testing. 13 unit tests covering startup discovery, the skip, readiness and cleared-IP transitions missed during a disconnect, sleep-label changes, End to end on a scale-to-zero GPU cluster, same protocol as before — block the router's egress to the API server, wait for the watch to drop, delete the worker while blind, restore: Against the unpatched image, the same sequence leaves |
| return ( | ||
| known.url == f"http://{pod.status.pod_ip}:{self.port}" | ||
| and known.model_label == self._get_model_label(pod) | ||
| and known.sleep == (labels.get("sleeping") == "true") | ||
| ) |
There was a problem hiding this comment.
Should we also compare the pod's lora-modified annotation (store it on EndpointInfo)?
There was a problem hiding this comment.
Yes — good catch, it was a real hole. Done in 7adf99a.
lora-modified is the annotation the operator stamps on the pod (triggerPodEvent, loraadapter_controller.go) for the sole purpose of forcing the router to re-read /v1/models. It is the only pod-side signal that the adapter list moved: when an adapter is loaded or unloaded, the URL, the readiness, the model label and the sleeping label are all identical before and after. So _is_unchanged returned True and, on reconnect, a lora-modified event lost during the disconnect was swallowed — the router kept serving a stale adapter list until the next real event on that pod. Same class of bug as the one this PR fixes, applied to adapters.
EndpointInfo now carries lora_modified, _handle_pod reads it through _get_lora_modified() and passes it down to _add_engine, and _is_unchanged compares it. New test test_reconcile_rehandles_pod_whose_lora_annotation_changed, checked against a mutation removing the comparison.
K8sServiceNameServiceDiscovery is deliberately untouched on this point: the operator annotates pods, not Services, so there is no equivalent signal to compare there.
The same commit also applies black to the two previous ones, which is what CI was failing on.
The LoRA operator stamps lora-modified on the pod every time it loads or unloads an adapter, for the sole purpose of forcing the router to re-read /v1/models. Nothing else about the pod moves: URL, readiness, model label and sleep label are all identical before and after. _is_unchanged therefore skipped such a pod on reconnect, and an adapter change that happened while the watch was down was never picked up: the router kept serving a stale adapter list until the next real event. Same class of bug as the one this series fixes, applied to adapters. EndpointInfo now carries the annotation and the reconciliation compares it. K8sServiceNameServiceDiscovery is untouched: the operator annotates pods, not Services, so there is no equivalent signal there. Also reformats the two previous commits to black, which CI rejected. Signed-off-by: Julien Renaud <julien.renaud@sancare.fr>
FIX #1091
K8sPodIPServiceDiscovery._watch_engines()treats the watch stream as its only source of engine removals. When the stream drops, the loop reconnects but never re-lists, so aDELETEDevent delivered while the connection is down is lost permanently: the engine stays inavailable_enginesand keeps receiving traffic.Requests then round-robin onto a pod IP that no longer exists and hang until
ClientConnectorError, while/v1/modelsand/healthstay green — both are served from the router's own state, so health checks never notice.vllm:num_requests_runningfor that server climbs and never drains, which also skews least-loaded routing.This is distinct from #1008 / #1012, which fixed the eviction case (
MODIFIEDwithpodIP=None). Here the pod is deleted normally; the event is simply never delivered.The change
_reconcile_engines()lists the pods and drops the engines whose pod is gone. It is called at the top of thewhile self.runningloop, so it runs before each stream is opened — that is, exactly when events may have been missed, and only once at startup in steady state. No new thread, no new flag, no periodic cost.Removal happens under
available_engines_lockrather than through_delete_engine(), which takes the same lock and woulddelwithout a guard — avoiding both a deadlock and a race with the watch thread.Testing
Two unit tests in the existing
test_k8s_pod_ip_service_discovery.py: an engine whose pod is gone is dropped, and live engines are never dropped. Both fail if the detection is neutralised (verified by mutatingstaleto an empty set).Validated end to end on a scale-to-zero GPU cluster (KEDA, workers 0↔1), same protocol both times — block the router's egress to the API server, wait for the watch to drop on its own, delete the worker while blind, restore the network:
no longer exists but was still registered: dropping it/v1/modelsadvertises the model, real request times out at 45 sNote for anyone reproducing this: the NetworkPolicy must exclude the real control-plane IP, not the
kubernetes.defaultClusterIP (kube-proxy DNATs before the policy applies), and an already-established watch survives the policy — the router is only truly blind once its next reconnect fails.🤖 Generated with Claude Code