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
31pub 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
108pub 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#[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}