Skip to main content

alien_core/resources/
kubernetes_compute.rs

1//! Compute pools on Kubernetes use administrator-managed node capacity.
2//!
3//! Hardware sizes and machine counts are advisory. Execution requirements must
4//! either be translated into Pod placement or rejected before workloads start.
5
6use crate::{
7    instance_catalog::Architecture, CapacityGroup, ComputeCluster, Container, Daemon, Stack,
8};
9use std::collections::{BTreeMap, HashSet};
10
11/// Validate logical pools without querying or provisioning Kubernetes nodes.
12pub 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
85/// Architecture shared by source images built for this Kubernetes stack.
86/// Unspecified pools and Workers inherit this build constraint, so their Pods
87/// cannot land on incompatible nodes in a cluster with mixed architectures.
88pub 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
104/// Match Pod placement to the architecture used by the stack's source images.
105/// This does not associate Workers with a compute pool or reserve node capacity.
106pub 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
127/// Resolve placement, inferring a single declared cluster and its general or
128/// sole pool. Legacy stacks without a declaration keep default scheduling.
129pub 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
141/// Resolve a daemon's pool without changing its per-node DaemonSet behavior.
142pub 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
203/// Resolve the one pool admitted for release-independent containers.
204pub 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}