use anyhow::{Context, Result};
use colored::Colorize;
use oxo_flow_core::config::WorkflowConfig;
use oxo_flow_core::dag::WorkflowDag;
use std::path::{Path, PathBuf};
use crate::commands::print_banner;
pub async fn validate_command(
workflow: PathBuf,
as_include: bool,
json: bool,
ai: bool,
) -> Result<()> {
if let Some(provider) = crate::commands::ai_template::try_resolve_ai(Some(&workflow), ai) {
crate::commands::ai_check::analyze_workflow(&workflow, &provider, "validate", "").await?;
println!();
}
let config_res = WorkflowConfig::from_file(&workflow);
match config_res {
Ok(cfg) => {
if cfg.rules.is_empty() {
if json {
let output = serde_json::json!({
"command": "validate",
"workflow": workflow.display().to_string(),
"valid": true,
"rules": 0,
"dependencies": 0,
"errors": [],
"missing_inputs": [],
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
eprintln!("{} {} — 0 rules", "✓".green().bold(), workflow.display());
eprintln!(
" {} Workflow has no rules. Add [[rules]] sections to define pipeline steps.",
"⚠ Warning:".yellow().bold()
);
}
return Ok(());
}
let validation = oxo_flow_core::format::validate_format(&cfg);
let mut error_count = 0usize;
let mut errors_json: Vec<serde_json::Value> = Vec::new();
for d in &validation.diagnostics {
if d.severity == oxo_flow_core::format::Severity::Error {
if as_include && (d.code == "E010" || d.code == "W020") {
continue;
}
error_count += 1;
if json {
errors_json.push(serde_json::json!({
"code": d.code,
"message": d.message,
"rule": d.rule,
"suggestion": d.suggestion,
}));
} else {
eprintln!(" {} [{}]: {}", "error".red().bold(), d.code, d.message);
if let Some(ref rule) = d.rule {
eprintln!(" rule: {}", rule);
}
if let Some(ref suggestion) = d.suggestion {
eprintln!(" hint: {}", suggestion);
}
}
}
}
let workflow_dir = oxo_flow_core::parent_dir(&workflow);
let mut missing_inputs = Vec::new();
if !as_include {
for rule in &cfg.rules {
for input in &rule.input {
if !input.contains('{')
&& !input.contains('}')
&& !workflow_dir.join(input).exists()
{
let is_generated =
cfg.rules.iter().any(|r| r.output.to_vec().contains(input));
if !is_generated {
missing_inputs.push(input);
}
}
}
}
}
let (rules, dependencies) = if as_include {
(cfg.rules.len(), 0)
} else {
match WorkflowDag::from_rules(&cfg.rules) {
Ok(dag) => (dag.node_count(), dag.edge_count()),
Err(e) => {
if json {
let output = serde_json::json!({
"command": "validate",
"workflow": workflow.display().to_string(),
"valid": false,
"rules": cfg.rules.len(),
"dependencies": 0,
"errors": [{"code": "DAG", "message": e.to_string()}],
"missing_inputs": [],
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
eprintln!(
"{} {} — DAG error: {}",
"✗".red().bold(),
workflow.display(),
e
);
}
std::process::exit(1);
}
}
};
if json {
let output = serde_json::json!({
"command": "validate",
"workflow": workflow.display().to_string(),
"valid": error_count == 0,
"rules": rules,
"dependencies": dependencies,
"errors": errors_json,
"missing_inputs": missing_inputs,
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
if error_count == 0 {
eprintln!(
"{} {} — {} rules, {} dependencies",
"✓".green().bold(),
workflow.display(),
rules,
dependencies
);
} else {
eprintln!(
"{} {} — {} validation error(s)",
"✗".red().bold(),
workflow.display(),
error_count
);
}
if !missing_inputs.is_empty() {
eprintln!(
"\n {} The following input files do not exist:",
"⚠ Warning:".yellow().bold()
);
for input in missing_inputs {
eprintln!(" - {}", input);
}
}
}
if error_count > 0 {
std::process::exit(1);
}
}
Err(e) => {
if json {
let output = serde_json::json!({
"command": "validate",
"workflow": workflow.display().to_string(),
"valid": false,
"rules": 0,
"dependencies": 0,
"errors": [{"code": "PARSE", "message": e.to_string()}],
"missing_inputs": [],
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
eprintln!("{} {} — {}", "✗".red().bold(), workflow.display(), e);
}
std::process::exit(1);
}
}
Ok(())
}
pub async fn lint_command(workflow: PathBuf, strict: bool, json: bool, ai: bool) -> Result<()> {
print_banner();
if let Some(provider) = crate::commands::ai_template::try_resolve_ai(Some(&workflow), ai) {
crate::commands::ai_check::analyze_workflow(&workflow, &provider, "lint", "").await?;
println!();
}
let config = WorkflowConfig::from_file(&workflow)
.with_context(|| format!("failed to parse {}", workflow.display()))?;
let validation = oxo_flow_core::format::validate_format(&config);
let lint_diags = oxo_flow_core::format::lint_format(&config);
let raw_content = std::fs::read_to_string(&workflow).ok();
let secret_diags = if let Some(content) = raw_content {
oxo_flow_core::format::scan_for_secrets(&content)
} else {
Vec::new()
};
let mut error_count = 0usize;
let mut warning_count = 0usize;
let mut info_count = 0usize;
for d in validation
.diagnostics
.iter()
.chain(lint_diags.iter())
.chain(secret_diags.iter())
{
let prefix = match d.severity {
oxo_flow_core::format::Severity::Error => {
error_count += 1;
"error".red().bold().to_string()
}
oxo_flow_core::format::Severity::Warning => {
warning_count += 1;
"warning".yellow().bold().to_string()
}
oxo_flow_core::format::Severity::Info => {
info_count += 1;
"info".blue().to_string()
}
};
eprint!(" {} [{}]: {}", prefix, d.code, d.message);
if let Some(ref rule) = d.rule {
eprint!(" (rule: {})", rule);
}
eprintln!();
}
eprintln!(
"\n{} {} error(s), {} warning(s), {} info",
"Summary:".bold(),
error_count,
warning_count,
info_count
);
if json {
let diagnostics: Vec<serde_json::Value> = validation
.diagnostics
.iter()
.chain(lint_diags.iter())
.chain(secret_diags.iter())
.map(|d| {
serde_json::json!({
"severity": format!("{:?}", d.severity).to_lowercase(),
"code": d.code,
"message": d.message,
"rule": d.rule,
"suggestion": d.suggestion,
})
})
.collect();
let output = serde_json::json!({
"command": "lint",
"workflow": workflow.display().to_string(),
"strict": strict,
"diagnostics": diagnostics,
"error_count": error_count,
"warning_count": warning_count,
"info_count": info_count,
"passed": error_count == 0 && (!strict || warning_count == 0),
});
println!("{}", serde_json::to_string_pretty(&output).unwrap());
if error_count > 0 || (strict && warning_count > 0) {
std::process::exit(1);
}
return Ok(());
}
if error_count > 0 || (strict && warning_count > 0) {
std::process::exit(1);
}
Ok(())
}
pub fn deep_check_command(workflow: &Path, workdir: Option<&Path>, json: bool) -> Result<()> {
let config = WorkflowConfig::from_file(workflow)
.with_context(|| format!("failed to parse {}", workflow.display()))?;
let base_dir = workdir
.map(Path::to_path_buf)
.unwrap_or_else(|| oxo_flow_core::parent_dir(workflow).to_path_buf());
let report = oxo_flow_core::deep_check::compute_deep_check(&config, &base_dir);
print_deep_console(&report);
if json {
let diagnostics: Vec<serde_json::Value> = report
.findings
.iter()
.map(|f| {
serde_json::json!({
"severity": format!("{:?}", f.severity).to_lowercase(),
"code": f.code,
"message": f.message,
"rule": f.rule,
"suggestion": f.suggestion,
"path": f.path,
})
})
.collect();
let output = serde_json::json!({
"command": "deep-check",
"workflow": workflow.display().to_string(),
"diagnostics": diagnostics,
"error_count": report.error_count,
"warning_count": report.warning_count,
"passed": report.passed,
});
println!("{}", serde_json::to_string_pretty(&output)?);
}
if report.error_count > 0 {
std::process::exit(1);
}
Ok(())
}
fn print_deep_console(report: &oxo_flow_core::deep_check::DeepCheckReport) {
let finding_lines = |code: &str| {
for f in report.findings.iter().filter(|f| f.code == code) {
let icon = if f.severity == oxo_flow_core::format::Severity::Error {
"✗".red().bold()
} else {
"⚠".yellow().bold()
};
let rule_suffix = f
.rule
.as_deref()
.map(|rule| format!(" (rule: {rule})"))
.unwrap_or_default();
eprintln!(" {} {}{} [{}]", icon, f.message, rule_suffix, f.code);
if let Some(hint) = &f.suggestion {
eprintln!(" hint: {}", hint.dimmed());
}
}
};
let has = |code: &str| report.findings.iter().any(|f| f.code == code);
if report.scripts_checked > 0 || has("D001") {
eprintln!("{}", " Scripts:".bold());
if report.scripts_checked > 0 && !has("D001") {
eprintln!(
" {} {} script reference(s) found",
"✓".green().bold(),
report.scripts_checked
);
}
finding_lines("D001");
}
if report.envs_checked > 0 || has("D002") {
eprintln!("{}", " Environments:".bold());
if report.envs_checked > 0 && !has("D002") {
eprintln!(
" {} {} environment definition(s) found",
"✓".green().bold(),
report.envs_checked
);
}
finding_lines("D002");
}
if report.commands_probed > 0 || has("D003") {
eprintln!("{}", " Binaries:".bold());
if report.commands_probed > 0 && !has("D003") {
eprintln!(
" {} {} command(s) found in PATH",
"✓".green().bold(),
report.commands_probed
);
}
finding_lines("D003");
}
if report.references_checked > 0 || has("D004") {
eprintln!("{}", " References:".bold());
if report.references_checked > 0 && !has("D004") {
eprintln!(
" {} {} reference path(s) found",
"✓".green().bold(),
report.references_checked
);
}
finding_lines("D004");
}
eprintln!(
"\n{} {} error(s), {} warning(s)",
"Deep check summary:".bold(),
report.error_count,
report.warning_count
);
}
pub fn format_command(workflow: PathBuf, output: Option<PathBuf>, check: bool) -> Result<()> {
let config = WorkflowConfig::from_file(&workflow)
.with_context(|| format!("failed to parse {}", workflow.display()))?;
let formatted = oxo_flow_core::format::format_workflow(&config);
if check {
let original = std::fs::read_to_string(&workflow)?;
if original.trim() == formatted.trim() {
eprintln!(
"{} {} is already formatted",
"✓".green().bold(),
workflow.display()
);
} else {
eprintln!(
"{} {} needs formatting",
"✗".red().bold(),
workflow.display()
);
std::process::exit(1);
}
} else {
match output {
Some(path) => {
std::fs::write(&path, &formatted)?;
eprintln!("Formatted workflow written to {}", path.display());
}
None => {
print!("{formatted}");
}
}
}
Ok(())
}
pub fn touch_command(
workflow: PathBuf,
rules: Vec<String>,
workdir: Option<PathBuf>,
) -> Result<()> {
print_banner();
let mut config = WorkflowConfig::from_file(&workflow)
.with_context(|| format!("failed to parse {}", workflow.display()))?;
config.apply_defaults();
if let Err(e) = config.expand_wildcards() {
eprintln!(" {} Could not expand wildcards: {}", "Note:".yellow(), e);
eprintln!(
" {} Wildcard patterns in outputs will be skipped.",
"Info:".dimmed()
);
}
let rules_to_touch: Vec<&oxo_flow_core::rule::Rule> = if rules.is_empty() {
config.rules.iter().collect()
} else {
config
.rules
.iter()
.filter(|r| rules.contains(&r.name))
.collect()
};
let mut touched = 0usize;
let mut skipped = 0usize;
let mut skipped_patterns: Vec<(String, String)> = Vec::new();
let base_dir = workdir
.clone()
.unwrap_or_else(|| workflow.parent().map(Path::to_path_buf).unwrap_or_default());
for rule in &rules_to_touch {
for output in &rule.output {
let has_wildcard = output.contains('{') && output.contains('}');
if has_wildcard {
skipped += 1;
skipped_patterns.push((rule.name.clone(), output.clone()));
continue;
}
if output.contains("..") || output.starts_with('/') || output.starts_with('~') {
eprintln!(" {} {} (rejected: unsafe path)", "✗".red().bold(), output);
continue;
}
let path = base_dir.join(output);
if path.exists() {
match filetime::set_file_mtime(&path, filetime::FileTime::now()) {
Ok(()) => {
touched += 1;
eprintln!(" {} {}", "✓".green(), output);
}
Err(e) => {
eprintln!(" {} {} ({})", "✗".red(), output, e);
}
}
} else {
if let Some(parent) = path.parent()
&& let Err(e) = std::fs::create_dir_all(parent)
{
eprintln!(
" {} {} (cannot create directory: {})",
"✗".red(),
output,
e
);
continue;
}
match std::fs::write(&path, "") {
Ok(()) => {
touched += 1;
eprintln!(" {} {} (created)", "✓".green(), output);
}
Err(e) => {
eprintln!(" {} {} (failed: {})", "✗".red(), output, e);
}
}
}
}
}
eprintln!(
"\n{} {} file(s) touched, {} wildcard pattern(s) skipped",
"Done:".bold(),
touched,
skipped
);
if !skipped_patterns.is_empty() {
eprintln!();
for (rule_name, pattern) in &skipped_patterns {
eprintln!(
" {} {} → {} (wildcard pattern — not expanded)",
"Skipped:".yellow(),
rule_name,
pattern.dimmed()
);
}
eprintln!();
eprintln!(
" {} To touch expanded rules, use specific rule names after wildcard expansion.",
"Tip:".bold().cyan()
);
eprintln!(
" {} Run 'oxo-flow dry-run {}' to see expanded rule names, then use:",
" ".dimmed(),
workflow.display()
);
eprintln!(
" {} oxo-flow touch {} --rule <expanded_name>",
" ".dimmed(),
workflow.display()
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use assert_cmd::Command;
use std::io::Write;
use tempfile::NamedTempFile;
#[test]
fn test_as_include_skips_dag_validation() {
let fragment = r#"
[workflow]
name = "qc-fragment"
[[rules]]
name = "fastqc"
input = ["{sample}.fastq"]
output = ["{sample}_fastqc.html"]
shell = "fastqc {input}"
"#;
let mut file = NamedTempFile::with_suffix(".oxoflow").unwrap();
file.write_all(fragment.as_bytes()).unwrap();
Command::cargo_bin("oxo-flow")
.unwrap()
.arg("validate")
.arg("--as-include")
.arg(file.path())
.assert()
.success();
}
#[test]
fn test_as_include_validates_syntax() {
let fragment = r#"
[workflow]
name = "bad-fragment"
[[rules]]
# Missing required 'name' field
input = ["test.txt"]
"#;
let mut file = NamedTempFile::with_suffix(".oxoflow").unwrap();
file.write_all(fragment.as_bytes()).unwrap();
Command::cargo_bin("oxo-flow")
.unwrap()
.arg("validate")
.arg("--as-include")
.arg(file.path())
.assert()
.failure();
}
}