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