Skip to main content

ahri_tre_pgmeta/
lifecycle_guard.rs

1use ahri_tre_core::CoreError;
2use ahri_tre_types::StudyId;
3
4use crate::{MetadataExecutor, PgMetaError, PgMetadataRepository};
5
6/// A synchronous destructive workflow excludes new Dataset admissions in its
7/// Study and refuses to start while any durable output reservation remains.
8/// The connection owns the lock, so connection loss releases it automatically.
9pub 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        // A failed unlock fails closed: the connection still owns the lock.
77        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}