dynamo_runtime/discovery/kube/
utils.rs1use 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#[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#[derive(Debug, Clone, PartialEq, Eq)]
42pub(super) enum KubeDiscoveryTarget {
43 Pod(String),
44 Container(String, String),
45}
46
47impl KubeDiscoveryTarget {
48 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 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
75pub 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
84pub 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
98pub(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
129pub(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#[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 pub mode: KubeDiscoveryMode,
164 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 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 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 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 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}