Skip to main content

ahri_tre_pgmeta/
governance.rs

1//! Guarded Asset classification and operator-owned preserving conversion.
2//! Both direct and libpq Sessions use the same transaction/function boundary.
3use crate::{MetadataExecutor, PgMetaError, schema::SchemaStatusRow};
4use ahri_tre_types::AssetClassification;
5pub use ahri_tre_types::{DisclosureIntent, DisclosureSnapshot};
6use uuid::Uuid;
7
8const STRUCTURE: &str = include_str!("governance_schema.sql");
9
10/// Narrow governance maintenance over administrator authority already resolved
11/// for the Session. No configuration documents or Secret stores are reopened.
12#[derive(Clone)]
13pub struct GovernanceMaintenance {
14    repository: crate::PgMetadataRepository<'static>,
15}
16impl crate::PgDatasetMaintenance {
17    pub fn governance(&self) -> GovernanceMaintenance {
18        GovernanceMaintenance {
19            repository: self.repository.clone(),
20        }
21    }
22}
23impl GovernanceMaintenance {
24    pub fn bind_dataset(
25        &self,
26        location: &ahri_tre_types::GovernedDatasetLocation,
27    ) -> Result<(), ahri_tre_core::CoreError> {
28        self.repository.with_session_connection(
29            "bind immutable Dataset",
30            |db| bind_dataset_location(db, location),
31            |db| bind_dataset_location(db, location),
32        )
33    }
34    pub fn binding(&self) -> Result<GovernanceLedgerBinding, ahri_tre_core::CoreError> {
35        self.repository.with_session_connection(
36            "Governance binding",
37            ledger_binding,
38            ledger_binding,
39        )
40    }
41    pub fn activate_identity(
42        &self,
43        issuer: &str,
44        subject: &str,
45        principal: &str,
46    ) -> Result<(), ahri_tre_core::CoreError> {
47        self.repository.with_session_connection(
48            "activate Governance identity",
49            |db| activate_identity(db, issuer, subject, principal),
50            |db| activate_identity(db, issuer, subject, principal),
51        )
52    }
53    pub fn pending(&self, limit: u32) -> Result<Vec<Uuid>, ahri_tre_core::CoreError> {
54        self.repository.with_session_connection(
55            "pending Governance evidence",
56            |db| pending_evidence(db, limit),
57            |db| pending_evidence(db, limit),
58        )
59    }
60    pub fn project(
61        &self,
62        event: Uuid,
63        append: impl FnMut(Uuid, &str) -> Result<(), ()>,
64    ) -> Result<bool, ahri_tre_core::CoreError> {
65        let append = std::cell::RefCell::new(append);
66        self.repository.with_session_connection(
67            "project Governance evidence",
68            |db| project_evidence(db, event, |id, text| (append.borrow_mut())(id, text)),
69            |db| project_evidence(db, event, |id, text| (append.borrow_mut())(id, text)),
70        )
71    }
72    pub fn unfinished(
73        &self,
74        limit: u32,
75    ) -> Result<Vec<UnfinishedDisclosure>, ahri_tre_core::CoreError> {
76        self.repository.with_session_connection(
77            "unfinished disclosures",
78            |db| unfinished_disclosures(db, limit),
79            |db| unfinished_disclosures(db, limit),
80        )
81    }
82    pub fn abandon(&self, receipt: &UnfinishedDisclosure) -> Result<(), ahri_tre_core::CoreError> {
83        self.repository.with_session_connection(
84            "abandon disclosure",
85            |db| abandon_disclosure(db, receipt),
86            |db| abandon_disclosure(db, receipt),
87        )
88    }
89}
90fn activate_identity<E: MetadataExecutor>(
91    db: &mut E,
92    issuer: &str,
93    subject: &str,
94    principal: &str,
95) -> Result<(), PgMetaError> {
96    db.query_one(
97        "activate Governance identity",
98        "SELECT tre_governance.activate_identity($1,$2,$3)",
99        &[&issuer, &subject, &principal],
100    )?;
101    Ok(())
102}
103
104pub(crate) fn validate_derivation<E: MetadataExecutor>(
105    db: &mut E,
106    inputs: &[Uuid],
107    risk: ahri_tre_types::AssetRisk,
108    inherit: bool,
109) -> Result<(), PgMetaError> {
110    db.query_one(
111        "validate derivation",
112        "SELECT public.tre_validate_derivation($1,$2,$3)",
113        &[&inputs, &risk.as_str(), &inherit],
114    )?;
115    Ok(())
116}
117
118impl crate::PgMetadataRepository<'_> {
119    /// Pre-connection destination authority; final admission rechecks current state.
120    pub fn authorize_acquisition_destination(
121        &self,
122        study: ahri_tre_types::StudyId,
123        classification: &ahri_tre_types::IngestClassification,
124    ) -> Result<(), ahri_tre_core::CoreError> {
125        self.with_session_connection(
126            "authorize acquisition destination",
127            |db| acquisition_destination(db, study.0, classification),
128            |db| acquisition_destination(db, study.0, classification),
129        )
130    }
131    /// Authorizes the one freshly ingested Datafile pinned by the App upload
132    /// workflow. This does not confer trusted or arbitrary derivation authority.
133    pub fn validate_uploaded_table_source(
134        &self,
135        source: ahri_tre_types::VersionId,
136        study: ahri_tre_types::StudyId,
137        risk: ahri_tre_types::AssetRisk,
138    ) -> Result<(), ahri_tre_core::CoreError> {
139        self.with_session_connection(
140            "validate uploaded table source",
141            |db| uploaded_table_source(db, source.0, study.0, risk),
142            |db| uploaded_table_source(db, source.0, study.0, risk),
143        )
144    }
145}
146fn acquisition_destination<E: MetadataExecutor>(
147    db: &mut E,
148    study: Uuid,
149    classification: &ahri_tre_types::IngestClassification,
150) -> Result<(), PgMetaError>
151where
152    E::Row: SchemaStatusRow,
153{
154    let row = db.query_one("authorize acquisition destination",
155        "SELECT (EXISTS (SELECT 1 FROM public.studies WHERE study_id=$1) AND public.tre_user_has_study_access($1) AND ($2 <> 'low' OR (public.tre_user_can_administer_study_access($1) AND length(btrim($3)) > 0)))::text",
156        &[&study, &classification.risk.as_str(), &classification.justification])?;
157    if row.required_text(0, "permitted")? != "true" {
158        return Err(invalid("Acquisition destination is unavailable"));
159    }
160    Ok(())
161}
162fn uploaded_table_source<E: MetadataExecutor>(
163    db: &mut E,
164    source: Uuid,
165    study: Uuid,
166    risk: ahri_tre_types::AssetRisk,
167) -> Result<(), PgMetaError>
168where
169    E::Row: SchemaStatusRow,
170{
171    // The existing function locks and authorizes input identities. HIGH here is
172    // its conservative read-authority rule, not the output's classification.
173    validate_derivation(db, &[source], ahri_tre_types::AssetRisk::High, false)?;
174    let row = db.query_one("verify uploaded table identity",
175        "SELECT count(*)::text FROM public.datafiles d JOIN public.asset_versions v ON v.version_id=d.datafile_id JOIN public.assets a USING(asset_id) WHERE v.version_id=$1 AND a.study_id=$2 AND public.tre_asset_classification(a.asset_id)->>'risk'=$3",
176        &[&source, &study, &risk.as_str()])?;
177    if number(&row, 0)? != 1 {
178        return Err(invalid("uploaded source scope or classification changed"));
179    }
180    Ok(())
181}
182
183/// Explicit appointment. Never inferred from the first Study creator.
184#[derive(Debug, Clone)]
185pub struct GovernanceAppointment {
186    issuers: Vec<String>,
187    subjects: Vec<String>,
188}
189
190impl GovernanceAppointment {
191    pub fn new(issuer: &str, subject: &str) -> Result<Self, PgMetaError> {
192        Self::from_identities([(issuer, subject)])
193    }
194    pub fn from_identities<'a>(
195        identities: impl IntoIterator<Item = (&'a str, &'a str)>,
196    ) -> Result<Self, PgMetaError> {
197        let mut pairs = Vec::new();
198        for (issuer, subject) in identities {
199            if !issuer.starts_with("https://")
200                || issuer.len() > 2048
201                || issuer.bytes().any(|b| b.is_ascii_control())
202                || subject.trim().is_empty()
203                || subject.len() > 255
204                || !subject.is_ascii()
205                || subject.bytes().any(|b| b.is_ascii_control())
206                || pairs.contains(&(issuer, subject))
207            {
208                return Err(invalid("governance appointment is invalid"));
209            }
210            pairs.push((issuer, subject));
211        }
212        if pairs.is_empty() || pairs.len() > 32 {
213            return Err(invalid("explicit governance appointments required"));
214        }
215        let (issuers, subjects) = pairs
216            .into_iter()
217            .map(|(i, s)| (i.to_owned(), s.to_owned()))
218            .unzip();
219        Ok(Self { issuers, subjects })
220    }
221}
222
223pub struct GovernanceUpgrade;
224
225#[derive(Debug)]
226pub struct GovernanceLedgerBinding {
227    pub study_id: Uuid,
228    pub asset_id: Uuid,
229    pub version_id: Uuid,
230}
231
232/// Operator-only binding; the application uses this to open the system ledger
233/// through the Lake adapter. It contains no physical location or credential.
234pub fn ledger_binding<E: MetadataExecutor>(
235    db: &mut E,
236) -> Result<GovernanceLedgerBinding, PgMetaError>
237where
238    E::Row: SchemaStatusRow,
239{
240    let row = db.query_one("governance ledger binding", "SELECT study_id::text,ledger_asset_id::text,ledger_version_id::text FROM tre_governance.state", &[])?;
241    let id = |index| {
242        row.required_text(index, "governance identity")?
243            .parse()
244            .map_err(|_| invalid("invalid governance identity"))
245    };
246    Ok(GovernanceLedgerBinding {
247        study_id: id(0)?,
248        asset_id: id(1)?,
249        version_id: id(2)?,
250    })
251}
252
253#[derive(Debug)]
254pub struct GovernanceUpgradePlan {
255    pub unclassified_assets: u64,
256    pub immutable_versions: u64,
257    pub prepared: bool,
258}
259
260impl GovernanceUpgrade {
261    pub fn plan<E: MetadataExecutor>(db: &mut E) -> Result<GovernanceUpgradePlan, PgMetaError>
262    where
263        E::Row: SchemaStatusRow,
264    {
265        let status = crate::datastore_schema_status(db)?;
266        if !matches!(
267            status.status,
268            crate::DatastoreSchemaCompatibility::Current
269                | crate::DatastoreSchemaCompatibility::Pending
270        ) {
271            return Err(invalid("unsupported metadata baseline"));
272        }
273        let row = db.query_one(
274            "governance structure",
275            "SELECT to_regclass('tre_governance.classifications')::text",
276            &[],
277        )?;
278        let prepared = row.optional_text(0).is_some();
279        let sql = if prepared {
280            "SELECT count(*)::text, (SELECT count(*)::text FROM public.asset_versions) FROM public.assets a WHERE NOT EXISTS (SELECT 1 FROM tre_governance.classifications c WHERE c.asset_id=a.asset_id)"
281        } else {
282            "SELECT count(*)::text, (SELECT count(*)::text FROM public.asset_versions) FROM public.assets"
283        };
284        let row = db.query_one("governance inventory", sql, &[])?;
285        Ok(GovernanceUpgradePlan {
286            unclassified_assets: number(&row, 0)?,
287            immutable_versions: number(&row, 1)?,
288            prepared,
289        })
290    }
291
292    pub fn prepare<E: MetadataExecutor>(
293        db: &mut E,
294        appointment: &GovernanceAppointment,
295    ) -> Result<(), PgMetaError>
296    where
297        E::Row: SchemaStatusRow,
298    {
299        Self::plan(db)?;
300        db.with_transaction("prepare governance", |tx| {
301            tx.query_one(
302                "serialize upgrade",
303                "SELECT pg_advisory_xact_lock(190921,9)",
304                &[],
305            )?;
306            // FORCE RLS normally applies even to the metadata owner. The
307            // privileged operator's bounded DDL transaction temporarily removes
308            // FORCE only on that owner's tables, with AccessExclusive locks.
309            // Restore every original flag before commit; no session ever observes
310            // a committed policy gap, including on an error/rollback.
311            let privileged = tx.query_one("verify upgrade operator", "SELECT (rolsuper OR rolbypassrls)::text FROM pg_catalog.pg_roles WHERE rolname=SESSION_USER", &[])?;
312            if privileged.required_text(0,"operator authority")? != "true" { return Err(invalid("privileged upgrade operator required")); }
313            let forced = tx.query_many("owned forced policies", "SELECT n.nspname::text,c.relname::text FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace WHERE n.nspname='public' AND c.relowner=CURRENT_USER::regrole AND c.relforcerowsecurity ORDER BY c.oid", &[])?;
314            let relations = forced.iter().map(|row| Ok(format!("\"{}\".\"{}\"", row.required_text(0,"schema")?.replace('"',"\"\""), row.required_text(1,"relation")?.replace('"',"\"\"")))).collect::<Result<Vec<_>, PgMetaError>>()?;
315            for relation in &relations { tx.execute_command("lock owner migration table", &format!("ALTER TABLE {relation} NO FORCE ROW LEVEL SECURITY"), &[])?; }
316            tx.execute_command("governance structure", STRUCTURE, &[])?;
317            let owner=tx.query_one("retain metadata owner","SELECT CURRENT_USER::text",&[])?.required_text(0,"metadata owner")?;
318            tx.execute_command("operator-owned migration steps","SET LOCAL ROLE NONE",&[])?;
319            tx.execute_command("preserve isolated lookup and remove inherited ownership",include_str!("governance_operator.sql"),&[])?;
320            tx.execute_command("restore metadata owner",&format!("SET LOCAL ROLE \"{}\"",owner.replace('"',"\"\"")),&[])?;
321
322            tx.execute_command(
323                "disclosure structure",
324                include_str!("disclosure_schema.sql"),
325                &[],
326            )?;
327            tx.execute_command("semantic provenance", include_str!("semantic_provenance.sql"), &[])?;
328            tx.execute_command("semantic instance readbacks", include_str!("semantic_records.sql"), &[])?;
329            tx.query_one(
330                "appoint governance",
331                "SELECT tre_governance.provision_custodians($1,$2)",
332                &[&appointment.issuers, &appointment.subjects],
333            )?;
334            tx.execute_command(
335                "governance policy guards",
336                include_str!("governance_policy.sql"),
337                &[],
338            )?;
339            tx.execute_command("governed ingest admission",include_str!("governance_ingest.sql"),&[])?;
340            tx.execute_command("record governance structure migration","INSERT INTO public.ahri_tre_schema_migrations(version,description,scope) VALUES($1,'Guarded governance and disclosure admission','metadata') ON CONFLICT(version) DO NOTHING",&[&crate::schema::GOVERNANCE_SCHEMA_VERSION])?;
341            tx.execute_command("record semantic provenance migration", "INSERT INTO public.ahri_tre_schema_migrations(version,description,scope) VALUES($1,'Durable semantic field provenance','metadata') ON CONFLICT(version) DO NOTHING", &[&crate::schema::SEMANTIC_PROVENANCE_SCHEMA_VERSION])?;
342            tx.execute_command("record semantic instances migration", "INSERT INTO public.ahri_tre_schema_migrations(version,description,scope) VALUES($1,'Retained instance and link membership histories','metadata') ON CONFLICT(version) DO NOTHING", &[&crate::schema::SEMANTIC_INSTANCES_SCHEMA_VERSION])?;
343            tx.execute_command("record Governance custodians migration", "INSERT INTO public.ahri_tre_schema_migrations(version,description,scope) VALUES($1,'Multiple explicit Governance appointments','metadata') ON CONFLICT(version) DO NOTHING", &[&crate::CURRENT_DATASTORE_SCHEMA_VERSION])?;
344            for relation in relations { tx.execute_command("restore forced policy", &format!("ALTER TABLE {relation} FORCE ROW LEVEL SECURITY"), &[])?; }
345            Ok(())
346        })
347    }
348
349    /// Bounded resumable Q1 conversion. Row locking and insert-only state mean a
350    /// retry cannot overwrite an intervening custodian classification.
351    pub fn classify_existing<E: MetadataExecutor>(
352        db: &mut E,
353        batch: u32,
354    ) -> Result<u64, PgMetaError>
355    where
356        E::Row: SchemaStatusRow,
357    {
358        if batch == 0 || batch > 1000 {
359            return Err(invalid("invalid migration batch"));
360        }
361        let row = db.query_one(
362            "classify existing Assets",
363            "SELECT tre_governance.convert_existing($1)::text",
364            &[&(batch as i32)],
365        )?;
366        number(&row, 0)
367    }
368}
369
370pub fn classification<E: MetadataExecutor>(
371    db: &mut E,
372    asset: Uuid,
373) -> Result<Option<AssetClassification>, PgMetaError>
374where
375    E::Row: SchemaStatusRow,
376{
377    let row = db.query_one(
378        "read Asset classification",
379        "SELECT public.tre_asset_classification($1)::text",
380        &[&asset],
381    )?;
382    row.optional_text(0)
383        .map(|text| {
384            serde_json::from_str(&text).map_err(|_| invalid("invalid classification evidence"))
385        })
386        .transpose()
387}
388
389fn number<R: SchemaStatusRow>(row: &R, index: usize) -> Result<u64, PgMetaError> {
390    row.required_text(index, "count")?
391        .parse()
392        .map_err(|_| invalid("invalid governance count"))
393}
394
395/// The database rechecks current owning-Study custodianship and revision while
396/// holding the same locks used for admission. This is an all-version decision.
397pub fn reclassify<E: MetadataExecutor>(
398    db: &mut E,
399    asset: Uuid,
400    expected_revision: u64,
401    risk: ahri_tre_types::AssetRisk,
402    justification: &str,
403) -> Result<AssetClassification, PgMetaError>
404where
405    E::Row: SchemaStatusRow,
406{
407    let revision = i64::try_from(expected_revision).map_err(|_| invalid("invalid revision"))?;
408    let row = db.query_one(
409        "reclassify Asset",
410        "SELECT public.tre_reclassify_asset($1,$2,$3,$4)::text",
411        &[&asset, &revision, &risk.as_str(), &justification],
412    )?;
413    serde_json::from_str(&row.required_text(0, "classification")?)
414        .map_err(|_| invalid("invalid classification evidence"))
415}
416
417fn invalid(message: &str) -> PgMetaError {
418    PgMetaError::Decode {
419        field: "governance",
420        value: "<redacted>".into(),
421        message: message.into(),
422    }
423}
424
425pub fn admit_disclosure<E: MetadataExecutor>(
426    db: &mut E,
427    intent: &DisclosureIntent,
428) -> Result<DisclosureSnapshot, PgMetaError>
429where
430    E::Row: SchemaStatusRow,
431{
432    let encoded =
433        serde_json::to_string(intent).map_err(|_| invalid("invalid disclosure intent"))?;
434    let row = db.query_one(
435        "admit disclosure",
436        "SELECT public.tre_admit_disclosure($1::text::jsonb)::text",
437        &[&encoded],
438    )?;
439    serde_json::from_str(&row.required_text(0, "disclosure snapshot")?)
440        .map_err(|_| invalid("invalid disclosure snapshot"))
441}
442
443/// Keep the metadata policy/version locks until the Lake adapter has captured
444/// immutable input capabilities. A failed capture rolls the decision back.
445pub fn admit_disclosure_with_inputs<E: MetadataExecutor, T>(
446    db: &mut E,
447    intent: &DisclosureIntent,
448    capture: impl FnOnce(
449        &DisclosureSnapshot,
450        Vec<ahri_tre_types::GovernedDatasetLocation>,
451    ) -> Result<T, ()>,
452) -> Result<(DisclosureSnapshot, T), PgMetaError>
453where
454    E::Row: SchemaStatusRow,
455{
456    admit_disclosure_with_semantic_inputs(db, intent, None, capture)
457}
458
459pub fn admit_disclosure_with_semantic_inputs<E: MetadataExecutor, T>(
460    db: &mut E,
461    intent: &DisclosureIntent,
462    semantic: Option<&crate::semantic::SemanticState>,
463    capture: impl FnOnce(
464        &DisclosureSnapshot,
465        Vec<ahri_tre_types::GovernedDatasetLocation>,
466    ) -> Result<T, ()>,
467) -> Result<(DisclosureSnapshot, T), PgMetaError>
468where
469    E::Row: SchemaStatusRow,
470{
471    let encoded =
472        serde_json::to_string(intent).map_err(|_| invalid("invalid disclosure intent"))?;
473    db.with_transaction("admit and pin disclosure", |tx| {
474        if let Some(state) = semantic {
475            crate::semantic::verify_semantic_state(tx, state)?;
476        }
477        let row = tx.query_one(
478            "admit disclosure",
479            "SELECT public.tre_admit_disclosure($1::text::jsonb)::text",
480            &[&encoded],
481        )?;
482        let snapshot = serde_json::from_str(&row.required_text(0, "snapshot")?)
483            .map_err(|_| invalid("invalid disclosure snapshot"))?;
484        let origins = tx.query_one(
485            "resolve admitted origins",
486            "SELECT public.tre_disclosure_locations($1)::text",
487            &[&intent.admission_id],
488        )?;
489        let origins = serde_json::from_str(&origins.required_text(0, "input origins")?)
490            .map_err(|_| invalid("invalid input origins"))?;
491        let inputs =
492            capture(&snapshot, origins).map_err(|_| invalid("immutable inputs unavailable"))?;
493        Ok((snapshot, inputs))
494    })
495}
496
497pub fn start_disclosure<E: MetadataExecutor>(db: &mut E, id: Uuid) -> Result<Uuid, PgMetaError>
498where
499    E::Row: SchemaStatusRow,
500{
501    let row = db.query_one(
502        "start disclosure",
503        "SELECT public.tre_start_disclosure($1)::text",
504        &[&id],
505    )?;
506    row.required_text(0, "delivery evidence")?
507        .parse()
508        .map_err(|_| invalid("invalid event identity"))
509}
510
511pub fn finish_disclosure<E: MetadataExecutor>(
512    db: &mut E,
513    id: Uuid,
514    outcome: ahri_tre_types::DisclosureOutcome,
515    bytes: Option<u64>,
516    rows: Option<u64>,
517) -> Result<Uuid, PgMetaError>
518where
519    E::Row: SchemaStatusRow,
520{
521    let bytes = bytes
522        .map(i64::try_from)
523        .transpose()
524        .map_err(|_| invalid("invalid byte count"))?;
525    let rows = rows
526        .map(i64::try_from)
527        .transpose()
528        .map_err(|_| invalid("invalid row count"))?;
529    let row = db.query_one(
530        "finish disclosure",
531        "SELECT public.tre_finish_disclosure($1,$2,$3,$4)::text",
532        &[&id, &outcome.as_str(), &bytes, &rows],
533    )?;
534    row.required_text(0, "terminal evidence")?
535        .parse()
536        .map_err(|_| invalid("invalid event identity"))
537}
538
539/// Bounded operator recovery inventory. Pending older disclosure completions
540/// are independent of the evidence required by a new admission.
541pub fn pending_evidence<E: MetadataExecutor>(
542    db: &mut E,
543    limit: u32,
544) -> Result<Vec<Uuid>, PgMetaError>
545where
546    E::Row: SchemaStatusRow,
547{
548    if limit == 0 || limit > 1000 {
549        return Err(invalid("invalid evidence batch"));
550    }
551    db.query_many("pending governance evidence", "SELECT event_id::text FROM tre_governance.events WHERE NOT ledger_acknowledged ORDER BY occurred_at,event_id LIMIT $1", &[&(limit as i64)])?
552        .iter().map(|row| row.required_text(0,"event identity")?.parse().map_err(|_| invalid("invalid event identity"))).collect()
553}
554
555pub struct UnfinishedDisclosure {
556    pub admission_id: Uuid,
557    pub coordinator_id: Uuid,
558    pub generation_id: Uuid,
559}
560
561pub fn unfinished_disclosures<E: MetadataExecutor>(
562    db: &mut E,
563    limit: u32,
564) -> Result<Vec<UnfinishedDisclosure>, PgMetaError>
565where
566    E::Row: SchemaStatusRow,
567{
568    if limit == 0 || limit > 1000 {
569        return Err(invalid("invalid recovery batch"));
570    }
571    db.query_many("unfinished disclosures","SELECT admission_id::text,snapshot->>'coordinator_id',snapshot->>'worker_generation' FROM tre_governance.disclosures WHERE terminal_evidence IS NULL ORDER BY admission_id LIMIT $1",&[&(limit as i64)])?.iter().map(|row| {
572        let id=|index|row.required_text(index,"recovery identity")?.parse().map_err(|_|invalid("invalid recovery identity"));
573        Ok(UnfinishedDisclosure {admission_id:id(0)?,coordinator_id:id(1)?,generation_id:id(2)?})
574    }).collect()
575}
576
577/// The application must prove executor exclusion before calling this operator
578/// function. It terminalizes evidence without recreating user authority.
579pub fn abandon_disclosure<E: MetadataExecutor>(
580    db: &mut E,
581    receipt: &UnfinishedDisclosure,
582) -> Result<(), PgMetaError>
583where
584    E::Row: SchemaStatusRow,
585{
586    db.query_one(
587        "recover unknown disclosure",
588        "SELECT tre_governance.abandon_disclosure($1,$2,$3)",
589        &[
590            &receipt.admission_id,
591            &receipt.coordinator_id,
592            &receipt.generation_id,
593        ],
594    )?;
595    Ok(())
596}
597
598/// Serialize one event's projection across workers and both durable stores.
599/// The callback must append and read back the identical event in the ledger.
600/// A lost acknowledgement leaves an independently retryable outbox event.
601pub fn project_evidence<E: MetadataExecutor>(
602    db: &mut E,
603    event: Uuid,
604    append_and_verify: impl FnOnce(Uuid, &str) -> Result<(), ()>,
605) -> Result<bool, PgMetaError>
606where
607    E::Row: SchemaStatusRow,
608{
609    db.with_transaction("project governance evidence",|tx| {
610        let row=tx.query_one("lock governance event", "SELECT jsonb_build_object('event_id',event_id,'study_id',study_id,'asset_id',asset_id,'event_kind',event_kind,'actor',actor,'occurred_at',to_char(occurred_at AT TIME ZONE 'UTC','YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"'),'evidence',evidence)::text,ledger_acknowledged::text FROM tre_governance.events WHERE event_id=$1 FOR UPDATE", &[&event])?;
611        if row.required_text(1,"acknowledgement")?=="true" { return Ok(false); }
612        append_and_verify(event,&row.required_text(0,"event")?).map_err(|_|invalid("ledger evidence unavailable"))?;
613        tx.execute_command("acknowledge governance event", "UPDATE tre_governance.events SET ledger_acknowledged=TRUE WHERE event_id=$1", &[&event])?;
614        Ok(true)
615    })
616}
617
618/// Current owning-Study custodians and the appointed governance custodian may
619/// discover only acknowledged event identities through the ordinary Session.
620pub fn scoped_evidence_ids<E: MetadataExecutor>(
621    db: &mut E,
622    study: Uuid,
623    after: chrono::DateTime<chrono::Utc>,
624    after_id: Uuid,
625    limit: u32,
626) -> Result<Vec<Uuid>, PgMetaError>
627where
628    E::Row: SchemaStatusRow,
629{
630    if limit == 0 || limit > 100 {
631        return Err(invalid("invalid evidence page"));
632    }
633    let timestamp = after.to_rfc3339();
634    db.query_many("scoped governance history", "SELECT event_id::text FROM public.tre_governance_evidence_ids($1,$2::text::timestamptz,$3,$4::bigint::integer)", &[&study,&timestamp,&after_id,&(limit as i64)])?.iter()
635        .map(|row| row.required_text(0,"event identity")?.parse().map_err(|_|invalid("invalid event identity"))).collect()
636}
637
638/// Only the resolved operator capability may attest an actual Lake table.
639pub fn bind_dataset_location<E: MetadataExecutor>(
640    db: &mut E,
641    location: &ahri_tre_types::GovernedDatasetLocation,
642) -> Result<(), PgMetaError>
643where
644    E::Row: SchemaStatusRow,
645{
646    let encoded = serde_json::to_string(location).map_err(|_| invalid("invalid Dataset origin"))?;
647    db.with_transaction("bind Dataset origin", |tx| {
648        tx.execute_command("retain immutable origin", "INSERT INTO tre_governance.dataset_locations SELECT * FROM jsonb_populate_record(NULL::tre_governance.dataset_locations,$1::text::jsonb) ON CONFLICT(version_id) DO NOTHING", &[&encoded])?;
649        let row = tx.query_one("verify immutable origin", "SELECT ((to_jsonb(l)-'lake_snapshot')=($2::text::jsonb-'lake_snapshot'))::text FROM tre_governance.dataset_locations l WHERE version_id=$1", &[&location.version_id,&encoded])?;
650        if row.required_text(0,"origin equality")? != "true" { return Err(invalid("immutable Dataset origin conflicts")); }
651        Ok(())
652    })
653}
654
655#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
656pub struct LegacyDatasetOrigin {
657    pub version_id: Uuid,
658    pub asset_id: Uuid,
659    pub study_id: Uuid,
660    pub name: ahri_tre_types::NcName,
661    pub major: i32,
662    pub minor: i32,
663    pub patch: i32,
664}
665impl GovernanceMaintenance {
666    pub fn convert_existing(&self, batch: u32) -> Result<u64, ahri_tre_core::CoreError> {
667        self.repository.with_session_connection(
668            "convert legacy classifications",
669            |db| GovernanceUpgrade::classify_existing(db, batch),
670            |db| GovernanceUpgrade::classify_existing(db, batch),
671        )
672    }
673    pub fn upgrade_plan(&self) -> Result<GovernanceUpgradePlan, ahri_tre_core::CoreError> {
674        self.repository.with_session_connection(
675            "inspect governance conversion",
676            GovernanceUpgrade::plan,
677            GovernanceUpgrade::plan,
678        )
679    }
680    pub fn unbound_legacy_datasets(
681        &self,
682        batch: u32,
683    ) -> Result<Vec<LegacyDatasetOrigin>, ahri_tre_core::CoreError> {
684        self.repository.with_session_connection(
685            "legacy Dataset origins",
686            |db| unbound_legacy_datasets(db, batch),
687            |db| unbound_legacy_datasets(db, batch),
688        )
689    }
690}
691fn unbound_legacy_datasets<E: MetadataExecutor>(
692    db: &mut E,
693    batch: u32,
694) -> Result<Vec<LegacyDatasetOrigin>, PgMetaError>
695where
696    E::Row: SchemaStatusRow,
697{
698    if batch == 0 || batch > 1000 {
699        return Err(invalid("invalid origin batch"));
700    }
701    db.query_many("legacy Dataset origin batch", "SELECT jsonb_build_object('version_id',v.version_id,'asset_id',a.asset_id,'study_id',a.study_id,'name',a.name,'major',v.major,'minor',v.minor,'patch',v.patch)::text FROM public.assets a JOIN public.asset_versions v USING(asset_id) JOIN public.datasets d ON d.dataset_id=v.version_id JOIN tre_governance.legacy_assets legacy USING(asset_id,study_id) WHERE NOT EXISTS(SELECT 1 FROM tre_governance.dataset_locations l WHERE l.version_id=v.version_id) ORDER BY v.version_id LIMIT $1",&[&(batch as i64)])?.iter().map(|row|serde_json::from_str(&row.required_text(0,"Dataset origin")?).map_err(|_|invalid("invalid Dataset origin"))).collect()
702}
703
704/// Computation may consume high inputs, but has no content-delivery capability.
705/// The final output admission independently revalidates its forced-high policy.
706pub fn capture_derivation_inputs<E: MetadataExecutor, T>(
707    db: &mut E,
708    inputs: &[Uuid],
709    capture: impl FnOnce(Vec<ahri_tre_types::GovernedDatasetLocation>) -> Result<T, ()>,
710) -> Result<T, PgMetaError>
711where
712    E::Row: SchemaStatusRow,
713{
714    db.with_transaction("authorize and capture derivation", |tx| {
715        let row = tx.query_one(
716            "resolve derivation closure",
717            "SELECT public.tre_derivation_locations($1)::text",
718            &[&inputs],
719        )?;
720        let origins = serde_json::from_str(&row.required_text(0, "origins")?)
721            .map_err(|_| invalid("invalid immutable origins"))?;
722        capture(origins).map_err(|_| invalid("derivation inputs unavailable"))
723    })
724}
725
726impl GovernanceMaintenance {
727    pub fn bind_datafile(
728        &self,
729        binding: &ahri_tre_types::GovernedDatafileBinding,
730    ) -> Result<(), ahri_tre_core::CoreError> {
731        self.repository.with_session_connection(
732            "attest Datafile",
733            |db| bind_datafile(db, binding),
734            |db| bind_datafile(db, binding),
735        )
736    }
737}
738pub fn bind_datafile<E: MetadataExecutor>(
739    db: &mut E,
740    binding: &ahri_tre_types::GovernedDatafileBinding,
741) -> Result<(), PgMetaError>
742where
743    E::Row: SchemaStatusRow,
744{
745    let encoded =
746        serde_json::to_string(binding).map_err(|_| invalid("invalid Datafile attestation"))?;
747    db.with_transaction("retain Datafile origin",|tx| {
748        tx.execute_command("retain Datafile attestation","INSERT INTO tre_governance.datafile_bindings SELECT * FROM jsonb_populate_record(NULL::tre_governance.datafile_bindings,$1::text::jsonb) ON CONFLICT(version_id) DO NOTHING",&[&encoded])?;
749        let row=tx.query_one("verify Datafile attestation","SELECT (to_jsonb(b)=$2::text::jsonb)::text FROM tre_governance.datafile_bindings b WHERE version_id=$1",&[&binding.version_id,&encoded])?;
750        if row.required_text(0,"attestation equality")?!="true" {return Err(invalid("immutable Datafile origin conflicts"));}
751        Ok(())
752    })
753}
754pub fn admit_disclosure_with_datafile<E: MetadataExecutor, T>(
755    db: &mut E,
756    intent: &DisclosureIntent,
757    capture: impl FnOnce(&DisclosureSnapshot, ahri_tre_types::AdmittedDatafileInput) -> Result<T, ()>,
758) -> Result<(DisclosureSnapshot, T), PgMetaError>
759where
760    E::Row: SchemaStatusRow,
761{
762    let encoded =
763        serde_json::to_string(intent).map_err(|_| invalid("invalid disclosure intent"))?;
764    db.with_transaction("admit and capture Datafile", |tx| {
765        let row = tx.query_one(
766            "admit disclosure",
767            "SELECT public.tre_admit_disclosure($1::text::jsonb)::text",
768            &[&encoded],
769        )?;
770        let snapshot = serde_json::from_str(&row.required_text(0, "snapshot")?)
771            .map_err(|_| invalid("invalid disclosure snapshot"))?;
772        let row = tx.query_one(
773            "resolve admitted Datafile",
774            "SELECT public.tre_disclosure_datafile($1)::text",
775            &[&intent.admission_id],
776        )?;
777        let input = serde_json::from_str(&row.required_text(0, "Datafile input")?)
778            .map_err(|_| invalid("Datafile input unavailable"))?;
779        let captured = capture(&snapshot, input)
780            .map_err(|_| invalid("immutable Datafile capture unavailable"))?;
781        Ok((snapshot, captured))
782    })
783}
784
785#[derive(serde::Deserialize)]
786pub struct LegacyDatafileOrigin {
787    pub asset_id: Uuid,
788    pub study_id: Uuid,
789    pub datafile: ahri_tre_types::DataFileRecord,
790}
791impl GovernanceMaintenance {
792    pub fn unbound_legacy_datafiles(
793        &self,
794        batch: u32,
795    ) -> Result<Vec<LegacyDatafileOrigin>, ahri_tre_core::CoreError> {
796        self.repository.with_session_connection(
797            "legacy Datafile origins",
798            |db| unbound_legacy_datafiles(db, batch),
799            |db| unbound_legacy_datafiles(db, batch),
800        )
801    }
802}
803fn unbound_legacy_datafiles<E: MetadataExecutor>(
804    db: &mut E,
805    batch: u32,
806) -> Result<Vec<LegacyDatafileOrigin>, PgMetaError>
807where
808    E::Row: SchemaStatusRow,
809{
810    if batch == 0 || batch > 1000 {
811        return Err(invalid("invalid origin batch"));
812    }
813    db.query_many("legacy Datafile origin batch", "SELECT jsonb_build_object('asset_id',a.asset_id,'study_id',a.study_id,'datafile',tre_governance.datafile_descriptor(d))::text FROM public.assets a JOIN public.asset_versions v USING(asset_id) JOIN public.datafiles d ON d.datafile_id=v.version_id JOIN tre_governance.legacy_assets legacy USING(asset_id,study_id) WHERE NOT EXISTS(SELECT 1 FROM tre_governance.datafile_bindings b WHERE b.version_id=v.version_id) ORDER BY v.version_id LIMIT $1", &[&(batch as i64)])?.iter().map(|row|serde_json::from_str(&row.required_text(0,"Datafile origin")?).map_err(|_|invalid("invalid Datafile origin"))).collect()
814}
815
816impl GovernanceMaintenance {
817    /// Only evidence required by this input closure, never unrelated completions.
818    pub fn required_evidence(
819        &self,
820        inputs: &[Uuid],
821        batch: u32,
822    ) -> Result<Vec<Uuid>, ahri_tre_core::CoreError> {
823        self.repository.with_session_connection(
824            "required disclosure evidence",
825            |db| required_evidence(db, inputs, batch),
826            |db| required_evidence(db, inputs, batch),
827        )
828    }
829}
830pub fn required_evidence<E: MetadataExecutor>(
831    db: &mut E,
832    inputs: &[Uuid],
833    batch: u32,
834) -> Result<Vec<Uuid>, PgMetaError>
835where
836    E::Row: SchemaStatusRow,
837{
838    if inputs.is_empty() || inputs.len() > 128 || batch == 0 || batch > 1000 {
839        return Err(invalid("invalid prerequisite batch"));
840    }
841    db.query_many("required disclosure evidence", "SELECT e.event_id::text FROM tre_governance.events e WHERE NOT e.ledger_acknowledged AND (e.event_id IN (SELECT c.evidence_id FROM tre_governance.classifications c JOIN public.asset_versions v USING(asset_id) WHERE v.version_id=ANY($1)) OR (e.event_kind IN ('study_policy_baseline','study_policy_change') AND e.study_id IN (SELECT a.study_id FROM public.assets a JOIN public.asset_versions v USING(asset_id) WHERE v.version_id=ANY($1))) OR (e.event_kind='version_ingest' AND e.evidence->>'version_id'=ANY(SELECT unnest($1::uuid[])::text))) ORDER BY e.occurred_at,e.event_id LIMIT $2", &[&inputs,&(batch as i64)])?.iter().map(|row| row.required_text(0,"event identity")?.parse().map_err(|_|invalid("invalid event identity"))).collect()
842}