Skip to main content

ironflow_engine/
plan.rs

1//! Execution plans -- what a run *would* do, without doing it.
2//!
3//! Ironflow workflows are Rust-native handlers, not declarative graphs: the
4//! only way to know which steps a run would create is to execute the handler
5//! with every step method short-circuited. That is what *plan mode* is.
6//!
7//! A recorder is attached to a
8//! [`WorkflowContext`](crate::context::WorkflowContext) by
9//! [`Engine::plan_handler`](crate::engine::Engine::plan_handler). Every step
10//! entry point checks it first, records a [`PlannedStep`], and returns a
11//! synthetic [`StepOutput`] without touching the store, the provider, the
12//! event bus or the network.
13//!
14//! # The success-shaped output assumption
15//!
16//! Synthetic outputs are shaped like a *successful* step (`exit_code: 0`,
17//! `status: 200`), so a native `if build.is_success()` branch in the handler
18//! follows the happy path. A plan therefore shows the nominal branch, not
19//! every branch the run might take. Conditions declared with
20//! [`WorkflowContext::when`](crate::context::WorkflowContext::when) are
21//! evaluated against the run input and reported; those declared with
22//! [`WorkflowContext::when_dynamic`](crate::context::WorkflowContext::when_dynamic)
23//! are reported as [`ConditionResult::Unevaluable`].
24//!
25//! # Why the global dry-run flag is not used
26//!
27//! [`ironflow_core::dry_run`] exposes a process-wide switch. Plan mode does not
28//! flip it: nothing is executed while planning, so the flag would buy nothing,
29//! and flipping a process-wide flag would corrupt real runs executing
30//! concurrently in the same process.
31
32use std::collections::HashMap;
33use std::mem::replace;
34use std::sync::{Arc, Mutex, MutexGuard};
35use std::time::Duration;
36
37use rust_decimal::Decimal;
38use serde::{Deserialize, Serialize};
39use serde_json::{Value, json};
40
41use ironflow_store::entities::{RunFilter, RunStatus, StepKind, StepStatus};
42use ironflow_store::store::Store;
43
44use crate::config::StepConfig;
45use crate::error::EngineError;
46use crate::executor::{StepArtifacts, StepOutput};
47
48/// How deep sub-workflows are expanded when the caller does not say.
49pub const DEFAULT_PLAN_MAX_DEPTH: u32 = 3;
50
51/// Hard cap on the number of steps a single plan may record.
52///
53/// A handler that loops forever would otherwise plan forever. Once the cap is
54/// reached the plan is returned truncated.
55pub const MAX_PLANNED_STEPS: usize = 1000;
56
57/// How many past runs are sampled to estimate step durations.
58pub const DEFAULT_ESTIMATE_SAMPLE_RUNS: u32 = 20;
59
60/// Outcome of a branch condition as seen by the planner.
61///
62/// # Examples
63///
64/// ```
65/// use serde_json::{Error, to_value};
66///
67/// use ironflow_engine::plan::ConditionResult;
68///
69/// # fn example() -> Result<(), Error> {
70/// let condition = ConditionResult::Evaluated {
71///     expression: "production run".to_string(),
72///     value: true,
73/// };
74/// assert_eq!(to_value(&condition)?["state"], "evaluated");
75/// # Ok(())
76/// # }
77/// ```
78#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
79#[serde(tag = "state", rename_all = "snake_case")]
80pub enum ConditionResult {
81    /// Resolved against the run input.
82    Evaluated {
83        /// Label the handler gave the branch. A name for the operator, never
84        /// parsed nor evaluated.
85        expression: String,
86        /// What the predicate returned for this input.
87        value: bool,
88    },
89    /// The step is explicitly skipped (recorded by
90    /// [`WorkflowContext::skip`](crate::context::WorkflowContext::skip)).
91    Skipped {
92        /// Reason the handler gave for skipping.
93        reason: String,
94    },
95    /// Depends on a previous step's output; unknown before the run.
96    Unevaluable {
97        /// Label the handler gave the branch.
98        expression: String,
99        /// Why the planner cannot resolve it.
100        reason: String,
101    },
102}
103
104/// One step the planner expects the run to create.
105///
106/// # Examples
107///
108/// ```
109/// use ironflow_engine::plan::PlannedStep;
110/// use ironflow_store::entities::StepKind;
111///
112/// let step = PlannedStep {
113///     name: "build".to_string(),
114///     kind: StepKind::Shell,
115///     workflow: "deploy".to_string(),
116///     depth: 0,
117///     depends_on: Vec::new(),
118///     condition: None,
119///     parallel_group: None,
120///     estimated_duration: None,
121/// };
122/// assert_eq!(step.name, "build");
123/// ```
124#[derive(Debug, Clone, Serialize, Deserialize)]
125pub struct PlannedStep {
126    /// Step name as the handler declares it.
127    pub name: String,
128    /// Kind of operation this step performs.
129    pub kind: StepKind,
130    /// Workflow that owns this step (top-level name, or the sub-workflow's).
131    pub workflow: String,
132    /// Sub-workflow nesting depth; `0` for the top-level workflow.
133    pub depth: u32,
134    /// Names of the steps this one runs after.
135    ///
136    /// Names, not identifiers: the same name can appear twice when two
137    /// branches or two sub-workflows declare a step with the same name.
138    /// Consumers must key on position, not on name.
139    pub depends_on: Vec<String>,
140    /// Branch condition recorded just before this step, if the handler
141    /// declared one.
142    ///
143    /// A condition pending before a parallel wave attaches to the *first*
144    /// member of the wave only.
145    pub condition: Option<ConditionResult>,
146    /// Parallel wave this step belongs to, when it runs concurrently with
147    /// its siblings.
148    pub parallel_group: Option<String>,
149    /// Average duration of this step across past completed runs.
150    #[serde(with = "opt_duration_ms")]
151    pub estimated_duration: Option<Duration>,
152}
153
154/// The full plan for one workflow and one input payload.
155///
156/// # Examples
157///
158/// ```
159/// use ironflow_engine::plan::ExecutionPlan;
160///
161/// let plan = ExecutionPlan {
162///     workflow: "deploy".to_string(),
163///     steps: Vec::new(),
164///     estimated_duration: None,
165///     max_depth: 3,
166///     truncated: false,
167///     incomplete_reason: None,
168/// };
169/// assert!(plan.steps.is_empty());
170/// ```
171#[derive(Debug, Clone, Serialize, Deserialize)]
172pub struct ExecutionPlan {
173    /// Workflow the plan was built for.
174    pub workflow: String,
175    /// Steps the run is expected to create, in execution order.
176    pub steps: Vec<PlannedStep>,
177    /// Sum of the step estimates, counting each parallel wave once.
178    #[serde(with = "opt_duration_ms")]
179    pub estimated_duration: Option<Duration>,
180    /// Sub-workflow expansion depth used for this plan.
181    pub max_depth: u32,
182    /// `true` when the step cap or the depth limit cut the plan short.
183    pub truncated: bool,
184    /// Why the plan stopped early (handler error, cap, depth limit).
185    pub incomplete_reason: Option<String>,
186}
187
188/// Knobs for [`Engine::plan_handler`](crate::engine::Engine::plan_handler).
189///
190/// # Examples
191///
192/// ```
193/// use ironflow_engine::plan::{PlanOptions, DEFAULT_PLAN_MAX_DEPTH};
194///
195/// let options = PlanOptions::default();
196/// assert_eq!(options.max_depth, DEFAULT_PLAN_MAX_DEPTH);
197/// ```
198#[derive(Debug, Clone)]
199pub struct PlanOptions {
200    /// How deep sub-workflows are expanded. Must be at least 1.
201    pub max_depth: u32,
202    /// Whether step durations are estimated from run history.
203    pub estimate_durations: bool,
204    /// How many past runs are sampled when estimating durations.
205    pub sample_runs: u32,
206}
207
208impl Default for PlanOptions {
209    fn default() -> Self {
210        Self {
211            max_depth: DEFAULT_PLAN_MAX_DEPTH,
212            estimate_durations: true,
213            sample_runs: DEFAULT_ESTIMATE_SAMPLE_RUNS,
214        }
215    }
216}
217
218/// Serde adapter mapping `Option<Duration>` to an optional millisecond count.
219mod opt_duration_ms {
220    use std::time::Duration;
221
222    use serde::{Deserialize, Deserializer, Serialize, Serializer};
223
224    /// Serialize a duration as whole milliseconds.
225    pub(super) fn serialize<S: Serializer>(
226        value: &Option<Duration>,
227        serializer: S,
228    ) -> Result<S::Ok, S::Error> {
229        let millis = value.map(|d| d.as_millis() as u64);
230        millis.serialize(serializer)
231    }
232
233    /// Deserialize whole milliseconds back into a duration.
234    pub(super) fn deserialize<'de, D: Deserializer<'de>>(
235        deserializer: D,
236    ) -> Result<Option<Duration>, D::Error> {
237        let millis = Option::<u64>::deserialize(deserializer)?;
238        Ok(millis.map(Duration::from_millis))
239    }
240}
241
242/// A [`PlanRecorder`] shared by a context and every child context it spawns.
243pub(crate) type SharedPlanRecorder = Arc<Mutex<PlanRecorder>>;
244
245/// Accumulates the steps a handler declares while running in plan mode.
246pub(crate) struct PlanRecorder {
247    payload: Value,
248    workflow: String,
249    steps: Vec<PlannedStep>,
250    last_names: Vec<String>,
251    pending_condition: Option<ConditionResult>,
252    depth: u32,
253    max_depth: u32,
254    parallel_groups: u32,
255    estimates: HashMap<String, Duration>,
256    truncated: bool,
257    incomplete_reason: Option<String>,
258}
259
260impl PlanRecorder {
261    /// Create a recorder for one workflow and one input payload.
262    pub(crate) fn new(
263        workflow: String,
264        payload: Value,
265        max_depth: u32,
266        estimates: HashMap<String, Duration>,
267    ) -> Self {
268        Self {
269            payload,
270            workflow,
271            steps: Vec::new(),
272            last_names: Vec::new(),
273            pending_condition: None,
274            depth: 0,
275            max_depth,
276            parallel_groups: 0,
277            estimates,
278            truncated: false,
279            incomplete_reason: None,
280        }
281    }
282
283    /// The payload the plan is being computed for.
284    pub(crate) fn payload(&self) -> Value {
285        self.payload.clone()
286    }
287
288    /// Swap the payload, returning the previous one.
289    ///
290    /// A sub-workflow plans against its own payload; the parent's is restored
291    /// when the expansion returns.
292    pub(crate) fn swap_payload(&mut self, next: Value) -> Value {
293        replace(&mut self.payload, next)
294    }
295
296    /// Attach a condition to the next recorded step.
297    pub(crate) fn set_condition(&mut self, condition: ConditionResult) {
298        self.pending_condition = Some(condition);
299    }
300
301    /// The historical estimate for a step name, if any.
302    pub(crate) fn estimate_for(&self, name: &str) -> Option<Duration> {
303        self.estimates.get(name).copied()
304    }
305
306    /// Provide a fallback estimate for a step the history does not cover.
307    ///
308    /// Used by steps whose duration is declared rather than observed, such as
309    /// a delay. History always wins when it exists.
310    pub(crate) fn seed_estimate(&mut self, name: &str, duration: Duration) {
311        self.estimates.entry(name.to_string()).or_insert(duration);
312    }
313
314    /// Record a step, returning `false` when the step cap is already reached.
315    ///
316    /// Does not update the dependency frontier: a parallel wave sets every
317    /// member at once, so the caller decides via [`set_last`](Self::set_last).
318    pub(crate) fn record(
319        &mut self,
320        name: &str,
321        kind: StepKind,
322        workflow: &str,
323        parallel_group: Option<String>,
324    ) -> bool {
325        if self.steps.len() >= MAX_PLANNED_STEPS {
326            self.truncated = true;
327            if self.incomplete_reason.is_none() {
328                self.incomplete_reason = Some(format!("step cap of {MAX_PLANNED_STEPS} reached"));
329            }
330            return false;
331        }
332
333        let estimated_duration = self.estimates.get(name).copied();
334        self.steps.push(PlannedStep {
335            name: name.to_string(),
336            kind,
337            workflow: workflow.to_string(),
338            depth: self.depth,
339            depends_on: self.last_names.clone(),
340            condition: self.pending_condition.take(),
341            parallel_group,
342            estimated_duration,
343        });
344        true
345    }
346
347    /// Set the dependency frontier the next recorded step depends on.
348    pub(crate) fn set_last(&mut self, names: Vec<String>) {
349        self.last_names = names;
350    }
351
352    /// Allocate a fresh parallel group name.
353    pub(crate) fn next_group(&mut self) -> String {
354        self.parallel_groups += 1;
355        format!("parallel-{}", self.parallel_groups)
356    }
357
358    /// Enter a sub-workflow, unless that would cross the depth limit.
359    pub(crate) fn enter_workflow(&mut self) -> bool {
360        if self.depth + 1 > self.max_depth {
361            self.truncated = true;
362            if self.incomplete_reason.is_none() {
363                self.incomplete_reason = Some(format!(
364                    "sub-workflow expansion stopped at depth {}",
365                    self.max_depth
366                ));
367            }
368            return false;
369        }
370        self.depth += 1;
371        true
372    }
373
374    /// Leave the current sub-workflow.
375    pub(crate) fn leave_workflow(&mut self) {
376        self.depth = self.depth.saturating_sub(1);
377    }
378
379    /// Record why the plan stopped early. The first writer wins.
380    pub(crate) fn fail(&mut self, reason: String) {
381        self.truncated = true;
382        if self.incomplete_reason.is_none() {
383            self.incomplete_reason = Some(reason);
384        }
385    }
386
387    /// Build the plan without consuming the recorder.
388    pub(crate) fn snapshot(&self) -> ExecutionPlan {
389        ExecutionPlan {
390            workflow: self.workflow.clone(),
391            estimated_duration: total_estimate(&self.steps),
392            steps: self.steps.clone(),
393            max_depth: self.max_depth,
394            truncated: self.truncated,
395            incomplete_reason: self.incomplete_reason.clone(),
396        }
397    }
398
399    /// Consume the recorder and build the plan.
400    pub(crate) fn into_plan(self) -> ExecutionPlan {
401        ExecutionPlan {
402            workflow: self.workflow,
403            estimated_duration: total_estimate(&self.steps),
404            steps: self.steps,
405            max_depth: self.max_depth,
406            truncated: self.truncated,
407            incomplete_reason: self.incomplete_reason,
408        }
409    }
410}
411
412/// Lock a shared recorder, ignoring poisoning.
413///
414/// A panicking handler must not turn every later plan into a panic of its own.
415/// Never hold the returned guard across an `.await`.
416pub(crate) fn lock_plan(plan: &SharedPlanRecorder) -> MutexGuard<'_, PlanRecorder> {
417    plan.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
418}
419
420/// Sum of sequential estimates; each parallel group counts once, at its
421/// slowest member.
422///
423/// Returns `None` when no step carries an estimate, so an absent history is
424/// reported as "unknown" rather than as zero.
425fn total_estimate(steps: &[PlannedStep]) -> Option<Duration> {
426    let mut total = Duration::ZERO;
427    let mut seen_any = false;
428    let mut current_group: Option<&str> = None;
429    let mut group_max = Duration::ZERO;
430
431    for step in steps {
432        match step.parallel_group.as_deref() {
433            Some(group) if current_group == Some(group) => {
434                if let Some(estimate) = step.estimated_duration {
435                    seen_any = true;
436                    group_max = group_max.max(estimate);
437                }
438            }
439            Some(group) => {
440                if current_group.is_some() {
441                    total += group_max;
442                }
443                current_group = Some(group);
444                group_max = step.estimated_duration.unwrap_or(Duration::ZERO);
445                if step.estimated_duration.is_some() {
446                    seen_any = true;
447                }
448            }
449            None => {
450                if current_group.is_some() {
451                    total += group_max;
452                    current_group = None;
453                    group_max = Duration::ZERO;
454                }
455                if let Some(estimate) = step.estimated_duration {
456                    seen_any = true;
457                    total += estimate;
458                }
459            }
460        }
461    }
462
463    if current_group.is_some() {
464        total += group_max;
465    }
466
467    seen_any.then_some(total)
468}
469
470/// Synthetic, success-shaped output returned to the handler in plan mode.
471///
472/// Shaped so that `output.is_success()` is `true` for shell and HTTP steps:
473/// the plan follows the branch a successful run would take.
474pub(crate) fn planned_output(config: &StepConfig, estimate: Option<Duration>) -> StepOutput {
475    let output = match config {
476        StepConfig::Shell(_) => json!({"stdout": "", "stderr": "", "exit_code": 0}),
477        StepConfig::Http(_) => json!({"status": 200, "headers": {}, "body": ""}),
478        StepConfig::Agent(_) => json!({}),
479        StepConfig::Workflow(c) => json!({
480            "run_id": Value::Null,
481            "workflow_name": c.workflow_name,
482            "status": "completed",
483            "cost_usd": 0,
484            "duration_ms": 0,
485        }),
486        StepConfig::Approval(_) | StepConfig::Decision(_) | StepConfig::Delay(_) => Value::Null,
487    };
488
489    StepOutput {
490        output,
491        duration_ms: estimate.map(|d| d.as_millis() as u64).unwrap_or(0),
492        cost_usd: Decimal::ZERO,
493        input_tokens: None,
494        cache_read_input_tokens: None,
495        cache_creation_input_tokens: None,
496        output_tokens: None,
497        model: None,
498        debug_messages: None,
499        artifacts: StepArtifacts::default(),
500        account_id: None,
501        environment_id: None,
502    }
503}
504
505/// Synthetic output for a custom operation step in plan mode.
506pub(crate) fn planned_custom_output(estimate: Option<Duration>) -> StepOutput {
507    StepOutput {
508        output: json!({}),
509        duration_ms: estimate.map(|d| d.as_millis() as u64).unwrap_or(0),
510        cost_usd: Decimal::ZERO,
511        input_tokens: None,
512        cache_read_input_tokens: None,
513        cache_creation_input_tokens: None,
514        output_tokens: None,
515        model: None,
516        debug_messages: None,
517        artifacts: StepArtifacts::default(),
518        account_id: None,
519        environment_id: None,
520    }
521}
522
523/// Average completed-step durations for a workflow, keyed by step name.
524///
525/// Samples the most recent completed runs of `workflow_name` and averages the
526/// duration of every completed step, per name. Step names absent from the
527/// history are simply absent from the map.
528///
529/// # Errors
530///
531/// Returns [`EngineError::Store`] when the history query fails.
532///
533/// # Examples
534///
535/// ```no_run
536/// use std::sync::Arc;
537///
538/// use ironflow_engine::error::EngineError;
539/// use ironflow_engine::plan::estimate_durations;
540/// use ironflow_store::store::Store;
541///
542/// # async fn example(store: &Arc<dyn Store>) -> Result<(), EngineError> {
543/// let estimates = estimate_durations(store, "deploy", 20).await?;
544/// println!("{} steps have a history", estimates.len());
545/// # Ok(())
546/// # }
547/// ```
548pub async fn estimate_durations(
549    store: &Arc<dyn Store>,
550    workflow_name: &str,
551    sample_runs: u32,
552) -> Result<HashMap<String, Duration>, EngineError> {
553    let page = store
554        .list_runs(
555            RunFilter {
556                workflow_name: Some(workflow_name.to_string()),
557                status: Some(RunStatus::Completed),
558                has_steps: Some(true),
559                ..RunFilter::default()
560            },
561            1,
562            sample_runs.clamp(1, 100),
563        )
564        .await?;
565
566    let mut totals: HashMap<String, (u64, u64)> = HashMap::new();
567    for run in &page.items {
568        for step in store.list_steps(run.id).await? {
569            if step.status.state != StepStatus::Completed {
570                continue;
571            }
572            let entry = totals.entry(step.name).or_insert((0, 0));
573            entry.0 += step.duration_ms;
574            entry.1 += 1;
575        }
576    }
577
578    Ok(totals
579        .into_iter()
580        .filter(|(_, (_, count))| *count > 0)
581        .map(|(name, (sum, count))| (name, Duration::from_millis(sum / count)))
582        .collect())
583}
584
585impl ExecutionPlan {
586    /// Total estimate in whole milliseconds, when the plan has one.
587    ///
588    /// # Examples
589    ///
590    /// ```
591    /// use std::time::Duration;
592    ///
593    /// use ironflow_engine::plan::ExecutionPlan;
594    ///
595    /// let plan = ExecutionPlan {
596    ///     workflow: "deploy".to_string(),
597    ///     steps: Vec::new(),
598    ///     estimated_duration: Some(Duration::from_millis(1500)),
599    ///     max_depth: 3,
600    ///     truncated: false,
601    ///     incomplete_reason: None,
602    /// };
603    /// assert_eq!(plan.estimated_duration_ms(), Some(1500));
604    /// ```
605    pub fn estimated_duration_ms(&self) -> Option<u64> {
606        self.estimated_duration.map(|d| d.as_millis() as u64)
607    }
608}
609
610impl PlannedStep {
611    /// This step's estimate in whole milliseconds, when it has one.
612    ///
613    /// # Examples
614    ///
615    /// ```
616    /// use std::time::Duration;
617    ///
618    /// use ironflow_engine::plan::PlannedStep;
619    /// use ironflow_store::entities::StepKind;
620    ///
621    /// let step = PlannedStep {
622    ///     name: "build".to_string(),
623    ///     kind: StepKind::Shell,
624    ///     workflow: "deploy".to_string(),
625    ///     depth: 0,
626    ///     depends_on: Vec::new(),
627    ///     condition: None,
628    ///     parallel_group: None,
629    ///     estimated_duration: Some(Duration::from_millis(250)),
630    /// };
631    /// assert_eq!(step.estimated_duration_ms(), Some(250));
632    /// ```
633    pub fn estimated_duration_ms(&self) -> Option<u64> {
634        self.estimated_duration.map(|d| d.as_millis() as u64)
635    }
636}
637
638#[cfg(test)]
639mod tests {
640    use std::collections::HashMap;
641    use std::thread::spawn;
642    use std::time::Duration;
643
644    use serde_json::{from_value, json, to_value};
645
646    use crate::config::{HttpConfig, ShellConfig};
647
648    use super::*;
649
650    fn step(name: &str, group: Option<&str>, estimate: Option<u64>) -> PlannedStep {
651        PlannedStep {
652            name: name.to_string(),
653            kind: StepKind::Shell,
654            workflow: "wf".to_string(),
655            depth: 0,
656            depends_on: Vec::new(),
657            condition: None,
658            parallel_group: group.map(str::to_string),
659            estimated_duration: estimate.map(Duration::from_millis),
660        }
661    }
662
663    #[test]
664    fn total_estimate_sums_sequential_steps() {
665        let steps = vec![
666            step("a", None, Some(100)),
667            step("b", None, Some(250)),
668            step("c", None, Some(50)),
669        ];
670        assert_eq!(total_estimate(&steps), Some(Duration::from_millis(400)));
671    }
672
673    #[test]
674    fn total_estimate_counts_a_parallel_group_once_at_its_slowest() {
675        let steps = vec![
676            step("build", None, Some(100)),
677            step("t1", Some("parallel-1"), Some(300)),
678            step("t2", Some("parallel-1"), Some(700)),
679            step("t3", Some("parallel-1"), Some(200)),
680            step("deploy", None, Some(100)),
681        ];
682        assert_eq!(total_estimate(&steps), Some(Duration::from_millis(900)));
683    }
684
685    #[test]
686    fn total_estimate_handles_a_trailing_parallel_group() {
687        let steps = vec![
688            step("build", None, Some(100)),
689            step("t1", Some("parallel-1"), Some(300)),
690            step("t2", Some("parallel-1"), Some(700)),
691        ];
692        assert_eq!(total_estimate(&steps), Some(Duration::from_millis(800)));
693    }
694
695    #[test]
696    fn total_estimate_is_none_without_any_estimate() {
697        let steps = vec![step("a", None, None), step("b", None, None)];
698        assert_eq!(total_estimate(&steps), None);
699    }
700
701    #[test]
702    fn planned_step_duration_round_trips_as_milliseconds() {
703        let original = step("a", None, Some(1234));
704        let value = to_value(&original).expect("serialize");
705        assert_eq!(value["estimated_duration"], 1234);
706
707        let back: PlannedStep = from_value(value).expect("deserialize");
708        assert_eq!(back.estimated_duration, Some(Duration::from_millis(1234)));
709    }
710
711    #[test]
712    fn planned_step_duration_round_trips_when_absent() {
713        let original = step("a", None, None);
714        let value = to_value(&original).expect("serialize");
715        assert!(value["estimated_duration"].is_null());
716
717        let back: PlannedStep = from_value(value).expect("deserialize");
718        assert_eq!(back.estimated_duration, None);
719    }
720
721    #[test]
722    fn planned_shell_and_http_outputs_look_successful() {
723        let shell = planned_output(&StepConfig::Shell(ShellConfig::new("echo hi")), None);
724        assert!(shell.is_success());
725
726        let http = planned_output(
727            &StepConfig::Http(HttpConfig::get("https://example.com")),
728            None,
729        );
730        assert!(http.is_success());
731    }
732
733    #[test]
734    fn planned_output_carries_the_estimate_as_its_duration() {
735        let output = planned_output(
736            &StepConfig::Shell(ShellConfig::new("echo hi")),
737            Some(Duration::from_millis(900)),
738        );
739        assert_eq!(output.duration_ms, 900);
740        assert_eq!(output.cost_usd, Decimal::ZERO);
741    }
742
743    #[test]
744    fn planned_custom_output_is_an_empty_object() {
745        let output = planned_custom_output(None);
746        assert_eq!(output.output, json!({}));
747        assert_eq!(output.duration_ms, 0);
748    }
749
750    #[test]
751    fn condition_result_serializes_its_state_tag() {
752        let evaluated = to_value(ConditionResult::Evaluated {
753            expression: "env == prod".to_string(),
754            value: false,
755        })
756        .expect("serialize");
757        assert_eq!(evaluated["state"], "evaluated");
758        assert_eq!(evaluated["value"], false);
759
760        let skipped = to_value(ConditionResult::Skipped {
761            reason: "not prod".to_string(),
762        })
763        .expect("serialize");
764        assert_eq!(skipped["state"], "skipped");
765
766        let unevaluable = to_value(ConditionResult::Unevaluable {
767            expression: "build succeeded".to_string(),
768            reason: "depends on a step output".to_string(),
769        })
770        .expect("serialize");
771        assert_eq!(unevaluable["state"], "unevaluable");
772    }
773
774    #[test]
775    fn recorder_records_dependencies_and_conditions() {
776        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 3, HashMap::new());
777        assert!(recorder.record("a", StepKind::Shell, "wf", None));
778        recorder.set_last(vec!["a".to_string()]);
779        recorder.set_condition(ConditionResult::Skipped {
780            reason: "nope".to_string(),
781        });
782        assert!(recorder.record("b", StepKind::Shell, "wf", None));
783
784        let plan = recorder.into_plan();
785        assert_eq!(plan.steps.len(), 2);
786        assert_eq!(plan.steps[1].depends_on, vec!["a".to_string()]);
787        assert!(matches!(
788            plan.steps[1].condition,
789            Some(ConditionResult::Skipped { .. })
790        ));
791        assert!(plan.steps[0].condition.is_none());
792    }
793
794    #[test]
795    fn recorder_stops_at_the_step_cap() {
796        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 3, HashMap::new());
797        for index in 0..MAX_PLANNED_STEPS {
798            assert!(recorder.record(&format!("s{index}"), StepKind::Shell, "wf", None));
799        }
800        assert!(!recorder.record("overflow", StepKind::Shell, "wf", None));
801
802        let plan = recorder.into_plan();
803        assert_eq!(plan.steps.len(), MAX_PLANNED_STEPS);
804        assert!(plan.truncated);
805        assert!(
806            plan.incomplete_reason
807                .expect("a reason")
808                .contains("step cap")
809        );
810    }
811
812    #[test]
813    fn recorder_refuses_to_expand_past_the_depth_limit() {
814        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 1, HashMap::new());
815        assert!(recorder.enter_workflow());
816        assert!(!recorder.enter_workflow());
817
818        let plan = recorder.snapshot();
819        assert!(plan.truncated);
820        assert!(
821            plan.incomplete_reason
822                .expect("a reason")
823                .contains("depth 1")
824        );
825    }
826
827    #[test]
828    fn recorder_allocates_successive_parallel_group_names() {
829        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 3, HashMap::new());
830        assert_eq!(recorder.next_group(), "parallel-1");
831        assert_eq!(recorder.next_group(), "parallel-2");
832    }
833
834    #[test]
835    fn recorder_swaps_and_restores_the_payload() {
836        let mut recorder =
837            PlanRecorder::new("wf".to_string(), json!({"env": "prod"}), 3, HashMap::new());
838        let previous = recorder.swap_payload(json!({"env": "dev"}));
839        assert_eq!(previous, json!({"env": "prod"}));
840        assert_eq!(recorder.payload(), json!({"env": "dev"}));
841        recorder.swap_payload(previous);
842        assert_eq!(recorder.payload(), json!({"env": "prod"}));
843    }
844
845    #[test]
846    fn recorder_keeps_the_first_failure_reason() {
847        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 3, HashMap::new());
848        recorder.fail("first".to_string());
849        recorder.fail("second".to_string());
850        let plan = recorder.into_plan();
851        assert_eq!(plan.incomplete_reason.as_deref(), Some("first"));
852        assert!(plan.truncated);
853    }
854
855    #[test]
856    fn recorder_uses_the_history_estimate_for_a_known_step() {
857        let estimates = HashMap::from([("build".to_string(), Duration::from_millis(400))]);
858        let mut recorder = PlanRecorder::new("wf".to_string(), json!({}), 3, estimates);
859        assert_eq!(
860            recorder.estimate_for("build"),
861            Some(Duration::from_millis(400))
862        );
863        assert_eq!(recorder.estimate_for("unknown"), None);
864        recorder.record("build", StepKind::Shell, "wf", None);
865        let plan = recorder.into_plan();
866        assert_eq!(plan.estimated_duration, Some(Duration::from_millis(400)));
867    }
868
869    #[test]
870    fn lock_plan_recovers_from_poisoning() {
871        let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
872            "wf".to_string(),
873            json!({}),
874            3,
875            HashMap::new(),
876        )));
877        let poisoner = Arc::clone(&shared);
878        let _ = spawn(move || {
879            let _guard = poisoner.lock().expect("lock");
880            panic!("poison the mutex");
881        })
882        .join();
883
884        assert_eq!(lock_plan(&shared).payload(), json!({}));
885    }
886}