Skip to main content

dynamo_runtime/discovery/kube/
utils.rs

1// SPDX-FileCopyrightText: Copyright (c) 2024-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use anyhow::Result;
5use k8s_openapi::api::core::v1::Pod;
6use k8s_openapi::api::discovery::v1::EndpointSlice;
7use std::collections::hash_map::DefaultHasher;
8use std::fs;
9use std::hash::{Hash, Hasher};
10use std::path::Path;
11
12use crate::config::environment_names::discovery;
13
14const INSTANCE_ID_MASK: u64 = 0x001F_FFFF_FFFF_FFFFu64;
15const MAIN_CONTAINER_NAME: &str = "main";
16
17/// Kube discovery mode.
18///
19/// - `Pod`: default. One identity per pod.
20/// - `Container`: each container independently registers with the discovery plane.
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub(super) enum KubeDiscoveryMode {
23    Pod,
24    Container,
25}
26
27impl KubeDiscoveryMode {
28    pub fn from_env() -> Result<Self> {
29        match std::env::var(discovery::DYN_KUBE_DISCOVERY_MODE).as_deref() {
30            Ok("container") => Ok(Self::Container),
31            Ok("pod") | Err(_) => Ok(Self::Pod),
32            Ok(other) => anyhow::bail!(
33                "Invalid DYN_KUBE_DISCOVERY_MODE value '{}'. Valid values: 'pod', 'container'",
34                other
35            ),
36        }
37    }
38}
39
40/// A resolved discovery target identifying either a pod or a specific container within a pod.
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub(super) enum KubeDiscoveryTarget {
43    Pod(String),
44    Container(String, String),
45}
46
47impl KubeDiscoveryTarget {
48    /// CR name for this target, used as the DynamoWorkerMetadata resource name.
49    pub fn cr_name(&self) -> String {
50        match self {
51            Self::Pod(pod_name) => pod_name.clone(),
52            Self::Container(pod_name, container_name) if container_name == MAIN_CONTAINER_NAME => {
53                pod_name.clone()
54            }
55            Self::Container(pod_name, container_name) => {
56                format!("{}-{}", pod_name, container_name)
57            }
58        }
59    }
60
61    /// Deterministic instance ID derived from cr_name.
62    pub fn instance_id(&self) -> u64 {
63        let mut hasher = DefaultHasher::new();
64        self.cr_name().hash(&mut hasher);
65        hasher.finish() & INSTANCE_ID_MASK
66    }
67
68    pub fn pod_name(&self) -> &str {
69        match self {
70            Self::Pod(pod_name) | Self::Container(pod_name, _) => pod_name,
71        }
72    }
73}
74
75/// Hash a pod name to get a consistent instance ID (pod-level).
76///
77/// Used by C bindings (EPP) for pod-level worker ID mapping.
78pub fn hash_pod_name(pod_name: &str) -> u64 {
79    let mut hasher = DefaultHasher::new();
80    pod_name.hash(&mut hasher);
81    hasher.finish() & INSTANCE_ID_MASK
82}
83
84/// Hash a (pod, container) pair to get the same per-container instance ID a
85/// worker computes for itself under `DYN_KUBE_DISCOVERY_MODE=container` (see
86/// `KubeDiscoveryTarget::Container`). A container named `"main"` collapses to
87/// the pod-level identity (`hash_pod_name`), matching a container-mode
88/// frontend's ability to discover pod-mode workers.
89///
90/// Used by the Rust EPP (`deploy/inference-gateway/ext-proc`) to resolve a
91/// registered worker's per-container instance ID back to its pod's endpoint
92/// when `DYN_KUBE_DISCOVERY_MODE=container` (e.g. intra-pod GMS failover,
93/// where each engine container registers under its own container name).
94pub fn hash_container_name(pod_name: &str, container_name: &str) -> u64 {
95    KubeDiscoveryTarget::Container(pod_name.to_string(), container_name.to_string()).instance_id()
96}
97
98/// Extract (instance_id, pod_name) tuples from an EndpointSlice for ready endpoints.
99pub(super) fn extract_endpoint_info(slice: &EndpointSlice) -> Vec<(u64, String)> {
100    let mut result = Vec::new();
101
102    for endpoint in &slice.endpoints {
103        let is_ready = endpoint
104            .conditions
105            .as_ref()
106            .and_then(|c| c.ready)
107            .unwrap_or(false);
108
109        if !is_ready {
110            continue;
111        }
112
113        let pod_name = match endpoint.target_ref.as_ref() {
114            Some(target_ref) => target_ref.name.as_deref().unwrap_or(""),
115            None => continue,
116        };
117
118        if pod_name.is_empty() {
119            continue;
120        }
121
122        let target = KubeDiscoveryTarget::Pod(pod_name.to_string());
123        result.push((target.instance_id(), target.cr_name()));
124    }
125
126    result
127}
128
129/// Extract (instance_id, cr_name) tuples from a Pod for each ready container.
130pub(super) fn extract_ready_containers(pod: &Pod) -> Vec<(u64, String)> {
131    let pod_name = match pod.metadata.name.as_deref() {
132        Some(name) => name,
133        None => return vec![],
134    };
135
136    let container_statuses = match pod
137        .status
138        .as_ref()
139        .and_then(|s| s.container_statuses.as_ref())
140    {
141        Some(statuses) => statuses,
142        None => return vec![],
143    };
144
145    container_statuses
146        .iter()
147        .filter(|cs| cs.ready)
148        .map(|cs| {
149            let target = KubeDiscoveryTarget::Container(pod_name.to_string(), cs.name.clone());
150            (target.instance_id(), target.cr_name())
151        })
152        .collect()
153}
154
155/// Pod information extracted from environment.
156#[derive(Debug, Clone)]
157pub(super) struct PodInfo {
158    pub pod_name: String,
159    pub pod_namespace: String,
160    pub pod_uid: String,
161    pub system_port: u16,
162    /// Kube discovery mode for this process, read from DYN_KUBE_DISCOVERY_MODE.
163    pub mode: KubeDiscoveryMode,
164    /// Discovery target for this process, derived from mode + pod/container identity.
165    pub target: KubeDiscoveryTarget,
166}
167
168const DEFAULT_PODINFO_PATH: &str = "/etc/podinfo";
169
170impl PodInfo {
171    fn read_from_file_or_env(file_path: &Path, env_var: &str) -> Option<String> {
172        if let Ok(content) = fs::read_to_string(file_path) {
173            let value = content.trim().to_string();
174            if !value.is_empty() {
175                return Some(value);
176            }
177        }
178        std::env::var(env_var).ok()
179    }
180
181    pub fn from_env() -> Result<Self> {
182        let podinfo_path = Path::new(DEFAULT_PODINFO_PATH);
183
184        let pod_name = Self::read_from_file_or_env(&podinfo_path.join("pod_name"), "POD_NAME")
185            .ok_or_else(|| anyhow::anyhow!("POD_NAME not available from file or environment"))?;
186
187        let pod_uid = Self::read_from_file_or_env(&podinfo_path.join("pod_uid"), "POD_UID")
188            .ok_or_else(|| anyhow::anyhow!("POD_UID not available from file or environment"))?;
189
190        let pod_namespace =
191            Self::read_from_file_or_env(&podinfo_path.join("pod_namespace"), "POD_NAMESPACE")
192                .unwrap_or_else(|| {
193                    tracing::warn!("POD_NAMESPACE not set, defaulting to 'default'");
194                    "default".to_string()
195                });
196
197        let mode = KubeDiscoveryMode::from_env()?;
198
199        let target = match mode {
200            KubeDiscoveryMode::Pod => KubeDiscoveryTarget::Pod(pod_name.clone()),
201            KubeDiscoveryMode::Container => {
202                let container_name = std::env::var("CONTAINER_NAME").map_err(|_| {
203                    anyhow::anyhow!(
204                        "CONTAINER_NAME is required when DYN_KUBE_DISCOVERY_MODE=container"
205                    )
206                })?;
207                KubeDiscoveryTarget::Container(pod_name.clone(), container_name)
208            }
209        };
210
211        if podinfo_path.join("pod_name").exists() {
212            tracing::info!(
213                "Pod identity loaded from Downward API volume mount at {}",
214                DEFAULT_PODINFO_PATH
215            );
216        } else {
217            tracing::info!("Pod identity loaded from environment variables");
218        }
219
220        let config = crate::config::RuntimeConfig::from_settings().unwrap_or_default();
221        let system_port = config.system_port as u16;
222
223        Ok(Self {
224            pod_name,
225            pod_namespace,
226            pod_uid,
227            system_port,
228            mode,
229            target,
230        })
231    }
232}
233
234#[cfg(test)]
235mod tests {
236    use super::*;
237
238    #[test]
239    fn test_pod_mode_backward_compat() {
240        // Pod mode must produce the same instance_id as hash_pod_name
241        // so existing deployments see no identity change on upgrade.
242        let target = KubeDiscoveryTarget::Pod("worker-0".into());
243        assert_eq!(target.instance_id(), hash_pod_name("worker-0"));
244        assert_eq!(target.cr_name(), "worker-0");
245    }
246
247    #[test]
248    fn test_container_mode_main_uses_pod_identity() {
249        // A container named "main" uses pod-level identity so that
250        // container-mode frontends can discover pod-mode workers.
251        let target = KubeDiscoveryTarget::Container("worker-0".into(), "main".into());
252        assert_eq!(target.instance_id(), hash_pod_name("worker-0"));
253        assert_eq!(target.cr_name(), "worker-0");
254    }
255
256    #[test]
257    fn test_container_mode_engine_gets_unique_identity() {
258        // Non-main containers get per-container identity so that
259        // failover engine containers are independently discoverable.
260        let e0 = KubeDiscoveryTarget::Container("worker-0".into(), "engine-0".into());
261        let e1 = KubeDiscoveryTarget::Container("worker-0".into(), "engine-1".into());
262        assert_eq!(e0.cr_name(), "worker-0-engine-0");
263        assert_eq!(e1.cr_name(), "worker-0-engine-1");
264        assert_ne!(e0.instance_id(), e1.instance_id());
265        assert_ne!(e0.instance_id(), hash_pod_name("worker-0"));
266    }
267
268    #[test]
269    fn test_hash_container_name_matches_target_instance_id() {
270        // hash_container_name is the public entry point the Rust EPP uses; it
271        // must stay in lockstep with the KubeDiscoveryTarget a registering
272        // worker computes for itself, including the "main" pod-identity
273        // fallback and per-engine uniqueness.
274        assert_eq!(
275            hash_container_name("worker-0", "main"),
276            hash_pod_name("worker-0")
277        );
278        let e0 = hash_container_name("worker-0", "engine-0");
279        let e1 = hash_container_name("worker-0", "engine-1");
280        assert_ne!(e0, e1);
281        assert_ne!(e0, hash_pod_name("worker-0"));
282        assert_eq!(
283            e0,
284            KubeDiscoveryTarget::Container("worker-0".into(), "engine-0".into()).instance_id()
285        );
286    }
287}