Skip to main content

ahri_tre_app/service/
disclosure.rs

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    /// One trusted content-admission seam. Correlation and origin come only
35    /// from the retained configured Session, never caller-supplied UUIDs.
36    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        // Keep every captured input and output inside the recorded writer tree.
330        // Final admission and compensation both verify cleanup of that tree.
331        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    /// Non-content readback from the actual ledger Dataset. Metadata authorizes
388    /// the current custodian and supplies identities; it cannot replace the ledger.
389    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            // A concrete identity also names retained history after deletion.
400            // The metadata adapter authorizes this non-content scope directly.
401            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                    // A multi-Study disclosure does not reveal another Study's scope.
477                    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    /// Resolve provenance, admit under one policy snapshot and freeze the JSON
537    /// result. A transport can only consume the retained result; it cannot
538    /// reopen a selector or select a different set of contributing versions.
539    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        // Dictionary identity remains available without inferred content when
701        // admission, bounds or evidence prerequisites cannot be satisfied.
702        captured.or_else(|_| Ok(PreparedSemanticResponse::Schema(withheld)))
703    }
704
705    /// One lifecycle for bounded semantic JSON: reserve, admit the retained
706    /// origin, acknowledge evidence, then retain delivery until completion.
707    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}