1use std::future::Future;
17use std::path::PathBuf;
18use std::pin::Pin;
19use std::sync::Arc;
20
21use futures_util::StreamExt;
22use glob::glob;
23use tokio::fs::{File, create_dir_all};
24use tokio::io::AsyncWriteExt;
25use tracing::{info, warn};
26use uuid::Uuid;
27
28use ironflow_artifacts::blob_store::{BlobStore, ByteStream};
29use ironflow_artifacts::error::ArtifactError;
30use ironflow_artifacts::name::{guess_content_type, storage_key, validate_artifact_name};
31use ironflow_artifacts::stream_from_path;
32use ironflow_store::entities::{Artifact, ArtifactLookup, NewArtifact};
33use ironflow_store::store::Store;
34
35use crate::config::ShellConfig;
36use crate::error::EngineError;
37
38pub type ArtifactFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, EngineError>> + Send + 'a>>;
40
41#[derive(Debug, Clone)]
58pub struct ArtifactUpload {
59 pub run_id: Uuid,
61 pub step_id: Uuid,
63 pub name: String,
65 pub content_type: String,
67}
68
69pub trait ArtifactSink: Send + Sync {
99 fn put<'a>(
107 &'a self,
108 upload: ArtifactUpload,
109 content: ByteStream,
110 ) -> ArtifactFuture<'a, Artifact>;
111
112 fn get<'a>(&'a self, artifact: &'a Artifact) -> ArtifactFuture<'a, ByteStream>;
118}
119
120pub struct DirectArtifactSink {
141 blob: Arc<dyn BlobStore>,
142 store: Arc<dyn Store>,
143}
144
145impl DirectArtifactSink {
146 pub fn new(blob: Arc<dyn BlobStore>, store: Arc<dyn Store>) -> Self {
148 Self { blob, store }
149 }
150}
151
152impl std::fmt::Debug for DirectArtifactSink {
153 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
154 f.debug_struct("DirectArtifactSink").finish_non_exhaustive()
155 }
156}
157
158impl ArtifactSink for DirectArtifactSink {
159 fn put<'a>(
160 &'a self,
161 upload: ArtifactUpload,
162 content: ByteStream,
163 ) -> ArtifactFuture<'a, Artifact> {
164 Box::pin(async move {
165 validate_artifact_name(&upload.name)?;
166
167 let id = Uuid::now_v7();
168 let key = storage_key(upload.run_id, upload.step_id, id);
169 let digest = self.blob.put(&key, content).await?;
170
171 let recorded = self
172 .store
173 .create_artifact(NewArtifact {
174 id,
175 run_id: upload.run_id,
176 step_id: upload.step_id,
177 name: upload.name,
178 storage_key: key.clone(),
179 content_type: upload.content_type,
180 size_bytes: digest.size_bytes,
181 sha256: digest.sha256,
182 })
183 .await;
184
185 match recorded {
186 Ok(artifact) => Ok(artifact),
187 Err(err) => {
188 if let Err(cleanup) = self.blob.delete(&key).await {
191 warn!(
192 storage_key = %key,
193 error = %cleanup,
194 "failed to remove the blob of an unrecorded artifact"
195 );
196 }
197 Err(err.into())
198 }
199 }
200 })
201 }
202
203 fn get<'a>(&'a self, artifact: &'a Artifact) -> ArtifactFuture<'a, ByteStream> {
204 Box::pin(async move { Ok(self.blob.get(&artifact.storage_key).await?) })
205 }
206}
207
208#[derive(Debug, Clone, Copy)]
210pub(crate) struct StepLocation {
211 pub(crate) run_id: Uuid,
213 pub(crate) attempt: u32,
215 pub(crate) position: u32,
217}
218
219pub(crate) async fn materialize_inputs(
223 sink: &Arc<dyn ArtifactSink>,
224 store: &Arc<dyn Store>,
225 config: &ShellConfig,
226 location: StepLocation,
227) -> Result<(), EngineError> {
228 let work_dir = working_dir(config);
229
230 for input in &config.inputs {
231 let artifact = store
232 .find_artifact_for_input(ArtifactLookup {
233 run_id: location.run_id,
234 attempt: location.attempt,
235 before_position: location.position,
236 step_name: input.step.clone(),
237 name: input.name.clone(),
238 })
239 .await?
240 .ok_or_else(|| EngineError::ArtifactNotFound {
241 step: input.step.clone(),
242 name: input.name.clone(),
243 })?;
244
245 let destination = work_dir.join(input.destination());
246 if let Some(parent) = destination.parent() {
247 create_dir_all(parent).await.map_err(ArtifactError::from)?;
248 }
249
250 let mut content = sink.get(&artifact).await?;
251 let mut file = File::create(&destination)
252 .await
253 .map_err(ArtifactError::from)?;
254 while let Some(chunk) = content.next().await {
255 file.write_all(&chunk?).await.map_err(ArtifactError::from)?;
256 }
257 file.flush().await.map_err(ArtifactError::from)?;
258
259 info!(
260 run_id = %location.run_id,
261 artifact = %input.name,
262 produced_by = %input.step,
263 destination = %destination.display(),
264 "artifact input materialized"
265 );
266 }
267
268 Ok(())
269}
270
271pub(crate) async fn collect_outputs(
278 sink: &Arc<dyn ArtifactSink>,
279 config: &ShellConfig,
280 run_id: Uuid,
281 step_id: Uuid,
282 step_name: &str,
283 step_succeeded: bool,
284) -> Result<(), EngineError> {
285 let work_dir = working_dir(config);
286
287 for output in &config.outputs {
288 let pattern = work_dir.join(&output.pattern);
289 let pattern = pattern.to_str().ok_or_else(|| {
290 EngineError::StepConfig(format!(
291 "output pattern {:?} is not valid UTF-8",
292 output.pattern
293 ))
294 })?;
295
296 let matches = glob(pattern)
297 .map_err(|err| {
298 EngineError::StepConfig(format!(
299 "invalid output pattern {:?}: {err}",
300 output.pattern
301 ))
302 })?
303 .filter_map(Result::ok)
304 .filter(|path| path.is_file())
305 .collect::<Vec<_>>();
306
307 if matches.is_empty() {
308 if step_succeeded {
309 return Err(EngineError::MissingArtifact {
310 step: step_name.to_string(),
311 pattern: output.pattern.clone(),
312 });
313 }
314 warn!(
315 run_id = %run_id,
316 step = %step_name,
317 pattern = %output.pattern,
318 "declared output matched no file on a failed step"
319 );
320 continue;
321 }
322
323 for path in matches {
324 let name = path
325 .file_name()
326 .and_then(|name| name.to_str())
327 .ok_or_else(|| {
328 EngineError::StepConfig(format!(
329 "output file {:?} has no valid UTF-8 name",
330 path.display()
331 ))
332 })?
333 .to_string();
334
335 let content_type = output
336 .content_type
337 .clone()
338 .unwrap_or_else(|| guess_content_type(&name));
339
340 let content = stream_from_path(&path).await?;
341 let artifact = sink
342 .put(
343 ArtifactUpload {
344 run_id,
345 step_id,
346 name: name.clone(),
347 content_type,
348 },
349 content,
350 )
351 .await?;
352
353 info!(
354 run_id = %run_id,
355 step = %step_name,
356 artifact = %artifact.name,
357 size_bytes = artifact.size_bytes,
358 "artifact output stored"
359 );
360 }
361 }
362
363 Ok(())
364}
365
366fn working_dir(config: &ShellConfig) -> PathBuf {
368 PathBuf::from(config.dir.as_deref().unwrap_or("."))
369}
370
371#[cfg(test)]
372mod tests {
373 use std::collections::HashMap;
374
375 use futures_util::TryStreamExt;
376 use ironflow_artifacts::local::LocalBlobStore;
377 use ironflow_artifacts::stream_from_bytes;
378 use ironflow_store::entities::{NewRun, NewStep, StepKind, TriggerKind, step_trace_id};
379 use ironflow_store::memory::InMemoryStore;
380 use serde_json::json;
381 use tempfile::TempDir;
382
383 use super::*;
384
385 async fn sink_with_step() -> (TempDir, DirectArtifactSink, Uuid, Uuid) {
386 let dir = TempDir::new().expect("temp dir");
387 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
388 let blob: Arc<dyn BlobStore> = Arc::new(LocalBlobStore::new(dir.path()));
389
390 let run = store
391 .create_run(NewRun {
392 workflow_name: "artifacts".to_string(),
393 trigger: TriggerKind::Manual,
394 payload: json!({}),
395 max_retries: 0,
396 handler_version: None,
397 labels: HashMap::new(),
398 scheduled_at: None,
399 created_by: None,
400 idempotency_key: None,
401 max_cost_usd: None,
402 })
403 .await
404 .expect("create run")
405 .into_run();
406
407 let step = store
408 .create_step(NewStep {
409 run_id: run.id,
410 trace_id: step_trace_id(run.id, "build", 0),
411 name: "build".to_string(),
412 kind: StepKind::Shell,
413 position: 0,
414 input: None,
415 is_error_handler: false,
416 })
417 .await
418 .expect("create step");
419
420 let sink = DirectArtifactSink::new(blob, store);
421 (dir, sink, run.id, step.id)
422 }
423
424 fn upload(run_id: Uuid, step_id: Uuid, name: &str) -> ArtifactUpload {
425 ArtifactUpload {
426 run_id,
427 step_id,
428 name: name.to_string(),
429 content_type: "text/plain".to_string(),
430 }
431 }
432
433 #[tokio::test]
434 async fn put_records_size_and_hash() {
435 let (_dir, sink, run_id, step_id) = sink_with_step().await;
436
437 let artifact = sink
438 .put(
439 upload(run_id, step_id, "report.txt"),
440 stream_from_bytes(b"abc".to_vec()),
441 )
442 .await
443 .expect("put");
444
445 assert_eq!(artifact.size_bytes, 3);
446 assert_eq!(
447 artifact.sha256,
448 "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad"
449 );
450 }
451
452 #[tokio::test]
453 async fn put_then_get_roundtrips_the_bytes() {
454 let (_dir, sink, run_id, step_id) = sink_with_step().await;
455
456 let artifact = sink
457 .put(
458 upload(run_id, step_id, "report.txt"),
459 stream_from_bytes(b"hello".to_vec()),
460 )
461 .await
462 .expect("put");
463
464 let chunks: Vec<bytes::Bytes> = sink
465 .get(&artifact)
466 .await
467 .expect("get")
468 .try_collect()
469 .await
470 .expect("collect");
471
472 assert_eq!(chunks.concat(), b"hello");
473 }
474
475 #[tokio::test]
476 async fn the_storage_key_never_embeds_the_name() {
477 let (_dir, sink, run_id, step_id) = sink_with_step().await;
478
479 let artifact = sink
480 .put(
481 upload(run_id, step_id, "report.txt"),
482 stream_from_bytes(b"x".to_vec()),
483 )
484 .await
485 .expect("put");
486
487 assert!(!artifact.storage_key.contains("report"));
488 assert!(artifact.storage_key.ends_with(&artifact.id.to_string()));
489 }
490
491 #[tokio::test]
492 async fn an_invalid_name_is_rejected_before_anything_is_written() {
493 let (dir, sink, run_id, step_id) = sink_with_step().await;
494
495 let err = sink
496 .put(
497 upload(run_id, step_id, "../escape"),
498 stream_from_bytes(b"x".to_vec()),
499 )
500 .await
501 .expect_err("invalid name");
502
503 assert!(matches!(err, EngineError::Artifact(_)));
504 assert!(!dir.path().join("artifacts").exists());
505 }
506
507 #[tokio::test]
508 async fn a_duplicate_name_fails_and_leaves_no_orphan_blob() {
509 let (dir, sink, run_id, step_id) = sink_with_step().await;
510
511 sink.put(
512 upload(run_id, step_id, "report.txt"),
513 stream_from_bytes(b"first".to_vec()),
514 )
515 .await
516 .expect("first");
517
518 let err = sink
519 .put(
520 upload(run_id, step_id, "report.txt"),
521 stream_from_bytes(b"second".to_vec()),
522 )
523 .await
524 .expect_err("duplicate");
525
526 assert!(matches!(err, EngineError::Store(_)));
527
528 let stored: Vec<_> = std::fs::read_dir(
529 dir.path()
530 .join("artifacts")
531 .join(run_id.to_string())
532 .join(step_id.to_string()),
533 )
534 .expect("read dir")
535 .filter_map(Result::ok)
536 .collect();
537 assert_eq!(stored.len(), 1, "the rejected blob was not cleaned up");
538 }
539}