Skip to main content

ahri_tre_app/
disclosure.rs

1//! Trusted disclosure workflow and evidence coordination. SQL and ledger
2//! mechanics remain with metadata and Lake adapters, respectively.
3use crate::AppError;
4use ahri_tre_core::{DatasetExecutorIdentity, DatasetExecutorLease};
5
6pub(crate) fn acquire_disclosure_permit() -> Result<tokio::sync::OwnedSemaphorePermit, AppError> {
7    static PERMITS: std::sync::LazyLock<std::sync::Arc<tokio::sync::Semaphore>> =
8        std::sync::LazyLock::new(|| std::sync::Arc::new(tokio::sync::Semaphore::new(4)));
9    PERMITS
10        .clone()
11        .try_acquire_owned()
12        .map_err(|_| AppError::Conflict("Trusted content transfer capacity is occupied".into()))
13}
14
15pub struct DisclosureQueryInput {
16    pub alias: String,
17    pub asset: ahri_tre_protocol::asset::AssetSelector,
18    pub version: Option<String>,
19}
20pub struct DisclosureQueryRequest {
21    pub preview_limit: Option<u64>,
22    pub representation: ahri_tre_types::ContentRepresentation,
23    pub inputs: Vec<DisclosureQueryInput>,
24    pub views: std::collections::BTreeMap<String, String>,
25    pub sql: String,
26    pub budgets: ahri_tre_types::DisclosureBudgetOverrides,
27}
28
29pub enum DisclosureRequest {
30    Query(DisclosureQueryRequest),
31    Datafile {
32        compress: bool,
33        representation: (bool, bool),
34        asset: ahri_tre_protocol::asset::AssetSelector,
35        version: Option<String>,
36        budgets: ahri_tre_types::DisclosureBudgetOverrides,
37    },
38}
39
40pub(crate) enum PreparedDisclosure {
41    Semantic {
42        bytes: Vec<u8>,
43        rows: u64,
44    },
45    Query(
46        ahri_tre_lake::PreparedDisclosureQuery,
47        ahri_tre_types::ContentRepresentation,
48        Option<u64>,
49    ),
50    Datafile(ahri_tre_lake::PreparedDisclosureDatafile, bool),
51}
52
53enum DisclosureStream {
54    Semantic(std::io::Cursor<Vec<u8>>, u64),
55    Query(ahri_tre_lake::RestrictedQueryStream),
56    Datafile(Box<dyn std::io::Read + Send>),
57}
58impl std::io::Read for DisclosureStream {
59    fn read(&mut self, bytes: &mut [u8]) -> std::io::Result<usize> {
60        match self {
61            Self::Semantic(stream, _) => stream.read(bytes),
62            Self::Query(stream) => stream.read(bytes),
63            Self::Datafile(stream) => stream.read(bytes),
64        }
65    }
66}
67
68/// Owns both immutable inputs and executor exclusion for its entire lifetime.
69/// The mutable Session borrow prevents content outliving its admitted Session.
70pub struct SessionDisclosure<'session> {
71    session: &'session mut crate::DataStoreSession,
72    operator: governance::GovernanceMaintenance,
73    snapshot: ahri_tre_types::DisclosureSnapshot,
74    prepared: Option<PreparedDisclosure>,
75    stream: Option<DisclosureStream>,
76    last_progress: std::time::Instant,
77    bytes: u64,
78    eof: bool,
79    finished: bool,
80    failure: Option<ahri_tre_types::DisclosureOutcome>,
81    _guard: Box<dyn ahri_tre_core::DatasetAttemptExecutor>,
82    _permit: tokio::sync::OwnedSemaphorePermit,
83}
84
85/// Schema metadata remains independently readable. Content is carried only by
86/// the same retained Session capability used for Dataset/Datafile disclosure.
87pub enum PreparedSemanticResponse<'session> {
88    Schema(serde_json::Value),
89    Content(Box<SessionDisclosure<'session>>),
90}
91impl<'session> SessionDisclosure<'session> {
92    pub(crate) fn new(
93        session: &'session mut crate::DataStoreSession,
94        operator: governance::GovernanceMaintenance,
95        snapshot: ahri_tre_types::DisclosureSnapshot,
96        prepared: PreparedDisclosure,
97        guard: Box<dyn ahri_tre_core::DatasetAttemptExecutor>,
98        permit: tokio::sync::OwnedSemaphorePermit,
99    ) -> Result<Self, AppError> {
100        if chrono::Utc::now() >= snapshot.deadline {
101            return Err(evidence_unavailable(()));
102        }
103        Ok(Self {
104            session,
105            operator,
106            snapshot,
107            prepared: Some(prepared),
108            stream: None,
109            last_progress: std::time::Instant::now(),
110            bytes: 0,
111            eof: false,
112            finished: false,
113            failure: None,
114            _guard: guard,
115            _permit: permit,
116        })
117    }
118    pub fn snapshot(&self) -> &ahri_tre_types::DisclosureSnapshot {
119        &self.snapshot
120    }
121    pub fn failure_outcome(&self) -> Option<ahri_tre_types::DisclosureOutcome> {
122        self.failure
123    }
124    pub fn rows(&self) -> Option<u64> {
125        match self.stream.as_ref() {
126            Some(DisclosureStream::Semantic(_, rows)) if self.eof => Some(*rows),
127            Some(DisclosureStream::Query(stream)) if self.eof => stream.rows(),
128            _ => None,
129        }
130    }
131
132    /// Internal batch preparation. Reading the captured source is computation,
133    /// not delivery; only the resulting semantic JSON can cross this capability.
134    pub(crate) fn prepare_semantic_batch(
135        &mut self,
136        state: &ahri_tre_pgmeta::semantic::SemanticState,
137        apply: impl FnOnce(
138            ahri_tre_pgmeta::PgMetadataRepository<'_>,
139            Vec<std::collections::BTreeMap<String, Option<String>>>,
140        ) -> Result<(Vec<u8>, u64), ahri_tre_core::CoreError>,
141    ) -> Result<(), AppError> {
142        use std::io::Read;
143        let Some(PreparedDisclosure::Query(query, _, None)) = self.prepared.take() else {
144            return Err(evidence_unavailable(()));
145        };
146        let remaining = (self.snapshot.deadline - chrono::Utc::now()).num_seconds();
147        if remaining < 1 || self.snapshot.max_bytes > 1024 * 1024 {
148            return Err(evidence_unavailable(()));
149        }
150        let worker = std::env::current_exe()
151            .map_err(evidence_unavailable)?
152            .with_file_name("ahri-tre-query-worker");
153        let mut stream = query
154            .execute_representation(
155                &worker,
156                ahri_tre_types::ContentRepresentation::Json,
157                None,
158                remaining as u32,
159                self.snapshot.max_bytes,
160                self.snapshot.max_rows,
161            )
162            .map_err(evidence_unavailable)?;
163        let mut bytes = Vec::new();
164        (&mut stream)
165            .take(self.snapshot.max_bytes + 1)
166            .read_to_end(&mut bytes)
167            .map_err(evidence_unavailable)?;
168        if bytes.len() as u64 > self.snapshot.max_bytes {
169            return Err(evidence_unavailable(()));
170        }
171        let source_rows: Vec<std::collections::BTreeMap<String, Option<String>>> =
172            serde_json::from_slice(&bytes).map_err(evidence_unavailable)?;
173        if source_rows.len() as u64 > self.snapshot.max_rows {
174            return Err(evidence_unavailable(()));
175        }
176        let repository = self
177            .session
178            .dataset_admission_repository()
179            .ok_or_else(|| evidence_unavailable(()))?;
180        let (bytes, rows) = repository
181            .with_admitted_semantic_mutation(state, self.snapshot.admission_id, |repository| {
182                let result = apply(repository, source_rows)?;
183                if result.0.len() as u64 > self.snapshot.max_bytes
184                    || result.1 > self.snapshot.max_rows
185                    || chrono::Utc::now() >= self.snapshot.deadline
186                {
187                    return Err(ahri_tre_core::CoreError::Infrastructure(
188                        "Semantic result exceeds its limits".into(),
189                    ));
190                }
191                Ok(result)
192            })
193            .map_err(evidence_unavailable)?;
194        self.prepared = Some(PreparedDisclosure::Semantic { bytes, rows });
195        Ok(())
196    }
197
198    fn start(&mut self) -> Result<(), AppError> {
199        let event = match &self.session.store {
200            crate::StoreSessionConnection::Direct(db) => governance::start_disclosure(
201                &mut *db.lock().map_err(evidence_unavailable)?,
202                self.snapshot.admission_id,
203            ),
204            crate::StoreSessionConnection::OAuth(db) => governance::start_disclosure(
205                &mut *db.lock().map_err(evidence_unavailable)?,
206                self.snapshot.admission_id,
207            ),
208            #[cfg(test)]
209            crate::StoreSessionConnection::TestUnavailable => return Err(evidence_unavailable(())),
210        }
211        .map_err(evidence_unavailable)?;
212        project_session_evidence(&self.operator, &mut self.session.lake, event)?;
213        let remaining = (self.snapshot.deadline - chrono::Utc::now()).num_seconds();
214        if remaining < 1 {
215            return Err(evidence_unavailable(()));
216        }
217        let executable = std::env::current_exe().map_err(evidence_unavailable)?;
218        let worker = executable
219            .parent()
220            .ok_or_else(|| evidence_unavailable(()))?
221            .join("ahri-tre-query-worker");
222        let prepared = self
223            .prepared
224            .take()
225            .ok_or_else(|| evidence_unavailable(()))?;
226        self.stream = Some(match prepared {
227            PreparedDisclosure::Semantic { bytes, rows } => {
228                DisclosureStream::Semantic(std::io::Cursor::new(bytes), rows)
229            }
230            PreparedDisclosure::Query(query, representation, preview_limit) => {
231                DisclosureStream::Query(
232                    query
233                        .execute_representation(
234                            &worker,
235                            representation,
236                            preview_limit,
237                            remaining as u32,
238                            self.snapshot.max_bytes,
239                            self.snapshot.max_rows,
240                        )
241                        .map_err(evidence_unavailable)?,
242                )
243            }
244            PreparedDisclosure::Datafile(file, compress) => DisclosureStream::Datafile(
245                file.open_representation(compress)
246                    .map_err(evidence_unavailable)?,
247            ),
248        });
249        Ok(())
250    }
251    /// Completion counts describe server transport, never confirmed client receipt.
252    pub fn finish(
253        mut self,
254        outcome: ahri_tre_types::DisclosureOutcome,
255        rows: Option<u64>,
256    ) -> Result<(), AppError> {
257        if outcome == ahri_tre_types::DisclosureOutcome::Complete && !self.eof {
258            return Err(AppError::Conflict("Disclosure did not complete".into()));
259        }
260        let rows = rows.or(self.rows());
261        self.stream.take();
262        let result = self.record_terminal(outcome, rows);
263        // A failed acknowledgement is recovered from the outbox/receipt; Drop
264        // must not overwrite the truthful outcome with a conflicting event.
265        self.finished = true;
266        result.map_err(|_| {
267            AppError::Infrastructure(
268                "Delivery may be partial or complete; completion evidence is pending".into(),
269            )
270        })
271    }
272    fn record_terminal(
273        &mut self,
274        outcome: ahri_tre_types::DisclosureOutcome,
275        rows: Option<u64>,
276    ) -> Result<(), AppError> {
277        let event = match &self.session.store {
278            crate::StoreSessionConnection::Direct(db) => governance::finish_disclosure(
279                &mut *db.lock().map_err(evidence_unavailable)?,
280                self.snapshot.admission_id,
281                outcome,
282                Some(self.bytes),
283                rows,
284            ),
285            crate::StoreSessionConnection::OAuth(db) => governance::finish_disclosure(
286                &mut *db.lock().map_err(evidence_unavailable)?,
287                self.snapshot.admission_id,
288                outcome,
289                Some(self.bytes),
290                rows,
291            ),
292            #[cfg(test)]
293            crate::StoreSessionConnection::TestUnavailable => return Err(evidence_unavailable(())),
294        }
295        .map_err(evidence_unavailable)?;
296        project_session_evidence(&self.operator, &mut self.session.lake, event)
297    }
298}
299impl std::io::Read for SessionDisclosure<'_> {
300    fn read(&mut self, buffer: &mut [u8]) -> std::io::Result<usize> {
301        if buffer.is_empty() {
302            return Ok(0);
303        }
304        if chrono::Utc::now() >= self.snapshot.deadline
305            || self.last_progress.elapsed() >= std::time::Duration::from_secs(30)
306        {
307            self.failure = Some(ahri_tre_types::DisclosureOutcome::TimedOut);
308            return Err(std::io::Error::other("Disclosure expired"));
309        }
310        if self.stream.is_none() {
311            self.start()
312                .map_err(|_| std::io::Error::other("Disclosure could not start"))?;
313        }
314        let read = self.stream.as_mut().unwrap().read(buffer)?;
315        if chrono::Utc::now() >= self.snapshot.deadline {
316            self.failure = Some(ahri_tre_types::DisclosureOutcome::TimedOut);
317            return Err(std::io::Error::other("Disclosure expired"));
318        }
319        let total = self
320            .bytes
321            .checked_add(read as u64)
322            .ok_or_else(|| std::io::Error::other("Disclosure byte limit"))?;
323        if total > self.snapshot.max_bytes {
324            self.failure = Some(ahri_tre_types::DisclosureOutcome::LimitExceeded);
325            return Err(std::io::Error::other("Disclosure byte limit"));
326        }
327        self.bytes = total;
328        if read > 0 {
329            self.last_progress = std::time::Instant::now();
330        }
331        self.eof = read == 0;
332        Ok(read)
333    }
334}
335impl Drop for SessionDisclosure<'_> {
336    fn drop(&mut self) {
337        self.stream.take();
338        if !self.finished {
339            let _ = self.record_terminal(
340                self.failure
341                    .unwrap_or(ahri_tre_types::DisclosureOutcome::Interrupted),
342                None,
343            );
344        }
345    }
346}
347
348pub(crate) fn project_session_evidence(
349    operator: &governance::GovernanceMaintenance,
350    lake: &mut duckdb::Connection,
351    event: uuid::Uuid,
352) -> Result<(), AppError> {
353    let binding = operator.binding().map_err(evidence_unavailable)?;
354    let mut ledger =
355        ahri_tre_lake::GovernanceLedger::open(lake, ahri_tre_types::StudyId(binding.study_id))
356            .map_err(evidence_unavailable)?;
357    operator
358        .project(event, |id, text| ledger.append(id, text).map_err(|_| ()))
359        .map_err(evidence_unavailable)?;
360    Ok(())
361}
362use ahri_tre_pgmeta::{MetadataExecutor, governance, schema::SchemaStatusRow};
363
364/// Reconcile only receipts whose executor has stopped. The persistent
365/// coordinator lock prevents a replacement process taking over a live worker.
366pub fn reconcile_abandoned_disclosures<E: MetadataExecutor>(
367    metadata: &mut E,
368    executor: &ahri_tre_lake::DatasetExecutor,
369    limit: u32,
370) -> Result<u64, AppError>
371where
372    E::Row: SchemaStatusRow,
373{
374    let mut recovered = 0;
375    for receipt in
376        governance::unfinished_disclosures(metadata, limit).map_err(evidence_unavailable)?
377    {
378        if executor.can_recover(
379            DatasetExecutorIdentity {
380                coordinator_id: receipt.coordinator_id,
381                generation_id: receipt.generation_id,
382            },
383            receipt.admission_id,
384        ) {
385            governance::abandon_disclosure(metadata, &receipt).map_err(evidence_unavailable)?;
386            recovered += 1;
387        }
388    }
389    Ok(recovered)
390}
391
392/// Operator recovery over already resolved metadata and Lake capabilities.
393/// This does not recreate Sessions, replay content, or pause other admissions.
394pub fn reconcile_governance_evidence<E: MetadataExecutor>(
395    metadata: &mut E,
396    lake: &mut duckdb::Connection,
397    limit: u32,
398) -> Result<u64, AppError>
399where
400    E::Row: SchemaStatusRow,
401{
402    let binding = governance::ledger_binding(metadata).map_err(evidence_unavailable)?;
403    let mut ledger =
404        ahri_tre_lake::GovernanceLedger::open(lake, ahri_tre_types::StudyId(binding.study_id))
405            .map_err(evidence_unavailable)?;
406    let pending = governance::pending_evidence(metadata, limit).map_err(evidence_unavailable)?;
407    let mut acknowledged = 0;
408    for event in pending {
409        if governance::project_evidence(metadata, event, |id, evidence| {
410            ledger.append(id, evidence).map_err(|_| ())
411        })
412        .map_err(evidence_unavailable)?
413        {
414            acknowledged += 1;
415        }
416    }
417    Ok(acknowledged)
418}
419
420fn evidence_unavailable<E>(_: E) -> AppError {
421    AppError::Infrastructure("Governance evidence is unavailable".into())
422}
423
424/// Only this application seam constructs a delivery capability. A metadata
425/// snapshot alone never becomes an `AdmittedDisclosure`.
426#[cfg(test)]
427pub struct AdmittedDisclosure {
428    snapshot: ahri_tre_types::DisclosureSnapshot,
429    started: bool,
430}
431
432#[cfg(test)]
433pub fn admit_disclosure<E: MetadataExecutor, O: MetadataExecutor>(
434    session: &mut E,
435    operator: &mut O,
436    lake: &mut duckdb::Connection,
437    intent: &ahri_tre_types::DisclosureIntent,
438) -> Result<AdmittedDisclosure, AppError>
439where
440    E::Row: SchemaStatusRow,
441    O::Row: SchemaStatusRow,
442{
443    let snapshot = governance::admit_disclosure(session, intent).map_err(evidence_unavailable)?;
444    project_required_evidence(operator, lake, snapshot.admission_evidence)?;
445    let admitted = AdmittedDisclosure {
446        snapshot,
447        started: false,
448    };
449    admitted.ensure_live()?;
450    Ok(admitted)
451}
452
453#[cfg(test)]
454impl AdmittedDisclosure {
455    pub fn start_delivery<E: MetadataExecutor, O: MetadataExecutor>(
456        &mut self,
457        session: &mut E,
458        operator: &mut O,
459        lake: &mut duckdb::Connection,
460    ) -> Result<(), AppError>
461    where
462        E::Row: SchemaStatusRow,
463        O::Row: SchemaStatusRow,
464    {
465        self.ensure_live()?;
466        if self.started {
467            return Err(AppError::Conflict(
468                "Disclosure delivery already started".into(),
469            ));
470        }
471        let event = governance::start_disclosure(session, self.snapshot.admission_id)
472            .map_err(evidence_unavailable)?;
473        project_required_evidence(operator, lake, event)?;
474        self.ensure_live()?;
475        self.started = true;
476        Ok(())
477    }
478
479    /// Counts describe the stated server transport boundary, never client receipt.
480    /// Failure after delivery preserves pending/unknown evidence for recovery.
481    pub fn finish<E: MetadataExecutor, O: MetadataExecutor>(
482        self,
483        session: &mut E,
484        operator: &mut O,
485        lake: &mut duckdb::Connection,
486        outcome: ahri_tre_types::DisclosureOutcome,
487        bytes: Option<u64>,
488        rows: Option<u64>,
489    ) -> Result<(), AppError>
490    where
491        E::Row: SchemaStatusRow,
492        O::Row: SchemaStatusRow,
493    {
494        let pending = |_| {
495            AppError::Infrastructure(
496                "Delivery may be partial or complete; completion evidence is pending".into(),
497            )
498        };
499        let event = governance::finish_disclosure(
500            session,
501            self.snapshot.admission_id,
502            outcome,
503            bytes,
504            rows,
505        )
506        .map_err(pending)?;
507        project_required_evidence(operator, lake, event).map_err(|_| {
508            AppError::Infrastructure(
509                "Delivery may be partial or complete; completion evidence is pending".into(),
510            )
511        })
512    }
513
514    fn ensure_live(&self) -> Result<(), AppError> {
515        if chrono::Utc::now() >= self.snapshot.deadline {
516            return Err(AppError::Conflict("Disclosure admission expired".into()));
517        }
518        Ok(())
519    }
520}
521
522#[cfg(test)]
523fn project_required_evidence<O: MetadataExecutor>(
524    operator: &mut O,
525    lake: &mut duckdb::Connection,
526    event: uuid::Uuid,
527) -> Result<(), AppError>
528where
529    O::Row: SchemaStatusRow,
530{
531    let binding = governance::ledger_binding(operator).map_err(evidence_unavailable)?;
532    let mut ledger =
533        ahri_tre_lake::GovernanceLedger::open(lake, ahri_tre_types::StudyId(binding.study_id))
534            .map_err(evidence_unavailable)?;
535    governance::project_evidence(operator, event, |id, evidence| {
536        ledger.append(id, evidence).map_err(|_| ())
537    })
538    .map_err(evidence_unavailable)?;
539    Ok(())
540}
541
542#[cfg(test)]
543mod tests;
544
545/// A readback target resolves before admission. The resulting SessionDisclosure
546/// retains the authorized records and is the only content-returning interface.
547pub struct SemanticReadbackRequest {
548    pub study: ahri_tre_protocol::study::StudySelector,
549    pub dataset: ahri_tre_protocol::dataset::DatasetSelector,
550    pub target: SemanticReadbackTarget,
551}
552pub enum SemanticReadbackTarget {
553    Entity(ahri_tre_protocol::model::EntitySelector),
554    Relation(ahri_tre_protocol::model::RelationSelector),
555}
556impl SessionDisclosure<'_> {
557    pub(crate) fn prepare_semantic_readback(
558        &mut self,
559        state: &ahri_tre_pgmeta::semantic::SemanticState,
560        read: impl FnOnce(
561            ahri_tre_pgmeta::PgMetadataRepository<'_>,
562        ) -> Result<(Vec<u8>, u64), ahri_tre_core::CoreError>,
563    ) -> Result<(), AppError> {
564        let repository = self
565            .session
566            .dataset_admission_repository()
567            .ok_or_else(|| evidence_unavailable(()))?;
568        let (bytes, rows) = repository
569            .with_admitted_semantic_mutation(state, self.snapshot.admission_id, |repository| {
570                let (bytes, rows) = read(repository)?;
571                if bytes.len() as u64 > self.snapshot.max_bytes
572                    || rows > self.snapshot.max_rows
573                    || chrono::Utc::now() >= self.snapshot.deadline
574                {
575                    return Err(ahri_tre_core::CoreError::Infrastructure(
576                        "Semantic readback exceeds its limits".into(),
577                    ));
578                }
579                Ok((bytes, rows))
580            })
581            .map_err(evidence_unavailable)?;
582        self.prepared = Some(PreparedDisclosure::Semantic { bytes, rows });
583        Ok(())
584    }
585}