1use std::fs;
5use std::io::{Write as _, stdout};
6use std::path::PathBuf;
7use std::slice;
8use std::time::Duration;
9
10use anyhow::{Context, Result, anyhow, bail};
11use clap::{Args, Subcommand, value_parser};
12use futures_util::StreamExt;
13use humantime::format_duration;
14use ironflow_sdk::IronflowClient;
15use ironflow_sdk::client::ListRunsFilter;
16use ironflow_sdk::types::{ConcurrencyLimit, CreateRunRequest, PlanWorkflowRequest, RunStatus};
17use ironflow_types::parse_concurrency_limit as shared_parse_concurrency_limit;
18use serde_json::{Map, Value, from_str, json, to_string};
19use tokio::time::timeout as tokio_timeout;
20use uuid::Uuid;
21
22use crate::output;
23
24#[derive(Debug, Args)]
26pub struct RunArgs {
27 #[command(subcommand)]
29 pub command: RunCommands,
30}
31
32#[derive(Debug, Subcommand)]
34pub enum RunCommands {
35 Create {
37 workflow: String,
39 #[arg(long, group = "payload_source")]
41 payload: Option<String>,
42 #[arg(long, group = "payload_source")]
44 payload_file: Option<PathBuf>,
45 #[arg(long)]
48 max_retries: Option<u32>,
49 #[arg(long)]
55 idempotency_key: Option<String>,
56 #[arg(long = "max-cost", value_name = "USD")]
59 max_cost: Option<f64>,
60 #[arg(long)]
64 concurrency_key: Option<String>,
65 #[arg(
69 long = "concurrency-limit",
70 value_name = "GROUP=N",
71 value_parser = parse_concurrency_limit
72 )]
73 concurrency_limits: Vec<ConcurrencyLimit>,
74 #[arg(
78 long,
79 allow_negative_numbers = true,
80 value_parser = value_parser!(i16).range(-100..=100)
81 )]
82 priority: Option<i16>,
83 #[arg(long = "worker-tag", value_name = "TAG")]
86 worker_tags: Vec<String>,
87 },
88 List {
90 #[arg(long)]
92 status: Option<String>,
93 #[arg(long)]
95 workflow: Option<String>,
96 #[arg(long)]
100 created_by: Option<Uuid>,
101 #[arg(long)]
103 concurrency_group: Option<String>,
104 #[arg(
106 long,
107 allow_negative_numbers = true,
108 value_parser = value_parser!(i16).range(-100..=100)
109 )]
110 priority: Option<i16>,
111 #[arg(long)]
113 page: Option<u32>,
114 #[arg(long)]
116 per_page: Option<u32>,
117 },
118 Get {
120 id: Uuid,
122 },
123 Cancel {
125 id: Uuid,
127 },
128 Approve {
130 id: Uuid,
132 },
133 Reject {
135 id: Uuid,
137 },
138 Input {
140 id: Uuid,
142 step_id: Uuid,
144 #[arg(long, group = "value_source")]
146 value: Option<String>,
147 #[arg(long, group = "value_source")]
149 value_file: Option<PathBuf>,
150 },
151 RejectInput {
153 id: Uuid,
155 step_id: Uuid,
157 #[arg(long)]
159 reason: Option<String>,
160 },
161 Retry {
163 id: Uuid,
165 #[arg(long)]
168 force: bool,
169 },
170 Replay {
172 id: Uuid,
174 },
175 Watch {
177 id: Uuid,
179 #[arg(long)]
181 no_logs: bool,
182 #[arg(long, value_parser = parse_humantime)]
184 timeout: Option<Duration>,
185 },
186 Plan {
188 workflow: String,
190 #[arg(long, group = "plan_input_source")]
192 input: Option<String>,
193 #[arg(long, group = "plan_input_source")]
195 input_file: Option<PathBuf>,
196 #[arg(long)]
198 max_depth: Option<u32>,
199 #[arg(long)]
201 no_estimates: bool,
202 },
203 Diff {
205 run_a: Uuid,
207 run_b: Uuid,
209 },
210}
211
212fn parse_humantime(s: &str) -> Result<Duration, String> {
214 humantime::parse_duration(s).map_err(|e| e.to_string())
215}
216
217fn parse_concurrency_limit(s: &str) -> Result<ConcurrencyLimit, String> {
222 let (group, limit) = shared_parse_concurrency_limit(s)?;
223 let limit =
224 i32::try_from(limit).map_err(|e| format!("invalid limit '{limit}' in '{s}': {e}"))?;
225 Ok(ConcurrencyLimit { group, limit })
226}
227
228const TERMINAL_EVENTS: &[&str] = &["run_completed", "run_failed", "run_cancelled"];
230
231fn resolve_payload(payload: Option<&str>, payload_file: Option<&PathBuf>) -> Result<Value> {
233 match (payload, payload_file) {
234 (Some(raw), _) => from_str(raw).context("invalid JSON in --payload"),
235 (_, Some(path)) => {
236 let content = fs::read_to_string(path)
237 .with_context(|| format!("cannot read payload file: {}", path.display()))?;
238 from_str(&content).with_context(|| format!("invalid JSON in {}", path.display()))
239 }
240 (None, None) => Ok(Value::Object(Map::new())),
241 }
242}
243
244fn validate_max_cost(max_cost: Option<f64>) -> Result<()> {
253 match max_cost {
254 Some(value) if !value.is_finite() => {
255 anyhow::bail!("--max-cost must be a finite number, got {value}")
256 }
257 Some(value) if value < 0.0 => {
258 anyhow::bail!("--max-cost must be zero or positive, got {value}")
259 }
260 _ => Ok(()),
261 }
262}
263
264pub async fn execute(
270 client: &IronflowClient,
271 args: &RunArgs,
272 json_mode: bool,
273 _verbose: bool,
274) -> Result<()> {
275 match &args.command {
276 RunCommands::Create {
277 workflow,
278 payload,
279 payload_file,
280 max_retries,
281 idempotency_key,
282 max_cost,
283 concurrency_key,
284 concurrency_limits,
285 priority,
286 worker_tags,
287 } => {
288 validate_max_cost(*max_cost)?;
289 let payload_value = resolve_payload(payload.as_deref(), payload_file.as_ref())?;
290 let payload_map = payload_value
291 .as_object()
292 .context("payload must be a JSON object")?
293 .clone();
294 let request: CreateRunRequest = CreateRunRequest::builder()
295 .workflow(workflow.clone())
296 .payload(Some(payload_map))
297 .max_retries(max_retries.map(|n| n as i32))
300 .max_cost_usd(*max_cost)
301 .concurrency_key(concurrency_key.clone())
302 .concurrency_limits(concurrency_limits.clone())
303 .priority(priority.map(i32::from))
304 .worker_tags(worker_tags.clone())
305 .try_into()
306 .context("failed to build CreateRunRequest")?;
307
308 let response = match idempotency_key {
309 Some(key) => client.create_run_idempotent(&request, key).await?,
310 None => client.create_run(&request).await?,
311 };
312 output::print_output(json_mode, &response, || {
313 output::runs_table(slice::from_ref(&response.data))
314 })?;
315 }
316 RunCommands::List {
317 status,
318 workflow,
319 created_by,
320 concurrency_group,
321 priority,
322 page,
323 per_page,
324 } => {
325 let filter = ListRunsFilter {
326 status: status.as_deref(),
327 workflow: workflow.as_deref(),
328 created_by: *created_by,
329 concurrency_group: concurrency_group.as_deref(),
330 priority: *priority,
331 page: *page,
332 per_page: *per_page,
333 ..Default::default()
334 };
335 let response = client.list_runs_filtered(&filter).await?;
336 output::print_output(json_mode, &response, || output::runs_table(&response.data))?;
337 }
338 RunCommands::Get { id } => {
339 let response = client.get_run(*id).await?;
340 output::print_output(json_mode, &response, || {
341 output::run_detail_table(&response.data)
342 })?;
343
344 if !json_mode && !response.data.steps.is_empty() {
345 let mut out = stdout().lock();
346 writeln!(out)?;
347 writeln!(out, "Steps:")?;
348 writeln!(out, "{}", output::steps_table(&response.data.steps))?;
349 }
350 }
351 RunCommands::Cancel { id } => {
352 let response = client.cancel_run(*id).await?;
353 output::print_output(json_mode, &response, || {
354 output::cancelled_table(&response.data)
355 })?;
356 }
357 RunCommands::Approve { id } => {
358 let response = client.approve_run(*id).await?;
359 output::print_output(json_mode, &response, || {
360 output::runs_table(slice::from_ref(&response.data))
361 })?;
362 if !json_mode && matches!(response.data.status, RunStatus::AwaitingApproval) {
364 println!("Approval recorded; more approvals are required.");
365 }
366 }
367 RunCommands::Reject { id } => {
368 let response = client.reject_run(*id).await?;
369 output::print_output(json_mode, &response, || {
370 output::runs_table(slice::from_ref(&response.data))
371 })?;
372 }
373 RunCommands::Input {
374 id,
375 step_id,
376 value,
377 value_file,
378 } => {
379 let answer = resolve_payload(value.as_deref(), value_file.as_ref())?;
380 let response = client.submit_input(*id, *step_id, &answer).await?;
381 output::print_output(json_mode, &response, || {
382 output::runs_table(slice::from_ref(&response.data))
383 })?;
384 }
385 RunCommands::RejectInput {
386 id,
387 step_id,
388 reason,
389 } => {
390 let response = client
391 .reject_input(*id, *step_id, reason.as_deref())
392 .await?;
393 output::print_output(json_mode, &response, || {
394 output::runs_table(slice::from_ref(&response.data))
395 })?;
396 }
397 RunCommands::Retry { id, force } => {
398 let response = client.retry_run(*id, *force).await?;
399 output::print_output(json_mode, &response, || {
400 output::runs_table(slice::from_ref(&response.data))
401 })?;
402 }
403 RunCommands::Replay { id } => {
404 let response = client.replay_run(*id).await?;
405 output::print_output(json_mode, &response, || {
406 output::runs_table(slice::from_ref(&response.data))
407 })?;
408 }
409 RunCommands::Watch {
410 id,
411 no_logs,
412 timeout,
413 } => {
414 execute_watch(client, *id, *no_logs, *timeout, json_mode).await?;
415 }
416 RunCommands::Plan {
417 workflow,
418 input,
419 input_file,
420 max_depth,
421 no_estimates,
422 } => {
423 let payload = resolve_payload(input.as_deref(), input_file.as_ref())?;
424 let payload_map = payload
425 .as_object()
426 .context("input must be a JSON object")?
427 .clone();
428 let request: PlanWorkflowRequest = PlanWorkflowRequest::builder()
429 .payload(Some(payload_map))
430 .max_depth(max_depth.map(|d| d as i32))
433 .estimate_durations(Some(!*no_estimates))
434 .try_into()
435 .context("failed to build PlanWorkflowRequest")?;
436 let response = client.plan_workflow(workflow, &request).await?;
437 output::render_execution_plan(&mut stdout().lock(), json_mode, &response)?;
438 }
439 RunCommands::Diff { run_a, run_b } => {
440 execute_diff(client, *run_a, *run_b, json_mode).await?;
441 }
442 }
443 Ok(())
444}
445
446async fn execute_watch(
448 client: &IronflowClient,
449 run_id: Uuid,
450 no_logs: bool,
451 timeout: Option<Duration>,
452 json_mode: bool,
453) -> Result<()> {
454 let run = client.get_run(run_id).await?;
455 let status = run.data.run.status;
456 if matches!(
457 status,
458 RunStatus::Completed | RunStatus::Failed | RunStatus::Cancelled
459 ) {
460 if json_mode {
461 output::print_output(json_mode, &run, || output::run_detail_table(&run.data))?;
462 } else {
463 let mut out = stdout().lock();
464 writeln!(out, "Run {run_id} already in terminal state: {status}")?;
465 }
466 return Ok(());
467 }
468
469 let watch_fut = async {
470 let mut stream = client.events(Some(run_id), None).await?;
471 let mut out = stdout().lock();
472
473 while let Some(event) = stream.next().await {
474 match event {
475 Ok(ev) => {
476 if no_logs
477 && !ev.event_type.starts_with("run_")
478 && !ev.event_type.starts_with("step_")
479 {
480 continue;
481 }
482
483 if json_mode {
484 let obj = json!({
485 "event": ev.event_type,
486 "data": ev.data,
487 });
488 writeln!(out, "{}", to_string(&obj)?)?;
489 } else {
490 writeln!(out, "[{}] {}", ev.event_type, ev.data)?;
491 }
492
493 if TERMINAL_EVENTS.contains(&ev.event_type.as_str()) {
494 break;
495 }
496 }
497 Err(e) => {
498 return Err(anyhow!("SSE stream error: {e}"));
499 }
500 }
501 }
502
503 Ok::<(), anyhow::Error>(())
504 };
505
506 match timeout {
507 Some(dur) => {
508 tokio_timeout(dur, watch_fut).await.unwrap_or_else(|_| {
509 eprintln!("Timeout reached after {}", format_duration(dur));
510 Ok(())
511 })?;
512 }
513 None => {
514 watch_fut.await?;
515 }
516 }
517
518 Ok(())
519}
520
521async fn execute_diff(
523 client: &IronflowClient,
524 run_a_id: Uuid,
525 run_b_id: Uuid,
526 json_mode: bool,
527) -> Result<()> {
528 if run_a_id == run_b_id {
529 bail!("both run IDs are the same; nothing to diff");
530 }
531
532 let (a, b) = tokio::try_join!(client.get_run(run_a_id), client.get_run(run_b_id))?;
533
534 if a.data.run.workflow_name != b.data.run.workflow_name {
535 bail!(
536 "cannot diff runs from different workflows: '{}' vs '{}'",
537 a.data.run.workflow_name,
538 b.data.run.workflow_name
539 );
540 }
541
542 if json_mode {
543 let diff = json!({
544 "run_a": a.data,
545 "run_b": b.data,
546 });
547 output::print_json(&diff)?;
548 } else {
549 let table = output::run_diff_table(&a.data, &b.data);
550 let mut out = stdout().lock();
551 writeln!(out, "{table}")?;
552 }
553
554 Ok(())
555}
556
557#[cfg(test)]
558mod tests {
559 use std::io::Write;
560
561 use tempfile::NamedTempFile;
562
563 use super::*;
564
565 #[test]
566 fn parse_concurrency_limit_reads_group_and_limit() {
567 let limit = parse_concurrency_limit("repo:acme=2").unwrap();
568 assert_eq!(limit.group, "repo:acme");
569 assert_eq!(limit.limit, 2);
570 }
571
572 #[test]
573 fn parse_concurrency_limit_rejects_a_limit_above_i32() {
574 assert!(parse_concurrency_limit("repo:acme=4294967295").is_err());
575 }
576
577 #[test]
578 fn parse_concurrency_limit_propagates_shape_errors() {
579 assert!(parse_concurrency_limit("repo:acme").is_err());
580 }
581
582 #[test]
583 fn resolve_payload_none_returns_empty_object() {
584 let value = resolve_payload(None, None).unwrap();
585 assert!(value.is_object());
586 assert!(value.as_object().unwrap().is_empty());
587 }
588
589 #[test]
590 fn resolve_payload_inline_valid_json() {
591 let value = resolve_payload(Some(r#"{"key": "value"}"#), None).unwrap();
592 assert_eq!(value["key"], "value");
593 }
594
595 #[test]
596 fn resolve_payload_inline_invalid_json() {
597 let result = resolve_payload(Some("not json"), None);
598 assert!(result.is_err());
599 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
600 }
601
602 #[test]
603 fn resolve_payload_file_valid() {
604 let mut tmp = NamedTempFile::new().unwrap();
605 write!(tmp, r#"{{"workflow": "test"}}"#).unwrap();
606 let path = tmp.path().to_path_buf();
607
608 let value = resolve_payload(None, Some(&path)).unwrap();
609 assert_eq!(value["workflow"], "test");
610 }
611
612 #[test]
613 fn resolve_payload_file_not_found() {
614 let path = PathBuf::from("/nonexistent/payload.json");
615 let result = resolve_payload(None, Some(&path));
616 assert!(result.is_err());
617 assert!(result.unwrap_err().to_string().contains("cannot read"));
618 }
619
620 #[test]
621 fn resolve_payload_file_invalid_json() {
622 let mut tmp = NamedTempFile::new().unwrap();
623 write!(tmp, "not valid json").unwrap();
624 let path = tmp.path().to_path_buf();
625
626 let result = resolve_payload(None, Some(&path));
627 assert!(result.is_err());
628 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
629 }
630}