1use super::*;
2use crate::StoreSessionConnection;
3use crate::disclosure::{
4 DisclosureQueryInput, DisclosureQueryRequest, DisclosureRequest, PreparedDisclosure,
5 SessionDisclosure,
6};
7use ahri_tre_core::DatasetExecutorLease;
8use ahri_tre_lake::{
9 DisclosureDatasetInput, PreparedDisclosureQuery, QueryBinding, analyze_disclosure_query,
10};
11use ahri_tre_pgmeta::governance;
12use sha2::{Digest, Sha256};
13use uuid::Uuid;
14
15impl AppService {
16 pub(super) async fn bind_staged_datafile(
17 operator: Option<governance::GovernanceMaintenance>,
18 root: String,
19 study: ahri_tre_types::StudyId,
20 asset: ahri_tre_types::AssetId,
21 datafile: ahri_tre_types::DataFileRecord,
22 ) -> Result<(), AppError> {
23 let Some(operator) = operator else {
24 return Ok(());
25 };
26 run_blocking_app_work(move || {
27 let binding = ahri_tre_lake::attest_datafile(&root, study, asset, &datafile)
28 .map_err(|_| unavailable())?;
29 operator.bind_datafile(&binding).map_err(|_| unavailable())
30 })
31 .await
32 }
33
34 pub async fn admit_disclosure<'session>(
37 session: &'session mut DataStoreSession,
38 request: DisclosureRequest,
39 ) -> Result<SessionDisclosure<'session>, AppError> {
40 Self::admit_disclosure_with_semantic_state(session, request, None, Vec::new()).await
41 }
42
43 pub(super) async fn admit_disclosure_with_semantic_state<'session>(
44 session: &'session mut DataStoreSession,
45 request: DisclosureRequest,
46 semantic: Option<Arc<ahri_tre_pgmeta::semantic::SemanticState>>,
47 required_latest: Vec<Uuid>,
48 ) -> Result<SessionDisclosure<'session>, AppError> {
49 let (request, datafile_representation) = match request {
50 DisclosureRequest::Query(request) => (request, None),
51 DisclosureRequest::Datafile {
52 compress,
53 representation,
54 asset,
55 version,
56 budgets,
57 } => (
58 DisclosureQueryRequest {
59 preview_limit: None,
60 representation: if compress {
61 ahri_tre_types::ContentRepresentation::FileZstd
62 } else {
63 ahri_tre_types::ContentRepresentation::File
64 },
65 inputs: vec![DisclosureQueryInput {
66 alias: "content".into(),
67 asset,
68 version,
69 }],
70 views: Default::default(),
71 sql: "SELECT * FROM content".into(),
72 budgets,
73 },
74 Some(representation),
75 ),
76 };
77 let budgets = session
78 .runtime()
79 .disclosure
80 .resolve(&request.budgets)
81 .ok_or_else(|| {
82 AppError::Validation(
83 "Disclosure budgets exceed the configured ceiling or are invalid".into(),
84 )
85 })?;
86 let service = Self::from_datastore_session(session)?;
87 let (session_id, authentication) = session
88 .disclosure_context
89 .as_ref()
90 .ok_or_else(unavailable)?;
91 let operator = session.governance_maintenance().ok_or_else(unavailable)?;
92 let executor = session.dataset_executor().ok_or_else(unavailable)?;
93 let scratch = session.shared_scratch_attempt().ok_or_else(unavailable)?;
94 let datastore_id = service.datastore_id.ok_or_else(unavailable)?;
95 let actor = session.authenticated_actor().ok_or_else(unavailable)?;
96 if authentication
97 .runtime_expires_at()
98 .is_some_and(|deadline| deadline <= Utc::now())
99 {
100 return Err(unavailable());
101 }
102 let activate = operator.clone();
103 let issuer = authentication.identity().issuer().to_owned();
104 let subject = authentication.identity().subject().to_owned();
105 let principal = actor.to_owned();
106 run_blocking_app_work(move || {
107 activate
108 .activate_identity(&issuer, &subject, &principal)
109 .map_err(|_| unavailable())
110 })
111 .await?;
112 if request.inputs.is_empty() || request.inputs.len() > 128 {
113 return Err(AppError::Validation(
114 "Invalid disclosure input count".into(),
115 ));
116 }
117 let permit = crate::disclosure::acquire_disclosure_permit()?;
118 let mut bindings = std::collections::BTreeMap::new();
119 let mut latest_inputs = required_latest;
120 for input in request.inputs {
121 let (_, catalog, pinned) = service.resolve_catalogue_asset(input.asset).await?;
122 let version = service.catalogue_version(&catalog, pinned, input.version.as_deref())?;
123 if pinned.is_none()
124 && input
125 .version
126 .as_deref()
127 .is_none_or(|value| value.eq_ignore_ascii_case("latest"))
128 {
129 latest_inputs.push(version.version_id.0);
130 }
131 if bindings
132 .insert(input.alias, QueryBinding::Dataset(version.version_id.0))
133 .is_some()
134 {
135 return Err(AppError::Conflict("Duplicate query input alias".into()));
136 }
137 }
138 for (name, sql) in request.views {
139 if bindings.insert(name, QueryBinding::View(sql)).is_some() {
140 return Err(AppError::Conflict("Duplicate query binding".into()));
141 }
142 }
143 let plan = analyze_disclosure_query(&request.sql, &bindings)
144 .map_err(|_| AppError::Validation("Query dependencies are not admissible".into()))?;
145 let admission_id = Uuid::new_v4();
146 let identity = executor.identity();
147 let guard = executor
148 .begin_attempt(admission_id)
149 .map_err(|_| unavailable())?;
150 if authentication.runtime_expires_at().is_some_and(|at| {
151 at < Utc::now() + chrono::Duration::seconds(i64::from(budgets.total_seconds))
152 }) {
153 return Err(AppError::Validation(
154 "Disclosure lifetime exceeds this Session's remaining lifetime".into(),
155 ));
156 }
157 let mut required_inputs = plan.input_versions().clone();
158 if let Some(state) = &semantic {
159 required_inputs.extend(state.inputs().iter().copied());
160 }
161 let intent = ahri_tre_types::DisclosureIntent {
162 admission_id,
163 session_id: *session_id,
164 login_id: authentication.owner().as_uuid(),
165 datastore_id,
166 worker_generation: identity.generation_id,
167 coordinator_id: identity.coordinator_id,
168 inputs: required_inputs.into_iter().collect(),
169 latest_inputs: latest_inputs
170 .into_iter()
171 .filter(|id| plan.input_versions().contains(id))
172 .collect(),
173 query_fingerprint: Sha256::digest(plan.restricted_sql().as_bytes())
174 .iter()
175 .map(|byte| format!("{byte:02x}"))
176 .collect(),
177 configuration_fingerprint: session
178 .runtime()
179 .session_provenance
180 .as_ref()
181 .ok_or_else(unavailable)?
182 .configuration_fingerprint()
183 .to_string(),
184 representation: match datafile_representation {
185 Some((false, false)) => "datafile_stored",
186 Some((true, false)) => "datafile_decrypted",
187 _ => request.representation.as_str(),
188 }
189 .into(),
190 lifetime_seconds: budgets.total_seconds,
191 max_bytes: budgets.payload_bytes,
192 max_decoded_bytes: budgets.decoded_file_bytes,
193 max_rows: budgets.rows,
194 };
195 let session_deadline = authentication.runtime_expires_at();
196 let lake_root = session.runtime().lake.data_path.clone();
197 let authority = session.disclosure_capture.clone().ok_or_else(unavailable)?;
198 let store = session.store.clone();
199 let mut lake = session.lake.try_clone().map_err(|_| unavailable())?;
200 let projection = operator.clone();
201 let executable = std::env::current_exe()
202 .map_err(|_| unavailable())?
203 .with_file_name("ahri-tre-query-worker");
204 let (snapshot, prepared, guard, permit) = run_blocking_app_work(move || {
205 for event in projection
206 .required_evidence(&intent.inputs, 1000)
207 .map_err(|_| unavailable())?
208 {
209 crate::disclosure::project_session_evidence(&projection, &mut lake, event)?;
210 }
211
212 let (snapshot, prepared) = if let Some(representation) = datafile_representation {
213 let capture = |snapshot: &ahri_tre_types::DisclosureSnapshot, input| {
214 if session_deadline.is_some_and(|expires| snapshot.deadline > expires) {
215 return Err(());
216 }
217 ahri_tre_lake::PreparedDisclosureDatafile::capture_representation(
218 &lake_root,
219 input,
220 &scratch,
221 snapshot.deadline,
222 snapshot.max_decoded_bytes,
223 &executable,
224 representation,
225 )
226 .map(|file| {
227 PreparedDisclosure::Datafile(
228 file,
229 request.representation
230 == ahri_tre_types::ContentRepresentation::FileZstd,
231 )
232 })
233 .map_err(|_| ())
234 };
235 match store {
236 StoreSessionConnection::Direct(db) => {
237 governance::admit_disclosure_with_datafile(
238 &mut *db.lock().map_err(|_| unavailable())?,
239 &intent,
240 capture,
241 )
242 }
243 StoreSessionConnection::OAuth(db) => {
244 governance::admit_disclosure_with_datafile(
245 &mut *db.lock().map_err(|_| unavailable())?,
246 &intent,
247 capture,
248 )
249 }
250 #[cfg(test)]
251 StoreSessionConnection::TestUnavailable => return Err(unavailable()),
252 }
253 } else {
254 let capture =
255 |snapshot: &ahri_tre_types::DisclosureSnapshot,
256 origins: Vec<ahri_tre_types::GovernedDatasetLocation>| {
257 if session_deadline.is_some_and(|expires| snapshot.deadline > expires) {
258 return Err(());
259 }
260 let inputs = origins
261 .into_iter()
262 .filter(|origin| plan.input_versions().contains(&origin.version_id))
263 .map(DisclosureDatasetInput::from)
264 .collect::<Vec<_>>();
265 PreparedDisclosureQuery::capture(
266 &authority,
267 plan,
268 &inputs,
269 &scratch,
270 snapshot.deadline,
271 &executable,
272 )
273 .map(|query| {
274 PreparedDisclosure::Query(
275 query,
276 request.representation,
277 request.preview_limit,
278 )
279 })
280 .map_err(|_| ())
281 };
282 match store {
283 StoreSessionConnection::Direct(db) => {
284 governance::admit_disclosure_with_semantic_inputs(
285 &mut *db.lock().map_err(|_| unavailable())?,
286 &intent,
287 semantic.as_deref(),
288 capture,
289 )
290 }
291 StoreSessionConnection::OAuth(db) => {
292 governance::admit_disclosure_with_semantic_inputs(
293 &mut *db.lock().map_err(|_| unavailable())?,
294 &intent,
295 semantic.as_deref(),
296 capture,
297 )
298 }
299 #[cfg(test)]
300 StoreSessionConnection::TestUnavailable => return Err(unavailable()),
301 }
302 }
303 .map_err(|_| unavailable())?;
304 if Utc::now() >= snapshot.deadline {
305 return Err(unavailable());
306 }
307 crate::disclosure::project_session_evidence(
308 &projection,
309 &mut lake,
310 snapshot.admission_evidence,
311 )?;
312 if Utc::now() >= snapshot.deadline {
313 return Err(unavailable());
314 }
315 Ok((snapshot, prepared, guard, permit))
316 })
317 .await?;
318 SessionDisclosure::new(session, operator, snapshot, prepared, guard, permit)
319 }
320
321 pub(super) async fn prepare_dataset_sql_transform(
322 session: &DataStoreSession,
323 scratch: &ahri_tre_lake::ScratchAttempt,
324 plan: ahri_tre_lake::RestrictedQueryPlan,
325 budgets: ahri_tre_types::DisclosureBudgets,
326 deadline: chrono::DateTime<Utc>,
327 ) -> Result<ahri_tre_lake::PreparedQueryOutput, AppError> {
328 let authority = session.disclosure_capture.clone().ok_or_else(unavailable)?;
329 let scratch = scratch
332 .create_child(
333 ahri_tre_lake::ScratchAttemptId::new(&Uuid::new_v4().simple().to_string())
334 .map_err(|_| unavailable())?,
335 )
336 .map_err(|_| unavailable())?;
337 let store = session.store.clone();
338 let permit = crate::disclosure::acquire_disclosure_permit()?;
339 let executor = session.dataset_executor().ok_or_else(unavailable)?;
340 let guard = executor
341 .begin_attempt(Uuid::new_v4())
342 .map_err(|_| unavailable())?;
343 let executable = std::env::current_exe()
344 .map_err(|_| unavailable())?
345 .with_file_name("ahri-tre-query-worker");
346 run_blocking_app_work(move || {
347 let _permit = permit;
348 let _guard = guard;
349 let inputs = plan.input_versions().iter().copied().collect::<Vec<_>>();
350 let capture = |origins: Vec<ahri_tre_types::GovernedDatasetLocation>| {
351 let inputs = origins
352 .into_iter()
353 .map(DisclosureDatasetInput::from)
354 .collect::<Vec<_>>();
355 PreparedDisclosureQuery::capture_owned(
356 &authority,
357 plan,
358 &inputs,
359 scratch,
360 deadline,
361 &executable,
362 )
363 .map_err(|_| ())
364 };
365 let prepared = match store {
366 StoreSessionConnection::Direct(db) => governance::capture_derivation_inputs(
367 &mut *db.lock().map_err(|_| unavailable())?,
368 &inputs,
369 capture,
370 ),
371 StoreSessionConnection::OAuth(db) => governance::capture_derivation_inputs(
372 &mut *db.lock().map_err(|_| unavailable())?,
373 &inputs,
374 capture,
375 ),
376 #[cfg(test)]
377 StoreSessionConnection::TestUnavailable => return Err(unavailable()),
378 }
379 .map_err(|_| unavailable())?;
380 prepared
381 .materialize(&executable, deadline, budgets.payload_bytes, budgets.rows)
382 .map_err(|_| unavailable())
383 })
384 .await
385 }
386
387 pub async fn governance_history(
390 session: &mut DataStoreSession,
391 study: ahri_tre_protocol::study::StudySelector,
392 after: chrono::DateTime<chrono::Utc>,
393 after_id: Uuid,
394 limit: u32,
395 ) -> Result<Vec<serde_json::Value>, AppError> {
396 let service = Self::from_datastore_session(session)?;
397 let selector = service.catalogue_study_selector(study)?;
398 let study = match selector {
399 StudySelector::Id { study_id } => study_id.0,
402 selector => {
403 service
404 .resolve_study_selector(selector, "Governance history Study")
405 .await?
406 .study_id
407 .0
408 }
409 };
410 let (_, authentication) = session
411 .disclosure_context
412 .as_ref()
413 .ok_or_else(unavailable)?;
414 if authentication
415 .runtime_expires_at()
416 .is_some_and(|at| at <= Utc::now())
417 {
418 return Err(unavailable());
419 }
420 let issuer = authentication.identity().issuer().to_owned();
421 let subject = authentication.identity().subject().to_owned();
422 let actor = session
423 .authenticated_actor()
424 .ok_or_else(unavailable)?
425 .to_owned();
426 let store = session.store.clone();
427 let operator = session.governance_maintenance().ok_or_else(unavailable)?;
428 let mut lake = session.lake.try_clone().map_err(|_| unavailable())?;
429 run_blocking_app_work(move || {
430 operator
431 .activate_identity(&issuer, &subject, &actor)
432 .map_err(|_| unavailable())?;
433 let events = match store {
434 StoreSessionConnection::Direct(db) => governance::scoped_evidence_ids(
435 &mut *db.lock().map_err(|_| unavailable())?,
436 study,
437 after,
438 after_id,
439 limit,
440 ),
441 StoreSessionConnection::OAuth(db) => governance::scoped_evidence_ids(
442 &mut *db.lock().map_err(|_| unavailable())?,
443 study,
444 after,
445 after_id,
446 limit,
447 ),
448 #[cfg(test)]
449 StoreSessionConnection::TestUnavailable => return Err(unavailable()),
450 }
451 .map_err(|_| unavailable())?;
452 let binding = operator.binding().map_err(|_| unavailable())?;
453 let ledger = ahri_tre_lake::GovernanceLedger::open(
454 &mut lake,
455 ahri_tre_types::StudyId(binding.study_id),
456 )
457 .map_err(|_| unavailable())?;
458 events
459 .into_iter()
460 .map(|event| {
461 let body = ledger
462 .read(event)
463 .map_err(|_| unavailable())?
464 .ok_or_else(unavailable)?;
465 let mut value: serde_json::Value =
466 serde_json::from_str(&body).map_err(|_| unavailable())?;
467 if let Some(inputs) = value
468 .pointer_mut("/evidence/inputs")
469 .and_then(serde_json::Value::as_array_mut)
470 {
471 inputs.retain(|input| {
472 input.get("study_id").and_then(serde_json::Value::as_str)
473 == Some(study.to_string().as_str())
474 });
475 }
476 value["study_id"] = serde_json::json!(study);
478 Ok(value)
479 })
480 .collect()
481 })
482 .await
483 }
484
485 pub async fn reclassify_catalogue_asset(
486 session: &mut DataStoreSession,
487 selector: ahri_tre_protocol::asset::AssetSelector,
488 revision: u64,
489 risk: ahri_tre_types::AssetRisk,
490 justification: &str,
491 ) -> Result<ahri_tre_types::AssetClassification, AppError> {
492 let service = Self::from_datastore_session(session)?;
493 let (_, catalog, _) = service.resolve_catalogue_asset(selector).await?;
494 let store = session.store.clone();
495 let operator = session.governance_maintenance().ok_or_else(unavailable)?;
496 let mut lake = session.lake.try_clone().map_err(|_| unavailable())?;
497 let justification = justification.to_owned();
498 run_blocking_app_work(move || {
499 let classification = match store {
500 StoreSessionConnection::Direct(db) => governance::reclassify(
501 &mut *db.lock().map_err(|_| unavailable())?,
502 catalog.asset.asset_id.0,
503 revision,
504 risk,
505 &justification,
506 ),
507 StoreSessionConnection::OAuth(db) => governance::reclassify(
508 &mut *db.lock().map_err(|_| unavailable())?,
509 catalog.asset.asset_id.0,
510 revision,
511 risk,
512 &justification,
513 ),
514 #[cfg(test)]
515 StoreSessionConnection::TestUnavailable => return Err(unavailable()),
516 }
517 .map_err(|_| unavailable())?;
518 crate::disclosure::project_session_evidence(
519 &operator,
520 &mut lake,
521 classification.evidence_id,
522 )?;
523 Ok(ahri_tre_types::AssetClassification {
524 evidence_durable: true,
525 ..classification
526 })
527 })
528 .await
529 }
530}
531fn unavailable() -> AppError {
532 AppError::Infrastructure("Disclosure prerequisites are unavailable".into())
533}
534
535impl AppService {
536 pub async fn prepare_semantic_response<'session>(
540 session: &'session mut DataStoreSession,
541 status: ahri_tre_protocol::session::SessionStatusPayload,
542 request: ahri_tre_protocol::request::ProtocolRequest,
543 ) -> Result<crate::disclosure::PreparedSemanticResponse<'session>, AppError> {
544 use crate::disclosure::PreparedSemanticResponse;
545 use ahri_tre_pgmeta::semantic::{SemanticOrigin, SemanticTarget};
546 use ahri_tre_protocol::request::ProtocolRequest;
547 let request = match request {
548 request @ (ProtocolRequest::EntityInstanceAdd(_)
549 | ProtocolRequest::RelationInstanceAdd(_)
550 | ProtocolRequest::EntityInstanceMapAdd(_)
551 | ProtocolRequest::RelationInstanceMapAdd(_)
552 | ProtocolRequest::EntityInstanceAssetLinkAdd(_)
553 | ProtocolRequest::RelationInstanceAssetLinkAdd(_)
554 | ProtocolRequest::EntityInstanceDatasetLinkAdd(_)
555 | ProtocolRequest::RelationInstanceDatasetLinkAdd(_)) => {
556 return Self::prepare_instance_mutation(session, request).await;
557 }
558 request @ (ProtocolRequest::EntityInstanceGet(_)
559 | ProtocolRequest::EntityInstanceList(_)
560 | ProtocolRequest::RelationInstanceGet(_)
561 | ProtocolRequest::RelationInstanceList(_)
562 | ProtocolRequest::EntityInstanceMapGet(_)
563 | ProtocolRequest::RelationInstanceMapGet(_)
564 | ProtocolRequest::EntityInstanceMapList(_)
565 | ProtocolRequest::RelationInstanceMapList(_)
566 | ProtocolRequest::EntityInstanceAssetLinkList(_)
567 | ProtocolRequest::RelationInstanceAssetLinkList(_)
568 | ProtocolRequest::EntityInstanceDatasets(_)
569 | ProtocolRequest::EntityInstanceDatasetLinkGet(_)
570 | ProtocolRequest::RelationInstanceDatasetLinkGet(_)
571 | ProtocolRequest::EntityInstanceDatasetLinkList(_)
572 | ProtocolRequest::RelationInstanceDatasetLinkList(_)) => {
573 return Self::prepare_instance_read(session, request).await;
574 }
575 ProtocolRequest::EntityInstanceEnsureFromDataset(request) => {
576 return Self::prepare_entity_batch(session, request).await;
577 }
578 ProtocolRequest::RelationInstanceEnsureFromDataset(request) => {
579 return Self::prepare_relation_batch(session, request).await;
580 }
581 request => request,
582 };
583 let service = Self::from_datastore_session(session)?;
584 let datastore = service.catalogue_datastore_id()?;
585 let fingerprint = Sha256::digest(serde_json::to_vec(&request).map_err(|_| unavailable())?)
586 .iter()
587 .map(|b| format!("{b:02x}"))
588 .collect();
589 let repository = session
590 .dataset_admission_repository()
591 .ok_or_else(unavailable)?;
592 let mut slots = Vec::new();
593 let mut safe = match request {
594 ProtocolRequest::VocabularyGet(request) => {
595 let response = service.get_catalogue_vocabulary(request).await?;
596 slots.push(SemanticVocabularySlot::new(
597 &service,
598 response.vocabulary.summary.vocabulary,
599 response.vocabulary.summary.domain,
600 "/vocabulary",
601 VocabularyShape::Vocabulary,
602 )?);
603 serde_json::to_value(response)
604 }
605 ProtocolRequest::VariableGet(request) => {
606 let response = service.get_catalogue_variable(request).await?;
607 if let Some(vocabulary) = &response.variable.summary.vocabulary {
608 slots.push(SemanticVocabularySlot::new(
609 &service,
610 vocabulary.vocabulary,
611 vocabulary.domain,
612 "/variable",
613 VocabularyShape::Variable,
614 )?);
615 }
616 serde_json::to_value(response)
617 }
618 ProtocolRequest::DatasetMetadata(request) => {
619 let result = service
620 .read_catalogue_dataset(request.dataset, request.with_variables)
621 .await?;
622 let mut declarations = Vec::new();
623 for (index, variable) in result.metadata.variables.iter().enumerate() {
624 if let Some(vocabulary) = &variable.vocabulary {
625 let summary =
626 crate::projections::vocabulary_summary(datastore, vocabulary.clone());
627 let slot = SemanticVocabularySlot::new(
628 &service,
629 summary.vocabulary_id,
630 summary.domain_id,
631 &format!("/metadata/variables/{index}"),
632 VocabularyShape::Dataset,
633 )?;
634 if let Some(items) = service.declared_vocabulary_items(slot.id) {
635 declarations.push((slot.clone(), items));
636 }
637 slots.push(slot);
638 }
639 }
640 let mut value = serde_json::to_value(
641 crate::projections::dataset_metadata_response(status, result),
642 )
643 .map_err(|_| unavailable())?;
644 for (slot, items) in declarations {
645 slot.project(&mut value, datastore, items)?;
646 }
647 Ok(value)
648 }
649 _ => return Err(AppError::Validation("Unsupported semantic response".into())),
650 }
651 .map_err(|_| unavailable())?;
652 let mut inputs = std::collections::BTreeSet::new();
653 let mut admitted = Vec::new();
654 for slot in slots {
655 if let SemanticOrigin::Derived { inputs: sources } = repository
656 .semantic_origin(&SemanticTarget::Vocabulary(slot.id))
657 .map_err(|_| unavailable())?
658 {
659 inputs.extend(sources);
660 admitted.push(slot);
661 }
662 }
663 if admitted.is_empty() {
664 return Ok(PreparedSemanticResponse::Schema(safe));
665 }
666 let inputs: Vec<_> = inputs.into_iter().collect();
667 let targets: Vec<_> = admitted
668 .iter()
669 .map(|slot| slot.id.0)
670 .collect::<std::collections::BTreeSet<_>>()
671 .into_iter()
672 .map(ahri_tre_types::VocabularyId)
673 .collect();
674 let withheld = safe.clone();
675 let max_rows = session.runtime().disclosure.defaults.rows;
676 let captured = Self::prepare_semantic_delivery(
677 session,
678 inputs,
679 fingerprint,
680 max_rows,
681 move |repository, intent| {
682 let (snapshot, values) = repository
683 .capture_semantic_vocabularies(intent, &targets)
684 .map_err(|_| unavailable())?;
685 let rows = values.iter().map(|(_, items)| items.len() as u64).sum();
686 for slot in admitted {
687 let (vocabulary, items) = values
688 .iter()
689 .find(|(v, _)| v.vocabulary_id == slot.id)
690 .ok_or_else(unavailable)?;
691 if vocabulary.domain_id != slot.domain {
692 return Err(unavailable());
693 }
694 slot.project(&mut safe, datastore, items.clone())?;
695 }
696 Ok((snapshot, safe, rows))
697 },
698 )
699 .await;
700 captured.or_else(|_| Ok(PreparedSemanticResponse::Schema(withheld)))
703 }
704
705 pub(super) async fn prepare_semantic_delivery<'s>(
708 session: &'s mut DataStoreSession,
709 inputs: Vec<Uuid>,
710 fingerprint: String,
711 max_rows: u64,
712 capture: impl FnOnce(
713 &ahri_tre_pgmeta::PgMetadataRepository<'_>,
714 &ahri_tre_types::DisclosureIntent,
715 ) -> Result<
716 (ahri_tre_types::DisclosureSnapshot, serde_json::Value, u64),
717 AppError,
718 > + Send
719 + 'static,
720 ) -> Result<crate::disclosure::PreparedSemanticResponse<'s>, AppError> {
721 use crate::disclosure::PreparedSemanticResponse;
722 let budgets = session.runtime().disclosure.defaults;
723 let lifetime = budgets.total_seconds.min(600);
724 let (session_id, authentication) = session
725 .disclosure_context
726 .as_ref()
727 .ok_or_else(unavailable)?;
728 if authentication.runtime_expires_at().is_some_and(|expires| {
729 expires < Utc::now() + chrono::Duration::seconds(i64::from(lifetime))
730 }) {
731 return Err(unavailable());
732 }
733 let operator = session.governance_maintenance().ok_or_else(unavailable)?;
734 let repository = session
735 .dataset_admission_repository()
736 .ok_or_else(unavailable)?;
737 let executor = session.dataset_executor().ok_or_else(unavailable)?;
738 let admission_id = Uuid::new_v4();
739 let identity = executor.identity();
740 let guard = executor
741 .begin_attempt(admission_id)
742 .map_err(|_| unavailable())?;
743 let permit = crate::disclosure::acquire_disclosure_permit()?;
744 let principal = session
745 .authenticated_actor()
746 .ok_or_else(unavailable)?
747 .to_owned();
748 let issuer = authentication.identity().issuer().to_owned();
749 let subject = authentication.identity().subject().to_owned();
750 let intent = ahri_tre_types::DisclosureIntent {
751 admission_id,
752 session_id: *session_id,
753 login_id: authentication.owner().as_uuid(),
754 datastore_id: Self::from_datastore_session(session)?
755 .catalogue_datastore_id()?
756 .as_uuid(),
757 worker_generation: identity.generation_id,
758 coordinator_id: identity.coordinator_id,
759 inputs: inputs.clone(),
760 latest_inputs: Vec::new(),
761 query_fingerprint: fingerprint,
762 configuration_fingerprint: session
763 .runtime()
764 .session_provenance
765 .as_ref()
766 .ok_or_else(unavailable)?
767 .configuration_fingerprint()
768 .to_string(),
769 representation: "json".into(),
770 lifetime_seconds: lifetime,
771 max_bytes: budgets.payload_bytes.min(1024 * 1024),
772 max_decoded_bytes: budgets.decoded_file_bytes,
773 max_rows,
774 };
775 let projection = operator.clone();
776 let mut lake = session.lake.try_clone().map_err(|_| unavailable())?;
777 let (snapshot, prepared, guard, permit) = run_blocking_app_work(move || {
778 projection
779 .activate_identity(&issuer, &subject, &principal)
780 .map_err(|_| unavailable())?;
781 for event in projection
782 .required_evidence(&inputs, 1000)
783 .map_err(|_| unavailable())?
784 {
785 crate::disclosure::project_session_evidence(&projection, &mut lake, event)?;
786 }
787 let (snapshot, response, rows) = capture(&repository, &intent)?;
788 let bytes = serde_json::to_vec(&response).map_err(|_| unavailable())?;
789 if bytes.len() as u64 > snapshot.max_bytes
790 || rows > snapshot.max_rows
791 || Utc::now() >= snapshot.deadline
792 {
793 return Err(unavailable());
794 }
795 crate::disclosure::project_session_evidence(
796 &projection,
797 &mut lake,
798 snapshot.admission_evidence,
799 )?;
800 Ok((
801 snapshot,
802 PreparedDisclosure::Semantic { bytes, rows },
803 guard,
804 permit,
805 ))
806 })
807 .await?;
808 SessionDisclosure::new(session, operator, snapshot, prepared, guard, permit)
809 .map(|disclosure| PreparedSemanticResponse::Content(Box::new(disclosure)))
810 }
811}
812
813#[derive(Clone, Copy)]
814enum VocabularyShape {
815 Vocabulary,
816 Variable,
817 Dataset,
818}
819#[derive(Clone)]
820struct SemanticVocabularySlot {
821 id: ahri_tre_types::VocabularyId,
822 domain: ahri_tre_types::DomainId,
823 pointer: String,
824 shape: VocabularyShape,
825}
826impl SemanticVocabularySlot {
827 fn new(
828 service: &AppService,
829 vocabulary: ahri_tre_protocol::refs::ObjectRef,
830 domain: ahri_tre_protocol::refs::ObjectRef,
831 pointer: &str,
832 shape: VocabularyShape,
833 ) -> Result<Self, AppError> {
834 use ahri_tre_protocol::refs::ObjectKind;
835 Ok(Self {
836 id: ahri_tre_types::VocabularyId(service.semantic_integer_id(
837 vocabulary,
838 ObjectKind::Vocabulary,
839 "vocabulary",
840 )?),
841 domain: ahri_tre_types::DomainId(service.semantic_integer_id(
842 domain,
843 ObjectKind::Domain,
844 "domain",
845 )?),
846 pointer: pointer.into(),
847 shape,
848 })
849 }
850 fn project(
851 &self,
852 response: &mut serde_json::Value,
853 datastore: ahri_tre_protocol::PublicUuid,
854 items: Vec<ahri_tre_types::VocabularyItemRecord>,
855 ) -> Result<(), AppError> {
856 let target = response
857 .pointer_mut(&self.pointer)
858 .ok_or_else(unavailable)?;
859 let count = items.len();
860 let values = match self.shape {
861 VocabularyShape::Dataset => serde_json::to_value(
862 items
863 .into_iter()
864 .map(|item| crate::projections::vocabulary_item_summary(datastore, item))
865 .collect::<Vec<_>>(),
866 ),
867 _ => serde_json::to_value(
868 items
869 .into_iter()
870 .map(
871 |item| ahri_tre_protocol::dictionary::VocabularyItemSummary {
872 value: item.value,
873 code: item.code,
874 description: item.description,
875 },
876 )
877 .collect::<Vec<_>>(),
878 ),
879 }
880 .map_err(|_| unavailable())?;
881 match self.shape {
882 VocabularyShape::Vocabulary => {
883 target["items"] = values;
884 target["items_withheld"] = false.into();
885 target["summary"]["item_count"] = count.into();
886 }
887 VocabularyShape::Variable => {
888 target["vocabulary_items"] = values;
889 target["vocabulary_items_withheld"] = false.into();
890 target["summary"]["vocabulary"]["item_count"] = count.into();
891 }
892 VocabularyShape::Dataset => {
893 target["vocabulary_items"] = values;
894 target["vocabulary_items_withheld"] = false.into();
895 }
896 }
897 Ok(())
898 }
899}