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 (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#[derive(Debug, Clone, Copy)]
230pub(crate) struct StepLocation {
231 pub(crate) run_id: Uuid,
233 pub(crate) attempt: u32,
235 pub(crate) position: u32,
237}
238
239pub(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
291pub(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
386fn 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 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}