Skip to main content

camel_function/
service.rs

1use crate::config::FunctionConfig;
2use crate::invoker::DefaultFunctionInvoker;
3use crate::pool::{RunnerPool, RunnerPoolKey, RunnerState};
4use crate::provider::{FunctionHealthStatus, FunctionProvider, ProviderError};
5use camel_api::function::{FunctionDefinition, FunctionId, FunctionInvoker};
6use camel_api::{CamelError, Lifecycle, ServiceStatus};
7use std::sync::Arc;
8use std::sync::atomic::{AtomicU8, Ordering};
9
10const STATUS_STOPPED: u8 = 0;
11const STATUS_STARTED: u8 = 1;
12const STATUS_FAILED: u8 = 2;
13
14pub struct FunctionRuntimeService {
15    config: FunctionConfig,
16    provider: Arc<dyn FunctionProvider>,
17    container_provider: Option<Arc<crate::provider::container::ContainerProvider>>,
18    pub(crate) invoker: Arc<DefaultFunctionInvoker>,
19    status: Arc<AtomicU8>,
20}
21
22impl FunctionRuntimeService {
23    pub(crate) fn new(config: FunctionConfig, provider: Arc<dyn FunctionProvider>) -> Self {
24        let pool = Arc::new(RunnerPool::new());
25        let invoker = Arc::new(DefaultFunctionInvoker::new(
26            Arc::clone(&pool),
27            Arc::clone(&provider),
28            config.clone(),
29        ));
30        Self {
31            config,
32            provider,
33            container_provider: None,
34            invoker,
35            status: Arc::new(AtomicU8::new(STATUS_STOPPED)),
36        }
37    }
38
39    pub fn with_fake_provider(
40        config: FunctionConfig,
41        provider: Arc<crate::provider::fake::FakeProvider>,
42    ) -> Self {
43        Self::new(config, provider as Arc<dyn FunctionProvider>)
44    }
45
46    pub fn with_container_provider(
47        config: FunctionConfig,
48        provider: crate::provider::container::ContainerProvider,
49    ) -> Self {
50        let arc = Arc::new(provider);
51        let mut svc = Self::new(config, arc.clone() as Arc<dyn FunctionProvider>);
52        svc.container_provider = Some(arc);
53        svc
54    }
55
56    pub fn with_default_container_provider(
57        config: FunctionConfig,
58    ) -> Result<Self, crate::provider::ProviderError> {
59        let provider = crate::provider::container::ContainerProvider::builder()
60            .egress_allowlist(config.egress_allowlist.clone())
61            .build()?;
62        Ok(Self::with_container_provider(config, provider))
63    }
64
65    pub fn invoker(&self) -> Arc<dyn FunctionInvoker> {
66        self.invoker.clone() as Arc<dyn FunctionInvoker>
67    }
68
69    pub fn provider(&self) -> Result<&crate::provider::container::ContainerProvider, CamelError> {
70        self.container_provider
71            .as_ref()
72            .map(|arc| arc.as_ref())
73            .ok_or_else(|| {
74                CamelError::Config("unsupported provider type: not a container provider".into())
75            })
76    }
77
78    pub fn runner_state(&self, runtime: &str) -> Option<RunnerState> {
79        let key = RunnerPoolKey {
80            runtime: runtime.to_string(),
81        };
82        self.invoker
83            .pool
84            .handles
85            .get(&key)
86            .map(|h| h.state.lock().expect("state").clone()) // allow-unwrap
87    }
88
89    pub fn force_runner_failed(&self, runtime: &str, reason: &str) {
90        let key = RunnerPoolKey {
91            runtime: runtime.to_string(),
92        };
93        if let Some(handle) = self.invoker.pool.handles.get(&key) {
94            *handle.state.lock().expect("state") = RunnerState::Failed {
95                // allow-unwrap
96                reason: reason.to_string(),
97            };
98        }
99    }
100
101    pub(crate) async fn wait_until_healthy(
102        &self,
103        handle: &crate::pool::RunnerHandle,
104    ) -> Result<(), ProviderError> {
105        let deadline = tokio::time::Instant::now() + self.config.boot_timeout;
106        loop {
107            if tokio::time::Instant::now() > deadline {
108                return Err(ProviderError::BootTimeout);
109            }
110            match self.provider.health(handle).await {
111                Ok(FunctionHealthStatus::Healthy) => {
112                    *handle.state.lock().expect("state") = RunnerState::Healthy; // allow-unwrap
113                    return Ok(());
114                }
115                Ok(FunctionHealthStatus::Unhealthy(reason)) => {
116                    *handle.state.lock().expect("state") = RunnerState::Unhealthy {
117                        // allow-unwrap
118                        since: std::time::Instant::now(),
119                        reason,
120                    };
121                    tokio::time::sleep(self.config.health_interval).await;
122                }
123                Err(_) => {
124                    tokio::time::sleep(self.config.health_interval).await;
125                }
126            }
127        }
128    }
129
130    pub(crate) fn spawn_health_task(&self, handle: crate::pool::RunnerHandle) {
131        let provider = Arc::clone(&self.provider);
132        let interval = self.config.health_interval;
133        tokio::spawn(async move {
134            let mut ticks = tokio::time::interval(interval);
135            let mut unhealthy_count = 0u8;
136            loop {
137                tokio::select! {
138                    _ = handle.cancel.cancelled() => break,
139                    _ = ticks.tick() => {
140                        match provider.health(&handle).await {
141                            Ok(FunctionHealthStatus::Healthy) => {
142                                unhealthy_count = 0;
143                                *handle.state.lock().expect("state") = RunnerState::Healthy; // allow-unwrap
144                            }
145                            Ok(FunctionHealthStatus::Unhealthy(reason)) => {
146                                unhealthy_count = unhealthy_count.saturating_add(1);
147                                if unhealthy_count >= 2 {
148                                    *handle.state.lock().expect("state") = RunnerState::Unhealthy { since: std::time::Instant::now(), reason }; // allow-unwrap
149                                }
150                            }
151                            Err(err) => {
152                                *handle.state.lock().expect("state") = RunnerState::Failed { reason: err.to_string() }; // allow-unwrap
153                            }
154                        }
155                    }
156                }
157            }
158        });
159    }
160
161    pub(crate) async fn rollback_start(
162        &self,
163        spawned: &[(RunnerPoolKey, crate::pool::RunnerHandle)],
164        registered_refs: &[((FunctionId, Option<String>), RunnerPoolKey)],
165        pending: &[(FunctionDefinition, Option<String>)],
166    ) {
167        for (ref_key, _pool_key) in registered_refs {
168            self.invoker.pool.ref_counts.remove(ref_key);
169            self.invoker.pool.function_to_key.remove(ref_key);
170            let still_used = self
171                .invoker
172                .pool
173                .function_to_key
174                .iter()
175                .any(|kv| kv.key().0 == ref_key.0);
176            if !still_used {
177                self.invoker
178                    .function_timeouts
179                    .lock()
180                    .expect("function_timeouts") // allow-unwrap
181                    .remove(&ref_key.0);
182            }
183        }
184        for (key, handle) in spawned {
185            self.invoker.pool.handles.remove(key);
186            handle.cancel.cancel();
187            let _ = self.provider.shutdown(handle.clone()).await;
188        }
189        self.invoker
190            .pending
191            .lock()
192            .expect("pending") // allow-unwrap
193            .extend(pending.iter().cloned());
194    }
195}
196
197#[async_trait::async_trait]
198impl Lifecycle for FunctionRuntimeService {
199    fn name(&self) -> &str {
200        "function-runtime"
201    }
202
203    fn status(&self) -> ServiceStatus {
204        match self.status.load(Ordering::SeqCst) {
205            STATUS_STOPPED => ServiceStatus::Stopped,
206            STATUS_STARTED => ServiceStatus::Started,
207            _ => ServiceStatus::Failed,
208        }
209    }
210
211    async fn start(&mut self) -> Result<(), CamelError> {
212        self.config.validate()?;
213        if self.status.load(Ordering::SeqCst) == STATUS_STARTED {
214            return Ok(());
215        }
216        let pending = {
217            let mut lock = self.invoker.pending.lock().expect("pending"); // allow-unwrap
218            std::mem::take(&mut *lock)
219        };
220        let mut grouped: std::collections::HashMap<
221            RunnerPoolKey,
222            Vec<(FunctionDefinition, Option<String>)>,
223        > = std::collections::HashMap::new();
224        for (def, route_id) in pending.iter().cloned() {
225            grouped
226                .entry(RunnerPoolKey {
227                    runtime: def.runtime.clone(),
228                })
229                .or_default()
230                .push((def, route_id));
231        }
232        let mut spawned: Vec<(RunnerPoolKey, crate::pool::RunnerHandle)> = Vec::new();
233        let mut registered_refs: Vec<((FunctionId, Option<String>), RunnerPoolKey)> = Vec::new();
234        for (key, defs) in grouped {
235            let handle = match self.provider.spawn(&key).await {
236                Ok(h) => h,
237                Err(e) => {
238                    self.rollback_start(&spawned, &registered_refs, &pending)
239                        .await;
240                    self.status.store(STATUS_FAILED, Ordering::SeqCst);
241                    return Err(CamelError::Config(format!("function: spawn failed: {e}")));
242                }
243            };
244            match self.wait_until_healthy(&handle).await {
245                Ok(()) => {}
246                Err(e) => {
247                    handle.cancel.cancel();
248                    let _ = self.provider.shutdown(handle).await;
249                    self.rollback_start(&spawned, &registered_refs, &pending)
250                        .await;
251                    self.status.store(STATUS_FAILED, Ordering::SeqCst);
252                    return Err(CamelError::Config(format!("function: boot timeout: {e}")));
253                }
254            }
255            self.invoker
256                .pool
257                .handles
258                .insert(key.clone(), handle.clone());
259            self.spawn_health_task(handle.clone());
260            spawned.push((key.clone(), handle.clone()));
261            for (def, route_id) in defs {
262                if let Err(err) = self.provider.register(&handle, &def).await {
263                    self.rollback_start(&spawned, &registered_refs, &pending)
264                        .await;
265                    self.status.store(STATUS_FAILED, Ordering::SeqCst);
266                    return Err(CamelError::Config(format!(
267                        "function: register failed: {err}"
268                    )));
269                }
270                let ref_key = (def.id.clone(), route_id.clone());
271                self.invoker.pool.ref_counts.insert(ref_key.clone(), 1);
272                self.invoker
273                    .pool
274                    .function_to_key
275                    .insert(ref_key.clone(), key.clone());
276                registered_refs.push((ref_key, key.clone()));
277            }
278        }
279        self.invoker.started.store(true, Ordering::SeqCst);
280        self.status.store(STATUS_STARTED, Ordering::SeqCst);
281        Ok(())
282    }
283
284    async fn stop(&mut self) -> Result<(), CamelError> {
285        let handles: Vec<_> = self
286            .invoker
287            .pool
288            .handles
289            .iter()
290            .map(|h| h.clone())
291            .collect();
292        self.invoker.pool.handles.clear();
293        self.invoker.pool.ref_counts.clear();
294        self.invoker.pool.function_to_key.clear();
295        self.invoker
296            .function_timeouts
297            .lock()
298            .expect("function_timeouts") // allow-unwrap
299            .clear();
300        let mut first_err: Option<ProviderError> = None;
301        for handle in handles {
302            handle.cancel.cancel();
303            if let Err(e) = self.provider.shutdown(handle).await
304                && first_err.is_none()
305            {
306                first_err = Some(e);
307            }
308        }
309        self.invoker.started.store(false, Ordering::SeqCst);
310        self.status.store(STATUS_STOPPED, Ordering::SeqCst);
311        if let Some(e) = first_err {
312            return Err(CamelError::ProcessorError(e.to_string()));
313        }
314        Ok(())
315    }
316
317    fn as_function_invoker(&self) -> Option<Arc<dyn FunctionInvoker>> {
318        Some(self.invoker.clone() as Arc<dyn FunctionInvoker>)
319    }
320}