Skip to main content

ahri_tre_app/service/
semantic_batch.rs

1//! Metadata-only batch application. Inputs and selectors are resolved before
2//! admission; this code runs with an adapter-owned mutation transaction.
3use 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(&registration.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(&registration.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}