Skip to main content

ahri_tre_app/
session.rs

1use ahri_tre_libpq_oauth::LibpqOAuthConnection;
2use ahri_tre_pgmeta::{PgMetadataConnection, PgMetadataRepository};
3use ahri_tre_runtime::DataStoreRuntime;
4use std::sync::{Arc, Mutex, MutexGuard};
5
6use crate::{AppError, LakeQueryAuthorizer, LakeQueryRequest, LakeQueryResult, execute_lake_query};
7use ahri_tre_core::{
8    AssetRepository, CatalogueRegistrationRepository, DomainRepository, EntityRepository,
9    StudyDomainRepository, StudyGovernanceRepository, StudyRepository, TagRepository,
10    TransformationRepository, VariableRepository, VocabularyRepository,
11};
12
13#[derive(Clone)]
14pub enum StoreSessionConnection {
15    Direct(Arc<Mutex<PgMetadataConnection>>),
16    OAuth(Arc<Mutex<LibpqOAuthConnection>>),
17    #[cfg(test)]
18    TestUnavailable,
19}
20
21pub type SessionMetadataRepositories = ScopedSessionMetadataRepositories<'static>;
22
23#[derive(Clone)]
24pub struct ScopedSessionMetadataRepositories<'repository> {
25    pub registrations: Arc<dyn CatalogueRegistrationRepository + 'repository>,
26    pub domains: Arc<dyn DomainRepository + 'repository>,
27    pub studies: Arc<dyn StudyRepository + 'repository>,
28    pub study_domains: Arc<dyn StudyDomainRepository + 'repository>,
29    pub study_governance: Arc<dyn StudyGovernanceRepository + 'repository>,
30    pub assets: Arc<dyn AssetRepository + 'repository>,
31    pub variables: Arc<dyn VariableRepository + 'repository>,
32    pub vocabularies: Arc<dyn VocabularyRepository + 'repository>,
33    pub entities: Arc<dyn EntityRepository + 'repository>,
34    pub transformations: Arc<dyn TransformationRepository + 'repository>,
35    pub tags: Arc<dyn TagRepository + 'repository>,
36}
37
38impl<'repository> ScopedSessionMetadataRepositories<'repository> {
39    fn from_direct_connection(connection: Arc<Mutex<PgMetadataConnection>>) -> Self {
40        let repository = Arc::new(PgMetadataRepository::from_shared(connection));
41        Self::from_repository(repository)
42    }
43
44    fn from_oauth_connection(connection: Arc<Mutex<LibpqOAuthConnection>>) -> Self {
45        let repository = Arc::new(PgMetadataRepository::from_oauth_shared(connection));
46        Self::from_repository(repository)
47    }
48
49    pub(crate) fn from_repository(repository: Arc<PgMetadataRepository<'repository>>) -> Self {
50        let registrations: Arc<dyn CatalogueRegistrationRepository + 'repository> =
51            repository.clone();
52        let domains: Arc<dyn DomainRepository + 'repository> = repository.clone();
53        let studies: Arc<dyn StudyRepository + 'repository> = repository.clone();
54        let study_domains: Arc<dyn StudyDomainRepository + 'repository> = repository.clone();
55        let study_governance: Arc<dyn StudyGovernanceRepository + 'repository> = repository.clone();
56        let assets: Arc<dyn AssetRepository + 'repository> = repository.clone();
57        let variables: Arc<dyn VariableRepository + 'repository> = repository.clone();
58        let vocabularies: Arc<dyn VocabularyRepository + 'repository> = repository.clone();
59        let entities: Arc<dyn EntityRepository + 'repository> = repository.clone();
60        let transformations: Arc<dyn TransformationRepository + 'repository> = repository.clone();
61        let tags: Arc<dyn TagRepository + 'repository> = repository;
62        Self {
63            registrations,
64            domains,
65            studies,
66            study_domains,
67            study_governance,
68            assets,
69            variables,
70            vocabularies,
71            entities,
72            transformations,
73            tags,
74        }
75    }
76}
77
78pub struct DataStoreSession {
79    pub(crate) disclosure_capture: Option<ahri_tre_lake::DisclosureCaptureAuthority>,
80    pub(crate) disclosure_context:
81        Option<(uuid::Uuid, ahri_tre_runtime::SessionAuthenticationRecord)>,
82    authenticated_actor: Option<String>,
83    current_study: Option<(ahri_tre_protocol::refs::ObjectRef, String)>,
84    observations: crate::SessionObservations,
85    runtime: DataStoreRuntime,
86    pub(crate) store: StoreSessionConnection,
87    pub(crate) lake: duckdb::Connection,
88    _lake_tls_guard: Option<ahri_tre_lake::DuckLakeTlsGuard>,
89    scratch_attempt: Option<Arc<ahri_tre_lake::ScratchAttempt>>,
90    dataset_executor: Option<Arc<ahri_tre_lake::DatasetExecutor>>,
91    dataset_maintenance: Option<ahri_tre_pgmeta::PgDatasetMaintenance>,
92}
93
94impl DataStoreSession {
95    pub fn new(
96        runtime: DataStoreRuntime,
97        store: StoreSessionConnection,
98        lake: duckdb::Connection,
99    ) -> Self {
100        Self {
101            disclosure_capture: None,
102            disclosure_context: None,
103            authenticated_actor: None,
104            current_study: None,
105            observations: crate::SessionObservations::default(),
106            runtime,
107            store,
108            lake,
109            _lake_tls_guard: None,
110            scratch_attempt: None,
111            dataset_executor: None,
112            dataset_maintenance: None,
113        }
114    }
115
116    pub fn observations(&self) -> &crate::SessionObservations {
117        &self.observations
118    }
119
120    /// The database principal admitted by the retained OAuth capability.
121    /// Runtime-login subjects and caller-supplied actors are not substitutes.
122    pub fn authenticated_actor(&self) -> Option<&str> {
123        self.authenticated_actor.as_deref()
124    }
125
126    pub(crate) fn with_disclosure_context(
127        mut self,
128        id: uuid::Uuid,
129        authentication: ahri_tre_runtime::SessionAuthenticationRecord,
130    ) -> Self {
131        self.disclosure_context = Some((id, authentication));
132        self
133    }
134
135    pub(crate) fn governance_maintenance(
136        &self,
137    ) -> Option<ahri_tre_pgmeta::governance::GovernanceMaintenance> {
138        self.dataset_maintenance
139            .as_ref()
140            .map(|authority| authority.governance())
141    }
142
143    pub(crate) fn with_authenticated_actor(mut self, actor: String) -> Self {
144        self.authenticated_actor = Some(actor);
145        self
146    }
147
148    pub fn current_study(&self) -> Option<&(ahri_tre_protocol::refs::ObjectRef, String)> {
149        self.current_study.as_ref()
150    }
151
152    pub fn set_current_study(
153        &mut self,
154        study: Option<(ahri_tre_protocol::refs::ObjectRef, String)>,
155    ) {
156        self.current_study = study;
157    }
158
159    pub fn with_observations(mut self, observations: crate::SessionObservations) -> Self {
160        self.observations = observations;
161        self
162    }
163
164    pub fn with_scratch_attempt(mut self, attempt: Option<ahri_tre_lake::ScratchAttempt>) -> Self {
165        self.scratch_attempt = attempt.map(Arc::new);
166        self
167    }
168
169    /// Supplies executor exclusion obtained from the configured Lake scratch
170    /// capability. Unconfigured direct callers cannot materialize Datasets.
171    pub fn with_dataset_executor(mut self, executor: Arc<ahri_tre_lake::DatasetExecutor>) -> Self {
172        self.dataset_executor = Some(executor);
173        self
174    }
175
176    pub(crate) fn with_dataset_maintenance(
177        mut self,
178        maintenance: Option<ahri_tre_pgmeta::PgDatasetMaintenance>,
179    ) -> Self {
180        self.dataset_maintenance = maintenance;
181        self
182    }
183
184    /// Whether configured opening established the operation storage/recovery capability.
185    pub fn has_operation_recovery(&self) -> bool {
186        self.dataset_maintenance.is_some()
187    }
188
189    pub(crate) fn require_operation_recovery(&self) -> Result<(), AppError> {
190        self.dataset_maintenance
191            .as_ref()
192            .ok_or_else(|| AppError::Infrastructure("Operation recovery is unavailable".into()))?
193            .verify_operation_recovery()
194            .map_err(|_| AppError::Infrastructure("Operation recovery is unavailable".into()))
195    }
196
197    /// Retry known stopped attempts with the capability resolved for this Session.
198    /// It cannot reopen user authority or admit output.
199    pub fn reconcile_dataset_attempts(&mut self) -> Result<(), AppError> {
200        let (Some(maintenance), Some(executor), Some(scratch)) = (
201            &self.dataset_maintenance,
202            &self.dataset_executor,
203            &self.scratch_attempt,
204        ) else {
205            return Ok(());
206        };
207        crate::configured_datastore::reconcile_dataset_attempts(
208            maintenance,
209            executor,
210            scratch,
211            &mut self.lake,
212            &ahri_tre_lake::DuckLakeAdapter::new(&self.runtime.lake.data_path),
213        )
214    }
215
216    /// Recovery owns no login credential and never resumes content delivery.
217    /// One opportunistic pass is bounded so older evidence cannot consume a
218    /// newly authorized Session's admission deadline. Unprocessed receipts stay
219    /// durable for a later pass; new disclosures project their own evidence.
220    pub fn reconcile_disclosures(&mut self) -> Result<(), AppError> {
221        use ahri_tre_core::{DatasetExecutorIdentity, DatasetExecutorLease};
222        const BATCH_SIZE: u32 = 16;
223        const BUDGET: std::time::Duration = std::time::Duration::from_millis(250);
224        let (Some(operator), Some(executor)) =
225            (self.governance_maintenance(), self.dataset_executor())
226        else {
227            return Ok(());
228        };
229        if let (Some((_, authentication)), Some(principal)) =
230            (&self.disclosure_context, self.authenticated_actor())
231        {
232            operator
233                .activate_identity(
234                    authentication.identity().issuer(),
235                    authentication.identity().subject(),
236                    principal,
237                )
238                .map_err(|_| {
239                    AppError::Infrastructure("Governance identity activation is unavailable".into())
240                })?;
241        }
242        for receipt in operator
243            .unfinished(1000)
244            .map_err(|_| AppError::Infrastructure("Disclosure recovery is unavailable".into()))?
245        {
246            if executor.can_recover(
247                DatasetExecutorIdentity {
248                    coordinator_id: receipt.coordinator_id,
249                    generation_id: receipt.generation_id,
250                },
251                receipt.admission_id,
252            ) {
253                operator.abandon(&receipt).map_err(|_| {
254                    AppError::Infrastructure("Disclosure recovery is unavailable".into())
255                })?;
256            }
257        }
258        let started = std::time::Instant::now();
259        for event in operator
260            .pending(BATCH_SIZE)
261            .map_err(|_| AppError::Infrastructure("Governance evidence is unavailable".into()))?
262        {
263            if started.elapsed() >= BUDGET {
264                break;
265            }
266            crate::disclosure::project_session_evidence(&operator, &mut self.lake, event)?;
267        }
268        Ok(())
269    }
270
271    pub(crate) fn dataset_executor(&self) -> Option<Arc<ahri_tre_lake::DatasetExecutor>> {
272        self.dataset_executor.clone()
273    }
274
275    /// Returns the trusted scratch attempt retained by this live Session.
276    ///
277    /// The returned capability is runtime-internal and its restricted path
278    /// must never be serialized or projected to a client.
279    pub fn scratch_attempt(&self) -> Option<&ahri_tre_lake::ScratchAttempt> {
280        self.scratch_attempt.as_deref()
281    }
282
283    pub(crate) fn shared_scratch_attempt(&self) -> Option<Arc<ahri_tre_lake::ScratchAttempt>> {
284        self.scratch_attempt.clone()
285    }
286
287    pub(crate) fn with_lake_tls_guard(
288        mut self,
289        guard: Option<ahri_tre_lake::DuckLakeTlsGuard>,
290    ) -> Self {
291        self._lake_tls_guard = guard;
292        self
293    }
294    #[cfg(test)]
295    pub(crate) fn new_lake_only_for_test(
296        runtime: DataStoreRuntime,
297        lake: duckdb::Connection,
298    ) -> Self {
299        Self::new(runtime, StoreSessionConnection::TestUnavailable, lake)
300    }
301
302    pub fn runtime(&self) -> &DataStoreRuntime {
303        &self.runtime
304    }
305
306    pub fn into_runtime(self) -> DataStoreRuntime {
307        self.runtime.clone()
308    }
309
310    pub fn store_connection(&self) -> MutexGuard<'_, PgMetadataConnection> {
311        match &self.store {
312            StoreSessionConnection::Direct(connection) => connection
313                .lock()
314                .expect("metadata connection lock should not be poisoned"),
315            StoreSessionConnection::OAuth(_) => {
316                panic!("store_connection is only available for direct PostgreSQL sessions")
317            }
318            #[cfg(test)]
319            StoreSessionConnection::TestUnavailable => {
320                panic!("store_connection is not available for lake-only test sessions")
321            }
322        }
323    }
324
325    pub fn direct_store_connection(&self) -> Option<MutexGuard<'_, PgMetadataConnection>> {
326        match &self.store {
327            StoreSessionConnection::Direct(connection) => Some(
328                connection
329                    .lock()
330                    .expect("metadata connection lock should not be poisoned"),
331            ),
332            StoreSessionConnection::OAuth(_) => None,
333            #[cfg(test)]
334            StoreSessionConnection::TestUnavailable => None,
335        }
336    }
337
338    pub fn oauth_store_connection(&self) -> Option<MutexGuard<'_, LibpqOAuthConnection>> {
339        match &self.store {
340            StoreSessionConnection::Direct(_) => None,
341            StoreSessionConnection::OAuth(connection) => Some(
342                connection
343                    .lock()
344                    .expect("OAuth metadata connection lock should not be poisoned"),
345            ),
346            #[cfg(test)]
347            StoreSessionConnection::TestUnavailable => None,
348        }
349    }
350
351    pub(crate) fn dataset_admission_repository(&self) -> Option<PgMetadataRepository<'static>> {
352        match &self.store {
353            StoreSessionConnection::Direct(connection) => {
354                Some(PgMetadataRepository::from_shared(connection.clone()))
355            }
356            StoreSessionConnection::OAuth(connection) => {
357                Some(PgMetadataRepository::from_oauth_shared(connection.clone()))
358            }
359            #[cfg(test)]
360            StoreSessionConnection::TestUnavailable => None,
361        }
362    }
363
364    pub fn metadata_repositories(&self) -> Result<SessionMetadataRepositories, AppError> {
365        match &self.store {
366            StoreSessionConnection::Direct(connection) => Ok(
367                SessionMetadataRepositories::from_direct_connection(Arc::clone(connection)),
368            ),
369            StoreSessionConnection::OAuth(connection) => Ok(
370                SessionMetadataRepositories::from_oauth_connection(Arc::clone(connection)),
371            ),
372            #[cfg(test)]
373            StoreSessionConnection::TestUnavailable => Err(AppError::Validation(
374                "metadata repositories are not available for lake-only test sessions".to_string(),
375            )),
376        }
377    }
378
379    pub(crate) fn acquire_study_lifecycle(
380        &self,
381        study_id: ahri_tre_types::StudyId,
382    ) -> Result<ahri_tre_pgmeta::PgStudyLifecycleGuard, AppError> {
383        let repository = match &self.store {
384            StoreSessionConnection::Direct(connection) => {
385                PgMetadataRepository::from_shared(Arc::clone(connection))
386            }
387            StoreSessionConnection::OAuth(connection) => {
388                PgMetadataRepository::from_oauth_shared(Arc::clone(connection))
389            }
390            #[cfg(test)]
391            StoreSessionConnection::TestUnavailable => {
392                return Err(AppError::Validation(
393                    "Lifecycle metadata capability is unavailable".into(),
394                ));
395            }
396        };
397        repository
398            .acquire_study_lifecycle(study_id)
399            .map_err(|error| match error {
400                ahri_tre_core::CoreError::Conflict(_) => AppError::Conflict(
401                    "Study lifecycle conflicts with an active writer or deletion".into(),
402                ),
403                _ => AppError::Infrastructure("Study lifecycle exclusion is unavailable".into()),
404            })
405    }
406
407    pub fn lake_connection(&self) -> &duckdb::Connection {
408        &self.lake
409    }
410
411    pub fn lake_connection_mut(&mut self) -> &mut duckdb::Connection {
412        &mut self.lake
413    }
414
415    pub fn query_lake<A>(
416        &self,
417        request: &LakeQueryRequest,
418        authorizer: &A,
419    ) -> Result<LakeQueryResult, AppError>
420    where
421        A: LakeQueryAuthorizer + ?Sized,
422    {
423        execute_lake_query(&self.lake, request, authorizer)
424    }
425}
426
427impl std::fmt::Debug for DataStoreSession {
428    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
429        f.debug_struct("DataStoreSession")
430            .field("runtime", &"DataStoreRuntime(..)")
431            .field("store", &"live store handle")
432            .field("lake", &"duckdb::Connection(..)")
433            .finish()
434    }
435}