Skip to main content

ahri_tre_lake/
governance_ledger.rs

1//! System-managed Governance provenance ledger Dataset. No update, delete,
2//! arbitrary relation, or content-query interface is exposed by this adapter.
3use 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    /// Only datastore provisioning calls this, before publishing readiness.
16    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    /// Ordinary recovery never creates a missing ledger or guesses its origin.
29    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    /// The metadata adapter retains the event's row lock across this call and
43    /// its acknowledgement. An ambiguous commit is retried with the same event
44    /// identity and exact evidence; a conflicting retry never overwrites it.
45    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        // A successful send/execute is insufficient: read back the committed
77        // event before its outbox acknowledgement becomes durable.
78        if self.read(event)?.as_deref() != Some(evidence) {
79            return Err(LakeError::GovernanceEvidenceConflict);
80        }
81        Ok(())
82    }
83
84    /// Operator/reconciliation read; scoped custodian projection is app-owned.
85    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    // Querying snapshots also verifies that the DuckLake extension can reach
116    // its committed catalog; an ordinary or in-memory DuckDB is never a ledger.
117    connection.query_row(
118        "SELECT count(*) FROM ducklake_snapshots(?)",
119        [crate::LAKE_ALIAS],
120        |row| row.get::<_, i64>(0),
121    )?;
122    Ok(())
123}