1mod history;
3use super::datafile_materialization::{
4 DatafileLifecycle, DatafileStage, ResolvedDatafileMaterialization,
5};
6use super::*;
7use ahri_tre_core::{OperationRepository, OperationResourceScope, StoredOperation};
8use ahri_tre_pgmeta::PgMetadataRepository;
9use ahri_tre_protocol::operation::*;
10use ahri_tre_protocol::pagination::Page;
11use ahri_tre_protocol::session::SessionStatusPayload;
12use ahri_tre_protocol::{ProtocolError, ProtocolErrorCode, RequestId};
13use chrono::{Duration, Utc};
14use uuid::Uuid;
15
16#[derive(Debug, Clone, Copy)]
18pub enum OperationCancellationPolicy {
19 Cooperative,
20 Unsupported,
21}
22impl OperationCancellationPolicy {
23 pub fn request(
24 self,
25 status: OperationStatus,
26 cancellable: bool,
27 ) -> Result<OperationStatus, Box<ProtocolError>> {
28 if matches!(self, Self::Unsupported) {
29 return Err(Box::new(ProtocolError::new(
30 ProtocolErrorCode::UnsupportedOperation,
31 "Operation does not support cancellation",
32 )));
33 }
34 if status == OperationStatus::CancelRequested {
35 return Ok(status);
36 }
37 if status.is_terminal() || !cancellable {
38 return Err(Box::new(cancellation_conflict()));
39 }
40 Ok(OperationStatus::CancelRequested)
41 }
42}
43
44#[derive(Clone, Copy, PartialEq, Eq)]
46pub enum OperationSessionAuthority {
47 Active,
48 Closed,
49 Lost,
50}
51
52pub struct OperationContext {
54 pub request_id: RequestId,
55 pub observation: ahri_tre_observability::CorrelationContext,
56 pub owner: String,
57 pub session: SessionStatusPayload,
58 pub authority: Box<dyn Fn() -> OperationSessionAuthority + Send>,
59}
60
61pub struct OperationControl {
63 repository: PgMetadataRepository<'static>,
64 deployment_id: Uuid,
65 datastore_id: Uuid,
66}
67
68impl OperationControl {
69 pub(crate) fn new(
70 repository: PgMetadataRepository<'static>,
71 deployment_id: Uuid,
72 datastore_id: Uuid,
73 ) -> Self {
74 Self {
75 repository,
76 deployment_id,
77 datastore_id,
78 }
79 }
80
81 fn service(&self) -> AppService {
82 let mut service = AppService::new(
83 crate::session::ScopedSessionMetadataRepositories::from_repository(
84 std::sync::Arc::new(self.repository.clone()),
85 )
86 .into(),
87 );
88 service.datastore_id = Some(self.datastore_id);
89 service
90 }
91
92 pub async fn get(
93 &self,
94 owner: &str,
95 operation: &OperationRef,
96 events: Option<&ahri_tre_protocol::pagination::PageRequest>,
97 ) -> Result<OperationDetail, ProtocolError> {
98 self.service()
99 .get_operation(self, owner, operation, events)
100 .await
101 }
102
103 pub async fn cancel(
104 &self,
105 owner: &str,
106 operation: &OperationRef,
107 ) -> Result<OperationSummary, ProtocolError> {
108 self.service()
109 .cancel_operation(self, owner, operation)
110 .await
111 }
112
113 pub async fn result(
114 &self,
115 owner: &str,
116 operation: &OperationRef,
117 ) -> Result<OperationResult, ProtocolError> {
118 self.service()
119 .get_operation_result(self, owner, operation)
120 .await
121 }
122}
123
124pub enum OperationAdmission {
127 Retained(Box<OperationSummary>),
128 New(OperationSubmission),
129}
130
131pub struct OperationSubmission {
132 key: Option<ahri_tre_pgmeta::PgOperationKey>,
133 preconditions: Vec<ahri_tre_protocol::mutation::Precondition>,
134}
135
136impl OperationControl {
137 pub async fn admit(
138 &self,
139 owner: &str,
140 request: &mut ahri_tre_protocol::ingest::DatasetCreateFromDatafileRequest,
141 current_study: Option<ahri_tre_protocol::refs::ObjectRef>,
142 ) -> Result<OperationAdmission, ProtocolError> {
143 use ahri_tre_protocol::mutation::Precondition;
144 let service = self.service();
147 let domain = service
148 .get_domain(GetDomainRequest {
149 domain: service
150 .catalogue_domain_selector(request.domain.clone())
151 .map_err(crate::projections::catalogue_protocol_error)?,
152 })
153 .await
154 .map_err(crate::projections::catalogue_protocol_error)?
155 .ok_or_else(|| {
156 ProtocolError::new(ProtocolErrorCode::NotFound, "Domain was not found")
157 })?;
158 request.domain = ahri_tre_protocol::domain::DomainSelector::Id {
159 domain: crate::projections::domain_ref(
160 ahri_tre_protocol::PublicUuid::from_uuid(self.datastore_id),
161 domain.domain_id,
162 ),
163 };
164 let output = service
165 .resolve_session_study(self.datastore_id, request.study.take(), current_study)
166 .await?;
167 let canonical = |reference: ahri_tre_protocol::refs::ObjectRef| {
168 ahri_tre_protocol::study::StudySelector::Id { study: reference }
169 };
170 request.study = Some(canonical(output));
171 let (source_study, catalog, pinned) = service
175 .resolve_catalogue_asset(request.source.asset.clone())
176 .await
177 .map_err(crate::projections::catalogue_protocol_error)?;
178 if source_study.study_id.0 != output.id.as_uuid() {
179 return Err(ProtocolError::new(
180 ProtocolErrorCode::ValidationFailed,
181 "This materialization requires source and Output in the same Study",
182 ));
183 }
184 if catalog.asset.asset_type != ahri_tre_types::AssetType::File {
185 return Err(ProtocolError::new(
186 ProtocolErrorCode::ValidationFailed,
187 "A Datafile source is required",
188 ));
189 }
190 let pinned_latest = pinned.is_some()
193 && request
194 .source
195 .version
196 .as_deref()
197 .is_some_and(|value| value.trim().eq_ignore_ascii_case("latest"));
198 let pinned_version = if pinned.is_some() {
199 Some(
200 service
201 .catalogue_version(
202 &catalog,
203 pinned,
204 if pinned_latest {
205 None
206 } else {
207 request.source.version.as_deref()
208 },
209 )
210 .map_err(crate::projections::catalogue_protocol_error)?,
211 )
212 } else {
213 None
214 };
215 let source_asset = ahri_tre_protocol::ingest::datafile_ref(
216 ahri_tre_protocol::PublicUuid::from_uuid(self.datastore_id),
217 catalog.asset.asset_id,
218 );
219 request.source.asset = ahri_tre_protocol::asset::AssetSelector::Id {
220 asset: pinned
221 .map(|version| {
222 ahri_tre_protocol::ingest::version_ref(
223 ahri_tre_protocol::PublicUuid::from_uuid(self.datastore_id),
224 version,
225 )
226 })
227 .unwrap_or(source_asset),
228 study: Some(canonical(output)),
229 asset_type: Some(ahri_tre_types::AssetType::File),
230 };
231 let mut preconditions = request.mutation.preconditions.clone();
232 for condition in &mut preconditions {
233 let supported = match condition {
234 Precondition::Exists { target } => matches!(target.as_str(), "source" | "dataset"),
235 Precondition::DoesNotExist { target } => target == "dataset",
236 Precondition::ResourceVersion { target, version } => {
237 if let Some((major, minor, patch)) = parse_semver_selector(version) {
238 *version = format!("{major}.{minor}.{patch}");
239 target == "source"
240 } else {
241 false
242 }
243 }
244 _ => false,
245 };
246 if !supported {
247 return Err(ProtocolError::new(
248 ProtocolErrorCode::ValidationFailed,
249 "Unsupported materialization precondition",
250 ));
251 }
252 }
253 let key = if let Some(key) = &request.mutation.idempotency_key {
254 use sha2::{Digest, Sha256};
255 let mut intent = request.clone();
258 intent.session = None;
259 intent.mutation.idempotency_key = None;
260 intent.mutation.preconditions = preconditions.clone();
261 if let Some(version) = &pinned_version {
264 intent.source.asset = ahri_tre_protocol::asset::AssetSelector::Id {
265 asset: source_asset,
266 study: Some(canonical(output)),
267 asset_type: Some(ahri_tre_types::AssetType::File),
268 };
269 intent.source.version = Some(format!(
270 "{}.{}.{}",
271 version.major, version.minor, version.patch
272 ));
273 }
274 intent.source.version = match intent.source.version.as_deref().map(str::trim) {
275 None | Some("") => None,
276 Some(value) if value.eq_ignore_ascii_case("latest") => None,
277 Some(value) => Some(
278 parse_semver_selector(value)
279 .map(|(major, minor, patch)| format!("{major}.{minor}.{patch}"))
280 .unwrap_or_else(|| value.to_string()),
281 ),
282 };
283 let canonical = serde_json::to_value(intent).map_err(|_| unavailable())?;
284 let canonical = serde_json::to_vec(&canonical).map_err(|_| unavailable())?;
285 let scope = serde_json::to_vec(&(
286 owner,
287 self.deployment_id,
288 self.datastore_id,
289 ahri_tre_protocol::request::kind::INGEST_DATASET_FROM_DATAFILE,
290 key.as_str(),
291 ))
292 .map_err(|_| unavailable())?;
293 let key = self
294 .repository
295 .lock_operation_key(
296 Sha256::digest(scope)
297 .iter()
298 .map(|byte| format!("{byte:02x}"))
299 .collect(),
300 Sha256::digest(canonical)
301 .iter()
302 .map(|byte| format!("{byte:02x}"))
303 .collect(),
304 )
305 .map_err(|_| unavailable())?;
306 if let Some((record, same_intent)) = key.retained().map_err(|_| unavailable())? {
307 let record = self
309 .service()
310 .authorized_operation(self, owner, &record.detail.summary.operation)
311 .await?;
312 if !same_intent {
313 let mut error = ProtocolError::new(
314 ProtocolErrorCode::Conflict,
315 "Idempotency key was accepted with different intent",
316 );
317 error.details = Some(
318 ahri_tre_protocol::public_error::ProtocolErrorDetails::IdempotencyConflict,
319 );
320 return Err(error);
321 }
322 return Ok(OperationAdmission::Retained(Box::new(
323 record.detail.summary,
324 )));
325 }
326 Some(key)
327 } else {
328 None
329 };
330 if pinned_latest
333 && pinned
334 != catalog
335 .latest_version
336 .as_ref()
337 .map(|version| version.version_id)
338 {
339 return Err(ProtocolError::new(
340 ProtocolErrorCode::Conflict,
341 "Version selectors disagree",
342 ));
343 }
344 if pinned.is_some() {
345 request.source.version = None;
346 }
347 Ok(OperationAdmission::New(OperationSubmission {
348 key,
349 preconditions,
350 }))
351 }
352}
353
354struct TrackedDatafileLifecycle {
355 submission: OperationSubmission,
356 rejection: Option<ProtocolError>,
357 context: OperationContext,
358 repository: PgMetadataRepository<'static>,
359 id: Uuid,
360 datastore_id: Uuid,
361 cleanup_succeeded: Option<bool>,
362 cancellation_observed: bool,
363 authority_lost: bool,
364 source: Option<(ahri_tre_types::StudyRecord, DataFileMetadata)>,
365 acceptance: Option<OperationSummary>,
366 notify: Option<Box<dyn FnOnce(OperationSummary) + Send>>,
367}
368
369impl TrackedDatafileLifecycle {
370 fn check_session_authority(&mut self) -> Result<(), AppError> {
371 match (self.context.authority)() {
372 OperationSessionAuthority::Active => {
373 if self
374 .repository
375 .operation_stage_authorized(self.id)
376 .unwrap_or(false)
377 {
378 return Ok(());
379 }
380 self.authority_lost = true;
381 Err(AppError::Conflict("Session authority has ended".into()))
382 }
383 OperationSessionAuthority::Lost => {
384 self.authority_lost = true;
385 Err(AppError::Conflict("Session authority has ended".into()))
386 }
387 OperationSessionAuthority::Closed => {
388 self.cancellation_observed = true;
389 update_operation(&self.repository, self.id, |record| {
390 if record.detail.summary.status.is_terminal()
391 || record.detail.summary.status == OperationStatus::CancelRequested
392 {
393 return Ok(false);
394 }
395 record.detail.summary.status = OperationStatus::CancelRequested;
396 record.detail.summary.cancellable = false;
397 let mut entry = event(
398 record.detail.events.items.len() as u64 + 1,
399 Utc::now().to_rfc3339(),
400 "Session close requested cancellation",
401 );
402 entry.kind = OperationEventKind::CancellationRequested;
403 record.detail.events.items.push(entry);
404 Ok(true)
405 })?;
406 Err(AppError::Conflict(
407 "Operation cancellation requested".into(),
408 ))
409 }
410 }
411 }
412}
413
414impl AppService {
415 pub async fn derive_dataset_from_managed_datafile_tracked(
417 &self,
418 session: &mut DataStoreSession,
419 request: DeriveDatasetFromManagedDataFileRequest,
420 mut context: OperationContext,
421 submission: OperationSubmission,
422 notify: impl FnOnce(OperationSummary) + Send + 'static,
423 ) -> Result<StartResult<ahri_tre_protocol::ingest::DatasetMaterializationResponse>, ProtocolError>
424 {
425 let id = Uuid::new_v4();
426 let span = context
427 .observation
428 .with_operation_id(id)
429 .span(ahri_tre_observability::Stage::Application);
430 context.observation = span.context();
431 span.record_start();
432 let result = async {
433 session
434 .require_operation_recovery()
435 .map_err(crate::projections::dataset_materialization_protocol_error)?;
436 if (context.authority)() != OperationSessionAuthority::Active {
437 return Err(ProtocolError::new(
438 ProtocolErrorCode::Forbidden,
439 "Session authority has ended",
440 ));
441 }
442 let repository = session
443 .dataset_admission_repository()
444 .ok_or_else(unavailable)?;
445 let mut lifecycle = TrackedDatafileLifecycle {
446 submission,
447 rejection: None,
448 context,
449 repository,
450 id,
451 datastore_id: session
452 .runtime()
453 .identity_binding
454 .as_ref()
455 .ok_or_else(unavailable)?
456 .datastore_id,
457 source: None,
458 cleanup_succeeded: None,
459 cancellation_observed: false,
460 authority_lost: false,
461 acceptance: None,
462 notify: Some(Box::new(notify)),
463 };
464 let result = self
465 .derive_dataset_with_lifecycle(session, request, &mut lifecycle)
466 .await;
467 if let Some(error) = lifecycle.rejection.take() {
468 return Err(error);
469 }
470 let Some(mut record) = lifecycle
471 .repository
472 .get_operation(lifecycle.id)
473 .map_err(|_| unavailable())?
474 else {
475 return result
476 .and_then(|_| Err(persistence_error()))
477 .map_err(crate::projections::dataset_materialization_protocol_error);
478 };
479 if let Err(error) = result {
480 if matches!(error, AppError::DatasetAdmissionOutcomeUnknown(_))
482 && record.detail.summary.status != OperationStatus::Completed
483 {
484 return Err(crate::projections::dataset_materialization_protocol_error(
485 error,
486 ));
487 }
488 let cancelled =
489 lifecycle.cancellation_observed && lifecycle.cleanup_succeeded == Some(true);
490 let mut failure = crate::projections::dataset_materialization_protocol_error(error);
491 if lifecycle.authority_lost {
492 failure = ProtocolError::new(
493 ProtocolErrorCode::OperationFailed,
494 "Operation interrupted after Session authority ended",
495 )
496 .with_target("interrupted");
497 } else if lifecycle.cleanup_succeeded == Some(false) {
498 failure = ProtocolError::new(
499 ProtocolErrorCode::OperationFailed,
500 "Dataset cleanup requires reconciliation",
501 );
502 failure.target = Some("cleanup_failed".into());
503 }
504 failure.adapter_observations = session.observations().snapshot(None);
506 record = update_operation(&lifecycle.repository, lifecycle.id, |record| {
507 if record.detail.summary.status.is_terminal() {
508 return Ok(false);
509 }
510 terminalize(
511 record,
512 if cancelled {
513 OperationStatus::Cancelled
514 } else {
515 OperationStatus::Failed
516 },
517 if cancelled {
518 None
519 } else {
520 Some(failure.clone())
521 },
522 );
523 Ok(true)
524 })
525 .map_err(crate::projections::dataset_materialization_protocol_error)?;
526 }
527 Ok(StartResult::Operation {
528 operation: Box::new(record.detail.summary),
529 })
530 }
531 .await;
532 use ahri_tre_observability::{FailureCategory, Outcome};
533 let (outcome, category) = match &result {
534 Ok(StartResult::Operation { operation }) => match operation.status {
535 OperationStatus::Completed => (Outcome::Success, None),
536 OperationStatus::Cancelled => (Outcome::Cancelled, None),
537 OperationStatus::Failed => (
538 Outcome::Unavailable,
539 Some(
540 match operation
541 .final_error
542 .as_ref()
543 .and_then(|error| error.target.as_deref())
544 {
545 Some("cleanup_failed") => FailureCategory::Cleanup,
546 Some("interrupted") => FailureCategory::Authorization,
547 _ => FailureCategory::Dependency,
548 },
549 ),
550 ),
551 _ => (Outcome::Unavailable, Some(FailureCategory::CommitUnknown)),
552 },
553 Err(error) => match error.code {
554 ProtocolErrorCode::Internal => {
555 (Outcome::InternalFailure, Some(FailureCategory::Internal))
556 }
557 ProtocolErrorCode::OperationFailed => {
558 (Outcome::Unavailable, Some(FailureCategory::Dependency))
559 }
560 ProtocolErrorCode::Forbidden => {
561 (Outcome::Rejected, Some(FailureCategory::Authorization))
562 }
563 _ => (Outcome::Rejected, None),
564 },
565 Ok(StartResult::Completed { .. }) => (Outcome::Success, None),
566 };
567 span.finish(outcome, category);
568 result
569 }
570
571 async fn authorized_operation(
572 &self,
573 inspection: &OperationControl,
574 owner: &str,
575 id: &OperationRef,
576 ) -> Result<StoredOperation, ProtocolError> {
577 if id.kind != ahri_tre_protocol::refs::ObjectKind::Operation {
578 return Err(ProtocolError::new(
579 ProtocolErrorCode::ValidationFailed,
580 "Expected an Operation reference",
581 ));
582 }
583 if id.datastore_id.as_uuid() != inspection.datastore_id {
584 return Err(ProtocolError::new(
585 ProtocolErrorCode::Conflict,
586 "Operation reference belongs to another Datastore",
587 ));
588 }
589 let id = id.id.as_uuid();
590 let mut record = inspection
591 .repository
592 .get_operation(id)
593 .map_err(|_| unavailable())?
594 .ok_or_else(not_found)?;
595 let now = inspection
596 .repository
597 .operation_now()
598 .map_err(|_| unavailable())?;
599 if operation_expired(&record.detail.summary, now).map_err(|_| unavailable())? {
600 return Err(not_found());
601 }
602 self.authorize_operation_scope(inspection, owner, &record.scope)
603 .await?;
604 let dataset_id = record.result.as_ref().map(|result| match result {
605 OperationResult::DatasetFromDatafile { data, .. } => data.dataset.dataset_id,
606 });
607 self.project_result_availability(&mut record.detail.summary, dataset_id, now)
608 .await?;
609 if record.detail.summary.status.is_terminal()
610 && deadline_reached(
611 record.detail.summary.retention.events_expires_at.as_deref(),
612 now,
613 )
614 .map_err(|_| unavailable())?
615 {
616 record.detail.events = Page {
617 items: vec![],
618 next_cursor: None,
619 };
620 }
621 Ok(record)
622 }
623
624 async fn authorize_operation_scope(
625 &self,
626 inspection: &OperationControl,
627 owner: &str,
628 scope: &OperationResourceScope,
629 ) -> Result<(), ProtocolError> {
630 if scope.owner != owner
631 || scope.datastore_id != inspection.datastore_id
632 || scope.deployment_id != inspection.deployment_id
633 {
634 return Err(not_found());
635 }
636 if self
638 .studies
639 .get_study_by_id(scope.study_id)
640 .await
641 .map_err(|_| unavailable())?
642 .is_none()
643 || self
644 .assets
645 .get_datafile(scope.source_version_id)
646 .await
647 .map_err(|_| unavailable())?
648 .is_none()
649 {
650 return Err(not_found());
651 }
652 if !self
653 .study_domains
654 .list_domains_for_study(scope.study_id)
655 .await
656 .map_err(|_| unavailable())?
657 .iter()
658 .any(|domain| domain.domain_id == scope.metadata_domain_id)
659 {
660 return Err(not_found());
661 }
662 Ok(())
663 }
664
665 pub async fn get_operation(
666 &self,
667 inspection: &OperationControl,
668 owner: &str,
669 operation: &OperationRef,
670 events: Option<&ahri_tre_protocol::pagination::PageRequest>,
671 ) -> Result<OperationDetail, ProtocolError> {
672 let mut detail = self
673 .authorized_operation(inspection, owner, operation)
674 .await?
675 .detail;
676 history::page_events(inspection, owner, &mut detail, events)?;
677 Ok(detail)
678 }
679
680 async fn cancel_operation(
681 &self,
682 control: &OperationControl,
683 owner: &str,
684 operation: &OperationRef,
685 ) -> Result<OperationSummary, ProtocolError> {
686 loop {
687 let mut record = self.authorized_operation(control, owner, operation).await?;
689 let summary = &record.detail.summary;
690 if matches!(&record.detail.progress, Some(OperationProgress::Stage { name, .. }) if name == "committing")
691 {
692 return Err(cancellation_conflict());
693 }
694 let policy = if summary.operation_kind.is_supported() {
695 OperationCancellationPolicy::Cooperative
696 } else {
697 OperationCancellationPolicy::Unsupported
698 };
699 let next = policy
700 .request(summary.status, summary.cancellable)
701 .map_err(|e| *e)?;
702 if summary.status == next {
703 return Ok(summary.clone());
704 }
705 let expected = summary.status;
706 let sequence = record.detail.events.items.len() as u64;
707 record.detail.summary.status = next;
708 record.detail.summary.cancellable = false;
709 let mut requested = event(
710 sequence + 1,
711 Utc::now().to_rfc3339(),
712 "cancellation requested",
713 );
714 requested.kind = OperationEventKind::CancellationRequested;
715 record.detail.events.items.push(requested);
716 match control
717 .repository
718 .transition_operation(expected, sequence, &record)
719 {
720 Ok(()) => return Ok(record.detail.summary),
721 Err(CoreError::Conflict(_)) => continue,
722 Err(_) => return Err(unavailable()),
723 }
724 }
725 }
726
727 pub async fn get_operation_result(
728 &self,
729 inspection: &OperationControl,
730 owner: &str,
731 operation: &OperationRef,
732 ) -> Result<OperationResult, ProtocolError> {
733 let record = self
734 .authorized_operation(inspection, owner, operation)
735 .await?;
736 match record.detail.summary.status {
737 OperationStatus::Pending
738 | OperationStatus::Running
739 | OperationStatus::CancelRequested => Err(ProtocolError::new(
740 ProtocolErrorCode::OperationNotReady,
741 "Operation result is not ready",
742 )),
743 OperationStatus::Failed => {
744 let mut error = record.detail.summary.final_error.unwrap_or_else(|| {
745 ProtocolError::new(
746 ProtocolErrorCode::OperationFailed,
747 "Dataset materialization failed",
748 )
749 });
750 error.code = ProtocolErrorCode::OperationFailed;
751 Err(error)
752 }
753 OperationStatus::Cancelled => Err(ProtocolError::new(
754 ProtocolErrorCode::Conflict,
755 "Operation was cancelled",
756 )),
757 OperationStatus::Completed => {
758 let availability = record.detail.summary.result.as_ref();
759 if let Some(reason) = unavailable_result_kind(availability) {
760 return Err(result_unavailable(reason));
761 }
762 match record.result {
763 Some(OperationResult::DatasetFromDatafile { mut data, .. }) => {
764 for variable in &mut data.variables {
767 variable.vocabulary_items.clear();
768 variable.vocabulary_items_withheld = variable.vocabulary.is_some();
769 }
770 Ok(OperationResult::DatasetFromDatafile {
771 data,
772 retention: record.detail.summary.retention,
773 availability: availability
774 .cloned()
775 .ok_or_else(|| result_unavailable("unavailable"))?,
776 })
777 }
778 None => Err(result_unavailable("unavailable")),
779 }
780 }
781 }
782 }
783}
784impl DatafileLifecycle for TrackedDatafileLifecycle {
785 fn observation(&self) -> Option<ahri_tre_observability::CorrelationContext> {
786 Some(self.context.observation)
787 }
788
789 fn accepted(&mut self) {
790 self.submission.key.take();
792 if let (Some(summary), Some(notify)) = (self.acceptance.take(), self.notify.take()) {
793 notify(summary);
794 }
795 }
796 fn cleanup_finished(&mut self, succeeded: bool) {
797 self.cleanup_succeeded = Some(succeeded);
798 }
799 fn resolved_source(&mut self, study: &ahri_tre_types::StudyRecord, source: &DataFileMetadata) {
800 self.source = Some((study.clone(), source.clone()));
801 }
802 fn record_acceptance(
803 &mut self,
804 repository: &PgMetadataRepository<'_>,
805 attempt: &ahri_tre_pgmeta::RecoverableDatasetAttempt,
806 resolved: &ResolvedDatafileMaterialization,
807 ) -> Result<(), CoreError> {
808 if (self.context.authority)() != OperationSessionAuthority::Active {
811 return Err(CoreError::Conflict("Session authority has ended".into()));
812 }
813 let fail = || CoreError::Validation("Operation context is unavailable".into());
814 let status = &self.context.session;
815 let scope = OperationResourceScope {
816 deployment_id: status
817 .session
818 .as_ref()
819 .ok_or_else(fail)?
820 .deployment_id
821 .as_uuid(),
822 datastore_id: self.datastore_id,
823 owner: self.context.owner.clone(),
824 session_id: status
825 .session
826 .as_ref()
827 .ok_or_else(fail)?
828 .session_id
829 .as_uuid(),
830 study_id: resolved.request.study_id,
831 metadata_domain_id: resolved.request.metadata_domain_id.ok_or_else(fail)?,
832 source_version_id: resolved.source_version.version_id,
833 dataset_name: resolved.request.dataset_name.as_str().into(),
834 };
835 for condition in &self.submission.preconditions {
836 if !repository.operation_precondition_satisfied(&scope, condition)? {
837 let mut error = ProtocolError::new(
838 ProtocolErrorCode::PreconditionFailed,
839 "Materialization precondition failed",
840 );
841 error.details = Some(
842 ahri_tre_protocol::public_error::ProtocolErrorDetails::PreconditionFailed {
843 precondition: serde_json::to_string(condition).map_err(|_| fail())?,
844 },
845 );
846 self.rejection = Some(error);
847 return Err(CoreError::Validation(
848 "Materialization precondition failed".into(),
849 ));
850 }
851 }
852 let now = Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Micros, true);
853 let summary = OperationSummary {
854 operation: OperationRef {
855 datastore_id: ahri_tre_protocol::PublicUuid::from_uuid(self.datastore_id),
856 kind: ahri_tre_protocol::refs::ObjectKind::Operation,
857 id: ahri_tre_protocol::PublicUuid::from_uuid(self.id),
858 },
859 operation_kind: OperationKind::new(
860 ahri_tre_protocol::request::kind::INGEST_DATASET_FROM_DATAFILE,
861 ),
862 operation_scope: OperationScope::Datastore,
863 status: OperationStatus::Pending,
864 cancellable: true,
865 started_by_request_id: self.context.request_id,
866 timestamps: OperationTimestamps {
867 created_at: now.clone(),
868 started_at: None,
869 finished_at: None,
870 },
871 retention: OperationRetention::default(),
872 result: None,
873 final_error: None,
874 warnings: vec![],
875 };
876 let record = StoredOperation {
877 scope: scope.clone(),
878 detail: OperationDetail {
879 summary,
880 progress: Some(Stage::Queued.progress()),
881 events: Page {
882 items: vec![Stage::Queued.event(1, now, OperationEventKind::StatusChanged)],
883 next_cursor: None,
884 },
885 },
886 result: None,
887 };
888 repository.create_operation(attempt.attempt_id, &record)?;
889 if let Some(key) = &self.submission.key {
890 key.bind(repository, self.id)?;
891 }
892 self.acceptance = Some(record.detail.summary);
893 Ok(())
894 }
895 fn checkpoint(
896 &mut self,
897 stage: DatafileStage,
898 _: &ResolvedDatafileMaterialization,
899 ) -> Result<(), AppError> {
900 let record = self
901 .repository
902 .get_operation(self.id)
903 .map_err(|_| persistence_error())?
904 .ok_or_else(persistence_error)?;
905 self.check_session_authority()?;
906 if record.detail.summary.status == OperationStatus::CancelRequested {
907 self.cancellation_observed = true;
908 return Err(AppError::Conflict(
909 "Operation cancellation requested".into(),
910 ));
911 }
912 if matches!(
913 stage,
914 DatafileStage::BeforeLoad | DatafileStage::AfterPreparation
915 ) {
916 update_operation(&self.repository, self.id, |record| {
917 if record.detail.summary.status == OperationStatus::CancelRequested {
919 self.cancellation_observed = true;
920 return Err(AppError::Conflict(
921 "Operation cancellation requested".into(),
922 ));
923 }
924 stage_event(
925 record,
926 if stage == DatafileStage::BeforeLoad {
927 Stage::Materializing
928 } else {
929 Stage::Preparing
930 },
931 OperationEventKind::ProgressUpdated,
932 );
933 Ok(true)
934 })?;
935 }
936 if stage == DatafileStage::BeforeMetadataCommit {
937 update_operation(&self.repository, self.id, |record| {
941 if record.detail.summary.status.is_terminal() {
942 return Err(persistence_error());
943 }
944 record.detail.summary.cancellable = false;
945 stage_event(
946 record,
947 Stage::Committing,
948 OperationEventKind::ProgressUpdated,
949 );
950 Ok(true)
951 })?;
952 }
953 Ok(())
954 }
955 fn claim_execution(&mut self) -> Result<(), AppError> {
956 let repository = self.repository.clone();
957 update_operation(&repository, self.id, |record| {
958 self.check_session_authority()?;
959 if record.detail.summary.status == OperationStatus::CancelRequested {
960 self.cancellation_observed = true;
961 return Err(AppError::Conflict(
962 "Operation cancellation requested".into(),
963 ));
964 }
965 if record.detail.summary.status != OperationStatus::Pending {
966 return Err(persistence_error());
967 }
968 let now = Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Micros, true);
969 record.detail.summary.status = OperationStatus::Running;
970 record.detail.summary.timestamps.started_at = Some(now.clone());
971 record.detail.progress = Some(Stage::Preparing.progress());
972 record.detail.events.items.push(Stage::Preparing.event(
973 2,
974 now,
975 OperationEventKind::StatusChanged,
976 ));
977 Ok(true)
978 })
979 .map(|_| ())
980 }
981 fn cleanup_started(&mut self) {
982 let _ = update_operation(&self.repository, self.id, |record| {
984 if record.detail.summary.status.is_terminal() {
985 return Ok(false);
986 }
987 record.detail.summary.cancellable = false;
988 stage_event(
989 record,
990 Stage::CleaningUp,
991 OperationEventKind::CleanupStarted,
992 );
993 Ok(true)
994 });
995 }
996 fn record_completion(
997 &mut self,
998 repository: &PgMetadataRepository<'_>,
999 result: &DatasetMaterialization,
1000 ) -> Result<(), AppError> {
1001 let mut record = repository
1002 .get_operation(self.id)
1003 .map_err(|_| persistence_error())?
1004 .ok_or_else(persistence_error)?;
1005 let (study, source) = self.source.clone().ok_or_else(persistence_error)?;
1006 let mut response = crate::projections::dataset_materialization_response(
1007 self.context.session.clone(),
1008 study.name.as_str().into(),
1009 ManagedDataFileDatasetMaterialization {
1010 source,
1011 dataset: result.clone(),
1012 },
1013 );
1014 response.lineage.transformation.description =
1015 "materialize Dataset from managed Datafile".into();
1016 response.metadata_warnings = response
1017 .metadata_warnings
1018 .iter()
1019 .take(100)
1020 .map(|_| "Dataset metadata inference notice".to_string())
1021 .collect();
1022 let expected = record.detail.summary.status;
1023 let sequence = record.detail.events.items.len() as u64;
1024 terminalize(&mut record, OperationStatus::Completed, None);
1025 record.detail.summary.result = Some(OperationResultRef::Durable {
1026 ref_id: self.id.to_string(),
1027 });
1028 record.result = Some(OperationResult::DatasetFromDatafile {
1029 data: Box::new(response),
1030 retention: record.detail.summary.retention.clone(),
1031 availability: record
1032 .detail
1033 .summary
1034 .result
1035 .clone()
1036 .ok_or_else(persistence_error)?,
1037 });
1038 repository
1039 .transition_operation(expected, sequence, &record)
1040 .map_err(|_| persistence_error())
1041 }
1042}
1043fn update_operation(
1044 repository: &PgMetadataRepository<'_>,
1045 id: Uuid,
1046 mut update: impl FnMut(&mut StoredOperation) -> Result<bool, AppError>,
1047) -> Result<StoredOperation, AppError> {
1048 loop {
1049 let mut record = repository
1050 .get_operation(id)
1051 .map_err(|_| persistence_error())?
1052 .ok_or_else(persistence_error)?;
1053 let expected = record.detail.summary.status;
1054 let sequence = record.detail.events.items.len() as u64;
1055 if !update(&mut record)? {
1056 return Ok(record);
1057 }
1058 match repository.transition_operation(expected, sequence, &record) {
1059 Ok(()) => return Ok(record),
1060 Err(CoreError::Conflict(_)) => continue,
1061 Err(_) => return Err(persistence_error()),
1062 }
1063 }
1064}
1065#[derive(Clone, Copy)]
1066enum Stage {
1067 Queued,
1068 Preparing,
1069 Materializing,
1070 Committing,
1071 CleaningUp,
1072}
1073impl Stage {
1074 fn name(self) -> &'static str {
1075 match self {
1076 Self::Queued => "queued",
1077 Self::Preparing => "preparing",
1078 Self::Materializing => "materializing",
1079 Self::Committing => "committing",
1080 Self::CleaningUp => "cleaning up",
1081 }
1082 }
1083 fn progress(self) -> OperationProgress {
1084 OperationProgress::Stage {
1085 name: self.name().into(),
1086 message: None,
1087 }
1088 }
1089 fn event(self, sequence: u64, at: String, kind: OperationEventKind) -> OperationEvent {
1090 let mut entry = event(sequence, at, self.name());
1091 entry.kind = kind;
1092 entry.progress = Some(self.progress());
1093 entry
1094 }
1095}
1096fn stage_event(record: &mut StoredOperation, stage: Stage, kind: OperationEventKind) {
1097 record.detail.progress = Some(stage.progress());
1098 record.detail.events.items.push(stage.event(
1099 record.detail.events.items.len() as u64 + 1,
1100 Utc::now().to_rfc3339(),
1101 kind,
1102 ));
1103}
1104fn cancellation_conflict() -> ProtocolError {
1105 ProtocolError::new(
1106 ProtocolErrorCode::Conflict,
1107 "Operation cannot be cancelled in its current state",
1108 )
1109}
1110fn terminalize(
1111 record: &mut StoredOperation,
1112 status: OperationStatus,
1113 error: Option<ProtocolError>,
1114) {
1115 let now = Utc::now();
1116 let expiry = (now + Duration::days(30)).to_rfc3339();
1117 record.detail.summary.status = status;
1118 record.detail.summary.cancellable = false;
1119 record.detail.progress = None;
1120 record.detail.summary.final_error = error;
1121 record.detail.summary.timestamps.finished_at = Some(now.to_rfc3339());
1122 record.detail.summary.retention = OperationRetention {
1123 operation_expires_at: Some(expiry.clone()),
1124 events_expires_at: Some(expiry.clone()),
1125 result_expires_at: Some(expiry),
1126 };
1127 let sequence = record.detail.events.items.len() as u64 + 1;
1128 record.detail.events.items.push(event(
1129 sequence,
1130 now.to_rfc3339(),
1131 if status == OperationStatus::Completed {
1132 "completed"
1133 } else if status == OperationStatus::Cancelled {
1134 "cancelled"
1135 } else {
1136 "failed"
1137 },
1138 ));
1139}
1140fn event(sequence: u64, at: String, message: &str) -> OperationEvent {
1141 OperationEvent {
1142 event_sequence: sequence,
1143 event_at: at,
1144 kind: OperationEventKind::StatusChanged,
1145 message: Some(message.into()),
1146 target: None,
1147 progress: None,
1148 warning: None,
1149 }
1150}
1151fn persistence_error() -> AppError {
1152 AppError::Infrastructure("Operation persistence is unavailable".into())
1153}
1154fn unavailable() -> ProtocolError {
1155 ProtocolError::new(
1156 ProtocolErrorCode::OperationFailed,
1157 "Operation persistence is unavailable",
1158 )
1159}
1160fn not_found() -> ProtocolError {
1161 ProtocolError::new(ProtocolErrorCode::NotFound, "Operation not found")
1162}
1163fn result_unavailable(availability: &str) -> ProtocolError {
1164 let mut error = ProtocolError::new(
1165 ProtocolErrorCode::OperationResultUnavailable,
1166 format!("Operation result is {availability}"),
1167 );
1168 error.details = Some(
1169 ahri_tre_protocol::public_error::ProtocolErrorDetails::OperationResultUnavailable {
1170 availability: availability.into(),
1171 },
1172 );
1173 error
1174}
1175
1176pub(crate) fn interrupt_operation(record: &mut StoredOperation) {
1178 terminalize(
1179 record,
1180 OperationStatus::Failed,
1181 Some(
1182 ProtocolError::new(
1183 ProtocolErrorCode::OperationFailed,
1184 "Operation interrupted after Session authority ended",
1185 )
1186 .with_target("interrupted"),
1187 ),
1188 );
1189}
1190
1191fn deadline_reached(
1192 deadline: Option<&str>,
1193 now: chrono::DateTime<Utc>,
1194) -> Result<bool, chrono::ParseError> {
1195 deadline
1196 .map(|value| chrono::DateTime::parse_from_rfc3339(value).map(|value| now >= value))
1197 .transpose()
1198 .map(|expired| expired.unwrap_or(false))
1199}
1200fn operation_expired(
1201 summary: &OperationSummary,
1202 now: chrono::DateTime<Utc>,
1203) -> Result<bool, chrono::ParseError> {
1204 Ok(summary.status.is_terminal()
1205 && deadline_reached(summary.retention.operation_expires_at.as_deref(), now)?)
1206}
1207
1208fn unavailable_result_kind(reference: Option<&OperationResultRef>) -> Option<&'static str> {
1209 match reference {
1210 Some(OperationResultRef::Durable { .. } | OperationResultRef::Temporary { .. }) => None,
1211 Some(OperationResultRef::Expired) => Some("expired"),
1212 Some(OperationResultRef::Removed) => Some("removed"),
1213 Some(OperationResultRef::Unavailable) | None => Some("unavailable"),
1214 }
1215}
1216impl AppService {
1217 async fn project_result_availability(
1218 &self,
1219 summary: &mut OperationSummary,
1220 dataset_id: Option<ahri_tre_types::VersionId>,
1221 now: chrono::DateTime<Utc>,
1222 ) -> Result<(), ProtocolError> {
1223 if summary.status != OperationStatus::Completed {
1224 return Ok(());
1225 }
1226 if deadline_reached(summary.retention.result_expires_at.as_deref(), now)
1227 .map_err(|_| unavailable())?
1228 {
1229 summary.result = Some(OperationResultRef::Expired);
1230 } else if unavailable_result_kind(summary.result.as_ref()).is_some() {
1231 summary
1232 .result
1233 .get_or_insert(OperationResultRef::Unavailable);
1234 } else if let Some(id) = dataset_id {
1235 if self
1236 .assets
1237 .get_dataset(id)
1238 .await
1239 .map_err(|_| unavailable())?
1240 .is_none()
1241 || self
1242 .assets
1243 .get_dataset_version_withdrawal(id)
1244 .await
1245 .map_err(|_| unavailable())?
1246 .is_some()
1247 {
1248 summary.result = Some(OperationResultRef::Unavailable);
1249 }
1250 } else {
1251 summary.result = Some(OperationResultRef::Unavailable);
1252 }
1253 Ok(())
1254 }
1255}