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}