Skip to main content

omni_dev/gmail/
threads_api.rs

1//! Gmail Threads API wrapper.
2//!
3//! Same cursor-pagination shape as [`crate::gmail::messages_api`].
4//! Read-only in Phase 1 — `threads.modify`/`.trash` are real Gmail
5//! endpoints but outside this issue's stated surface, an explicit non-goal
6//! rather than an oversight.
7
8use anyhow::Result;
9use url::Url;
10
11use crate::gmail::client::GmailClient;
12use crate::gmail::types::{Thread, ThreadListResponse};
13
14/// Maximum page size accepted by `GET /gmail/v1/users/{userId}/threads`.
15pub const MAX_PAGE_LIMIT: usize = 500;
16
17/// Per-call upper bound on the number of threads returned by
18/// [`ThreadsApi::search_all`], even when the caller passes `limit = 0`.
19pub const HARD_CAP: usize = 10_000;
20
21/// The `format` query parameter accepted by `threads.get`.
22///
23/// Deliberately has no `Raw` variant — meaningless for a thread (a thread's
24/// whole point is showing the conversation's parsed messages), unlike
25/// [`crate::gmail::messages_api::MessageFormat`].
26#[derive(Debug, Clone, Copy, Default)]
27pub enum ThreadFormat {
28    /// Only `id`/`historyId` per message — no headers or body.
29    Minimal,
30    /// The full parsed MIME structure for every message. Default.
31    #[default]
32    Full,
33    /// Headers and snippet only per message, no body.
34    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/// Threads API façade.
48#[derive(Debug)]
49pub struct ThreadsApi<'a> {
50    client: &'a GmailClient,
51}
52
53impl<'a> ThreadsApi<'a> {
54    /// Wraps an existing [`GmailClient`] for thread operations.
55    #[must_use]
56    pub fn new(client: &'a GmailClient) -> Self {
57        Self { client }
58    }
59
60    /// Searches threads matching `query`, returning a single page.
61    ///
62    /// `limit` is rejected client-side when it exceeds [`MAX_PAGE_LIMIT`];
63    /// use [`Self::search_all`] to auto-paginate across pages.
64    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    /// Searches threads, auto-paginating via cursor as needed.
85    ///
86    /// Same "cursor-only, never a short page" termination rule as
87    /// [`crate::gmail::messages_api::MessagesApi::search_all`].
88    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    /// Fetches a single thread (with its messages) by id.
124    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    // Only touch `query_pairs_mut()` when there's something to append —
142    // calling it unconditionally leaves a bare trailing `?` even with zero
143    // pairs appended.
144    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
168/// Clamps a caller-supplied limit to [`HARD_CAP`], treating `0` as "fetch
169/// as many as the cap allows".
170fn 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        // Routes the session's token endpoint to the same dead address —
196        // otherwise `GmailSession` would try to refresh against the real
197        // Google token endpoint before the API call is ever attempted.
198        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    // ── URL builders (pure) ──────────────────────────────────────────
242
243    #[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    // ── Standard error paths ─────────────────────────────────────────
287
288    #[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        // `dead_client()` also points the session's token endpoint at the
320        // dead address, so the failure surfaces during token acquisition
321        // before the threads.list request is ever attempted.
322        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    // ── Pagination ────────────────────────────────────────────────────
350
351    #[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    // ── get ───────────────────────────────────────────────────────────
493
494    #[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    // ── effective_cap ─────────────────────────────────────────────────
540
541    #[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}