Skip to main content

ahri_tre_app/service/
upload.rs

1use super::*;
2use ahri_tre_protocol::{ingest, request::ProtocolRequest, upload::UploadRequest};
3use ahri_tre_types::NcName;
4use std::time::Instant;
5
6struct UploadLifecycle<'a> {
7    source_cleanup_safe: &'a mut bool,
8    output_intent: Option<DatasetOutputIntent>,
9    source: ahri_tre_types::VersionId,
10    deadline: Instant,
11    classification: ahri_tre_types::IngestClassification,
12    cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
13    monitor: Option<UploadMonitor>,
14    metadata_deadline: Option<ahri_tre_pgmeta::UploadMetadataDeadline>,
15}
16impl datafile_materialization::DatafileLifecycle for UploadLifecycle<'_> {
17    fn accepted(&mut self) {
18        *self.source_cleanup_safe = false;
19    }
20    fn cleanup_finished(&mut self, succeeded: bool) {
21        *self.source_cleanup_safe = succeeded;
22    }
23    fn output_intent(&self) -> Option<DatasetOutputIntent> {
24        self.output_intent
25    }
26    fn ordinary_upload_source(&self) -> Option<ahri_tre_types::VersionId> {
27        Some(self.source)
28    }
29    fn cleanup_started(&mut self) {
30        self.monitor.take();
31        self.metadata_deadline.take();
32    }
33    fn classification(&self) -> ahri_tre_types::IngestClassification {
34        self.classification.clone()
35    }
36    fn record_completion(
37        &mut self,
38        _: &ahri_tre_pgmeta::PgMetadataRepository<'_>,
39        _: &DatasetMaterialization,
40    ) -> Result<(), AppError> {
41        self.check_deadline()
42    }
43    fn checkpoint(
44        &mut self,
45        _: datafile_materialization::DatafileStage,
46        _: &datafile_materialization::ResolvedDatafileMaterialization,
47    ) -> Result<(), AppError> {
48        self.check_deadline()
49    }
50}
51impl UploadLifecycle<'_> {
52    fn check_deadline(&self) -> Result<(), AppError> {
53        if Instant::now() >= self.deadline
54            || self.cancelled.load(std::sync::atomic::Ordering::Acquire)
55        {
56            return Err(AppError::Validation("Upload deadline elapsed".into()));
57        }
58        Ok(())
59    }
60}
61impl AppService {
62    /// Receives a path-free upload under retained Session authority. All output
63    /// naming, classification and materialization decisions remain app-owned.
64    pub async fn ingest_upload(
65        &self,
66        session: &mut DataStoreSession,
67        status: ahri_tre_protocol::session::SessionStatusPayload,
68        upload: UploadRequest,
69        input: ahri_tre_runtime::upload::UploadInput,
70    ) -> Result<serde_json::Value, AppError> {
71        if matches!(upload.request.request, ProtocolRequest::IngestRedcap(_)) {
72            return self
73                .ingest_redcap_upload(session, status, upload, input, None)
74                .await;
75        }
76        self.ingest_upload_with_acquisition(session, status, upload, input, None)
77            .await
78    }
79    pub(super) async fn ingest_upload_with_acquisition(
80        &self,
81        session: &mut DataStoreSession,
82        status: ahri_tre_protocol::session::SessionStatusPayload,
83        mut upload: UploadRequest,
84        input: ahri_tre_runtime::upload::UploadInput,
85        acquisition: Option<ahri_tre_types::AcquisitionPath>,
86    ) -> Result<serde_json::Value, AppError> {
87        let output_intent = match &upload.request.request {
88            ProtocolRequest::IngestSql(request) => Some(DatasetOutputIntent::parse(
89                request.materialization.replace,
90                request.materialization.new_version.as_deref(),
91            )?),
92            _ => None,
93        };
94        let sql_source = if let ProtocolRequest::IngestSql(request) = &upload.request.request {
95            if acquisition.is_none()
96                || upload.metadata.filename != "source.parquet"
97                || upload.metadata.media_type != "application/vnd.apache.parquet"
98            {
99                return Err(AppError::Validation(
100                    "SQL results require admitted acquisition".into(),
101                ));
102            }
103            use sha2::Digest;
104            let identity = match &request.source {
105                ingest::SqlSource::Stream { engine } => format!("{engine:?}"),
106                ingest::SqlSource::Remote { endpoint } => format!(
107                    "{} endpoint sha256:{}",
108                    endpoint.as_str().split(':').next().unwrap_or("sql"),
109                    sha2::Sha256::digest(endpoint.as_str().as_bytes())
110                        .iter()
111                        .map(|byte| format!("{byte:02x}"))
112                        .collect::<String>()
113                ),
114                _ => {
115                    return Err(AppError::Validation(
116                        "SQL source must not contain a local path".into(),
117                    ));
118                }
119            };
120            let query_digest = sha2::Sha256::digest(request.sql.as_bytes())
121                .iter()
122                .map(|byte| format!("{byte:02x}"))
123                .collect::<String>();
124            let table = ingest::IngestDatasetFileRequest {
125                classification: request.classification.clone(),
126                session: request.session.clone(),
127                study: request.study.clone(),
128                domain: request.domain.clone(),
129                dataset: request.dataset.clone(),
130                source: ingest::IngestDatasetFileSource::Stream,
131                materialization: ingest::DatasetMaterializationOptions {
132                    description: request.materialization.description.clone(),
133                    replace: request.materialization.replace,
134                    new_version: request.materialization.new_version.clone(),
135                },
136                parse: ingest::DatasetFileParseOptions {
137                    format: Some("parquet".into()),
138                    sheet: None,
139                    header: None,
140                    delimiter: None,
141                    null_strings: Vec::new(),
142                    json_format: None,
143                    sample_rows: None,
144                    // SQL carries an explicit schema; sampled values must not change it.
145                    vocabulary_threshold: Some(0),
146                    force_metadata_updates: false,
147                },
148            };
149            upload.request.request = ProtocolRequest::IngestDatasetFile(table);
150            Some(format!("SQL {identity} query sha256:{query_digest}"))
151        } else {
152            None
153        };
154        let ahri_tre_runtime::upload::UploadInput {
155            source,
156            budgets,
157            deadline,
158            cancelled,
159        } = input;
160        let preflight_deadline = session
161            .dataset_admission_repository()
162            .map(|repository| repository.upload_deadline(deadline, cancelled.clone()))
163            .transpose()
164            .map_err(|error| core_error("bound upload preflight", error))?;
165        let explicit_study = match &upload.request.request {
166            ProtocolRequest::IngestFile(value) => value.study.clone(),
167            ProtocolRequest::IngestDatasetFile(value) => Some(value.study.clone()),
168            _ => None,
169        };
170        let selected = self
171            .resolve_session_study(
172                self.catalogue_datastore_id()?.as_uuid(),
173                explicit_study,
174                session.current_study().map(|study| study.0),
175            )
176            .await
177            .map_err(|error| match error.code {
178                ahri_tre_protocol::ProtocolErrorCode::Conflict => {
179                    AppError::Conflict("Upload Study selectors disagree".into())
180                }
181                ahri_tre_protocol::ProtocolErrorCode::NotFound => {
182                    AppError::NotFound("Upload Study is unavailable".into())
183                }
184                _ => AppError::Validation("Upload Study is invalid".into()),
185            })?;
186        let study = StudySelector::Id {
187            study_id: ahri_tre_types::StudyId(selected.id.as_uuid()),
188        };
189        let invalid = || AppError::Validation("Invalid upload request".into());
190        let (file, table) = match upload.request.request {
191            ProtocolRequest::IngestFile(request)
192                if matches!(request.source, ingest::IngestFileSource::Stream) =>
193            {
194                (request, None)
195            }
196            ProtocolRequest::IngestDatasetFile(request)
197                if matches!(request.source, ingest::IngestDatasetFileSource::Stream) =>
198            {
199                if sql_source.is_none()
200                    && (request.materialization.replace
201                        || request.materialization.new_version.is_some())
202                {
203                    return Err(AppError::Validation("Dataset output versions are allocated automatically; replacement and explicit versions are not supported".into()));
204                }
205                // Resolve the metadata Domain before acquiring source content.
206                let domain = self.catalogue_domain_selector(request.domain.clone())?;
207                self.get_domain(GetDomainRequest { domain })
208                    .await?
209                    .ok_or_else(invalid)?;
210                let name = derived_source_file_asset_name(
211                    &NcName::parse(request.dataset.as_str().to_string()).map_err(|_| invalid())?,
212                );
213                (
214                    ingest::IngestFileRequest {
215                        classification: request.classification.clone(),
216                        session: request.session.clone(),
217                        study: Some(request.study.clone()),
218                        asset: ahri_tre_protocol::PublicName::new(name.as_str())
219                            .map_err(|_| invalid())?,
220                        source: ingest::IngestFileSource::Stream,
221                        format: request.parse.format.clone().ok_or_else(invalid)?,
222                        description: Some(request.materialization.description.clone()),
223                        new_version: false,
224                        bump_major: false,
225                        bump_minor: false,
226                        compress: true,
227                        encrypt: None,
228                        hash_source_local: false,
229                        verify_copy: true,
230                    },
231                    Some(request),
232                )
233            }
234            _ => return Err(invalid()),
235        };
236        let format = ahri_tre_protocol::upload::upload_format(
237            &upload.metadata.filename,
238            &file.format,
239            &upload.metadata.media_type,
240        )
241        .ok_or_else(invalid)?;
242        let source_format = resolve_dataset_file_format(None, Some(format))?;
243        let options = match &table {
244            Some(table) => dataset_file_read_options(
245                table.parse.header,
246                table.parse.delimiter,
247                table.parse.null_strings.clone(),
248                table.parse.json_format.clone(),
249                table.parse.sheet.clone(),
250            ),
251            None => DatasetFileReadOptions::default(),
252        };
253        let metadata = upload.metadata;
254        let mut content = IngestFileContent::new(
255            metadata.filename.clone(),
256            Some(metadata.media_type.clone()),
257            Some(metadata.content_length),
258            Some(metadata.digest.clone()),
259            CancelledReader {
260                source,
261                cancelled: cancelled.clone(),
262                deadline,
263            },
264        );
265        content.acquisition = acquisition;
266        content.upload_validation = Some(ahri_tre_lake::UploadValidation {
267            format: source_format,
268            options,
269            budgets,
270            deadline,
271            cancelled: cancelled.clone(),
272        });
273        drop(preflight_deadline);
274        let ingested = self
275            .ingest_datafile(
276                session,
277                IngestDataFileRequest {
278                    classification: file.classification,
279                    study: study.clone(),
280                    asset_name: NcName::parse(file.asset.as_str().to_string())
281                        .map_err(|_| invalid())?,
282                    content,
283                    format: format.to_owned(),
284                    description: file.description,
285                    new_version: file.new_version,
286                    bump_major: file.bump_major,
287                    bump_minor: file.bump_minor,
288                    compress: file.compress,
289                    encrypt: file.encrypt,
290                },
291            )
292            .await?;
293        let Some(table) = table else {
294            let response = crate::projections::datafile_ingest_response(
295                status,
296                ingested,
297                ingest::IngestSourceSummary {
298                    kind: if acquisition.is_some() {
299                        ingest::IngestSourceKind::HttpsUri
300                    } else {
301                        ingest::IngestSourceKind::UploadedContent
302                    },
303                    request_only: true,
304                    upload: Some(ingest::UploadSummary {
305                        acquisition,
306                        budgets: Some(budgets),
307                        filename: metadata.filename,
308                        media_type: metadata.media_type,
309                        content_length: metadata.content_length,
310                        digest: metadata.digest,
311                        outcome: ingest::UploadOutcome::Ingested,
312                    }),
313                },
314            );
315            return serde_json::to_value(response).map_err(|_| invalid());
316        };
317        let mut source_cleanup_safe = true;
318        let materialized = async {
319            let source_version = ingested.datafile.datafile_id;
320            let metadata_domain = self.catalogue_domain_selector(table.domain)?;
321            let datastore_id = self.catalogue_datastore_id()?;
322            let request = DeriveDatasetFromManagedDataFileRequest {
323                study,
324                metadata_domain,
325                dataset_name: NcName::parse(table.dataset.as_str().to_string())
326                    .map_err(|_| invalid())?,
327                source: ingest::ManagedDataFileSource {
328                    asset: ahri_tre_protocol::asset::AssetSelector::Id {
329                        asset: ingest::version_ref(datastore_id, source_version),
330                        study: None,
331                        asset_type: Some(ahri_tre_types::AssetType::File),
332                    },
333                    version: None,
334                },
335                input_format: table.parse.format,
336                header: table.parse.header,
337                delimiter: table.parse.delimiter,
338                null_strings: table.parse.null_strings,
339                json_format: table.parse.json_format,
340                sheet: table.parse.sheet,
341                description: Some(table.materialization.description.clone()),
342                version_note: Some(table.materialization.description),
343                replace: table.materialization.replace,
344                new_version: table.materialization.new_version,
345                transformation: ahri_tre_types::NewTransformationRecord {
346                    transformation_type: ahri_tre_types::TransformationType::Transform,
347                    description: format!(
348                        "materialize ordinary-ingest table from managed Datafile; {}; {}",
349                        sql_source.as_deref().unwrap_or("uploaded table"),
350                        if sql_source.is_some() {
351                            match acquisition {
352                                Some(ahri_tre_types::AcquisitionPath::Trusted) => {
353                                    "Trusted-observed acquisition"
354                                }
355                                _ => "client-asserted acquisition",
356                            }
357                        } else {
358                            acquisition_attribution(acquisition)
359                        }
360                    ),
361                    repository_url: None,
362                    commit_hash: None,
363                    file_path: None,
364                    date_created: None,
365                    created_by: session.authenticated_actor().map(str::to_owned),
366                },
367                sample_rows: table.parse.sample_rows,
368                vocabulary_threshold: table.parse.vocabulary_threshold,
369                force_metadata_updates: table.parse.force_metadata_updates,
370            };
371            let metadata_deadline = session
372                .dataset_admission_repository()
373                .map(|repository| repository.upload_deadline(deadline, cancelled.clone()))
374                .transpose()
375                .map_err(|error| core_error("bound upload metadata", error))?;
376            let monitor = UploadMonitor::new(session, deadline, cancelled.clone());
377            let result = self
378                .derive_dataset_with_lifecycle(
379                    session,
380                    request,
381                    &mut UploadLifecycle {
382                        source_cleanup_safe: &mut source_cleanup_safe,
383                        output_intent,
384                        source: source_version,
385                        deadline,
386                        classification: table.classification,
387                        cancelled,
388                        monitor: Some(monitor),
389                        metadata_deadline,
390                    },
391                )
392                .await?;
393            Ok::<_, AppError>(result)
394        }
395        .await;
396        let result = match materialized {
397            Ok(result) => result,
398            Err(error) => {
399                if sql_source.is_some()
400                    && source_cleanup_safe
401                    && !matches!(error, AppError::DatasetAdmissionOutcomeUnknown(_))
402                {
403                    self.rollback_uploaded_table_source(session, ingested)
404                        .await?;
405                }
406                return Err(error);
407            }
408        };
409        let mut response = crate::projections::dataset_materialization_response(
410            status,
411            result.study.name.as_str().to_string(),
412            result.materialization,
413        );
414        response.source = ingest::DatasetMaterializationSourceSummary {
415            upload: Some(ingest::UploadSummary {
416                acquisition,
417                budgets: Some(budgets),
418                filename: metadata.filename,
419                media_type: metadata.media_type,
420                content_length: metadata.content_length,
421                digest: metadata.digest,
422                outcome: ingest::UploadOutcome::Ingested,
423            }),
424            kind: if sql_source.is_some() {
425                ingest::DatasetMaterializationSourceKind::SqlSource
426            } else if acquisition.is_some() {
427                ingest::DatasetMaterializationSourceKind::Https
428            } else {
429                ingest::DatasetMaterializationSourceKind::LocalDatasetFile
430            },
431            request_only: true,
432        };
433        serde_json::to_value(response).map_err(|_| invalid())
434    }
435
436    /// A source is admitted before table parsing. A refused Dataset must restore
437    /// the source catalogue as well as the Dataset attempt's Lake/metadata state.
438    async fn rollback_uploaded_table_source(
439        &self,
440        session: &DataStoreSession,
441        ingested: crate::DataFileIngestResult,
442    ) -> Result<(), AppError> {
443        let _guard = Arc::clone(&FILE_INGEST_WORKFLOW_LOCK).lock_owned().await;
444        let asset = ingested.catalog.asset.asset_id;
445        let version = ingested.datafile.datafile_id;
446        let datafile = self
447            .assets
448            .get_datafile(version)
449            .await
450            .map_err(|error| core_error("read failed upload source", error))?
451            .ok_or_else(|| AppError::Infrastructure("Failed upload source is missing".into()))?;
452        let versions = self
453            .assets
454            .list_asset_versions(asset)
455            .await
456            .map_err(|error| core_error("read failed upload versions", error))?;
457        let latest = versions.iter().find(|item| item.is_latest == Some(true));
458        let promotion = if latest.is_some_and(|item| item.version_id == version) {
459            Some((
460                asset,
461                versions
462                    .iter()
463                    .filter(|item| item.version_id != version)
464                    .max_by_key(|item| (item.major, item.minor, item.patch))
465                    .map(|item| item.version_id),
466            ))
467        } else {
468            None
469        };
470        let lake = DuckLakeAdapter::new(&session.runtime().lake.data_path);
471        run_blocking_lake_work(move || lake.delete_datafile(&DeleteDataFileRequest { datafile }))
472            .await
473            .map_err(|error| infrastructure_error("remove failed upload source payload", error))?;
474        for transformation in self
475            .transformations
476            .list_transformations_producing(version)
477            .await
478            .map_err(|error| core_error("read failed upload provenance", error))?
479        {
480            self.transformations
481                .delete_transformation(transformation.transformation_id)
482                .await
483                .map_err(|error| core_error("remove failed upload provenance", error))?;
484        }
485        self.assets
486            .delete_datafile_version_metadata(version, promotion)
487            .await
488            .map_err(|error| core_error("restore failed upload source catalogue", error))?;
489        if versions.iter().all(|item| item.version_id == version) {
490            self.assets
491                .delete_asset_record(asset)
492                .await
493                .map_err(|error| core_error("remove failed upload source Asset", error))?;
494        }
495        Ok(())
496    }
497}
498
499struct CancelledReader {
500    source: Box<dyn std::io::Read + Send>,
501    cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
502    deadline: Instant,
503}
504impl std::io::Read for CancelledReader {
505    fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
506        if Instant::now() >= self.deadline
507            || self.cancelled.load(std::sync::atomic::Ordering::Acquire)
508        {
509            return Err(std::io::Error::other("Upload stopped"));
510        }
511        self.source.read(out)
512    }
513}
514pub(super) struct UploadMonitor {
515    stop: std::sync::mpsc::Sender<()>,
516    worker: Option<std::thread::JoinHandle<()>>,
517}
518impl UploadMonitor {
519    pub(super) fn new(
520        session: &DataStoreSession,
521        deadline: Instant,
522        cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
523    ) -> Self {
524        let interrupt = session.lake_connection().interrupt_handle();
525        let (stop, stopped) = std::sync::mpsc::channel();
526        let worker = std::thread::spawn(move || {
527            while stopped
528                .recv_timeout(std::time::Duration::from_millis(25))
529                .is_err()
530            {
531                if Instant::now() >= deadline
532                    || cancelled.load(std::sync::atomic::Ordering::Acquire)
533                {
534                    interrupt.interrupt();
535                }
536            }
537        });
538        Self {
539            stop,
540            worker: Some(worker),
541        }
542    }
543}
544impl Drop for UploadMonitor {
545    fn drop(&mut self) {
546        let _ = self.stop.send(());
547        if let Some(worker) = self.worker.take() {
548            let _ = worker.join();
549        }
550    }
551}
552
553fn acquisition_attribution(path: Option<ahri_tre_types::AcquisitionPath>) -> &'static str {
554    match path {
555        Some(ahri_tre_types::AcquisitionPath::Trusted) => {
556            "HTTPS source, Trusted-observed acquisition"
557        }
558        Some(ahri_tre_types::AcquisitionPath::Client) => {
559            "HTTPS source, client-asserted acquisition"
560        }
561        None => "client-uploaded source",
562    }
563}