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}