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};
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(long = "worker-tag", value_name = "TAG")]
77 worker_tags: Vec<String>,
78 },
79 List {
81 #[arg(long)]
83 status: Option<String>,
84 #[arg(long)]
86 workflow: Option<String>,
87 #[arg(long)]
91 created_by: Option<Uuid>,
92 #[arg(long)]
94 concurrency_group: Option<String>,
95 #[arg(long)]
97 page: Option<u32>,
98 #[arg(long)]
100 per_page: Option<u32>,
101 },
102 Get {
104 id: Uuid,
106 },
107 Cancel {
109 id: Uuid,
111 },
112 Approve {
114 id: Uuid,
116 },
117 Reject {
119 id: Uuid,
121 },
122 Input {
124 id: Uuid,
126 step_id: Uuid,
128 #[arg(long, group = "value_source")]
130 value: Option<String>,
131 #[arg(long, group = "value_source")]
133 value_file: Option<PathBuf>,
134 },
135 RejectInput {
137 id: Uuid,
139 step_id: Uuid,
141 #[arg(long)]
143 reason: Option<String>,
144 },
145 Retry {
147 id: Uuid,
149 #[arg(long)]
152 force: bool,
153 },
154 Replay {
156 id: Uuid,
158 },
159 Watch {
161 id: Uuid,
163 #[arg(long)]
165 no_logs: bool,
166 #[arg(long, value_parser = parse_humantime)]
168 timeout: Option<Duration>,
169 },
170 Plan {
172 workflow: String,
174 #[arg(long, group = "plan_input_source")]
176 input: Option<String>,
177 #[arg(long, group = "plan_input_source")]
179 input_file: Option<PathBuf>,
180 #[arg(long)]
182 max_depth: Option<u32>,
183 #[arg(long)]
185 no_estimates: bool,
186 },
187 Diff {
189 run_a: Uuid,
191 run_b: Uuid,
193 },
194}
195
196fn parse_humantime(s: &str) -> Result<Duration, String> {
198 humantime::parse_duration(s).map_err(|e| e.to_string())
199}
200
201fn parse_concurrency_limit(s: &str) -> Result<ConcurrencyLimit, String> {
206 let (group, limit) = shared_parse_concurrency_limit(s)?;
207 let limit =
208 i32::try_from(limit).map_err(|e| format!("invalid limit '{limit}' in '{s}': {e}"))?;
209 Ok(ConcurrencyLimit { group, limit })
210}
211
212const TERMINAL_EVENTS: &[&str] = &["run_completed", "run_failed", "run_cancelled"];
214
215fn resolve_payload(payload: Option<&str>, payload_file: Option<&PathBuf>) -> Result<Value> {
217 match (payload, payload_file) {
218 (Some(raw), _) => from_str(raw).context("invalid JSON in --payload"),
219 (_, Some(path)) => {
220 let content = fs::read_to_string(path)
221 .with_context(|| format!("cannot read payload file: {}", path.display()))?;
222 from_str(&content).with_context(|| format!("invalid JSON in {}", path.display()))
223 }
224 (None, None) => Ok(Value::Object(Map::new())),
225 }
226}
227
228fn validate_max_cost(max_cost: Option<f64>) -> Result<()> {
237 match max_cost {
238 Some(value) if !value.is_finite() => {
239 anyhow::bail!("--max-cost must be a finite number, got {value}")
240 }
241 Some(value) if value < 0.0 => {
242 anyhow::bail!("--max-cost must be zero or positive, got {value}")
243 }
244 _ => Ok(()),
245 }
246}
247
248pub async fn execute(
254 client: &IronflowClient,
255 args: &RunArgs,
256 json_mode: bool,
257 _verbose: bool,
258) -> Result<()> {
259 match &args.command {
260 RunCommands::Create {
261 workflow,
262 payload,
263 payload_file,
264 max_retries,
265 idempotency_key,
266 max_cost,
267 concurrency_key,
268 concurrency_limits,
269 worker_tags,
270 } => {
271 validate_max_cost(*max_cost)?;
272 let payload_value = resolve_payload(payload.as_deref(), payload_file.as_ref())?;
273 let payload_map = payload_value
274 .as_object()
275 .context("payload must be a JSON object")?
276 .clone();
277 let request: CreateRunRequest = CreateRunRequest::builder()
278 .workflow(workflow.clone())
279 .payload(Some(payload_map))
280 .max_retries(max_retries.map(|n| n as i32))
283 .max_cost_usd(*max_cost)
284 .concurrency_key(concurrency_key.clone())
285 .concurrency_limits(concurrency_limits.clone())
286 .worker_tags(worker_tags.clone())
287 .try_into()
288 .context("failed to build CreateRunRequest")?;
289
290 let response = match idempotency_key {
291 Some(key) => client.create_run_idempotent(&request, key).await?,
292 None => client.create_run(&request).await?,
293 };
294 output::print_output(json_mode, &response, || {
295 output::runs_table(slice::from_ref(&response.data))
296 })?;
297 }
298 RunCommands::List {
299 status,
300 workflow,
301 created_by,
302 concurrency_group,
303 page,
304 per_page,
305 } => {
306 let filter = ListRunsFilter {
307 status: status.as_deref(),
308 workflow: workflow.as_deref(),
309 created_by: *created_by,
310 concurrency_group: concurrency_group.as_deref(),
311 page: *page,
312 per_page: *per_page,
313 ..Default::default()
314 };
315 let response = client.list_runs_filtered(&filter).await?;
316 output::print_output(json_mode, &response, || output::runs_table(&response.data))?;
317 }
318 RunCommands::Get { id } => {
319 let response = client.get_run(*id).await?;
320 output::print_output(json_mode, &response, || {
321 output::run_detail_table(&response.data)
322 })?;
323
324 if !json_mode && !response.data.steps.is_empty() {
325 let mut out = stdout().lock();
326 writeln!(out)?;
327 writeln!(out, "Steps:")?;
328 writeln!(out, "{}", output::steps_table(&response.data.steps))?;
329 }
330 }
331 RunCommands::Cancel { id } => {
332 let response = client.cancel_run(*id).await?;
333 output::print_output(json_mode, &response, || {
334 output::cancelled_table(&response.data)
335 })?;
336 }
337 RunCommands::Approve { id } => {
338 let response = client.approve_run(*id).await?;
339 output::print_output(json_mode, &response, || {
340 output::runs_table(slice::from_ref(&response.data))
341 })?;
342 if !json_mode && matches!(response.data.status, RunStatus::AwaitingApproval) {
344 println!("Approval recorded; more approvals are required.");
345 }
346 }
347 RunCommands::Reject { id } => {
348 let response = client.reject_run(*id).await?;
349 output::print_output(json_mode, &response, || {
350 output::runs_table(slice::from_ref(&response.data))
351 })?;
352 }
353 RunCommands::Input {
354 id,
355 step_id,
356 value,
357 value_file,
358 } => {
359 let answer = resolve_payload(value.as_deref(), value_file.as_ref())?;
360 let response = client.submit_input(*id, *step_id, &answer).await?;
361 output::print_output(json_mode, &response, || {
362 output::runs_table(slice::from_ref(&response.data))
363 })?;
364 }
365 RunCommands::RejectInput {
366 id,
367 step_id,
368 reason,
369 } => {
370 let response = client
371 .reject_input(*id, *step_id, reason.as_deref())
372 .await?;
373 output::print_output(json_mode, &response, || {
374 output::runs_table(slice::from_ref(&response.data))
375 })?;
376 }
377 RunCommands::Retry { id, force } => {
378 let response = client.retry_run(*id, *force).await?;
379 output::print_output(json_mode, &response, || {
380 output::runs_table(slice::from_ref(&response.data))
381 })?;
382 }
383 RunCommands::Replay { id } => {
384 let response = client.replay_run(*id).await?;
385 output::print_output(json_mode, &response, || {
386 output::runs_table(slice::from_ref(&response.data))
387 })?;
388 }
389 RunCommands::Watch {
390 id,
391 no_logs,
392 timeout,
393 } => {
394 execute_watch(client, *id, *no_logs, *timeout, json_mode).await?;
395 }
396 RunCommands::Plan {
397 workflow,
398 input,
399 input_file,
400 max_depth,
401 no_estimates,
402 } => {
403 let payload = resolve_payload(input.as_deref(), input_file.as_ref())?;
404 let payload_map = payload
405 .as_object()
406 .context("input must be a JSON object")?
407 .clone();
408 let request: PlanWorkflowRequest = PlanWorkflowRequest::builder()
409 .payload(Some(payload_map))
410 .max_depth(max_depth.map(|d| d as i32))
413 .estimate_durations(Some(!*no_estimates))
414 .try_into()
415 .context("failed to build PlanWorkflowRequest")?;
416 let response = client.plan_workflow(workflow, &request).await?;
417 output::render_execution_plan(&mut stdout().lock(), json_mode, &response)?;
418 }
419 RunCommands::Diff { run_a, run_b } => {
420 execute_diff(client, *run_a, *run_b, json_mode).await?;
421 }
422 }
423 Ok(())
424}
425
426async fn execute_watch(
428 client: &IronflowClient,
429 run_id: Uuid,
430 no_logs: bool,
431 timeout: Option<Duration>,
432 json_mode: bool,
433) -> Result<()> {
434 let run = client.get_run(run_id).await?;
435 let status = run.data.run.status;
436 if matches!(
437 status,
438 RunStatus::Completed | RunStatus::Failed | RunStatus::Cancelled
439 ) {
440 if json_mode {
441 output::print_output(json_mode, &run, || output::run_detail_table(&run.data))?;
442 } else {
443 let mut out = stdout().lock();
444 writeln!(out, "Run {run_id} already in terminal state: {status}")?;
445 }
446 return Ok(());
447 }
448
449 let watch_fut = async {
450 let mut stream = client.events(Some(run_id), None).await?;
451 let mut out = stdout().lock();
452
453 while let Some(event) = stream.next().await {
454 match event {
455 Ok(ev) => {
456 if no_logs
457 && !ev.event_type.starts_with("run_")
458 && !ev.event_type.starts_with("step_")
459 {
460 continue;
461 }
462
463 if json_mode {
464 let obj = json!({
465 "event": ev.event_type,
466 "data": ev.data,
467 });
468 writeln!(out, "{}", to_string(&obj)?)?;
469 } else {
470 writeln!(out, "[{}] {}", ev.event_type, ev.data)?;
471 }
472
473 if TERMINAL_EVENTS.contains(&ev.event_type.as_str()) {
474 break;
475 }
476 }
477 Err(e) => {
478 return Err(anyhow!("SSE stream error: {e}"));
479 }
480 }
481 }
482
483 Ok::<(), anyhow::Error>(())
484 };
485
486 match timeout {
487 Some(dur) => {
488 tokio_timeout(dur, watch_fut).await.unwrap_or_else(|_| {
489 eprintln!("Timeout reached after {}", format_duration(dur));
490 Ok(())
491 })?;
492 }
493 None => {
494 watch_fut.await?;
495 }
496 }
497
498 Ok(())
499}
500
501async fn execute_diff(
503 client: &IronflowClient,
504 run_a_id: Uuid,
505 run_b_id: Uuid,
506 json_mode: bool,
507) -> Result<()> {
508 if run_a_id == run_b_id {
509 bail!("both run IDs are the same; nothing to diff");
510 }
511
512 let (a, b) = tokio::try_join!(client.get_run(run_a_id), client.get_run(run_b_id))?;
513
514 if a.data.run.workflow_name != b.data.run.workflow_name {
515 bail!(
516 "cannot diff runs from different workflows: '{}' vs '{}'",
517 a.data.run.workflow_name,
518 b.data.run.workflow_name
519 );
520 }
521
522 if json_mode {
523 let diff = json!({
524 "run_a": a.data,
525 "run_b": b.data,
526 });
527 output::print_json(&diff)?;
528 } else {
529 let table = output::run_diff_table(&a.data, &b.data);
530 let mut out = stdout().lock();
531 writeln!(out, "{table}")?;
532 }
533
534 Ok(())
535}
536
537#[cfg(test)]
538mod tests {
539 use std::io::Write;
540
541 use tempfile::NamedTempFile;
542
543 use super::*;
544
545 #[test]
546 fn parse_concurrency_limit_reads_group_and_limit() {
547 let limit = parse_concurrency_limit("repo:acme=2").unwrap();
548 assert_eq!(limit.group, "repo:acme");
549 assert_eq!(limit.limit, 2);
550 }
551
552 #[test]
553 fn parse_concurrency_limit_rejects_a_limit_above_i32() {
554 assert!(parse_concurrency_limit("repo:acme=4294967295").is_err());
555 }
556
557 #[test]
558 fn parse_concurrency_limit_propagates_shape_errors() {
559 assert!(parse_concurrency_limit("repo:acme").is_err());
560 }
561
562 #[test]
563 fn resolve_payload_none_returns_empty_object() {
564 let value = resolve_payload(None, None).unwrap();
565 assert!(value.is_object());
566 assert!(value.as_object().unwrap().is_empty());
567 }
568
569 #[test]
570 fn resolve_payload_inline_valid_json() {
571 let value = resolve_payload(Some(r#"{"key": "value"}"#), None).unwrap();
572 assert_eq!(value["key"], "value");
573 }
574
575 #[test]
576 fn resolve_payload_inline_invalid_json() {
577 let result = resolve_payload(Some("not json"), None);
578 assert!(result.is_err());
579 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
580 }
581
582 #[test]
583 fn resolve_payload_file_valid() {
584 let mut tmp = NamedTempFile::new().unwrap();
585 write!(tmp, r#"{{"workflow": "test"}}"#).unwrap();
586 let path = tmp.path().to_path_buf();
587
588 let value = resolve_payload(None, Some(&path)).unwrap();
589 assert_eq!(value["workflow"], "test");
590 }
591
592 #[test]
593 fn resolve_payload_file_not_found() {
594 let path = PathBuf::from("/nonexistent/payload.json");
595 let result = resolve_payload(None, Some(&path));
596 assert!(result.is_err());
597 assert!(result.unwrap_err().to_string().contains("cannot read"));
598 }
599
600 #[test]
601 fn resolve_payload_file_invalid_json() {
602 let mut tmp = NamedTempFile::new().unwrap();
603 write!(tmp, "not valid json").unwrap();
604 let path = tmp.path().to_path_buf();
605
606 let result = resolve_payload(None, Some(&path));
607 assert!(result.is_err());
608 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
609 }
610}