Skip to main content

pidge_client/graph/
delta.rs

1//! Graph delta queries: "what changed since last time".
2//!
3//! Bootstrap walks `@odata.nextLink` pages to the terminal `@odata.deltaLink`
4//! (establishing the sync state); subsequent calls at the deltaLink return
5//! only created/updated messages plus `@removed` tombstones and a fresh
6//! deltaLink. An expired token (HTTP 410) surfaces as
7//! [`ClientError::DeltaExpired`] so callers can re-bootstrap.
8
9use chrono::{DateTime, Utc};
10use pidge_core::{Event, Message};
11use serde_json::Value;
12
13use crate::error::ClientError;
14
15/// A delta change carries only the properties that changed (plus id), so
16/// everything except `graph_id` is optional.
17#[derive(Debug, Clone, serde::Deserialize)]
18pub struct DeltaMessage {
19    #[serde(rename = "id")]
20    pub graph_id: String,
21    pub subject: Option<String>,
22    #[serde(rename = "isRead")]
23    pub is_read: Option<bool>,
24    #[serde(rename = "receivedDateTime")]
25    pub received_at: Option<DateTime<Utc>>,
26    #[serde(rename = "bodyPreview")]
27    pub preview: Option<String>,
28    #[serde(rename = "hasAttachments")]
29    pub has_attachments: Option<bool>,
30    #[serde(rename = "conversationId")]
31    pub conversation_id: Option<String>,
32    #[serde(rename = "from")]
33    pub from: Option<serde_json::Value>,
34}
35
36#[derive(Debug, Clone)]
37pub enum MailDeltaEvent {
38    Changed(Box<DeltaMessage>),
39    Deleted { graph_id: String },
40}
41
42#[derive(Debug, Clone)]
43pub enum CalendarDeltaEvent {
44    CreatedOrUpdated(Box<Event>),
45    Deleted { graph_id: String },
46}
47
48/// Follow a delta stream from `url` to its terminal deltaLink, collecting
49/// raw item values along the way. Returns (items, delta_link).
50async fn drain(
51    http: &reqwest::Client,
52    access_token: &str,
53    url: &str,
54    prefer: Option<&str>,
55) -> Result<(Vec<Value>, String), ClientError> {
56    let mut items: Vec<Value> = Vec::new();
57    let mut next = url.to_string();
58    loop {
59        let mut req = http.get(&next).bearer_auth(access_token);
60        if let Some(prefer) = prefer {
61            req = req.header("Prefer", prefer);
62        }
63        let resp = super::send_with_retry(req).await?;
64        let status = resp.status();
65        if status.as_u16() == 410 {
66            return Err(ClientError::DeltaExpired);
67        }
68        if !status.is_success() {
69            let text = resp.text().await.unwrap_or_default();
70            return Err(ClientError::Graph {
71                status: status.as_u16(),
72                message: text,
73            });
74        }
75        let body: Value = resp.json().await?;
76        if let Some(page_items) = body.get("value").and_then(Value::as_array) {
77            items.extend(page_items.iter().cloned());
78        }
79        if let Some(delta_link) = body.get("@odata.deltaLink").and_then(Value::as_str) {
80            return Ok((items, delta_link.to_string()));
81        }
82        match body.get("@odata.nextLink").and_then(Value::as_str) {
83            Some(link) => next = link.to_string(),
84            None => {
85                return Err(ClientError::Graph {
86                    status: 500,
87                    message: "delta response carried neither nextLink nor deltaLink".into(),
88                });
89            }
90        }
91    }
92}
93
94/// Establish mail delta state for a folder. Returns the deltaLink (and the
95/// current messages, which `--full` bootstraps replay as `created` events).
96pub async fn mail_delta_bootstrap(
97    http: &reqwest::Client,
98    base_url: &str,
99    access_token: &str,
100    account: &str,
101    folder: &str,
102) -> Result<(Vec<Message>, String), ClientError> {
103    let url = format!(
104        "{base_url}/me/mailFolders/{folder}/messages/delta?$select=id,subject,from,receivedDateTime,isRead,bodyPreview,hasAttachments,flag"
105    );
106    let (items, delta_link) = drain(http, access_token, &url, None).await?;
107    let messages = items
108        .into_iter()
109        .filter(|v| v.get("@removed").is_none())
110        .filter_map(|v| super::mail::message_from_delta_value(v, account))
111        .collect();
112    Ok((messages, delta_link))
113}
114
115/// Poll a mail deltaLink: changes since the link was minted + a fresh link.
116pub async fn mail_delta(
117    http: &reqwest::Client,
118    access_token: &str,
119    account: &str,
120    delta_link: &str,
121) -> Result<(Vec<MailDeltaEvent>, String), ClientError> {
122    let (items, next_link) = drain(http, access_token, delta_link, None).await?;
123    let events = items
124        .into_iter()
125        .filter_map(|v| {
126            if v.get("@removed").is_some() {
127                return v
128                    .get("id")
129                    .and_then(Value::as_str)
130                    .map(|id| MailDeltaEvent::Deleted {
131                        graph_id: id.to_string(),
132                    });
133            }
134            // Graph delta change items carry only the properties that
135            // changed (plus id), so parse leniently and let consumers fetch
136            // details on demand.
137            let _ = account;
138            serde_json::from_value::<DeltaMessage>(v)
139                .ok()
140                .map(|m| MailDeltaEvent::Changed(Box::new(m)))
141        })
142        .collect();
143    Ok((events, next_link))
144}
145
146/// Establish calendar delta state over a rolling window.
147pub async fn calendar_delta_bootstrap(
148    http: &reqwest::Client,
149    base_url: &str,
150    access_token: &str,
151    account: &str,
152    start: DateTime<Utc>,
153    end: DateTime<Utc>,
154) -> Result<(Vec<Event>, String), ClientError> {
155    let url = format!(
156        "{base_url}/me/calendarView/delta?startDateTime={}&endDateTime={}",
157        start.to_rfc3339(),
158        end.to_rfc3339()
159    );
160    let (items, delta_link) =
161        drain(http, access_token, &url, Some("outlook.timezone=\"UTC\"")).await?;
162    let events = items
163        .into_iter()
164        .filter(|v| v.get("@removed").is_none())
165        .filter_map(|v| super::events::event_from_delta_value(v, account))
166        .collect();
167    Ok((events, delta_link))
168}
169
170/// Poll a calendar deltaLink.
171pub async fn calendar_delta(
172    http: &reqwest::Client,
173    access_token: &str,
174    account: &str,
175    delta_link: &str,
176) -> Result<(Vec<CalendarDeltaEvent>, String), ClientError> {
177    let (items, next_link) = drain(
178        http,
179        access_token,
180        delta_link,
181        Some("outlook.timezone=\"UTC\""),
182    )
183    .await?;
184    let events = items
185        .into_iter()
186        .filter_map(|v| {
187            if v.get("@removed").is_some() {
188                return v
189                    .get("id")
190                    .and_then(Value::as_str)
191                    .map(|id| CalendarDeltaEvent::Deleted {
192                        graph_id: id.to_string(),
193                    });
194            }
195            super::events::event_from_delta_value(v, account)
196                .map(|e| CalendarDeltaEvent::CreatedOrUpdated(Box::new(e)))
197        })
198        .collect();
199    Ok((events, next_link))
200}
201
202#[cfg(test)]
203mod tests {
204    use super::*;
205    use wiremock::matchers::{method, path, query_param};
206    use wiremock::{Mock, MockServer, ResponseTemplate};
207
208    fn graph_msg(id: &str) -> Value {
209        serde_json::json!({
210            "id": id,
211            "subject": format!("s-{id}"),
212            "from": {"emailAddress": {"name": "N", "address": "n@x.se"}},
213            "receivedDateTime": "2026-07-05T10:00:00Z",
214            "isRead": false,
215            "bodyPreview": "p",
216            "hasAttachments": false
217        })
218    }
219
220    #[tokio::test]
221    async fn bootstrap_follows_next_links_to_delta_link() {
222        let server = MockServer::start().await;
223        let page2 = format!("{}/delta-page-2", server.uri());
224        let final_delta = format!("{}/delta-final?token=abc", server.uri());
225        Mock::given(method("GET"))
226            .and(path("/me/mailFolders/inbox/messages/delta"))
227            .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
228                "value": [graph_msg("m1")],
229                "@odata.nextLink": page2
230            })))
231            .mount(&server)
232            .await;
233        Mock::given(method("GET"))
234            .and(path("/delta-page-2"))
235            .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
236                "value": [graph_msg("m2")],
237                "@odata.deltaLink": final_delta
238            })))
239            .mount(&server)
240            .await;
241
242        let http = reqwest::Client::new();
243        let (messages, delta_link) =
244            mail_delta_bootstrap(&http, &server.uri(), "tok", "a@b.se", "inbox")
245                .await
246                .unwrap();
247        assert_eq!(messages.len(), 2);
248        assert!(delta_link.contains("token=abc"));
249    }
250
251    #[tokio::test]
252    async fn poll_parses_changes_and_removals() {
253        let server = MockServer::start().await;
254        let next_delta = format!("{}/delta-final?token=next", server.uri());
255        Mock::given(method("GET"))
256            .and(path("/delta-poll"))
257            .and(query_param("token", "prev"))
258            .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
259                "value": [
260                    graph_msg("m9"),
261                    {"id": "gone-1", "@removed": {"reason": "deleted"}}
262                ],
263                "@odata.deltaLink": next_delta
264            })))
265            .mount(&server)
266            .await;
267
268        let http = reqwest::Client::new();
269        let url = format!("{}/delta-poll?token=prev", server.uri());
270        let (events, new_link) = mail_delta(&http, "tok", "a@b.se", &url).await.unwrap();
271        assert_eq!(events.len(), 2);
272        assert!(matches!(&events[0], MailDeltaEvent::Changed(m) if m.graph_id == "m9"));
273        assert!(matches!(&events[1], MailDeltaEvent::Deleted { graph_id } if graph_id == "gone-1"));
274        assert!(new_link.contains("token=next"));
275    }
276
277    #[tokio::test]
278    async fn expired_delta_is_typed() {
279        let server = MockServer::start().await;
280        Mock::given(method("GET"))
281            .and(path("/delta-poll"))
282            .respond_with(ResponseTemplate::new(410))
283            .mount(&server)
284            .await;
285
286        let http = reqwest::Client::new();
287        let url = format!("{}/delta-poll", server.uri());
288        let err = mail_delta(&http, "tok", "a@b.se", &url).await.unwrap_err();
289        assert!(matches!(err, ClientError::DeltaExpired));
290    }
291}