1use anyhow::Result;
9use url::Url;
10
11use crate::gmail::client::GmailClient;
12use crate::gmail::types::{Thread, ThreadListResponse};
13
14pub const MAX_PAGE_LIMIT: usize = 500;
16
17pub const HARD_CAP: usize = 10_000;
20
21#[derive(Debug, Clone, Copy, Default)]
27pub enum ThreadFormat {
28 Minimal,
30 #[default]
32 Full,
33 Metadata,
35}
36
37impl ThreadFormat {
38 fn as_str(self) -> &'static str {
39 match self {
40 Self::Minimal => "minimal",
41 Self::Full => "full",
42 Self::Metadata => "metadata",
43 }
44 }
45}
46
47#[derive(Debug)]
49pub struct ThreadsApi<'a> {
50 client: &'a GmailClient,
51}
52
53impl<'a> ThreadsApi<'a> {
54 #[must_use]
56 pub fn new(client: &'a GmailClient) -> Self {
57 Self { client }
58 }
59
60 pub async fn search(
65 &self,
66 query: Option<&str>,
67 label_ids: &[&str],
68 limit: usize,
69 page_token: Option<&str>,
70 ) -> Result<ThreadListResponse> {
71 if limit > MAX_PAGE_LIMIT {
72 return Err(anyhow::anyhow!(
73 "`limit` must be <= {MAX_PAGE_LIMIT} (Gmail threads.list per-page cap; use \
74 `search_all` to auto-paginate)"
75 ));
76 }
77 let url =
78 build_threads_list_url(self.client.base_url(), query, label_ids, limit, page_token)?;
79 self.client
80 .get_parsed(url.as_str(), "Failed to parse threads.list response")
81 .await
82 }
83
84 pub async fn search_all(
89 &self,
90 query: Option<&str>,
91 label_ids: &[&str],
92 limit: usize,
93 ) -> Result<ThreadListResponse> {
94 let cap = effective_cap(limit);
95 let mut acc: Option<ThreadListResponse> = None;
96 let mut page_token: Option<String> = None;
97 loop {
98 let collected = acc.as_ref().map_or(0, |r| r.threads.len());
99 let page_size = (cap - collected).min(MAX_PAGE_LIMIT);
100 let page = self
101 .search(query, label_ids, page_size, page_token.as_deref())
102 .await?;
103 let next_token = page.next_page_token.clone();
104 match acc.as_mut() {
105 Some(existing) => {
106 existing.threads.extend(page.threads);
107 existing.next_page_token = page.next_page_token;
108 existing.result_size_estimate = page.result_size_estimate;
109 }
110 None => acc = Some(page),
111 }
112 let collected = acc.as_ref().map_or(0, |r| r.threads.len());
113 if collected >= cap || next_token.is_none() {
114 break;
115 }
116 page_token = next_token;
117 }
118 let mut result = acc.unwrap_or_default();
119 result.threads.truncate(cap);
120 Ok(result)
121 }
122
123 pub async fn get(&self, id: &str, format: ThreadFormat) -> Result<Thread> {
125 let url = build_thread_get_url(self.client.base_url(), id, format)?;
126 self.client
127 .get_parsed(url.as_str(), "Failed to parse threads.get response")
128 .await
129 }
130}
131
132fn build_threads_list_url(
133 base_url: &str,
134 query: Option<&str>,
135 label_ids: &[&str],
136 limit: usize,
137 page_token: Option<&str>,
138) -> Result<Url> {
139 let mut url = GmailClient::api_url(base_url, "/gmail/v1/users/me/threads")?;
140 let query = query.filter(|q| !q.is_empty());
141 if query.is_some() || !label_ids.is_empty() || limit > 0 || page_token.is_some() {
145 let mut pairs = url.query_pairs_mut();
146 if let Some(q) = query {
147 pairs.append_pair("q", q);
148 }
149 for label in label_ids {
150 pairs.append_pair("labelIds", label);
151 }
152 if limit > 0 {
153 pairs.append_pair("maxResults", &limit.to_string());
154 }
155 if let Some(token) = page_token {
156 pairs.append_pair("pageToken", token);
157 }
158 }
159 Ok(url)
160}
161
162fn build_thread_get_url(base_url: &str, id: &str, format: ThreadFormat) -> Result<Url> {
163 let mut url = GmailClient::api_url(base_url, &format!("/gmail/v1/users/me/threads/{id}"))?;
164 url.query_pairs_mut().append_pair("format", format.as_str());
165 Ok(url)
166}
167
168fn effective_cap(limit: usize) -> usize {
171 if limit == 0 {
172 HARD_CAP
173 } else {
174 limit.min(HARD_CAP)
175 }
176}
177
178#[cfg(test)]
179#[allow(clippy::unwrap_used, clippy::expect_used)]
180mod tests {
181 use super::*;
182 use crate::gmail::auth::{GmailCredentials, GmailScope};
183 use crate::utils::secret::Secret;
184
185 fn test_credentials() -> GmailCredentials {
186 GmailCredentials {
187 client_id: "client-1".to_string(),
188 client_secret: Secret::new("secret-1"),
189 refresh_token: Secret::new("refresh-1"),
190 scope: GmailScope::ReadOnly,
191 }
192 }
193
194 fn dead_client() -> GmailClient {
195 let mut client = GmailClient::new("http://127.0.0.1:1", &test_credentials()).unwrap();
199 crate::gmail::client::test_support::replace_session(
200 &mut client,
201 &test_credentials(),
202 "http://127.0.0.1:1",
203 );
204 client
205 }
206
207 async fn client_with_bootstrapped_token(server: &wiremock::MockServer) -> GmailClient {
208 wiremock::Mock::given(wiremock::matchers::method("POST"))
209 .and(wiremock::matchers::path("/token"))
210 .respond_with(
211 wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({
212 "access_token": "test-token",
213 "expires_in": 3600,
214 })),
215 )
216 .mount(server)
217 .await;
218
219 let mut client = GmailClient::new(&server.uri(), &test_credentials()).unwrap();
220 crate::gmail::client::test_support::replace_session(
221 &mut client,
222 &test_credentials(),
223 &format!("{}/token", server.uri()),
224 );
225 client
226 }
227
228 fn thread_ref_json(id: &str) -> serde_json::Value {
229 serde_json::json!({"id": id})
230 }
231
232 fn page_body(ids: &[&str], next_token: Option<&str>) -> serde_json::Value {
233 let threads: Vec<serde_json::Value> = ids.iter().map(|id| thread_ref_json(id)).collect();
234 let mut body = serde_json::json!({"threads": threads});
235 if let Some(token) = next_token {
236 body["nextPageToken"] = serde_json::json!(token);
237 }
238 body
239 }
240
241 #[test]
244 fn build_threads_list_url_with_only_provided_filters() {
245 let url =
246 build_threads_list_url("https://gmail.googleapis.com", None, &[], 0, None).unwrap();
247 assert_eq!(
248 url.as_str(),
249 "https://gmail.googleapis.com/gmail/v1/users/me/threads"
250 );
251 }
252
253 #[test]
254 fn build_threads_list_url_with_full_filter_set() {
255 let url = build_threads_list_url(
256 "https://gmail.googleapis.com",
257 Some("label:finance"),
258 &["INBOX"],
259 25,
260 Some("cursor-1"),
261 )
262 .unwrap();
263 let query: Vec<_> = url.query_pairs().collect();
264 assert!(query.contains(&("q".into(), "label:finance".into())));
265 assert!(query.contains(&("labelIds".into(), "INBOX".into())));
266 assert!(query.contains(&("maxResults".into(), "25".into())));
267 assert!(query.contains(&("pageToken".into(), "cursor-1".into())));
268 }
269
270 #[test]
271 fn build_threads_list_url_rejects_invalid_base_url() {
272 let err = build_threads_list_url("not a url", None, &[], 0, None).unwrap_err();
273 assert!(err.to_string().contains("Invalid Gmail base URL"));
274 }
275
276 #[test]
277 fn build_thread_get_url_uses_thread_format_query_param_not_message_format() {
278 let url =
279 build_thread_get_url("https://gmail.googleapis.com", "t1", ThreadFormat::Metadata)
280 .unwrap();
281 assert!(url
282 .query_pairs()
283 .any(|pair| pair == ("format".into(), "metadata".into())));
284 }
285
286 #[tokio::test]
289 async fn search_propagates_api_errors() {
290 let server = wiremock::MockServer::start().await;
291 let client = client_with_bootstrapped_token(&server).await;
292 wiremock::Mock::given(wiremock::matchers::method("GET"))
293 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
294 .respond_with(wiremock::ResponseTemplate::new(400).set_body_string("bad query"))
295 .mount(&server)
296 .await;
297
298 let err = ThreadsApi::new(&client)
299 .search(Some("???"), &[], 10, None)
300 .await
301 .unwrap_err();
302 assert!(err.to_string().contains("400"));
303 }
304
305 #[tokio::test]
306 async fn search_rejects_limit_above_max_page_limit_client_side() {
307 let client = dead_client();
308 let err = ThreadsApi::new(&client)
309 .search(None, &[], MAX_PAGE_LIMIT + 1, None)
310 .await
311 .unwrap_err();
312 let msg = err.to_string();
313 assert!(msg.contains("limit"));
314 assert!(msg.contains("search_all"));
315 }
316
317 #[tokio::test]
318 async fn search_propagates_network_errors() {
319 let client = dead_client();
323 let err = ThreadsApi::new(&client)
324 .search(None, &[], 10, None)
325 .await
326 .unwrap_err();
327 assert!(err
328 .to_string()
329 .contains("Failed to obtain a Gmail access token"));
330 }
331
332 #[tokio::test]
333 async fn search_errors_on_malformed_response() {
334 let server = wiremock::MockServer::start().await;
335 let client = client_with_bootstrapped_token(&server).await;
336 wiremock::Mock::given(wiremock::matchers::method("GET"))
337 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
338 .respond_with(wiremock::ResponseTemplate::new(200).set_body_string("not json"))
339 .mount(&server)
340 .await;
341
342 let err = ThreadsApi::new(&client)
343 .search(None, &[], 10, None)
344 .await
345 .unwrap_err();
346 assert!(err.to_string().contains("Failed to parse"));
347 }
348
349 #[tokio::test]
352 async fn search_all_single_page_when_no_next_token() {
353 let server = wiremock::MockServer::start().await;
354 let client = client_with_bootstrapped_token(&server).await;
355 wiremock::Mock::given(wiremock::matchers::method("GET"))
356 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
357 .respond_with(
358 wiremock::ResponseTemplate::new(200).set_body_json(page_body(&["a", "b"], None)),
359 )
360 .expect(1)
361 .mount(&server)
362 .await;
363
364 let result = ThreadsApi::new(&client)
365 .search_all(None, &[], 100)
366 .await
367 .unwrap();
368 assert_eq!(result.threads.len(), 2);
369 }
370
371 #[tokio::test]
372 async fn search_all_follows_next_page_token_to_exhaustion() {
373 let server = wiremock::MockServer::start().await;
374 let client = client_with_bootstrapped_token(&server).await;
375 wiremock::Mock::given(wiremock::matchers::method("GET"))
376 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
377 .and(wiremock::matchers::query_param_is_missing("pageToken"))
378 .respond_with(
379 wiremock::ResponseTemplate::new(200)
380 .set_body_json(page_body(&["a", "b"], Some("c1"))),
381 )
382 .expect(1)
383 .mount(&server)
384 .await;
385 wiremock::Mock::given(wiremock::matchers::method("GET"))
386 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
387 .and(wiremock::matchers::query_param("pageToken", "c1"))
388 .respond_with(
389 wiremock::ResponseTemplate::new(200).set_body_json(page_body(&["c"], None)),
390 )
391 .expect(1)
392 .mount(&server)
393 .await;
394
395 let result = ThreadsApi::new(&client)
396 .search_all(None, &[], 0)
397 .await
398 .unwrap();
399 let ids: Vec<&str> = result.threads.iter().map(|t| t.id.as_str()).collect();
400 assert_eq!(ids, ["a", "b", "c"]);
401 }
402
403 #[tokio::test]
404 async fn search_all_stops_at_explicit_limit() {
405 let server = wiremock::MockServer::start().await;
406 let client = client_with_bootstrapped_token(&server).await;
407 wiremock::Mock::given(wiremock::matchers::method("GET"))
408 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
409 .respond_with(
410 wiremock::ResponseTemplate::new(200)
411 .set_body_json(page_body(&["a", "b", "c"], Some("more"))),
412 )
413 .expect(1)
414 .mount(&server)
415 .await;
416
417 let result = ThreadsApi::new(&client)
418 .search_all(None, &[], 3)
419 .await
420 .unwrap();
421 assert_eq!(result.threads.len(), 3);
422 }
423
424 #[tokio::test]
425 async fn search_all_truncates_to_hard_cap() {
426 let server = wiremock::MockServer::start().await;
427 let client = client_with_bootstrapped_token(&server).await;
428 let full_page: Vec<serde_json::Value> = (0..MAX_PAGE_LIMIT)
429 .map(|i| thread_ref_json(&format!("t{i}")))
430 .collect();
431 let body = serde_json::json!({"threads": full_page, "nextPageToken": "always-more"});
432 wiremock::Mock::given(wiremock::matchers::method("GET"))
433 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
434 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body))
435 .mount(&server)
436 .await;
437
438 let result = ThreadsApi::new(&client)
439 .search_all(None, &[], 0)
440 .await
441 .unwrap();
442 assert_eq!(result.threads.len(), HARD_CAP);
443 }
444
445 #[tokio::test]
446 async fn search_all_continues_past_empty_page_with_a_valid_next_page_token() {
447 let server = wiremock::MockServer::start().await;
448 let client = client_with_bootstrapped_token(&server).await;
449 wiremock::Mock::given(wiremock::matchers::method("GET"))
450 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
451 .and(wiremock::matchers::query_param_is_missing("pageToken"))
452 .respond_with(
453 wiremock::ResponseTemplate::new(200).set_body_json(page_body(&[], Some("p2"))),
454 )
455 .expect(1)
456 .mount(&server)
457 .await;
458 wiremock::Mock::given(wiremock::matchers::method("GET"))
459 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
460 .and(wiremock::matchers::query_param("pageToken", "p2"))
461 .respond_with(
462 wiremock::ResponseTemplate::new(200).set_body_json(page_body(&["a"], None)),
463 )
464 .expect(1)
465 .mount(&server)
466 .await;
467
468 let result = ThreadsApi::new(&client)
469 .search_all(Some("rare-query"), &[], 0)
470 .await
471 .unwrap();
472 assert_eq!(result.threads.len(), 1);
473 }
474
475 #[tokio::test]
476 async fn search_all_propagates_api_errors_on_first_page() {
477 let server = wiremock::MockServer::start().await;
478 let client = client_with_bootstrapped_token(&server).await;
479 wiremock::Mock::given(wiremock::matchers::method("GET"))
480 .and(wiremock::matchers::path("/gmail/v1/users/me/threads"))
481 .respond_with(wiremock::ResponseTemplate::new(403).set_body_string("nope"))
482 .mount(&server)
483 .await;
484
485 let err = ThreadsApi::new(&client)
486 .search_all(None, &[], 0)
487 .await
488 .unwrap_err();
489 assert!(err.to_string().contains("403"));
490 }
491
492 #[tokio::test]
495 async fn get_returns_thread_with_messages() {
496 let server = wiremock::MockServer::start().await;
497 let client = client_with_bootstrapped_token(&server).await;
498 wiremock::Mock::given(wiremock::matchers::method("GET"))
499 .and(wiremock::matchers::path("/gmail/v1/users/me/threads/t1"))
500 .respond_with(
501 wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({
502 "id": "t1",
503 "messages": [{"id": "m1", "threadId": "t1"}],
504 })),
505 )
506 .expect(1)
507 .mount(&server)
508 .await;
509
510 let thread = ThreadsApi::new(&client)
511 .get("t1", ThreadFormat::Full)
512 .await
513 .unwrap();
514 assert_eq!(thread.id, "t1");
515 assert_eq!(thread.messages.len(), 1);
516 }
517
518 #[tokio::test]
519 async fn get_sends_minimal_format_query_param() {
520 let server = wiremock::MockServer::start().await;
521 let client = client_with_bootstrapped_token(&server).await;
522 wiremock::Mock::given(wiremock::matchers::method("GET"))
523 .and(wiremock::matchers::path("/gmail/v1/users/me/threads/t1"))
524 .and(wiremock::matchers::query_param("format", "minimal"))
525 .respond_with(
526 wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": "t1"})),
527 )
528 .expect(1)
529 .mount(&server)
530 .await;
531
532 let thread = ThreadsApi::new(&client)
533 .get("t1", ThreadFormat::Minimal)
534 .await
535 .unwrap();
536 assert_eq!(thread.id, "t1");
537 }
538
539 #[test]
542 fn effective_cap_zero_is_hard_cap() {
543 assert_eq!(effective_cap(0), HARD_CAP);
544 }
545
546 #[test]
547 fn effective_cap_clamps_above_hard_cap() {
548 assert_eq!(effective_cap(HARD_CAP + 5), HARD_CAP);
549 }
550}