ahri_tre_lake/
governance_ledger.rs1use ahri_tre_types::StudyId;
4use duckdb::{Connection, OptionalExt};
5use uuid::Uuid;
6
7use crate::{LakeError, dataset_loading::dataset_table_relation};
8
9pub struct GovernanceLedger<'connection> {
10 connection: &'connection mut Connection,
11 relation: String,
12}
13
14impl<'connection> GovernanceLedger<'connection> {
15 pub fn provision(
17 connection: &'connection mut Connection,
18 study: StudyId,
19 ) -> Result<Self, LakeError> {
20 require_persistent_lake(connection)?;
21 let (schema, relation) = relation(study);
22 let tx = connection.transaction()?;
23 tx.execute_batch(&format!("CREATE SCHEMA IF NOT EXISTS {schema}; CREATE TABLE IF NOT EXISTS {relation} (event_id VARCHAR NOT NULL, evidence VARCHAR NOT NULL)"))?;
24 tx.commit()?;
25 Self::open(connection, study)
26 }
27
28 pub fn open(
30 connection: &'connection mut Connection,
31 study: StudyId,
32 ) -> Result<Self, LakeError> {
33 require_persistent_lake(connection)?;
34 let (_, relation) = relation(study);
35 connection.prepare(&format!("SELECT event_id,evidence FROM {relation} LIMIT 0"))?;
36 Ok(Self {
37 connection,
38 relation,
39 })
40 }
41
42 pub fn append(&mut self, event: Uuid, evidence: &str) -> Result<(), LakeError> {
46 if evidence.len() > 65536
47 || !serde_json::from_str::<serde_json::Value>(evidence)
48 .is_ok_and(|value| value.is_object())
49 {
50 return Err(LakeError::GovernanceEvidenceConflict);
51 }
52 let tx = self.connection.transaction()?;
53 let existing: Option<String> = tx
54 .query_row(
55 &format!(
56 "SELECT evidence FROM {} WHERE event_id=? LIMIT 2",
57 self.relation
58 ),
59 [event.to_string()],
60 |row| row.get(0),
61 )
62 .optional()?;
63 match existing {
64 Some(existing) if existing != evidence => {
65 return Err(LakeError::GovernanceEvidenceConflict);
66 }
67 Some(_) => {}
68 None => {
69 tx.execute(
70 &format!("INSERT INTO {} VALUES (?,?)", self.relation),
71 [event.to_string(), evidence.to_string()],
72 )?;
73 }
74 }
75 tx.commit()?;
76 if self.read(event)?.as_deref() != Some(evidence) {
79 return Err(LakeError::GovernanceEvidenceConflict);
80 }
81 Ok(())
82 }
83
84 pub fn read(&self, event: Uuid) -> Result<Option<String>, LakeError> {
86 let mut statement = self.connection.prepare(&format!(
87 "SELECT evidence FROM {} WHERE event_id=? LIMIT 2",
88 self.relation
89 ))?;
90 let rows = statement
91 .query_map([event.to_string()], |row| row.get::<_, String>(0))?
92 .collect::<Result<Vec<_>, _>>()?;
93 if rows.len() > 1 {
94 return Err(LakeError::GovernanceEvidenceConflict);
95 }
96 Ok(rows.into_iter().next())
97 }
98}
99
100fn relation(study: StudyId) -> (String, String) {
101 let relation = dataset_table_relation(study, "Governance_provenance_ledger", 1, 0, 0);
102 let schema = format!("\"{}\"", relation.schema_name.replace('"', "\"\""));
103 let table = format!("\"{}\"", relation.table_name.replace('"', "\"\""));
104 let schema = format!("\"{}\".{schema}", crate::LAKE_ALIAS);
105 (schema.clone(), format!("{schema}.{table}"))
106}
107
108fn require_persistent_lake(connection: &Connection) -> Result<(), LakeError> {
109 let persistent:bool=connection.query_row(
110 "SELECT count(*)=1 FROM duckdb_databases() WHERE database_name=? AND type='ducklake' AND NOT readonly",
111 [crate::LAKE_ALIAS],|row|row.get(0))?;
112 if !persistent {
113 return Err(LakeError::GovernanceEvidenceConflict);
114 }
115 connection.query_row(
118 "SELECT count(*) FROM ducklake_snapshots(?)",
119 [crate::LAKE_ALIAS],
120 |row| row.get::<_, i64>(0),
121 )?;
122 Ok(())
123}