1use super::*;
2
3impl<'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 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(®istration.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 ®istered_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 ®istration.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 ®istration {
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 ®istration {
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 ®istration.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 ®istration.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}