1use super::*;
4
5pub(super) struct ResolvedEntityBatch {
6 pub request: EnsureEntityInstancesFromDatasetRequest,
7 pub study: ahri_tre_types::StudyRecord,
8 pub entity: ahri_tre_types::EntityRecord,
9 pub version: ahri_tre_types::AssetVersionRecord,
10 pub external_variable_id: ahri_tre_types::VariableId,
11 pub rows: Vec<EntityBatchSourceRow>,
12 pub inputs: Vec<ahri_tre_types::VersionId>,
13}
14pub(super) struct ResolvedRelationBatch {
15 pub request: EnsureRelationInstancesFromDatasetRequest,
16 pub study: ahri_tre_types::StudyRecord,
17 pub relation: ahri_tre_types::EntityRelationRecord,
18 pub version: ahri_tre_types::AssetVersionRecord,
19 pub relation_variable_id: ahri_tre_types::VariableId,
20 pub subject_variable_id: ahri_tre_types::VariableId,
21 pub object_variable_id: ahri_tre_types::VariableId,
22 pub rows: Vec<EntityBatchSourceRow>,
23 pub inputs: Vec<ahri_tre_types::VersionId>,
24}
25impl<'repository> ScopedAppService<'repository> {
26 pub(super) async fn apply_entity_batch(
27 &self,
28 batch: ResolvedEntityBatch,
29 ) -> Result<EntityInstanceDatasetBatchResult, AppError> {
30 let ResolvedEntityBatch {
31 request,
32 study,
33 entity,
34 version,
35 external_variable_id,
36 rows,
37 inputs,
38 } = batch;
39 let dataset_name = request.dataset.clone();
40 let mut rejected_row_count = 0;
41 let mut rejected_rows = Vec::new();
42 let mut external_ids = BTreeMap::<String, Option<String>>::new();
43 for row in rows {
44 let external_id = row
45 .values
46 .get(&request.external_id_variable)
47 .and_then(|value| value.as_deref())
48 .map(str::trim)
49 .filter(|value| !value.is_empty())
50 .map(str::to_string);
51 let Some(external_id) = external_id else {
52 rejected_row_count += 1;
53 if rejected_rows.len() < 10 {
54 rejected_rows.push(EntityInstanceDatasetBatchRejectedRow {
55 row_number: row.row_number,
56 reason: "external ID is null or empty".to_string(),
57 });
58 }
59 continue;
60 };
61 external_ids.entry(external_id).or_insert_with(|| {
62 request
63 .label_template
64 .as_deref()
65 .map(|template| render_label_template(template, &row.values))
66 });
67 }
68
69 let existing_links = self
70 .entities
71 .list_dataset_version_entities(version.version_id)
72 .await
73 .map_err(|error| core_error("list dataset version entity links", error))?;
74 let existing_links_by_instance = existing_links
75 .into_iter()
76 .map(|link| (link.entity_instance_id, link))
77 .collect::<BTreeMap<_, _>>();
78
79 let mut counts = EntityInstanceDatasetBatchCounts {
80 created_instances: 0,
81 reused_instances: 0,
82 created_study_mappings: 0,
83 reused_study_mappings: 0,
84 created_dataset_links: 0,
85 reused_dataset_links: 0,
86 rejected_rows: rejected_row_count,
87 conflicts: 0,
88 };
89 let mut conflicts = Vec::new();
90 let mut plans = Vec::new();
91
92 for (external_id, label) in external_ids {
93 let mapping = self
94 .entities
95 .get_study_entity_instance(study.study_id, entity.entity_id, &external_id)
96 .await
97 .map_err(|error| core_error("lookup study entity instance", error))?;
98 let instance = match mapping.as_ref() {
99 Some(mapping) => {
100 let instance = self
101 .entities
102 .get_entity_instance(mapping.entity_instance_id)
103 .await
104 .map_err(|error| core_error("lookup mapped entity instance", error))?
105 .ok_or_else(|| {
106 AppError::Validation(format!(
107 "study entity mapping references missing entity instance: {}",
108 mapping.entity_instance_id
109 ))
110 })?;
111 if instance.entity_id != entity.entity_id {
112 conflicts.push(EntityInstanceDatasetBatchConflict {
113 external_id: external_id.clone(),
114 message: format!(
115 "study mapping points at entity {}, not {}",
116 instance.entity_id.0, entity.entity_id.0
117 ),
118 });
119 }
120 Some(instance)
121 }
122 None => None,
123 };
124
125 if let Some(instance) = instance.as_ref()
126 && let Some(link) = existing_links_by_instance.get(&instance.instance_id)
127 && link.entity_variable_id != external_variable_id
128 {
129 conflicts.push(EntityInstanceDatasetBatchConflict {
130 external_id: external_id.clone(),
131 message: format!(
132 "dataset link for instance {} uses variable {}, not {}",
133 instance.instance_id, link.entity_variable_id.0, external_variable_id.0
134 ),
135 });
136 }
137
138 match (&mapping, &instance) {
139 (Some(_), Some(instance)) => {
140 counts.reused_instances += 1;
141 counts.reused_study_mappings += 1;
142 if existing_links_by_instance
143 .get(&instance.instance_id)
144 .is_some_and(|link| link.entity_variable_id == external_variable_id)
145 {
146 counts.reused_dataset_links += 1;
147 } else {
148 counts.created_dataset_links += 1;
149 }
150 }
151 _ => {
152 counts.created_instances += 1;
153 counts.created_study_mappings += 1;
154 counts.created_dataset_links += 1;
155 }
156 }
157 plans.push(EntityInstanceDatasetBatchPlan { external_id, label });
158 }
159 counts.conflicts = conflicts.len();
160
161 if !conflicts.is_empty() {
162 if !request.dry_run {
163 return Err(AppError::Validation(format!(
164 "entity batch has {} conflicts; rerun with --dry-run for details",
165 conflicts.len()
166 )));
167 }
168 return Ok(entity_batch_result(
169 &request,
170 &version,
171 None,
172 None,
173 counts,
174 conflicts,
175 rejected_rows,
176 ));
177 }
178
179 if request.dry_run {
180 return Ok(entity_batch_result(
181 &request,
182 &version,
183 None,
184 None,
185 counts,
186 conflicts,
187 rejected_rows,
188 ));
189 }
190
191 let lineage = self
192 .record_transformation(RecordTransformationRequest {
193 transformation: ahri_tre_types::NewTransformationRecord {
194 transformation_type: ahri_tre_types::TransformationType::Entity,
195 description: format!(
196 "ensure {} entity instances from dataset {} variable {}",
197 entity.name.as_str(),
198 dataset_name.as_str(),
199 request.external_id_variable
200 ),
201 repository_url: None,
202 commit_hash: None,
203 file_path: None,
204 date_created: None,
205 created_by: None,
206 },
207 input_version_ids: inputs,
208 output_version_ids: Vec::new(),
209 })
210 .await?;
211 let transformation_id = lineage.transformation.transformation_id;
212
213 let mut links = Vec::new();
214 for plan in plans {
215 let registration = self
216 .ensure_study_entity_instance(EnsureStudyEntityInstanceRequest {
217 study_id: study.study_id,
218 entity_id: entity.entity_id,
219 external_id: plan.external_id,
220 canonical_uuid: None,
221 label: plan.label,
222 note: None,
223 transformation_id,
224 })
225 .await?;
226 if existing_links_by_instance
227 .get(®istration.instance.instance_id)
228 .is_none_or(|link| link.entity_variable_id != external_variable_id)
229 {
230 links.push(ahri_tre_types::DatasetVersionEntityRecord {
231 version_id: version.version_id,
232 entity_instance_id: registration.instance.instance_id,
233 entity_variable_id: external_variable_id,
234 transformation_id,
235 });
236 }
237 }
238 if !links.is_empty() {
239 self.entities
240 .link_dataset_version_entities(&links)
241 .await
242 .map_err(|error| core_error("link dataset version entities", error))?;
243 }
244
245 Ok(entity_batch_result(
246 &request,
247 &version,
248 Some(transformation_id),
249 Some(lineage),
250 counts,
251 conflicts,
252 rejected_rows,
253 ))
254 }
255
256 pub(super) async fn apply_relation_batch(
257 &self,
258 batch: ResolvedRelationBatch,
259 ) -> Result<RelationInstanceDatasetBatchResult, AppError> {
260 let ResolvedRelationBatch {
261 request,
262 study,
263 relation,
264 version,
265 relation_variable_id,
266 subject_variable_id,
267 object_variable_id,
268 rows,
269 inputs,
270 } = batch;
271 let dataset_name = request.dataset.clone();
272 let mut rejected_row_count = 0;
273 let mut rejected_rows = Vec::new();
274 let mut raw_plans = BTreeMap::<String, RelationInstanceDatasetBatchRawPlan>::new();
275 let mut conflicts = Vec::new();
276
277 for row in rows {
278 let relation_external_id =
279 trimmed_row_value(&row.values, &request.relation_external_id_variable);
280 let subject_external_id =
281 trimmed_row_value(&row.values, &request.subject_external_id_variable);
282 let object_external_id =
283 trimmed_row_value(&row.values, &request.object_external_id_variable);
284 let (Some(relation_external_id), Some(subject_external_id), Some(object_external_id)) = (
285 relation_external_id,
286 subject_external_id,
287 object_external_id,
288 ) else {
289 rejected_row_count += 1;
290 if rejected_rows.len() < 10 {
291 rejected_rows.push(RelationInstanceDatasetBatchRejectedRow {
292 row_number: row.row_number,
293 reason: "relation, subject, or object external ID is null or empty"
294 .to_string(),
295 });
296 }
297 continue;
298 };
299 let rejected_before = rejected_row_count;
300 let valid_from = parse_relation_batch_date(
301 &row.values,
302 request.valid_from_variable.as_deref(),
303 row.row_number,
304 &mut rejected_row_count,
305 &mut rejected_rows,
306 )?;
307 if rejected_row_count != rejected_before {
308 continue;
309 }
310 let valid_to = parse_relation_batch_date(
311 &row.values,
312 request.valid_to_variable.as_deref(),
313 row.row_number,
314 &mut rejected_row_count,
315 &mut rejected_rows,
316 )?;
317 if rejected_row_count != rejected_before {
318 continue;
319 }
320 let raw_plan = RelationInstanceDatasetBatchRawPlan {
321 external_id: relation_external_id.clone(),
322 subject_external_id,
323 object_external_id,
324 valid_from,
325 valid_to,
326 };
327 match raw_plans.get(&relation_external_id) {
328 Some(existing) if existing != &raw_plan => {
329 conflicts.push(RelationInstanceDatasetBatchConflict {
330 external_id: relation_external_id,
331 message:
332 "dataset rows disagree on subject, object, or validity for relation external ID"
333 .to_string(),
334 });
335 }
336 Some(_) => {}
337 None => {
338 raw_plans.insert(relation_external_id, raw_plan);
339 }
340 }
341 }
342
343 let existing_links = self
344 .entities
345 .list_dataset_version_relation_instances(version.version_id)
346 .await
347 .map_err(|error| core_error("list dataset version relation links", error))?;
348 let existing_links_by_instance = existing_links
349 .into_iter()
350 .map(|link| (link.relation_instance_id, link))
351 .collect::<BTreeMap<_, _>>();
352
353 let mut counts = RelationInstanceDatasetBatchCounts {
354 created_instances: 0,
355 reused_instances: 0,
356 created_study_mappings: 0,
357 reused_study_mappings: 0,
358 created_dataset_links: 0,
359 reused_dataset_links: 0,
360 missing_endpoints: 0,
361 rejected_rows: rejected_row_count,
362 conflicts: 0,
363 };
364 let mut missing_endpoints = Vec::new();
365 let mut plans = Vec::new();
366
367 for raw_plan in raw_plans.into_values() {
368 let subject_mapping = self
369 .entities
370 .get_study_entity_instance(
371 study.study_id,
372 relation.subject_entity_id,
373 &raw_plan.subject_external_id,
374 )
375 .await
376 .map_err(|error| core_error("lookup subject entity mapping", error))?;
377 let object_mapping = self
378 .entities
379 .get_study_entity_instance(
380 study.study_id,
381 relation.object_entity_id,
382 &raw_plan.object_external_id,
383 )
384 .await
385 .map_err(|error| core_error("lookup object entity mapping", error))?;
386
387 let subject_instance = self
388 .resolve_relation_endpoint_instance(
389 RelationEndpointLookup {
390 mapping: subject_mapping,
391 expected_entity_id: relation.subject_entity_id,
392 relation_external_id: &raw_plan.external_id,
393 endpoint_role: "subject",
394 endpoint_external_id: &raw_plan.subject_external_id,
395 },
396 &mut missing_endpoints,
397 &mut conflicts,
398 )
399 .await?;
400 let object_instance = self
401 .resolve_relation_endpoint_instance(
402 RelationEndpointLookup {
403 mapping: object_mapping,
404 expected_entity_id: relation.object_entity_id,
405 relation_external_id: &raw_plan.external_id,
406 endpoint_role: "object",
407 endpoint_external_id: &raw_plan.object_external_id,
408 },
409 &mut missing_endpoints,
410 &mut conflicts,
411 )
412 .await?;
413 let (Some(subject_instance), Some(object_instance)) =
414 (subject_instance, object_instance)
415 else {
416 continue;
417 };
418
419 let mapping = self
420 .entities
421 .get_study_relation_instance(
422 study.study_id,
423 relation.entity_relation_id,
424 &raw_plan.external_id,
425 )
426 .await
427 .map_err(|error| core_error("lookup study relation instance", error))?;
428 let instance = match mapping.as_ref() {
429 Some(mapping) => {
430 let instance = self
431 .entities
432 .get_relation_instance(mapping.relation_instance_id)
433 .await
434 .map_err(|error| core_error("lookup mapped relation instance", error))?
435 .ok_or_else(|| {
436 AppError::Validation(format!(
437 "study relation mapping references missing relation instance: {}",
438 mapping.relation_instance_id
439 ))
440 })?;
441 if instance.entity_relation_id != relation.entity_relation_id
442 || instance.entity_instance_id_1 != subject_instance.instance_id
443 || instance.entity_instance_id_2 != object_instance.instance_id
444 || instance.valid_from != raw_plan.valid_from
445 || instance.valid_to != raw_plan.valid_to
446 {
447 conflicts.push(RelationInstanceDatasetBatchConflict {
448 external_id: raw_plan.external_id.clone(),
449 message: format!(
450 "study mapping points at relation instance {} with different relation endpoints or validity",
451 instance.relation_instance_id
452 ),
453 });
454 }
455 Some(instance)
456 }
457 None => None,
458 };
459
460 if let Some(instance) = instance.as_ref()
461 && let Some(link) = existing_links_by_instance.get(&instance.relation_instance_id)
462 && (link.subject_variable_id != subject_variable_id
463 || link.object_variable_id != object_variable_id
464 || link.relation_variable_id != Some(relation_variable_id))
465 {
466 conflicts.push(RelationInstanceDatasetBatchConflict {
467 external_id: raw_plan.external_id.clone(),
468 message: format!(
469 "dataset link for relation instance {} uses different subject, object, or relation variables",
470 instance.relation_instance_id
471 ),
472 });
473 }
474
475 match (&mapping, &instance) {
476 (Some(_), Some(instance)) => {
477 counts.reused_instances += 1;
478 counts.reused_study_mappings += 1;
479 if existing_links_by_instance
480 .get(&instance.relation_instance_id)
481 .is_some_and(|link| {
482 link.subject_variable_id == subject_variable_id
483 && link.object_variable_id == object_variable_id
484 && link.relation_variable_id == Some(relation_variable_id)
485 })
486 {
487 counts.reused_dataset_links += 1;
488 } else {
489 counts.created_dataset_links += 1;
490 }
491 }
492 _ => {
493 counts.created_instances += 1;
494 counts.created_study_mappings += 1;
495 counts.created_dataset_links += 1;
496 }
497 }
498 plans.push(RelationInstanceDatasetBatchPlan {
499 external_id: raw_plan.external_id,
500 subject_entity_instance_id: subject_instance.instance_id,
501 object_entity_instance_id: object_instance.instance_id,
502 valid_from: raw_plan.valid_from,
503 valid_to: raw_plan.valid_to,
504 });
505 }
506 counts.missing_endpoints = missing_endpoints.len();
507 counts.conflicts = conflicts.len();
508
509 if !missing_endpoints.is_empty() {
510 if !request.dry_run {
511 return Err(AppError::Validation(format!(
512 "relation batch has {} missing endpoint mappings; rerun with --dry-run for details",
513 missing_endpoints.len()
514 )));
515 }
516 return Ok(relation_batch_result(
517 &request,
518 &version,
519 RelationBatchLineage::none(),
520 counts,
521 conflicts,
522 missing_endpoints,
523 rejected_rows,
524 ));
525 }
526
527 if !conflicts.is_empty() {
528 if !request.dry_run {
529 return Err(AppError::Validation(format!(
530 "relation batch has {} conflicts; rerun with --dry-run for details",
531 conflicts.len()
532 )));
533 }
534 return Ok(relation_batch_result(
535 &request,
536 &version,
537 RelationBatchLineage::none(),
538 counts,
539 conflicts,
540 missing_endpoints,
541 rejected_rows,
542 ));
543 }
544
545 if request.dry_run {
546 return Ok(relation_batch_result(
547 &request,
548 &version,
549 RelationBatchLineage::none(),
550 counts,
551 conflicts,
552 missing_endpoints,
553 rejected_rows,
554 ));
555 }
556
557 let lineage = self
558 .record_transformation(RecordTransformationRequest {
559 transformation: ahri_tre_types::NewTransformationRecord {
560 transformation_type: ahri_tre_types::TransformationType::Entity,
561 description: format!(
562 "ensure {} relation instances from dataset {} variable {}",
563 relation.name.as_str(),
564 dataset_name.as_str(),
565 request.relation_external_id_variable
566 ),
567 repository_url: None,
568 commit_hash: None,
569 file_path: None,
570 date_created: None,
571 created_by: None,
572 },
573 input_version_ids: inputs,
574 output_version_ids: Vec::new(),
575 })
576 .await?;
577 let transformation_id = lineage.transformation.transformation_id;
578
579 let mut links = Vec::new();
580 for plan in plans {
581 let registration = self
582 .ensure_study_relation_instance(EnsureStudyRelationInstanceRequest {
583 study_id: study.study_id,
584 entity_relation_id: relation.entity_relation_id,
585 external_id: plan.external_id,
586 subject_entity_instance_id: plan.subject_entity_instance_id,
587 object_entity_instance_id: plan.object_entity_instance_id,
588 valid_from: plan.valid_from,
589 valid_to: plan.valid_to,
590 note: None,
591 transformation_id,
592 })
593 .await?;
594 if existing_links_by_instance
595 .get(®istration.instance.relation_instance_id)
596 .is_none_or(|link| {
597 link.subject_variable_id != subject_variable_id
598 || link.object_variable_id != object_variable_id
599 || link.relation_variable_id != Some(relation_variable_id)
600 })
601 {
602 links.push(ahri_tre_types::DatasetVersionRelationInstanceRecord {
603 version_id: version.version_id,
604 relation_instance_id: registration.instance.relation_instance_id,
605 subject_variable_id,
606 object_variable_id,
607 relation_variable_id: Some(relation_variable_id),
608 transformation_id,
609 });
610 }
611 }
612 if !links.is_empty() {
613 self.entities
614 .link_dataset_version_relation_instances(&links)
615 .await
616 .map_err(|error| core_error("link dataset version relation instances", error))?;
617 }
618
619 Ok(relation_batch_result(
620 &request,
621 &version,
622 RelationBatchLineage {
623 transformation_id: Some(transformation_id),
624 transformation: Some(lineage),
625 },
626 counts,
627 conflicts,
628 missing_endpoints,
629 rejected_rows,
630 ))
631 }
632
633 async fn resolve_relation_endpoint_instance(
634 &self,
635 lookup: RelationEndpointLookup<'_>,
636 missing_endpoints: &mut Vec<RelationInstanceDatasetBatchMissingEndpoint>,
637 conflicts: &mut Vec<RelationInstanceDatasetBatchConflict>,
638 ) -> Result<Option<ahri_tre_types::EntityInstanceRecord>, AppError> {
639 let Some(mapping) = lookup.mapping else {
640 missing_endpoints.push(RelationInstanceDatasetBatchMissingEndpoint {
641 relation_external_id: lookup.relation_external_id.to_string(),
642 endpoint_role: lookup.endpoint_role.to_string(),
643 endpoint_external_id: lookup.endpoint_external_id.to_string(),
644 entity_id: lookup.expected_entity_id,
645 });
646 return Ok(None);
647 };
648 let instance = self
649 .entities
650 .get_entity_instance(mapping.entity_instance_id)
651 .await
652 .map_err(|error| core_error("lookup mapped endpoint entity instance", error))?;
653 let Some(instance) = instance else {
654 conflicts.push(RelationInstanceDatasetBatchConflict {
655 external_id: lookup.relation_external_id.to_string(),
656 message: format!(
657 "{} mapping references missing entity instance {}",
658 lookup.endpoint_role, mapping.entity_instance_id
659 ),
660 });
661 return Ok(None);
662 };
663 if instance.entity_id != lookup.expected_entity_id {
664 conflicts.push(RelationInstanceDatasetBatchConflict {
665 external_id: lookup.relation_external_id.to_string(),
666 message: format!(
667 "{} mapping points at entity {}, not {}",
668 lookup.endpoint_role, instance.entity_id.0, lookup.expected_entity_id.0
669 ),
670 });
671 return Ok(None);
672 }
673 Ok(Some(instance))
674 }
675
676 pub async fn ensure_study_entity_instance(
677 &self,
678 request: EnsureStudyEntityInstanceRequest,
679 ) -> Result<StudyEntityInstanceRegistration, AppError> {
680 self.studies
681 .get_study_by_id(request.study_id)
682 .await
683 .map_err(|error| core_error("lookup study for entity instance", error))?
684 .ok_or_else(|| {
685 AppError::Validation(format!(
686 "study not found for entity instance: {}",
687 request.study_id.0
688 ))
689 })?;
690
691 let existing_mapping = self
692 .entities
693 .get_study_entity_instance(
694 request.study_id,
695 request.entity_id,
696 request.external_id.as_str(),
697 )
698 .await
699 .map_err(|error| core_error("lookup study entity instance", error))?;
700 if let Some(mapping) = existing_mapping {
701 let instance = self
702 .entities
703 .get_entity_instance(mapping.entity_instance_id)
704 .await
705 .map_err(|error| core_error("lookup mapped entity instance", error))?
706 .ok_or_else(|| {
707 AppError::Validation(format!(
708 "study entity mapping references missing entity instance: {}",
709 mapping.entity_instance_id
710 ))
711 })?;
712 if let Some(canonical_uuid) = request.canonical_uuid
713 && instance.uuid != canonical_uuid
714 {
715 return Err(AppError::Validation(format!(
716 "study entity mapping for entity {} external_id {} has canonical UUID {}, expected {}",
717 request.entity_id.0, request.external_id, instance.uuid, canonical_uuid
718 )));
719 }
720 if let Some(repository) = &self.semantic_repository {
721 repository
722 .retain_instance_contribution(
723 ahri_tre_pgmeta::semantic::SemanticInstance::Entity(instance.instance_id),
724 )
725 .map_err(|error| core_error("retain reused entity origins", error))?;
726 }
727 return Ok(StudyEntityInstanceRegistration { instance, mapping });
728 }
729
730 let instance = if let Some(canonical_uuid) = request.canonical_uuid {
731 match self
732 .entities
733 .get_entity_instance_by_uuid(canonical_uuid)
734 .await
735 .map_err(|error| core_error("lookup entity instance by UUID", error))?
736 {
737 Some(instance) => {
738 if instance.entity_id != request.entity_id {
739 return Err(AppError::Validation(format!(
740 "canonical UUID {} belongs to entity {}, not {}",
741 canonical_uuid, instance.entity_id.0, request.entity_id.0
742 )));
743 }
744 instance
745 }
746 None => self
747 .entities
748 .save_entity_instance(&ahri_tre_types::EntityInstanceWriteRecord {
749 instance_id: None,
750 uuid: Some(canonical_uuid),
751 entity_id: request.entity_id,
752 label: request.label.clone(),
753 note: request.note.clone(),
754 transformation_id: request.transformation_id,
755 })
756 .await
757 .map_err(|error| core_error("save entity instance", error))?,
758 }
759 } else {
760 self.entities
761 .create_entity_instance(&ahri_tre_types::NewEntityInstanceRecord {
762 entity_id: request.entity_id,
763 label: request.label.clone(),
764 note: request.note.clone(),
765 transformation_id: request.transformation_id,
766 })
767 .await
768 .map_err(|error| core_error("create entity instance", error))?
769 };
770
771 let mapping = self
772 .entities
773 .save_study_entity_instance(&ahri_tre_types::StudyEntityInstanceWriteRecord {
774 study_id: request.study_id,
775 entity_instance_id: instance.instance_id,
776 entity_id: request.entity_id,
777 external_id: request.external_id,
778 transformation_id: request.transformation_id,
779 })
780 .await
781 .map_err(|error| core_error("save study entity instance", error))?;
782 Ok(StudyEntityInstanceRegistration { instance, mapping })
783 }
784
785 pub async fn ensure_study_relation_instance(
786 &self,
787 request: EnsureStudyRelationInstanceRequest,
788 ) -> Result<StudyRelationInstanceRegistration, AppError> {
789 let existing_mapping = self
790 .entities
791 .get_study_relation_instance(
792 request.study_id,
793 request.entity_relation_id,
794 request.external_id.as_str(),
795 )
796 .await
797 .map_err(|error| core_error("lookup study relation instance", error))?;
798 if let Some(mapping) = existing_mapping {
799 let instance = self
800 .entities
801 .get_relation_instance(mapping.relation_instance_id)
802 .await
803 .map_err(|error| core_error("lookup mapped relation instance", error))?
804 .ok_or_else(|| {
805 AppError::Validation(format!(
806 "study relation mapping references missing relation instance: {}",
807 mapping.relation_instance_id
808 ))
809 })?;
810 if let Some(repository) = &self.semantic_repository {
811 repository
812 .retain_instance_contribution(
813 ahri_tre_pgmeta::semantic::SemanticInstance::Relation(
814 instance.relation_instance_id,
815 ),
816 )
817 .map_err(|error| core_error("retain reused relation origins", error))?;
818 }
819 return Ok(StudyRelationInstanceRegistration { instance, mapping });
820 }
821
822 self.entities
823 .get_entity_instance(request.subject_entity_instance_id)
824 .await
825 .map_err(|error| core_error("lookup relation subject instance", error))?
826 .ok_or_else(|| {
827 AppError::Validation(format!(
828 "relation subject entity instance not found: {}",
829 request.subject_entity_instance_id
830 ))
831 })?;
832 self.entities
833 .get_entity_instance(request.object_entity_instance_id)
834 .await
835 .map_err(|error| core_error("lookup relation object instance", error))?
836 .ok_or_else(|| {
837 AppError::Validation(format!(
838 "relation object entity instance not found: {}",
839 request.object_entity_instance_id
840 ))
841 })?;
842
843 let instance = self
844 .entities
845 .create_relation_instance(&ahri_tre_types::NewRelationInstanceRecord {
846 entity_relation_id: request.entity_relation_id,
847 entity_instance_id_1: request.subject_entity_instance_id,
848 entity_instance_id_2: request.object_entity_instance_id,
849 valid_from: request.valid_from,
850 valid_to: request.valid_to,
851 note: request.note.clone(),
852 transformation_id: request.transformation_id,
853 })
854 .await
855 .map_err(|error| core_error("save relation instance", error))?;
856 let mapping = self
857 .entities
858 .save_study_relation_instance(&ahri_tre_types::StudyRelationInstanceWriteRecord {
859 study_id: request.study_id,
860 relation_instance_id: instance.relation_instance_id,
861 entity_relation_id: request.entity_relation_id,
862 external_id: request.external_id,
863 transformation_id: request.transformation_id,
864 })
865 .await
866 .map_err(|error| core_error("save study relation instance", error))?;
867 Ok(StudyRelationInstanceRegistration { instance, mapping })
868 }
869
870 pub async fn record_transformation(
871 &self,
872 request: RecordTransformationRequest,
873 ) -> Result<ahri_tre_types::TransformationLineageRecord, AppError> {
874 let policy = DefaultProvenancePolicy;
875 let transformation = enrich_transformation_source(&request.transformation);
876 ahri_tre_core::record_transformation_provenance(
877 self.transformations.as_ref(),
878 &policy,
879 &transformation,
880 &request.input_version_ids,
881 &request.output_version_ids,
882 )
883 .await
884 .map_err(|error| core_error("record transformation provenance", error))
885 }
886}