1use crate::{MetadataExecutor, PgMetaError, PgMetadataRepository, schema::SchemaStatusRow};
4use ahri_tre_core::CoreError;
5use ahri_tre_core::VocabularyRepository;
6use ahri_tre_types::VocabularyId;
7use futures::FutureExt;
8use serde::{Deserialize, Serialize};
9
10pub type SemanticVocabulary = (
11 ahri_tre_types::VocabularyRecord,
12 Vec<ahri_tre_types::VocabularyItemRecord>,
13);
14
15pub struct CapturedVocabulary {
18 snapshot: ahri_tre_types::DisclosureSnapshot,
19 vocabulary: ahri_tre_types::VocabularyRecord,
20 items: Vec<ahri_tre_types::VocabularyItemRecord>,
21}
22
23impl CapturedVocabulary {
24 pub fn into_parts(
25 self,
26 ) -> (
27 ahri_tre_types::DisclosureSnapshot,
28 ahri_tre_types::VocabularyRecord,
29 Vec<ahri_tre_types::VocabularyItemRecord>,
30 ) {
31 (self.snapshot, self.vocabulary, self.items)
32 }
33}
34
35#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
37#[serde(tag = "kind", rename_all = "snake_case")]
38pub enum SemanticOrigin {
39 Unknown,
40 Declared,
41 Derived { inputs: Vec<uuid::Uuid> },
42}
43
44#[derive(Debug, Clone)]
45pub enum SemanticTarget {
46 Vocabulary(VocabularyId),
47 EntityInstance(i64),
48}
49
50impl PgMetadataRepository<'_> {
51 pub fn capture_semantic_vocabulary(
54 &self,
55 intent: &ahri_tre_types::DisclosureIntent,
56 target: VocabularyId,
57 ) -> Result<CapturedVocabulary, CoreError> {
58 let (snapshot, mut values) = self.capture_semantic_vocabularies(intent, &[target])?;
59 let (vocabulary, items) = values.pop().ok_or_else(unavailable)?;
60 Ok(CapturedVocabulary {
61 snapshot,
62 vocabulary,
63 items,
64 })
65 }
66
67 pub fn capture_semantic_vocabularies(
69 &self,
70 intent: &ahri_tre_types::DisclosureIntent,
71 targets: &[VocabularyId],
72 ) -> Result<(ahri_tre_types::DisclosureSnapshot, Vec<SemanticVocabulary>), CoreError> {
73 if intent.lifetime_seconds == 0
74 || intent.lifetime_seconds > 600
75 || intent.max_bytes == 0
76 || intent.max_rows == 0
77 || targets.is_empty()
78 || targets.len() > 128
79 {
80 return Err(unavailable());
81 }
82 self.with_dataset_acceptance(|repository| {
83 let mut inputs = std::collections::BTreeSet::new();
84 for target in targets {
85 repository.with_session_connection(
86 "capture semantic fields",
87 |db| lock_vocabulary(db, *target, intent),
88 |db| lock_vocabulary(db, *target, intent),
89 )?;
90 let SemanticOrigin::Derived { inputs: sources } =
91 repository.semantic_origin(&SemanticTarget::Vocabulary(*target))?
92 else {
93 return Err(unavailable());
94 };
95 inputs.extend(sources);
96 }
97 let mut intent = intent.clone();
98 intent.inputs = inputs.into_iter().collect();
99 intent.latest_inputs.clear();
100 let snapshot = repository.with_session_connection(
101 "admit semantic content",
102 |db| crate::governance::admit_disclosure(db, &intent),
103 |db| crate::governance::admit_disclosure(db, &intent),
104 )?;
105 let mut values = Vec::new();
106 let mut rows = 0;
107 for target in targets {
108 let vocabulary = repository
109 .get_vocabulary(*target)
110 .now_or_never()
111 .ok_or_else(unavailable)??
112 .ok_or_else(unavailable)?;
113 let items = repository
114 .list_vocabulary_items(*target)
115 .now_or_never()
116 .ok_or_else(unavailable)??;
117 rows += items.len() as u64;
118 if rows > intent.max_rows {
119 return Err(unavailable());
120 }
121 values.push((vocabulary, items));
122 }
123 Ok((snapshot, values))
124 })
125 }
126
127 pub fn catalogue_vocabulary_items(
130 &self,
131 target: VocabularyId,
132 ) -> Result<Option<Vec<ahri_tre_types::VocabularyItemRecord>>, CoreError> {
133 self.with_dataset_acceptance(|repository| {
134 if repository.semantic_origin(&SemanticTarget::Vocabulary(target))?
135 != SemanticOrigin::Declared
136 {
137 return Ok(None);
138 }
139 repository
140 .with_session_connection(
141 "capture declared dictionary schema",
142 |db| declared_vocabulary_items(db, target),
143 |db| declared_vocabulary_items(db, target),
144 )
145 .map(Some)
146 })
147 }
148
149 pub fn retain_vocabulary_contribution(&self, target: VocabularyId) -> Result<(), CoreError> {
151 self.with_session_connection(
152 "retain vocabulary contribution",
153 |db| contribution(db, target),
154 |db| contribution(db, target),
155 )
156 }
157
158 pub fn with_semantic_derivation<T>(
159 &self,
160 inputs: &[ahri_tre_types::VersionId],
161 write: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
162 ) -> Result<T, CoreError> {
163 self.with_dataset_acceptance(|repository| {
164 repository.begin_semantic_derivation(inputs)?;
165 write(repository)
166 })
167 }
168
169 pub fn begin_semantic_derivation(
172 &self,
173 inputs: &[ahri_tre_types::VersionId],
174 ) -> Result<(), CoreError> {
175 use crate::repository::PgMetadataRepositoryConnection;
176 if !matches!(
177 &self.connection,
178 PgMetadataRepositoryConnection::ScopedDirect(_)
179 | PgMetadataRepositoryConnection::ScopedOAuth(_)
180 ) {
181 return Err(CoreError::Conflict(
182 "Semantic provenance requires a mutation transaction".into(),
183 ));
184 }
185 let inputs = inputs.iter().map(|id| id.0).collect::<Vec<_>>();
186 self.with_session_connection(
187 "retain semantic sources",
188 |db| begin_derivation(db, &inputs),
189 |db| begin_derivation(db, &inputs),
190 )
191 }
192
193 pub fn with_semantic_declarations<T>(
196 &self,
197 write: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
198 ) -> Result<T, CoreError> {
199 self.with_dataset_acceptance(|repository| {
200 repository.with_session_connection(
201 "declare dictionary input",
202 begin_declaration,
203 begin_declaration,
204 )?;
205 write(repository)
206 })
207 }
208
209 pub fn semantic_origin(&self, target: &SemanticTarget) -> Result<SemanticOrigin, CoreError> {
210 let (statement, id) = match target {
211 SemanticTarget::Vocabulary(id) => {
212 ("SELECT public.tre_vocabulary_origin($1)::text", id.0)
213 }
214 SemanticTarget::EntityInstance(id) => {
215 ("SELECT public.tre_entity_instance_origin($1)::text", *id)
216 }
217 };
218 self.with_session_connection(
219 "read semantic provenance",
220 |db| read_origin(db, statement, id),
221 |db| read_origin(db, statement, id),
222 )
223 }
224}
225
226fn unavailable() -> CoreError {
227 CoreError::Infrastructure("Semantic provenance or admission is unavailable".into())
228}
229
230fn lock_vocabulary<E: MetadataExecutor>(
231 db: &mut E,
232 target: VocabularyId,
233 intent: &ahri_tre_types::DisclosureIntent,
234) -> Result<(), PgMetaError> {
235 let milliseconds = i64::from(intent.lifetime_seconds) * 1000;
236 let rows = i64::try_from(intent.max_rows).unwrap_or(i64::MAX);
237 let bytes = intent.max_bytes.min(1024 * 1024) as i64;
238 db.query_one(
239 "bound semantic capture deadline",
240 "SELECT set_config('statement_timeout',$1,TRUE)",
241 &[&milliseconds.to_string()],
242 )?;
243 db.query_one(
244 "capture bounded semantic fields",
245 "SELECT public.tre_lock_semantic_vocabulary($1,$2,$3,$4)",
246 &[&target.0, &milliseconds, &rows, &bytes],
247 )?;
248 Ok(())
249}
250
251fn begin_derivation<E: MetadataExecutor>(
252 db: &mut E,
253 inputs: &[uuid::Uuid],
254) -> Result<(), PgMetaError> {
255 db.query_one(
256 "retain semantic sources",
257 "SELECT public.tre_semantic_derivation($1)",
258 &[&inputs],
259 )?;
260 Ok(())
261}
262
263fn begin_declaration<E: MetadataExecutor>(db: &mut E) -> Result<(), PgMetaError> {
264 db.query_one(
265 "declare dictionary input",
266 "SELECT public.tre_semantic_declaration()",
267 &[],
268 )?;
269 Ok(())
270}
271
272fn read_origin<E: MetadataExecutor>(
273 db: &mut E,
274 statement: &str,
275 id: i64,
276) -> Result<SemanticOrigin, PgMetaError>
277where
278 E::Row: SchemaStatusRow,
279{
280 let row = db.query_one("read semantic provenance", statement, &[&id])?;
281 serde_json::from_str(&row.required_text(0, "semantic origin")?).map_err(|_| {
282 PgMetaError::Decode {
283 field: "semantic origin",
284 value: "<redacted>".into(),
285 message: "invalid semantic provenance".into(),
286 }
287 })
288}
289
290fn contribution<E: MetadataExecutor>(db: &mut E, target: VocabularyId) -> Result<(), PgMetaError> {
291 db.query_one(
292 "retain vocabulary contribution",
293 "SELECT public.tre_semantic_vocabulary_contribution($1)",
294 &[&target.0],
295 )?;
296 Ok(())
297}
298
299pub fn declared_vocabulary_items<E: MetadataExecutor>(
302 db: &mut E,
303 target: VocabularyId,
304) -> Result<Vec<ahri_tre_types::VocabularyItemRecord>, PgMetaError>
305where
306 E::Row: SchemaStatusRow,
307{
308 let row = db.query_one(
309 "capture declared vocabulary",
310 "SELECT public.tre_declared_vocabulary_items($1)::text",
311 &[&target.0],
312 )?;
313 serde_json::from_str(&row.required_text(0, "declared vocabulary")?).map_err(|_| {
314 PgMetaError::Decode {
315 field: "declared vocabulary",
316 value: "<redacted>".into(),
317 message: "invalid declaration capture".into(),
318 }
319 })
320}
321
322#[derive(Debug, Clone)]
324pub enum SemanticScope {
325 Entity {
326 study_id: ahri_tre_types::StudyId,
327 entity_id: ahri_tre_types::EntityId,
328 version_id: ahri_tre_types::VersionId,
329 },
330 Relation {
331 study_id: ahri_tre_types::StudyId,
332 relation_id: ahri_tre_types::EntityRelationId,
333 version_id: ahri_tre_types::VersionId,
334 },
335}
336
337pub struct SemanticResolution {
339 revision: i64,
340}
341impl SemanticResolution {
342 pub(crate) fn revision(&self) -> i64 {
343 self.revision
344 }
345}
346
347#[derive(Clone, Copy)]
348pub enum SemanticInstance {
349 Entity(i64),
350 Relation(i64),
351}
352
353#[derive(Debug)]
356pub struct SemanticState {
357 scope: SemanticScope,
358 revision: i64,
359 inputs: Vec<uuid::Uuid>,
360}
361impl SemanticState {
362 pub fn inputs(&self) -> &[uuid::Uuid] {
363 &self.inputs
364 }
365 pub fn scope(&self) -> &SemanticScope {
366 &self.scope
367 }
368 pub fn revision(&self) -> i64 {
369 self.revision
370 }
371}
372
373impl PgMetadataRepository<'_> {
374 pub fn begin_semantic_resolution(&self) -> Result<SemanticResolution, CoreError> {
375 self.with_session_connection(
376 "begin semantic resolution",
377 resolution_revision,
378 resolution_revision,
379 )
380 }
381
382 pub fn retain_instance_contribution(
383 &self,
384 instance: SemanticInstance,
385 ) -> Result<(), CoreError> {
386 let (collection, id) = match instance {
387 SemanticInstance::Entity(id) => ("entity_instance", id),
388 SemanticInstance::Relation(id) => ("relation_instances", id),
389 };
390 self.with_session_connection(
391 "retain instance contribution",
392 |db| instance_contribution(db, collection, id),
393 |db| instance_contribution(db, collection, id),
394 )
395 }
396
397 pub fn capture_semantic_state(
398 &self,
399 resolution: &SemanticResolution,
400 scope: &SemanticScope,
401 limits: &ahri_tre_types::DisclosureBudgets,
402 ) -> Result<SemanticState, CoreError> {
403 self.with_dataset_acceptance(|repository| {
404 let state = repository.with_session_connection(
405 "capture semantic state",
406 |db| capture_state(db, scope, limits),
407 |db| capture_state(db, scope, limits),
408 )?;
409 if state.revision != resolution.revision {
410 return Err(unavailable());
411 }
412 Ok(state)
413 })
414 }
415}
416fn resolution_revision<E: MetadataExecutor>(db: &mut E) -> Result<SemanticResolution, PgMetaError>
417where
418 E::Row: SchemaStatusRow,
419{
420 let row = db.query_one(
421 "begin semantic resolution",
422 "SELECT public.tre_semantic_revision()::text",
423 &[],
424 )?;
425 let revision = row
426 .required_text(0, "semantic revision")?
427 .parse()
428 .map_err(|_| PgMetaError::Decode {
429 field: "semantic revision",
430 value: "<redacted>".into(),
431 message: "invalid revision".into(),
432 })?;
433 Ok(SemanticResolution { revision })
434}
435fn instance_contribution<E: MetadataExecutor>(
436 db: &mut E,
437 collection: &str,
438 id: i64,
439) -> Result<(), PgMetaError> {
440 db.query_one(
441 "retain instance contribution",
442 "SELECT public.tre_semantic_instance_contribution($1,$2)",
443 &[&collection, &id],
444 )?;
445 Ok(())
446}
447
448fn capture_state<E: MetadataExecutor>(
449 db: &mut E,
450 scope: &SemanticScope,
451 limits: &ahri_tre_types::DisclosureBudgets,
452) -> Result<SemanticState, PgMetaError>
453where
454 E::Row: SchemaStatusRow,
455{
456 let (study, definition, relation, version) = match scope {
457 SemanticScope::Entity {
458 study_id,
459 entity_id,
460 version_id,
461 } => (study_id.0, entity_id.0, false, version_id.0),
462 SemanticScope::Relation {
463 study_id,
464 relation_id,
465 version_id,
466 } => (study_id.0, relation_id.0, true, version_id.0),
467 };
468 if limits.total_seconds == 0 || limits.total_seconds > 600 {
469 return Err(PgMetaError::Decode {
470 field: "semantic limits",
471 value: "<redacted>".into(),
472 message: "invalid deadline".into(),
473 });
474 }
475 db.query_one(
476 "bound semantic state capture",
477 "SELECT set_config('statement_timeout',$1,TRUE)",
478 &[&(u64::from(limits.total_seconds) * 1000).to_string()],
479 )?;
480 let rows = i64::try_from(limits.rows).unwrap_or(i64::MAX);
481 let bytes = limits.payload_bytes.min(1024 * 1024) as i64;
482 let row = db.query_one(
483 "capture semantic state",
484 "SELECT public.tre_capture_semantic_state($1,$2,$3,$4,$5,$6)::text",
485 &[&study, &definition, &relation, &version, &rows, &bytes],
486 )?;
487 #[derive(Deserialize)]
488 struct Capture {
489 revision: i64,
490 inputs: Vec<uuid::Uuid>,
491 }
492 let capture: Capture =
493 serde_json::from_str(&row.required_text(0, "semantic state")?).map_err(|_| {
494 PgMetaError::Decode {
495 field: "semantic state",
496 value: "<redacted>".into(),
497 message: "invalid semantic capture".into(),
498 }
499 })?;
500 Ok(SemanticState {
501 scope: scope.clone(),
502 revision: capture.revision,
503 inputs: capture.inputs,
504 })
505}
506
507pub fn verify_semantic_state<E: crate::MetadataTransaction + ?Sized>(
509 db: &mut E,
510 state: &SemanticState,
511) -> Result<(), PgMetaError> {
512 db.query_one(
513 "retain semantic snapshot",
514 "SELECT public.tre_lock_semantic_state($1)",
515 &[&state.revision],
516 )?;
517 Ok(())
518}
519impl PgMetadataRepository<'_> {
520 pub fn with_admitted_semantic_mutation<T>(
523 &self,
524 state: &SemanticState,
525 admission: uuid::Uuid,
526 write: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
527 ) -> Result<T, CoreError> {
528 self.with_dataset_acceptance(|repository| {
529 repository.with_session_connection(
530 "authorize semantic mutation",
531 |db| mutation_context(db, state, admission),
532 |db| mutation_context(db, state, admission),
533 )?;
534 write(repository)
535 })
536 }
537}
538fn mutation_context<E: MetadataExecutor>(
539 db: &mut E,
540 state: &SemanticState,
541 admission: uuid::Uuid,
542) -> Result<(), PgMetaError> {
543 db.query_one(
544 "authorize semantic mutation",
545 "SELECT public.tre_semantic_mutation_context($1,$2,$3)",
546 &[&admission, &state.revision, &state.inputs],
547 )?;
548 Ok(())
549}
550
551pub enum SemanticMappings {
552 Entity(Vec<ahri_tre_types::StudyEntityInstanceRecord>),
553 Relation(Vec<ahri_tre_types::StudyRelationInstanceRecord>),
554}
555impl PgMetadataRepository<'_> {
556 pub fn retained_semantic_mappings(
559 &self,
560 state: &SemanticState,
561 ) -> Result<SemanticMappings, CoreError> {
562 use crate::repository::PgMetadataRepositoryConnection;
563 if !matches!(
564 self.connection,
565 PgMetadataRepositoryConnection::ScopedDirect(_)
566 | PgMetadataRepositoryConnection::ScopedOAuth(_)
567 ) {
568 return Err(unavailable());
569 }
570 self.with_session_connection(
571 "read retained semantic mappings",
572 |db| read_mappings(db, state),
573 |db| read_mappings(db, state),
574 )
575 }
576}
577fn read_mappings<E: MetadataExecutor>(
578 db: &mut E,
579 state: &SemanticState,
580) -> Result<SemanticMappings, PgMetaError>
581where
582 E::Row: SchemaStatusRow,
583{
584 db.query_one(
585 "retain semantic mapping revision",
586 "SELECT public.tre_lock_semantic_state($1)",
587 &[&state.revision],
588 )?;
589 let (study, definition, sql, relation) = match state.scope {
590 SemanticScope::Entity {
591 study_id,
592 entity_id,
593 ..
594 } => (
595 study_id.0,
596 entity_id.0,
597 "SELECT to_jsonb(m)::text FROM public.study_entity_instances m WHERE study_id=$1 AND entity_id::bigint=$2 ORDER BY entity_instance_id",
598 false,
599 ),
600 SemanticScope::Relation {
601 study_id,
602 relation_id,
603 ..
604 } => (
605 study_id.0,
606 relation_id.0,
607 "SELECT (to_jsonb(m)-'entityrelation_id' || jsonb_build_object('entity_relation_id',m.entityrelation_id))::text FROM public.study_relation_instances m WHERE study_id=$1 AND entityrelation_id::bigint=$2 ORDER BY relation_instance_id",
608 true,
609 ),
610 };
611 let rows = db.query_many(
612 "read retained semantic mappings",
613 sql,
614 &[&study, &definition],
615 )?;
616 let mut payloads = Vec::new();
617 for row in rows {
618 payloads.push(row.required_text(0, "semantic mapping")?);
619 }
620 let decode_error = |_| PgMetaError::Decode {
621 field: "semantic mapping",
622 value: "<redacted>".into(),
623 message: "invalid stored mapping".into(),
624 };
625 if relation {
626 Ok(SemanticMappings::Relation(
627 payloads
628 .iter()
629 .map(|v| serde_json::from_str(v).map_err(decode_error))
630 .collect::<Result<_, _>>()?,
631 ))
632 } else {
633 Ok(SemanticMappings::Entity(
634 payloads
635 .iter()
636 .map(|v| serde_json::from_str(v).map_err(decode_error))
637 .collect::<Result<_, _>>()?,
638 ))
639 }
640}