Skip to main content

ahri_tre_app/service/operations/
history.rs

1//! Authorized, bounded history traversal. Cursors contain no owner or resource scope.
2use 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        // Datastore history survives explicitly reopening a Session. Session
52        // scope binds its selected Session, though this producer never uses it.
53        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                // PostgreSQL timestamps have microsecond precision. Truncate
78                // rather than round so the traversal never includes future work.
79                (
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                // Filter bounded actor-owned candidates here: JSON predicates in
110                // the RLS query can cause repeated sorting of the remaining history.
111                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                // Look ahead only through authorized rows. A cursor cannot reveal
133                // whether any invisible operations remain after the visible page.
134                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}