ahri_tre_app/service/operations/
history.rs1use super::*;
3use ahri_tre_core::{OperationHistoryQuery, OperationPosition};
4use ahri_tre_protocol::{pagination::PageRequest, request::OperationListRequest};
5use chrono::{DateTime, Timelike};
6use serde::{Deserialize, Serialize};
7use std::sync::OnceLock;
8
9#[derive(Serialize, Deserialize)]
10#[serde(tag = "kind")]
11enum Position {
12 History {
13 upper: DateTime<Utc>,
14 after: OperationPosition,
15 },
16 Events {
17 after: u64,
18 },
19}
20#[derive(Serialize, Deserialize)]
21struct Cursor {
22 position: Position,
23 signature: Vec<u8>,
24}
25
26impl OperationControl {
27 pub async fn list(
28 &self,
29 owner: &str,
30 request: OperationListRequest,
31 ) -> Result<Page<OperationSummary>, ProtocolError> {
32 request.page.validate().map_err(|_| invalid_page())?;
33 if request
34 .operation_kind
35 .as_ref()
36 .is_some_and(|kind| !kind.is_supported())
37 {
38 return Err(invalid_page().into());
39 }
40 let mut statuses = request.statuses;
41 statuses.sort_by_key(|status| *status as u8);
42 statuses.dedup();
43 let lower = timestamp(request.created_after.as_deref())?;
44 let before = timestamp(request.created_before.as_deref())?;
45 if lower
46 .zip(before)
47 .is_some_and(|(lower, before)| lower >= before)
48 {
49 return Err(invalid_page().into());
50 }
51 let context = serde_json::to_vec(&(
54 self.deployment_id,
55 self.datastore_id,
56 owner,
57 request.scope,
58 (request.scope == OperationScope::Session).then_some(&request.session),
59 &statuses,
60 &request.operation_kind,
61 lower,
62 before,
63 ))
64 .map_err(|_| unavailable())?;
65 let now = self
66 .repository
67 .operation_now()
68 .map_err(|_| super::unavailable())?;
69 let (upper, after) = match request.page.cursor.as_deref() {
70 Some(raw) => match decode(raw, &context)? {
71 Position::History { upper, after } if after.created_at <= upper => {
72 (upper, Some(after))
73 }
74 _ => return Err(invalid_page().into()),
75 },
76 None => {
77 (
80 now.with_nanosecond(now.nanosecond() / 1000 * 1000)
81 .ok_or_else(unavailable)?,
82 None,
83 )
84 }
85 };
86 let mut query = OperationHistoryQuery {
87 created_after: lower,
88 created_before: before,
89 upper_bound: upper,
90 after,
91 limit: 100,
92 };
93 let service = self.service();
94 let mut items = Vec::new();
95 let mut last_visible = None;
96 loop {
97 let batch = self
98 .repository
99 .operation_history(&query)
100 .map_err(|_| unavailable())?;
101 let exhausted = batch.len() < usize::from(query.limit);
102 for mut record in batch {
103 let position = OperationPosition {
104 created_at: timestamp(Some(&record.summary.timestamps.created_at))?
105 .ok_or_else(unavailable)?,
106 id: record.summary.operation.id.as_uuid(),
107 };
108 query.after = Some(position.clone());
109 if operation_expired(&record.summary, now).map_err(|_| super::unavailable())?
112 || record.summary.operation_scope != request.scope
113 || (!statuses.is_empty() && !statuses.contains(&record.summary.status))
114 || request
115 .operation_kind
116 .as_ref()
117 .is_some_and(|kind| kind != &record.summary.operation_kind)
118 {
119 continue;
120 }
121 match service
122 .authorize_operation_scope(self, owner, &record.scope)
123 .await
124 {
125 Ok(()) => {}
126 Err(error) if error.code == ProtocolErrorCode::NotFound => continue,
127 Err(error) => return Err(error),
128 }
129 service
130 .project_result_availability(&mut record.summary, record.result_dataset_id, now)
131 .await?;
132 if items.len() == usize::from(request.page.limit) {
135 let after = last_visible.ok_or_else(unavailable)?;
136 return Ok(Page {
137 items,
138 next_cursor: Some(encode(Position::History { upper, after }, &context)?),
139 });
140 }
141 last_visible = Some(position);
142 items.push(record.summary);
143 }
144 if exhausted {
145 return Ok(Page {
146 items,
147 next_cursor: None,
148 });
149 }
150 }
151 }
152}
153
154pub(super) fn page_events(
155 control: &OperationControl,
156 owner: &str,
157 detail: &mut OperationDetail,
158 page: Option<&PageRequest>,
159) -> Result<(), PageError> {
160 let page = page.cloned().unwrap_or_default();
161 page.validate().map_err(|_| invalid_page())?;
162 let context = serde_json::to_vec(&(
163 control.deployment_id,
164 control.datastore_id,
165 owner,
166 &detail.summary.operation,
167 ))
168 .map_err(|_| unavailable())?;
169 let after = match page.cursor.as_deref() {
170 Some(raw) => match decode(raw, &context)? {
171 Position::Events { after } => after,
172 _ => return Err(invalid_page()),
173 },
174 None => 0,
175 };
176 detail
177 .events
178 .items
179 .retain(|event| event.event_sequence > after);
180 detail
181 .events
182 .items
183 .sort_by_key(|event| event.event_sequence);
184 let has_more = detail.events.items.len() > usize::from(page.limit);
185 detail.events.items.truncate(usize::from(page.limit));
186 detail.events.next_cursor = if has_more {
187 let after = detail
188 .events
189 .items
190 .last()
191 .ok_or_else(unavailable)?
192 .event_sequence;
193 Some(encode(Position::Events { after }, &context)?)
194 } else {
195 None
196 };
197 Ok(())
198}
199
200fn timestamp(value: Option<&str>) -> Result<Option<DateTime<Utc>>, PageError> {
201 value
202 .map(|value| {
203 DateTime::parse_from_rfc3339(value)
204 .map(|value| value.with_timezone(&Utc))
205 .map_err(|_| invalid_page())
206 })
207 .transpose()
208}
209fn signature(position: &Position, context: &[u8]) -> Result<Vec<u8>, PageError> {
210 static KEY: OnceLock<[u8; 16]> = OnceLock::new();
211 let key = KEY.get_or_init(|| *Uuid::new_v4().as_bytes());
212 let bytes = serde_json::to_vec(&("operation-pagination-v1", context, position))
213 .map_err(|_| unavailable())?;
214 Ok(crate::cursor_integrity::hmac_sha256(key, &bytes).to_vec())
215}
216fn encode(position: Position, context: &[u8]) -> Result<String, PageError> {
217 let cursor = Cursor {
218 signature: signature(&position, context)?,
219 position,
220 };
221 Ok(serde_json::to_vec(&cursor)
222 .map_err(|_| unavailable())?
223 .iter()
224 .map(|byte| format!("{byte:02x}"))
225 .collect())
226}
227fn decode(raw: &str, context: &[u8]) -> Result<Position, PageError> {
228 if raw.len() > 2048 || !raw.len().is_multiple_of(2) || !raw.is_ascii() {
229 return Err(invalid_page());
230 }
231 let bytes = (0..raw.len())
232 .step_by(2)
233 .map(|i| u8::from_str_radix(&raw[i..i + 2], 16).map_err(|_| invalid_page()))
234 .collect::<Result<Vec<_>, _>>()?;
235 let cursor: Cursor = serde_json::from_slice(&bytes).map_err(|_| invalid_page())?;
236 let expected = signature(&cursor.position, context)?;
237 if cursor.signature.len() != expected.len()
238 || cursor
239 .signature
240 .iter()
241 .zip(&expected)
242 .fold(0u8, |diff, (a, b)| diff | (a ^ b))
243 != 0
244 {
245 return Err(invalid_page());
246 }
247 Ok(cursor.position)
248}
249#[derive(Debug)]
250pub(super) enum PageError {
251 Invalid,
252 Unavailable,
253}
254impl From<PageError> for ProtocolError {
255 fn from(error: PageError) -> Self {
256 match error {
257 PageError::Invalid => ProtocolError::new(
258 ProtocolErrorCode::ValidationFailed,
259 "Invalid operation page, cursor, or filters",
260 ),
261 PageError::Unavailable => super::unavailable(),
262 }
263 }
264}
265fn invalid_page() -> PageError {
266 PageError::Invalid
267}
268fn unavailable() -> PageError {
269 PageError::Unavailable
270}