Skip to main content

ahri_tre_pgmeta/
semantic_records.rs

1//! Bounded, revision-bound instance and composite-link readbacks.
2use crate::{MetadataExecutor, PgMetaError, PgMetadataRepository, schema::SchemaStatusRow};
3use ahri_tre_core::CoreError;
4use ahri_tre_types::{DisclosureBudgets, DisclosureIntent, DisclosureSnapshot};
5use serde::{Deserialize, Serialize};
6
7#[derive(Debug, Clone, Copy, Serialize)]
8#[serde(rename_all = "snake_case")]
9pub enum SemanticCollection {
10    EntityInstance,
11    RelationInstances,
12    StudyEntityInstances,
13    StudyRelationInstances,
14    AssetVersionEntities,
15    AssetVersionRelationInstances,
16    DatasetVersionEntities,
17    DatasetVersionRelationInstances,
18}
19
20/// Resolved keys, supplied only after application selector authorization.
21#[derive(Debug, Clone, Serialize)]
22pub struct SemanticRecordSelection {
23    pub collection: SemanticCollection,
24    pub definition: Option<i64>,
25    pub instance: Option<i64>,
26    pub study: Option<uuid::Uuid>,
27    pub version: Option<uuid::Uuid>,
28    pub external_id: Option<String>,
29}
30
31/// Adapter-produced capture. Construction is private so callers cannot supply
32/// their own origin closure, revision or retained authorization scope.
33///
34/// ```compile_fail
35/// use ahri_tre_pgmeta::semantic_records::CapturedSemanticRecords;
36/// let forged: CapturedSemanticRecords = serde_json::from_str("{}").unwrap();
37/// ```
38#[derive(Debug)]
39pub struct CapturedSemanticRecords {
40    stored: StoredSemanticRecords,
41    studies: Vec<uuid::Uuid>,
42}
43
44#[derive(Debug, Deserialize)]
45struct StoredSemanticRecords {
46    revision: i64,
47    rows: Vec<serde_json::Value>,
48    instances: Vec<serde_json::Value>,
49    origin: crate::semantic::SemanticOrigin,
50    versions: Vec<serde_json::Value>,
51    known_inputs: Vec<uuid::Uuid>,
52}
53impl CapturedSemanticRecords {
54    pub fn rows(&self) -> &[serde_json::Value] {
55        &self.stored.rows
56    }
57    pub fn instances(&self) -> &[serde_json::Value] {
58        &self.stored.instances
59    }
60    pub fn versions(&self) -> &[serde_json::Value] {
61        &self.stored.versions
62    }
63    /// Known contributors remain descriptive even when another field's origin
64    /// is Unknown. They cannot authorize disclosure or repair missing history.
65    pub fn known_inputs(&self) -> &[uuid::Uuid] {
66        &self.stored.known_inputs
67    }
68    pub fn origin(&self) -> &crate::semantic::SemanticOrigin {
69        &self.stored.origin
70    }
71}
72
73impl PgMetadataRepository<'_> {
74    pub fn capture_semantic_records(
75        &self,
76        resolution: &crate::semantic::SemanticResolution,
77        selection: &SemanticRecordSelection,
78        budgets: &DisclosureBudgets,
79    ) -> Result<CapturedSemanticRecords, CoreError> {
80        let selected_study = selection.study;
81        let selection = serde_json::to_string(selection).map_err(|_| unavailable())?;
82        self.with_dataset_acceptance(|repository| {
83            let mut captured = repository.with_session_connection(
84                "capture semantic records",
85                |db| capture(db, &selection, budgets),
86                |db| capture(db, &selection, budgets),
87            )?;
88            if captured.stored.revision != resolution.revision() {
89                return Err(unavailable());
90            }
91            let mut studies = std::collections::BTreeSet::new();
92            studies.extend(selected_study);
93            for row in &captured.stored.rows {
94                if let Some(id) = row["payload"]["study_id"].as_str() {
95                    studies.insert(id.parse().map_err(|_| unavailable())?);
96                }
97            }
98            for row in &captured.stored.versions {
99                if let Some(id) = row["study_id"].as_str() {
100                    studies.insert(id.parse().map_err(|_| unavailable())?);
101                }
102            }
103            captured.studies = studies.into_iter().collect();
104            Ok(captured)
105        })
106    }
107
108    /// Revalidate the retained revision under the same locks as the single policy
109    /// decision. Inputs come from the stored origin closure, never the caller.
110    pub fn admit_semantic_records(
111        &self,
112        captured: &CapturedSemanticRecords,
113        intent: &DisclosureIntent,
114    ) -> Result<DisclosureSnapshot, CoreError> {
115        let crate::semantic::SemanticOrigin::Derived { inputs } = &captured.stored.origin else {
116            return Err(unavailable());
117        };
118        if inputs.is_empty() {
119            return Err(unavailable());
120        }
121        let mut intent = intent.clone();
122        intent.inputs = inputs.clone();
123        intent.latest_inputs.clear();
124        self.with_dataset_acceptance(|repository| {
125            repository.with_session_connection(
126                "admit retained semantic records",
127                |db| admit(db, captured.stored.revision, &captured.studies, &intent),
128                |db| admit(db, captured.stored.revision, &captured.studies, &intent),
129            )
130        })
131    }
132}
133fn capture<E: MetadataExecutor>(
134    db: &mut E,
135    selection: &str,
136    budgets: &DisclosureBudgets,
137) -> Result<CapturedSemanticRecords, PgMetaError>
138where
139    E::Row: SchemaStatusRow,
140{
141    db.query_one(
142        "bound semantic readback",
143        "SELECT set_config('statement_timeout',$1,TRUE)",
144        &[&(u64::from(budgets.total_seconds.min(600)) * 1000)
145            .max(1)
146            .to_string()],
147    )?;
148    let rows = i64::try_from(budgets.rows.min(10_000)).unwrap_or(10_000);
149    let bytes = budgets.payload_bytes.min(1024 * 1024) as i64;
150    let row = db.query_one(
151        "capture semantic records",
152        "SELECT public.tre_capture_semantic_records($1::text::jsonb,$2,$3)::text",
153        &[&selection, &rows, &bytes],
154    )?;
155    serde_json::from_str(&row.required_text(0, "semantic records")?)
156        .map(|stored| CapturedSemanticRecords {
157            stored,
158            studies: Vec::new(),
159        })
160        .map_err(|_| PgMetaError::Decode {
161            field: "semantic records",
162            value: "<redacted>".into(),
163            message: "invalid semantic capture".into(),
164        })
165}
166fn admit<E: MetadataExecutor>(
167    db: &mut E,
168    revision: i64,
169    studies: &[uuid::Uuid],
170    intent: &DisclosureIntent,
171) -> Result<DisclosureSnapshot, PgMetaError>
172where
173    E::Row: SchemaStatusRow,
174{
175    db.query_one(
176        "retain semantic records",
177        "SELECT public.tre_lock_semantic_state($1)",
178        &[&revision],
179    )?;
180    for study in studies {
181        let row = db.query_one("retain semantic Study authorization", "SELECT public.tre_user_has_study_access(study_id)::text FROM public.studies WHERE study_id=$1 FOR SHARE", &[study])?;
182        if row.required_text(0, "Study authorization")? != "true" {
183            return Err(PgMetaError::Decode {
184                field: "semantic Study",
185                value: "<redacted>".into(),
186                message: "Study is not visible".into(),
187            });
188        }
189    }
190    crate::governance::admit_disclosure(db, intent)
191}
192fn unavailable() -> CoreError {
193    CoreError::Infrastructure("Semantic readback is unavailable".into())
194}
195
196impl PgMetadataRepository<'_> {
197    /// Hold Study/version authorization until the enclosing mutation commits,
198    /// including provenance attached to an existing no-op mapping or link.
199    pub fn retain_semantic_version_access(
200        &self,
201        versions: &[ahri_tre_types::VersionId],
202    ) -> Result<(), CoreError> {
203        if versions.len() > 256 {
204            return Err(unavailable());
205        }
206        self.with_session_connection(
207            "retain semantic version authorization",
208            |db| retain_versions(db, versions),
209            |db| retain_versions(db, versions),
210        )
211    }
212
213    /// Manual declarations do not acquire a derived-content origin. Mutations
214    /// and their Transformation links commit together under the retained revision.
215    pub fn with_semantic_record_mutation<T>(
216        &self,
217        resolution: &crate::semantic::SemanticResolution,
218        studies: &[ahri_tre_types::StudyId],
219        budgets: &DisclosureBudgets,
220        write: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
221    ) -> Result<T, CoreError> {
222        self.with_dataset_acceptance(|repository| {
223            repository.with_session_connection(
224                "begin semantic mutation",
225                |db| begin_mutation(db, resolution.revision(), studies, budgets),
226                |db| begin_mutation(db, resolution.revision(), studies, budgets),
227            )?;
228            write(repository)
229        })
230    }
231}
232fn begin_mutation<E: MetadataExecutor>(
233    db: &mut E,
234    revision: i64,
235    studies: &[ahri_tre_types::StudyId],
236    budgets: &DisclosureBudgets,
237) -> Result<(), PgMetaError>
238where
239    E::Row: SchemaStatusRow,
240{
241    let milliseconds = (u64::from(budgets.total_seconds.min(600)) * 1000)
242        .max(1)
243        .to_string();
244    db.query_one(
245        "bound semantic mutation",
246        "SELECT set_config('statement_timeout',$1,TRUE),set_config('transaction_timeout',$1,TRUE)",
247        &[&milliseconds],
248    )?;
249    db.query_one(
250        "serialize semantic mutation",
251        "SELECT public.tre_begin_semantic_record_mutation($1)",
252        &[&revision],
253    )?;
254    for study in studies {
255        let row=db.query_one("retain mutation Study authorization","SELECT public.tre_user_has_study_access(study_id)::text FROM public.studies WHERE study_id=$1 FOR SHARE",&[&study.0])?;
256        if row.required_text(0, "Study authorization")? != "true" {
257            return Err(PgMetaError::Decode {
258                field: "semantic Study",
259                value: "<redacted>".into(),
260                message: "Study is not visible".into(),
261            });
262        }
263    }
264    Ok(())
265}
266
267fn retain_versions<E: MetadataExecutor>(
268    db: &mut E,
269    versions: &[ahri_tre_types::VersionId],
270) -> Result<(), PgMetaError>
271where
272    E::Row: SchemaStatusRow,
273{
274    for version in versions {
275        let row = db.query_one("retain semantic version authorization", "SELECT public.tre_user_has_study_access(s.study_id)::text FROM public.asset_versions v JOIN public.assets a USING(asset_id) JOIN public.studies s USING(study_id) WHERE v.version_id=$1 FOR SHARE OF v,a,s", &[&version.0])?;
276        if row.required_text(0, "version authorization")? != "true" {
277            return Err(PgMetaError::Decode {
278                field: "semantic version",
279                value: "<redacted>".into(),
280                message: "version is not visible".into(),
281            });
282        }
283    }
284    Ok(())
285}