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