1#![cfg_attr(docsrs, feature(doc_cfg))]
2
3pub mod auth_catalog;
15pub mod cli;
16pub mod commands;
17pub mod compose;
18pub mod config;
19pub mod env_config;
20pub mod env_loader;
21pub mod error;
22pub mod executor;
23pub mod expand;
24pub mod init_template;
25pub mod interpolate;
26#[cfg(feature = "lineage")]
27pub mod lineage_glue;
28pub mod merge;
29pub mod obs;
30pub mod registry;
31#[cfg(feature = "schedule")]
32pub mod schedule;
33pub mod secrets;
34#[cfg(feature = "serve")]
35pub mod serve;
36pub mod state;
37pub mod transforms;
38
39pub use error::{CliError, CliResult};
40
41pub async fn run_from_yaml_str(yaml: &str) -> CliResult<executor::RunSummary> {
48 let interpolated = interpolate::interpolate(yaml)?;
51 let mut cfg: config::PipelineConfig =
52 serde_yaml::from_str(&interpolated).map_err(|e| CliError::ParseConfig {
53 path: std::path::PathBuf::from("<yaml-string>"),
54 message: e.to_string(),
55 })?;
56 if cfg.version != 1 {
57 return Err(CliError::ParseConfig {
58 path: std::path::PathBuf::from("<yaml-string>"),
59 message: format!(
60 "unsupported pipeline version {}, only version 1 is recognised",
61 cfg.version
62 ),
63 });
64 }
65 crate::secrets::resolve_secrets(&mut cfg).await?;
66 let pipeline_name = cfg.name.clone().unwrap_or_else(|| "unnamed".to_string());
67 let auth = auth_catalog::build_auth_catalog(cfg.auth.as_ref())?;
68 let nodes = expand::expand(&cfg)?;
69 executor::run_expanded(
70 nodes,
71 executor::ExecuteOptions {
72 pipeline_name,
73 execution: cfg.execution.clone(),
74 dry_run: false,
75 limit: None,
76 state_path_override: None,
77 auth,
78 clock: chrono::Utc::now().fixed_offset(),
79 cancel: None,
80 #[cfg(feature = "lineage")]
81 lineage: None,
82 #[cfg(feature = "lineage")]
83 lineage_cfg: None,
84 },
85 )
86 .await
87}