1use super::*;
3use ahri_tre_observability::{CorrelationContext, FailureCategory, Outcome, Stage};
4
5pub(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
43pub(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 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 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}