1use crate::{
7 instance_catalog::Architecture, CapacityGroup, ComputeCluster, Container, Daemon, Stack,
8};
9use std::collections::{BTreeMap, HashSet};
10
11pub fn validate_kubernetes_compute(stack: &Stack) -> Vec<String> {
13 let mut errors = Vec::new();
14 if let Err(message) = kubernetes_compute_architecture(stack) {
15 errors.push(message);
16 }
17 for (_, entry) in stack.resources() {
18 let Some(cluster) = entry.config.downcast_ref::<ComputeCluster>() else {
19 continue;
20 };
21 let mut names = HashSet::new();
22 if cluster.capacity_groups.is_empty() {
23 errors.push(format!(
24 "ComputeCluster '{}' must declare at least one pool on Kubernetes",
25 cluster.id
26 ));
27 }
28 for pool in &cluster.capacity_groups {
29 if pool.group_id.is_empty() || !names.insert(pool.group_id.as_str()) {
30 errors.push(format!(
31 "ComputeCluster '{}' has an empty or duplicate pool '{}'",
32 cluster.id, pool.group_id
33 ));
34 }
35 if pool.min_size > pool.max_size {
36 errors.push(format!(
37 "ComputeCluster '{}' pool '{}' has min greater than max",
38 cluster.id, pool.group_id
39 ));
40 }
41 if pool.instance_type.is_some()
42 || pool.nested_virtualization == Some(true)
43 || pool
44 .profile
45 .as_ref()
46 .is_some_and(|profile| profile.gpu.is_some())
47 {
48 errors.push(format!("ComputeCluster '{}' pool '{}' requests instance type, GPU, or nested virtualization; Kubernetes cannot enforce these node requirements", cluster.id, pool.group_id));
49 }
50 }
51 if cluster
52 .dynamic_container_pool
53 .as_ref()
54 .is_some_and(|name| !names.contains(name.as_str()))
55 {
56 errors.push(format!(
57 "ComputeCluster '{}' dynamic container pool does not exist",
58 cluster.id
59 ));
60 }
61 if !cluster.selected_failure_domains.is_empty()
62 || cluster
63 .container_cidr
64 .as_deref()
65 .is_some_and(|cidr| cidr != "10.244.0.0/16")
66 {
67 errors.push(format!("ComputeCluster '{}' configures provider failure domains or a container CIDR; existing Kubernetes networking and node topology are administrator-managed", cluster.id));
68 }
69 }
70 for (_, entry) in stack.resources() {
71 let placement = if let Some(container) = entry.config.downcast_ref::<Container>() {
72 kubernetes_container_pool(stack, container)
73 } else if let Some(daemon) = entry.config.downcast_ref::<Daemon>() {
74 kubernetes_daemon_pool(stack, daemon)
75 } else {
76 continue;
77 };
78 if let Err(message) = placement {
79 errors.push(message);
80 }
81 }
82 errors
83}
84
85pub fn kubernetes_compute_architecture(stack: &Stack) -> Result<Option<Architecture>, String> {
89 let mut selected = None;
90 for architecture in stack
91 .resources()
92 .filter_map(|(_, entry)| entry.config.downcast_ref::<ComputeCluster>())
93 .flat_map(|cluster| &cluster.capacity_groups)
94 .filter_map(|pool| pool.profile.as_ref()?.architecture)
95 {
96 if selected.is_some_and(|selected| selected != architecture) {
97 return Err("Kubernetes compute pools require mixed CPU architectures; one platform image cannot satisfy both".to_string());
98 }
99 selected = Some(architecture);
100 }
101 Ok(selected)
102}
103
104pub fn kubernetes_compute_node_selector(
107 stack: &Stack,
108 pool: Option<&CapacityGroup>,
109) -> Result<Option<BTreeMap<String, String>>, String> {
110 let built_architecture = kubernetes_compute_architecture(stack)?;
111 let architecture = pool
112 .and_then(|pool| pool.profile.as_ref())
113 .and_then(|profile| profile.architecture)
114 .or(built_architecture);
115 Ok(architecture.map(|architecture| {
116 BTreeMap::from([(
117 "kubernetes.io/arch".to_string(),
118 match architecture {
119 Architecture::Arm64 => "arm64",
120 Architecture::X86_64 => "amd64",
121 }
122 .to_string(),
123 )])
124 }))
125}
126
127pub fn kubernetes_container_pool<'a>(
130 stack: &'a Stack,
131 container: &Container,
132) -> Result<Option<&'a CapacityGroup>, String> {
133 kubernetes_workload_pool(
134 stack,
135 &container.id,
136 container.cluster.as_deref(),
137 container.pool.as_deref(),
138 )
139}
140
141pub fn kubernetes_daemon_pool<'a>(
143 stack: &'a Stack,
144 daemon: &Daemon,
145) -> Result<Option<&'a CapacityGroup>, String> {
146 kubernetes_workload_pool(
147 stack,
148 &daemon.id,
149 daemon.cluster.as_deref(),
150 daemon.pool.as_deref(),
151 )
152}
153
154fn kubernetes_workload_pool<'a>(
155 stack: &'a Stack,
156 resource_id: &str,
157 cluster: Option<&str>,
158 pool: Option<&str>,
159) -> Result<Option<&'a CapacityGroup>, String> {
160 let cluster = match cluster {
161 Some(cluster_id) => stack
162 .resources
163 .get(cluster_id)
164 .and_then(|entry| entry.config.downcast_ref::<ComputeCluster>())
165 .ok_or_else(|| {
166 format!("Workload '{resource_id}' references missing ComputeCluster '{cluster_id}'")
167 })?,
168 None => {
169 let mut clusters = stack
170 .resources()
171 .filter_map(|(_, entry)| entry.config.downcast_ref::<ComputeCluster>());
172 let Some(cluster) = clusters.next() else {
173 return if pool.is_some() {
174 Err(format!(
175 "Workload '{resource_id}' selects a pool without a ComputeCluster"
176 ))
177 } else {
178 Ok(None)
179 };
180 };
181 if clusters.next().is_some() {
182 return Err(format!(
183 "Workload '{resource_id}' must select a ComputeCluster explicitly"
184 ));
185 }
186 cluster
187 }
188 };
189 let selected = match pool {
190 Some(name) => cluster
191 .capacity_groups
192 .iter()
193 .find(|pool| pool.group_id == name),
194 None => cluster
195 .capacity_groups
196 .iter()
197 .find(|pool| pool.group_id == "general")
198 .or_else(|| (cluster.capacity_groups.len() == 1).then(|| &cluster.capacity_groups[0])),
199 };
200 selected.map(Some).ok_or_else(|| format!("Workload '{resource_id}' references missing or ambiguous pool '{}' in ComputeCluster '{}'", pool.unwrap_or("general"), cluster.id))
201}
202
203pub fn kubernetes_dynamic_pool(stack: &Stack) -> Result<Option<&CapacityGroup>, String> {
205 let mut clusters = stack
206 .resources()
207 .filter_map(|(_, entry)| entry.config.downcast_ref::<ComputeCluster>());
208 let Some(cluster) = clusters.next() else {
209 return Ok(None);
210 };
211 if clusters.next().is_some() {
212 return Err("Dynamic containers require exactly one ComputeCluster".to_string());
213 }
214 let name = cluster
215 .dynamic_container_pool
216 .as_deref()
217 .unwrap_or("general");
218 cluster
219 .capacity_groups
220 .iter()
221 .find(|pool| pool.group_id == name)
222 .map(Some)
223 .ok_or_else(|| {
224 format!(
225 "ComputeCluster '{}' has no dynamic container pool '{name}'",
226 cluster.id
227 )
228 })
229}
230
231#[cfg(test)]
232mod tests {
233 use super::*;
234 use crate::{instance_catalog::Architecture, ContainerCode, ResourceLifecycle};
235
236 fn cluster() -> ComputeCluster {
237 ComputeCluster::new("compute".to_string())
238 .capacity_group(CapacityGroup {
239 group_id: "apps".to_string(),
240 instance_type: None,
241 profile: Some(crate::MachineProfile {
242 cpu: "2".to_string(),
243 memory_bytes: 4 << 30,
244 ephemeral_storage_bytes: 20 << 30,
245 architecture: Some(Architecture::X86_64),
246 gpu: None,
247 }),
248 min_size: 1,
249 max_size: 3,
250 scale_policy: None,
251 nested_virtualization: None,
252 })
253 .dynamic_container_pool("apps".to_string())
254 .build()
255 }
256
257 fn stack(cluster: ComputeCluster) -> Stack {
258 let container = Container::new("api".to_string())
259 .code(ContainerCode::Image {
260 image: "example.test/api:1".to_string(),
261 })
262 .cpu(crate::ResourceSpec {
263 min: "1".to_string(),
264 desired: "1".to_string(),
265 })
266 .memory(crate::ResourceSpec {
267 min: "128Mi".to_string(),
268 desired: "128Mi".to_string(),
269 })
270 .permissions("default".to_string())
271 .cluster("compute".to_string())
272 .pool("apps".to_string())
273 .build();
274 Stack::new("test".to_string())
275 .add(cluster, ResourceLifecycle::Frozen)
276 .add(container, ResourceLifecycle::Live)
277 .build()
278 }
279
280 #[test]
281 fn portable_pool_preserves_placement_and_dynamic_admission() {
282 let stack = stack(cluster());
283 assert!(validate_kubernetes_compute(&stack).is_empty());
284 let pool = kubernetes_dynamic_pool(&stack).unwrap().unwrap();
285 assert_eq!(pool.group_id, "apps");
286 assert_eq!(
287 pool.profile.as_ref().unwrap().architecture,
288 Some(Architecture::X86_64)
289 );
290 }
291
292 #[test]
293 fn rejects_execution_requirements_and_invalid_references() {
294 let mut cluster = cluster();
295 cluster.capacity_groups[0].nested_virtualization = Some(true);
296 cluster.dynamic_container_pool = Some("missing".to_string());
297 let errors = validate_kubernetes_compute(&stack(cluster));
298 assert!(errors
299 .iter()
300 .any(|message| message.contains("nested virtualization")));
301 assert!(errors
302 .iter()
303 .any(|message| message.contains("dynamic container pool")));
304 let mut cluster = self::cluster();
305 cluster.capacity_groups[0].group_id = "other".to_string();
306 assert!(validate_kubernetes_compute(&stack(cluster))
307 .iter()
308 .any(|message| message.contains("pool 'apps'")));
309 }
310
311 #[test]
312 fn implicit_placement_uses_the_declared_pool_architecture() {
313 let stack = stack(cluster());
314 let mut container = stack
315 .resources
316 .get("api")
317 .unwrap()
318 .config
319 .downcast_ref::<Container>()
320 .unwrap()
321 .clone();
322 container.cluster = None;
323 container.pool = None;
324 let pool = kubernetes_container_pool(&stack, &container)
325 .unwrap()
326 .unwrap();
327 assert_eq!(pool.group_id, "apps");
328 assert_eq!(
329 pool.profile.as_ref().unwrap().architecture,
330 Some(Architecture::X86_64)
331 );
332 }
333
334 #[test]
335 fn unspecified_pool_and_unpooled_worker_use_the_build_architecture() {
336 let mut compute = cluster();
337 let mut unspecified = compute.capacity_groups[0].clone();
338 unspecified.group_id = "other".to_string();
339 unspecified.profile.as_mut().unwrap().architecture = None;
340 compute.capacity_groups.push(unspecified.clone());
341 let stack = stack(compute);
342 let expected = Some(BTreeMap::from([(
343 "kubernetes.io/arch".to_string(),
344 "amd64".to_string(),
345 )]));
346 assert_eq!(
347 kubernetes_compute_node_selector(&stack, Some(&unspecified)).unwrap(),
348 expected
349 );
350 assert_eq!(
351 kubernetes_compute_node_selector(&stack, None).unwrap(),
352 expected
353 );
354 }
355
356 #[test]
357 fn legacy_stack_keeps_default_placement() {
358 let stack = Stack::new("legacy".to_string()).build();
359 assert!(validate_kubernetes_compute(&stack).is_empty());
360 assert!(kubernetes_dynamic_pool(&stack).unwrap().is_none());
361 }
362}