1#![deny(unsafe_code)]
2use std::collections::{HashMap, HashSet};
10use std::fs;
11use std::path::PathBuf;
12use std::sync::Arc;
13use std::time::Duration;
14
15use crate::secrets::SecretsBackend;
16use anyhow::{Context, Result, anyhow};
17use greentic_config::ResolvedConfig;
18use greentic_config_types::TelemetryExporterKind;
19use greentic_config_types::{
20 NetworkConfig, PackSourceConfig, PacksConfig, PathsConfig, TelemetryConfig,
21};
22use greentic_telemetry::export::{ExportConfig as TelemetryExportConfig, ExportMode, Sampling};
23use runner_core::env::PackConfig;
24use serde_json::json;
25use tokio::signal;
26
27pub mod boot;
28pub mod cache;
29pub mod component_api;
30pub mod config;
31pub mod engine;
32pub mod fault;
33pub mod gtbind;
34pub mod http;
35pub mod metrics;
36pub mod operator_metrics;
37pub mod operator_registry;
38pub mod pack;
39pub mod provider;
40pub mod provider_core;
41pub mod provider_core_only;
42pub mod routing;
43pub mod runner;
44pub mod runtime;
45pub mod runtime_wasmtime;
46pub mod secrets;
47pub mod storage;
48pub mod telemetry;
49pub mod telemetry_scan;
50#[cfg(feature = "fault-injection")]
51pub mod testing;
52pub mod trace;
53pub mod validate;
54pub mod verify;
55pub mod wasi;
56pub mod watcher;
57
58mod activity;
59mod host;
60pub mod oauth;
61
62pub use activity::{Activity, ActivityKind};
63pub use config::HostConfig;
64pub use gtbind::{PackBinding, TenantBindings};
65pub use host::TelemetryCfg;
66pub use host::{HostBuilder, RunnerHost, TenantHandle};
67pub use wasi::{PreopenSpec, RunnerWasiPolicy};
68
69pub use greentic_types::{EnvId, FlowId, PackId, TenantCtx, TenantId};
70
71pub use http::auth::AdminAuth;
72pub use routing::RoutingConfig;
73use routing::TenantRouting;
74pub use runner::HostServer;
75
76#[cfg(test)]
77pub(crate) mod test_support {
78 use super::*;
79 use crate::config::{OperatorPolicy, SecretsPolicy};
80 use crate::runtime::TenantRuntime;
81 use crate::secrets::default_manager;
82 use crate::storage::{new_session_store, new_state_store, session_host_from, state_host_from};
83 use crate::trace::TraceConfig;
84 use crate::validate::ValidationConfig;
85 use tempfile::TempDir;
86
87 pub(crate) fn fixture_pack_path() -> PathBuf {
88 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
89 .join("../../examples/packs/demo.gtpack")
90 .canonicalize()
91 .expect("fixture pack path")
92 }
93
94 fn minimal_config(workspace: &std::path::Path) -> Result<Arc<HostConfig>> {
95 let bindings_path = workspace.join("bindings.yaml");
96 std::fs::write(
97 &bindings_path,
98 r#"
99tenant: demo
100flow_type_bindings: {}
101rate_limits: {}
102retry: {}
103timers: []
104"#,
105 )?;
106 let mut config =
107 HostConfig::load_from_path(&bindings_path).context("load minimal host bindings")?;
108 config.secrets_policy = SecretsPolicy::allow_all();
109 config.operator_policy = OperatorPolicy::allow_all();
110 config.trace = TraceConfig::from_env();
111 config.validation = ValidationConfig::from_env();
112 Ok(Arc::new(config))
113 }
114
115 pub(crate) async fn build_test_runtime() -> Result<(TempDir, Arc<TenantRuntime>)> {
116 let workspace = TempDir::new().context("temp workspace")?;
117 let config = minimal_config(workspace.path())?;
118 let session_store = new_session_store();
119 let session_host = session_host_from(Arc::clone(&session_store));
120 let state_store = new_state_store();
121 let state_host = state_host_from(Arc::clone(&state_store));
122 let secrets = default_manager()?;
123 let pack_path = fixture_pack_path();
124 let runtime = TenantRuntime::load(
125 &pack_path,
126 config,
127 None,
128 Some(&pack_path),
129 None,
130 Arc::new(RunnerWasiPolicy::new()),
131 session_host,
132 Arc::clone(&session_store),
133 Arc::clone(&state_store),
134 state_host,
135 secrets,
136 )
137 .await?;
138 Ok((workspace, runtime))
139 }
140}
141
142#[derive(Clone)]
144pub struct RunnerConfig {
145 pub tenant_bindings: HashMap<String, TenantBindings>,
146 pub pack: PackConfig,
147 pub port: u16,
148 pub refresh_interval: Duration,
149 pub routing: RoutingConfig,
150 pub admin: AdminAuth,
151 pub telemetry: Option<TelemetryCfg>,
152 pub secrets_backend: SecretsBackend,
153 pub wasi_policy: RunnerWasiPolicy,
154 pub resolved_config: ResolvedConfig,
155 pub trace: trace::TraceConfig,
156 pub validation: validate::ValidationConfig,
157}
158
159impl RunnerConfig {
160 pub fn from_config(resolved_config: ResolvedConfig, bindings: Vec<PathBuf>) -> Result<Self> {
162 if bindings.is_empty() {
163 anyhow::bail!("at least one gtbind file is required");
164 }
165 let tenant_bindings = gtbind::load_gtbinds(&bindings)?;
166 if tenant_bindings.is_empty() {
167 anyhow::bail!("no gtbind files loaded");
168 }
169 let mut pack = pack_config_from(
170 &resolved_config.config.packs,
171 &resolved_config.config.paths,
172 &resolved_config.config.network,
173 )?;
174 maybe_write_gtbind_index(&tenant_bindings, &resolved_config.config.paths, &mut pack)?;
175 let refresh = parse_refresh_interval(std::env::var("PACK_REFRESH_INTERVAL").ok())?;
176 let port = std::env::var("PORT")
177 .ok()
178 .and_then(|value| value.parse().ok())
179 .unwrap_or(8080);
180 let default_tenant = resolved_config
181 .config
182 .dev
183 .as_ref()
184 .map(|dev| dev.default_tenant.clone())
185 .unwrap_or_else(|| "demo".into());
186 let routing = RoutingConfig::from_env_with_default(default_tenant);
187 let paths = &resolved_config.config.paths;
188 ensure_paths_exist(paths)?;
189 let mut wasi_policy = default_wasi_policy(paths);
190 let mut env_allow = HashSet::new();
191 for binding in tenant_bindings.values() {
192 env_allow.extend(binding.env_passthrough.iter().cloned());
193 }
194 for key in env_allow {
195 wasi_policy = wasi_policy.allow_env(key);
196 }
197
198 let admin = AdminAuth::new(resolved_config.config.services.as_ref().and_then(|s| {
199 s.events
200 .as_ref()
201 .and_then(|svc| svc.headers.as_ref())
202 .and_then(|headers| headers.get("x-admin-token").cloned())
203 }));
204 let secrets_backend = SecretsBackend::from_config(&resolved_config.config.secrets)?;
205 Ok(Self {
206 tenant_bindings,
207 pack,
208 port,
209 refresh_interval: refresh,
210 routing,
211 admin,
212 telemetry: telemetry_from(&resolved_config.config.telemetry),
213 secrets_backend,
214 wasi_policy,
215 resolved_config,
216 trace: trace::TraceConfig::from_env(),
217 validation: validate::ValidationConfig::from_env(),
218 })
219 }
220
221 pub fn with_port(mut self, port: u16) -> Self {
223 self.port = port;
224 self
225 }
226
227 pub fn with_wasi_policy(mut self, policy: RunnerWasiPolicy) -> Self {
228 self.wasi_policy = policy;
229 self
230 }
231}
232
233fn maybe_write_gtbind_index(
234 tenant_bindings: &HashMap<String, TenantBindings>,
235 paths: &PathsConfig,
236 pack: &mut PackConfig,
237) -> Result<()> {
238 let mut uses_locators = false;
239 for binding in tenant_bindings.values() {
240 for pack_binding in &binding.packs {
241 if pack_binding.pack_locator.is_some() {
242 uses_locators = true;
243 }
244 }
245 }
246 if !uses_locators {
247 return Ok(());
248 }
249
250 let mut entries = serde_json::Map::new();
251 for binding in tenant_bindings.values() {
252 let mut packs = Vec::new();
253 for pack_binding in &binding.packs {
254 let locator = pack_binding.pack_locator.as_ref().ok_or_else(|| {
255 anyhow::anyhow!(
256 "gtbind {} missing pack_locator for pack {}",
257 binding.tenant,
258 pack_binding.pack_id
259 )
260 })?;
261 let (name, version_or_digest) =
262 pack_binding.pack_ref.split_once('@').ok_or_else(|| {
263 anyhow::anyhow!(
264 "gtbind {} invalid pack_ref {} (expected name@version)",
265 binding.tenant,
266 pack_binding.pack_ref
267 )
268 })?;
269 if name != pack_binding.pack_id {
270 anyhow::bail!(
271 "gtbind {} pack_ref {} does not match pack_id {}",
272 binding.tenant,
273 pack_binding.pack_ref,
274 pack_binding.pack_id
275 );
276 }
277 let mut entry = serde_json::Map::new();
278 entry.insert("name".to_string(), json!(name));
279 if version_or_digest.contains(':') {
280 entry.insert("digest".to_string(), json!(version_or_digest));
281 } else {
282 entry.insert("version".to_string(), json!(version_or_digest));
283 }
284 entry.insert("locator".to_string(), json!(locator));
285 packs.push(serde_json::Value::Object(entry));
286 }
287 let main_pack = packs
288 .first()
289 .cloned()
290 .ok_or_else(|| anyhow::anyhow!("gtbind {} has no packs", binding.tenant))?;
291 let overlays = packs.into_iter().skip(1).collect::<Vec<_>>();
292 entries.insert(
293 binding.tenant.clone(),
294 json!({
295 "main_pack": main_pack,
296 "overlays": overlays,
297 }),
298 );
299 }
300
301 let index_path = paths.greentic_root.join("packs").join("gtbind.index.json");
302 if let Some(parent) = index_path.parent() {
303 fs::create_dir_all(parent)
304 .with_context(|| format!("failed to create {}", parent.display()))?;
305 }
306 let serialized = serde_json::to_vec_pretty(&serde_json::Value::Object(entries))?;
307 fs::write(&index_path, serialized)
308 .with_context(|| format!("failed to write {}", index_path.display()))?;
309 pack.index_location = runner_core::env::IndexLocation::File(index_path);
310 Ok(())
311}
312
313fn parse_refresh_interval(value: Option<String>) -> Result<Duration> {
314 let raw = value.unwrap_or_else(|| "30s".into());
315 humantime::parse_duration(&raw).map_err(|err| anyhow!("invalid PACK_REFRESH_INTERVAL: {err}"))
316}
317
318fn default_wasi_policy(paths: &PathsConfig) -> RunnerWasiPolicy {
319 let mut policy = RunnerWasiPolicy::default()
320 .with_env("GREENTIC_ROOT", paths.greentic_root.display().to_string())
321 .with_env("GREENTIC_STATE_DIR", paths.state_dir.display().to_string())
322 .with_env("GREENTIC_CACHE_DIR", paths.cache_dir.display().to_string())
323 .with_env("GREENTIC_LOGS_DIR", paths.logs_dir.display().to_string());
324 policy = policy
325 .with_preopen(PreopenSpec::new(&paths.state_dir, "/state"))
326 .with_preopen(PreopenSpec::new(&paths.cache_dir, "/cache"))
327 .with_preopen(PreopenSpec::new(&paths.logs_dir, "/logs"));
328 policy
329}
330
331fn ensure_paths_exist(paths: &PathsConfig) -> Result<()> {
332 for dir in [
333 &paths.greentic_root,
334 &paths.state_dir,
335 &paths.cache_dir,
336 &paths.logs_dir,
337 ] {
338 fs::create_dir_all(dir)
339 .with_context(|| format!("failed to ensure directory {}", dir.display()))?;
340 }
341 Ok(())
342}
343
344fn pack_config_from(
345 packs: &Option<PacksConfig>,
346 paths: &PathsConfig,
347 network: &NetworkConfig,
348) -> Result<PackConfig> {
349 if let Some(cfg) = packs {
350 let cache_dir = cfg.cache_dir.clone();
351 let index_location = match &cfg.source {
352 PackSourceConfig::LocalIndex { path } => {
353 runner_core::env::IndexLocation::File(path.clone())
354 }
355 PackSourceConfig::HttpIndex { url } => {
356 runner_core::env::IndexLocation::from_value(url)?
357 }
358 PackSourceConfig::OciRegistry { reference } => {
359 runner_core::env::IndexLocation::from_value(reference)?
360 }
361 };
362 let public_key = cfg
363 .trust
364 .as_ref()
365 .and_then(|trust| trust.public_keys.first().cloned());
366 return Ok(PackConfig {
367 source: runner_core::env::PackSource::Fs,
368 index_location,
369 cache_dir,
370 public_key,
371 network: Some(network.clone()),
372 });
373 }
374 let mut cfg = PackConfig::default_for_paths(paths)?;
375 cfg.network = Some(network.clone());
376 Ok(cfg)
377}
378
379fn telemetry_from(cfg: &TelemetryConfig) -> Option<TelemetryCfg> {
380 if !cfg.enabled || matches!(cfg.exporter, TelemetryExporterKind::None) {
381 return None;
382 }
383 let mut export = TelemetryExportConfig::json_default();
384 export.mode = match cfg.exporter {
385 TelemetryExporterKind::Otlp => ExportMode::OtlpGrpc,
386 TelemetryExporterKind::Stdout => ExportMode::JsonStdout,
387 TelemetryExporterKind::Gcp => ExportMode::GcpCloudTrace,
388 TelemetryExporterKind::Azure => ExportMode::AzureAppInsights,
389 TelemetryExporterKind::Aws => ExportMode::AwsXRay,
390 TelemetryExporterKind::None => return None,
391 };
392 export.endpoint = cfg.endpoint.clone();
393 export.sampling = Sampling::TraceIdRatio(cfg.sampling as f64);
394 Some(TelemetryCfg {
395 config: greentic_telemetry::TelemetryConfig {
396 service_name: "greentic-runner".into(),
397 },
398 export,
399 })
400}
401
402pub async fn run(cfg: RunnerConfig) -> Result<()> {
404 let RunnerConfig {
405 tenant_bindings,
406 pack,
407 port,
408 refresh_interval,
409 routing,
410 admin,
411 telemetry,
412 secrets_backend,
413 wasi_policy,
414 resolved_config: _resolved_config,
415 trace,
416 validation,
417 } = cfg;
418
419 let mut builder = HostBuilder::new();
420 for bindings in tenant_bindings.into_values() {
421 let mut host_config = HostConfig::from_gtbind(bindings);
422 host_config.trace = trace.clone();
423 host_config.validation = validation.clone();
424 builder = builder.with_config(host_config);
425 }
426 if let Some(telemetry) = telemetry.clone() {
427 builder = builder.with_telemetry(telemetry);
428 }
429 builder = builder
430 .with_wasi_policy(wasi_policy.clone())
431 .with_secrets_manager(
432 secrets_backend
433 .build_manager()
434 .context("failed to initialise secrets backend")?,
435 );
436
437 let host = Arc::new(builder.build()?);
438 host.start().await?;
439
440 let (watcher, reload_handle) =
441 watcher::start_pack_watcher(Arc::clone(&host), pack.clone(), refresh_interval).await?;
442
443 let routing = TenantRouting::new(routing.clone());
444 let server = HostServer::new(
445 port,
446 host.active_packs(),
447 routing,
448 host.health_state(),
449 Some(reload_handle),
450 admin.clone(),
451 )?;
452
453 tokio::select! {
454 result = server.serve() => {
455 result?;
456 }
457 _ = signal::ctrl_c() => {
458 tracing::info!("received shutdown signal");
459 }
460 }
461
462 drop(watcher);
463 host.stop().await?;
464 Ok(())
465}
466
467#[cfg(test)]
468mod tests {
469 use super::*;
470 use crate::gtbind::{PackBinding, TenantBindings};
471 use greentic_config_types::PackTrustConfig;
472 use tempfile::TempDir;
473
474 fn paths(temp: &TempDir) -> PathsConfig {
475 PathsConfig {
476 greentic_root: temp.path().join("greentic"),
477 state_dir: temp.path().join("state"),
478 cache_dir: temp.path().join("cache"),
479 logs_dir: temp.path().join("logs"),
480 }
481 }
482
483 #[test]
484 fn parse_refresh_interval_uses_default_and_rejects_invalid_values() {
485 assert_eq!(
486 parse_refresh_interval(None).expect("default interval"),
487 Duration::from_secs(30)
488 );
489 assert!(parse_refresh_interval(Some("not-a-duration".into())).is_err());
490 }
491
492 #[test]
493 fn ensure_paths_exist_creates_expected_directories() {
494 let temp = TempDir::new().expect("tempdir");
495 let paths = paths(&temp);
496 ensure_paths_exist(&paths).expect("create directories");
497
498 assert!(paths.greentic_root.is_dir());
499 assert!(paths.state_dir.is_dir());
500 assert!(paths.cache_dir.is_dir());
501 assert!(paths.logs_dir.is_dir());
502 }
503
504 #[test]
505 fn maybe_write_gtbind_index_writes_locator_index() {
506 let temp = TempDir::new().expect("tempdir");
507 let paths = paths(&temp);
508 fs::create_dir_all(&paths.greentic_root).expect("greentic root");
509 let mut pack = PackConfig::default_for_paths(&paths).expect("pack config");
510 let tenant_bindings = HashMap::from([(
511 "demo".to_string(),
512 TenantBindings {
513 tenant: "demo".into(),
514 packs: vec![
515 PackBinding {
516 pack_id: "pack.main".into(),
517 pack_ref: "pack.main@1.0.0".into(),
518 pack_locator: Some("fs:///packs/main.gtpack".into()),
519 flows: vec!["main".into()],
520 },
521 PackBinding {
522 pack_id: "pack.overlay".into(),
523 pack_ref: "pack.overlay@sha256:abcd".into(),
524 pack_locator: Some("fs:///packs/overlay.gtpack".into()),
525 flows: vec![],
526 },
527 ],
528 env_passthrough: vec![],
529 },
530 )]);
531
532 maybe_write_gtbind_index(&tenant_bindings, &paths, &mut pack).expect("write gtbind index");
533
534 let index_path = paths.greentic_root.join("packs").join("gtbind.index.json");
535 let json: serde_json::Value =
536 serde_json::from_slice(&fs::read(&index_path).expect("read index")).expect("json");
537 assert_eq!(
538 json["demo"]["main_pack"]["locator"],
539 "fs:///packs/main.gtpack"
540 );
541 assert_eq!(json["demo"]["overlays"][0]["digest"], "sha256:abcd");
542 match pack.index_location {
543 runner_core::env::IndexLocation::File(path) => assert_eq!(path, index_path),
544 runner_core::env::IndexLocation::Remote(_) => panic!("expected generated file index"),
545 }
546 }
547
548 #[test]
549 fn pack_config_from_prefers_configured_pack_source() {
550 let temp = TempDir::new().expect("tempdir");
551 let paths = paths(&temp);
552 let packs = Some(PacksConfig {
553 source: PackSourceConfig::HttpIndex {
554 url: "https://example.com/index.json".into(),
555 },
556 cache_dir: temp.path().join("packs-cache"),
557 index_cache_ttl_secs: None,
558 trust: Some(PackTrustConfig {
559 public_keys: vec!["ed25519:test".into()],
560 require_signatures: true,
561 }),
562 });
563 let config = pack_config_from(&packs, &paths, &NetworkConfig::default())
564 .expect("pack config from packs");
565
566 match config.index_location {
567 runner_core::env::IndexLocation::Remote(url) => {
568 assert_eq!(url.as_str(), "https://example.com/index.json");
569 }
570 runner_core::env::IndexLocation::File(_) => panic!("expected remote index"),
571 }
572 assert_eq!(config.public_key.as_deref(), Some("ed25519:test"));
573 assert!(config.network.is_some());
574 }
575
576 #[test]
577 fn default_wasi_policy_sets_expected_env_and_preopens() {
578 let temp = TempDir::new().expect("tempdir");
579 let paths = paths(&temp);
580 let policy = default_wasi_policy(&paths);
581
582 assert_eq!(
583 policy.env_set.get("GREENTIC_ROOT"),
584 Some(&paths.greentic_root.display().to_string())
585 );
586 assert_eq!(policy.preopens.len(), 3);
587 assert_eq!(policy.preopens[0].guest_path, "/state");
588 assert_eq!(policy.preopens[1].guest_path, "/cache");
589 assert_eq!(policy.preopens[2].guest_path, "/logs");
590 }
591
592 #[test]
593 fn telemetry_from_is_disabled_without_exporter() {
594 assert!(telemetry_from(&TelemetryConfig::default()).is_none());
595 }
596}