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