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#[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 #[error("Dataset admission outcome is unresolved; retain output for reconciliation")]
19 OutcomeUnknown,
20}
21
22#[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,
38 Admitted,
39}
40
41impl PgMetadataRepository<'_> {
42 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 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 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 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 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}