rskit_workload/
component.rs1use 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
19pub struct WorkloadComponent {
27 config: WorkloadConfig,
28 providers: WorkloadRegistry,
29 manager: Mutex<Option<Arc<dyn Manager>>>,
30}
31
32impl WorkloadComponent {
33 #[must_use]
38 pub fn new(config: WorkloadConfig) -> Self {
39 Self::with_registry(config, WorkloadRegistry::new())
40 }
41
42 #[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 #[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)] 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}