1use chrono::{DateTime, Utc};
10use pidge_core::{Event, Message};
11use serde_json::Value;
12
13use crate::error::ClientError;
14
15#[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
48async 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
94pub 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
115pub 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 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
146pub 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
170pub 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}