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::{CreateRunRequest, PlanWorkflowRequest, RunStatus};
17use serde_json::{Map, Value, from_str, json, to_string};
18use tokio::time::timeout as tokio_timeout;
19use uuid::Uuid;
20
21use crate::output;
22
23#[derive(Debug, Args)]
25pub struct RunArgs {
26 #[command(subcommand)]
28 pub command: RunCommands,
29}
30
31#[derive(Debug, Subcommand)]
33pub enum RunCommands {
34 Create {
36 workflow: String,
38 #[arg(long, group = "payload_source")]
40 payload: Option<String>,
41 #[arg(long, group = "payload_source")]
43 payload_file: Option<PathBuf>,
44 #[arg(long)]
47 max_retries: Option<u32>,
48 #[arg(long)]
54 idempotency_key: Option<String>,
55 #[arg(long = "max-cost", value_name = "USD")]
58 max_cost: Option<f64>,
59 #[arg(long)]
63 concurrency_key: Option<String>,
64 },
65 List {
67 #[arg(long)]
69 status: Option<String>,
70 #[arg(long)]
72 workflow: Option<String>,
73 #[arg(long)]
77 created_by: Option<Uuid>,
78 #[arg(long)]
80 page: Option<u32>,
81 #[arg(long)]
83 per_page: Option<u32>,
84 },
85 Get {
87 id: Uuid,
89 },
90 Cancel {
92 id: Uuid,
94 },
95 Approve {
97 id: Uuid,
99 },
100 Reject {
102 id: Uuid,
104 },
105 Input {
107 id: Uuid,
109 step_id: Uuid,
111 #[arg(long, group = "value_source")]
113 value: Option<String>,
114 #[arg(long, group = "value_source")]
116 value_file: Option<PathBuf>,
117 },
118 RejectInput {
120 id: Uuid,
122 step_id: Uuid,
124 #[arg(long)]
126 reason: Option<String>,
127 },
128 Retry {
130 id: Uuid,
132 #[arg(long)]
135 force: bool,
136 },
137 Replay {
139 id: Uuid,
141 },
142 Watch {
144 id: Uuid,
146 #[arg(long)]
148 no_logs: bool,
149 #[arg(long, value_parser = parse_humantime)]
151 timeout: Option<Duration>,
152 },
153 Plan {
155 workflow: String,
157 #[arg(long, group = "plan_input_source")]
159 input: Option<String>,
160 #[arg(long, group = "plan_input_source")]
162 input_file: Option<PathBuf>,
163 #[arg(long)]
165 max_depth: Option<u32>,
166 #[arg(long)]
168 no_estimates: bool,
169 },
170 Diff {
172 run_a: Uuid,
174 run_b: Uuid,
176 },
177}
178
179fn parse_humantime(s: &str) -> Result<Duration, String> {
181 humantime::parse_duration(s).map_err(|e| e.to_string())
182}
183
184const TERMINAL_EVENTS: &[&str] = &["run_completed", "run_failed", "run_cancelled"];
186
187fn resolve_payload(payload: Option<&str>, payload_file: Option<&PathBuf>) -> Result<Value> {
189 match (payload, payload_file) {
190 (Some(raw), _) => from_str(raw).context("invalid JSON in --payload"),
191 (_, Some(path)) => {
192 let content = fs::read_to_string(path)
193 .with_context(|| format!("cannot read payload file: {}", path.display()))?;
194 from_str(&content).with_context(|| format!("invalid JSON in {}", path.display()))
195 }
196 (None, None) => Ok(Value::Object(Map::new())),
197 }
198}
199
200fn validate_max_cost(max_cost: Option<f64>) -> Result<()> {
209 match max_cost {
210 Some(value) if !value.is_finite() => {
211 anyhow::bail!("--max-cost must be a finite number, got {value}")
212 }
213 Some(value) if value < 0.0 => {
214 anyhow::bail!("--max-cost must be zero or positive, got {value}")
215 }
216 _ => Ok(()),
217 }
218}
219
220pub async fn execute(
226 client: &IronflowClient,
227 args: &RunArgs,
228 json_mode: bool,
229 _verbose: bool,
230) -> Result<()> {
231 match &args.command {
232 RunCommands::Create {
233 workflow,
234 payload,
235 payload_file,
236 max_retries,
237 idempotency_key,
238 max_cost,
239 concurrency_key,
240 } => {
241 validate_max_cost(*max_cost)?;
242 let payload_value = resolve_payload(payload.as_deref(), payload_file.as_ref())?;
243 let payload_map = payload_value
244 .as_object()
245 .context("payload must be a JSON object")?
246 .clone();
247 let request: CreateRunRequest = CreateRunRequest::builder()
248 .workflow(workflow.clone())
249 .payload(Some(payload_map))
250 .max_retries(max_retries.map(|n| n as i32))
253 .max_cost_usd(*max_cost)
254 .concurrency_key(concurrency_key.clone())
255 .try_into()
256 .context("failed to build CreateRunRequest")?;
257
258 let response = match idempotency_key {
259 Some(key) => client.create_run_idempotent(&request, key).await?,
260 None => client.create_run(&request).await?,
261 };
262 output::print_output(json_mode, &response, || {
263 output::runs_table(slice::from_ref(&response.data))
264 })?;
265 }
266 RunCommands::List {
267 status,
268 workflow,
269 created_by,
270 page,
271 per_page,
272 } => {
273 let filter = ListRunsFilter {
274 status: status.as_deref(),
275 workflow: workflow.as_deref(),
276 created_by: *created_by,
277 page: *page,
278 per_page: *per_page,
279 ..Default::default()
280 };
281 let response = client.list_runs_filtered(&filter).await?;
282 output::print_output(json_mode, &response, || output::runs_table(&response.data))?;
283 }
284 RunCommands::Get { id } => {
285 let response = client.get_run(*id).await?;
286 output::print_output(json_mode, &response, || {
287 output::run_detail_table(&response.data)
288 })?;
289
290 if !json_mode && !response.data.steps.is_empty() {
291 let mut out = stdout().lock();
292 writeln!(out)?;
293 writeln!(out, "Steps:")?;
294 writeln!(out, "{}", output::steps_table(&response.data.steps))?;
295 }
296 }
297 RunCommands::Cancel { id } => {
298 let response = client.cancel_run(*id).await?;
299 output::print_output(json_mode, &response, || {
300 output::runs_table(slice::from_ref(&response.data))
301 })?;
302 }
303 RunCommands::Approve { id } => {
304 let response = client.approve_run(*id).await?;
305 output::print_output(json_mode, &response, || {
306 output::runs_table(slice::from_ref(&response.data))
307 })?;
308 if !json_mode && matches!(response.data.status, RunStatus::AwaitingApproval) {
310 println!("Approval recorded; more approvals are required.");
311 }
312 }
313 RunCommands::Reject { id } => {
314 let response = client.reject_run(*id).await?;
315 output::print_output(json_mode, &response, || {
316 output::runs_table(slice::from_ref(&response.data))
317 })?;
318 }
319 RunCommands::Input {
320 id,
321 step_id,
322 value,
323 value_file,
324 } => {
325 let answer = resolve_payload(value.as_deref(), value_file.as_ref())?;
326 let response = client.submit_input(*id, *step_id, &answer).await?;
327 output::print_output(json_mode, &response, || {
328 output::runs_table(slice::from_ref(&response.data))
329 })?;
330 }
331 RunCommands::RejectInput {
332 id,
333 step_id,
334 reason,
335 } => {
336 let response = client
337 .reject_input(*id, *step_id, reason.as_deref())
338 .await?;
339 output::print_output(json_mode, &response, || {
340 output::runs_table(slice::from_ref(&response.data))
341 })?;
342 }
343 RunCommands::Retry { id, force } => {
344 let response = client.retry_run(*id, *force).await?;
345 output::print_output(json_mode, &response, || {
346 output::runs_table(slice::from_ref(&response.data))
347 })?;
348 }
349 RunCommands::Replay { id } => {
350 let response = client.replay_run(*id).await?;
351 output::print_output(json_mode, &response, || {
352 output::runs_table(slice::from_ref(&response.data))
353 })?;
354 }
355 RunCommands::Watch {
356 id,
357 no_logs,
358 timeout,
359 } => {
360 execute_watch(client, *id, *no_logs, *timeout, json_mode).await?;
361 }
362 RunCommands::Plan {
363 workflow,
364 input,
365 input_file,
366 max_depth,
367 no_estimates,
368 } => {
369 let payload = resolve_payload(input.as_deref(), input_file.as_ref())?;
370 let payload_map = payload
371 .as_object()
372 .context("input must be a JSON object")?
373 .clone();
374 let request: PlanWorkflowRequest = PlanWorkflowRequest::builder()
375 .payload(Some(payload_map))
376 .max_depth(max_depth.map(|d| d as i32))
379 .estimate_durations(Some(!*no_estimates))
380 .try_into()
381 .context("failed to build PlanWorkflowRequest")?;
382 let response = client.plan_workflow(workflow, &request).await?;
383 output::render_execution_plan(&mut stdout().lock(), json_mode, &response)?;
384 }
385 RunCommands::Diff { run_a, run_b } => {
386 execute_diff(client, *run_a, *run_b, json_mode).await?;
387 }
388 }
389 Ok(())
390}
391
392async fn execute_watch(
394 client: &IronflowClient,
395 run_id: Uuid,
396 no_logs: bool,
397 timeout: Option<Duration>,
398 json_mode: bool,
399) -> Result<()> {
400 let run = client.get_run(run_id).await?;
401 let status = run.data.run.status;
402 if matches!(
403 status,
404 RunStatus::Completed | RunStatus::Failed | RunStatus::Cancelled
405 ) {
406 if json_mode {
407 output::print_output(json_mode, &run, || output::run_detail_table(&run.data))?;
408 } else {
409 let mut out = stdout().lock();
410 writeln!(out, "Run {run_id} already in terminal state: {status}")?;
411 }
412 return Ok(());
413 }
414
415 let watch_fut = async {
416 let mut stream = client.events(Some(run_id), None).await?;
417 let mut out = stdout().lock();
418
419 while let Some(event) = stream.next().await {
420 match event {
421 Ok(ev) => {
422 if no_logs
423 && !ev.event_type.starts_with("run_")
424 && !ev.event_type.starts_with("step_")
425 {
426 continue;
427 }
428
429 if json_mode {
430 let obj = json!({
431 "event": ev.event_type,
432 "data": ev.data,
433 });
434 writeln!(out, "{}", to_string(&obj)?)?;
435 } else {
436 writeln!(out, "[{}] {}", ev.event_type, ev.data)?;
437 }
438
439 if TERMINAL_EVENTS.contains(&ev.event_type.as_str()) {
440 break;
441 }
442 }
443 Err(e) => {
444 return Err(anyhow!("SSE stream error: {e}"));
445 }
446 }
447 }
448
449 Ok::<(), anyhow::Error>(())
450 };
451
452 match timeout {
453 Some(dur) => {
454 tokio_timeout(dur, watch_fut).await.unwrap_or_else(|_| {
455 eprintln!("Timeout reached after {}", format_duration(dur));
456 Ok(())
457 })?;
458 }
459 None => {
460 watch_fut.await?;
461 }
462 }
463
464 Ok(())
465}
466
467async fn execute_diff(
469 client: &IronflowClient,
470 run_a_id: Uuid,
471 run_b_id: Uuid,
472 json_mode: bool,
473) -> Result<()> {
474 if run_a_id == run_b_id {
475 bail!("both run IDs are the same; nothing to diff");
476 }
477
478 let (a, b) = tokio::try_join!(client.get_run(run_a_id), client.get_run(run_b_id))?;
479
480 if a.data.run.workflow_name != b.data.run.workflow_name {
481 bail!(
482 "cannot diff runs from different workflows: '{}' vs '{}'",
483 a.data.run.workflow_name,
484 b.data.run.workflow_name
485 );
486 }
487
488 if json_mode {
489 let diff = json!({
490 "run_a": a.data,
491 "run_b": b.data,
492 });
493 output::print_json(&diff)?;
494 } else {
495 let table = output::run_diff_table(&a.data, &b.data);
496 let mut out = stdout().lock();
497 writeln!(out, "{table}")?;
498 }
499
500 Ok(())
501}
502
503#[cfg(test)]
504mod tests {
505 use std::io::Write;
506
507 use tempfile::NamedTempFile;
508
509 use super::*;
510
511 #[test]
512 fn resolve_payload_none_returns_empty_object() {
513 let value = resolve_payload(None, None).unwrap();
514 assert!(value.is_object());
515 assert!(value.as_object().unwrap().is_empty());
516 }
517
518 #[test]
519 fn resolve_payload_inline_valid_json() {
520 let value = resolve_payload(Some(r#"{"key": "value"}"#), None).unwrap();
521 assert_eq!(value["key"], "value");
522 }
523
524 #[test]
525 fn resolve_payload_inline_invalid_json() {
526 let result = resolve_payload(Some("not json"), None);
527 assert!(result.is_err());
528 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
529 }
530
531 #[test]
532 fn resolve_payload_file_valid() {
533 let mut tmp = NamedTempFile::new().unwrap();
534 write!(tmp, r#"{{"workflow": "test"}}"#).unwrap();
535 let path = tmp.path().to_path_buf();
536
537 let value = resolve_payload(None, Some(&path)).unwrap();
538 assert_eq!(value["workflow"], "test");
539 }
540
541 #[test]
542 fn resolve_payload_file_not_found() {
543 let path = PathBuf::from("/nonexistent/payload.json");
544 let result = resolve_payload(None, Some(&path));
545 assert!(result.is_err());
546 assert!(result.unwrap_err().to_string().contains("cannot read"));
547 }
548
549 #[test]
550 fn resolve_payload_file_invalid_json() {
551 let mut tmp = NamedTempFile::new().unwrap();
552 write!(tmp, "not valid json").unwrap();
553 let path = tmp.path().to_path_buf();
554
555 let result = resolve_payload(None, Some(&path));
556 assert!(result.is_err());
557 assert!(result.unwrap_err().to_string().contains("invalid JSON"));
558 }
559}