Skip to main content

ironflow_engine/
artifact.rs

1//! Artifact production and consumption from a running workflow.
2//!
3//! [`ArtifactSink`] is the seam between the engine and wherever artifact bytes
4//! actually live. The engine never talks to a blob store directly, because the
5//! two execution topologies reach storage differently:
6//!
7//! - in-process (API server, tests): [`DirectArtifactSink`] writes to the blob
8//!   store and records the metadata itself;
9//! - remote worker: the worker's sink uploads over the internal HTTP API, which
10//!   keeps storage credentials on the API side only.
11//!
12//! Ordering is always "bytes first, metadata second". The metadata row is the
13//! source of truth, so a crash in between leaves an unreferenced blob rather
14//! than a record pointing at nothing.
15
16use 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
38/// Boxed future returned by [`ArtifactSink`] methods -- keeps the trait object safe.
39pub type ArtifactFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, EngineError>> + Send + 'a>>;
40
41/// Everything needed to record an artifact, minus the bytes.
42///
43/// # Examples
44///
45/// ```
46/// use ironflow_engine::artifact::ArtifactUpload;
47/// use uuid::Uuid;
48///
49/// let upload = ArtifactUpload {
50///     run_id: Uuid::now_v7(),
51///     step_id: Uuid::now_v7(),
52///     name: "report.html".to_string(),
53///     content_type: "text/html".to_string(),
54/// };
55/// assert_eq!(upload.name, "report.html");
56/// ```
57#[derive(Debug, Clone)]
58pub struct ArtifactUpload {
59    /// The run the artifact belongs to.
60    pub run_id: Uuid,
61    /// The step that produced it.
62    pub step_id: Uuid,
63    /// User-facing file name, unique within the step.
64    pub name: String,
65    /// MIME type to serve on download.
66    pub content_type: String,
67}
68
69/// Where a running workflow reads and writes artifact bytes.
70///
71/// # Examples
72///
73/// ```no_run
74/// use std::sync::Arc;
75///
76/// use ironflow_artifacts::stream_from_bytes;
77/// use ironflow_engine::artifact::{ArtifactSink, ArtifactUpload};
78/// use uuid::Uuid;
79///
80/// # async fn example(sink: Arc<dyn ArtifactSink>, run_id: Uuid, step_id: Uuid)
81/// # -> Result<(), ironflow_engine::error::EngineError> {
82/// let artifact = sink
83///     .put(
84///         ArtifactUpload {
85///             run_id,
86///             step_id,
87///             name: "report.json".to_string(),
88///             content_type: "application/json".to_string(),
89///         },
90///         stream_from_bytes(b"{}".to_vec()),
91///     )
92///     .await?;
93///
94/// assert_eq!(artifact.size_bytes, 2);
95/// # Ok(())
96/// # }
97/// ```
98pub trait ArtifactSink: Send + Sync {
99    /// Store the bytes of an artifact and record its metadata.
100    ///
101    /// # Errors
102    ///
103    /// Returns [`EngineError::Artifact`] when the name is invalid or storage
104    /// fails, and [`EngineError::Store`] when the metadata cannot be recorded
105    /// (for instance when the step already owns that name).
106    fn put<'a>(
107        &'a self,
108        upload: ArtifactUpload,
109        content: ByteStream,
110    ) -> ArtifactFuture<'a, Artifact>;
111
112    /// Open the bytes of a recorded artifact for reading.
113    ///
114    /// # Errors
115    ///
116    /// Returns [`EngineError::Artifact`] when the blob is missing or storage fails.
117    fn get<'a>(&'a self, artifact: &'a Artifact) -> ArtifactFuture<'a, ByteStream>;
118}
119
120/// [`ArtifactSink`] backed by a blob store and a run store in the same process.
121///
122/// Used by the API server and by any in-process engine. A remote worker uses an
123/// HTTP-backed sink instead.
124///
125/// # Examples
126///
127/// ```no_run
128/// use std::sync::Arc;
129///
130/// use ironflow_artifacts::local::LocalBlobStore;
131/// use ironflow_engine::artifact::DirectArtifactSink;
132/// use ironflow_store::memory::InMemoryStore;
133/// use ironflow_store::store::Store;
134///
135/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
136/// let blob = Arc::new(LocalBlobStore::new("/var/lib/ironflow/artifacts"));
137/// let sink = DirectArtifactSink::new(blob, store);
138/// # drop(sink);
139/// ```
140pub struct DirectArtifactSink {
141    blob: Arc<dyn BlobStore>,
142    store: Arc<dyn Store>,
143}
144
145impl DirectArtifactSink {
146    /// Build a sink over a blob store and the run store holding the metadata.
147    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            // Dedup: reuse an existing blob if one with the same SHA-256 exists.
172            // Dedup: reuse an existing blob if one with the same SHA-256 exists.
173            let (final_key, dedup_hit) =
174                match self.store.find_artifact_by_sha256(&digest.sha256).await {
175                    Ok(Some(existing)) => {
176                        if let Err(cleanup) = self.blob.delete(&key).await {
177                            warn!(
178                                storage_key = %key,
179                                error = %cleanup,
180                                "failed to remove deduplicated blob"
181                            );
182                        }
183                        info!(
184                            sha256 = %digest.sha256,
185                            reused_key = %existing.storage_key,
186                            "artifact deduplicated"
187                        );
188                        (existing.storage_key, true)
189                    }
190                    _ => (key.clone(), false),
191                };
192
193            let recorded = self
194                .store
195                .create_artifact(NewArtifact {
196                    id,
197                    run_id: upload.run_id,
198                    step_id: upload.step_id,
199                    name: upload.name,
200                    storage_key: final_key.clone(),
201                    content_type: upload.content_type,
202                    size_bytes: digest.size_bytes,
203                    sha256: digest.sha256,
204                })
205                .await;
206
207            match recorded {
208                Ok(artifact) => Ok(artifact),
209                Err(err) => {
210                    if !dedup_hit && let Err(cleanup) = self.blob.delete(&key).await {
211                        warn!(
212                            storage_key = %key,
213                            error = %cleanup,
214                            "failed to remove the blob of an unrecorded artifact"
215                        );
216                    }
217                    Err(err.into())
218                }
219            }
220        })
221    }
222
223    fn get<'a>(&'a self, artifact: &'a Artifact) -> ArtifactFuture<'a, ByteStream> {
224        Box::pin(async move { Ok(self.blob.get(&artifact.storage_key).await?) })
225    }
226}
227
228/// Where a step sits in its run, for resolving declared inputs.
229#[derive(Debug, Clone, Copy)]
230pub(crate) struct StepLocation {
231    /// Run being executed.
232    pub(crate) run_id: Uuid,
233    /// Attempt being executed.
234    pub(crate) attempt: u32,
235    /// Position of the step within the attempt.
236    pub(crate) position: u32,
237}
238
239/// Write every artifact a step declared as an input into its working directory.
240///
241/// Runs before the command starts, so the command sees the files it expects.
242pub(crate) async fn materialize_inputs(
243    sink: &Arc<dyn ArtifactSink>,
244    store: &Arc<dyn Store>,
245    config: &ShellConfig,
246    location: StepLocation,
247) -> Result<(), EngineError> {
248    let work_dir = working_dir(config);
249
250    for input in &config.inputs {
251        let artifact = store
252            .find_artifact_for_input(ArtifactLookup {
253                run_id: location.run_id,
254                attempt: location.attempt,
255                before_position: location.position,
256                step_name: input.step.clone(),
257                name: input.name.clone(),
258            })
259            .await?
260            .ok_or_else(|| EngineError::ArtifactNotFound {
261                step: input.step.clone(),
262                name: input.name.clone(),
263            })?;
264
265        let destination = work_dir.join(input.destination());
266        if let Some(parent) = destination.parent() {
267            create_dir_all(parent).await.map_err(ArtifactError::from)?;
268        }
269
270        let mut content = sink.get(&artifact).await?;
271        let mut file = File::create(&destination)
272            .await
273            .map_err(ArtifactError::from)?;
274        while let Some(chunk) = content.next().await {
275            file.write_all(&chunk?).await.map_err(ArtifactError::from)?;
276        }
277        file.flush().await.map_err(ArtifactError::from)?;
278
279        info!(
280            run_id = %location.run_id,
281            artifact = %input.name,
282            produced_by = %input.step,
283            destination = %destination.display(),
284            "artifact input materialized"
285        );
286    }
287
288    Ok(())
289}
290
291/// Store every file a step declared as an output.
292///
293/// `step_succeeded` decides how strict the collection is. On success, a pattern
294/// that matches nothing fails the step: a declared output that never appeared
295/// is a broken contract. On failure the collection is best-effort, because the
296/// partial files a failing step leaves behind are usually the useful ones.
297pub(crate) async fn collect_outputs(
298    sink: &Arc<dyn ArtifactSink>,
299    config: &ShellConfig,
300    run_id: Uuid,
301    step_id: Uuid,
302    step_name: &str,
303    step_succeeded: bool,
304) -> Result<(), EngineError> {
305    let work_dir = working_dir(config);
306
307    for output in &config.outputs {
308        let pattern = work_dir.join(&output.pattern);
309        let pattern = pattern.to_str().ok_or_else(|| {
310            EngineError::StepConfig(format!(
311                "output pattern {:?} is not valid UTF-8",
312                output.pattern
313            ))
314        })?;
315
316        let matches = glob(pattern)
317            .map_err(|err| {
318                EngineError::StepConfig(format!(
319                    "invalid output pattern {:?}: {err}",
320                    output.pattern
321                ))
322            })?
323            .filter_map(Result::ok)
324            .filter(|path| path.is_file())
325            .collect::<Vec<_>>();
326
327        if matches.is_empty() {
328            if step_succeeded {
329                return Err(EngineError::MissingArtifact {
330                    step: step_name.to_string(),
331                    pattern: output.pattern.clone(),
332                });
333            }
334            warn!(
335                run_id = %run_id,
336                step = %step_name,
337                pattern = %output.pattern,
338                "declared output matched no file on a failed step"
339            );
340            continue;
341        }
342
343        for path in matches {
344            let name = path
345                .file_name()
346                .and_then(|name| name.to_str())
347                .ok_or_else(|| {
348                    EngineError::StepConfig(format!(
349                        "output file {:?} has no valid UTF-8 name",
350                        path.display()
351                    ))
352                })?
353                .to_string();
354
355            let content_type = output
356                .content_type
357                .clone()
358                .unwrap_or_else(|| guess_content_type(&name));
359
360            let content = stream_from_path(&path).await?;
361            let artifact = sink
362                .put(
363                    ArtifactUpload {
364                        run_id,
365                        step_id,
366                        name: name.clone(),
367                        content_type,
368                    },
369                    content,
370                )
371                .await?;
372
373            info!(
374                run_id = %run_id,
375                step = %step_name,
376                artifact = %artifact.name,
377                size_bytes = artifact.size_bytes,
378                "artifact output stored"
379            );
380        }
381    }
382
383    Ok(())
384}
385
386/// Directory a shell step runs in, and the base for its artifact paths.
387fn working_dir(config: &ShellConfig) -> PathBuf {
388    PathBuf::from(config.dir.as_deref().unwrap_or("."))
389}
390
391#[cfg(test)]
392mod tests {
393    use std::collections::HashMap;
394
395    use futures_util::TryStreamExt;
396    use ironflow_artifacts::local::LocalBlobStore;
397    use ironflow_artifacts::stream_from_bytes;
398    use ironflow_store::entities::{NewRun, NewStep, StepKind, TriggerKind, step_trace_id};
399    use ironflow_store::memory::InMemoryStore;
400    use serde_json::json;
401    use tempfile::TempDir;
402
403    use super::*;
404
405    fn count_files_recursive(dir: &std::path::Path) -> usize {
406        let mut count = 0;
407        if let Ok(entries) = std::fs::read_dir(dir) {
408            for entry in entries.flatten() {
409                let path = entry.path();
410                if path.is_dir() {
411                    count += count_files_recursive(&path);
412                } else if path.is_file() {
413                    count += 1;
414                }
415            }
416        }
417        count
418    }
419
420    async fn sink_with_step() -> (TempDir, DirectArtifactSink, Uuid, Uuid) {
421        let dir = TempDir::new().expect("temp dir");
422        let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
423        let blob: Arc<dyn BlobStore> = Arc::new(LocalBlobStore::new(dir.path()));
424
425        let run = store
426            .create_run(NewRun {
427                workflow_name: "artifacts".to_string(),
428                trigger: TriggerKind::Manual,
429                payload: json!({}),
430                max_retries: 0,
431                handler_version: None,
432                labels: HashMap::new(),
433                scheduled_at: None,
434                created_by: None,
435                idempotency_key: None,
436                concurrency_key: None,
437                concurrency_limits: Vec::new(),
438                max_cost_usd: None,
439            })
440            .await
441            .expect("create run")
442            .into_run();
443
444        let step = store
445            .create_step(NewStep {
446                run_id: run.id,
447                trace_id: step_trace_id(run.id, "build", 0),
448                name: "build".to_string(),
449                kind: StepKind::Shell,
450                position: 0,
451                input: None,
452                is_error_handler: false,
453            })
454            .await
455            .expect("create step");
456
457        let sink = DirectArtifactSink::new(blob, store);
458        (dir, sink, run.id, step.id)
459    }
460
461    fn upload(run_id: Uuid, step_id: Uuid, name: &str) -> ArtifactUpload {
462        ArtifactUpload {
463            run_id,
464            step_id,
465            name: name.to_string(),
466            content_type: "text/plain".to_string(),
467        }
468    }
469
470    #[tokio::test]
471    async fn put_records_size_and_hash() {
472        let (_dir, sink, run_id, step_id) = sink_with_step().await;
473
474        let artifact = sink
475            .put(
476                upload(run_id, step_id, "report.txt"),
477                stream_from_bytes(b"abc".to_vec()),
478            )
479            .await
480            .expect("put");
481
482        assert_eq!(artifact.size_bytes, 3);
483        assert_eq!(
484            artifact.sha256,
485            "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad"
486        );
487    }
488
489    #[tokio::test]
490    async fn put_then_get_roundtrips_the_bytes() {
491        let (_dir, sink, run_id, step_id) = sink_with_step().await;
492
493        let artifact = sink
494            .put(
495                upload(run_id, step_id, "report.txt"),
496                stream_from_bytes(b"hello".to_vec()),
497            )
498            .await
499            .expect("put");
500
501        let chunks: Vec<bytes::Bytes> = sink
502            .get(&artifact)
503            .await
504            .expect("get")
505            .try_collect()
506            .await
507            .expect("collect");
508
509        assert_eq!(chunks.concat(), b"hello");
510    }
511
512    #[tokio::test]
513    async fn the_storage_key_never_embeds_the_name() {
514        let (_dir, sink, run_id, step_id) = sink_with_step().await;
515
516        let artifact = sink
517            .put(
518                upload(run_id, step_id, "report.txt"),
519                stream_from_bytes(b"x".to_vec()),
520            )
521            .await
522            .expect("put");
523
524        assert!(!artifact.storage_key.contains("report"));
525        assert!(artifact.storage_key.ends_with(&artifact.id.to_string()));
526    }
527
528    #[tokio::test]
529    async fn an_invalid_name_is_rejected_before_anything_is_written() {
530        let (dir, sink, run_id, step_id) = sink_with_step().await;
531
532        let err = sink
533            .put(
534                upload(run_id, step_id, "../escape"),
535                stream_from_bytes(b"x".to_vec()),
536            )
537            .await
538            .expect_err("invalid name");
539
540        assert!(matches!(err, EngineError::Artifact(_)));
541        assert!(!dir.path().join("artifacts").exists());
542    }
543
544    #[tokio::test]
545    async fn dedup_same_sha256_reuses_storage_key() {
546        let (dir, sink, run_id, step_id) = sink_with_step().await;
547
548        let second_step_id = sink
549            .store
550            .create_step(NewStep {
551                run_id,
552                trace_id: step_trace_id(run_id, "test", 1),
553                name: "test".to_string(),
554                kind: StepKind::Shell,
555                position: 1,
556                input: None,
557                is_error_handler: false,
558            })
559            .await
560            .expect("create step")
561            .id;
562
563        let first = sink
564            .put(
565                upload(run_id, step_id, "report.txt"),
566                stream_from_bytes(b"identical content".to_vec()),
567            )
568            .await
569            .expect("first put");
570
571        let second = sink
572            .put(
573                upload(run_id, second_step_id, "report-copy.txt"),
574                stream_from_bytes(b"identical content".to_vec()),
575            )
576            .await
577            .expect("second put");
578
579        assert_eq!(first.sha256, second.sha256);
580        assert_eq!(
581            first.storage_key, second.storage_key,
582            "dedup should reuse the same storage_key"
583        );
584
585        // Only one blob file should exist on disk: the dedup blob was deleted
586        let file_count = count_files_recursive(dir.path());
587        assert_eq!(file_count, 1, "dedup should not store a second copy");
588    }
589
590    #[tokio::test]
591    async fn dedup_different_sha256_gets_separate_storage_key() {
592        let (_dir, sink, run_id, step_id) = sink_with_step().await;
593
594        let second_step_id = sink
595            .store
596            .create_step(NewStep {
597                run_id,
598                trace_id: step_trace_id(run_id, "test", 1),
599                name: "test".to_string(),
600                kind: StepKind::Shell,
601                position: 1,
602                input: None,
603                is_error_handler: false,
604            })
605            .await
606            .expect("create step")
607            .id;
608
609        let first = sink
610            .put(
611                upload(run_id, step_id, "a.txt"),
612                stream_from_bytes(b"content A".to_vec()),
613            )
614            .await
615            .expect("first put");
616
617        let second = sink
618            .put(
619                upload(run_id, second_step_id, "b.txt"),
620                stream_from_bytes(b"content B".to_vec()),
621            )
622            .await
623            .expect("second put");
624
625        assert_ne!(first.sha256, second.sha256);
626        assert_ne!(
627            first.storage_key, second.storage_key,
628            "different content should get different storage keys"
629        );
630    }
631
632    #[tokio::test]
633    async fn a_duplicate_name_fails_and_leaves_no_orphan_blob() {
634        let (dir, sink, run_id, step_id) = sink_with_step().await;
635
636        sink.put(
637            upload(run_id, step_id, "report.txt"),
638            stream_from_bytes(b"first".to_vec()),
639        )
640        .await
641        .expect("first");
642
643        let err = sink
644            .put(
645                upload(run_id, step_id, "report.txt"),
646                stream_from_bytes(b"second".to_vec()),
647            )
648            .await
649            .expect_err("duplicate");
650
651        assert!(matches!(err, EngineError::Store(_)));
652
653        let stored: Vec<_> = std::fs::read_dir(
654            dir.path()
655                .join("artifacts")
656                .join(run_id.to_string())
657                .join(step_id.to_string()),
658        )
659        .expect("read dir")
660        .filter_map(Result::ok)
661        .collect();
662        assert_eq!(stored.len(), 1, "the rejected blob was not cleaned up");
663    }
664}