use orion::config;
#[derive(Clone, Copy, clap::ValueEnum)]
pub(crate) enum ConfigFormat {
Toml,
Json,
Summary,
}
pub(crate) fn handle_validate_config(
config: &config::AppConfig,
format: ConfigFormat,
) -> Result<(), Box<dyn std::error::Error>> {
match format {
ConfigFormat::Summary => print_config_summary(config),
ConfigFormat::Toml => {
eprintln!("Configuration is valid.");
let masked = masked_effective_config(config)?;
print!(
"{}",
toml::to_string_pretty(&toml::Value::try_from(&masked)?)?
);
}
ConfigFormat::Json => {
eprintln!("Configuration is valid.");
let masked = masked_effective_config(config)?;
println!("{}", serde_json::to_string_pretty(&masked)?);
}
}
Ok(())
}
fn masked_effective_config(
config: &config::AppConfig,
) -> Result<serde_json::Value, Box<dyn std::error::Error>> {
let without_unset = toml::Value::try_from(config)?;
let mut tree = serde_json::to_value(&without_unset)?;
orion::connector::mask_secrets(&mut tree);
Ok(tree)
}
fn redacted(value: &str) -> String {
orion::connector::redact_url_secrets_or_raw(value)
}
fn print_config_summary(config: &config::AppConfig) {
println!("Configuration is valid.\n");
println!(" environment: {}", config.environment);
println!(
" server: {}:{}",
config.server.host, config.server.port
);
println!(
" tls: {}",
if config.server.tls.enabled {
format!("enabled (cert={})", config.server.tls.cert_path)
} else {
"disabled".to_string()
}
);
println!(" storage: {}", redacted(&config.storage.url));
println!(
" logging: level={}, format={}",
config.logging.level,
match config.logging.format {
config::LogFormat::Json => "json",
config::LogFormat::Pretty => "pretty",
}
);
println!(
" admin_auth: {}",
if config.admin_auth.enabled {
"enabled"
} else {
"disabled"
}
);
println!(
" cors: {}{}",
config.cors.allowed_origins.join(", "),
if config.cors.allow_credentials {
" (credentials allowed)"
} else {
""
}
);
println!(
" rate_limiting: {}",
if config.rate_limit.enabled {
format!(
"enabled (rps={}, burst={})",
config.rate_limit.default_rps, config.rate_limit.default_burst
)
} else {
"disabled".to_string()
}
);
println!(
" queue: workers={}, buffer={}",
config.trace_queue.workers, config.trace_queue.buffer_size
);
println!(
" metrics: {}",
if config.metrics.enabled {
"enabled"
} else {
"disabled"
}
);
println!(
" tracing: {}",
if config.tracing.enabled {
format!("enabled (endpoint={})", config.tracing.otlp_endpoint)
} else {
"disabled".to_string()
}
);
println!(
" cluster: {}",
if config.cluster.enabled {
format!("enabled (instance_id={})", config.cluster.instance_id)
} else {
"disabled".to_string()
}
);
println!(
" kafka: {}",
if config.kafka.enabled {
let brokers: Vec<String> = config.kafka.brokers.iter().map(|b| redacted(b)).collect();
format!("enabled (brokers={})", brokers.join(","))
} else {
"disabled".to_string()
}
);
}
pub(crate) async fn handle_migrate(
config: &config::AppConfig,
dry_run: bool,
) -> Result<(), Box<dyn std::error::Error>> {
let pool = orion::storage::init_pool_no_migrate(&config.storage).await?;
let backend = orion::storage::get_backend();
let pending = orion::storage::pending_migrations(&pool).await?;
if pending.is_empty() {
println!("No pending migrations ({backend}).");
return Ok(());
}
if dry_run {
println!("Pending migrations on {backend} ({}):", pending.len());
} else {
println!("Applying {} migration(s) on {backend}...", pending.len());
}
for (version, description) in &pending {
println!(" {backend} {version:03} — {description}");
}
if dry_run {
println!(
"\nMigration numbers are per-backend and are not comparable across \
sqlite/postgres/mysql. Refer to a migration by its name."
);
} else {
orion::storage::run_migrations(&pool).await?;
println!("Migrations applied successfully.");
}
Ok(())
}
pub(crate) fn run_lint(
workflow_path: &str,
deny_warnings: bool,
boundary: orion::definitions::Boundary,
definitions: Option<&str>,
) -> Result<(), Box<dyn std::error::Error>> {
use orion::storage::repositories::workflows::CreateWorkflowRequest;
if std::path::Path::new(workflow_path).is_dir() {
return run_lint_set(workflow_path, deny_warnings, boundary);
}
let catalog = Catalog::load_opt(definitions)?;
let doc = read_expanded_workflow(workflow_path, catalog.as_ref())?;
let req: CreateWorkflowRequest = serde_json::from_value(doc)
.map_err(|e| format!("'{workflow_path}' is not a valid workflow JSON: {e}"))?;
let loop_cap = orion::config::EngineConfig::default().max_loop_iterations;
if let Err(err) = orion::validation::validate_create_workflow(&req, loop_cap) {
return Err(format_lint_error(workflow_path, err).into());
}
let warnings: Vec<orion::definitions::Finding> =
orion::validation::unresolvable_logic_warnings(&req.tasks)
.into_iter()
.map(|(path, message)| {
orion::definitions::Finding::warning(
"logic.unresolvable",
format!("workflow '{}' {path}", req.name),
message,
)
})
.collect();
for finding in &warnings {
eprintln!("{finding}");
}
if deny_warnings && !warnings.is_empty() {
return Err(format!(
"'{workflow_path}' has {} warning(s) and --deny-warnings is set",
warnings.len()
)
.into());
}
println!("'{workflow_path}' is valid.");
Ok(())
}
pub(crate) fn read_expanded_workflow(
path: &str,
definitions: Option<&Catalog>,
) -> Result<serde_json::Value, Box<dyn std::error::Error>> {
let raw = std::fs::read_to_string(path).map_err(|e| format!("Failed to read '{path}': {e}"))?;
let mut doc: serde_json::Value =
serde_json::from_str(&raw).map_err(|e| format!("'{path}' is not valid JSON: {e}"))?;
let Some(catalog) = definitions else {
if let Some(reference) = orion::definitions::first_reference(&doc) {
return Err(format!(
"'{path}' contains {reference}, but no --definitions directory was \
given to resolve it against"
)
.into());
}
return Ok(doc);
};
let mut findings = Vec::new();
catalog.shared.expand(&mut doc, path, &mut findings);
let errors = findings.iter().filter(|f| f.is_error()).count();
for finding in &findings {
eprintln!("{finding}");
}
if errors > 0 {
return Err(format!(
"{errors} unresolved reference(s) expanding '{path}' against '{}'",
catalog.dir
)
.into());
}
Ok(doc)
}
pub(crate) struct Catalog {
dir: String,
shared: orion::definitions::SharedDefinitions,
}
impl Catalog {
pub(crate) fn load(dir: &str) -> Result<Self, Box<dyn std::error::Error>> {
let (shared, findings) =
orion::definitions::SharedDefinitions::from_directory(std::path::Path::new(dir))?;
let errors = findings.iter().filter(|f| f.is_error()).count();
for finding in &findings {
eprintln!("{finding}");
}
if errors > 0 {
return Err(format!("{errors} error(s) in the definitions under '{dir}'").into());
}
Ok(Self {
dir: dir.to_string(),
shared,
})
}
pub(crate) fn load_opt(dir: Option<&str>) -> Result<Option<Self>, Box<dyn std::error::Error>> {
dir.map(Self::load).transpose()
}
}
fn run_lint_set(
dir: &str,
deny_warnings: bool,
boundary: orion::definitions::Boundary,
) -> Result<(), Box<dyn std::error::Error>> {
load_and_gate(dir, boundary, false, deny_warnings)?;
Ok(())
}
fn load_and_gate(
dir: &str,
boundary: orion::definitions::Boundary,
require_ids: bool,
deny_warnings: bool,
) -> Result<orion::definitions::DefinitionSet, Box<dyn std::error::Error>> {
let (set, report) =
orion::definitions::DefinitionSet::from_directory(std::path::Path::new(dir))?;
for (path, error) in &report.unparseable {
eprintln!("warning: {} is not readable JSON: {error}", path.display());
}
for path in &report.skipped {
eprintln!(
"note: {} is not a channel, workflow or connector — skipped",
path.display()
);
}
if set.is_empty() {
return Err(format!(
"no definitions found under '{dir}'. A definition is a JSON object with \
'tasks' (workflow), 'channel_type' (channel) or 'connector_type' (connector)."
)
.into());
}
let mut findings = report.findings;
findings.extend(orion::definitions::check(&set, &boundary, require_ids));
let errors = findings.iter().filter(|f| f.is_error()).count();
let warnings = findings.iter().filter(|f| f.is_warning()).count();
for finding in &findings {
eprintln!("{finding}");
}
for (pass, count) in &report.compiled {
println!("compiled: {pass} rewrote {count} document(s)");
}
use orion::definitions::Entity;
let shared = if report.shared.is_empty() {
String::new()
} else {
format!(
", {} shared value(s), {} fragment(s)",
report
.shared
.namespaces
.values()
.map(|n| n.len())
.sum::<usize>(),
report.shared.fragments.len(),
)
};
println!(
"{dir}: {} connector(s), {} workflow(s), {} channel(s){shared} — {errors} error(s), \
{warnings} warning(s)",
set.count(Entity::Connector),
set.count(Entity::Workflow),
set.count(Entity::Channel),
);
if errors > 0 {
return Err(format!("{errors} error(s) in '{dir}'").into());
}
if deny_warnings && warnings > 0 {
return Err(format!("{warnings} warning(s) in '{dir}' and --deny-warnings is set").into());
}
Ok(set)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, clap::ValueEnum)]
pub(crate) enum CompileFormat {
Artifact,
Dir,
Bulk,
}
pub(crate) struct CompileRequest<'a> {
pub(crate) dir: &'a str,
pub(crate) output: Option<&'a str>,
pub(crate) format: CompileFormat,
pub(crate) name: Option<&'a str>,
pub(crate) version: Option<&'a str>,
pub(crate) boundary: orion::definitions::Boundary,
pub(crate) deny_warnings: bool,
pub(crate) no_activate: bool,
}
pub(crate) fn run_compile(req: CompileRequest<'_>) -> Result<(), Box<dyn std::error::Error>> {
let requires_ids = req.format == CompileFormat::Artifact;
let (name, version) = match req.format {
CompileFormat::Artifact => match (req.name, req.version) {
(Some(n), Some(v)) => (n, v),
_ => return Err("--name and --version are required for --format artifact".into()),
},
_ => ("", ""),
};
let requires = orion::definitions::Boundary {
channels: req.boundary.channels.clone(),
connectors: req.boundary.connectors.clone(),
};
let set = load_and_gate(req.dir, req.boundary, requires_ids, req.deny_warnings)?;
match req.format {
CompileFormat::Artifact => emit_artifact(
&set,
req.dir,
name,
version,
requires,
req.no_activate,
req.output,
),
CompileFormat::Dir => emit_dir(&set, req.dir, require_output(req.output, "--format dir")?),
CompileFormat::Bulk => emit_bulk(&set, require_output(req.output, "--format bulk")?),
}
}
fn require_output<'a>(
output: Option<&'a str>,
what: &str,
) -> Result<&'a str, Box<dyn std::error::Error>> {
output.ok_or_else(|| format!("-o <DIR> is required for {what}").into())
}
fn emit_artifact(
set: &orion::definitions::DefinitionSet,
dir: &str,
name: &str,
version: &str,
requires: orion::definitions::Boundary,
no_activate: bool,
output: Option<&str>,
) -> Result<(), Box<dyn std::error::Error>> {
use orion::definitions::Entity;
let collect = |kind: Entity| -> Vec<serde_json::Value> {
set.iter(kind).map(|d| d.doc.clone()).collect()
};
let mut workflows = collect(Entity::Workflow);
let mut channels = collect(Entity::Channel);
if !no_activate {
for entity in workflows.iter_mut().chain(channels.iter_mut()) {
if let Some(obj) = entity.as_object_mut() {
obj.entry("activate")
.or_insert_with(|| serde_json::Value::Bool(true));
}
}
}
let mut artifact = crate::package_cli::PackageArtifact {
package: crate::package_cli::PackageMeta {
name: name.to_string(),
version: version.to_string(),
orion: env!("CARGO_PKG_VERSION").to_string(),
content_hash: String::new(),
exported_from: dir.to_string(),
exported_at: chrono::Utc::now().to_rfc3339(),
},
requires: crate::package_cli::Requires {
channels: requires.channels,
connectors: requires.connectors,
},
connectors: collect(Entity::Connector),
workflows,
channels,
};
artifact.package.content_hash = crate::package_cli::artifact_content_hash(&artifact)?;
let rendered = serde_json::to_string_pretty(&artifact)?;
match output {
Some(path) => {
if let Some(parent) = std::path::Path::new(path).parent()
&& !parent.as_os_str().is_empty()
{
std::fs::create_dir_all(parent)
.map_err(|e| format!("create '{}': {e}", parent.display()))?;
}
std::fs::write(path, rendered).map_err(|e| format!("write '{path}': {e}"))?;
println!(
"wrote {}@{} ({} connectors, {} workflows, {} channels) to {path}",
artifact.package.name,
artifact.package.version,
artifact.connectors.len(),
artifact.workflows.len(),
artifact.channels.len(),
);
}
None => println!("{rendered}"),
}
Ok(())
}
fn emit_dir(
set: &orion::definitions::DefinitionSet,
dir: &str,
out: &str,
) -> Result<(), Box<dyn std::error::Error>> {
let root = std::path::Path::new(dir);
let out_root = std::path::Path::new(out);
for def in &set.definitions {
let origin = std::path::Path::new(&def.origin);
let relative = origin.strip_prefix(root).unwrap_or(origin);
let target = out_root.join(relative);
if let Some(parent) = target.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("create '{}': {e}", parent.display()))?;
}
std::fs::write(&target, serde_json::to_string_pretty(&def.doc)?)
.map_err(|e| format!("write '{}': {e}", target.display()))?;
}
println!(
"wrote {} compiled definition(s) to {out}",
set.definitions.len()
);
Ok(())
}
fn emit_bulk(
set: &orion::definitions::DefinitionSet,
out: &str,
) -> Result<(), Box<dyn std::error::Error>> {
use orion::definitions::Entity;
std::fs::create_dir_all(out).map_err(|e| format!("create '{out}': {e}"))?;
for (kind, file) in [
(Entity::Connector, "connectors.json"),
(Entity::Workflow, "workflows.json"),
(Entity::Channel, "channels.json"),
] {
let entries: Vec<&serde_json::Value> = set.iter(kind).map(|d| &d.doc).collect();
let path = std::path::Path::new(out).join(file);
std::fs::write(&path, serde_json::to_string_pretty(&entries)?)
.map_err(|e| format!("write '{}': {e}", path.display()))?;
println!(
"wrote {} {}(s) to {}",
entries.len(),
kind.as_str(),
path.display()
);
}
Ok(())
}
pub(crate) fn run_dump_openapi() -> Result<(), Box<dyn std::error::Error>> {
println!("{}", orion::server::routes::openapi::pretty_json());
Ok(())
}
fn format_lint_error(workflow_path: &str, err: orion::errors::OrionError) -> String {
use orion::errors::OrionError;
match err {
OrionError::Validation {
code: _,
message,
details,
} => {
let mut out = format!("'{workflow_path}' is invalid: {message}\n");
for d in &details {
out.push_str(&format!(" - {} [{}]: {}\n", d.path, d.code, d.message));
}
out
}
other => format!("'{workflow_path}' is invalid: {other}"),
}
}
pub(crate) fn build_dry_run_engine(
workflow_path: &str,
stubs_path: Option<&str>,
definitions: Option<&Catalog>,
secrets: &orion::engine::ResolvedSecrets,
) -> Result<OfflineRun, Box<dyn std::error::Error>> {
let stubs = match stubs_path {
Some(path) => {
let raw = std::fs::read_to_string(path)
.map_err(|e| format!("Failed to read stubs '{path}': {e}"))?;
orion::engine::functions::stub::parse_stubs(&raw, path)?
}
None => orion::engine::functions::stub::StubTable::new(),
};
build_dry_run_engine_with_stubs(workflow_path, stubs, definitions, secrets)
}
pub(crate) fn offline_secrets(
value: &serde_json::Value,
source: &str,
) -> Result<orion::engine::ResolvedSecrets, String> {
match value {
serde_json::Value::Object(map) => {
for (name, value) in map {
if !value.is_string() {
return Err(format!(
"{source}: secrets.{name} must be a string, got {}",
orion::engine::utils::json_kind(value)
));
}
}
Ok(orion::engine::ResolvedSecrets::from_values(map.clone()))
}
other => Err(format!(
"{source}: secrets must be a JSON object of name -> value, got {}",
orion::engine::utils::json_kind(other)
)),
}
}
pub(crate) struct OfflineRun {
pub engine: dataflow_rs::Engine,
pub log: std::sync::Arc<orion::engine::functions::stub::CallLog>,
}
pub(crate) fn build_dry_run_engine_with_stubs(
workflow_path: &str,
stubs: orion::engine::functions::stub::StubTable,
definitions: Option<&Catalog>,
secrets: &orion::engine::ResolvedSecrets,
) -> Result<OfflineRun, Box<dyn std::error::Error>> {
use orion::storage::repositories::workflows::{CreateWorkflowRequest, workflow_to_dataflow};
let doc = read_expanded_workflow(workflow_path, definitions)?;
let req: CreateWorkflowRequest = serde_json::from_value(doc)
.map_err(|e| format!("'{workflow_path}' is not a valid workflow JSON: {e}"))?;
orion::validation::validate_create_workflow(
&req,
orion::config::EngineConfig::default().max_loop_iterations,
)
.map_err(|e| format_lint_error(workflow_path, e))?;
let synthetic = orion::storage::repositories::workflows::synthetic_workflow(
&req,
req.workflow_id.as_deref().unwrap_or("dry-run"),
)?;
let df_workflow = workflow_to_dataflow(&synthetic, "__dry_run__")?;
let log = std::sync::Arc::new(orion::engine::functions::stub::CallLog::new());
let functions =
orion::engine::functions::stub::build_stub_functions_with_log(stubs, log.clone());
let engine = orion::engine::operators::with_orion_engine_defaults(
dataflow_rs::Engine::builder(),
secrets,
)
.with_workflow(df_workflow)
.with_handlers(functions)
.build()?;
Ok(OfflineRun { engine, log })
}
pub(crate) async fn run_dry_run(
workflow_path: &str,
input_path: &str,
stubs_path: Option<&str>,
metadata_path: Option<&str>,
secrets_path: Option<&str>,
definitions: Option<&str>,
) -> Result<(), Box<dyn std::error::Error>> {
let input_raw = std::fs::read_to_string(input_path)
.map_err(|e| format!("Failed to read input '{input_path}': {e}"))?;
let input: serde_json::Value = serde_json::from_str(&input_raw)
.map_err(|e| format!("'{input_path}' is not valid JSON: {e}"))?;
let metadata = match metadata_path {
Some(path) => {
let raw = std::fs::read_to_string(path)
.map_err(|e| format!("Failed to read metadata '{path}': {e}"))?;
serde_json::from_str(&raw)
.map_err(|e| format!("'{path}' is not valid JSON: {e}"))
.and_then(|v| {
orion::engine::utils::prepare_offline_metadata(v)
.map_err(|e| format!("'{path}': {e}"))
})?
}
None => serde_json::json!({}),
};
let secrets = match secrets_path {
Some(path) => {
let raw = std::fs::read_to_string(path)
.map_err(|e| format!("Failed to read secrets '{path}': {e}"))?;
let value: serde_json::Value = serde_json::from_str(&raw)
.map_err(|e| format!("'{path}' is not valid JSON: {e}"))?;
offline_secrets(&value, path)?
}
None => orion::engine::ResolvedSecrets::empty(),
};
let catalog = Catalog::load_opt(definitions)?;
let run = build_dry_run_engine(workflow_path, stubs_path, catalog.as_ref(), &secrets)?;
let mut message = dataflow_rs::Message::builder()
.payload_json(&input)
.metadata_json(&metadata)
.build();
let mut trace = dataflow_rs::ExecutionTrace::new();
let run_error = run
.engine
.process_message_tracing(&mut message, &mut trace)
.await
.err();
let mut output = serde_json::json!({
"matched": !trace.steps.is_empty(),
"trace": trace,
"output": message.data(),
"errors": message.errors().iter().filter_map(|e| serde_json::to_value(e).ok()).collect::<Vec<_>>(),
});
for (name, document) in orion::engine::functions::stub::run_documents(&message, &run.log) {
output[name] = document;
}
if let Some(ref e) = run_error {
output["error"] = serde_json::json!(e.to_string());
}
println!("{}", serde_json::to_string_pretty(&output)?);
match run_error {
Some(e) => Err(orion::errors::OrionError::Engine(e).into()),
None => Ok(()),
}
}
pub(crate) async fn run_preflight(
config: &config::AppConfig,
) -> Result<(), Box<dyn std::error::Error>> {
use orion::storage::repositories::channels::SqlChannelRepository;
use orion::storage::repositories::workflows::SqlWorkflowRepository;
let pool = orion::storage::init_pool_no_migrate(&config.storage)
.await
.map_err(|e| format!("storage: connection failed: {e}"))?;
println!("Config and environment: OK (checked while loading).");
eprintln!("Scanning stored channels and workflows ...");
let channels = SqlChannelRepository::new(pool.clone());
let workflows = SqlWorkflowRepository::new(pool);
let findings = orion::preflight::scan(&channels, &workflows).await?;
if findings.is_empty() {
println!("Stored channels and workflows: OK — nothing to migrate.");
return Ok(());
}
println!(
"\n{} item(s) need attention before upgrading. Numbers in brackets are \
checklist rows at https://docs.goplasmatic.io/operate/upgrading-to-1.0.html.\n",
findings.len()
);
for finding in &findings {
println!("{finding}\n");
}
Err(format!("preflight found {} item(s) to fix", findings.len()).into())
}
pub(crate) async fn run_test_connectivity(
config: &config::AppConfig,
) -> Result<(), Box<dyn std::error::Error>> {
eprintln!("Probing storage at {} ...", redacted(&config.storage.url));
let pool = orion::storage::init_pool_no_migrate(&config.storage)
.await
.map_err(|e| format!("storage: connection failed: {e}"))?;
let pending = orion::storage::pending_migrations(&pool)
.await
.map_err(|e| format!("storage: pending_migrations query failed: {e}"))?;
println!(
" storage: OK ({} pending migrations)",
pending.len()
);
if config.kafka.enabled {
let broker_list: Vec<String> = config.kafka.brokers.iter().map(|b| redacted(b)).collect();
eprintln!("Probing Kafka brokers {} ...", broker_list.join(","));
let kafka_config = config.kafka.clone();
let brokers = tokio::task::spawn_blocking(move || {
orion::kafka::probe_brokers(&kafka_config, std::time::Duration::from_secs(5))
})
.await
.map_err(|e| format!("kafka: probe task failed: {e}"))?
.map_err(|e| format!("kafka: {e}"))?;
println!(" kafka: OK ({brokers} brokers visible)");
} else {
println!(" kafka: disabled");
}
Ok(())
}
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct TestCase {
#[serde(default)]
name: Option<String>,
workflow: String,
input: serde_json::Value,
#[serde(default)]
metadata: serde_json::Value,
#[serde(default)]
stubs: Option<serde_json::Value>,
#[serde(default)]
secrets: Option<serde_json::Value>,
#[serde(default)]
stubs_file: Option<String>,
#[serde(default)]
expect: std::collections::BTreeMap<String, serde_json::Value>,
#[serde(default)]
expect_errors: Vec<String>,
#[serde(default)]
expect_calls: std::collections::BTreeMap<String, Vec<serde_json::Value>>,
#[serde(default)]
expect_tasks: Option<Vec<String>>,
}
struct CaseResult {
name: String,
failures: Vec<String>,
}
pub(crate) async fn run_test(
path: &str,
definitions: Option<&str>,
) -> Result<(), Box<dyn std::error::Error>> {
let cases = collect_case_files(path)?;
if cases.is_empty() {
return Err(format!(
"no test cases found under '{path}' (looking for *{CASE_SUFFIX}). \
Name a case file explicitly to run one that does not follow the convention."
)
.into());
}
let catalog = Catalog::load_opt(definitions)?;
let mut results = Vec::new();
for case_path in &cases {
results.push(run_case(case_path, catalog.as_ref()).await);
}
let failed: Vec<&CaseResult> = results.iter().filter(|r| !r.failures.is_empty()).collect();
for result in &results {
if result.failures.is_empty() {
println!(" ok {}", result.name);
} else {
println!(" FAIL {}", result.name);
for failure in &result.failures {
println!(" {failure}");
}
}
}
println!(
"\n{} passed, {} failed ({} case(s))",
results.len() - failed.len(),
failed.len(),
results.len()
);
if failed.is_empty() {
Ok(())
} else {
Err(format!("{} test case(s) failed", failed.len()).into())
}
}
pub(crate) const CASE_SUFFIX: &str = ".case.json";
fn collect_case_files(path: &str) -> Result<Vec<std::path::PathBuf>, Box<dyn std::error::Error>> {
let p = std::path::Path::new(path);
if p.is_file() {
return Ok(vec![p.to_path_buf()]);
}
if !p.is_dir() {
return Err(format!("'{path}' is neither a file nor a directory").into());
}
let mut out: Vec<std::path::PathBuf> = std::fs::read_dir(p)
.map_err(|e| format!("Failed to read '{path}': {e}"))?
.filter_map(Result::ok)
.map(|entry| entry.path())
.filter(|p| {
p.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.ends_with(CASE_SUFFIX))
})
.collect();
out.sort();
Ok(out)
}
async fn run_case(case_path: &std::path::Path, definitions: Option<&Catalog>) -> CaseResult {
let display = case_path.display().to_string();
let stem = case_path
.file_name()
.and_then(|s| s.to_str())
.map(|n| n.strip_suffix(CASE_SUFFIX).unwrap_or(n).to_string())
.unwrap_or_else(|| display.clone());
let fail = |name: &str, message: String| CaseResult {
name: name.to_string(),
failures: vec![message],
};
let raw = match std::fs::read_to_string(case_path) {
Ok(raw) => raw,
Err(e) => return fail(&stem, format!("cannot read case: {e}")),
};
let case: TestCase = match serde_json::from_str(&raw) {
Ok(case) => case,
Err(e) => return fail(&stem, format!("not a valid test case: {e}")),
};
let name = case.name.clone().unwrap_or(stem);
let base = case_path.parent().unwrap_or(std::path::Path::new("."));
let workflow_path = base.join(&case.workflow);
let workflow_path = workflow_path.to_string_lossy().to_string();
let stubs = match (&case.stubs, &case.stubs_file) {
(Some(inline), _) => orion::engine::functions::stub::parse_stub_value(inline, "stubs"),
(None, Some(file)) => match std::fs::read_to_string(base.join(file)) {
Ok(raw) => orion::engine::functions::stub::parse_stubs(&raw, file),
Err(e) => Err(format!("cannot read stubs '{file}': {e}")),
},
(None, None) => Ok(orion::engine::functions::stub::StubTable::new()),
};
let stubs = match stubs {
Ok(stubs) => stubs,
Err(e) => return fail(&name, e),
};
let unrooted: Vec<String> = case
.expect
.keys()
.filter(|path| !orion::engine::functions::stub::is_rooted(path))
.map(|path| unrooted_message(path))
.collect();
if !unrooted.is_empty() {
return CaseResult {
name,
failures: unrooted,
};
}
let metadata = match orion::engine::utils::prepare_offline_metadata(case.metadata.clone()) {
Ok(metadata) => metadata,
Err(e) => return fail(&name, e),
};
let secrets = match case.secrets.as_ref() {
Some(value) => match offline_secrets(value, &name) {
Ok(secrets) => secrets,
Err(e) => return fail(&name, e),
},
None => orion::engine::ResolvedSecrets::empty(),
};
let run = match build_dry_run_engine_with_stubs(&workflow_path, stubs, definitions, &secrets) {
Ok(run) => run,
Err(e) => return fail(&name, e.to_string()),
};
let mut message = dataflow_rs::Message::builder()
.payload_json(&case.input)
.metadata_json(&metadata)
.build();
let mut trace = dataflow_rs::ExecutionTrace::new();
let run_error = run
.engine
.process_message_tracing(&mut message, &mut trace)
.await
.err();
let mut failures = Vec::new();
if let Some(e) = run_error {
failures.push(format!("workflow failed: {e}"));
}
let roots = serde_json::Value::Object(orion::engine::functions::stub::run_documents(
&message, &run.log,
));
for (path, expected) in &case.expect {
let actual = lookup_path(&roots, path);
let matched = match actual {
None => expected.is_null(),
Some(ref actual) => actual == expected,
};
if !matched {
failures.push(format!(
"{path}: expected {expected}, got {}",
actual.map_or("<absent>".to_string(), |v| v.to_string())
));
}
}
failures.extend(check_expected_calls(&case.expect_calls, &run.log));
if let Some(ref expected) = case.expect_tasks {
let actual = executed_task_ids(&trace);
if &actual != expected {
failures.push(format!("tasks: expected {expected:?}, ran {actual:?}"));
}
}
let actual_errors: Vec<String> = message
.errors()
.iter()
.map(|e| e.code.to_string())
.collect();
if actual_errors != case.expect_errors {
failures.push(format!(
"task errors: expected {:?}, got {:?}",
case.expect_errors, actual_errors
));
}
CaseResult { name, failures }
}
fn executed_task_ids(trace: &dataflow_rs::ExecutionTrace) -> Vec<String> {
trace
.steps
.iter()
.filter(|step| matches!(step.result, dataflow_rs::StepResult::Executed))
.filter_map(|step| step.task_id.clone())
.collect()
}
fn unrooted_message(path: &str) -> String {
format!(
"expect path '{path}' has no root — did you mean 'data.{path}'? \
roots: {}",
orion::engine::functions::stub::RUN_DOCUMENTS.join(", ")
)
}
enum Segment<'a> {
Key(&'a str),
Index(usize),
}
fn path_segments(path: &str) -> Vec<Segment<'_>> {
let mut out = Vec::new();
for part in path.split('.') {
let (head, mut rest) = match part.find('[') {
Some(i) => (&part[..i], &part[i..]),
None => (part, ""),
};
if let Ok(index) = head.parse::<usize>() {
out.push(Segment::Index(index));
} else if !head.is_empty() {
out.push(Segment::Key(head));
}
while let Some(close) = rest.find(']') {
if let Ok(index) = rest[1..close].parse::<usize>() {
out.push(Segment::Index(index));
}
rest = &rest[close + 1..];
}
}
out
}
fn lookup_path(roots: &serde_json::Value, path: &str) -> Option<serde_json::Value> {
path_segments(path)
.into_iter()
.try_fold(roots, |acc, segment| match segment {
Segment::Key(key) => acc.get(key),
Segment::Index(i) => acc.get(i),
})
.cloned()
}
fn check_expected_calls(
expected: &std::collections::BTreeMap<String, Vec<serde_json::Value>>,
log: &orion::engine::functions::stub::CallLog,
) -> Vec<String> {
if expected.is_empty() {
return Vec::new();
}
let recorded = log.calls();
let mut failures = Vec::new();
for (function, expected_calls) in expected {
let actual: Vec<&orion::engine::functions::stub::RecordedCall> = recorded
.iter()
.filter(|call| call.function == function.as_str())
.collect();
if actual.len() != expected_calls.len() {
failures.push(format!(
"calls.{function}: expected {} call(s), recorded {}",
expected_calls.len(),
actual.len()
));
continue;
}
for (i, want) in expected_calls.iter().enumerate() {
failures.extend(subset_mismatch(
want,
&actual[i].input,
&format!("calls.{function}[{i}].input"),
));
}
}
failures
}
fn subset_mismatch(
expected: &serde_json::Value,
actual: &serde_json::Value,
path: &str,
) -> Vec<String> {
match (expected, actual) {
(serde_json::Value::Object(want), serde_json::Value::Object(got)) => want
.iter()
.flat_map(|(key, want_value)| match got.get(key) {
Some(got_value) => subset_mismatch(want_value, got_value, &format!("{path}.{key}")),
None => vec![format!("{path}.{key}: expected {want_value}, not written")],
})
.collect(),
(serde_json::Value::Array(want), serde_json::Value::Array(got))
if want.len() == got.len() =>
{
want.iter()
.zip(got)
.enumerate()
.flat_map(|(i, (want_value, got_value))| {
subset_mismatch(want_value, got_value, &format!("{path}[{i}]"))
})
.collect()
}
_ if expected == actual => Vec::new(),
_ => vec![format!("{path}: expected {expected}, got {actual}")],
}
}