Skip to main content

ahri_tre_app/service/
operations.rs

1//! Application-owned operation lifecycle and authorization-aware projection.
2mod 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/// Cancellation policy is separate from the current execution window.
17#[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/// Live Session authority is checked at each governed stage; cleanup needs no user capability.
45#[derive(Clone, Copy, PartialEq, Eq)]
46pub enum OperationSessionAuthority {
47    Active,
48    Closed,
49    Lost,
50}
51
52/// Newly authenticated context. Only safe correlation fields are persisted.
53pub 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
61/// Independent lifecycle capability exposing only authorized inspection and cancellation.
62pub 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
124/// Newly authorized admission or a retained result, resolved before worker and
125/// output ownership acquisition. The private key guard lives only to acceptance.
126pub 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        // Resolve on the independent operation-control connection, so retries
145        // and polling do not wait for a producer's live datastore lock.
146        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        // Resolve identity and every supplied scope before checking agreement.
172        // Keep an omitted/latest selection moving in the idempotency intent;
173        // the producer pins it once before reservation and source preparation.
174        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        // A pinned reference stays immutable across the worker handoff. Its
191        // semantic version is used only to compare equivalent retry intents.
192        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            // Value serialization sorts object keys and applies DTO defaults;
256            // arrays (including preconditions) preserve order and duplicates.
257            let mut intent = request.clone();
258            intent.session = None;
259            intent.mutation.idempotency_key = None;
260            intent.mutation.preconditions = preconditions.clone();
261            // These source-selector spellings resolve identically on admission.
262            // Normalize the selector, never resolve latest again for a retry.
263            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                // Recheck pinned resources before disclosing either replay or conflict.
308                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        // A retained retry uses its original pinned resources. Only a new
331        // admission must establish agreement with a supplied moving latest.
332        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    /// Executes in the caller's worker, notifying only after durable acceptance.
416    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                // Unknown commit cannot be terminalized or compensated until resolved.
481                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                // Existing authorized failure readback retains the acquisition times.
505                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        // Study reads enforce current grants through the authenticated metadata adapter.
637        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            // Reauthorize on every competing transition; IDs confer no authority.
688            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                        // Stored operation JSON is never a retained disclosure
765                        // capability. Historical inferred dictionaries stay withheld.
766                        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        // Release only after the acceptance transaction has committed.
791        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        // Source resolution can block before this transaction. A close or
809        // revocation in that interval must roll back acceptance itself.
810        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                // Progress must not overwrite a cancellation accepted at this checkpoint.
918                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            // The last cooperative checkpoint has passed. A cancellation accepted
938            // concurrently may lose to authoritative admission. Close the window
939            // in a short conditional update before starting the admission transaction.
940            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        // Failure to publish progress never interrupts compensation or releases ownership.
983        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
1176/// Projection policy supplied to the adapter's stopped-attempt recovery boundary.
1177pub(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}