#![cfg_attr(docsrs, feature(doc_cfg))]
pub mod auth_catalog;
#[cfg(feature = "catalog")]
pub mod catalog;
pub mod cli;
pub mod commands;
pub mod compose;
pub mod config;
pub mod dlq_replay;
pub mod env_config;
pub mod env_loader;
pub mod error;
pub mod executor;
pub mod expand;
pub mod init_template;
pub mod interpolate;
#[cfg(feature = "lineage")]
pub mod lineage_glue;
pub mod merge;
#[cfg(feature = "notify")]
pub mod notify;
pub mod obs;
pub mod pipeline_test;
pub mod registry;
pub mod replication;
#[cfg(feature = "schedule")]
pub mod schedule;
pub mod secrets;
#[cfg(feature = "serve")]
pub mod serve;
pub mod sla;
pub mod state;
pub mod transforms;
pub use error::{CliError, CliResult};
pub async fn run_from_yaml_str(yaml: &str) -> CliResult<executor::RunSummary> {
let mut value: serde_json::Value =
serde_yaml::from_str(yaml).map_err(|e| CliError::ParseConfig {
path: std::path::PathBuf::from("<yaml-string>"),
message: e.to_string(),
})?;
interpolate::interpolate_value(&mut value)?;
let interpolated = serde_yaml::to_string(&value).map_err(|e| CliError::ParseConfig {
path: std::path::PathBuf::from("<yaml-string>"),
message: e.to_string(),
})?;
let mut cfg: config::PipelineConfig =
serde_yaml::from_str(&interpolated).map_err(|e| CliError::ParseConfig {
path: std::path::PathBuf::from("<yaml-string>"),
message: e.to_string(),
})?;
if cfg.version != 1 {
return Err(CliError::ParseConfig {
path: std::path::PathBuf::from("<yaml-string>"),
message: format!(
"unsupported pipeline version {}, only version 1 is recognised",
cfg.version
),
});
}
crate::secrets::resolve_secrets(&mut cfg).await?;
let pipeline_name = cfg.name.clone().unwrap_or_else(|| "unnamed".to_string());
let auth = auth_catalog::build_auth_catalog(cfg.auth.as_ref())?;
let resilience = match &cfg.resilience {
Some(spec) => Some(spec.to_policy()?),
None => None,
};
#[cfg(feature = "catalog")]
let catalog = match cfg.catalog.as_ref() {
Some(spec) => Some(catalog::connect_from_spec(spec).await?),
None => None,
};
let nodes = expand::expand(&cfg)?;
executor::run_expanded(
nodes,
executor::ExecuteOptions {
pipeline_name,
execution: cfg.execution.clone(),
dry_run: false,
limit: None,
state_path_override: None,
shard: None,
auth,
clock: chrono::Utc::now().fixed_offset(),
cancel: None,
resilience,
sla: cfg.sla.clone(),
#[cfg(feature = "lineage")]
lineage: None,
#[cfg(feature = "lineage")]
lineage_cfg: None,
#[cfg(feature = "notify")]
notifier: None,
#[cfg(feature = "catalog")]
catalog,
},
)
.await
}