Skip to main content

ahri_tre_pgmeta/
dataset_admission.rs

1use std::sync::{Arc, Mutex};
2
3use ahri_tre_core::{CoreError, DatasetOutputAuthority};
4use ahri_tre_libpq_oauth::LibpqOAuthConnection;
5use ahri_tre_types::{AssetRecord, AssetVersionRecord, NcName, StudyId, VersionId};
6
7use crate::repository::PgMetadataRepositoryConnection;
8use crate::{MetadataExecutor, PgMetaError, PgMetadataConnection, PgMetadataRepository};
9
10/// A failed admission distinguishes a known rollback from an unresolved outcome.
11#[derive(Debug, thiserror::Error)]
12pub enum DatasetAdmissionError<E> {
13    #[error("Dataset metadata workflow rejected admission")]
14    Rejected(#[source] E),
15    #[error("Dataset metadata admission failed")]
16    Metadata(#[source] CoreError),
17    /// Keep Lake output and cleanup evidence until authoritative reconciliation.
18    #[error("Dataset admission outcome is unresolved; retain output for reconciliation")]
19    OutcomeUnknown,
20}
21
22/// Private evidence retained independently of public operation history.
23#[derive(Debug, Clone, PartialEq, Eq)]
24pub struct DatasetAdmissionReceipt {
25    pub version_id: VersionId,
26    pub study_id: StudyId,
27    pub dataset_name: NcName,
28    pub major: i32,
29    pub minor: i32,
30    pub patch: i32,
31    pub state: DatasetAdmissionState,
32}
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35pub enum DatasetAdmissionState {
36    /// Pending evidence alone does not establish that an executor has stopped.
37    Pending,
38    Admitted,
39}
40
41impl PgMetadataRepository<'_> {
42    /// Admits Dataset metadata and consumes shared output ownership atomically.
43    /// The exclusive mutable writer borrow excludes concurrent Lake effects.
44    pub fn with_reserved_dataset_admission<T, E>(
45        &self,
46        reservation: &mut crate::PgDatasetReservation<'_>,
47        asset: &AssetRecord,
48        version: &AssetVersionRecord,
49        admit: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, E>,
50    ) -> Result<T, DatasetAdmissionError<E>> {
51        let span = self
52            .observation
53            .map(|context| context.span(ahri_tre_observability::Stage::Metadata));
54        let result = (|| {
55            reservation
56                .ensure_current()
57                .map_err(DatasetAdmissionError::Metadata)?;
58            if reservation
59                .output_version()
60                .map_err(DatasetAdmissionError::Metadata)?
61                .version_id
62                != version.version_id
63            {
64                return Err(DatasetAdmissionError::Metadata(CoreError::Conflict(
65                    "Dataset admission does not own this version".into(),
66                )));
67            }
68            self.admit_dataset(reservation, asset, version, admit)
69        })();
70        if let Some(span) = span {
71            use ahri_tre_observability::{FailureCategory, Outcome};
72            let (outcome, category) = match &result {
73                Ok(_) => (Outcome::Success, None),
74                Err(DatasetAdmissionError::Rejected(_)) => (Outcome::Rejected, None),
75                Err(DatasetAdmissionError::Metadata(error)) => metadata_outcome(Some(error)),
76                Err(DatasetAdmissionError::OutcomeUnknown) => {
77                    (Outcome::Unavailable, Some(FailureCategory::CommitUnknown))
78                }
79            };
80            span.finish_observed(outcome, category);
81        }
82        result
83    }
84
85    fn admit_dataset<T, E>(
86        &self,
87        reservation: &crate::PgDatasetReservation<'_>,
88        asset: &AssetRecord,
89        version: &AssetVersionRecord,
90        admit: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, E>,
91    ) -> Result<T, DatasetAdmissionError<E>> {
92        let lock_error = |_| {
93            DatasetAdmissionError::Metadata(CoreError::Infrastructure(
94                "Dataset admission connection is unavailable".into(),
95            ))
96        };
97        match &self.connection {
98            PgMetadataRepositoryConnection::Direct(connection) => {
99                let mut connection = connection.lock().map_err(lock_error)?;
100                AdmissionTransaction::new(Connection::Direct(&mut connection), asset, version)?.run(
101                    version.version_id,
102                    reservation,
103                    admit,
104                )
105            }
106            PgMetadataRepositoryConnection::OAuth(connection) => {
107                let mut connection = connection.lock().map_err(|_| {
108                    DatasetAdmissionError::Metadata(CoreError::Infrastructure(
109                        "Dataset admission connection is unavailable".into(),
110                    ))
111                })?;
112                AdmissionTransaction::new(Connection::OAuth(&mut connection), asset, version)?.run(
113                    version.version_id,
114                    reservation,
115                    admit,
116                )
117            }
118            PgMetadataRepositoryConnection::ScopedDirect(_)
119            | PgMetadataRepositoryConnection::ScopedOAuth(_) => {
120                Err(DatasetAdmissionError::Metadata(CoreError::Conflict(
121                    "Dataset admission is already in progress".into(),
122                )))
123            }
124        }
125    }
126    pub(super) fn with_dataset_acceptance<T>(
127        &self,
128        accept: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
129    ) -> Result<T, CoreError> {
130        let span = self
131            .observation
132            .map(|context| context.span(ahri_tre_observability::Stage::Metadata));
133        let result = (|| {
134            let unavailable =
135                || CoreError::Infrastructure("Dataset acceptance is unavailable".into());
136            match &self.connection {
137                PgMetadataRepositoryConnection::Direct(connection) => {
138                    let mut connection = connection.lock().map_err(|_| unavailable())?;
139                    AdmissionTransaction::accept(Connection::Direct(&mut connection), accept)
140                }
141                PgMetadataRepositoryConnection::OAuth(connection) => {
142                    let mut connection = connection.lock().map_err(|_| unavailable())?;
143                    AdmissionTransaction::accept(Connection::OAuth(&mut connection), accept)
144                }
145                _ => Err(unavailable()),
146            }
147        })();
148        if let Some(span) = span {
149            let (outcome, category) = metadata_outcome(result.as_ref().err());
150            span.finish_observed(outcome, category);
151        }
152        result
153    }
154
155    /// Reads this actor's private receipt without relying on catalogue visibility.
156    pub fn dataset_admission_receipt(
157        &self,
158        version_id: VersionId,
159    ) -> Result<Option<DatasetAdmissionReceipt>, CoreError> {
160        self.with_session_connection(
161            "read Dataset admission receipt",
162            |connection| read_receipt(connection, version_id),
163            |connection| read_receipt(connection, version_id),
164        )
165    }
166
167    /// Clears pending evidence only after the caller has completed Lake cleanup.
168    /// Admitted receipts cannot be removed through this cleanup boundary.
169    pub fn confirm_dataset_cleanup(&self, version_id: VersionId) -> Result<(), CoreError> {
170        self.with_session_connection(
171            "confirm Dataset cleanup",
172            |connection| clear_pending_receipt(connection, version_id),
173            |connection| clear_pending_receipt(connection, version_id),
174        )
175    }
176}
177
178enum Connection<'a> {
179    Direct(&'a mut PgMetadataConnection),
180    OAuth(&'a mut LibpqOAuthConnection),
181}
182
183struct AdmissionTransaction<'a> {
184    connection: Connection<'a>,
185    finished: bool,
186}
187
188impl<'a> AdmissionTransaction<'a> {
189    fn accept<T>(
190        connection: Connection<'a>,
191        accept: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, CoreError>,
192    ) -> Result<T, CoreError> {
193        let unavailable =
194            |_| CoreError::Infrastructure("Dataset acceptance outcome is unavailable".into());
195        let active = match &connection {
196            Connection::Direct(c) => c.admission_active,
197            Connection::OAuth(c) => c.is_in_transaction(),
198        };
199        if active {
200            return Err(CoreError::Conflict(
201                "Dataset acceptance is already in progress".into(),
202            ));
203        }
204        let mut transaction = Self {
205            connection,
206            finished: true,
207        };
208        transaction
209            .command("BEGIN ISOLATION LEVEL READ COMMITTED")
210            .map_err(unavailable)?;
211        transaction.finished = false;
212        if let Connection::Direct(c) = &mut transaction.connection {
213            c.admission_active = true;
214        }
215        let scoped = match &mut transaction.connection {
216            Connection::Direct(c) => {
217                PgMetadataRepositoryConnection::ScopedDirect(Arc::new(Mutex::new(&mut **c)))
218            }
219            Connection::OAuth(c) => {
220                PgMetadataRepositoryConnection::ScopedOAuth(Arc::new(Mutex::new(&mut **c)))
221            }
222        };
223        match accept(PgMetadataRepository {
224            connection: scoped,
225            observation: None,
226        }) {
227            Ok(value) => {
228                // Failed acknowledgement never starts execution. A committed
229                // private attempt, if any, remains owned until reconciliation.
230                transaction.command("COMMIT").map_err(unavailable)?;
231                transaction.finished = true;
232                Ok(value)
233            }
234            Err(error) => {
235                transaction.command("ROLLBACK").map_err(unavailable)?;
236                transaction.finished = true;
237                Err(error)
238            }
239        }
240    }
241
242    fn new<E>(
243        connection: Connection<'a>,
244        asset: &AssetRecord,
245        version: &AssetVersionRecord,
246    ) -> Result<Self, DatasetAdmissionError<E>> {
247        if asset.asset_id != version.asset_id {
248            return Err(DatasetAdmissionError::Metadata(CoreError::Validation(
249                "Dataset admission Asset/version mismatch".into(),
250            )));
251        }
252        let active = match &connection {
253            Connection::Direct(connection) => connection.admission_active,
254            Connection::OAuth(connection) => connection.is_in_transaction(),
255        };
256        if active {
257            return Err(DatasetAdmissionError::Metadata(CoreError::Conflict(
258                "Dataset admission is already in progress".into(),
259            )));
260        }
261        let mut transaction = Self {
262            connection,
263            finished: true,
264        };
265        match &mut transaction.connection {
266            Connection::Direct(connection) => create_receipt(*connection, asset, version),
267            Connection::OAuth(connection) => create_receipt(*connection, asset, version),
268        }
269        .map_err(metadata_error)?;
270        transaction
271            .command("BEGIN ISOLATION LEVEL SERIALIZABLE")
272            .map_err(metadata_error)?;
273        transaction.finished = false;
274        if let Connection::Direct(connection) = &mut transaction.connection {
275            connection.admission_active = true;
276        }
277        Ok(transaction)
278    }
279
280    fn command(&mut self, command: &str) -> Result<(), PgMetaError> {
281        match &mut self.connection {
282            Connection::Direct(connection) => connection
283                .client()
284                .batch_execute(command)
285                .map_err(PgMetaError::Transaction),
286            Connection::OAuth(connection) => connection
287                .execute_raw(command, &[])
288                .map(|_| ())
289                .map_err(|source| PgMetaError::OAuth {
290                    operation: "Dataset admission",
291                    source,
292                }),
293        }
294    }
295
296    fn admitted(
297        &mut self,
298        version_id: VersionId,
299    ) -> Result<Option<DatasetAdmissionState>, PgMetaError> {
300        let receipt = match &mut self.connection {
301            Connection::Direct(connection) => read_receipt(*connection, version_id),
302            Connection::OAuth(connection) => read_receipt(*connection, version_id),
303        }?;
304        Ok(receipt.map(|receipt| receipt.state))
305    }
306
307    fn record_admission(&mut self, version_id: VersionId) -> Result<(), PgMetaError> {
308        match &mut self.connection {
309            Connection::Direct(connection) => record_admission(*connection, version_id),
310            Connection::OAuth(connection) => record_admission(*connection, version_id),
311        }
312    }
313
314    fn run<T, E>(
315        mut self,
316        version_id: VersionId,
317        reservation: &crate::PgDatasetReservation<'_>,
318        admit: impl FnOnce(PgMetadataRepository<'_>) -> Result<T, E>,
319    ) -> Result<T, DatasetAdmissionError<E>> {
320        let scoped = match &mut self.connection {
321            Connection::Direct(connection) => PgMetadataRepositoryConnection::ScopedDirect(
322                Arc::new(Mutex::new(&mut **connection)),
323            ),
324            Connection::OAuth(connection) => {
325                PgMetadataRepositoryConnection::ScopedOAuth(Arc::new(Mutex::new(&mut **connection)))
326            }
327        };
328        let result = admit(PgMetadataRepository {
329            connection: scoped,
330            observation: None,
331        })
332        .map_err(DatasetAdmissionError::Rejected)
333        .and_then(|value| {
334            self.record_admission(version_id).map_err(metadata_error)?;
335            {
336                match &mut self.connection {
337                    Connection::Direct(connection) => reservation.consume(*connection),
338                    Connection::OAuth(connection) => reservation.consume(*connection),
339                }
340                .map_err(metadata_error)?;
341            }
342            Ok(value)
343        });
344        match result {
345            Err(error) => {
346                self.command("ROLLBACK")
347                    .map_err(|_| DatasetAdmissionError::OutcomeUnknown)?;
348                self.finished = true;
349                Err(error)
350            }
351            Ok(value) => {
352                match self.command("COMMIT") {
353                    Ok(()) => {
354                        self.finished = true;
355                        Ok(value)
356                    }
357                    Err(error) => {
358                        // A lost COMMIT response must not authorize destructive Lake
359                        // compensation. Clear any aborted transaction before readback.
360                        let _ = self.command("ROLLBACK");
361                        match self.admitted(version_id) {
362                            Ok(Some(DatasetAdmissionState::Admitted)) => {
363                                self.finished = true;
364                                Ok(value)
365                            }
366                            Ok(Some(DatasetAdmissionState::Pending)) => {
367                                self.finished = true;
368                                Err(metadata_error(error))
369                            }
370                            Ok(None) | Err(_) => Err(DatasetAdmissionError::OutcomeUnknown),
371                        }
372                    }
373                }
374            }
375        }
376    }
377}
378
379impl Drop for AdmissionTransaction<'_> {
380    fn drop(&mut self) {
381        if !self.finished {
382            let _ = self.command("ROLLBACK");
383        }
384        if let Connection::Direct(connection) = &mut self.connection {
385            connection.admission_active = false;
386        }
387    }
388}
389
390fn metadata_error<E>(_error: PgMetaError) -> DatasetAdmissionError<E> {
391    DatasetAdmissionError::Metadata(CoreError::Infrastructure(
392        "Dataset metadata admission failed".into(),
393    ))
394}
395
396fn create_receipt(
397    connection: &mut impl MetadataExecutor,
398    asset: &AssetRecord,
399    version: &AssetVersionRecord,
400) -> Result<(), PgMetaError> {
401    connection.execute_command("prepare Dataset admission receipt", "
402        INSERT INTO public.dataset_admission_receipts (version_id, study_id, dataset_name, major, minor, patch)
403        VALUES ($1, $2, $3, $4, $5, $6)",
404        &[&version.version_id.0, &asset.study_id.0, &asset.name.as_str(), &version.major, &version.minor, &version.patch])?;
405    Ok(())
406}
407
408fn record_admission(
409    connection: &mut impl MetadataExecutor,
410    version_id: VersionId,
411) -> Result<(), PgMetaError> {
412    let changed = connection.execute_command(
413        "complete Dataset admission receipt",
414        "
415        UPDATE public.dataset_admission_receipts SET admitted = TRUE
416        WHERE version_id = $1 AND actor = CURRENT_USER AND NOT admitted
417          AND EXISTS (SELECT 1 FROM public.datasets WHERE dataset_id = $1)
418          AND EXISTS (SELECT 1 FROM public.transformation_outputs WHERE version_id = $1)",
419        &[&version_id.0],
420    )?;
421    if changed != 1 {
422        return Err(receipt_decode_error(
423            "admission",
424            "complete Dataset and provenance required",
425        ));
426    }
427    Ok(())
428}
429
430fn clear_pending_receipt(
431    connection: &mut impl MetadataExecutor,
432    version_id: VersionId,
433) -> Result<(), PgMetaError> {
434    connection.execute_command(
435        "clear Dataset cleanup evidence",
436        "
437        DELETE FROM public.dataset_admission_receipts
438        WHERE version_id = $1 AND actor = CURRENT_USER AND NOT admitted",
439        &[&version_id.0],
440    )?;
441    Ok(())
442}
443
444pub(super) trait ReceiptRow {
445    fn text(&self, index: usize) -> Result<&str, PgMetaError>;
446}
447impl ReceiptRow for postgres::Row {
448    fn text(&self, index: usize) -> Result<&str, PgMetaError> {
449        self.try_get(index).map_err(PgMetaError::Query)
450    }
451}
452impl ReceiptRow for ahri_tre_libpq_oauth::LibpqOAuthRow {
453    fn text(&self, index: usize) -> Result<&str, PgMetaError> {
454        self.require(index).map_err(|source| PgMetaError::OAuth {
455            operation: "read Dataset admission receipt",
456            source,
457        })
458    }
459}
460fn read_receipt<E: MetadataExecutor>(
461    connection: &mut E,
462    version_id: VersionId,
463) -> Result<Option<DatasetAdmissionReceipt>, PgMetaError>
464where
465    E::Row: ReceiptRow,
466{
467    let row = connection.query_optional(
468        "read Dataset admission receipt",
469        "
470        SELECT study_id::text, dataset_name, major::text, minor::text, patch::text, admitted::text
471        FROM public.dataset_admission_receipts WHERE version_id = $1 AND actor = CURRENT_USER",
472        &[&version_id.0],
473    )?;
474    row.map(|row| {
475        Ok(DatasetAdmissionReceipt {
476            version_id,
477            study_id: StudyId(
478                row.text(0)?
479                    .parse()
480                    .map_err(|_| receipt_decode_error("study_id", "invalid UUID"))?,
481            ),
482            dataset_name: NcName::parse(row.text(1)?)
483                .map_err(|_| receipt_decode_error("dataset_name", "invalid Dataset name"))?,
484            major: row
485                .text(2)?
486                .parse()
487                .map_err(|_| receipt_decode_error("major", "invalid version"))?,
488            minor: row
489                .text(3)?
490                .parse()
491                .map_err(|_| receipt_decode_error("minor", "invalid version"))?,
492            patch: row
493                .text(4)?
494                .parse()
495                .map_err(|_| receipt_decode_error("patch", "invalid version"))?,
496            state: match row.text(5)? {
497                "true" => DatasetAdmissionState::Admitted,
498                "false" => DatasetAdmissionState::Pending,
499                _ => return Err(receipt_decode_error("admitted", "invalid outcome")),
500            },
501        })
502    })
503    .transpose()
504}
505fn receipt_decode_error(field: &'static str, message: &str) -> PgMetaError {
506    PgMetaError::Decode {
507        field,
508        value: "Dataset admission receipt".into(),
509        message: message.into(),
510    }
511}
512
513fn metadata_outcome(
514    error: Option<&CoreError>,
515) -> (
516    ahri_tre_observability::Outcome,
517    Option<ahri_tre_observability::FailureCategory>,
518) {
519    use ahri_tre_observability::{FailureCategory, Outcome};
520    match error {
521        None => (Outcome::Success, None),
522        Some(CoreError::Infrastructure(_)) => {
523            (Outcome::Unavailable, Some(FailureCategory::Metadata))
524        }
525        Some(_) => (Outcome::Rejected, None),
526    }
527}