1use 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
68pub 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
85pub 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 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 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 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
364pub 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
392pub 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#[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 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
545pub 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}