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()) }
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 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; return Ok(());
114 }
115 Ok(FunctionHealthStatus::Unhealthy(reason)) => {
116 *handle.state.lock().expect("state") = RunnerState::Unhealthy {
117 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; }
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 }; }
150 }
151 Err(err) => {
152 *handle.state.lock().expect("state") = RunnerState::Failed { reason: err.to_string() }; }
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") .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") .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"); 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, ®istered_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, ®istered_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, ®istered_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") .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}