Skip to main content

ahri_tre_lake/
dataset_source.rs

1use std::path::{Path, PathBuf};
2
3use ahri_tre_types::DataFileRecord;
4
5use crate::{
6    DatasetFileFormat, DuckLakeAdapter, LakeError, MaterializeDataFileRequest, ScratchAttempt,
7    ScratchAttemptId,
8};
9
10/// Restored source bytes and intermediate files owned by one private attempt.
11/// The restricted path remains inside the trusted workflow and Lake adapter.
12pub struct PreparedDatasetSource {
13    attempt: ScratchAttempt,
14    path: PathBuf,
15}
16
17impl std::fmt::Debug for PreparedDatasetSource {
18    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
19        formatter
20            .debug_struct("PreparedDatasetSource")
21            .field("attempt", &self.attempt)
22            .finish_non_exhaustive()
23    }
24}
25
26impl PreparedDatasetSource {
27    pub fn path(&self) -> &Path {
28        &self.path
29    }
30
31    /// Completes source cleanup before Dataset metadata can be admitted.
32    pub fn cleanup(self) -> Result<(), LakeError> {
33        self.attempt.cleanup().map_err(Into::into)
34    }
35
36    fn prepare(
37        scratch: &ScratchAttempt,
38        format: DatasetFileFormat,
39        restore: impl FnOnce(&Path) -> Result<(), LakeError>,
40    ) -> Result<Self, LakeError> {
41        let identifier = ScratchAttemptId::new(&uuid::Uuid::new_v4().simple().to_string())?;
42        let attempt = scratch.create_child(identifier)?;
43        let path = attempt
44            .path()
45            .join(format!("source.{}", format.extension()));
46        if let Err(error) = restore(&path) {
47            return match attempt.cleanup() {
48                Ok(()) => Err(error),
49                Err(cleanup) => Err(LakeError::PhysicalRollback {
50                    operation: "Dataset source restoration failed".into(),
51                    rollback: cleanup.to_string(),
52                }),
53            };
54        }
55        Ok(Self { attempt, path })
56    }
57}
58
59impl DuckLakeAdapter {
60    /// Restores a managed Datafile using only the configured scratch capability.
61    /// Failed restoration removes the source and all partial intermediate files.
62    pub fn restore_dataset_source(
63        &self,
64        scratch: &ScratchAttempt,
65        datafile: &DataFileRecord,
66        format: DatasetFileFormat,
67    ) -> Result<PreparedDatasetSource, LakeError> {
68        PreparedDatasetSource::prepare(scratch, format, |path| {
69            self.materialize_datafile(&MaterializeDataFileRequest {
70                datafile: datafile.clone(),
71                destination_path: path.to_path_buf(),
72            })
73            .map(|_| ())
74        })
75    }
76
77    /// Observes the existing storage read without formatting content or paths.
78    pub fn restore_dataset_source_observed(
79        &self,
80        context: Option<ahri_tre_observability::CorrelationContext>,
81        scratch: &ScratchAttempt,
82        datafile: &DataFileRecord,
83        format: DatasetFileFormat,
84    ) -> Result<PreparedDatasetSource, LakeError> {
85        let span = context.map(|context| context.span(ahri_tre_observability::Stage::LakeStorage));
86        let result = self.restore_dataset_source(scratch, datafile, format);
87        crate::error::finish_dataset_observation(span, &result, &[]);
88        result
89    }
90
91    /// Stages an external source whose reader needs writable conversion scratch.
92    pub fn stage_dataset_source_file(
93        &self,
94        scratch: &ScratchAttempt,
95        source: &Path,
96        format: DatasetFileFormat,
97    ) -> Result<PreparedDatasetSource, LakeError> {
98        PreparedDatasetSource::prepare(scratch, format, |path| {
99            std::fs::copy(source, path).map(|_| ()).map_err(Into::into)
100        })
101    }
102}