Skip to main content

ahri_tre_app/service/
semantic_batch_admission.rs

1//! Resolve batch inputs before admission, then apply only retained data inside
2//! the metadata adapter's transaction. The returned capability owns delivery.
3use super::*;
4use crate::disclosure::{
5    DisclosureQueryInput, DisclosureQueryRequest, DisclosureRequest, PreparedSemanticResponse,
6};
7use ahri_tre_pgmeta::semantic::{SemanticResolution, SemanticScope, SemanticState};
8use ahri_tre_protocol::{
9    asset, dataset, dictionary,
10    model::{self, SemanticInput},
11    refs::{ObjectKind, ObjectRef},
12    study,
13};
14use ahri_tre_types::{AssetType, DisclosureBudgets, VersionId};
15use futures::FutureExt;
16
17struct BatchSource {
18    study: ahri_tre_types::StudyRecord,
19    asset: ahri_tre_types::AssetRecord,
20    version: ahri_tre_types::AssetVersionRecord,
21    variables: Vec<ahri_tre_types::VariableRecord>,
22    latest: bool,
23}
24impl AppService {
25    async fn resolve_batch_source(
26        &self,
27        study: SemanticInput<study::StudySelector>,
28        dataset: SemanticInput<dataset::DatasetSelector>,
29        version: Option<&str>,
30    ) -> Result<BatchSource, AppError> {
31        let study_selector = match study {
32            SemanticInput::Name(name) => study::StudySelector::Name { domain: None, name },
33            SemanticInput::Reference(study) => study::StudySelector::Id { study },
34            SemanticInput::Selector(selector) => selector,
35        };
36        let study = self
37            .resolve_study_selector(
38                self.catalogue_study_selector(study_selector.clone())?,
39                "batch Study",
40            )
41            .await?;
42        let dataset = match dataset {
43            SemanticInput::Name(name) => dataset::DatasetSelector::Name {
44                study: study_selector,
45                name,
46            },
47            SemanticInput::Reference(dataset) => dataset::DatasetSelector::Id {
48                study: Some(study_selector),
49                dataset,
50            },
51            SemanticInput::Selector(selector) => selector,
52        };
53        let selector = match dataset {
54            dataset::DatasetSelector::Id { dataset, study } => asset::AssetSelector::Id {
55                asset: dataset,
56                study,
57                asset_type: Some(AssetType::Dataset),
58            },
59            dataset::DatasetSelector::Name { study, name } => asset::AssetSelector::Name {
60                study,
61                name,
62                asset_type: Some(AssetType::Dataset),
63            },
64        };
65        let (parent, catalog, pinned) = self.resolve_catalogue_asset(selector).await?;
66        if parent.study_id != study.study_id {
67            return Err(AppError::Conflict("Batch Study selectors disagree".into()));
68        }
69        let latest = pinned.is_none() && version.is_none_or(|v| v.eq_ignore_ascii_case("latest"));
70        let catalog = self.complete_dataset_catalog(catalog.asset).await?;
71        let selected = self.catalogue_version(&catalog, pinned, version)?;
72        let links = self
73            .variables
74            .list_dataset_variables(selected.version_id)
75            .await
76            .map_err(|e| core_error("read batch variable links", e))?;
77        let mut variables = Vec::new();
78        for link in links {
79            variables.push(
80                self.variables
81                    .get_variable(link.variable_id)
82                    .await
83                    .map_err(|e| core_error("read batch variable", e))?
84                    .ok_or_else(batch_unavailable)?,
85            );
86        }
87        Ok(BatchSource {
88            study,
89            asset: catalog.asset,
90            version: selected,
91            variables,
92            latest,
93        })
94    }
95    async fn batch_variable(
96        &self,
97        source: &BatchSource,
98        input: SemanticInput<dictionary::VariableSelector>,
99    ) -> Result<ahri_tre_types::VariableRecord, AppError> {
100        let variable = match input {
101            SemanticInput::Name(name) => {
102                let mut matches = source
103                    .variables
104                    .iter()
105                    .filter(|v| v.name.as_str() == name.as_str());
106                let value = matches.next().cloned().ok_or_else(|| {
107                    AppError::NotFound(
108                        "Batch variable is not present in the selected Dataset".into(),
109                    )
110                })?;
111                if matches.next().is_some() {
112                    return Err(AppError::Conflict(
113                        "Batch variable name is ambiguous".into(),
114                    ));
115                }
116                value
117            }
118            input => {
119                let selector = match input {
120                    SemanticInput::Reference(variable) => dictionary::VariableSelector::Id {
121                        variable,
122                        domain: None,
123                    },
124                    SemanticInput::Selector(selector) => selector,
125                    SemanticInput::Name(_) => unreachable!(),
126                };
127                self.resolve_variable_selector_optional(
128                    self.catalogue_variable_selector(selector).await?,
129                )
130                .await?
131                .ok_or_else(|| AppError::NotFound("Batch variable was not found".into()))?
132            }
133        };
134        if !source
135            .variables
136            .iter()
137            .any(|v| v.variable_id == variable.variable_id)
138        {
139            return Err(AppError::Conflict(
140                "Batch variable does not belong to the selected Dataset".into(),
141            ));
142        }
143        Ok(variable)
144    }
145    pub(super) async fn prepare_entity_batch<'s>(
146        session: &'s mut DataStoreSession,
147        request: model::EntityInstanceEnsureFromDatasetRequest,
148    ) -> Result<PreparedSemanticResponse<'s>, AppError> {
149        let service = Self::from_datastore_session(session)?;
150        let resolution = service
151            .semantic_repository
152            .as_ref()
153            .ok_or_else(batch_unavailable)?
154            .begin_semantic_resolution()
155            .map_err(|_| batch_unavailable())?;
156        let datastore = service.catalogue_datastore_id()?;
157        let source = service
158            .resolve_batch_source(request.study, request.dataset, request.version.as_deref())
159            .await?;
160        let entity = service
161            .resolve_entity_selector_optional(
162                service.catalogue_entity_selector(request.entity).await?,
163            )
164            .await?
165            .ok_or_else(|| AppError::NotFound("Entity was not found".into()))?;
166        let domain = service
167            .domains
168            .get_domain_by_id(entity.domain_id)
169            .await
170            .map_err(|e| core_error("read batch Domain", e))?
171            .ok_or_else(batch_unavailable)?;
172        let variable = service
173            .batch_variable(&source, request.external_id_variable)
174            .await?;
175        if request
176            .label_template
177            .as_ref()
178            .is_some_and(|template| template.len() > 4096)
179        {
180            return Err(AppError::Validation(
181                "Label template exceeds the batch limit".into(),
182            ));
183        }
184        let mut columns = label_template_columns(request.label_template.as_deref())?;
185        for column in &columns {
186            service
187                .batch_variable(
188                    &source,
189                    SemanticInput::Name(
190                        ahri_tre_protocol::PublicName::new(column.clone())
191                            .map_err(|_| batch_unavailable())?,
192                    ),
193                )
194                .await?;
195        }
196        columns.push(variable.name.as_str().to_owned());
197        columns.sort();
198        columns.dedup();
199        let budgets = batch_budgets(session);
200        let state = service
201            .capture_batch_state(
202                session,
203                resolution,
204                SemanticScope::Entity {
205                    study_id: source.study.study_id,
206                    entity_id: entity.entity_id,
207                    version_id: source.version.version_id,
208                },
209                budgets,
210            )
211            .await?;
212        let query = batch_query(datastore, &source, &columns, budgets);
213        let latest = if source.latest {
214            vec![source.version.version_id.0]
215        } else {
216            Vec::new()
217        };
218        let mut disclosure = Self::admit_disclosure_with_semantic_state(
219            session,
220            DisclosureRequest::Query(query),
221            Some(Arc::clone(&state)),
222            latest,
223        )
224        .await?;
225        let inputs = disclosure
226            .snapshot()
227            .inputs
228            .iter()
229            .map(|i| VersionId(i.version_id))
230            .collect();
231        let refs = model::SemanticBatchReferences {
232            study: crate::projections::object_ref(
233                datastore,
234                ObjectKind::Study,
235                source.study.study_id.0,
236            ),
237            definition: integer_ref(datastore, ObjectKind::Entity, "entity", entity.entity_id.0),
238            dataset: crate::projections::object_ref(
239                datastore,
240                ObjectKind::Asset,
241                source.asset.asset_id.0,
242            ),
243            version: crate::projections::object_ref(
244                datastore,
245                ObjectKind::AssetVersion,
246                source.version.version_id.0,
247            ),
248            variables: vec![integer_ref(
249                datastore,
250                ObjectKind::Variable,
251                "variable",
252                variable.variable_id.0,
253            )],
254        };
255        let batch = ResolvedEntityBatch {
256            request: EnsureEntityInstancesFromDatasetRequest {
257                study: source.study.name.as_str().into(),
258                domain: domain.name.as_str().into(),
259                entity: entity.name.as_str().into(),
260                dataset: source.asset.name.as_str().into(),
261                version: request.version,
262                external_id_variable: variable.name.as_str().into(),
263                label_template: request.label_template,
264                dry_run: request.intent.dry_run,
265            },
266            study: source.study,
267            entity,
268            version: source.version,
269            external_variable_id: variable.variable_id,
270            rows: Vec::new(),
271            inputs,
272        };
273        let read_scope = state.scope().clone();
274        disclosure.prepare_semantic_batch(&state, move |repository, rows| {
275            let service = semantic_transaction_service(repository);
276            let mut batch = batch;
277            batch.rows = batch_source_rows(rows);
278            let count = batch.rows.len() as u64;
279            let result = service
280                .apply_entity_batch(batch)
281                .now_or_never()
282                .ok_or_else(batch_core_unavailable)?
283                .map_err(|_| batch_core_unavailable())?;
284            let instances = service
285                .batch_instances(datastore, &read_scope)
286                .now_or_never()
287                .ok_or_else(batch_core_unavailable)??;
288            let response = model::EntityInstanceEnsureFromDatasetResponse {
289                summary: model::EntityInstanceEnsureFromDatasetSummary {
290                    schema_version: result.schema_version,
291                    operation: result.operation,
292                    dry_run: result.dry_run,
293                    study: result.study,
294                    dataset: result.dataset,
295                    requested_version: result.requested_version,
296                    resolved_version: result.resolved_version,
297                    external_id_variable: result.external_id_variable,
298                    label_template: result.label_template,
299                    entity: model::EntitySelectorSummary {
300                        domain: result.domain,
301                        name: result.entity,
302                    },
303                    transformation: batch_transformation(datastore, result.transformation),
304                    references: Some(refs),
305                },
306                counts: model::EntityInstanceEnsureFromDatasetCounts {
307                    created_instances: result.counts.created_instances,
308                    reused_instances: result.counts.reused_instances,
309                    created_study_mappings: result.counts.created_study_mappings,
310                    reused_study_mappings: result.counts.reused_study_mappings,
311                    created_dataset_links: result.counts.created_dataset_links,
312                    reused_dataset_links: result.counts.reused_dataset_links,
313                    rejected_rows: result.counts.rejected_rows,
314                    conflicts: result.counts.conflicts,
315                },
316                conflicts: result
317                    .conflicts
318                    .into_iter()
319                    .map(|row| model::SemanticWorkflowConflict {
320                        external_id: row.external_id,
321                        message: row.message,
322                    })
323                    .collect(),
324                rejected_rows: result
325                    .rejected_rows
326                    .into_iter()
327                    .map(|row| model::SemanticWorkflowRejectedRow {
328                        row_number: row.row_number,
329                        reason: row.reason,
330                    })
331                    .collect(),
332                warnings: Vec::new(),
333                instances,
334            };
335            Ok((
336                serde_json::to_vec(&response).map_err(|_| batch_core_unavailable())?,
337                count.max(response.instances.len() as u64),
338            ))
339        })?;
340        Ok(PreparedSemanticResponse::Content(Box::new(disclosure)))
341    }
342    pub(super) async fn prepare_relation_batch<'s>(
343        session: &'s mut DataStoreSession,
344        request: model::RelationInstanceEnsureFromDatasetRequest,
345    ) -> Result<PreparedSemanticResponse<'s>, AppError> {
346        let service = Self::from_datastore_session(session)?;
347        let resolution = service
348            .semantic_repository
349            .as_ref()
350            .ok_or_else(batch_unavailable)?
351            .begin_semantic_resolution()
352            .map_err(|_| batch_unavailable())?;
353        let datastore = service.catalogue_datastore_id()?;
354        let source = service
355            .resolve_batch_source(request.study, request.dataset, request.version.as_deref())
356            .await?;
357        let relation = service
358            .resolve_entity_relation_selector_optional(
359                service
360                    .catalogue_relation_selector(request.relation)
361                    .await?,
362            )
363            .await?
364            .ok_or_else(|| AppError::NotFound("Relation was not found".into()))?;
365        let domain = service
366            .domains
367            .get_domain_by_id(relation.domain_id)
368            .await
369            .map_err(|e| core_error("read batch Domain", e))?
370            .ok_or_else(batch_unavailable)?;
371        let relation_variable = service
372            .batch_variable(&source, request.relation_external_id_variable)
373            .await?;
374        let subject_variable = service
375            .batch_variable(&source, request.subject_external_id_variable)
376            .await?;
377        let object_variable = service
378            .batch_variable(&source, request.object_external_id_variable)
379            .await?;
380        let from = match request.valid_from_variable {
381            Some(input) => Some(service.batch_variable(&source, input).await?),
382            None => None,
383        };
384        let to = match request.valid_to_variable {
385            Some(input) => Some(service.batch_variable(&source, input).await?),
386            None => None,
387        };
388        let variables = std::iter::once(&relation_variable)
389            .chain([&subject_variable, &object_variable])
390            .chain(from.iter())
391            .chain(to.iter())
392            .collect::<Vec<_>>();
393        let mut columns = variables
394            .iter()
395            .map(|v| v.name.as_str().to_owned())
396            .collect::<Vec<_>>();
397        columns.sort();
398        columns.dedup();
399        let refs = model::SemanticBatchReferences {
400            study: crate::projections::object_ref(
401                datastore,
402                ObjectKind::Study,
403                source.study.study_id.0,
404            ),
405            definition: integer_ref(
406                datastore,
407                ObjectKind::Relation,
408                "relation",
409                relation.entity_relation_id.0,
410            ),
411            dataset: crate::projections::object_ref(
412                datastore,
413                ObjectKind::Asset,
414                source.asset.asset_id.0,
415            ),
416            version: crate::projections::object_ref(
417                datastore,
418                ObjectKind::AssetVersion,
419                source.version.version_id.0,
420            ),
421            variables: variables
422                .iter()
423                .map(|v| integer_ref(datastore, ObjectKind::Variable, "variable", v.variable_id.0))
424                .collect(),
425        };
426        let budgets = batch_budgets(session);
427        let state = service
428            .capture_batch_state(
429                session,
430                resolution,
431                SemanticScope::Relation {
432                    study_id: source.study.study_id,
433                    relation_id: relation.entity_relation_id,
434                    version_id: source.version.version_id,
435                },
436                budgets,
437            )
438            .await?;
439        let query = batch_query(datastore, &source, &columns, budgets);
440        let latest = if source.latest {
441            vec![source.version.version_id.0]
442        } else {
443            Vec::new()
444        };
445        let mut disclosure = Self::admit_disclosure_with_semantic_state(
446            session,
447            DisclosureRequest::Query(query),
448            Some(Arc::clone(&state)),
449            latest,
450        )
451        .await?;
452        let inputs = disclosure
453            .snapshot()
454            .inputs
455            .iter()
456            .map(|i| VersionId(i.version_id))
457            .collect();
458        let batch = ResolvedRelationBatch {
459            request: EnsureRelationInstancesFromDatasetRequest {
460                study: source.study.name.as_str().into(),
461                domain: domain.name.as_str().into(),
462                relation: relation.name.as_str().into(),
463                dataset: source.asset.name.as_str().into(),
464                version: request.version,
465                relation_external_id_variable: relation_variable.name.as_str().into(),
466                subject_external_id_variable: subject_variable.name.as_str().into(),
467                object_external_id_variable: object_variable.name.as_str().into(),
468                valid_from_variable: from.map(|v| v.name.as_str().to_owned()),
469                valid_to_variable: to.map(|v| v.name.as_str().to_owned()),
470                dry_run: request.intent.dry_run,
471            },
472            study: source.study,
473            relation,
474            version: source.version,
475            relation_variable_id: relation_variable.variable_id,
476            subject_variable_id: subject_variable.variable_id,
477            object_variable_id: object_variable.variable_id,
478            rows: Vec::new(),
479            inputs,
480        };
481        let read_scope = state.scope().clone();
482        disclosure.prepare_semantic_batch(&state, move |repository, rows| {
483            let service = semantic_transaction_service(repository);
484            let mut batch = batch;
485            batch.rows = batch_source_rows(rows);
486            let count = batch.rows.len() as u64;
487            let result = service
488                .apply_relation_batch(batch)
489                .now_or_never()
490                .ok_or_else(batch_core_unavailable)?
491                .map_err(|_| batch_core_unavailable())?;
492            let instances = service
493                .batch_instances(datastore, &read_scope)
494                .now_or_never()
495                .ok_or_else(batch_core_unavailable)??;
496            let response = model::RelationInstanceEnsureFromDatasetResponse {
497                summary: model::RelationInstanceEnsureFromDatasetSummary {
498                    schema_version: result.schema_version,
499                    operation: result.operation,
500                    dry_run: result.dry_run,
501                    study: result.study,
502                    dataset: result.dataset,
503                    requested_version: result.requested_version,
504                    resolved_version: result.resolved_version,
505                    relation_external_id_variable: result.relation_external_id_variable,
506                    subject_external_id_variable: result.subject_external_id_variable,
507                    object_external_id_variable: result.object_external_id_variable,
508                    valid_from_variable: result.valid_from_variable,
509                    valid_to_variable: result.valid_to_variable,
510                    relation: model::RelationSelectorSummary {
511                        domain: result.domain,
512                        name: result.relation,
513                    },
514                    transformation: batch_transformation(datastore, result.transformation),
515                    references: Some(refs),
516                },
517                counts: model::RelationInstanceEnsureFromDatasetCounts {
518                    created_instances: result.counts.created_instances,
519                    reused_instances: result.counts.reused_instances,
520                    created_study_mappings: result.counts.created_study_mappings,
521                    reused_study_mappings: result.counts.reused_study_mappings,
522                    created_dataset_links: result.counts.created_dataset_links,
523                    reused_dataset_links: result.counts.reused_dataset_links,
524                    rejected_rows: result.counts.rejected_rows,
525                    conflicts: result.counts.conflicts,
526                    missing_endpoints: result.counts.missing_endpoints,
527                },
528                conflicts: result
529                    .conflicts
530                    .into_iter()
531                    .map(|row| model::SemanticWorkflowConflict {
532                        external_id: row.external_id,
533                        message: row.message,
534                    })
535                    .collect(),
536                rejected_rows: result
537                    .rejected_rows
538                    .into_iter()
539                    .map(|row| model::SemanticWorkflowRejectedRow {
540                        row_number: row.row_number,
541                        reason: row.reason,
542                    })
543                    .collect(),
544                missing_endpoints: result
545                    .missing_endpoints
546                    .into_iter()
547                    .map(|row| model::RelationInstanceMissingEndpointSummary {
548                        relation_external_id: row.relation_external_id,
549                        endpoint_role: row.endpoint_role,
550                        endpoint_external_id: row.endpoint_external_id,
551                        entity: integer_ref(
552                            datastore,
553                            ObjectKind::Entity,
554                            "entity",
555                            row.entity_id.0,
556                        ),
557                    })
558                    .collect(),
559                warnings: Vec::new(),
560                instances,
561            };
562            Ok((
563                serde_json::to_vec(&response).map_err(|_| batch_core_unavailable())?,
564                count.max(response.instances.len() as u64),
565            ))
566        })?;
567        Ok(PreparedSemanticResponse::Content(Box::new(disclosure)))
568    }
569    /// Ticket 15 consumes this retained capability. All scope resolution and
570    /// content projection remain in the application; transports receive bytes
571    /// only after the usual admission and delivery evidence acknowledgements.
572    pub async fn admit_semantic_readback<'s>(
573        session: &'s mut DataStoreSession,
574        request: crate::disclosure::SemanticReadbackRequest,
575    ) -> Result<crate::disclosure::SessionDisclosure<'s>, AppError> {
576        let service = Self::from_datastore_session(session)?;
577        let resolution = service
578            .semantic_repository
579            .as_ref()
580            .ok_or_else(batch_unavailable)?
581            .begin_semantic_resolution()
582            .map_err(|_| batch_unavailable())?;
583        let datastore = service.catalogue_datastore_id()?;
584        let source = service
585            .resolve_batch_source(
586                SemanticInput::Selector(request.study),
587                SemanticInput::Selector(request.dataset),
588                None,
589            )
590            .await?;
591        let scope = match request.target {
592            crate::disclosure::SemanticReadbackTarget::Entity(selector) => SemanticScope::Entity {
593                study_id: source.study.study_id,
594                entity_id: service
595                    .resolve_entity_selector_optional(
596                        service.catalogue_entity_selector(selector).await?,
597                    )
598                    .await?
599                    .ok_or_else(batch_unavailable)?
600                    .entity_id,
601                version_id: source.version.version_id,
602            },
603            crate::disclosure::SemanticReadbackTarget::Relation(selector) => {
604                SemanticScope::Relation {
605                    study_id: source.study.study_id,
606                    relation_id: service
607                        .resolve_entity_relation_selector_optional(
608                            service.catalogue_relation_selector(selector).await?,
609                        )
610                        .await?
611                        .ok_or_else(batch_unavailable)?
612                        .entity_relation_id,
613                    version_id: source.version.version_id,
614                }
615            }
616        };
617        let budgets = batch_budgets(session);
618        let state = service
619            .capture_batch_state(session, resolution, scope, budgets)
620            .await?;
621        let columns = source
622            .variables
623            .iter()
624            .take(1)
625            .map(|v| v.name.as_str().to_string())
626            .collect::<Vec<_>>();
627        if columns.is_empty() {
628            return Err(batch_unavailable());
629        }
630        let query = batch_query(datastore, &source, &columns, budgets);
631        let latest = if source.latest {
632            vec![source.version.version_id.0]
633        } else {
634            Vec::new()
635        };
636        let mut disclosure = Self::admit_disclosure_with_semantic_state(
637            session,
638            DisclosureRequest::Query(query),
639            Some(Arc::clone(&state)),
640            latest,
641        )
642        .await?;
643        let read_state = Arc::clone(&state);
644        disclosure.prepare_semantic_readback(&state,move |repository| {
645            let retained_mappings=repository.retained_semantic_mappings(&read_state)?;
646            let service = semantic_transaction_service(repository);
647            async {
648                let instances=service.batch_instances(datastore,read_state.scope()).await?;
649                let mut mappings=Vec::new();
650                let mut links=Vec::new();
651                let study_ref=crate::projections::object_ref(datastore,ObjectKind::Study,source.study.study_id.0);
652                let version_ref=crate::projections::object_ref(datastore,ObjectKind::AssetVersion,source.version.version_id.0);
653                match (read_state.scope(), retained_mappings) {
654                  (SemanticScope::Relation { relation_id, .. }, ahri_tre_pgmeta::semantic::SemanticMappings::Relation(relation_mappings)) => {
655                    let definition = relation_id.0;
656                    for row in relation_mappings {
657                        if row.entity_relation_id.0!=definition {continue;}
658                        let instance=service.entities.get_relation_instance(row.relation_instance_id).await?.ok_or_else(batch_core_unavailable)?;
659                        mappings.push(serde_json::json!({"study":study_ref,"definition":integer_ref(datastore,ObjectKind::Relation,"relation",definition),"instance":crate::projections::object_ref(datastore,ObjectKind::RelationInstance,instance.uuid),"external_id":row.external_id,"transformation":integer_ref(datastore,ObjectKind::Transformation,"transformation",row.transformation_id.0)}));
660                    }
661                    for row in service.entities.list_dataset_version_relation_instances(source.version.version_id).await? {
662                        let instance=service.entities.get_relation_instance(row.relation_instance_id).await?.ok_or_else(batch_core_unavailable)?;
663                        if instance.entity_relation_id.0!=definition {continue;}
664                        links.push(serde_json::json!({"version":version_ref,"instance":crate::projections::object_ref(datastore,ObjectKind::RelationInstance,instance.uuid),"subject_variable":integer_ref(datastore,ObjectKind::Variable,"variable",row.subject_variable_id.0),"object_variable":integer_ref(datastore,ObjectKind::Variable,"variable",row.object_variable_id.0),"relation_variable":row.relation_variable_id.map(|id|integer_ref(datastore,ObjectKind::Variable,"variable",id.0)),"transformation":integer_ref(datastore,ObjectKind::Transformation,"transformation",row.transformation_id.0)}));
665                    }
666                  }
667                  (SemanticScope::Entity { entity_id, .. }, ahri_tre_pgmeta::semantic::SemanticMappings::Entity(entity_mappings)) => {
668                    let definition = entity_id.0;
669                    for row in entity_mappings {
670                        if row.entity_id.0!=definition {continue;}
671                        let instance=service.entities.get_entity_instance(row.entity_instance_id).await?.ok_or_else(batch_core_unavailable)?;
672                        mappings.push(serde_json::json!({"study":study_ref,"definition":integer_ref(datastore,ObjectKind::Entity,"entity",definition),"instance":crate::projections::object_ref(datastore,ObjectKind::EntityInstance,instance.uuid),"external_id":row.external_id,"transformation":integer_ref(datastore,ObjectKind::Transformation,"transformation",row.transformation_id.0)}));
673                    }
674                    for row in service.entities.list_dataset_version_entities(source.version.version_id).await? {
675                        let instance=service.entities.get_entity_instance(row.entity_instance_id).await?.ok_or_else(batch_core_unavailable)?;
676                        if instance.entity_id.0!=definition {continue;}
677                        links.push(serde_json::json!({"version":version_ref,"instance":crate::projections::object_ref(datastore,ObjectKind::EntityInstance,instance.uuid),"variable":integer_ref(datastore,ObjectKind::Variable,"variable",row.entity_variable_id.0),"transformation":integer_ref(datastore,ObjectKind::Transformation,"transformation",row.transformation_id.0)}));
678                    }
679                }
680                  _ => return Err(batch_core_unavailable()),
681                }
682                let rows=(instances.len()+mappings.len()+links.len()) as u64;
683                let bytes=serde_json::to_vec(&serde_json::json!({"instances":instances,"mappings":mappings,"links":links})).map_err(|_|batch_core_unavailable())?;
684                Ok((bytes,rows))
685            }.now_or_never().ok_or_else(batch_core_unavailable)?
686        })?;
687        Ok(disclosure)
688    }
689    async fn capture_batch_state(
690        &self,
691        session: &DataStoreSession,
692        resolution: SemanticResolution,
693        scope: SemanticScope,
694        budgets: DisclosureBudgets,
695    ) -> Result<Arc<SemanticState>, AppError> {
696        let repository = session
697            .dataset_admission_repository()
698            .ok_or_else(batch_unavailable)?;
699        run_blocking_app_work(move || {
700            repository
701                .capture_semantic_state(&resolution, &scope, &budgets)
702                .map(Arc::new)
703                .map_err(|_| batch_unavailable())
704        })
705        .await
706    }
707}
708fn batch_budgets(session: &DataStoreSession) -> DisclosureBudgets {
709    let mut budgets = session.runtime().disclosure.defaults;
710    budgets.payload_bytes = budgets.payload_bytes.min(1024 * 1024);
711    budgets.total_seconds = budgets.total_seconds.min(600);
712    budgets.rows = budgets.rows.min(10_000);
713    budgets
714}
715fn batch_query(
716    datastore: ahri_tre_protocol::PublicUuid,
717    source: &BatchSource,
718    columns: &[String],
719    budgets: DisclosureBudgets,
720) -> DisclosureQueryRequest {
721    DisclosureQueryRequest {
722        preview_limit: None,
723        representation: ahri_tre_types::ContentRepresentation::Json,
724        inputs: vec![DisclosureQueryInput {
725            alias: "content".into(),
726            asset: asset::AssetSelector::Id {
727                asset: crate::projections::object_ref(
728                    datastore,
729                    ObjectKind::AssetVersion,
730                    source.version.version_id.0,
731                ),
732                study: None,
733                asset_type: Some(AssetType::Dataset),
734            },
735            version: None,
736        }],
737        views: Default::default(),
738        sql: format!(
739            "SELECT {} FROM content",
740            columns
741                .iter()
742                .map(|c| format!("CAST({0} AS VARCHAR) AS {0}", sql_identifier(c)))
743                .collect::<Vec<_>>()
744                .join(", ")
745        ),
746        budgets: ahri_tre_types::DisclosureBudgetOverrides {
747            payload_bytes: Some(budgets.payload_bytes),
748            decoded_file_bytes: Some(budgets.decoded_file_bytes),
749            rows: Some(budgets.rows),
750            total_seconds: Some(budgets.total_seconds),
751        },
752    }
753}
754fn integer_ref(
755    datastore: ahri_tre_protocol::PublicUuid,
756    kind: ObjectKind,
757    scope: &str,
758    id: i64,
759) -> ObjectRef {
760    ObjectRef {
761        datastore_id: datastore,
762        kind,
763        id: ahri_tre_protocol::refs::encode_scoped_integer_ref(scope, id),
764    }
765}
766fn semantic_transaction_service(
767    repository: ahri_tre_pgmeta::PgMetadataRepository<'_>,
768) -> ScopedAppService<'_> {
769    let repository = Arc::new(repository);
770    let repositories =
771        crate::session::ScopedSessionMetadataRepositories::from_repository(Arc::clone(&repository));
772    let mut service = ScopedAppService::new(repositories.into());
773    service.semantic_repository = Some(repository.as_ref().clone());
774    service
775}
776fn batch_source_rows(rows: Vec<BTreeMap<String, Option<String>>>) -> Vec<EntityBatchSourceRow> {
777    rows.into_iter()
778        .enumerate()
779        .map(|(i, values)| EntityBatchSourceRow {
780            row_number: i + 1,
781            values,
782        })
783        .collect()
784}
785fn batch_transformation(
786    datastore: ahri_tre_protocol::PublicUuid,
787    lineage: Option<ahri_tre_types::TransformationLineageRecord>,
788) -> Option<model::SemanticWorkflowTransformationSummary> {
789    lineage.map(|lineage| model::SemanticWorkflowTransformationSummary {
790        transformation_id: lineage.transformation.transformation_id.0,
791        transformation: integer_ref(
792            datastore,
793            ObjectKind::Transformation,
794            "transformation",
795            lineage.transformation.transformation_id.0,
796        ),
797        transformation_type: match lineage.transformation.transformation_type {
798            ahri_tre_types::TransformationType::Ingest => "ingest",
799            ahri_tre_types::TransformationType::Transform => "transform",
800            ahri_tre_types::TransformationType::Entity => "entity",
801            ahri_tre_types::TransformationType::Export => "export",
802            ahri_tre_types::TransformationType::Repository => "repository",
803        }
804        .into(),
805        description: Some(lineage.transformation.description),
806        input_count: lineage.inputs.len(),
807        output_count: lineage.outputs.len(),
808        inputs: lineage
809            .inputs
810            .iter()
811            .map(|i| {
812                crate::projections::object_ref(datastore, ObjectKind::AssetVersion, i.version_id.0)
813            })
814            .collect(),
815        outputs: lineage
816            .outputs
817            .iter()
818            .map(|i| {
819                crate::projections::object_ref(datastore, ObjectKind::AssetVersion, i.version_id.0)
820            })
821            .collect(),
822    })
823}
824fn batch_unavailable() -> AppError {
825    AppError::Infrastructure("Semantic batch prerequisites are unavailable".into())
826}
827fn batch_core_unavailable() -> ahri_tre_core::CoreError {
828    ahri_tre_core::CoreError::Infrastructure("Semantic batch could not be completed".into())
829}
830
831impl ScopedAppService<'_> {
832    async fn batch_instances(
833        &self,
834        datastore: ahri_tre_protocol::PublicUuid,
835        scope: &SemanticScope,
836    ) -> Result<Vec<model::SemanticInstanceReadback>, CoreError> {
837        let mut values = Vec::new();
838        match *scope {
839            SemanticScope::Relation {
840                version_id: version,
841                relation_id,
842                ..
843            } => {
844                let definition = relation_id.0;
845                for link in self
846                    .entities
847                    .list_dataset_version_relation_instances(version)
848                    .await?
849                {
850                    let row = self
851                        .entities
852                        .get_relation_instance(link.relation_instance_id)
853                        .await?
854                        .ok_or_else(batch_core_unavailable)?;
855                    if row.entity_relation_id.0 != definition {
856                        continue;
857                    }
858                    let subject = self
859                        .entities
860                        .get_entity_instance(row.entity_instance_id_1)
861                        .await?
862                        .ok_or_else(batch_core_unavailable)?;
863                    let object = self
864                        .entities
865                        .get_entity_instance(row.entity_instance_id_2)
866                        .await?
867                        .ok_or_else(batch_core_unavailable)?;
868                    values.push(model::SemanticInstanceReadback {
869                        instance: crate::projections::object_ref(
870                            datastore,
871                            ObjectKind::RelationInstance,
872                            row.uuid,
873                        ),
874                        definition: integer_ref(
875                            datastore,
876                            ObjectKind::Relation,
877                            "relation",
878                            definition,
879                        ),
880                        transformation: integer_ref(
881                            datastore,
882                            ObjectKind::Transformation,
883                            "transformation",
884                            row.transformation_id.0,
885                        ),
886                        label: None,
887                        note: row.note,
888                        subject: Some(crate::projections::object_ref(
889                            datastore,
890                            ObjectKind::EntityInstance,
891                            subject.uuid,
892                        )),
893                        object: Some(crate::projections::object_ref(
894                            datastore,
895                            ObjectKind::EntityInstance,
896                            object.uuid,
897                        )),
898                        valid_from: row.valid_from,
899                        valid_to: row.valid_to,
900                    });
901                }
902            }
903            SemanticScope::Entity {
904                version_id: version,
905                entity_id,
906                ..
907            } => {
908                let definition = entity_id.0;
909                for link in self.entities.list_dataset_version_entities(version).await? {
910                    let row = self
911                        .entities
912                        .get_entity_instance(link.entity_instance_id)
913                        .await?
914                        .ok_or_else(batch_core_unavailable)?;
915                    if row.entity_id.0 != definition {
916                        continue;
917                    }
918                    values.push(model::SemanticInstanceReadback {
919                        instance: crate::projections::object_ref(
920                            datastore,
921                            ObjectKind::EntityInstance,
922                            row.uuid,
923                        ),
924                        definition: integer_ref(
925                            datastore,
926                            ObjectKind::Entity,
927                            "entity",
928                            definition,
929                        ),
930                        transformation: integer_ref(
931                            datastore,
932                            ObjectKind::Transformation,
933                            "transformation",
934                            row.transformation_id.0,
935                        ),
936                        label: row.label,
937                        note: row.note,
938                        subject: None,
939                        object: None,
940                        valid_from: None,
941                        valid_to: None,
942                    });
943                }
944            }
945        }
946        values.sort_by_key(|row| row.instance.id.as_uuid());
947        Ok(values)
948    }
949}