1use 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 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}