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#[derive(Debug, Clone)]
22pub struct RecoverableDatasetAttempt {
23 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
33pub 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
51pub 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 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 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 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 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 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 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 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 establish_generation(transaction, owner)?;
485 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 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
531pub 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
568pub 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}