Skip to main content

ahri_tre_app/service/
metadata_admission.rs

1use super::*;
2
3// Reuses the same registration rules with owned Session repositories or with
4// repositories borrowed for the duration of one adapter-owned admission.
5impl<'repository> ScopedAppService<'repository> {
6    pub fn new(repositories: ScopedAppRepositories<'repository>) -> Self {
7        Self {
8            semantic_repository: None,
9            datastore_id: None,
10            registrations: repositories.registrations,
11            domains: repositories.domains,
12            studies: repositories.studies,
13            study_domains: repositories.study_domains,
14            study_governance: repositories.study_governance,
15            assets: repositories.assets,
16            variables: repositories.variables,
17            vocabularies: repositories.vocabularies,
18            entities: repositories.entities,
19            transformations: repositories.transformations,
20            tags: repositories.tags,
21        }
22    }
23
24    pub async fn get_asset(
25        &self,
26        request: GetAssetRequest,
27    ) -> Result<AssetVersionCatalog, AppError> {
28        let asset = match (request.asset_id, request.study_id, request.asset_name) {
29            (Some(asset_id), None, None) => self
30                .assets
31                .get_asset_by_id(asset_id)
32                .await
33                .map_err(|error| core_error("lookup asset by id", error))?
34                .ok_or_else(|| AppError::Validation(format!("asset not found: {}", asset_id.0)))?,
35            (None, Some(study_id), Some(asset_name)) => self
36                .assets
37                .get_asset_by_name(study_id, asset_name.as_str())
38                .await
39                .map_err(|error| core_error("lookup asset by study and name", error))?
40                .ok_or_else(|| {
41                    AppError::Validation(format!(
42                        "asset not found in study {}: {}",
43                        study_id.0,
44                        asset_name.as_str()
45                    ))
46                })?,
47            (None, None, None) => {
48                return Err(AppError::Validation(
49                    "asset lookup requires an asset_id or study_id plus asset_name".to_string(),
50                ));
51            }
52            _ => {
53                return Err(AppError::Validation(
54                    "asset lookup accepts either asset_id or study_id plus asset_name".to_string(),
55                ));
56            }
57        };
58
59        self.asset_version_catalog(asset).await
60    }
61
62    pub(super) async fn asset_version_catalog(
63        &self,
64        asset: ahri_tre_types::AssetRecord,
65    ) -> Result<AssetVersionCatalog, AppError> {
66        let versions = self
67            .assets
68            .list_asset_versions(asset.asset_id)
69            .await
70            .map_err(|error| core_error("list asset versions", error))?;
71        let latest_version = self
72            .latest_catalogue_version(&asset, &versions)
73            .await
74            .map_err(|error| core_error("resolve latest asset version", error))?;
75
76        Ok(AssetVersionCatalog {
77            asset,
78            versions,
79            latest_version,
80        })
81    }
82
83    pub(super) async fn latest_catalogue_version(
84        &self,
85        asset: &ahri_tre_types::AssetRecord,
86        versions: &[ahri_tre_types::AssetVersionRecord],
87    ) -> Result<Option<ahri_tre_types::AssetVersionRecord>, CoreError> {
88        // Final withdrawal retains authorized history but has no usable latest.
89        // Other missing/multiple markers remain catalogue inconsistencies.
90        if asset.asset_type == ahri_tre_types::AssetType::Dataset
91            && !versions.is_empty()
92            && versions
93                .iter()
94                .all(|version| version.is_latest != Some(true))
95        {
96            let withdrawals = self
97                .assets
98                .list_dataset_version_withdrawals(asset.asset_id)
99                .await?;
100            if versions.iter().all(|version| {
101                withdrawals
102                    .iter()
103                    .any(|withdrawal| withdrawal.dataset_id == version.version_id)
104            }) {
105                return Ok(None);
106            }
107        }
108        ahri_tre_core::latest_asset_version(versions)
109    }
110
111    pub(super) async fn write_dataset_metadata_in_admission(
112        &self,
113        prepared: &PreparedDatasetVersion,
114        dataset: ahri_tre_types::DatasetRecord,
115        transformation: &ahri_tre_types::NewTransformationRecord,
116        input_version_ids: &[ahri_tre_types::VersionId],
117        registration: DatasetMetadataRegistration<'_>,
118    ) -> Result<
119        (
120            AssetVersionCatalog,
121            ahri_tre_types::DatasetRecord,
122            ahri_tre_types::TransformationLineageRecord,
123            Vec<RegisteredDatasetVariable>,
124        ),
125        AppError,
126    > {
127        if let Some(source) = prepared.ordinary_upload_source {
128            if input_version_ids != [source] {
129                return Err(AppError::Validation(
130                    "Uploaded source provenance changed".into(),
131                ));
132            }
133            self.semantic_repository
134                .as_ref()
135                .ok_or_else(|| AppError::Infrastructure("Upload admission is unavailable".into()))?
136                .validate_uploaded_table_source(
137                    source,
138                    prepared.asset.study_id,
139                    prepared.classification.risk,
140                )
141                .map_err(|error| core_error("validate uploaded source classification", error))?;
142        } else if !input_version_ids.is_empty() {
143            self.assets
144                .validate_dataset_derivation(
145                    input_version_ids,
146                    prepared.classification.risk,
147                    prepared.inherit_risk,
148                )
149                .await
150                .map_err(|error| core_error("validate derivation classification", error))?;
151        }
152        if !prepared.asset_persisted {
153            self.assets
154                .save_asset(&prepared.asset)
155                .await
156                .map_err(|error| core_error("save dataset asset", error))?;
157        }
158        self.attach_tags(
159            TagAttachmentTarget::Asset(prepared.asset.asset_id),
160            prepared.asset_tags.clone(),
161        )
162        .await?;
163
164        if !prepared.version_persisted {
165            self.assets
166                .admit_asset_version(&prepared.version, &prepared.classification)
167                .await
168                .map_err(|error| core_error("save dataset asset version", error))?;
169        }
170        self.attach_tags(
171            TagAttachmentTarget::AssetVersion(prepared.version.version_id),
172            prepared.version_tags.clone(),
173        )
174        .await?;
175        let dataset = self
176            .assets
177            .save_dataset(&dataset)
178            .await
179            .map_err(|error| core_error("save dataset record", error))?;
180        if let Some(repository) = &self.semantic_repository {
181            let mut sources = input_version_ids.to_vec();
182            sources.push(prepared.version.version_id);
183            sources.sort_by_key(|id| id.0);
184            sources.dedup();
185            repository
186                .begin_semantic_derivation(&sources)
187                .map_err(|error| core_error("retain Dataset dictionary sources", error))?;
188        }
189        if let Some((domain_id, forms)) = registration.redcap_dictionary {
190            self.register_redcap_dictionary_variables(
191                domain_id,
192                forms,
193                registration.force_metadata_updates,
194            )
195            .await?;
196        }
197        let registered_variables = self
198            .register_dataset_variables_internal(
199                dataset.dataset_id,
200                first_registration_domain(&registration.variable_registrations),
201                registration.variable_registrations,
202                registration.loaded_column_names,
203                registration.force_metadata_updates,
204            )
205            .await?;
206
207        if let Some(repository) = &self.semantic_repository {
208            for variable in &registered_variables {
209                if let Some(vocabulary) = &variable.vocabulary {
210                    repository
211                        .retain_vocabulary_contribution(vocabulary.vocabulary_id)
212                        .map_err(|error| core_error("retain reused dictionary sources", error))?;
213                }
214            }
215        }
216
217        let policy = DefaultProvenancePolicy;
218        let transformation = enrich_transformation_source(transformation);
219        let outputs = [prepared.version.version_id];
220        let lineage = if input_version_ids.is_empty() {
221            ahri_tre_core::record_ingest_provenance(
222                self.transformations.as_ref(),
223                &policy,
224                &transformation,
225                &outputs,
226            )
227            .await
228            .map_err(|error| core_error("record dataset ingest provenance", error))?
229        } else {
230            ahri_tre_core::transform_assets(
231                self.transformations.as_ref(),
232                &policy,
233                &transformation,
234                input_version_ids,
235                &outputs,
236            )
237            .await
238            .map_err(|error| core_error("record dataset transform provenance", error))?
239        };
240        self.assets
241            .mark_asset_version_latest(prepared.asset.asset_id, prepared.version.version_id)
242            .await
243            .map_err(|error| core_error("mark completed dataset version latest", error))?;
244        let catalog = self
245            .get_asset(GetAssetRequest {
246                asset_id: Some(prepared.asset.asset_id),
247                study_id: None,
248                asset_name: None,
249            })
250            .await?;
251
252        Ok((catalog, dataset, lineage, registered_variables))
253    }
254
255    async fn register_redcap_dictionary_variables(
256        &self,
257        domain_id: ahri_tre_types::DomainId,
258        forms: &[RedcapForm],
259        force_metadata_updates: bool,
260    ) -> Result<Vec<String>, AppError> {
261        let mut warnings = Vec::new();
262        for form in forms {
263            for field in &form.fields {
264                let (registration, mut field_warnings) =
265                    redcap_variable_registration(domain_id, field)?;
266                warnings.append(&mut field_warnings);
267                let DatasetVariableRegistrationInput::New(registration) = registration else {
268                    unreachable!("REDCap registration always creates new metadata inputs");
269                };
270                let vocabulary = self
271                    .register_or_validate_new_vocabulary(
272                        domain_id,
273                        registration.variable.vocabulary_id,
274                        registration.vocabulary,
275                        &registration.vocabulary_items,
276                        force_metadata_updates,
277                    )
278                    .await?;
279                let mut variable = registration.variable;
280                if let Some(vocabulary) = vocabulary {
281                    variable.vocabulary_id = Some(vocabulary.vocabulary_id);
282                }
283                self.register_or_validate_new_variable(
284                    domain_id,
285                    &variable,
286                    force_metadata_updates,
287                )
288                .await?;
289            }
290        }
291        Ok(warnings)
292    }
293
294    pub(super) async fn register_dataset_variables_internal(
295        &self,
296        dataset_id: ahri_tre_types::VersionId,
297        request_domain_id: Option<ahri_tre_types::DomainId>,
298        variable_registrations: Vec<DatasetVariableRegistrationInput>,
299        loaded_column_names: Option<&[String]>,
300        force_metadata_updates: bool,
301    ) -> Result<Vec<RegisteredDatasetVariable>, AppError> {
302        if variable_registrations.is_empty() {
303            return Ok(Vec::new());
304        }
305        let domain_id = request_domain_id.ok_or_else(|| {
306            AppError::Validation("variable registration requires a domain_id".to_string())
307        })?;
308        let loaded_column_names = loaded_column_names.map(|columns| {
309            columns
310                .iter()
311                .map(|column| column.to_ascii_lowercase())
312                .collect::<std::collections::BTreeSet<_>>()
313        });
314
315        let mut registered = Vec::with_capacity(variable_registrations.len());
316        for registration in variable_registrations {
317            let (variable_name, validation_column_name, variable_domain_id, variable_key_role) =
318                match &registration {
319                    DatasetVariableRegistrationInput::Explicit(registration) => (
320                        registration.variable.name.as_str(),
321                        registration.variable.name.as_str(),
322                        registration.variable.domain_id,
323                        registration.variable.key_role,
324                    ),
325                    DatasetVariableRegistrationInput::New(registration) => (
326                        registration.variable.name.as_str(),
327                        registration
328                            .source_column_name
329                            .as_deref()
330                            .unwrap_or_else(|| registration.variable.name.as_str()),
331                        registration.variable.domain_id,
332                        registration.variable.key_role,
333                    ),
334                };
335            if variable_domain_id != domain_id {
336                return Err(AppError::Validation(format!(
337                    "variable {} belongs to domain {}, not {}",
338                    variable_name, variable_domain_id.0, domain_id.0
339                )));
340            }
341            if let Some(column_names) = &loaded_column_names
342                && !column_names.contains(&validation_column_name.to_ascii_lowercase())
343            {
344                return Err(AppError::Validation(format!(
345                    "dataset column not found for variable {}",
346                    variable_name
347                )));
348            }
349            validate_metadata_key_role("variable key_role", variable_key_role)?;
350            let row_role = match &registration {
351                DatasetVariableRegistrationInput::Explicit(registration) => {
352                    registration.row_role.unwrap_or(variable_key_role)
353                }
354                DatasetVariableRegistrationInput::New(registration) => {
355                    registration.row_role.unwrap_or(variable_key_role)
356                }
357            };
358            validate_metadata_key_role("dataset row_role", row_role)?;
359
360            let (variable, vocabulary) = match registration {
361                DatasetVariableRegistrationInput::Explicit(registration) => {
362                    let vocabulary = self
363                        .register_or_validate_vocabulary(
364                            domain_id,
365                            registration.variable.vocabulary_id,
366                            registration.vocabulary,
367                            &registration.vocabulary_items,
368                            force_metadata_updates,
369                        )
370                        .await?;
371                    let mut variable_record = registration.variable;
372                    if let Some(vocabulary) = &vocabulary {
373                        variable_record.vocabulary_id = Some(vocabulary.vocabulary_id);
374                    }
375                    let variable = self
376                        .register_or_validate_variable(
377                            domain_id,
378                            &variable_record,
379                            force_metadata_updates,
380                        )
381                        .await?;
382                    (variable, vocabulary)
383                }
384                DatasetVariableRegistrationInput::New(registration) => {
385                    let vocabulary = self
386                        .register_or_validate_new_vocabulary(
387                            domain_id,
388                            registration.variable.vocabulary_id,
389                            registration.vocabulary,
390                            &registration.vocabulary_items,
391                            force_metadata_updates,
392                        )
393                        .await?;
394                    let mut variable_record = registration.variable;
395                    if let Some(vocabulary) = &vocabulary {
396                        variable_record.vocabulary_id = Some(vocabulary.vocabulary_id);
397                    }
398                    let variable = self
399                        .register_or_validate_new_variable(
400                            domain_id,
401                            &variable_record,
402                            force_metadata_updates,
403                        )
404                        .await?;
405                    (variable, vocabulary)
406                }
407            };
408            let vocabulary_items = match &vocabulary {
409                Some(vocabulary) => self
410                    .vocabularies
411                    .list_vocabulary_items(vocabulary.vocabulary_id)
412                    .await
413                    .map_err(|error| core_error("list registered vocabulary items", error))?,
414                None => Vec::new(),
415            };
416            let link = self
417                .variables
418                .link_dataset_variable(dataset_id, variable.variable_id, row_role)
419                .await
420                .map_err(|error| core_error("link dataset variable", error))?;
421            registered.push(RegisteredDatasetVariable {
422                variable,
423                vocabulary,
424                vocabulary_items,
425                link,
426            });
427        }
428        Ok(registered)
429    }
430
431    pub(super) async fn resolve_value_type(
432        &self,
433        value_type_id: Option<ahri_tre_types::ValueTypeId>,
434        value_type_name: Option<&str>,
435    ) -> Result<ahri_tre_types::ValueTypeRecord, AppError> {
436        let value_types = self
437            .variables
438            .list_value_types()
439            .await
440            .map_err(|error| core_error("list value types", error))?;
441        let by_id = value_type_id.map(|id| {
442            value_types
443                .iter()
444                .find(|candidate| candidate.value_type_id == id)
445                .cloned()
446                .ok_or_else(|| AppError::Validation(format!("value_type_id {} not found", id.0)))
447        });
448        let by_name = value_type_name.map(|name| {
449            value_types
450                .iter()
451                .find(|candidate| candidate.value_type.eq_ignore_ascii_case(name))
452                .cloned()
453                .ok_or_else(|| AppError::Validation(format!("value_type {name} not found")))
454        });
455        match (by_id, by_name) {
456            (Some(by_id), Some(by_name)) => {
457                let by_id = by_id?;
458                let by_name = by_name?;
459                if by_id.value_type_id != by_name.value_type_id {
460                    return Err(AppError::Validation(format!(
461                        "value_type_id {} is {}, not {}",
462                        by_id.value_type_id.0, by_id.value_type, by_name.value_type
463                    )));
464                }
465                Ok(by_id)
466            }
467            (Some(by_id), None) => by_id,
468            (None, Some(by_name)) => by_name,
469            (None, None) => Err(AppError::Validation(
470                "variable registration requires value_type_id or value_type".to_string(),
471            )),
472        }
473    }
474
475    pub(super) async fn value_type_for_id(
476        &self,
477        value_type_id: ahri_tre_types::ValueTypeId,
478    ) -> Result<ahri_tre_types::ValueTypeRecord, AppError> {
479        self.resolve_value_type(Some(value_type_id), None).await
480    }
481
482    pub(super) async fn register_or_validate_vocabulary(
483        &self,
484        domain_id: ahri_tre_types::DomainId,
485        variable_vocabulary_id: Option<ahri_tre_types::VocabularyId>,
486        vocabulary: Option<ahri_tre_types::VocabularyRecord>,
487        vocabulary_items: &[ahri_tre_types::VocabularyItemRecord],
488        force_metadata_updates: bool,
489    ) -> Result<Option<ahri_tre_types::VocabularyRecord>, AppError> {
490        let Some(vocabulary_id) = vocabulary
491            .as_ref()
492            .map(|v| v.vocabulary_id)
493            .or(variable_vocabulary_id)
494        else {
495            if !vocabulary_items.is_empty() {
496                return Err(AppError::Validation(
497                    "vocabulary items require a variable vocabulary_id or vocabulary record"
498                        .to_string(),
499                ));
500            }
501            return Ok(None);
502        };
503
504        if let Some(variable_vocabulary_id) = variable_vocabulary_id
505            && variable_vocabulary_id != vocabulary_id
506        {
507            return Err(AppError::Validation(format!(
508                "variable vocabulary_id {} does not match vocabulary record {}",
509                variable_vocabulary_id.0, vocabulary_id.0
510            )));
511        }
512
513        let vocabulary = match vocabulary {
514            Some(vocabulary) => {
515                if vocabulary.domain_id != domain_id {
516                    return Err(AppError::Validation(format!(
517                        "vocabulary {} belongs to domain {}, not {}",
518                        vocabulary.name.as_str(),
519                        vocabulary.domain_id.0,
520                        domain_id.0
521                    )));
522                }
523                let domain_vocabularies = self
524                    .vocabularies
525                    .list_vocabularies(domain_id)
526                    .await
527                    .map_err(|error| core_error("list domain vocabularies", error))?;
528                if let Some(existing) = domain_vocabularies.into_iter().find(|candidate| {
529                    candidate.vocabulary_id == vocabulary.vocabulary_id
530                        || candidate.name == vocabulary.name
531                }) {
532                    if existing.name != vocabulary.name
533                        && existing.vocabulary_id == vocabulary.vocabulary_id
534                    {
535                        return Err(AppError::Validation(format!(
536                            "vocabulary {} is named {}, not {}",
537                            vocabulary.vocabulary_id.0,
538                            existing.name.as_str(),
539                            vocabulary.name.as_str()
540                        )));
541                    }
542                    if force_metadata_updates && existing.description != vocabulary.description {
543                        let mut updated = vocabulary.clone();
544                        updated.vocabulary_id = existing.vocabulary_id;
545                        self.vocabularies
546                            .update_vocabulary(&updated)
547                            .await
548                            .map_err(|error| core_error("update variable vocabulary", error))?
549                    } else {
550                        existing
551                    }
552                } else {
553                    self.vocabularies
554                        .save_vocabulary(&vocabulary)
555                        .await
556                        .map_err(|error| core_error("save variable vocabulary", error))?
557                }
558            }
559            None => self
560                .vocabularies
561                .list_vocabularies(domain_id)
562                .await
563                .map_err(|error| core_error("list domain vocabularies", error))?
564                .into_iter()
565                .find(|candidate| candidate.vocabulary_id == vocabulary_id)
566                .ok_or_else(|| {
567                    AppError::Validation(format!(
568                        "vocabulary {} not found in domain {}",
569                        vocabulary_id.0, domain_id.0
570                    ))
571                })?,
572        };
573
574        let vocabulary_id = vocabulary.vocabulary_id;
575        if vocabulary.domain_id != domain_id {
576            return Err(AppError::Validation(format!(
577                "vocabulary {} belongs to domain {}, not {}",
578                vocabulary.vocabulary_id.0, vocabulary.domain_id.0, domain_id.0
579            )));
580        }
581        if !vocabulary_items.is_empty() {
582            let normalized_items: Vec<_> = vocabulary_items
583                .iter()
584                .map(|item| {
585                    let mut item = item.clone();
586                    item.vocabulary_id = vocabulary_id;
587                    item
588                })
589                .collect();
590            let existing_items = self
591                .vocabularies
592                .list_vocabulary_items(vocabulary_id)
593                .await
594                .map_err(|error| core_error("list vocabulary items", error))?;
595            let mut seen_items: Vec<_> = existing_items
596                .iter()
597                .map(SeenVocabularyItem::from_record)
598                .collect();
599            let mut items_to_save = Vec::new();
600            let mut items_to_update = Vec::new();
601            for item in &normalized_items {
602                let candidate = vocabulary_item_candidate(
603                    Some(item.vocabulary_item_id),
604                    item.value,
605                    item.code.as_str(),
606                    item.description.as_deref(),
607                );
608                match classify_vocabulary_item(&seen_items, &candidate) {
609                    VocabularyItemReconciliation::ExactMatch => {}
610                    VocabularyItemReconciliation::Conflict(existing_index) => {
611                        let existing = &seen_items[existing_index];
612                        if force_metadata_updates {
613                            let mut updated = item.clone();
614                            updated.vocabulary_item_id = existing
615                                .vocabulary_item_id
616                                .expect("persisted vocabulary item should have an id");
617                            seen_items[existing_index] = SeenVocabularyItem::from_record(&updated);
618                            items_to_update.push(updated);
619                        } else {
620                            return Err(vocabulary_item_conflict_error(
621                                Some(item.vocabulary_item_id),
622                                existing,
623                                &candidate,
624                            ));
625                        }
626                    }
627                    VocabularyItemReconciliation::Missing => {
628                        seen_items.push(SeenVocabularyItem::from_record(item));
629                        items_to_save.push(item.clone());
630                    }
631                }
632            }
633            if !items_to_save.is_empty() {
634                self.vocabularies
635                    .save_vocabulary_items(&items_to_save)
636                    .await
637                    .map_err(|error| core_error("save vocabulary items", error))?;
638            }
639            if !items_to_update.is_empty() {
640                self.vocabularies
641                    .update_vocabulary_items(&items_to_update)
642                    .await
643                    .map_err(|error| core_error("update vocabulary items", error))?;
644            }
645        }
646        Ok(Some(vocabulary))
647    }
648
649    pub(super) async fn register_or_validate_new_vocabulary(
650        &self,
651        domain_id: ahri_tre_types::DomainId,
652        variable_vocabulary_id: Option<ahri_tre_types::VocabularyId>,
653        vocabulary: Option<ahri_tre_types::NewVocabularyRecord>,
654        vocabulary_items: &[ahri_tre_types::NewVocabularyItemRecord],
655        force_metadata_updates: bool,
656    ) -> Result<Option<ahri_tre_types::VocabularyRecord>, AppError> {
657        let Some(vocabulary) = vocabulary else {
658            let Some(vocabulary_id) = variable_vocabulary_id else {
659                if !vocabulary_items.is_empty() {
660                    return Err(AppError::Validation(
661                        "vocabulary items require a variable vocabulary_id or vocabulary record"
662                            .to_string(),
663                    ));
664                }
665                return Ok(None);
666            };
667            return self
668                .vocabularies
669                .list_vocabularies(domain_id)
670                .await
671                .map_err(|error| core_error("list domain vocabularies", error))?
672                .into_iter()
673                .find(|candidate| candidate.vocabulary_id == vocabulary_id)
674                .ok_or_else(|| {
675                    AppError::Validation(format!(
676                        "vocabulary {} not found in domain {}",
677                        vocabulary_id.0, domain_id.0
678                    ))
679                })
680                .map(Some);
681        };
682
683        if vocabulary.domain_id != domain_id {
684            return Err(AppError::Validation(format!(
685                "vocabulary {} belongs to domain {}, not {}",
686                vocabulary.name.as_str(),
687                vocabulary.domain_id.0,
688                domain_id.0
689            )));
690        }
691        let domain_vocabularies = self
692            .vocabularies
693            .list_vocabularies(domain_id)
694            .await
695            .map_err(|error| core_error("list domain vocabularies", error))?;
696        let vocabulary = if let Some(existing) = domain_vocabularies
697            .into_iter()
698            .find(|candidate| candidate.name == vocabulary.name)
699        {
700            if force_metadata_updates && existing.description != vocabulary.description {
701                let updated = ahri_tre_types::VocabularyRecord {
702                    vocabulary_id: existing.vocabulary_id,
703                    domain_id: vocabulary.domain_id,
704                    name: vocabulary.name.clone(),
705                    description: vocabulary.description.clone(),
706                };
707                self.vocabularies
708                    .update_vocabulary(&updated)
709                    .await
710                    .map_err(|error| core_error("update variable vocabulary", error))?
711            } else {
712                existing
713            }
714        } else {
715            self.vocabularies
716                .create_vocabulary(&vocabulary)
717                .await
718                .map_err(|error| core_error("create variable vocabulary", error))?
719        };
720
721        if let Some(variable_vocabulary_id) = variable_vocabulary_id
722            && variable_vocabulary_id != vocabulary.vocabulary_id
723        {
724            return Err(AppError::Validation(format!(
725                "variable vocabulary_id {} does not match vocabulary record {}",
726                variable_vocabulary_id.0, vocabulary.vocabulary_id.0
727            )));
728        }
729
730        self.register_new_vocabulary_items(
731            vocabulary.vocabulary_id,
732            vocabulary_items,
733            force_metadata_updates,
734        )
735        .await?;
736        Ok(Some(vocabulary))
737    }
738
739    pub(super) async fn register_new_vocabulary_items(
740        &self,
741        vocabulary_id: ahri_tre_types::VocabularyId,
742        vocabulary_items: &[ahri_tre_types::NewVocabularyItemRecord],
743        force_metadata_updates: bool,
744    ) -> Result<(), AppError> {
745        if vocabulary_items.is_empty() {
746            return Ok(());
747        }
748        let existing_items = self
749            .vocabularies
750            .list_vocabulary_items(vocabulary_id)
751            .await
752            .map_err(|error| core_error("list vocabulary items", error))?;
753        let mut seen_items: Vec<_> = existing_items
754            .iter()
755            .map(SeenVocabularyItem::from_record)
756            .collect();
757        let mut items_to_create = Vec::new();
758        let mut items_to_update = Vec::new();
759        for item in vocabulary_items {
760            let candidate = vocabulary_item_candidate(
761                None,
762                item.value,
763                item.code.as_str(),
764                item.description.as_deref(),
765            );
766            match classify_vocabulary_item(&seen_items, &candidate) {
767                VocabularyItemReconciliation::ExactMatch => {}
768                VocabularyItemReconciliation::Conflict(existing_index) => {
769                    let existing = &seen_items[existing_index];
770                    if force_metadata_updates {
771                        if let Some(vocabulary_item_id) = existing.vocabulary_item_id {
772                            let updated = build_vocabulary_item_record(
773                                vocabulary_item_id,
774                                vocabulary_id,
775                                item.value,
776                                item.code.clone(),
777                                item.description.clone(),
778                            );
779                            seen_items[existing_index] = SeenVocabularyItem::from_record(&updated);
780                            items_to_update.push(updated);
781                        } else if let Some(pending_index) = existing.pending_create_index {
782                            items_to_create[pending_index] = item.clone();
783                            seen_items[existing_index] =
784                                SeenVocabularyItem::from_pending_create(pending_index, item);
785                        }
786                    } else {
787                        return Err(vocabulary_item_conflict_error(
788                            existing.vocabulary_item_id,
789                            existing,
790                            &candidate,
791                        ));
792                    }
793                }
794                VocabularyItemReconciliation::Missing => {
795                    items_to_create.push(item.clone());
796                    let pending_index = items_to_create.len() - 1;
797                    seen_items.push(SeenVocabularyItem::from_pending_create(pending_index, item));
798                }
799            }
800        }
801        if !items_to_create.is_empty() {
802            self.vocabularies
803                .create_vocabulary_items(vocabulary_id, &items_to_create)
804                .await
805                .map_err(|error| core_error("create vocabulary items", error))?;
806        }
807        if !items_to_update.is_empty() {
808            self.vocabularies
809                .update_vocabulary_items(&items_to_update)
810                .await
811                .map_err(|error| core_error("update vocabulary items", error))?;
812        }
813        Ok(())
814    }
815
816    pub(super) async fn register_or_validate_variable(
817        &self,
818        domain_id: ahri_tre_types::DomainId,
819        variable: &ahri_tre_types::VariableRecord,
820        force_metadata_updates: bool,
821    ) -> Result<ahri_tre_types::VariableRecord, AppError> {
822        let variables = self
823            .variables
824            .list_domain_variables(domain_id)
825            .await
826            .map_err(|error| core_error("list domain variables", error))?;
827        if let Some(existing) = variables
828            .iter()
829            .find(|candidate| candidate.variable_id == variable.variable_id)
830        {
831            if existing.name != variable.name {
832                return Err(AppError::Validation(format!(
833                    "variable {} is named {}, not {}",
834                    variable.variable_id.0,
835                    existing.name.as_str(),
836                    variable.name.as_str()
837                )));
838            }
839            if existing.vocabulary_id != variable.vocabulary_id && !force_metadata_updates {
840                return Err(AppError::Validation(format!(
841                    "variable {} vocabulary_id does not match existing variable; set force_metadata_updates to update metadata",
842                    variable.variable_id.0
843                )));
844            }
845            return self
846                .merge_or_update_variable(existing, variable, force_metadata_updates)
847                .await;
848        }
849        if let Some(existing) = variables
850            .iter()
851            .find(|candidate| candidate.name == variable.name)
852        {
853            if existing.value_type_id != variable.value_type_id && !force_metadata_updates {
854                return Err(AppError::Validation(format!(
855                    "variable {} already exists in domain {} with incompatible value type; set force_metadata_updates to update metadata",
856                    variable.name.as_str(),
857                    domain_id.0
858                )));
859            }
860            return self
861                .merge_or_update_variable(existing, variable, force_metadata_updates)
862                .await;
863        }
864        self.variables
865            .save_variable(variable)
866            .await
867            .map_err(|error| core_error("save dataset variable", error))
868    }
869
870    pub(super) async fn register_or_validate_new_variable(
871        &self,
872        domain_id: ahri_tre_types::DomainId,
873        variable: &ahri_tre_types::NewVariableRecord,
874        force_metadata_updates: bool,
875    ) -> Result<ahri_tre_types::VariableRecord, AppError> {
876        let variables = self
877            .variables
878            .list_domain_variables(domain_id)
879            .await
880            .map_err(|error| core_error("list domain variables", error))?;
881        if let Some(existing) = variables
882            .iter()
883            .find(|candidate| candidate.name == variable.name)
884        {
885            let incoming = ahri_tre_types::VariableRecord {
886                variable_id: existing.variable_id,
887                domain_id: variable.domain_id,
888                name: variable.name.clone(),
889                value_type_id: variable.value_type_id,
890                value_format: variable.value_format.clone(),
891                vocabulary_id: variable.vocabulary_id,
892                key_role: variable.key_role,
893                description: variable.description.clone(),
894                note: variable.note.clone(),
895                ontology_namespace: variable.ontology_namespace.clone(),
896                ontology_class: variable.ontology_class.clone(),
897            };
898            if existing.value_type_id != incoming.value_type_id && !force_metadata_updates {
899                return Err(AppError::Validation(format!(
900                    "variable {} already exists in domain {} with incompatible value type; set force_metadata_updates to update metadata",
901                    variable.name.as_str(),
902                    domain_id.0
903                )));
904            }
905            return self
906                .merge_or_update_variable(existing, &incoming, force_metadata_updates)
907                .await;
908        }
909        self.variables
910            .create_variable(variable)
911            .await
912            .map_err(|error| core_error("create dataset variable", error))
913    }
914
915    pub(super) async fn merge_or_update_variable(
916        &self,
917        existing: &ahri_tre_types::VariableRecord,
918        incoming: &ahri_tre_types::VariableRecord,
919        force_metadata_updates: bool,
920    ) -> Result<ahri_tre_types::VariableRecord, AppError> {
921        let mut merged = incoming.clone();
922        merged.variable_id = existing.variable_id;
923        let value_type_changed = existing.value_type_id != incoming.value_type_id;
924        let value_format_changed = existing.value_format != incoming.value_format;
925        let storage_affecting_change = value_format_changed
926            || self
927                .storage_affecting_value_type_change(existing.value_type_id, incoming.value_type_id)
928                .await?;
929        if value_type_changed || value_format_changed {
930            if !force_metadata_updates {
931                return Err(AppError::Validation(format!(
932                    "variable {} value type or format differs; set force_metadata_updates to update metadata",
933                    existing.name.as_str()
934                )));
935            }
936            if storage_affecting_change {
937                let links = self
938                    .variables
939                    .list_dataset_variables_for_variable(existing.variable_id)
940                    .await
941                    .map_err(|error| core_error("list variable dataset links", error))?;
942                if !links.is_empty() {
943                    return Err(AppError::Validation(format!(
944                        "variable {} is linked to {} dataset(s); changing value type or format requires a storage-aware migration workflow",
945                        existing.name.as_str(),
946                        links.len()
947                    )));
948                }
949            }
950        }
951        if !force_metadata_updates {
952            if existing.vocabulary_id != incoming.vocabulary_id {
953                return Err(AppError::Validation(format!(
954                    "variable {} vocabulary differs; set force_metadata_updates to update metadata",
955                    existing.name.as_str()
956                )));
957            }
958            return Ok(existing.clone());
959        }
960        self.variables
961            .update_variable_metadata(&merged)
962            .await
963            .map_err(|error| core_error("update dataset variable", error))
964    }
965
966    pub(super) async fn storage_affecting_value_type_change(
967        &self,
968        existing: ahri_tre_types::ValueTypeId,
969        incoming: ahri_tre_types::ValueTypeId,
970    ) -> Result<bool, AppError> {
971        if existing == incoming {
972            return Ok(false);
973        }
974        let existing = self.value_type_for_id(existing).await?;
975        let incoming = self.value_type_for_id(incoming).await?;
976        Ok(!integer_to_enumeration_value_type_names(
977            &existing.value_type,
978            &incoming.value_type,
979        ))
980    }
981
982    pub(super) async fn attach_tags(
983        &self,
984        target: TagAttachmentTarget,
985        tags: Vec<String>,
986    ) -> Result<(), AppError> {
987        for name in tags {
988            let tag = self
989                .tags
990                .create_tag(&ahri_tre_types::NewTagRecord { name })
991                .await
992                .map_err(|error| core_error("create tag", error))?;
993            self.tags
994                .attach_tag(tag.tag_id, target.clone())
995                .await
996                .map_err(|error| core_error("attach tag", error))?;
997        }
998        Ok(())
999    }
1000}
1001
1002impl ScopedAppService<'_> {
1003    pub async fn register_vocabulary(
1004        &self,
1005        request: RegisterVocabularyRequest,
1006    ) -> Result<RegisteredVocabulary, AppError> {
1007        self.domains
1008            .get_domain_by_id(request.domain_id)
1009            .await
1010            .map_err(|error| core_error("lookup vocabulary domain", error))?
1011            .ok_or_else(|| {
1012                AppError::Validation(format!(
1013                    "vocabulary domain not found: {}",
1014                    request.domain_id.0
1015                ))
1016            })?;
1017
1018        let domain_vocabularies = self
1019            .vocabularies
1020            .list_vocabularies(request.domain_id)
1021            .await
1022            .map_err(|error| core_error("list domain vocabularies", error))?;
1023        let existing = domain_vocabularies.into_iter().find(|candidate| {
1024            request
1025                .vocabulary_id
1026                .is_some_and(|id| candidate.vocabulary_id == id)
1027                || candidate.name == request.name
1028        });
1029        let vocabulary = match (existing, request.vocabulary_id) {
1030            (Some(existing), _) => {
1031                if existing.name != request.name {
1032                    return Err(AppError::Validation(format!(
1033                        "vocabulary {} is named {}, not {}",
1034                        existing.vocabulary_id.0,
1035                        existing.name.as_str(),
1036                        request.name.as_str()
1037                    )));
1038                }
1039                if request.force_metadata_updates && existing.description != request.description {
1040                    self.vocabularies
1041                        .update_vocabulary(&ahri_tre_types::VocabularyRecord {
1042                            vocabulary_id: existing.vocabulary_id,
1043                            domain_id: request.domain_id,
1044                            name: request.name,
1045                            description: request.description,
1046                        })
1047                        .await
1048                        .map_err(|error| core_error("update vocabulary", error))?
1049                } else {
1050                    existing
1051                }
1052            }
1053            (None, Some(vocabulary_id)) => self
1054                .vocabularies
1055                .save_vocabulary(&ahri_tre_types::VocabularyRecord {
1056                    vocabulary_id,
1057                    domain_id: request.domain_id,
1058                    name: request.name,
1059                    description: request.description,
1060                })
1061                .await
1062                .map_err(|error| core_error("save vocabulary", error))?,
1063            (None, None) => self
1064                .vocabularies
1065                .create_vocabulary(&ahri_tre_types::NewVocabularyRecord {
1066                    domain_id: request.domain_id,
1067                    name: request.name,
1068                    description: request.description,
1069                })
1070                .await
1071                .map_err(|error| core_error("create vocabulary", error))?,
1072        };
1073        self.register_vocabulary_item_requests(
1074            vocabulary.vocabulary_id,
1075            request.items,
1076            request.force_metadata_updates,
1077        )
1078        .await?;
1079        let items = self
1080            .vocabularies
1081            .list_vocabulary_items(vocabulary.vocabulary_id)
1082            .await
1083            .map_err(|error| core_error("list registered vocabulary items", error))?;
1084        Ok(RegisteredVocabulary { vocabulary, items })
1085    }
1086
1087    async fn register_vocabulary_item_requests(
1088        &self,
1089        vocabulary_id: ahri_tre_types::VocabularyId,
1090        vocabulary_items: Vec<RegisterVocabularyItemRequest>,
1091        force_metadata_updates: bool,
1092    ) -> Result<(), AppError> {
1093        if vocabulary_items.is_empty() {
1094            return Ok(());
1095        }
1096        let existing_items = self
1097            .vocabularies
1098            .list_vocabulary_items(vocabulary_id)
1099            .await
1100            .map_err(|error| core_error("list vocabulary items", error))?;
1101        let mut seen_items: Vec<_> = existing_items
1102            .iter()
1103            .map(SeenVocabularyItem::from_record)
1104            .collect();
1105        let mut items_to_create = Vec::new();
1106        let mut items_to_save = Vec::new();
1107        let mut items_to_update = Vec::new();
1108        for item in vocabulary_items {
1109            let candidate = vocabulary_item_candidate(
1110                item.vocabulary_item_id,
1111                item.value,
1112                item.code.as_str(),
1113                item.description.as_deref(),
1114            );
1115            match classify_vocabulary_item(&seen_items, &candidate) {
1116                VocabularyItemReconciliation::ExactMatch => {}
1117                VocabularyItemReconciliation::Conflict(existing_index) => {
1118                    let existing = &seen_items[existing_index];
1119                    if force_metadata_updates {
1120                        if let Some(vocabulary_item_id) = existing.vocabulary_item_id {
1121                            let updated = build_vocabulary_item_record(
1122                                vocabulary_item_id,
1123                                vocabulary_id,
1124                                item.value,
1125                                item.code,
1126                                item.description,
1127                            );
1128                            seen_items[existing_index] = SeenVocabularyItem::from_record(&updated);
1129                            items_to_update.push(updated);
1130                        } else if let Some(pending_index) = existing.pending_create_index {
1131                            let created = ahri_tre_types::NewVocabularyItemRecord {
1132                                value: item.value,
1133                                code: item.code,
1134                                description: item.description,
1135                            };
1136                            items_to_create[pending_index] = created;
1137                            seen_items[existing_index] = SeenVocabularyItem::from_pending_create(
1138                                pending_index,
1139                                &items_to_create[pending_index],
1140                            );
1141                        }
1142                    } else {
1143                        return Err(vocabulary_item_conflict_error(
1144                            existing.vocabulary_item_id,
1145                            existing,
1146                            &candidate,
1147                        ));
1148                    }
1149                }
1150                VocabularyItemReconciliation::Missing => {
1151                    if let Some(vocabulary_item_id) = item.vocabulary_item_id {
1152                        let saved = build_vocabulary_item_record(
1153                            vocabulary_item_id,
1154                            vocabulary_id,
1155                            item.value,
1156                            item.code,
1157                            item.description,
1158                        );
1159                        seen_items.push(SeenVocabularyItem::from_record(&saved));
1160                        items_to_save.push(saved);
1161                    } else {
1162                        let created = ahri_tre_types::NewVocabularyItemRecord {
1163                            value: item.value,
1164                            code: item.code,
1165                            description: item.description,
1166                        };
1167                        items_to_create.push(created);
1168                        let pending_index = items_to_create.len() - 1;
1169                        seen_items.push(SeenVocabularyItem::from_pending_create(
1170                            pending_index,
1171                            &items_to_create[pending_index],
1172                        ));
1173                    }
1174                }
1175            }
1176        }
1177        if !items_to_create.is_empty() {
1178            self.vocabularies
1179                .create_vocabulary_items(vocabulary_id, &items_to_create)
1180                .await
1181                .map_err(|error| core_error("create vocabulary items", error))?;
1182        }
1183        if !items_to_save.is_empty() {
1184            self.vocabularies
1185                .save_vocabulary_items(&items_to_save)
1186                .await
1187                .map_err(|error| core_error("save vocabulary items", error))?;
1188        }
1189        if !items_to_update.is_empty() {
1190            self.vocabularies
1191                .update_vocabulary_items(&items_to_update)
1192                .await
1193                .map_err(|error| core_error("update vocabulary items", error))?;
1194        }
1195        Ok(())
1196    }
1197}