Skip to main content

cargo_hammerwork/commands/
workflow.rs

1use anyhow::Result;
2use clap::Subcommand;
3use hammerwork::{FailurePolicy, JobGroup};
4use serde_json::Value;
5use std::collections::{HashMap, HashSet, VecDeque};
6use tracing::{info, warn};
7use uuid::Uuid;
8
9use crate::config::Config;
10use crate::utils::database::DatabasePool;
11use crate::utils::validation::validate_json_payload;
12
13#[derive(Debug, Clone)]
14pub struct JobNode {
15    pub id: String,
16    pub queue_name: String,
17    pub status: String,
18    pub dependency_status: String,
19    pub depends_on: Vec<String>,
20    pub dependents: Vec<String>,
21    pub workflow_id: Option<String>,
22    pub workflow_name: Option<String>,
23}
24
25#[derive(Subcommand)]
26pub enum WorkflowCommand {
27    #[command(about = "List workflows")]
28    List {
29        #[arg(short = 'u', long, help = "Database connection URL")]
30        database_url: Option<String>,
31        #[arg(short, long, help = "Maximum number of workflows to display")]
32        limit: Option<u32>,
33        #[arg(long, help = "Show only running workflows")]
34        running: bool,
35        #[arg(long, help = "Show only completed workflows")]
36        completed: bool,
37        #[arg(long, help = "Show only failed workflows")]
38        failed: bool,
39    },
40    #[command(about = "Show details of a specific workflow")]
41    Show {
42        #[arg(short = 'u', long, help = "Database connection URL")]
43        database_url: Option<String>,
44        #[arg(help = "Workflow ID")]
45        workflow_id: String,
46        #[arg(long, help = "Show dependency graph")]
47        dependencies: bool,
48    },
49    #[command(about = "Create a new workflow")]
50    Create {
51        #[arg(short = 'u', long, help = "Database connection URL")]
52        database_url: Option<String>,
53        #[arg(short = 'n', long, help = "Workflow name")]
54        name: String,
55        #[arg(long, help = "Failure policy (fail_fast, continue_on_failure, manual)")]
56        failure_policy: Option<String>,
57        #[arg(long, help = "Workflow metadata as JSON")]
58        metadata: Option<String>,
59    },
60    #[command(about = "Cancel a running workflow")]
61    Cancel {
62        #[arg(short = 'u', long, help = "Database connection URL")]
63        database_url: Option<String>,
64        #[arg(help = "Workflow ID")]
65        workflow_id: String,
66        #[arg(long, help = "Force cancel even if jobs are running")]
67        force: bool,
68    },
69    #[command(about = "Show job dependencies")]
70    Dependencies {
71        #[arg(short = 'u', long, help = "Database connection URL")]
72        database_url: Option<String>,
73        #[arg(help = "Job ID")]
74        job_id: String,
75        #[arg(long, help = "Show dependency tree")]
76        tree: bool,
77        #[arg(long, help = "Show jobs that depend on this job")]
78        dependents: bool,
79    },
80    #[command(about = "Visualize workflow dependency graph")]
81    Graph {
82        #[arg(short = 'u', long, help = "Database connection URL")]
83        database_url: Option<String>,
84        #[arg(help = "Workflow ID")]
85        workflow_id: String,
86        #[arg(long, help = "Output format (text, dot, mermaid, json)")]
87        format: Option<String>,
88    },
89}
90
91impl WorkflowCommand {
92    pub async fn execute(&self, config: Config) -> Result<()> {
93        let db_url = self.get_database_url(&config)?;
94        let pool = DatabasePool::connect(&db_url, config.get_connection_pool_size()).await?;
95
96        match self {
97            WorkflowCommand::List {
98                limit,
99                running,
100                completed,
101                failed,
102                ..
103            } => {
104                self.list_workflows(pool, *limit, *running, *completed, *failed)
105                    .await
106            }
107            WorkflowCommand::Show {
108                workflow_id,
109                dependencies,
110                ..
111            } => self.show_workflow(pool, workflow_id, *dependencies).await,
112            WorkflowCommand::Create {
113                name,
114                failure_policy,
115                metadata,
116                ..
117            } => {
118                self.create_workflow(name, failure_policy.as_deref(), metadata.as_deref())
119                    .await
120            }
121            WorkflowCommand::Cancel {
122                workflow_id, force, ..
123            } => self.cancel_workflow(pool, workflow_id, *force).await,
124            WorkflowCommand::Dependencies {
125                job_id,
126                tree,
127                dependents,
128                ..
129            } => {
130                self.show_dependencies(pool, job_id, *tree, *dependents)
131                    .await
132            }
133            WorkflowCommand::Graph {
134                workflow_id,
135                format,
136                ..
137            } => self.show_graph(pool, workflow_id, format.as_deref()).await,
138        }
139    }
140
141    pub fn get_database_url(&self, config: &Config) -> Result<String> {
142        let url_option = match self {
143            WorkflowCommand::List { database_url, .. } => database_url,
144            WorkflowCommand::Show { database_url, .. } => database_url,
145            WorkflowCommand::Create { database_url, .. } => database_url,
146            WorkflowCommand::Cancel { database_url, .. } => database_url,
147            WorkflowCommand::Dependencies { database_url, .. } => database_url,
148            WorkflowCommand::Graph { database_url, .. } => database_url,
149        };
150
151        url_option
152            .as_ref()
153            .map(|s| s.as_str())
154            .or(config.get_database_url())
155            .map(|s| s.to_string())
156            .ok_or_else(|| anyhow::anyhow!("Database URL is required"))
157    }
158
159    async fn list_workflows(
160        &self,
161        _pool: DatabasePool,
162        limit: Option<u32>,
163        running: bool,
164        completed: bool,
165        failed: bool,
166    ) -> Result<()> {
167        warn!("Workflow listing is not fully implemented yet");
168        println!("⚠️  Workflow listing is not fully implemented yet.");
169        println!("    This would list workflows with filters:");
170        println!("    - Limit: {}", limit.unwrap_or(50));
171        println!("    - Running: {}", running);
172        println!("    - Completed: {}", completed);
173        println!("    - Failed: {}", failed);
174
175        Ok(())
176    }
177
178    async fn show_workflow(
179        &self,
180        _pool: DatabasePool,
181        workflow_id: &str,
182        show_dependencies: bool,
183    ) -> Result<()> {
184        warn!("Workflow details view is not fully implemented yet");
185        println!("⚠️  Workflow details view is not fully implemented yet.");
186        println!("    This would show details for workflow: {}", workflow_id);
187        println!("    Show dependencies: {}", show_dependencies);
188
189        Ok(())
190    }
191
192    async fn create_workflow(
193        &self,
194        name: &str,
195        failure_policy: Option<&str>,
196        metadata: Option<&str>,
197    ) -> Result<()> {
198        // Parse failure policy
199        let policy = match failure_policy {
200            Some("fail_fast") => FailurePolicy::FailFast,
201            Some("continue_on_failure") => FailurePolicy::ContinueOnFailure,
202            Some("manual") => FailurePolicy::Manual,
203            Some(p) => {
204                return Err(anyhow::anyhow!(
205                    "Invalid failure policy: {}. Use: fail_fast, continue_on_failure, manual",
206                    p
207                ));
208            }
209            None => FailurePolicy::FailFast,
210        };
211
212        // Parse metadata
213        let metadata_json = match metadata {
214            Some(m) => validate_json_payload(m)?,
215            None => serde_json::Value::Object(serde_json::Map::new()),
216        };
217
218        // Create workflow
219        let workflow = JobGroup::new(name)
220            .with_failure_policy(policy)
221            .with_metadata(metadata_json);
222
223        println!("Workflow created:");
224        println!("  ID: {}", workflow.id);
225        println!("  Name: {}", workflow.name);
226        println!("  Failure Policy: {:?}", workflow.failure_policy);
227
228        info!("Created workflow '{}' with ID {}", name, workflow.id);
229        println!("\nUse 'cargo hammerwork job enqueue' to add jobs to this workflow.");
230
231        Ok(())
232    }
233
234    async fn cancel_workflow(
235        &self,
236        _pool: DatabasePool,
237        workflow_id: &str,
238        force: bool,
239    ) -> Result<()> {
240        warn!("Workflow cancellation is not fully implemented yet");
241        println!("⚠️  Workflow cancellation is not fully implemented yet.");
242        println!("    This would cancel workflow: {}", workflow_id);
243        println!("    Force: {}", force);
244
245        Ok(())
246    }
247
248    async fn show_dependencies(
249        &self,
250        pool: DatabasePool,
251        job_id: &str,
252        show_tree: bool,
253        show_dependents: bool,
254    ) -> Result<()> {
255        let job_uuid = Uuid::parse_str(job_id)?;
256
257        // Get the target job details
258        let target_job = self.get_job_node(&pool, &job_uuid).await?;
259        let target_job = match target_job {
260            Some(job) => job,
261            None => {
262                println!("Job not found: {}", job_id);
263                return Ok(());
264            }
265        };
266
267        println!("Job Dependencies for {}", job_id);
268        println!("Queue: {}", target_job.queue_name);
269        println!("Status: {}", target_job.status);
270        println!("Dependency Status: {}", target_job.dependency_status);
271
272        if let Some(workflow_name) = &target_job.workflow_name {
273            println!("Workflow: {}", workflow_name);
274        }
275
276        // Show immediate dependencies
277        if !target_job.depends_on.is_empty() {
278            println!("\nDirect Dependencies:");
279            for dep_id in &target_job.depends_on {
280                if let Some(dep_job) = self.get_job_node_by_string(&pool, dep_id).await? {
281                    println!("  ├─ {} ({})", dep_id, dep_job.status);
282                } else {
283                    println!("  ├─ {} (not found)", dep_id);
284                }
285            }
286        } else {
287            println!("\nNo direct dependencies");
288        }
289
290        // Show immediate dependents if requested
291        if show_dependents {
292            if !target_job.dependents.is_empty() {
293                println!("\nDirect Dependents:");
294                for dep_id in &target_job.dependents {
295                    if let Some(dep_job) = self.get_job_node_by_string(&pool, dep_id).await? {
296                        println!("  ├─ {} ({})", dep_id, dep_job.status);
297                    } else {
298                        println!("  ├─ {} (not found)", dep_id);
299                    }
300                }
301            } else {
302                println!("\nNo direct dependents");
303            }
304        }
305
306        // Show full dependency tree if requested
307        if show_tree {
308            println!("\nDependency Tree:");
309
310            // Build dependency graph for the workflow or just this job
311            let jobs = if let Some(workflow_id) = &target_job.workflow_id {
312                self.get_workflow_jobs(&pool, workflow_id).await?
313            } else {
314                // Collect all related jobs by traversing dependencies
315                self.collect_related_jobs(&pool, &target_job).await?
316            };
317
318            if jobs.is_empty() {
319                println!("  No related jobs found");
320            } else {
321                self.print_dependency_tree(&jobs, &target_job.id);
322            }
323        }
324
325        Ok(())
326    }
327
328    async fn show_graph(
329        &self,
330        pool: DatabasePool,
331        workflow_id: &str,
332        format: Option<&str>,
333    ) -> Result<()> {
334        let format = format.unwrap_or("text");
335
336        // Get all jobs in the workflow
337        let jobs = self.get_workflow_jobs(&pool, workflow_id).await?;
338
339        if jobs.is_empty() {
340            println!("No jobs found in workflow: {}", workflow_id);
341            return Ok(());
342        }
343
344        println!("Workflow Graph: {}", workflow_id);
345        println!("Format: {}", format);
346        println!("Jobs: {}", jobs.len());
347
348        match format {
349            "text" => self.print_text_graph(&jobs),
350            "dot" => self.print_dot_graph(&jobs, workflow_id),
351            "mermaid" => self.print_mermaid_graph(&jobs, workflow_id),
352            "json" => self.print_json_graph(&jobs)?,
353            _ => {
354                println!(
355                    "Unsupported format: {}. Use: text, dot, mermaid, json",
356                    format
357                );
358                return Ok(());
359            }
360        }
361
362        Ok(())
363    }
364
365    fn print_text_graph(&self, jobs: &[JobNode]) {
366        println!("\nDependency Graph (Text Format):");
367        println!("{}", "=".repeat(50));
368
369        // Group jobs by dependency level
370        let mut levels = self.calculate_dependency_levels(jobs);
371        levels.sort_by_key(|level| level.0);
372
373        for (level, jobs_at_level) in levels {
374            println!("\nLevel {}: {} job(s)", level, jobs_at_level.len());
375            for job in jobs_at_level {
376                let deps_str = if job.depends_on.is_empty() {
377                    "none".to_string()
378                } else {
379                    format!("{} dependencies", job.depends_on.len())
380                };
381
382                println!("  ├─ [{}] {} ({})", &job.id[..8], job.status, deps_str);
383            }
384        }
385    }
386
387    fn print_dot_graph(&self, jobs: &[JobNode], workflow_id: &str) {
388        println!("\nDOT Graph Format:");
389        println!("{}", "=".repeat(50));
390        println!("digraph workflow_{} {{", workflow_id.replace('-', "_"));
391        println!("  rankdir=TB;");
392        println!("  node [shape=box];");
393
394        // Define nodes
395        for job in jobs {
396            let color = match job.status.as_str() {
397                "Completed" => "lightgreen",
398                "Failed" => "lightcoral",
399                "Running" => "lightblue",
400                "Pending" => "lightyellow",
401                _ => "lightgray",
402            };
403
404            println!(
405                "  \"{}\" [label=\"{}\\n{}\" fillcolor={} style=filled];",
406                job.id,
407                &job.id[..8],
408                job.status,
409                color
410            );
411        }
412
413        // Define edges (dependencies)
414        for job in jobs {
415            for dep_id in &job.depends_on {
416                println!("  \"{}\" -> \"{}\";", dep_id, job.id);
417            }
418        }
419
420        println!("}}");
421        println!(
422            "\nTo visualize: copy the above DOT code to https://dreampuf.github.io/GraphvizOnline/"
423        );
424    }
425
426    fn print_mermaid_graph(&self, jobs: &[JobNode], workflow_id: &str) {
427        println!("\nMermaid Graph Format:");
428        println!("{}", "=".repeat(50));
429        println!("---");
430        println!("title: Hammerwork Workflow Dependency Graph");
431        println!("---");
432        println!("graph TD");
433        println!("    subgraph \"📋 Workflow: {}\"", &workflow_id[..8]);
434
435        // Define nodes with styling
436        for job in jobs {
437            let short_id = &job.id[..8];
438            let status_class = match job.status.as_str() {
439                "Completed" => ":::completed",
440                "Failed" => ":::failed",
441                "Running" => ":::running",
442                "Pending" => ":::pending",
443                _ => ":::default",
444            };
445
446            // Node definition with label inside subgraph
447            let dependency_indicator = match job.dependency_status.as_str() {
448                "waiting" => "⏳",
449                "satisfied" => "✅",
450                "failed" => "❌",
451                _ => "🔵",
452            };
453
454            println!(
455                "        {}[\"{}<br/>{}<br/>{} {}\"]{}",
456                short_id, short_id, job.status, dependency_indicator, job.queue_name, status_class
457            );
458        }
459
460        println!();
461
462        // Define edges (dependencies) inside subgraph
463        for job in jobs {
464            let job_short = &job.id[..8];
465            for dep_id in &job.depends_on {
466                let dep_short = &dep_id[..8];
467                println!("        {} --> {}", dep_short, job_short);
468            }
469        }
470
471        println!("    end");
472        println!();
473
474        // Define styling classes
475        println!(
476            "    classDef completed fill:#d4edda,stroke:#155724,stroke-width:2px,color:#155724"
477        );
478        println!("    classDef failed fill:#f8d7da,stroke:#721c24,stroke-width:2px,color:#721c24");
479        println!("    classDef running fill:#cce7ff,stroke:#004085,stroke-width:2px,color:#004085");
480        println!("    classDef pending fill:#fff3cd,stroke:#856404,stroke-width:2px,color:#856404");
481        println!("    classDef default fill:#e2e3e5,stroke:#383d41,stroke-width:2px,color:#383d41");
482
483        println!("\nTo visualize: copy the above Mermaid code to:");
484        println!("- GitHub/GitLab markdown (```mermaid ... ```)");
485        println!("- https://mermaid.live/");
486        println!("- VS Code with Mermaid extension");
487    }
488
489    pub fn print_json_graph(&self, jobs: &[JobNode]) -> Result<()> {
490        println!("\nJSON Graph Format:");
491        println!("{}", "=".repeat(50));
492
493        let graph = serde_json::json!({
494            "nodes": jobs.iter().map(|job| {
495                serde_json::json!({
496                    "id": job.id,
497                    "queue": job.queue_name,
498                    "status": job.status,
499                    "dependency_status": job.dependency_status,
500                    "workflow_id": job.workflow_id,
501                    "workflow_name": job.workflow_name
502                })
503            }).collect::<Vec<_>>(),
504            "edges": jobs.iter().flat_map(|job| {
505                job.depends_on.iter().map(move |dep_id| {
506                    serde_json::json!({
507                        "from": dep_id,
508                        "to": job.id,
509                        "type": "dependency"
510                    })
511                })
512            }).collect::<Vec<_>>()
513        });
514
515        println!("{}", serde_json::to_string_pretty(&graph)?);
516        Ok(())
517    }
518
519    pub fn calculate_dependency_levels<'a>(
520        &self,
521        jobs: &'a [JobNode],
522    ) -> Vec<(usize, Vec<&'a JobNode>)> {
523        let job_map: HashMap<String, &JobNode> =
524            jobs.iter().map(|job| (job.id.clone(), job)).collect();
525
526        let mut levels = HashMap::new();
527        let mut visited = HashSet::new();
528
529        for job in jobs {
530            if !visited.contains(&job.id) {
531                Self::calculate_job_level(job, &job_map, &mut levels, &mut visited, 0);
532            }
533        }
534
535        // Group jobs by level
536        let mut result: HashMap<usize, Vec<&JobNode>> = HashMap::new();
537        for (job_id, level) in levels {
538            if let Some(job) = job_map.get(&job_id) {
539                result.entry(level).or_default().push(job);
540            }
541        }
542
543        result.into_iter().collect()
544    }
545
546    fn calculate_job_level(
547        job: &JobNode,
548        job_map: &HashMap<String, &JobNode>,
549        levels: &mut HashMap<String, usize>,
550        visited: &mut HashSet<String>,
551        _current_level: usize,
552    ) {
553        if visited.contains(&job.id) {
554            return;
555        }
556
557        visited.insert(job.id.clone());
558
559        // Calculate max dependency level
560        let max_dep_level = job
561            .depends_on
562            .iter()
563            .filter_map(|dep_id| {
564                if let Some(dep_job) = job_map.get(dep_id) {
565                    if !visited.contains(dep_id) {
566                        Self::calculate_job_level(
567                            dep_job,
568                            job_map,
569                            levels,
570                            visited,
571                            _current_level,
572                        );
573                    }
574                    levels.get(dep_id).copied()
575                } else {
576                    None
577                }
578            })
579            .max()
580            .unwrap_or(0);
581
582        let job_level = max_dep_level + if job.depends_on.is_empty() { 0 } else { 1 };
583        levels.insert(job.id.clone(), job_level);
584    }
585
586    // Helper methods for dependency tree visualization
587
588    async fn get_job_node(&self, pool: &DatabasePool, job_id: &Uuid) -> Result<Option<JobNode>> {
589        let query = r#"
590            SELECT id, queue_name, status, dependency_status, depends_on, dependents, workflow_id, workflow_name
591            FROM hammerwork_jobs 
592            WHERE id = $1
593        "#;
594
595        match pool {
596            DatabasePool::Postgres(pg_pool) => {
597                if let Some(row) = sqlx::query(query)
598                    .bind(job_id)
599                    .fetch_optional(pg_pool)
600                    .await?
601                {
602                    Ok(Some(self.postgres_row_to_job_node(&row)?))
603                } else {
604                    Ok(None)
605                }
606            }
607            DatabasePool::MySQL(mysql_pool) => {
608                let mysql_query = r#"
609                    SELECT id, queue_name, status, dependency_status, depends_on, dependents, workflow_id, workflow_name
610                    FROM hammerwork_jobs 
611                    WHERE id = ?
612                "#;
613                if let Some(row) = sqlx::query(mysql_query)
614                    .bind(job_id.to_string())
615                    .fetch_optional(mysql_pool)
616                    .await?
617                {
618                    Ok(Some(self.mysql_row_to_job_node(&row)?))
619                } else {
620                    Ok(None)
621                }
622            }
623        }
624    }
625
626    async fn get_job_node_by_string(
627        &self,
628        pool: &DatabasePool,
629        job_id: &str,
630    ) -> Result<Option<JobNode>> {
631        let uuid = Uuid::parse_str(job_id)?;
632        self.get_job_node(pool, &uuid).await
633    }
634
635    async fn get_workflow_jobs(
636        &self,
637        pool: &DatabasePool,
638        workflow_id: &str,
639    ) -> Result<Vec<JobNode>> {
640        let query = r#"
641            SELECT id, queue_name, status, dependency_status, depends_on, dependents, workflow_id, workflow_name
642            FROM hammerwork_jobs 
643            WHERE workflow_id = $1
644            ORDER BY created_at
645        "#;
646
647        match pool {
648            DatabasePool::Postgres(pg_pool) => {
649                let workflow_uuid = Uuid::parse_str(workflow_id)?;
650                let rows = sqlx::query(query)
651                    .bind(workflow_uuid)
652                    .fetch_all(pg_pool)
653                    .await?;
654
655                let mut jobs = Vec::new();
656                for row in rows {
657                    jobs.push(self.postgres_row_to_job_node(&row)?);
658                }
659                Ok(jobs)
660            }
661            DatabasePool::MySQL(mysql_pool) => {
662                let mysql_query = r#"
663                    SELECT id, queue_name, status, dependency_status, depends_on, dependents, workflow_id, workflow_name
664                    FROM hammerwork_jobs 
665                    WHERE workflow_id = ?
666                    ORDER BY created_at
667                "#;
668                let rows = sqlx::query(mysql_query)
669                    .bind(workflow_id)
670                    .fetch_all(mysql_pool)
671                    .await?;
672
673                let mut jobs = Vec::new();
674                for row in rows {
675                    jobs.push(self.mysql_row_to_job_node(&row)?);
676                }
677                Ok(jobs)
678            }
679        }
680    }
681
682    async fn collect_related_jobs(
683        &self,
684        pool: &DatabasePool,
685        target_job: &JobNode,
686    ) -> Result<Vec<JobNode>> {
687        let mut jobs = HashMap::new();
688        let mut to_visit = VecDeque::new();
689        let mut visited = HashSet::new();
690
691        // Start with the target job
692        jobs.insert(target_job.id.clone(), target_job.clone());
693        to_visit.push_back(target_job.id.clone());
694
695        // Traverse both dependencies and dependents
696        while let Some(job_id) = to_visit.pop_front() {
697            if visited.contains(&job_id) {
698                continue;
699            }
700            visited.insert(job_id.clone());
701
702            if let Some(job) = jobs.get(&job_id).cloned() {
703                // Visit all dependencies
704                for dep_id in &job.depends_on {
705                    if !jobs.contains_key(dep_id) {
706                        if let Some(dep_job) = self.get_job_node_by_string(pool, dep_id).await? {
707                            jobs.insert(dep_id.clone(), dep_job);
708                            to_visit.push_back(dep_id.clone());
709                        }
710                    }
711                }
712
713                // Visit all dependents
714                for dep_id in &job.dependents {
715                    if !jobs.contains_key(dep_id) {
716                        if let Some(dep_job) = self.get_job_node_by_string(pool, dep_id).await? {
717                            jobs.insert(dep_id.clone(), dep_job);
718                            to_visit.push_back(dep_id.clone());
719                        }
720                    }
721                }
722            }
723        }
724
725        Ok(jobs.into_values().collect())
726    }
727
728    fn postgres_row_to_job_node(&self, row: &sqlx::postgres::PgRow) -> Result<JobNode> {
729        use sqlx::Row;
730
731        let id: Uuid = row.get("id");
732        let depends_on = self.parse_json_array(row.try_get("depends_on").ok())?;
733        let dependents = self.parse_json_array(row.try_get("dependents").ok())?;
734
735        let workflow_id: Option<String> = match row.try_get::<Option<Uuid>, _>("workflow_id") {
736            Ok(Some(uuid)) => Some(uuid.to_string()),
737            Ok(None) => None,
738            Err(_) => None,
739        };
740
741        Ok(JobNode {
742            id: id.to_string(),
743            queue_name: row.get("queue_name"),
744            status: row.get("status"),
745            dependency_status: row
746                .try_get("dependency_status")
747                .unwrap_or_else(|_| "none".to_string()),
748            depends_on,
749            dependents,
750            workflow_id,
751            workflow_name: row.try_get("workflow_name").ok(),
752        })
753    }
754
755    fn mysql_row_to_job_node(&self, row: &sqlx::mysql::MySqlRow) -> Result<JobNode> {
756        use sqlx::Row;
757
758        let id: String = row.get("id");
759        let depends_on = self.parse_json_array(row.try_get("depends_on").ok())?;
760        let dependents = self.parse_json_array(row.try_get("dependents").ok())?;
761
762        let workflow_id: Option<String> = row.try_get("workflow_id").ok();
763
764        Ok(JobNode {
765            id,
766            queue_name: row.get("queue_name"),
767            status: row.get("status"),
768            dependency_status: row
769                .try_get("dependency_status")
770                .unwrap_or_else(|_| "none".to_string()),
771            depends_on,
772            dependents,
773            workflow_id,
774            workflow_name: row.try_get("workflow_name").ok(),
775        })
776    }
777
778    pub fn parse_json_array(&self, json_value: Option<Value>) -> Result<Vec<String>> {
779        match json_value {
780            Some(Value::Array(arr)) => Ok(arr
781                .into_iter()
782                .filter_map(|v| v.as_str().map(|s| s.to_string()))
783                .collect()),
784            _ => Ok(Vec::new()),
785        }
786    }
787
788    fn print_dependency_tree(&self, jobs: &[JobNode], target_job_id: &str) {
789        let job_map: HashMap<String, &JobNode> =
790            jobs.iter().map(|job| (job.id.clone(), job)).collect();
791
792        // Find root jobs (jobs with no dependencies)
793        let mut roots: Vec<&JobNode> = jobs
794            .iter()
795            .filter(|job| job.depends_on.is_empty())
796            .collect();
797
798        // If no natural roots, use all jobs as potential roots
799        if roots.is_empty() {
800            roots = jobs.iter().collect();
801        }
802
803        // Sort roots by creation order (assuming UUID ordering roughly correlates)
804        roots.sort_by(|a, b| a.id.cmp(&b.id));
805
806        println!("  Tree Structure:");
807        let mut visited = HashSet::new();
808
809        for root in &roots {
810            if !visited.contains(&root.id) {
811                Self::print_job_tree_node(root, &job_map, &mut visited, 0, target_job_id);
812            }
813        }
814
815        // Handle any remaining unvisited jobs (cycles or disconnected components)
816        for job in jobs {
817            if !visited.contains(&job.id) {
818                println!("  [Disconnected]");
819                Self::print_job_tree_node(job, &job_map, &mut visited, 0, target_job_id);
820            }
821        }
822    }
823
824    fn print_job_tree_node(
825        job: &JobNode,
826        job_map: &HashMap<String, &JobNode>,
827        visited: &mut HashSet<String>,
828        depth: usize,
829        target_job_id: &str,
830    ) {
831        if visited.contains(&job.id) {
832            return;
833        }
834        visited.insert(job.id.clone());
835
836        let indent = "  ".repeat(depth + 1);
837        let marker = if depth == 0 { "┌─" } else { "├─" };
838        let highlight = if job.id == target_job_id { " ⭐" } else { "" };
839
840        println!(
841            "{}{}[{}] {} ({}){}",
842            indent,
843            marker,
844            &job.id[..8], // Show first 8 chars of UUID
845            job.status,
846            job.dependency_status,
847            highlight
848        );
849
850        // Recursively print dependents (children in the tree)
851        for dependent_id in &job.dependents {
852            if let Some(dependent_job) = job_map.get(dependent_id) {
853                if !visited.contains(dependent_id) {
854                    Self::print_job_tree_node(
855                        dependent_job,
856                        job_map,
857                        visited,
858                        depth + 1,
859                        target_job_id,
860                    );
861                }
862            }
863        }
864    }
865}
866
867#[cfg(test)]
868mod tests {
869    use super::*;
870    use serde_json::json;
871
872    #[test]
873    fn test_workflow_command_structure() {
874        // Test that commands can be created
875        let list_cmd = WorkflowCommand::List {
876            database_url: None,
877            limit: Some(10),
878            running: true,
879            completed: false,
880            failed: false,
881        };
882
883        // This should compile without errors
884        assert!(matches!(list_cmd, WorkflowCommand::List { .. }));
885    }
886
887    #[test]
888    fn test_workflow_id_parsing() {
889        let test_uuid = "550e8400-e29b-41d4-a716-446655440000";
890        let parsed = Uuid::parse_str(test_uuid);
891        assert!(parsed.is_ok());
892    }
893
894    #[test]
895    fn test_job_node_creation() {
896        let job_node = JobNode {
897            id: "test-job-123".to_string(),
898            queue_name: "test-queue".to_string(),
899            status: "Pending".to_string(),
900            dependency_status: "waiting".to_string(),
901            depends_on: vec!["dep1".to_string(), "dep2".to_string()],
902            dependents: vec!["child1".to_string()],
903            workflow_id: Some("workflow-123".to_string()),
904            workflow_name: Some("test-workflow".to_string()),
905        };
906
907        assert_eq!(job_node.id, "test-job-123");
908        assert_eq!(job_node.depends_on.len(), 2);
909        assert_eq!(job_node.dependents.len(), 1);
910        assert!(job_node.workflow_id.is_some());
911    }
912
913    #[test]
914    fn test_parse_json_array() {
915        let workflow_cmd = WorkflowCommand::List {
916            database_url: None,
917            limit: None,
918            running: false,
919            completed: false,
920            failed: false,
921        };
922
923        // Test valid array
924        let json_array = json!(["job1", "job2", "job3"]);
925        let result = workflow_cmd.parse_json_array(Some(json_array)).unwrap();
926        assert_eq!(result, vec!["job1", "job2", "job3"]);
927
928        // Test empty array
929        let empty_array = json!([]);
930        let result = workflow_cmd.parse_json_array(Some(empty_array)).unwrap();
931        assert!(result.is_empty());
932
933        // Test non-array JSON
934        let non_array = json!({"not": "array"});
935        let result = workflow_cmd.parse_json_array(Some(non_array)).unwrap();
936        assert!(result.is_empty());
937
938        // Test None input
939        let result = workflow_cmd.parse_json_array(None).unwrap();
940        assert!(result.is_empty());
941
942        // Test array with mixed types (should filter non-strings)
943        let mixed_array = json!(["job1", 123, "job2", null, "job3"]);
944        let result = workflow_cmd.parse_json_array(Some(mixed_array)).unwrap();
945        assert_eq!(result, vec!["job1", "job2", "job3"]);
946    }
947
948    #[test]
949    fn test_dependency_level_calculation() {
950        let workflow_cmd = WorkflowCommand::List {
951            database_url: None,
952            limit: None,
953            running: false,
954            completed: false,
955            failed: false,
956        };
957
958        // Create test jobs with dependencies
959        let jobs = vec![
960            JobNode {
961                id: "job1".to_string(),
962                queue_name: "queue1".to_string(),
963                status: "Completed".to_string(),
964                dependency_status: "none".to_string(),
965                depends_on: vec![], // Root job
966                dependents: vec!["job2".to_string()],
967                workflow_id: Some("workflow1".to_string()),
968                workflow_name: Some("test".to_string()),
969            },
970            JobNode {
971                id: "job2".to_string(),
972                queue_name: "queue1".to_string(),
973                status: "Running".to_string(),
974                dependency_status: "satisfied".to_string(),
975                depends_on: vec!["job1".to_string()],
976                dependents: vec!["job3".to_string()],
977                workflow_id: Some("workflow1".to_string()),
978                workflow_name: Some("test".to_string()),
979            },
980            JobNode {
981                id: "job3".to_string(),
982                queue_name: "queue1".to_string(),
983                status: "Pending".to_string(),
984                dependency_status: "waiting".to_string(),
985                depends_on: vec!["job2".to_string()],
986                dependents: vec![],
987                workflow_id: Some("workflow1".to_string()),
988                workflow_name: Some("test".to_string()),
989            },
990        ];
991
992        let levels = workflow_cmd.calculate_dependency_levels(&jobs);
993
994        // Should have 3 levels (0, 1, 2)
995        assert!(levels.len() <= 3);
996
997        // Verify that we have some levels calculated
998        assert!(!levels.is_empty());
999
1000        // Check that levels are properly ordered
1001        let mut level_numbers: Vec<usize> = levels.iter().map(|(level, _)| *level).collect();
1002        level_numbers.sort();
1003
1004        // Should start from 0
1005        assert_eq!(level_numbers[0], 0);
1006    }
1007
1008    #[test]
1009    fn test_graph_format_validation() {
1010        let workflow_cmd = WorkflowCommand::Graph {
1011            database_url: None,
1012            workflow_id: "test-workflow".to_string(),
1013            format: Some("text".to_string()),
1014        };
1015
1016        // Test that command structure is correct
1017        if let WorkflowCommand::Graph { format, .. } = workflow_cmd {
1018            assert_eq!(format, Some("text".to_string()));
1019        } else {
1020            panic!("Expected Graph command");
1021        }
1022    }
1023
1024    #[test]
1025    fn test_dependencies_command_structure() {
1026        let deps_cmd = WorkflowCommand::Dependencies {
1027            database_url: None,
1028            job_id: "test-job-123".to_string(),
1029            tree: true,
1030            dependents: true,
1031        };
1032
1033        if let WorkflowCommand::Dependencies {
1034            job_id,
1035            tree,
1036            dependents,
1037            ..
1038        } = deps_cmd
1039        {
1040            assert_eq!(job_id, "test-job-123");
1041            assert!(tree);
1042            assert!(dependents);
1043        } else {
1044            panic!("Expected Dependencies command");
1045        }
1046    }
1047
1048    #[test]
1049    fn test_create_workflow_validation() {
1050        let create_cmd = WorkflowCommand::Create {
1051            database_url: None,
1052            name: "test-workflow".to_string(),
1053            failure_policy: Some("fail_fast".to_string()),
1054            metadata: Some(r#"{"key": "value"}"#.to_string()),
1055        };
1056
1057        if let WorkflowCommand::Create {
1058            name,
1059            failure_policy,
1060            metadata,
1061            ..
1062        } = create_cmd
1063        {
1064            assert_eq!(name, "test-workflow");
1065            assert_eq!(failure_policy, Some("fail_fast".to_string()));
1066            assert!(metadata.is_some());
1067        } else {
1068            panic!("Expected Create command");
1069        }
1070    }
1071
1072    #[test]
1073    fn test_workflow_command_variants() {
1074        // Test all command variants can be created
1075        let commands = vec![
1076            WorkflowCommand::List {
1077                database_url: None,
1078                limit: Some(10),
1079                running: false,
1080                completed: false,
1081                failed: false,
1082            },
1083            WorkflowCommand::Show {
1084                database_url: None,
1085                workflow_id: "test".to_string(),
1086                dependencies: false,
1087            },
1088            WorkflowCommand::Create {
1089                database_url: None,
1090                name: "test".to_string(),
1091                failure_policy: None,
1092                metadata: None,
1093            },
1094            WorkflowCommand::Cancel {
1095                database_url: None,
1096                workflow_id: "test".to_string(),
1097                force: false,
1098            },
1099            WorkflowCommand::Dependencies {
1100                database_url: None,
1101                job_id: "test".to_string(),
1102                tree: false,
1103                dependents: false,
1104            },
1105            WorkflowCommand::Graph {
1106                database_url: None,
1107                workflow_id: "test".to_string(),
1108                format: None,
1109            },
1110        ];
1111
1112        // All commands should be valid
1113        assert_eq!(commands.len(), 6);
1114    }
1115
1116    #[test]
1117    fn test_job_node_dependency_relationships() {
1118        let parent = JobNode {
1119            id: "parent".to_string(),
1120            queue_name: "queue1".to_string(),
1121            status: "Completed".to_string(),
1122            dependency_status: "none".to_string(),
1123            depends_on: vec![],
1124            dependents: vec!["child1".to_string(), "child2".to_string()],
1125            workflow_id: Some("workflow1".to_string()),
1126            workflow_name: Some("test".to_string()),
1127        };
1128
1129        let child1 = JobNode {
1130            id: "child1".to_string(),
1131            queue_name: "queue1".to_string(),
1132            status: "Running".to_string(),
1133            dependency_status: "satisfied".to_string(),
1134            depends_on: vec!["parent".to_string()],
1135            dependents: vec![],
1136            workflow_id: Some("workflow1".to_string()),
1137            workflow_name: Some("test".to_string()),
1138        };
1139
1140        let child2 = JobNode {
1141            id: "child2".to_string(),
1142            queue_name: "queue1".to_string(),
1143            status: "Pending".to_string(),
1144            dependency_status: "waiting".to_string(),
1145            depends_on: vec!["parent".to_string()],
1146            dependents: vec![],
1147            workflow_id: Some("workflow1".to_string()),
1148            workflow_name: Some("test".to_string()),
1149        };
1150
1151        // Verify relationships
1152        assert!(parent.depends_on.is_empty());
1153        assert_eq!(parent.dependents.len(), 2);
1154        assert!(parent.dependents.contains(&"child1".to_string()));
1155        assert!(parent.dependents.contains(&"child2".to_string()));
1156
1157        assert_eq!(child1.depends_on.len(), 1);
1158        assert!(child1.depends_on.contains(&"parent".to_string()));
1159        assert!(child1.dependents.is_empty());
1160
1161        assert_eq!(child2.depends_on.len(), 1);
1162        assert!(child2.depends_on.contains(&"parent".to_string()));
1163        assert!(child2.dependents.is_empty());
1164    }
1165
1166    #[test]
1167    fn test_database_url_extraction() {
1168        let config = Config {
1169            database_url: Some("postgres://localhost/test".to_string()),
1170            default_queue: None,
1171            default_limit: None,
1172            log_level: None,
1173            connection_pool_size: None,
1174        };
1175
1176        let workflow_cmd = WorkflowCommand::List {
1177            database_url: Some("postgres://override/test".to_string()),
1178            limit: None,
1179            running: false,
1180            completed: false,
1181            failed: false,
1182        };
1183
1184        // Test database URL extraction logic
1185        let result = workflow_cmd.get_database_url(&config);
1186        assert!(result.is_ok());
1187        assert_eq!(result.unwrap(), "postgres://override/test");
1188
1189        // Test fallback to config
1190        let workflow_cmd_no_url = WorkflowCommand::List {
1191            database_url: None,
1192            limit: None,
1193            running: false,
1194            completed: false,
1195            failed: false,
1196        };
1197
1198        let result = workflow_cmd_no_url.get_database_url(&config);
1199        assert!(result.is_ok());
1200        assert_eq!(result.unwrap(), "postgres://localhost/test");
1201    }
1202
1203    #[test]
1204    fn test_job_status_classification() {
1205        let statuses = vec!["Completed", "Failed", "Running", "Pending", "Unknown"];
1206
1207        for status in statuses {
1208            let job = JobNode {
1209                id: "test".to_string(),
1210                queue_name: "queue1".to_string(),
1211                status: status.to_string(),
1212                dependency_status: "none".to_string(),
1213                depends_on: vec![],
1214                dependents: vec![],
1215                workflow_id: None,
1216                workflow_name: None,
1217            };
1218
1219            // Verify status is preserved
1220            assert_eq!(job.status, status);
1221        }
1222    }
1223
1224    #[test]
1225    fn test_dependency_status_types() {
1226        let dep_statuses = vec!["none", "waiting", "satisfied", "failed"];
1227
1228        for dep_status in dep_statuses {
1229            let job = JobNode {
1230                id: "test".to_string(),
1231                queue_name: "queue1".to_string(),
1232                status: "Pending".to_string(),
1233                dependency_status: dep_status.to_string(),
1234                depends_on: vec![],
1235                dependents: vec![],
1236                workflow_id: None,
1237                workflow_name: None,
1238            };
1239
1240            // Verify dependency status is preserved
1241            assert_eq!(job.dependency_status, dep_status);
1242        }
1243    }
1244}