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 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 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 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 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 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 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}