Skip to main content

rskit_workload/
component.rs

1//! Lifecycle-managed workload component.
2//!
3//! Mirrors gokit's `workload.Component`: wraps a [`Manager`] built from an
4//! injected [`WorkloadRegistry`] and participates in ordered startup, shutdown,
5//! and health reporting through [`rskit_component::Component`].
6
7use std::sync::Arc;
8
9use async_trait::async_trait;
10use parking_lot::Mutex;
11use rskit_component::{Component, Health};
12use rskit_errors::AppResult;
13use tracing::{debug, info};
14
15use crate::config::WorkloadConfig;
16use crate::manager::Manager;
17use crate::registry::WorkloadRegistry;
18
19/// A lifecycle-managed workload component.
20///
21/// On start it builds the configured backend from the injected registry and
22/// probes it with a health check so an unreachable backend fails startup rather
23/// than reporting healthy; on stop it releases the manager. When
24/// [`WorkloadConfig::enabled`] is `false` the component is a healthy no-op and
25/// never touches the registry.
26pub struct WorkloadComponent {
27    config: WorkloadConfig,
28    providers: WorkloadRegistry,
29    manager: Mutex<Option<Arc<dyn Manager>>>,
30}
31
32impl WorkloadComponent {
33    /// Create a component with an empty provider registry.
34    ///
35    /// Useful when the component is disabled; register providers via
36    /// [`WorkloadComponent::with_registry`] to run an enabled backend.
37    #[must_use]
38    pub fn new(config: WorkloadConfig) -> Self {
39        Self::with_registry(config, WorkloadRegistry::new())
40    }
41
42    /// Create a component with an explicit provider registry.
43    #[must_use]
44    pub fn with_registry(config: WorkloadConfig, providers: WorkloadRegistry) -> Self {
45        Self {
46            config,
47            providers,
48            manager: Mutex::new(None),
49        }
50    }
51
52    /// Return the underlying manager once the component has started, or `None`.
53    #[must_use]
54    pub fn manager(&self) -> Option<Arc<dyn Manager>> {
55        self.manager.lock().clone()
56    }
57}
58
59#[async_trait]
60impl Component for WorkloadComponent {
61    #[allow(clippy::unnecessary_literal_bound)] // trait fixes the return type to &str
62    fn name(&self) -> &str {
63        "workload"
64    }
65
66    async fn start(&self) -> AppResult<()> {
67        let mut config = self.config.clone();
68        config.apply_defaults();
69
70        if !config.enabled {
71            debug!("workload component disabled — skipping backend init");
72            return Ok(());
73        }
74
75        let manager = self.providers.build(&config).await?;
76        manager.health_check().await?;
77        info!(provider = %config.provider, "workload manager initialized");
78        *self.manager.lock() = Some(manager);
79        Ok(())
80    }
81
82    async fn stop(&self) -> AppResult<()> {
83        debug!("workload component stopping");
84        *self.manager.lock() = None;
85        Ok(())
86    }
87
88    fn health(&self) -> Health {
89        if !self.config.enabled {
90            return Health::healthy("workload (disabled)");
91        }
92        if self.manager.lock().is_none() {
93            return Health::unhealthy("workload", "manager not initialized");
94        }
95        Health::healthy("workload")
96    }
97}
98
99#[cfg(test)]
100mod tests {
101    use super::*;
102    use crate::test_support::{FailingFactory, FakeFactory, UnhealthyFactory};
103    use rskit_errors::ErrorCode;
104
105    fn enabled_registry(
106        factory: Arc<dyn crate::registry::ManagerFactory>,
107    ) -> (WorkloadConfig, WorkloadRegistry) {
108        let mut registry = WorkloadRegistry::new();
109        registry.register("docker", factory).unwrap();
110        let config = WorkloadConfig {
111            enabled: true,
112            provider: "docker".to_string(),
113            ..Default::default()
114        };
115        (config, registry)
116    }
117
118    #[tokio::test]
119    async fn disabled_component_starts_healthy_without_manager() {
120        let component = WorkloadComponent::new(WorkloadConfig::default());
121        assert_eq!(component.name(), "workload");
122        component.start().await.unwrap();
123        assert!(component.manager().is_none());
124        assert!(component.health().is_healthy());
125        component.stop().await.unwrap();
126    }
127
128    #[tokio::test]
129    async fn enabled_component_builds_manager_and_releases_on_stop() {
130        let (config, registry) = enabled_registry(Arc::new(FakeFactory));
131        let component = WorkloadComponent::with_registry(config, registry);
132
133        component.start().await.unwrap();
134        assert!(component.manager().is_some());
135        assert!(component.health().is_healthy());
136
137        component.stop().await.unwrap();
138        assert!(component.manager().is_none());
139        assert!(!component.health().is_healthy());
140    }
141
142    #[tokio::test]
143    async fn enabled_component_before_start_is_unhealthy() {
144        let (config, registry) = enabled_registry(Arc::new(FakeFactory));
145        let component = WorkloadComponent::with_registry(config, registry);
146        assert!(!component.health().is_healthy());
147    }
148
149    #[tokio::test]
150    async fn start_propagates_backend_build_failure() {
151        let (config, registry) = enabled_registry(Arc::new(FailingFactory));
152        let component = WorkloadComponent::with_registry(config, registry);
153        let err = component.start().await.unwrap_err();
154        assert_eq!(err.code(), ErrorCode::Internal);
155        assert!(component.manager().is_none());
156    }
157
158    #[tokio::test]
159    async fn start_fails_when_backend_health_check_fails() {
160        let (config, registry) = enabled_registry(Arc::new(UnhealthyFactory));
161        let component = WorkloadComponent::with_registry(config, registry);
162        let err = component.start().await.unwrap_err();
163        assert_eq!(err.code(), ErrorCode::ServiceUnavailable);
164        assert!(component.manager().is_none());
165    }
166
167    #[tokio::test]
168    async fn enabled_component_with_unregistered_provider_fails_to_start() {
169        let config = WorkloadConfig {
170            enabled: true,
171            provider: "podman".to_string(),
172            ..Default::default()
173        };
174        let component = WorkloadComponent::with_registry(config, WorkloadRegistry::new());
175        assert_eq!(
176            component.start().await.unwrap_err().code(),
177            ErrorCode::NotFound
178        );
179    }
180}