Skip to main content

ahri_tre_app/service/
datafile_materialization.rs

1//! Lifecycle boundaries for the existing managed Datafile producer.
2use super::*;
3use ahri_tre_observability::{CorrelationContext, FailureCategory, Outcome, Stage};
4
5/// Observe an existing workflow boundary, preserving its exact result.
6pub(super) async fn observed<T>(
7    context: Option<CorrelationContext>,
8    stage: Stage,
9    work: impl std::future::Future<Output = Result<T, AppError>>,
10) -> Result<T, AppError> {
11    let span = context.map(|context| context.span(stage));
12    let result = work.await;
13    if let Some(span) = span {
14        let (outcome, category) = match &result {
15            Ok(_) => (Outcome::Success, None),
16            Err(AppError::DatasetAdmissionOutcomeUnknown(_)) => {
17                (Outcome::Unavailable, Some(FailureCategory::CommitUnknown))
18            }
19            Err(AppError::Infrastructure(_)) => (
20                Outcome::Unavailable,
21                Some(match stage {
22                    Stage::Metadata | Stage::MetadataCommit => FailureCategory::Metadata,
23                    Stage::Compensation => FailureCategory::Cleanup,
24                    _ => FailureCategory::Dependency,
25                }),
26            ),
27            Err(_) => (Outcome::Rejected, None),
28        };
29        span.finish(outcome, category);
30    }
31    result
32}
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35pub(super) enum DatafileStage {
36    BeforePreparation,
37    BeforeLoad,
38    AfterLoad,
39    AfterPreparation,
40    BeforeMetadataCommit,
41}
42
43/// Cooperative stops occur only before admission. Cleanup is always allowed to
44/// finish; it cannot be interrupted by a cancellation checkpoint.
45pub(super) trait DatafileLifecycle {
46    fn output_intent(&self) -> Option<DatasetOutputIntent> {
47        None
48    }
49    fn ordinary_upload_source(&self) -> Option<ahri_tre_types::VersionId> {
50        None
51    }
52    fn classification(&self) -> ahri_tre_types::IngestClassification {
53        ahri_tre_types::IngestClassification::derived_high()
54    }
55
56    fn observation(&self) -> Option<CorrelationContext> {
57        None
58    }
59
60    fn accepted(&mut self) {}
61    fn claim_execution(&mut self) -> Result<(), AppError> {
62        Ok(())
63    }
64    fn resolved_source(
65        &mut self,
66        _study: &ahri_tre_types::StudyRecord,
67        _source: &DataFileMetadata,
68    ) {
69    }
70    /// Joins the shared output reservation transaction before allocation or
71    /// preparation. The resolved source and intent are already pinned.
72    fn record_acceptance(
73        &mut self,
74        _repository: &ahri_tre_pgmeta::PgMetadataRepository<'_>,
75        _attempt: &ahri_tre_pgmeta::RecoverableDatasetAttempt,
76        _resolved: &ResolvedDatafileMaterialization,
77    ) -> Result<(), CoreError> {
78        Ok(())
79    }
80    fn checkpoint(
81        &mut self,
82        stage: DatafileStage,
83        resolved: &ResolvedDatafileMaterialization,
84    ) -> Result<(), AppError>;
85    /// Runs in the final metadata transaction. Later operation completion may
86    /// use typed repositories here; a rejection rolls back the complete admission.
87    fn record_completion(
88        &mut self,
89        _repository: &ahri_tre_pgmeta::PgMetadataRepository<'_>,
90        _result: &DatasetMaterialization,
91    ) -> Result<(), AppError> {
92        Ok(())
93    }
94    fn cleanup_started(&mut self) {}
95    fn cleanup_finished(&mut self, _succeeded: bool) {}
96}
97pub(super) struct SynchronousDatafileLifecycle;
98impl DatafileLifecycle for SynchronousDatafileLifecycle {
99    fn checkpoint(
100        &mut self,
101        _: DatafileStage,
102        _: &ResolvedDatafileMaterialization,
103    ) -> Result<(), AppError> {
104        Ok(())
105    }
106}
107
108pub(super) struct ResolvedDatafileMaterialization {
109    pub request: DataFileToDatasetRequest,
110    source_datafile: ahri_tre_types::DataFileRecord,
111    pub source_version: ahri_tre_types::AssetVersionRecord,
112    source_format: DatasetFileFormat,
113    read_options: DatasetFileReadOptions,
114}
115struct PreparedDatafileDataset {
116    loaded: ahri_tre_lake::LoadedDatasetTable,
117    transformation: ahri_tre_types::NewTransformationRecord,
118    variable_registrations: Vec<DatasetVariableRegistrationInput>,
119    metadata_warnings: Vec<String>,
120}
121
122impl AppService {
123    pub async fn datafile_to_dataset(
124        &self,
125        session: &mut DataStoreSession,
126        request: DataFileToDatasetRequest,
127    ) -> Result<DatasetMaterialization, AppError> {
128        self.datafile_to_dataset_with_lifecycle(session, request, &mut SynchronousDatafileLifecycle)
129            .await
130    }
131
132    pub(super) async fn datafile_to_dataset_with_lifecycle(
133        &self,
134        session: &mut DataStoreSession,
135        request: DataFileToDatasetRequest,
136        lifecycle: &mut impl DatafileLifecycle,
137    ) -> Result<DatasetMaterialization, AppError> {
138        let context = lifecycle.observation();
139        let resolved = observed(
140            context,
141            Stage::Metadata,
142            self.resolve_datafile_materialization(request),
143        )
144        .await;
145        if resolved.is_ok() || matches!(&resolved, Err(AppError::Infrastructure(_))) {
146            session.observations().record(
147                ahri_tre_protocol::diagnostics::DiagnosticComponent::Metadata,
148                ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
149                resolved.is_ok(),
150            );
151        }
152        let resolved = resolved?;
153        let request = &resolved.request;
154        let classification = lifecycle.classification();
155        let output_intent = lifecycle.output_intent();
156        let lifecycle = std::cell::RefCell::new(lifecycle);
157        let mut attempt = observed(
158            context,
159            Stage::Preparation,
160            self.prepare_dataset_attempt_with_notification(
161                session,
162                PrepareDatasetVersionRequest {
163                    output_intent,
164                    inherit_risk: false,
165                    classification,
166                    study_id: request.study_id,
167                    dataset_asset_id: request.dataset_asset_id,
168                    dataset_name: request.dataset_name.clone(),
169                    dataset_version_id: request.dataset_version_id,
170                    description: request.description.clone(),
171                    agent_instructions: request.agent_instructions.clone(),
172                    version_note: request.version_note.clone(),
173                    created_by: request.transformation.created_by.clone(),
174                },
175                context,
176                |repository, attempt| {
177                    lifecycle
178                        .borrow_mut()
179                        .record_acceptance(&repository, attempt, &resolved)
180                },
181                || {
182                    let mut lifecycle = lifecycle.borrow_mut();
183                    lifecycle.accepted();
184                    lifecycle.claim_execution()
185                },
186                |succeeded| lifecycle.borrow_mut().cleanup_finished(succeeded),
187            ),
188        )
189        .await?;
190        let lifecycle = lifecycle.into_inner();
191        attempt.prepared.ordinary_upload_source = lifecycle.ordinary_upload_source();
192        let result = async {
193            lifecycle.checkpoint(DatafileStage::BeforePreparation, &resolved)?;
194            let prepared = observed(
195                context,
196                Stage::Preparation,
197                self.prepare_datafile_dataset(session, &mut attempt, &resolved, lifecycle),
198            )
199            .await?;
200            lifecycle.checkpoint(DatafileStage::AfterPreparation, &resolved)?;
201            lifecycle.checkpoint(DatafileStage::BeforeMetadataCommit, &resolved)?;
202            let committed = observed(
203                context,
204                Stage::MetadataCommit,
205                self.commit_datafile_dataset(session, &mut attempt, &resolved, prepared, lifecycle),
206            )
207            .await;
208            if committed.is_ok() || matches!(&committed, Err(AppError::Infrastructure(_))) {
209                session.observations().record(
210                    ahri_tre_protocol::diagnostics::DiagnosticComponent::Metadata,
211                    ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
212                    committed.is_ok(),
213                );
214            }
215            committed
216        }
217        .await;
218        if result.is_err() && !matches!(result, Err(AppError::DatasetAdmissionOutcomeUnknown(_))) {
219            lifecycle.cleanup_started();
220        }
221        let cleanup = if result.is_err()
222            && !matches!(result, Err(AppError::DatasetAdmissionOutcomeUnknown(_)))
223        {
224            context.map(|context| context.span(Stage::Compensation))
225        } else {
226            None
227        };
228        attempt.finish_with_cleanup(session, result, |succeeded| {
229            lifecycle.cleanup_finished(succeeded);
230            if let Some(span) = cleanup {
231                span.finish(
232                    if succeeded {
233                        Outcome::Success
234                    } else {
235                        Outcome::Unavailable
236                    },
237                    (!succeeded).then_some(FailureCategory::Cleanup),
238                );
239            }
240        })
241    }
242
243    async fn resolve_datafile_materialization(
244        &self,
245        mut request: DataFileToDatasetRequest,
246    ) -> Result<ResolvedDatafileMaterialization, AppError> {
247        if request.transformation.transformation_type
248            != ahri_tre_types::TransformationType::Transform
249        {
250            return Err(AppError::Validation(
251                "datafile-to-dataset requires a transform transformation".to_string(),
252            ));
253        }
254
255        let source_datafile = self
256            .assets
257            .get_datafile(request.source_datafile_id)
258            .await
259            .map_err(|error| core_error("lookup source datafile", error))?
260            .ok_or_else(|| {
261                AppError::Validation(format!(
262                    "source datafile not found: {}",
263                    request.source_datafile_id.0
264                ))
265            })?;
266        let source_version = self
267            .assets
268            .get_asset_version(source_datafile.datafile_id)
269            .await
270            .map_err(|error| core_error("lookup source datafile asset version", error))?
271            .ok_or_else(|| {
272                AppError::Validation(format!(
273                    "source datafile version not found: {}",
274                    source_datafile.datafile_id.0
275                ))
276            })?;
277        let source_asset = self
278            .assets
279            .get_asset_by_id(source_version.asset_id)
280            .await
281            .map_err(|error| core_error("lookup source datafile asset", error))?
282            .ok_or_else(|| {
283                AppError::Validation(format!(
284                    "source datafile asset not found: {}",
285                    source_version.asset_id.0
286                ))
287            })?;
288        if source_asset.study_id != request.study_id {
289            return Err(AppError::Validation(format!(
290                "source datafile asset {} belongs to study {}, not {}",
291                source_asset.asset_id.0, source_asset.study_id.0, request.study_id.0
292            )));
293        }
294        if source_asset.asset_type != ahri_tre_types::AssetType::File {
295            return Err(AppError::Validation(format!(
296                "source datafile asset {} is not a file asset",
297                source_asset.asset_id.0
298            )));
299        }
300
301        let source_format = resolve_dataset_file_format(
302            None,
303            request
304                .input_format
305                .as_deref()
306                .or(Some(source_datafile.edam_format.as_str())),
307        )?;
308        if request.variable_registrations.is_empty()
309            && automatic_fallback_variable_format(source_format)
310        {
311            request.metadata_domain_id = Some(
312                self.resolve_automatic_metadata_domain(
313                    request.study_id,
314                    request.metadata_domain_id,
315                    source_format,
316                )
317                .await?,
318            );
319        }
320        let read_options = dataset_file_read_options(
321            request.header,
322            request.delimiter,
323            request.null_strings.clone(),
324            request.json_format.clone(),
325            request.sheet.clone(),
326        );
327
328        Ok(ResolvedDatafileMaterialization {
329            request,
330            source_datafile,
331            source_version,
332            source_format,
333            read_options,
334        })
335    }
336
337    async fn prepare_datafile_dataset(
338        &self,
339        session: &mut DataStoreSession,
340        attempt: &mut DatasetMaterializationAttempt,
341        resolved: &ResolvedDatafileMaterialization,
342        lifecycle: &mut impl DatafileLifecycle,
343    ) -> Result<PreparedDatafileDataset, AppError> {
344        let request = &resolved.request;
345        let prepared = &attempt.prepared;
346        let source_format = resolved.source_format;
347        let lake_root = session.runtime().lake.data_path.clone();
348        let lake = DuckLakeAdapter::new(&lake_root);
349        let source = lake.restore_dataset_source_observed(
350            lifecycle.observation(),
351            &attempt.scratch,
352            &resolved.source_datafile,
353            source_format,
354        );
355        session.observations().record(
356            ahri_tre_protocol::diagnostics::DiagnosticComponent::Lake,
357            ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
358            source.is_ok(),
359        );
360        let source =
361            source.map_err(|error| infrastructure_error("restore Dataset source", error))?;
362
363        lifecycle.checkpoint(DatafileStage::BeforeLoad, resolved)?;
364        let load_result = lake.load_dataset_table_from_file_reserved_observed(
365            lifecycle.observation(),
366            session.lake_connection(),
367            &mut attempt.authority,
368            &LoadDatasetTableFromFileRequest {
369                study_id: request.study_id,
370                dataset_name: prepared.asset.name.clone(),
371                major: prepared.version.major,
372                minor: prepared.version.minor,
373                patch: prepared.version.patch,
374                source_path: source.path().to_path_buf(),
375                format: source_format,
376                options: resolved.read_options.clone(),
377            },
378        );
379        session.observations().record(
380            ahri_tre_protocol::diagnostics::DiagnosticComponent::Lake,
381            ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
382            load_result.is_ok(),
383        );
384        let loaded = match load_result {
385            Ok(loaded) => loaded,
386            Err(error) => {
387                source
388                    .cleanup()
389                    .map_err(|error| infrastructure_error("clean Dataset source", error))?;
390                return Err(infrastructure_error("load datafile into DuckLake", error));
391            }
392        };
393        lifecycle.checkpoint(DatafileStage::AfterLoad, resolved)?;
394        let field_metadata_result = if metadata_bearing_dataset_file_format(source_format)
395            && request.variable_registrations.is_empty()
396        {
397            Self::field_metadata_proposals(source_format, source.path())
398        } else {
399            Ok((Vec::new(), Vec::new()))
400        };
401        let source_cleanup = source.cleanup();
402        let relation = loaded.relation.clone();
403
404        let transformation = enrich_transformation_source(&request.transformation);
405
406        source_cleanup.map_err(|error| infrastructure_error("clean Dataset source", error))?;
407        let (field_metadata_proposals, discovered_metadata_warnings) = field_metadata_result?;
408        let field_metadata_proposals =
409            (!field_metadata_proposals.is_empty()).then_some(field_metadata_proposals);
410        let profiled_vocabularies = if request.variable_registrations.is_empty()
411            && automatic_fallback_variable_format(source_format)
412        {
413            profile_string_like_column_vocabularies(
414                session.lake_connection(),
415                &relation,
416                &loaded.columns,
417                request.sample_rows,
418                request.vocabulary_threshold,
419                &request.null_strings,
420            )?
421        } else {
422            BTreeMap::new()
423        };
424        let (variable_registrations, metadata_warnings) =
425            if request.variable_registrations.is_empty()
426                && automatic_fallback_variable_format(source_format)
427            {
428                let (variable_registrations, mut metadata_warnings) = self
429                    .automatic_column_shape_variable_registrations(
430                        request.study_id,
431                        request.metadata_domain_id,
432                        source_format,
433                        &loaded.columns,
434                        field_metadata_proposals.as_deref(),
435                        &profiled_vocabularies,
436                    )
437                    .await?;
438                metadata_warnings.extend(discovered_metadata_warnings);
439                metadata_warnings.extend(
440                    self.automatic_metadata_conflict_warnings(
441                        request.study_id,
442                        request.metadata_domain_id,
443                        &variable_registrations,
444                        request.strict_metadata,
445                        request.force_metadata_updates,
446                    )
447                    .await?,
448                );
449                (variable_registrations, metadata_warnings)
450            } else {
451                let metadata_warnings = Self::explicit_registration_metadata_warnings(
452                    source_format,
453                    &loaded.columns,
454                    !request.variable_registrations.is_empty(),
455                );
456                (
457                    request
458                        .variable_registrations
459                        .iter()
460                        .cloned()
461                        .map(Into::into)
462                        .collect(),
463                    metadata_warnings,
464                )
465            };
466
467        Ok(PreparedDatafileDataset {
468            loaded,
469            transformation,
470            variable_registrations,
471            metadata_warnings,
472        })
473    }
474
475    async fn commit_datafile_dataset(
476        &self,
477        session: &DataStoreSession,
478        attempt: &mut DatasetMaterializationAttempt,
479        resolved: &ResolvedDatafileMaterialization,
480        prepared: PreparedDatafileDataset,
481        lifecycle: &mut impl DatafileLifecycle,
482    ) -> Result<DatasetMaterialization, AppError> {
483        let PreparedDatafileDataset {
484            loaded,
485            transformation,
486            variable_registrations,
487            metadata_warnings,
488        } = prepared;
489        let dataset = ahri_tre_types::DatasetRecord {
490            dataset_id: attempt.prepared.version.version_id,
491        };
492        let build_materialization =
493            |(catalog, dataset, lineage, variables),
494             relation: ahri_tre_lake::DatasetTableRelation,
495             metadata_warnings| DatasetMaterialization {
496                catalog,
497                dataset,
498                lineage,
499                variables,
500                lake_relation: relation.qualified_name(),
501                lake_schema: relation.schema_name,
502                lake_table: relation.table_name,
503                row_count: loaded.row_count,
504                column_count: loaded.column_count,
505                metadata_warnings,
506            };
507        let records = self
508            .save_dataset_metadata_with_completion(
509                session,
510                attempt,
511                DatasetAdmissionMetadata {
512                    dataset,
513                    transformation: &transformation,
514                    input_version_ids: &[resolved.source_version.version_id],
515                    registration: DatasetMetadataRegistration {
516                        redcap_dictionary: None,
517                        variable_registrations,
518                        loaded_column_names: Some(&loaded.column_names),
519                        force_metadata_updates: resolved.request.force_metadata_updates,
520                    },
521                },
522                |repository, records| {
523                    let materialization = build_materialization(
524                        records.clone(),
525                        loaded.relation.clone(),
526                        metadata_warnings.clone(),
527                    );
528                    lifecycle.record_completion(repository, &materialization)
529                },
530            )
531            .await?;
532        Ok(build_materialization(
533            records,
534            loaded.relation,
535            metadata_warnings,
536        ))
537    }
538
539    pub(super) async fn resolve_automatic_metadata_domain(
540        &self,
541        study_id: ahri_tre_types::StudyId,
542        metadata_domain_id: Option<ahri_tre_types::DomainId>,
543        format: DatasetFileFormat,
544    ) -> Result<ahri_tre_types::DomainId, AppError> {
545        let format_name = dataset_file_format_name(format);
546        let domain_id = if let Some(domain_id) = metadata_domain_id {
547            self.ensure_study_domain(study_id, domain_id, "automatic variable registration")
548                .await?;
549            domain_id
550        } else {
551            let domains = self
552                .study_domains
553                .list_domains_for_study(study_id)
554                .await
555                .map_err(|error| core_error("list study domains for fallback metadata", error))?;
556            match domains.as_slice() {
557                [domain] => domain.domain_id,
558                [] => {
559                    return Err(AppError::Validation(format!(
560                        "{format_name} automatic variable registration requires study {} to have one linked domain",
561                        study_id.0
562                    )));
563                }
564                _ => {
565                    return Err(AppError::Validation(format!(
566                        "{format_name} automatic variable registration requires study {} to have exactly one linked domain, found {}",
567                        study_id.0,
568                        domains.len()
569                    )));
570                }
571            }
572        };
573
574        Ok(domain_id)
575    }
576}