Skip to main content

ahri_tre_pgmeta/
dataset_reservation.rs

1use std::sync::Arc;
2
3use ahri_tre_core::{
4    CoreError, DatasetAttemptExecutor, DatasetExecutorIdentity, DatasetExecutorLease,
5    DatasetOutputAuthority, DatasetOutputKey, DatasetOutputVersion, DatasetWriteAuthority,
6};
7use uuid::Uuid;
8
9use crate::dataset_admission::ReceiptRow;
10use crate::{MetadataExecutor, MetadataTransaction, PgMetaError, PgMetadataRepository};
11
12#[derive(Debug, thiserror::Error)]
13pub enum DatasetReservationError {
14    #[error("Dataset output is already reserved")]
15    Conflict,
16    #[error("Dataset output ownership is unavailable")]
17    Unavailable(#[source] CoreError),
18}
19
20/// Private attempt identity; this is independent of any public operation record.
21#[derive(Debug, Clone)]
22pub struct RecoverableDatasetAttempt {
23    /// Existing opaque Operation join; untracked attempts have none.
24    pub operation_id: Option<Uuid>,
25    pub attempt_id: Uuid,
26    pub owner: DatasetExecutorIdentity,
27    pub output: DatasetOutputKey,
28    pub version: Option<DatasetOutputVersion>,
29    pub creation_intent: bool,
30    pub scratch_root_identity: Option<String>,
31}
32
33/// An accepted writer holds process exclusion until its last synchronous effect
34/// returns. Dropping it leaves the durable reservation for reconciliation.
35pub struct PgDatasetReservation<'connection> {
36    repository: PgMetadataRepository<'connection>,
37    executor: Arc<dyn DatasetExecutorLease>,
38    _active: Box<dyn DatasetAttemptExecutor>,
39    attempt: RecoverableDatasetAttempt,
40    version: Option<DatasetOutputVersion>,
41}
42
43impl std::fmt::Debug for PgDatasetReservation<'_> {
44    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
45        f.debug_struct("PgDatasetReservation")
46            .field("attempt", &self.attempt)
47            .finish_non_exhaustive()
48    }
49}
50
51/// Exclusive cleanup authority for one known stopped attempt. It has no writer
52/// or metadata-admission interface and cannot replay the interrupted operation.
53pub struct PgDatasetRecovery<'connection> {
54    reservation: PgDatasetReservation<'connection>,
55}
56
57impl PgDatasetRecovery<'_> {
58    pub fn attempt_id(&self) -> Uuid {
59        self.reservation.attempt_id()
60    }
61    pub fn requires_output_cleanup(&self) -> bool {
62        self.reservation.attempt.creation_intent
63    }
64    pub fn release_after_cleanup(self) -> Result<(), CoreError> {
65        self.reservation.release_after_cleanup()
66    }
67}
68impl DatasetOutputAuthority for PgDatasetRecovery<'_> {
69    fn authorize_scratch_root(&mut self, identity: &str) -> Result<(), CoreError> {
70        if self
71            .reservation
72            .attempt
73            .scratch_root_identity
74            .as_deref()
75            .is_some_and(|expected| expected != identity)
76        {
77            return Err(ownership_unavailable());
78        }
79        self.ensure_current()
80    }
81    fn attempt_id(&self) -> Uuid {
82        self.reservation.attempt_id()
83    }
84    fn output_version(&self) -> Result<&DatasetOutputVersion, CoreError> {
85        self.reservation.output_version()
86    }
87    fn ensure_current(&mut self) -> Result<(), CoreError> {
88        self.reservation.ensure_current()
89    }
90}
91
92impl<'connection> PgMetadataRepository<'connection> {
93    /// Reserves output before version allocation, source staging, or Lake effects.
94    /// Conflicting acceptance rolls back without an accepted private attempt.
95    pub fn reserve_dataset_output(
96        &self,
97        executor: Arc<dyn DatasetExecutorLease>,
98        output: DatasetOutputKey,
99    ) -> Result<PgDatasetReservation<'connection>, DatasetReservationError> {
100        self.reserve_dataset_output_with(executor, output, |_, _| Ok(()))
101            .map(|(reservation, ())| reservation)
102    }
103
104    /// Pending operation/idempotency repositories may participate in acceptance
105    /// through this scoped metadata callback. No operation is required by writers.
106    pub fn reserve_dataset_output_with<T>(
107        &self,
108        executor: Arc<dyn DatasetExecutorLease>,
109        output: DatasetOutputKey,
110        accept: impl FnOnce(
111            PgMetadataRepository<'_>,
112            &RecoverableDatasetAttempt,
113        ) -> Result<T, CoreError>,
114    ) -> Result<(PgDatasetReservation<'connection>, T), DatasetReservationError> {
115        if output.lake_name.is_empty() {
116            return Err(DatasetReservationError::Unavailable(CoreError::Validation(
117                "Dataset output identity is invalid".into(),
118            )));
119        }
120        let attempt = RecoverableDatasetAttempt {
121            operation_id: None,
122            attempt_id: Uuid::new_v4(),
123            owner: executor.identity(),
124            output,
125            version: None,
126            creation_intent: false,
127            scratch_root_identity: None,
128        };
129        let active = executor
130            .clone()
131            .begin_attempt(attempt.attempt_id)
132            .map_err(DatasetReservationError::Unavailable)?;
133        let accepted = self
134            .with_dataset_acceptance(|repository| {
135                let reserved = repository.with_session_connection(
136                    "reserve Dataset output",
137                    |c| reserve(c, &attempt),
138                    |c| reserve(c, &attempt),
139                )?;
140                if !reserved {
141                    return Err(CoreError::Conflict(
142                        "Dataset output is already reserved".into(),
143                    ));
144                }
145                accept(repository, &attempt)
146            })
147            .map_err(|error| match error {
148                CoreError::Conflict(_) => DatasetReservationError::Conflict,
149                error => DatasetReservationError::Unavailable(error),
150            })?;
151        Ok((
152            PgDatasetReservation {
153                repository: self.clone(),
154                executor,
155                _active: active,
156                attempt,
157                version: None,
158            },
159            accepted,
160        ))
161    }
162
163    /// Lists only attempts whose executor is excluded by the retained local
164    /// coordinator. Metadata RLS limits ordinary actors to their own evidence;
165    /// configured maintenance connections may see all actors' attempts.
166    pub fn recoverable_dataset_attempts(
167        &self,
168        executor: Arc<dyn DatasetExecutorLease>,
169    ) -> Result<Vec<RecoverableDatasetAttempt>, CoreError> {
170        let attempts =
171            self.with_session_connection("inspect Dataset attempts", read_attempts, read_attempts)?;
172        Ok(attempts
173            .into_iter()
174            .filter(|attempt| executor.can_recover(attempt.owner, attempt.attempt_id))
175            .collect())
176    }
177
178    /// Claims a known stopped attempt before any physical reconciliation. The
179    /// reservation row orders this claim against an in-flight metadata commit.
180    pub fn claim_dataset_recovery(
181        &self,
182        executor: Arc<dyn DatasetExecutorLease>,
183        attempt_id: Uuid,
184    ) -> Result<PgDatasetRecovery<'connection>, CoreError> {
185        let attempt = self
186            .recoverable_dataset_attempts(executor.clone())?
187            .into_iter()
188            .find(|attempt| attempt.attempt_id == attempt_id)
189            .ok_or_else(ownership_unavailable)?;
190        let active = executor.clone().begin_attempt(attempt_id)?;
191        let owner = executor.identity();
192        let claimed = self.with_session_connection(
193            "claim Dataset cleanup",
194            |connection| claim_recovery(connection, &attempt, owner),
195            |connection| claim_recovery(connection, &attempt, owner),
196        )?;
197        let version = claimed.version.clone();
198        Ok(PgDatasetRecovery {
199            reservation: PgDatasetReservation {
200                repository: self.clone(),
201                executor,
202                _active: active,
203                attempt: claimed,
204                version,
205            },
206        })
207    }
208}
209
210impl PgDatasetReservation<'_> {
211    pub fn attempt_id(&self) -> Uuid {
212        self.attempt.attempt_id
213    }
214
215    pub fn requires_output_cleanup(&self) -> bool {
216        self.attempt.creation_intent
217    }
218
219    /// Binds version coordinates chosen after reservation, before any Lake effect.
220    pub fn bind_version(
221        &mut self,
222        asset: &ahri_tre_types::AssetRecord,
223        version: &ahri_tre_types::AssetVersionRecord,
224    ) -> Result<(), CoreError> {
225        if self.version.is_some()
226            || asset.study_id != self.attempt.output.study_id
227            || asset.asset_id != version.asset_id
228        {
229            return Err(ownership_unavailable());
230        }
231        self.ensure_current()?;
232        self.repository.with_session_connection(
233            "bind Dataset output version",
234            |connection| bind_version(connection, &self.attempt, asset, version),
235            |connection| bind_version(connection, &self.attempt, asset, version),
236        )?;
237        self.version = Some(DatasetOutputVersion {
238            output: self.attempt.output.clone(),
239            version_id: version.version_id,
240            major: version.major,
241            minor: version.minor,
242            patch: version.patch,
243        });
244        Ok(())
245    }
246
247    pub(super) fn consume(
248        &self,
249        connection: &mut impl MetadataExecutor,
250    ) -> Result<(), PgMetaError> {
251        let version = self.version.as_ref().ok_or_else(ownership_error)?;
252        let changed = connection.execute_command("admit owned Dataset output", "
253            UPDATE public.dataset_output_attempts a SET admitted = TRUE
254            WHERE a.attempt_id = $1 AND a.version_id = $2 AND a.coordinator_id = $3 AND a.generation_id = $4
255              AND NOT a.admitted AND EXISTS (SELECT 1 FROM public.dataset_output_reservations r WHERE r.attempt_id = a.attempt_id)
256              AND EXISTS (SELECT 1 FROM public.dataset_executor_generation g WHERE g.coordinator_id = a.coordinator_id AND g.generation_id = a.generation_id)",
257            &[&self.attempt.attempt_id, &version.version_id.0, &self.attempt.owner.coordinator_id, &self.attempt.owner.generation_id])?;
258        if changed != 1 {
259            return Err(ownership_error());
260        }
261        let consumed = connection.execute_command(
262            "consume Dataset output reservation",
263            "DELETE FROM public.dataset_output_reservations WHERE attempt_id = $1",
264            &[&self.attempt.attempt_id],
265        )?;
266        if consumed != 1 {
267            return Err(ownership_error());
268        }
269        Ok(())
270    }
271
272    /// Called only after this attempt's source and physical output cleanup have
273    /// completed. Consuming the writer excludes subsequent effects through it.
274    pub fn release_after_cleanup(self) -> Result<(), CoreError> {
275        if self.executor.identity() != self.attempt.owner {
276            return Err(ownership_unavailable());
277        }
278        self.repository.with_session_connection(
279            "release Dataset output",
280            |connection| release(connection, &self.attempt),
281            |connection| release(connection, &self.attempt),
282        )
283    }
284}
285
286impl DatasetOutputAuthority for PgDatasetReservation<'_> {
287    fn authorize_scratch_root(&mut self, identity: &str) -> Result<(), CoreError> {
288        self.ensure_current()?;
289        if identity.len() != 64 {
290            return Err(ownership_unavailable());
291        }
292        self.repository.with_session_connection(
293            "bind Dataset source root",
294            |c| bind_scratch(c, self.attempt.attempt_id, identity),
295            |c| bind_scratch(c, self.attempt.attempt_id, identity),
296        )?;
297        self.attempt.scratch_root_identity = Some(identity.to_string());
298        Ok(())
299    }
300    fn attempt_id(&self) -> Uuid {
301        self.attempt.attempt_id
302    }
303    fn output_version(&self) -> Result<&DatasetOutputVersion, CoreError> {
304        self.version.as_ref().ok_or_else(ownership_unavailable)
305    }
306
307    fn ensure_current(&mut self) -> Result<(), CoreError> {
308        self.repository.with_session_connection(
309            "check Dataset output ownership",
310            |connection| ensure_current(connection, &self.attempt),
311            |connection| ensure_current(connection, &self.attempt),
312        )
313    }
314}
315
316impl DatasetWriteAuthority for PgDatasetReservation<'_> {
317    fn record_creation_intent(&mut self) -> Result<(), CoreError> {
318        self.output_version()?;
319        self.ensure_current()?;
320        self.repository.with_session_connection(
321            "record Dataset creation intent",
322            |connection| record_creation_intent(connection, &self.attempt),
323            |connection| record_creation_intent(connection, &self.attempt),
324        )?;
325        self.attempt.creation_intent = true;
326        Ok(())
327    }
328}
329
330fn bind_version(
331    connection: &mut impl MetadataExecutor,
332    attempt: &RecoverableDatasetAttempt,
333    asset: &ahri_tre_types::AssetRecord,
334    version: &ahri_tre_types::AssetVersionRecord,
335) -> Result<(), PgMetaError> {
336    let updated = connection.execute_command("bind Dataset output version", "
337        UPDATE public.dataset_output_attempts SET version_id = $2, dataset_name = $3, major = $4, minor = $5, patch = $6
338        WHERE attempt_id = $1 AND version_id IS NULL AND NOT admitted", &[&attempt.attempt_id, &version.version_id.0, &asset.name.as_str(), &version.major, &version.minor, &version.patch])?;
339    if updated != 1 {
340        return Err(ownership_error());
341    }
342    Ok(())
343}
344
345fn ensure_current(
346    connection: &mut impl MetadataExecutor,
347    attempt: &RecoverableDatasetAttempt,
348) -> Result<(), PgMetaError> {
349    let current = connection.query_optional("check Dataset output ownership", "
350        SELECT a.attempt_id FROM public.dataset_output_attempts a
351        JOIN public.dataset_output_reservations r USING (attempt_id)
352        JOIN public.dataset_executor_generation g ON g.coordinator_id = a.coordinator_id AND g.generation_id = a.generation_id
353        WHERE a.attempt_id = $1 AND a.coordinator_id = $2 AND a.generation_id = $3 AND NOT a.admitted
354          AND NOT EXISTS (SELECT 1 FROM public.dataset_admission_receipts d WHERE d.version_id = a.version_id AND d.admitted)",
355        &[&attempt.attempt_id, &attempt.owner.coordinator_id, &attempt.owner.generation_id])?;
356    if current.is_none() {
357        return Err(ownership_error());
358    }
359    Ok(())
360}
361
362fn record_creation_intent(
363    connection: &mut impl MetadataExecutor,
364    attempt: &RecoverableDatasetAttempt,
365) -> Result<(), PgMetaError> {
366    let changed = connection.execute_command("record Dataset creation intent", "UPDATE public.dataset_output_attempts SET creation_intent = TRUE WHERE attempt_id = $1 AND version_id IS NOT NULL AND NOT admitted", &[&attempt.attempt_id])?;
367    if changed != 1 {
368        return Err(ownership_error());
369    }
370    Ok(())
371}
372
373fn establish_generation<R>(
374    transaction: &mut dyn MetadataTransaction<Row = R>,
375    owner: DatasetExecutorIdentity,
376) -> Result<(), PgMetaError> {
377    transaction.execute_command("initialize Dataset executor", "
378            INSERT INTO public.dataset_executor_generation (singleton, coordinator_id, generation_id)
379            VALUES (TRUE, $1, $2) ON CONFLICT (singleton) DO NOTHING", &[&owner.coordinator_id, &owner.generation_id])?;
380    transaction.execute_command(
381        "replace excluded Dataset executor",
382        "
383            UPDATE public.dataset_executor_generation SET generation_id = $2
384            WHERE coordinator_id = $1 AND generation_id <> $2",
385        &[&owner.coordinator_id, &owner.generation_id],
386    )?;
387    if transaction.query_optional("verify Dataset executor", "SELECT singleton FROM public.dataset_executor_generation WHERE coordinator_id = $1 AND generation_id = $2", &[&owner.coordinator_id, &owner.generation_id])?.is_none() {
388            return Err(ownership_error());
389        }
390    Ok(())
391}
392
393fn reserve<E: MetadataExecutor>(
394    connection: &mut E,
395    attempt: &RecoverableDatasetAttempt,
396) -> Result<bool, PgMetaError> {
397    connection.with_transaction("reserve Dataset output", |transaction| {
398        // This transaction lock closes the race with lifecycle guard acquisition.
399        // The durable reservation then excludes deletion until terminal cleanup.
400        if transaction.query_optional("check Study lifecycle exclusion", "SELECT 1 WHERE pg_try_advisory_xact_lock_shared(hashtextextended(($1::uuid)::text, 606))", &[&attempt.output.study_id.0])?.is_none() {
401            return Ok(false);
402        }
403        establish_generation(transaction, attempt.owner)?;
404        let reserved = transaction.execute_command("reserve Dataset output", "
405            INSERT INTO public.dataset_output_reservations (study_id, lake_name, attempt_id)
406            VALUES ($1, $2, $3) ON CONFLICT DO NOTHING",
407            &[&attempt.output.study_id.0, &attempt.output.lake_name, &attempt.attempt_id])?;
408        if reserved == 0 { return Ok(false); }
409        transaction.execute_command("record private Dataset attempt", "
410            INSERT INTO public.dataset_output_attempts (attempt_id, coordinator_id, generation_id, study_id, lake_name)
411            VALUES ($1, $2, $3, $4, $5)", &[&attempt.attempt_id, &attempt.owner.coordinator_id, &attempt.owner.generation_id, &attempt.output.study_id.0, &attempt.output.lake_name])?;
412        Ok(true)
413    })
414}
415
416fn read_attempts<E: MetadataExecutor>(
417    connection: &mut E,
418) -> Result<Vec<RecoverableDatasetAttempt>, PgMetaError>
419where
420    E::Row: ReceiptRow,
421{
422    let rows = connection.query_many("inspect Dataset attempts", "
423        SELECT a.attempt_id::text, a.coordinator_id::text, a.generation_id::text, a.study_id::text, a.lake_name,
424            coalesce(a.version_id::text, ''), coalesce(a.major::text, ''), coalesce(a.minor::text, ''), coalesce(a.patch::text, ''), a.creation_intent::text, coalesce(a.scratch_root_identity, ''), coalesce(o.operation_id::text, '')
425        FROM public.dataset_output_attempts a JOIN public.dataset_output_reservations r USING (attempt_id)
426        LEFT JOIN public.operations o USING (attempt_id)
427        ORDER BY a.created_at, a.attempt_id", &[])?;
428    rows.into_iter()
429        .map(|row| {
430            let mut attempt = decode_attempt(&row)?;
431            attempt.operation_id = if row.text(11)?.is_empty() {
432                None
433            } else {
434                Some(row.text(11)?.parse().map_err(|_| ownership_error())?)
435            };
436            Ok(attempt)
437        })
438        .collect()
439}
440
441fn decode_attempt(row: &impl ReceiptRow) -> Result<RecoverableDatasetAttempt, PgMetaError> {
442    let output = DatasetOutputKey {
443        study_id: ahri_tre_types::StudyId(row.text(3)?.parse().map_err(|_| ownership_error())?),
444        lake_name: row.text(4)?.into(),
445    };
446    let version = if row.text(5)?.is_empty() {
447        None
448    } else {
449        Some(DatasetOutputVersion {
450            output: output.clone(),
451            version_id: ahri_tre_types::VersionId(
452                row.text(5)?.parse().map_err(|_| ownership_error())?,
453            ),
454            major: row.text(6)?.parse().map_err(|_| ownership_error())?,
455            minor: row.text(7)?.parse().map_err(|_| ownership_error())?,
456            patch: row.text(8)?.parse().map_err(|_| ownership_error())?,
457        })
458    };
459    Ok(RecoverableDatasetAttempt {
460        operation_id: None,
461        attempt_id: row.text(0)?.parse().map_err(|_| ownership_error())?,
462        owner: DatasetExecutorIdentity {
463            coordinator_id: row.text(1)?.parse().map_err(|_| ownership_error())?,
464            generation_id: row.text(2)?.parse().map_err(|_| ownership_error())?,
465        },
466        output,
467        version,
468        creation_intent: row.text(9)? == "true",
469        scratch_root_identity: (!row.text(10)?.is_empty())
470            .then(|| row.text(10).unwrap().to_string()),
471    })
472}
473
474fn claim_recovery<E: MetadataExecutor>(
475    connection: &mut E,
476    prior: &RecoverableDatasetAttempt,
477    owner: DatasetExecutorIdentity,
478) -> Result<RecoverableDatasetAttempt, PgMetaError>
479where
480    E::Row: ReceiptRow,
481{
482    connection.with_transaction("claim Dataset cleanup", |transaction| {
483        // Use the same generation-before-reservation lock order as acceptance.
484        establish_generation(transaction, owner)?;
485        // Lock before reading outcome; a final admission that already consumed ownership wins.
486        if transaction.query_optional("lock interrupted Dataset reservation", "SELECT attempt_id FROM public.dataset_output_reservations WHERE attempt_id = $1 FOR UPDATE", &[&prior.attempt_id])?.is_none() { return Err(ownership_error()); }
487        let row = transaction.query_optional("inspect interrupted Dataset attempt", "
488            SELECT a.attempt_id::text, a.coordinator_id::text, a.generation_id::text, a.study_id::text, a.lake_name,
489                coalesce(a.version_id::text, ''), coalesce(a.major::text, ''), coalesce(a.minor::text, ''), coalesce(a.patch::text, ''), a.creation_intent::text, coalesce(a.scratch_root_identity, '')
490            FROM public.dataset_output_attempts a WHERE a.attempt_id = $1 AND a.coordinator_id = $2 AND a.generation_id = $3
491                AND NOT a.admitted AND NOT EXISTS (SELECT 1 FROM public.dataset_admission_receipts d WHERE d.version_id = a.version_id AND d.admitted)",
492            &[&prior.attempt_id, &prior.owner.coordinator_id, &prior.owner.generation_id])?.ok_or_else(ownership_error)?;
493        let mut claimed = decode_attempt(&row)?;
494        transaction.execute_command("transfer stopped Dataset attempt", "UPDATE public.dataset_output_attempts SET generation_id = $2 WHERE attempt_id = $1", &[&prior.attempt_id, &owner.generation_id])?;
495        claimed.owner = owner;
496        Ok(claimed)
497    })
498}
499
500fn release<E: MetadataExecutor>(
501    connection: &mut E,
502    attempt: &RecoverableDatasetAttempt,
503) -> Result<(), PgMetaError> {
504    connection.with_transaction("release Dataset output", |transaction| {
505        // Row locking orders reconciliation after any in-flight final admission.
506        transaction.query_optional("lock Dataset reservation", "SELECT attempt_id FROM public.dataset_output_reservations WHERE attempt_id = $1 FOR UPDATE", &[&attempt.attempt_id])?;
507        let removed = transaction.execute_command("release Dataset output", "
508            DELETE FROM public.dataset_output_reservations r USING public.dataset_output_attempts a
509            WHERE r.attempt_id = a.attempt_id AND a.attempt_id = $1
510              AND a.coordinator_id = $2 AND a.generation_id = $3 AND NOT a.admitted
511              AND NOT EXISTS (SELECT 1 FROM public.dataset_admission_receipts d WHERE d.version_id = a.version_id AND d.admitted)",
512            &[&attempt.attempt_id, &attempt.owner.coordinator_id, &attempt.owner.generation_id])?;
513        if removed != 1 { return Err(ownership_error()); }
514        transaction.execute_command("remove cleaned admission evidence", "DELETE FROM public.dataset_admission_receipts d USING public.dataset_output_attempts a WHERE a.attempt_id = $1 AND d.version_id = a.version_id AND NOT d.admitted", &[&attempt.attempt_id])?;
515        transaction.execute_command("remove cleaned private Dataset attempt", "DELETE FROM public.dataset_output_attempts a WHERE attempt_id = $1 AND NOT EXISTS (SELECT 1 FROM public.operations o WHERE o.attempt_id=a.attempt_id AND o.record#>>'{detail,summary,status}' NOT IN ('completed','failed','cancelled'))", &[&attempt.attempt_id])?;
516        Ok(())
517    })
518}
519
520fn ownership_unavailable() -> CoreError {
521    CoreError::Infrastructure("Dataset output ownership could not be established".into())
522}
523fn ownership_error() -> PgMetaError {
524    PgMetaError::Decode {
525        field: "dataset_output",
526        value: String::new(),
527        message: "Dataset output ownership could not be established".into(),
528    }
529}
530
531/// Configured maintenance authority exposes only known-attempt reconciliation.
532/// The underlying administrator connection never escapes this adapter boundary.
533pub struct PgDatasetMaintenance {
534    pub(crate) repository: PgMetadataRepository<'static>,
535}
536impl PgDatasetMaintenance {
537    pub(crate) fn new(mut connection: crate::PgMetadataConnection) -> Result<Self, PgMetaError> {
538        let row = connection.client().query_one("SELECT count(*) = 4 AND coalesce(bool_and(r.rolsuper OR r.rolbypassrls OR (NOT c.relforcerowsecurity AND pg_has_role(c.relowner, 'USAGE'))), false) FROM pg_class c CROSS JOIN pg_roles r WHERE r.rolname = CURRENT_USER AND c.oid IN (to_regclass('public.dataset_output_attempts'), to_regclass('public.dataset_output_reservations'), to_regclass('public.dataset_admission_receipts'), to_regclass('public.dataset_executor_generation'))", &[]).map_err(PgMetaError::Query)?;
539        if !row.get::<_, bool>(0) {
540            return Err(ownership_error());
541        }
542        Ok(Self {
543            repository: PgMetadataRepository::new(connection),
544        })
545    }
546    pub fn recoverable_attempts(
547        &self,
548        executor: Arc<dyn DatasetExecutorLease>,
549    ) -> Result<Vec<RecoverableDatasetAttempt>, CoreError> {
550        self.repository.recoverable_dataset_attempts(executor)
551    }
552    pub fn claim(
553        &self,
554        executor: Arc<dyn DatasetExecutorLease>,
555        attempt_id: Uuid,
556    ) -> Result<PgDatasetRecovery<'static>, CoreError> {
557        self.repository.claim_dataset_recovery(executor, attempt_id)
558    }
559    pub fn has_unresolved_attempts(&self) -> Result<bool, CoreError> {
560        self.repository.with_session_connection(
561            "inspect unresolved Dataset outputs",
562            |c| read_attempts(c).map(|v| !v.is_empty()),
563            |c| read_attempts(c).map(|v| !v.is_empty()),
564        )
565    }
566}
567
568/// Administrative Lake replacement has no materialization/admission methods.
569pub struct PgDatasetReset {
570    repository: PgMetadataRepository<'static>,
571    _exclusion: Box<dyn DatasetAttemptExecutor>,
572}
573impl PgDatasetMaintenance {
574    pub fn authorize_namespace_reset(
575        &self,
576        executor: Arc<dyn DatasetExecutorLease>,
577    ) -> Result<PgDatasetReset, CoreError> {
578        ahri_tre_observability::LifecycleSpan::run(
579            ahri_tre_observability::Phase::Maintenance,
580            None,
581            ahri_tre_observability::FailureCategory::Admission,
582            || {
583                let exclusion = executor.clone().begin_maintenance()?;
584                self.repository.with_session_connection(
585                    "exclude Dataset writers for reset",
586                    |c| {
587                        c.with_transaction("establish Dataset reset executor", |t| {
588                            establish_generation(t, executor.identity())
589                        })
590                    },
591                    |c| {
592                        c.with_transaction("establish Dataset reset executor", |t| {
593                            establish_generation(t, executor.identity())
594                        })
595                    },
596                )?;
597                let mut reset = PgDatasetReset {
598                    repository: self.repository.clone(),
599                    _exclusion: exclusion,
600                };
601                ahri_tre_core::DatasetMaintenanceAuthority::ensure_exclusive(&mut reset)?;
602                Ok(reset)
603            },
604        )
605    }
606}
607impl ahri_tre_core::DatasetMaintenanceAuthority for PgDatasetReset {
608    fn ensure_exclusive(&mut self) -> Result<(), CoreError> {
609        let unresolved = self.repository.with_session_connection(
610            "verify exclusive Dataset reset",
611            |c| read_attempts(c).map(|v| !v.is_empty()),
612            |c| read_attempts(c).map(|v| !v.is_empty()),
613        )?;
614        if unresolved {
615            Err(CoreError::Conflict(
616                "Dataset output remains reserved".into(),
617            ))
618        } else {
619            Ok(())
620        }
621    }
622}
623
624fn bind_scratch(
625    connection: &mut impl MetadataExecutor,
626    attempt: Uuid,
627    identity: &str,
628) -> Result<(), PgMetaError> {
629    let changed = connection.execute_command("bind Dataset source root", "UPDATE public.dataset_output_attempts SET scratch_root_identity = $2 WHERE attempt_id = $1 AND NOT admitted AND (scratch_root_identity IS NULL OR scratch_root_identity = $2)", &[&attempt, &identity])?;
630    if changed != 1 {
631        return Err(ownership_error());
632    }
633    Ok(())
634}