Skip to main content

ironflow_engine/context/
artifacts.rs

1//! Artifact plumbing for [`WorkflowContext`].
2//!
3//! Covers the explicit [`put_artifact`](WorkflowContext::put_artifact) and
4//! [`get_artifact`](WorkflowContext::get_artifact) API plus the declarative
5//! input/output handling the step lifecycle calls around every shell step.
6
7use std::sync::Arc;
8
9use futures_util::StreamExt;
10use tracing::warn;
11use uuid::Uuid;
12
13use ironflow_artifacts::name::guess_content_type;
14use ironflow_artifacts::stream_from_bytes;
15use ironflow_store::entities::Artifact;
16use ironflow_store::models::ArtifactLookup;
17
18use crate::artifact::{
19    ArtifactSink, ArtifactUpload, StepLocation, collect_outputs, materialize_inputs,
20};
21use crate::config::StepConfig;
22use crate::error::EngineError;
23
24use super::WorkflowContext;
25
26impl WorkflowContext {
27    /// The artifact backend, or an explicit error when none is configured.
28    fn artifact_sink(&self) -> Result<&Arc<dyn ArtifactSink>, EngineError> {
29        self.artifact_sink.as_ref().ok_or_else(|| {
30            EngineError::ArtifactsUnavailable(
31                "no artifact storage is attached to this run".to_string(),
32            )
33        })
34    }
35
36    /// Store an in-memory payload as an artifact of the given step.
37    ///
38    /// The declarative [`ShellConfig::output`](crate::config::ShellConfig::output)
39    /// covers shell steps; this covers custom operations and agent steps, which
40    /// have no working directory to collect from.
41    ///
42    /// The MIME type is guessed from `name` unless `content_type` is set.
43    ///
44    /// # Errors
45    ///
46    /// Returns [`EngineError::ArtifactsUnavailable`] when no backend is
47    /// attached, [`EngineError::Artifact`] when the name is invalid or storage
48    /// fails, and [`EngineError::Store`] when the step already owns that name.
49    ///
50    /// # Examples
51    ///
52    /// ```no_run
53    /// use ironflow_engine::context::WorkflowContext;
54    /// use ironflow_engine::error::EngineError;
55    /// use uuid::Uuid;
56    ///
57    /// # async fn example(ctx: &WorkflowContext, step_id: Uuid) -> Result<(), EngineError> {
58    /// let artifact = ctx
59    ///     .put_artifact(step_id, "summary.json", None, br#"{"ok":true}"#.to_vec())
60    ///     .await?;
61    /// assert_eq!(artifact.content_type, "application/json");
62    /// # Ok(())
63    /// # }
64    /// ```
65    pub async fn put_artifact(
66        &self,
67        step_id: Uuid,
68        name: &str,
69        content_type: Option<&str>,
70        content: Vec<u8>,
71    ) -> Result<Artifact, EngineError> {
72        let sink = self.artifact_sink()?;
73        sink.put(
74            ArtifactUpload {
75                run_id: self.run_id,
76                step_id,
77                name: name.to_string(),
78                content_type: content_type
79                    .map(str::to_string)
80                    .unwrap_or_else(|| guess_content_type(name)),
81            },
82            stream_from_bytes(content),
83        )
84        .await
85    }
86
87    /// Read back an artifact produced earlier in this run.
88    ///
89    /// Resolution follows the same rule as a declared input: same run and
90    /// attempt, steps positioned strictly before the current one, closest
91    /// producer wins.
92    ///
93    /// # Errors
94    ///
95    /// Returns [`EngineError::ArtifactNotFound`] when nothing matches,
96    /// [`EngineError::ArtifactsUnavailable`] when no backend is attached, and
97    /// [`EngineError::Artifact`] when the bytes cannot be read.
98    ///
99    /// # Examples
100    ///
101    /// ```no_run
102    /// use ironflow_engine::context::WorkflowContext;
103    /// use ironflow_engine::error::EngineError;
104    ///
105    /// # async fn example(ctx: &WorkflowContext) -> Result<(), EngineError> {
106    /// let bytes = ctx.get_artifact("build", "report.html").await?;
107    /// println!("{} bytes", bytes.len());
108    /// # Ok(())
109    /// # }
110    /// ```
111    pub async fn get_artifact(&self, step: &str, name: &str) -> Result<Vec<u8>, EngineError> {
112        let sink = self.artifact_sink()?;
113
114        let artifact = self
115            .store
116            .find_artifact_for_input(ArtifactLookup {
117                run_id: self.run_id,
118                attempt: self.attempt,
119                before_position: self.position,
120                step_name: step.to_string(),
121                name: name.to_string(),
122            })
123            .await?
124            .ok_or_else(|| EngineError::ArtifactNotFound {
125                step: step.to_string(),
126                name: name.to_string(),
127            })?;
128
129        let mut content = sink.get(&artifact).await?;
130        let mut buffer = Vec::with_capacity(artifact.size_bytes as usize);
131        while let Some(chunk) = content.next().await {
132            let chunk = chunk?;
133            buffer.extend_from_slice(chunk.as_ref());
134        }
135
136        Ok(buffer)
137    }
138
139    /// Place a shell step's declared inputs in its working directory.
140    ///
141    /// A step that declares none needs no backend, so the check for one only
142    /// happens when there is something to materialize.
143    pub(super) async fn prepare_step_inputs(
144        &self,
145        config: &StepConfig,
146        position: u32,
147    ) -> Result<(), EngineError> {
148        let StepConfig::Shell(shell) = config else {
149            return Ok(());
150        };
151        if shell.inputs.is_empty() {
152            return Ok(());
153        }
154
155        materialize_inputs(
156            self.artifact_sink()?,
157            &self.store,
158            shell,
159            StepLocation {
160                run_id: self.run_id,
161                attempt: self.attempt,
162                position,
163            },
164        )
165        .await
166    }
167
168    /// Store a shell step's declared outputs.
169    ///
170    /// On a failed step this is best-effort: the collection error is logged and
171    /// swallowed so it never masks the failure that actually stopped the step.
172    pub(super) async fn store_step_outputs(
173        &self,
174        config: &StepConfig,
175        step_id: Uuid,
176        step_name: &str,
177        step_succeeded: bool,
178    ) -> Result<(), EngineError> {
179        let StepConfig::Shell(shell) = config else {
180            return Ok(());
181        };
182        if shell.outputs.is_empty() {
183            return Ok(());
184        }
185
186        let sink = match self.artifact_sink() {
187            Ok(sink) => sink,
188            Err(err) if step_succeeded => return Err(err),
189            Err(err) => {
190                warn!(
191                    run_id = %self.run_id,
192                    step = %step_name,
193                    error = %err,
194                    "cannot collect outputs of a failed step"
195                );
196                return Ok(());
197            }
198        };
199
200        let collected =
201            collect_outputs(sink, shell, self.run_id, step_id, step_name, step_succeeded).await;
202
203        match collected {
204            Ok(()) => Ok(()),
205            Err(err) if step_succeeded => Err(err),
206            Err(err) => {
207                warn!(
208                    run_id = %self.run_id,
209                    step = %step_name,
210                    error = %err,
211                    "failed to collect outputs of a failed step"
212                );
213                Ok(())
214            }
215        }
216    }
217}