ahri_tre_pgmeta/
lifecycle_guard.rs1use ahri_tre_core::CoreError;
2use ahri_tre_types::StudyId;
3
4use crate::{MetadataExecutor, PgMetaError, PgMetadataRepository};
5
6pub struct PgStudyLifecycleGuard {
10 repository: PgMetadataRepository<'static>,
11 study_id: StudyId,
12}
13
14impl PgMetadataRepository<'static> {
15 pub fn acquire_study_lifecycle(
16 &self,
17 study_id: StudyId,
18 ) -> Result<PgStudyLifecycleGuard, CoreError> {
19 let acquired = self.with_session_connection(
20 "exclude concurrent Study lifecycle work",
21 |c| acquire(c, study_id),
22 |c| acquire(c, study_id),
23 )?;
24 if !acquired {
25 return Err(CoreError::Conflict(
26 "Study lifecycle work is already in progress".into(),
27 ));
28 }
29 let guard = PgStudyLifecycleGuard {
30 repository: self.clone(),
31 study_id,
32 };
33 let busy = self.with_session_connection(
34 "inspect Study output reservations",
35 |c| reserved(c, study_id),
36 |c| reserved(c, study_id),
37 )?;
38 if busy {
39 return Err(CoreError::Conflict(
40 "Study has an active Dataset output reservation".into(),
41 ));
42 }
43 Ok(guard)
44 }
45}
46
47fn acquire<E: MetadataExecutor>(c: &mut E, study: StudyId) -> Result<bool, PgMetaError> {
48 c.query_optional(
49 "lock Study lifecycle",
50 "SELECT 1 WHERE pg_try_advisory_lock(hashtextextended(($1::uuid)::text, 606))",
51 &[&study.0],
52 )
53 .map(|row| row.is_some())
54}
55
56fn reserved<E: MetadataExecutor>(c: &mut E, study: StudyId) -> Result<bool, PgMetaError> {
57 c.query_optional(
58 "inspect Study output reservations",
59 "SELECT 1 WHERE public.tre_study_has_dataset_reservations($1)",
60 &[&study.0],
61 )
62 .map(|row| row.is_some())
63}
64
65fn release<E: MetadataExecutor>(c: &mut E, study: StudyId) -> Result<(), PgMetaError> {
66 c.query_optional(
67 "unlock Study lifecycle",
68 "SELECT pg_advisory_unlock(hashtextextended(($1::uuid)::text, 606))",
69 &[&study.0],
70 )
71 .map(|_| ())
72}
73
74impl Drop for PgStudyLifecycleGuard {
75 fn drop(&mut self) {
76 let _ = self.repository.with_session_connection(
78 "unlock Study lifecycle",
79 |c| release(c, self.study_id),
80 |c| release(c, self.study_id),
81 );
82 }
83}