Skip to main content

ahri_tre_protocol/
operation.rs

1use crate::{RequestId, pagination::Page, public_error::ProtocolError, warning::ProtocolWarning};
2use serde::{Deserialize, Serialize};
3
4pub type OperationRef = crate::refs::ObjectRef;
5
6#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
7#[serde(rename_all = "snake_case")]
8pub enum OperationScope {
9    Session,
10    Datastore,
11}
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
14#[serde(rename_all = "snake_case")]
15pub enum OperationStatus {
16    Pending,
17    Running,
18    CancelRequested,
19    Cancelled,
20    Completed,
21    Failed,
22}
23
24impl OperationStatus {
25    pub const fn is_terminal(self) -> bool {
26        matches!(self, Self::Cancelled | Self::Completed | Self::Failed)
27    }
28
29    pub const fn can_transition_to(self, next: Self) -> bool {
30        matches!(
31            (self, next),
32            (Self::Pending, Self::Running)
33                | (Self::Pending, Self::CancelRequested)
34                | (Self::Pending, Self::Completed)
35                | (Self::Pending, Self::Failed)
36                | (Self::Running, Self::CancelRequested)
37                | (Self::Running, Self::Completed)
38                | (Self::Running, Self::Failed)
39                | (Self::CancelRequested, Self::Cancelled)
40                | (Self::CancelRequested, Self::Completed)
41                | (Self::CancelRequested, Self::Failed)
42        )
43    }
44}
45
46#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
47#[serde(try_from = "String")]
48pub struct OperationKind(pub String);
49impl TryFrom<String> for OperationKind {
50    type Error = &'static str;
51    fn try_from(value: String) -> Result<Self, Self::Error> {
52        let kind = Self(value);
53        if kind.is_supported() {
54            Ok(kind)
55        } else {
56            Err("unsupported operation kind")
57        }
58    }
59}
60
61impl OperationKind {
62    pub fn new(value: impl Into<String>) -> Self {
63        Self(value.into())
64    }
65
66    pub fn is_supported(&self) -> bool {
67        self.0 == crate::request::kind::INGEST_DATASET_FROM_DATAFILE
68    }
69}
70
71#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
72pub struct OperationTimestamps {
73    pub created_at: String,
74    #[serde(default, skip_serializing_if = "Option::is_none")]
75    pub started_at: Option<String>,
76    #[serde(default, skip_serializing_if = "Option::is_none")]
77    pub finished_at: Option<String>,
78}
79
80#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
81pub struct OperationRetention {
82    #[serde(default, skip_serializing_if = "Option::is_none")]
83    pub operation_expires_at: Option<String>,
84    #[serde(default, skip_serializing_if = "Option::is_none")]
85    pub events_expires_at: Option<String>,
86    #[serde(default, skip_serializing_if = "Option::is_none")]
87    pub result_expires_at: Option<String>,
88}
89
90#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
91#[serde(tag = "kind", rename_all = "snake_case")]
92pub enum OperationResultRef {
93    Durable { ref_id: String },
94    Temporary { ref_id: String },
95    Expired,
96    Removed,
97    Unavailable,
98}
99
100#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101pub struct OperationSummary {
102    pub operation: OperationRef,
103    pub operation_kind: OperationKind,
104    pub operation_scope: OperationScope,
105    pub status: OperationStatus,
106    pub cancellable: bool,
107    pub started_by_request_id: RequestId,
108    pub timestamps: OperationTimestamps,
109    pub retention: OperationRetention,
110    #[serde(default, skip_serializing_if = "Option::is_none")]
111    pub result: Option<OperationResultRef>,
112    #[serde(default, skip_serializing_if = "Option::is_none")]
113    pub final_error: Option<ProtocolError>,
114    #[serde(default, skip_serializing_if = "Vec::is_empty")]
115    pub warnings: Vec<ProtocolWarning>,
116}
117
118#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
119pub struct OperationDetail {
120    pub summary: OperationSummary,
121    #[serde(default, skip_serializing_if = "Option::is_none")]
122    pub progress: Option<OperationProgress>,
123    pub events: Page<OperationEvent>,
124}
125
126#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
127#[serde(tag = "kind", rename_all = "snake_case")]
128pub enum OperationProgress {
129    Stage {
130        name: String,
131        message: Option<String>,
132    },
133    Count {
134        completed: u64,
135        total: Option<u64>,
136        unit: Option<String>,
137    },
138    Bytes {
139        completed: u64,
140        total: Option<u64>,
141    },
142    Fraction {
143        value: f64,
144    },
145}
146
147#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
148pub struct OperationEvent {
149    pub event_sequence: u64,
150    pub event_at: String,
151    pub kind: OperationEventKind,
152    #[serde(default, skip_serializing_if = "Option::is_none")]
153    pub message: Option<String>,
154    #[serde(default, skip_serializing_if = "Option::is_none")]
155    pub target: Option<String>,
156    #[serde(default, skip_serializing_if = "Option::is_none")]
157    pub progress: Option<OperationProgress>,
158    #[serde(default, skip_serializing_if = "Option::is_none")]
159    pub warning: Option<ProtocolWarning>,
160}
161
162#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
163#[serde(rename_all = "snake_case")]
164pub enum OperationEventKind {
165    StatusChanged,
166    ProgressUpdated,
167    Warning,
168    RetryScheduled,
169    CancellationRequested,
170    CleanupStarted,
171    CleanupCompleted,
172    ResultAvailable,
173    ResultUnavailable,
174    Message,
175}
176
177#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
178#[serde(tag = "kind", rename_all = "snake_case")]
179pub enum StartResult<T> {
180    Completed { data: T },
181    Operation { operation: Box<OperationSummary> },
182}
183
184/// Retained outcome of a supported producer. Session fields are correlation only.
185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
186#[serde(tag = "kind")]
187pub enum OperationResult {
188    #[serde(rename = "ingest.dataset.from_datafile")]
189    DatasetFromDatafile {
190        data: Box<crate::ingest::DatasetMaterializationResponse>,
191        retention: OperationRetention,
192        availability: OperationResultRef,
193    },
194}