Skip to main content

greentic_runner_host/
host.rs

1use std::collections::HashMap;
2use std::path::Path;
3use std::sync::Arc;
4
5use anyhow::{Context, Result, anyhow, bail};
6use serde_json::Value;
7
8use crate::activity::Activity;
9use crate::boot;
10use crate::config::HostConfig;
11use crate::engine::host::{SessionHost, StateHost};
12use crate::engine::runtime::IngressEnvelope;
13use crate::http::health::HealthState;
14use crate::pack::PackRuntime;
15use crate::runner::adapt_timer;
16use crate::runner::engine::FlowEngine;
17use crate::runtime::{ActivePacks, TenantRuntime};
18use crate::secrets::{DynSecretsManager, default_manager};
19use crate::storage::{
20    DynSessionStore, DynStateStore, new_session_store, new_state_store, session_host_from,
21    state_host_from,
22};
23use crate::wasi::RunnerWasiPolicy;
24
25#[derive(Clone, Debug)]
26pub struct TelemetryCfg {
27    pub config: greentic_telemetry::TelemetryConfig,
28    pub export: greentic_telemetry::export::ExportConfig,
29}
30
31/// Builder for composing multi-tenant host instances.
32pub struct HostBuilder {
33    configs: HashMap<String, HostConfig>,
34    telemetry: Option<TelemetryCfg>,
35    wasi_policy: RunnerWasiPolicy,
36    secrets: Option<DynSecretsManager>,
37}
38
39impl HostBuilder {
40    pub fn new() -> Self {
41        Self {
42            configs: HashMap::new(),
43            telemetry: None,
44            wasi_policy: RunnerWasiPolicy::default(),
45            secrets: None,
46        }
47    }
48
49    pub fn with_config(mut self, config: HostConfig) -> Self {
50        self.configs.insert(config.tenant.clone(), config);
51        self
52    }
53
54    pub fn with_telemetry(mut self, telemetry: TelemetryCfg) -> Self {
55        self.telemetry = Some(telemetry);
56        self
57    }
58
59    pub fn with_wasi_policy(mut self, policy: RunnerWasiPolicy) -> Self {
60        self.wasi_policy = policy;
61        self
62    }
63
64    pub fn with_secrets_manager(mut self, manager: DynSecretsManager) -> Self {
65        self.secrets = Some(manager);
66        self
67    }
68
69    pub fn build(self) -> Result<RunnerHost> {
70        if self.configs.is_empty() {
71            bail!("at least one tenant configuration is required");
72        }
73        let wasi_policy = Arc::new(self.wasi_policy);
74        let configs = self
75            .configs
76            .into_iter()
77            .map(|(tenant, cfg)| (tenant, Arc::new(cfg)))
78            .collect();
79        let session_store = new_session_store();
80        let session_host = session_host_from(Arc::clone(&session_store));
81        let state_store = new_state_store();
82        let state_host = state_host_from(Arc::clone(&state_store));
83        let secrets = match self.secrets {
84            Some(manager) => manager,
85            None => default_manager().context("failed to initialise default secrets backend")?,
86        };
87        Ok(RunnerHost {
88            configs,
89            active: Arc::new(ActivePacks::new()),
90            health: Arc::new(HealthState::new()),
91            session_store,
92            state_store,
93            session_host,
94            state_host,
95            wasi_policy,
96            secrets_manager: secrets,
97            telemetry: self.telemetry,
98        })
99    }
100}
101
102impl Default for HostBuilder {
103    fn default() -> Self {
104        Self::new()
105    }
106}
107
108/// Runtime host that manages tenant-bound packs and flow execution.
109pub struct RunnerHost {
110    configs: HashMap<String, Arc<HostConfig>>,
111    active: Arc<ActivePacks>,
112    health: Arc<HealthState>,
113    session_store: DynSessionStore,
114    state_store: DynStateStore,
115    session_host: Arc<dyn SessionHost>,
116    state_host: Arc<dyn StateHost>,
117    wasi_policy: Arc<RunnerWasiPolicy>,
118    secrets_manager: DynSecretsManager,
119    telemetry: Option<TelemetryCfg>,
120}
121
122/// Handle exposing tenant internals for embedding hosts (e.g. CLI server).
123#[derive(Clone)]
124pub struct TenantHandle {
125    runtime: Arc<TenantRuntime>,
126}
127
128impl RunnerHost {
129    pub async fn start(&self) -> Result<()> {
130        boot::init(&self.health, self.telemetry.as_ref())?;
131        Ok(())
132    }
133
134    pub async fn stop(&self) -> Result<()> {
135        self.active.replace(HashMap::new());
136        Ok(())
137    }
138
139    pub async fn load_pack(&self, tenant: &str, pack_path: &Path) -> Result<()> {
140        let archive_source = if is_pack_archive(pack_path) {
141            Some(pack_path)
142        } else {
143            None
144        };
145        let runtime = self
146            .prepare_runtime(tenant, pack_path, archive_source)
147            .await
148            .with_context(|| format!("failed to load tenant {tenant}"))?;
149        let mut next = (*self.active.snapshot()).clone();
150        next.insert(tenant.to_string(), runtime);
151        self.active.replace(next);
152        tracing::info!(tenant, pack = %pack_path.display(), "pack loaded");
153        Ok(())
154    }
155
156    pub async fn handle_activity(&self, tenant: &str, activity: Activity) -> Result<Vec<Activity>> {
157        let runtime = self
158            .active
159            .load(tenant)
160            .with_context(|| format!("tenant {tenant} not loaded"))?;
161        let (pack_id, flow_id) = resolve_flow_id(&runtime, &activity)?;
162        let action = activity.action().map(|value| value.to_string());
163        let session = activity.session_id().map(|value| value.to_string());
164        let provider = activity.provider_id().map(|value| value.to_string());
165        let channel = activity.channel().map(|value| value.to_string());
166        let conversation = activity.conversation().map(|value| value.to_string());
167        let user = activity.user().map(|value| value.to_string());
168        let flow_type = activity
169            .flow_type()
170            .map(|value| value.to_string())
171            .or_else(|| {
172                runtime
173                    .engine()
174                    .flow_by_key(&pack_id, &flow_id)
175                    .map(|desc| desc.flow_type.clone())
176            });
177        let payload = activity.into_payload();
178
179        let envelope = IngressEnvelope {
180            tenant: tenant.to_string(),
181            env: std::env::var("GREENTIC_ENV").ok(),
182            pack_id: Some(pack_id.clone()),
183            flow_id: flow_id.clone(),
184            flow_type,
185            action,
186            session_hint: session,
187            provider,
188            channel,
189            conversation,
190            user,
191            activity_id: None,
192            timestamp: None,
193            payload,
194            metadata: None,
195            reply_scope: None,
196        }
197        .canonicalize();
198
199        let result = runtime.state_machine().handle(envelope).await?;
200        Ok(normalize_replies(result, tenant))
201    }
202
203    pub async fn tenant(&self, tenant: &str) -> Option<TenantHandle> {
204        self.active
205            .load(tenant)
206            .map(|runtime| TenantHandle { runtime })
207    }
208
209    pub fn active_packs(&self) -> Arc<ActivePacks> {
210        Arc::clone(&self.active)
211    }
212
213    pub fn health_state(&self) -> Arc<HealthState> {
214        Arc::clone(&self.health)
215    }
216
217    pub fn wasi_policy(&self) -> Arc<RunnerWasiPolicy> {
218        Arc::clone(&self.wasi_policy)
219    }
220
221    pub fn session_store(&self) -> DynSessionStore {
222        Arc::clone(&self.session_store)
223    }
224
225    pub fn state_store(&self) -> DynStateStore {
226        Arc::clone(&self.state_store)
227    }
228
229    pub fn session_host(&self) -> Arc<dyn SessionHost> {
230        Arc::clone(&self.session_host)
231    }
232
233    pub fn state_host(&self) -> Arc<dyn StateHost> {
234        Arc::clone(&self.state_host)
235    }
236
237    pub fn secrets_manager(&self) -> DynSecretsManager {
238        Arc::clone(&self.secrets_manager)
239    }
240
241    pub fn tenant_configs(&self) -> HashMap<String, Arc<HostConfig>> {
242        self.configs.clone()
243    }
244
245    async fn prepare_runtime(
246        &self,
247        tenant: &str,
248        pack_path: &Path,
249        archive_source: Option<&Path>,
250    ) -> Result<Arc<TenantRuntime>> {
251        let config = self
252            .configs
253            .get(tenant)
254            .cloned()
255            .with_context(|| format!("tenant {tenant} not registered"))?;
256        if config.tenant != tenant {
257            bail!(
258                "tenant mismatch: config declares '{}' but '{tenant}' was requested",
259                config.tenant
260            );
261        }
262        let runtime = TenantRuntime::load(
263            pack_path,
264            Arc::clone(&config),
265            None,
266            archive_source,
267            None,
268            self.wasi_policy(),
269            self.session_host(),
270            self.session_store(),
271            self.state_store(),
272            self.state_host(),
273            self.secrets_manager(),
274        )
275        .await?;
276        let timers = adapt_timer::spawn_timers(Arc::clone(&runtime))?;
277        runtime.register_timers(timers);
278        Ok(runtime)
279    }
280}
281
282impl TenantHandle {
283    pub fn config(&self) -> Arc<HostConfig> {
284        Arc::clone(self.runtime.config())
285    }
286
287    pub fn pack(&self) -> Arc<PackRuntime> {
288        self.runtime.pack()
289    }
290
291    pub fn engine(&self) -> Arc<FlowEngine> {
292        Arc::clone(self.runtime.engine())
293    }
294
295    pub fn overlays(&self) -> Vec<Arc<PackRuntime>> {
296        self.runtime.overlays()
297    }
298
299    pub fn overlay_digests(&self) -> Vec<Option<String>> {
300        self.runtime.overlay_digests()
301    }
302}
303
304fn resolve_flow_id(runtime: &TenantRuntime, activity: &Activity) -> Result<(String, String)> {
305    let engine = runtime.engine();
306    if let Some(flow_id) = activity.flow_id() {
307        if let Some(pack_id) = activity.pack_id() {
308            if engine.flow_by_key(pack_id, flow_id).is_none() {
309                bail!("flow {flow_id} not registered for pack {pack_id}");
310            }
311            return Ok((pack_id.to_string(), flow_id.to_string()));
312        }
313        if let Some(flow) = engine.flow_by_id(flow_id) {
314            return Ok((flow.pack_id.clone(), flow.id.clone()));
315        }
316        bail!("flow {flow_id} is ambiguous; pack_id is required");
317    }
318
319    if let Some(flow_type) = activity.flow_type() {
320        if let Some(pack_id) = activity.pack_id() {
321            if let Some(flow) = engine
322                .flows()
323                .iter()
324                .find(|flow| flow.pack_id == pack_id && flow.flow_type == flow_type)
325            {
326                return Ok((pack_id.to_string(), flow.id.clone()));
327            }
328            bail!("flow type {flow_type} not registered for pack {pack_id}");
329        }
330        if let Some(flow) = engine.flow_by_type(flow_type) {
331            return Ok((flow.pack_id.clone(), flow.id.clone()));
332        }
333        bail!("flow type {flow_type} is ambiguous; pack_id is required");
334    }
335
336    let pack = runtime.pack();
337    let flow_id = pack
338        .metadata()
339        .entry_flows
340        .first()
341        .cloned()
342        .ok_or_else(|| anyhow!("no entry flows registered for tenant {}", runtime.tenant()))?;
343    Ok((pack.metadata().pack_id.clone(), flow_id))
344}
345
346fn normalize_replies(result: Value, tenant: &str) -> Vec<Activity> {
347    result
348        .as_array()
349        .cloned()
350        .unwrap_or_else(|| vec![result])
351        .into_iter()
352        .map(|payload| Activity::from_output(payload, tenant))
353        .collect()
354}
355
356fn is_pack_archive(path: &Path) -> bool {
357    path.extension()
358        .and_then(|ext| ext.to_str())
359        .map(|ext| ext.eq_ignore_ascii_case("gtpack"))
360        .unwrap_or(false)
361}