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 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 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 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 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 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 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 if show_tree {
308 println!("\nDependency Tree:");
309
310 let jobs = if let Some(workflow_id) = &target_job.workflow_id {
312 self.get_workflow_jobs(&pool, workflow_id).await?
313 } else {
314 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 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 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 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 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 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 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 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 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 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 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 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 jobs.insert(target_job.id.clone(), target_job.clone());
693 to_visit.push_back(target_job.id.clone());
694
695 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 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 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 let mut roots: Vec<&JobNode> = jobs
794 .iter()
795 .filter(|job| job.depends_on.is_empty())
796 .collect();
797
798 if roots.is_empty() {
800 roots = jobs.iter().collect();
801 }
802
803 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 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], job.status,
846 job.dependency_status,
847 highlight
848 );
849
850 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 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 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 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 let empty_array = json!([]);
930 let result = workflow_cmd.parse_json_array(Some(empty_array)).unwrap();
931 assert!(result.is_empty());
932
933 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 let result = workflow_cmd.parse_json_array(None).unwrap();
940 assert!(result.is_empty());
941
942 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 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![], 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 assert!(levels.len() <= 3);
996
997 assert!(!levels.is_empty());
999
1000 let mut level_numbers: Vec<usize> = levels.iter().map(|(level, _)| *level).collect();
1002 level_numbers.sort();
1003
1004 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 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 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 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 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 let result = workflow_cmd.get_database_url(&config);
1186 assert!(result.is_ok());
1187 assert_eq!(result.unwrap(), "postgres://override/test");
1188
1189 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 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 assert_eq!(job.dependency_status, dep_status);
1242 }
1243 }
1244}