Skip to main content

alien_core/resources/
daemon.rs

1use crate::error::{ErrorData, Result};
2use crate::resource::{ResourceDefinition, ResourceOutputsDefinition, ResourceRef, ResourceType};
3use crate::resources::{
4    ComputeCluster, ExposeProtocol, HealthCheck, PublicEndpoint, PublicEndpointOutput,
5    ResourceSpec, ToolchainConfig, APEX_HOST_LABEL,
6};
7use alien_error::AlienError;
8use bon::Builder;
9use serde::{Deserialize, Serialize};
10use std::any::Any;
11use std::collections::HashMap;
12use std::fmt::Debug;
13
14#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
15#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
16#[serde(rename_all = "camelCase", tag = "type")]
17pub enum DaemonCode {
18    #[serde(rename_all = "camelCase")]
19    Image { image: String },
20    #[serde(rename_all = "camelCase")]
21    Source {
22        src: String,
23        toolchain: ToolchainConfig,
24    },
25}
26
27#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
28#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
29#[serde(rename_all = "camelCase", deny_unknown_fields)]
30pub struct DaemonRuntimeMount {
31    /// Absolute host path to mount into the daemon container.
32    pub source: String,
33    /// Absolute container path where the source is mounted.
34    pub target: String,
35    /// Optional mount options understood by the backend runtime.
36    #[serde(skip_serializing_if = "Option::is_none")]
37    pub options: Option<String>,
38}
39
40#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
41#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
42#[serde(rename_all = "camelCase", deny_unknown_fields)]
43pub struct DaemonRuntime {
44    /// Run the daemon container with elevated host capabilities.
45    #[serde(skip_serializing_if = "Option::is_none")]
46    pub privileged: Option<bool>,
47    /// Process namespace mode. Supported values are `host` and `private`.
48    #[serde(skip_serializing_if = "Option::is_none")]
49    pub pid_namespace: Option<String>,
50    /// Network mode. Supported values are `host` and `appnet`.
51    #[serde(skip_serializing_if = "Option::is_none")]
52    pub network_mode: Option<String>,
53    /// Host mounts exposed to the daemon container.
54    #[serde(default, skip_serializing_if = "Vec::is_empty")]
55    pub mounts: Vec<DaemonRuntimeMount>,
56    /// Runtime user, as a numeric uid or uid:gid string.
57    #[serde(skip_serializing_if = "Option::is_none")]
58    pub user: Option<String>,
59}
60
61#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Builder)]
62#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
63#[serde(rename_all = "camelCase", deny_unknown_fields)]
64#[builder(start_fn = new)]
65pub struct Daemon {
66    #[builder(start_fn)]
67    pub id: String,
68    #[builder(field)]
69    pub links: Vec<ResourceRef>,
70    /// Public endpoints exposed by the daemon.
71    #[builder(field)]
72    #[serde(default, skip_serializing_if = "Vec::is_empty")]
73    pub public_endpoints: Vec<PublicEndpoint>,
74    /// HTTP health check for public daemon endpoint load balancers.
75    #[serde(skip_serializing_if = "Option::is_none")]
76    pub health_check: Option<HealthCheck>,
77    /// ComputeCluster resource ID that this daemon runs on for managed cloud
78    /// compute backends. Kubernetes and Local runtimes ignore this field.
79    #[serde(skip_serializing_if = "Option::is_none")]
80    pub cluster: Option<String>,
81    /// Named workload permission profile. Absent means no workload cloud identity.
82    #[serde(default, skip_serializing_if = "Option::is_none")]
83    pub permissions: Option<String>,
84    pub code: DaemonCode,
85    /// CPU resource requirements for each daemon instance.
86    #[builder(default = default_daemon_cpu())]
87    #[serde(default = "default_daemon_cpu")]
88    pub cpu: ResourceSpec,
89    /// Memory resource requirements for each daemon instance.
90    #[builder(default = default_daemon_memory())]
91    #[serde(default = "default_daemon_memory")]
92    pub memory: ResourceSpec,
93    /// Capacity group/pool to run on for backends that expose machine pools.
94    #[serde(skip_serializing_if = "Option::is_none")]
95    pub pool: Option<String>,
96    /// Command to override the image default.
97    #[serde(skip_serializing_if = "Option::is_none")]
98    pub command: Option<Vec<String>>,
99    /// Grace period in seconds for stopping daemon instances during updates, drains, and deletes.
100    ///
101    /// When omitted, the runtime backend applies its default. Valid values are
102    /// 1 second through 24 hours.
103    #[serde(skip_serializing_if = "Option::is_none")]
104    #[cfg_attr(feature = "openapi", schema(minimum = 1, maximum = 86400))]
105    pub stop_grace_period_seconds: Option<u32>,
106    /// Optional backend runtime settings for trusted daemons.
107    ///
108    /// These settings are intended for daemon-style infrastructure that must
109    /// operate on the host. Backends that do not support a setting may reject
110    /// it during provisioning.
111    #[serde(skip_serializing_if = "Option::is_none")]
112    pub runtime: Option<DaemonRuntime>,
113    #[builder(default)]
114    #[serde(default)]
115    pub environment: HashMap<String, String>,
116    /// Whether an app-owned command receiver can lease pending commands for
117    /// this Daemon and execute registered handlers.
118    #[builder(default = default_commands_enabled())]
119    #[serde(default = "default_commands_enabled")]
120    #[cfg_attr(feature = "openapi", schema(default = default_commands_enabled))]
121    pub commands_enabled: bool,
122}
123
124impl Daemon {
125    pub const RESOURCE_TYPE: ResourceType = ResourceType::from_static("daemon");
126
127    pub fn get_permissions(&self) -> Option<&str> {
128        self.permissions.as_deref()
129    }
130
131    fn validate_public_endpoints(&self) -> Result<()> {
132        let mut endpoint_names = std::collections::HashSet::new();
133        let mut backend_ports = std::collections::HashSet::new();
134        let mut apex_endpoint_name: Option<&str> = None;
135
136        for endpoint in &self.public_endpoints {
137            endpoint.validate_for_resource(&self.id)?;
138            if !endpoint_names.insert(endpoint.name.as_str()) {
139                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
140                    resource_id: self.id.clone(),
141                    reason: format!("duplicate public endpoint name '{}'", endpoint.name),
142                }));
143            }
144            if endpoint.host_label.as_deref() == Some(APEX_HOST_LABEL) {
145                if let Some(existing_name) = apex_endpoint_name {
146                    return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
147                        resource_id: self.id.clone(),
148                        reason: format!(
149                            "only one apex public endpoint is allowed per resource; '{}' already uses hostLabel '@'",
150                            existing_name
151                        ),
152                    }));
153                }
154                apex_endpoint_name = Some(endpoint.name.as_str());
155            }
156            backend_ports.insert(endpoint.port);
157            if endpoint.protocol != ExposeProtocol::Http {
158                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
159                    resource_id: self.id.clone(),
160                    reason: "daemon public endpoints currently support only HTTP".to_string(),
161                }));
162            }
163        }
164
165        if backend_ports.len() > 1 {
166            return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
167                resource_id: self.id.clone(),
168                reason:
169                    "public endpoints on one daemon must currently route to the same backend port"
170                        .to_string(),
171            }));
172        }
173
174        Ok(())
175    }
176
177    fn validate_runtime(&self) -> Result<()> {
178        let Some(runtime) = &self.runtime else {
179            return Ok(());
180        };
181
182        if let Some(pid_namespace) = &runtime.pid_namespace {
183            if pid_namespace != "host" && pid_namespace != "private" {
184                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
185                    resource_id: self.id.clone(),
186                    reason: "runtime.pidNamespace must be 'host' or 'private'".to_string(),
187                }));
188            }
189        }
190
191        if let Some(network_mode) = &runtime.network_mode {
192            if network_mode != "host" && network_mode != "appnet" {
193                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
194                    resource_id: self.id.clone(),
195                    reason: "runtime.networkMode must be 'host' or 'appnet'".to_string(),
196                }));
197            }
198        }
199
200        if let Some(user) = &runtime.user {
201            let valid = match user.split_once(':') {
202                Some((uid, gid)) => {
203                    !uid.is_empty()
204                        && !gid.is_empty()
205                        && uid.chars().all(|c| c.is_ascii_digit())
206                        && gid.chars().all(|c| c.is_ascii_digit())
207                }
208                None => !user.is_empty() && user.chars().all(|c| c.is_ascii_digit()),
209            };
210            if !valid {
211                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
212                    resource_id: self.id.clone(),
213                    reason: "runtime.user must be a numeric uid or uid:gid".to_string(),
214                }));
215            }
216        }
217
218        for mount in &runtime.mounts {
219            if mount.source.is_empty() || mount.target.is_empty() {
220                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
221                    resource_id: self.id.clone(),
222                    reason: "runtime.mounts source and target must be non-empty".to_string(),
223                }));
224            }
225            if !mount.source.starts_with('/') || !mount.target.starts_with('/') {
226                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
227                    resource_id: self.id.clone(),
228                    reason: "runtime.mounts source and target must be absolute paths".to_string(),
229                }));
230            }
231        }
232
233        Ok(())
234    }
235}
236
237fn default_commands_enabled() -> bool {
238    false
239}
240
241fn default_daemon_cpu() -> ResourceSpec {
242    ResourceSpec {
243        min: "0.1".to_string(),
244        desired: "0.1".to_string(),
245    }
246}
247
248fn default_daemon_memory() -> ResourceSpec {
249    ResourceSpec {
250        min: "128Mi".to_string(),
251        desired: "128Mi".to_string(),
252    }
253}
254
255impl<S: daemon_builder::State> DaemonBuilder<S> {
256    pub fn link<R: ?Sized>(mut self, resource: &R) -> Self
257    where
258        for<'a> &'a R: Into<ResourceRef>,
259    {
260        let resource_ref: ResourceRef = resource.into();
261        self.links.push(resource_ref);
262        self
263    }
264
265    pub fn public_endpoint(mut self, endpoint: PublicEndpoint) -> Self {
266        self.public_endpoints.push(endpoint);
267        self
268    }
269}
270
271impl ResourceDefinition for Daemon {
272    fn get_resource_type(&self) -> ResourceType {
273        Self::RESOURCE_TYPE
274    }
275
276    fn id(&self) -> &str {
277        &self.id
278    }
279
280    fn get_dependencies(&self) -> Vec<ResourceRef> {
281        let mut dependencies = self.links.clone();
282        if let Some(cluster) = &self.cluster {
283            dependencies.push(ResourceRef::new(
284                ComputeCluster::RESOURCE_TYPE,
285                cluster.clone(),
286            ));
287        }
288        dependencies
289    }
290
291    fn get_permissions(&self) -> Option<&str> {
292        self.permissions.as_deref()
293    }
294
295    fn validate_update(&self, new_config: &dyn ResourceDefinition) -> Result<()> {
296        let new_daemon = new_config
297            .as_any()
298            .downcast_ref::<Daemon>()
299            .ok_or_else(|| {
300                AlienError::new(ErrorData::UnexpectedResourceType {
301                    resource_id: self.id.clone(),
302                    expected: Self::RESOURCE_TYPE,
303                    actual: new_config.get_resource_type(),
304                })
305            })?;
306
307        if self.id != new_daemon.id {
308            return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
309                resource_id: self.id.clone(),
310                reason: "the 'id' field is immutable".to_string(),
311            }));
312        }
313
314        self.validate_public_endpoints()?;
315        new_daemon.validate_public_endpoints()?;
316        self.validate_runtime()?;
317        new_daemon.validate_runtime()?;
318
319        if self.public_endpoints != new_daemon.public_endpoints {
320            return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
321                resource_id: self.id.clone(),
322                reason: "the 'publicEndpoints' field is immutable".to_string(),
323            }));
324        }
325
326        Ok(())
327    }
328
329    fn as_any(&self) -> &dyn Any {
330        self
331    }
332
333    fn as_any_mut(&mut self) -> &mut dyn Any {
334        self
335    }
336
337    fn box_clone(&self) -> Box<dyn ResourceDefinition> {
338        Box::new(self.clone())
339    }
340
341    fn resource_eq(&self, other: &dyn ResourceDefinition) -> bool {
342        other.as_any().downcast_ref::<Daemon>() == Some(self)
343    }
344
345    fn to_json_value(&self) -> serde_json::Result<serde_json::Value> {
346        serde_json::to_value(self)
347    }
348}
349
350#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
351#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
352#[serde(rename_all = "camelCase")]
353pub struct DaemonOutputs {
354    pub daemon_name: String,
355    pub running: bool,
356    #[serde(default, skip_serializing_if = "HashMap::is_empty")]
357    pub public_endpoints: HashMap<String, PublicEndpointOutput>,
358}
359
360impl ResourceOutputsDefinition for DaemonOutputs {
361    fn get_resource_type(&self) -> ResourceType {
362        Daemon::RESOURCE_TYPE.clone()
363    }
364
365    fn as_any(&self) -> &dyn Any {
366        self
367    }
368
369    fn box_clone(&self) -> Box<dyn ResourceOutputsDefinition> {
370        Box::new(self.clone())
371    }
372
373    fn outputs_eq(&self, other: &dyn ResourceOutputsDefinition) -> bool {
374        other.as_any().downcast_ref::<DaemonOutputs>() == Some(self)
375    }
376
377    fn to_json_value(&self) -> serde_json::Result<serde_json::Value> {
378        serde_json::to_value(self)
379    }
380}
381
382#[cfg(test)]
383mod tests {
384    use super::*;
385
386    #[test]
387    fn daemon_without_permissions_preserves_links_without_an_identity() {
388        let storage = crate::Storage::new("objects".to_string()).build();
389        let daemon = Daemon::new("observer".to_string())
390            .code(DaemonCode::Image {
391                image: "observer:latest".to_string(),
392            })
393            .link(&storage)
394            .build();
395        assert_eq!(daemon.get_permissions(), None);
396        assert_eq!(ResourceDefinition::get_permissions(&daemon), None);
397        assert_eq!(daemon.get_dependencies(), daemon.links);
398        assert_eq!(daemon.links[0].id(), "objects");
399        let json = serde_json::to_value(&daemon).unwrap();
400        assert!(json.get("permissions").is_none());
401        let decoded: Daemon = serde_json::from_value(json).unwrap();
402        assert_eq!(decoded, daemon);
403
404        let mut explicit = serde_json::to_value(&daemon).unwrap();
405        explicit["permissions"] = serde_json::json!("reader");
406        let decoded: Daemon = serde_json::from_value(explicit.clone()).unwrap();
407        assert_eq!(decoded.get_permissions(), Some("reader"));
408        assert_eq!(
409            ResourceDefinition::get_permissions(&decoded),
410            Some("reader")
411        );
412        assert_eq!(serde_json::to_value(decoded).unwrap(), explicit);
413    }
414
415    #[test]
416    fn daemon_serializes_with_resource_type() {
417        let daemon = Daemon::new("endpoint-agent".to_string())
418            .code(DaemonCode::Source {
419                src: "./agent".to_string(),
420                toolchain: ToolchainConfig::Rust {
421                    binary_name: "agent".to_string(),
422                },
423            })
424            .permissions("execution".to_string())
425            .commands_enabled(true)
426            .build();
427
428        let resource = crate::Resource::new(daemon);
429        let json = serde_json::to_value(&resource).expect("daemon should serialize");
430        assert_eq!(json["type"], "daemon");
431
432        let roundtrip: crate::Resource =
433            serde_json::from_value(json).expect("daemon should deserialize");
434        assert_eq!(roundtrip.resource_type().as_ref(), "daemon");
435    }
436
437    #[test]
438    fn daemon_accepts_one_public_http_endpoint() {
439        let daemon = Daemon::new("gateway".to_string())
440            .code(DaemonCode::Image {
441                image: "gateway:latest".to_string(),
442            })
443            .public_endpoint(PublicEndpoint {
444                name: "public".to_string(),
445                port: 8080,
446                protocol: ExposeProtocol::Http,
447                host_label: Some("public".to_string()),
448                wildcard_subdomains: true,
449            })
450            .permissions("gateway".to_string())
451            .build();
452
453        assert!(daemon.validate_public_endpoints().is_ok());
454        assert_eq!(daemon.public_endpoints.len(), 1);
455        assert_eq!(
456            daemon.public_endpoints[0].host_label.as_deref(),
457            Some("public")
458        );
459        assert!(daemon.public_endpoints[0].wildcard_subdomains);
460    }
461
462    #[test]
463    fn daemon_serializes_stop_grace_period_when_set() {
464        let daemon = Daemon::new("gateway".to_string())
465            .code(DaemonCode::Image {
466                image: "gateway:latest".to_string(),
467            })
468            .permissions("gateway".to_string())
469            .stop_grace_period_seconds(21_600)
470            .build();
471
472        let json = serde_json::to_value(&daemon).expect("daemon should serialize");
473        assert_eq!(json["stopGracePeriodSeconds"], 21_600);
474    }
475
476    #[test]
477    fn daemon_omits_stop_grace_period_when_absent() {
478        let daemon = Daemon::new("gateway".to_string())
479            .code(DaemonCode::Image {
480                image: "gateway:latest".to_string(),
481            })
482            .permissions("gateway".to_string())
483            .build();
484
485        let json = serde_json::to_value(&daemon).expect("daemon should serialize");
486        assert!(json.get("stopGracePeriodSeconds").is_none());
487    }
488
489    #[test]
490    fn daemon_rejects_multiple_backend_ports_or_non_http_public_endpoints() {
491        let multiple = Daemon::new("gateway".to_string())
492            .code(DaemonCode::Image {
493                image: "gateway:latest".to_string(),
494            })
495            .public_endpoint(PublicEndpoint {
496                name: "api".to_string(),
497                port: 8080,
498                protocol: ExposeProtocol::Http,
499                host_label: None,
500                wildcard_subdomains: false,
501            })
502            .public_endpoint(PublicEndpoint {
503                name: "admin".to_string(),
504                port: 9090,
505                protocol: ExposeProtocol::Http,
506                host_label: None,
507                wildcard_subdomains: false,
508            })
509            .permissions("gateway".to_string())
510            .build();
511        assert!(multiple.validate_public_endpoints().is_err());
512
513        let tcp = Daemon::new("gateway".to_string())
514            .code(DaemonCode::Image {
515                image: "gateway:latest".to_string(),
516            })
517            .public_endpoint(PublicEndpoint {
518                name: "api".to_string(),
519                port: 8080,
520                protocol: ExposeProtocol::Tcp,
521                host_label: None,
522                wildcard_subdomains: false,
523            })
524            .permissions("gateway".to_string())
525            .build();
526        assert!(tcp.validate_public_endpoints().is_err());
527    }
528
529    #[test]
530    fn daemon_rejects_multiple_apex_public_endpoints() {
531        let daemon = Daemon::new("gateway".to_string())
532            .code(DaemonCode::Image {
533                image: "gateway:latest".to_string(),
534            })
535            .public_endpoint(PublicEndpoint {
536                name: "api".to_string(),
537                port: 8080,
538                protocol: ExposeProtocol::Http,
539                host_label: Some(APEX_HOST_LABEL.to_string()),
540                wildcard_subdomains: false,
541            })
542            .public_endpoint(PublicEndpoint {
543                name: "admin".to_string(),
544                port: 8080,
545                protocol: ExposeProtocol::Http,
546                host_label: Some(APEX_HOST_LABEL.to_string()),
547                wildcard_subdomains: false,
548            })
549            .permissions("gateway".to_string())
550            .build();
551
552        assert!(daemon.validate_public_endpoints().is_err());
553    }
554}