pub mod batch;
mod calendars;
pub mod delta;
pub mod events;
mod mail;
mod me;
mod people;
pub use calendars::list_calendars;
pub use events::{
EventsPage, NewEvent, ProposedTime, RsvpKind, cancel_event, create_event, delete_event,
get_event, list_calendar_view, list_events_at, move_event_to_calendar, move_time, rsvp_event,
update_event,
};
pub use mail::{
InboxPage, MailFolder, Outgoing, add_attachment, create_child_folder, create_draft,
create_forward_draft, create_mail_folder, create_reply_all_draft, create_reply_draft,
delete_attachment, delete_mail_folder, delete_message, fetch_message_headers, forward_message,
get_attachment_bytes, get_categories, get_message, list_attachments, list_child_folders,
list_drafts, list_folder_messages, list_inbox, list_mail_folders, list_messages_at, mark_read,
mark_unread, move_message, reply_all_message, reply_message, search_folder_messages,
search_messages, send_draft, send_mail, set_categories, set_flag, unsubscribe_one_click,
update_draft, update_draft_recipients,
};
pub use me::{Me, get_me};
pub use people::{Person, list_people};
use crate::auth::AuthClient;
use crate::auth::config;
use crate::error::ClientError;
use pidge_core::Message;
const MAX_ATTEMPTS: u32 = 4;
fn is_transient(status: reqwest::StatusCode) -> bool {
matches!(status.as_u16(), 429 | 503 | 504)
}
pub(crate) async fn send_with_retry(
req: reqwest::RequestBuilder,
) -> Result<reqwest::Response, ClientError> {
let mut attempt: u32 = 0;
loop {
let this_try = match req.try_clone() {
Some(clone) => clone,
None => return Ok(req.send().await?),
};
let resp = this_try.send().await?;
let status = resp.status();
if !is_transient(status) {
return Ok(resp);
}
let retry_after = resp
.headers()
.get(reqwest::header::RETRY_AFTER)
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse::<u64>().ok());
attempt += 1;
if attempt >= MAX_ATTEMPTS {
return Err(ClientError::Throttled { retry_after });
}
let backoff = retry_after
.map(std::time::Duration::from_secs)
.unwrap_or_else(|| {
let jitter = std::time::Duration::from_millis(u64::from(attempt) * 83 % 250);
std::time::Duration::from_secs(1u64 << attempt.min(4)) / 2 + jitter
});
tracing::debug!(
status = status.as_u16(),
attempt,
?backoff,
"retrying Graph request"
);
tokio::time::sleep(backoff).await;
}
}
pub(crate) fn check_continuation(url: &str, base_url: &str) -> Result<(), ClientError> {
let origin = |u: &str| url::Url::parse(u).ok().map(|u| u.origin());
let allowed = origin(url).is_some_and(|o| {
o.is_tuple()
&& (Some(&o) == origin(config::GRAPH_BASE).as_ref()
|| Some(&o) == origin(base_url).as_ref())
});
if allowed {
Ok(())
} else {
Err(ClientError::Graph {
status: 400,
message: "refusing to follow a continuation link off graph.microsoft.com".into(),
})
}
}
pub struct GraphClient {
auth: AuthClient,
http: reqwest::Client,
base_url: String,
one_click: mail::OneClickPolicy,
}
impl GraphClient {
pub fn new(auth: AuthClient) -> Result<Self, ClientError> {
Ok(Self {
auth,
http: reqwest::Client::builder()
.user_agent(format!("pidge/{}", env!("CARGO_PKG_VERSION")))
.build()?,
base_url: config::GRAPH_BASE.to_string(),
one_click: mail::OneClickPolicy::PublicOnly,
})
}
pub fn for_test(auth: AuthClient, base_url: impl Into<String>) -> Self {
Self {
auth,
http: reqwest::Client::new(),
base_url: base_url.into(),
one_click: mail::OneClickPolicy::AllowLoopback,
}
}
pub fn auth(&self) -> &AuthClient {
&self.auth
}
pub async fn me(&self, access_token: &str) -> Result<Me, ClientError> {
get_me(&self.http, &self.base_url, access_token).await
}
pub async fn list_inbox(
&self,
account: &str,
limit: usize,
skip: usize,
unread_only: bool,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
list_inbox(
&self.http,
&self.base_url,
&token,
account,
limit,
skip,
unread_only,
)
.await
}
pub async fn list_folder(
&self,
account: &str,
folder_id: &str,
limit: usize,
skip: usize,
unread_only: bool,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
list_folder_messages(
&self.http,
&self.base_url,
&token,
account,
folder_id,
limit,
skip,
unread_only,
)
.await
}
pub async fn mail_delta_bootstrap(
&self,
account: &str,
folder: &str,
) -> Result<(Vec<Message>, String), ClientError> {
let token = self.auth.get_valid_token(account).await?;
delta::mail_delta_bootstrap(&self.http, &self.base_url, &token, account, folder).await
}
pub async fn mail_delta(
&self,
account: &str,
delta_link: &str,
) -> Result<(Vec<delta::MailDeltaEvent>, String), ClientError> {
let token = self.auth.get_valid_token(account).await?;
delta::mail_delta(&self.http, &token, account, delta_link).await
}
pub async fn calendar_delta_bootstrap(
&self,
account: &str,
start: chrono::DateTime<chrono::Utc>,
end: chrono::DateTime<chrono::Utc>,
) -> Result<(Vec<pidge_core::Event>, String), ClientError> {
let token = self.auth.get_valid_token(account).await?;
delta::calendar_delta_bootstrap(&self.http, &self.base_url, &token, account, start, end)
.await
}
pub async fn calendar_delta(
&self,
account: &str,
delta_link: &str,
) -> Result<(Vec<delta::CalendarDeltaEvent>, String), ClientError> {
let token = self.auth.get_valid_token(account).await?;
delta::calendar_delta(&self.http, &token, account, delta_link).await
}
pub async fn batch_all(
&self,
account: &str,
requests: Vec<batch::BatchRequest>,
) -> Result<Vec<batch::BatchResponse>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
batch::batch_all(&self.http, &self.base_url, &token, requests).await
}
pub async fn list_conversation(
&self,
account: &str,
conversation_id: &str,
) -> Result<Vec<Message>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::list_conversation(&self.http, &self.base_url, &token, account, conversation_id).await
}
pub async fn list_messages_at(
&self,
account: &str,
url: &str,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
list_messages_at(&self.http, &self.base_url, &token, account, url).await
}
pub async fn search_messages(
&self,
account: &str,
query: &str,
limit: usize,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
search_messages(&self.http, &self.base_url, &token, account, query, limit).await
}
pub async fn search_folder(
&self,
account: &str,
folder: &str,
query: &str,
limit: usize,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
search_folder_messages(
&self.http,
&self.base_url,
&token,
account,
folder,
query,
limit,
)
.await
}
pub async fn mark_unread(&self, account: &str, message_id: &str) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::mark_unread(&self.http, &self.base_url, &token, message_id).await
}
pub async fn set_flag(
&self,
account: &str,
message_id: &str,
flagged: bool,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::set_flag(&self.http, &self.base_url, &token, message_id, flagged).await
}
pub async fn get_categories(
&self,
account: &str,
message_id: &str,
) -> Result<Vec<String>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::get_categories(&self.http, &self.base_url, &token, message_id).await
}
pub async fn set_categories(
&self,
account: &str,
message_id: &str,
categories: &[String],
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::set_categories(&self.http, &self.base_url, &token, message_id, categories).await
}
pub async fn move_message(
&self,
account: &str,
message_id: &str,
destination: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::move_message(&self.http, &self.base_url, &token, message_id, destination).await
}
pub async fn list_mail_folders(&self, account: &str) -> Result<Vec<MailFolder>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::list_mail_folders(&self.http, &self.base_url, &token).await
}
pub async fn create_mail_folder(
&self,
account: &str,
display_name: &str,
) -> Result<MailFolder, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_mail_folder(&self.http, &self.base_url, &token, display_name).await
}
pub async fn list_child_folders(
&self,
account: &str,
parent_id: &str,
) -> Result<Vec<MailFolder>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::list_child_folders(&self.http, &self.base_url, &token, parent_id).await
}
pub async fn create_child_folder(
&self,
account: &str,
parent_id: &str,
display_name: &str,
) -> Result<MailFolder, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_child_folder(&self.http, &self.base_url, &token, parent_id, display_name).await
}
pub async fn delete_mail_folder(
&self,
account: &str,
folder_id: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::delete_mail_folder(&self.http, &self.base_url, &token, folder_id).await
}
pub async fn send_mail(&self, account: &str, message: &Outgoing) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::send_mail(&self.http, &self.base_url, &token, message).await
}
pub async fn reply_message(
&self,
account: &str,
message_id: &str,
comment: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::reply_message(&self.http, &self.base_url, &token, message_id, comment).await
}
pub async fn reply_all_message(
&self,
account: &str,
message_id: &str,
comment: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::reply_all_message(&self.http, &self.base_url, &token, message_id, comment).await
}
pub async fn forward_message(
&self,
account: &str,
message_id: &str,
to: &[String],
comment: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::forward_message(&self.http, &self.base_url, &token, message_id, to, comment).await
}
pub async fn list_drafts(
&self,
account: &str,
limit: usize,
skip: usize,
) -> Result<InboxPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::list_drafts(&self.http, &self.base_url, &token, account, limit, skip).await
}
pub async fn create_draft(
&self,
account: &str,
message: &Outgoing,
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_draft(&self.http, &self.base_url, &token, message).await
}
pub async fn create_reply_draft(
&self,
account: &str,
message_id: &str,
comment: &str,
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_reply_draft(&self.http, &self.base_url, &token, message_id, comment).await
}
pub async fn create_reply_all_draft(
&self,
account: &str,
message_id: &str,
comment: &str,
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_reply_all_draft(&self.http, &self.base_url, &token, message_id, comment).await
}
pub async fn create_forward_draft(
&self,
account: &str,
message_id: &str,
to: &[String],
comment: &str,
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::create_forward_draft(&self.http, &self.base_url, &token, message_id, to, comment)
.await
}
pub async fn send_draft(&self, account: &str, message_id: &str) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::send_draft(&self.http, &self.base_url, &token, message_id).await
}
pub async fn update_draft(
&self,
account: &str,
message_id: &str,
message: &Outgoing,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::update_draft(&self.http, &self.base_url, &token, message_id, message).await
}
pub async fn update_draft_recipients(
&self,
account: &str,
message_id: &str,
to: Option<&[String]>,
cc: Option<&[String]>,
bcc: Option<&[String]>,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::update_draft_recipients(&self.http, &self.base_url, &token, message_id, to, cc, bcc)
.await
}
pub async fn delete_message(&self, account: &str, message_id: &str) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::delete_message(&self.http, &self.base_url, &token, message_id).await
}
pub async fn add_attachment(
&self,
account: &str,
message_id: &str,
name: &str,
content_type: &str,
bytes: &[u8],
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::add_attachment(
&self.http,
&self.base_url,
&token,
message_id,
name,
content_type,
bytes,
)
.await
}
pub async fn delete_attachment(
&self,
account: &str,
message_id: &str,
attachment_id: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::delete_attachment(
&self.http,
&self.base_url,
&token,
message_id,
attachment_id,
)
.await
}
pub async fn get_message(
&self,
account: &str,
message_id: &str,
) -> Result<pidge_core::FullMessage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::get_message(&self.http, &self.base_url, &token, account, message_id).await
}
pub async fn fetch_message_headers(
&self,
account: &str,
message_id: &str,
) -> Result<Vec<(String, String)>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::fetch_message_headers(&self.http, &self.base_url, &token, message_id).await
}
pub async fn list_attachments(
&self,
account: &str,
message_id: &str,
) -> Result<Vec<pidge_core::Attachment>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::list_attachments(&self.http, &self.base_url, &token, message_id).await
}
pub async fn unsubscribe_one_click(&self, url: &str) -> Result<(), ClientError> {
mail::post_one_click(url, self.one_click).await
}
pub async fn list_people(
&self,
account: &str,
top: usize,
) -> Result<Vec<people::Person>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
people::list_people(&self.http, &self.base_url, &token, top).await
}
pub async fn get_attachment_bytes(
&self,
account: &str,
message_id: &str,
attachment_id: &str,
) -> Result<Vec<u8>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::get_attachment_bytes(
&self.http,
&self.base_url,
&token,
message_id,
attachment_id,
)
.await
}
pub async fn mark_read(&self, account: &str, message_id: &str) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
mail::mark_read(&self.http, &self.base_url, &token, message_id).await
}
pub async fn list_calendars(
&self,
account: &str,
) -> Result<Vec<pidge_core::Calendar>, ClientError> {
let token = self.auth.get_valid_token(account).await?;
calendars::list_calendars(&self.http, &self.base_url, &token, account).await
}
pub async fn list_calendar_view(
&self,
account: &str,
calendar_id: Option<&str>,
start: chrono::DateTime<chrono::Utc>,
end: chrono::DateTime<chrono::Utc>,
limit: usize,
) -> Result<events::EventsPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::list_calendar_view(
&self.http,
&self.base_url,
&token,
account,
calendar_id,
start,
end,
limit,
)
.await
}
pub async fn list_events_at(
&self,
account: &str,
url: &str,
) -> Result<EventsPage, ClientError> {
let token = self.auth.get_valid_token(account).await?;
list_events_at(&self.http, &self.base_url, &token, account, url).await
}
pub async fn get_event(
&self,
account: &str,
event_id: &str,
) -> Result<pidge_core::Event, ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::get_event(&self.http, &self.base_url, &token, account, event_id).await
}
pub async fn create_event(
&self,
account: &str,
calendar_id: Option<&str>,
new_event: &events::NewEvent,
) -> Result<String, ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::create_event(&self.http, &self.base_url, &token, calendar_id, new_event).await
}
pub async fn update_event(
&self,
account: &str,
event_id: &str,
new_event: &events::NewEvent,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::update_event(&self.http, &self.base_url, &token, event_id, new_event).await
}
pub async fn move_time(
&self,
account: &str,
event_id: &str,
start: chrono::DateTime<chrono::Utc>,
end: chrono::DateTime<chrono::Utc>,
tz: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::move_time(&self.http, &self.base_url, &token, event_id, start, end, tz).await
}
pub async fn delete_event(&self, account: &str, event_id: &str) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::delete_event(&self.http, &self.base_url, &token, event_id).await
}
pub async fn cancel_event(
&self,
account: &str,
event_id: &str,
comment: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::cancel_event(&self.http, &self.base_url, &token, event_id, comment).await
}
#[allow(clippy::too_many_arguments)]
pub async fn rsvp_event(
&self,
account: &str,
event_id: &str,
kind: events::RsvpKind,
comment: &str,
send_response: bool,
proposed: Option<&events::ProposedTime>,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::rsvp_event(
&self.http,
&self.base_url,
&token,
event_id,
kind,
comment,
send_response,
proposed,
)
.await
}
pub async fn move_event_to_calendar(
&self,
account: &str,
event_id: &str,
destination_calendar_id: &str,
) -> Result<(), ClientError> {
let token = self.auth.get_valid_token(account).await?;
events::move_event_to_calendar(
&self.http,
&self.base_url,
&token,
event_id,
destination_calendar_id,
)
.await
}
}
#[cfg(test)]
mod retry_tests {
use super::*;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, Request, Respond, ResponseTemplate};
struct FlakyResponder {
failures: std::sync::atomic::AtomicU32,
}
impl Respond for FlakyResponder {
fn respond(&self, _req: &Request) -> ResponseTemplate {
let n = self
.failures
.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
if n > 0 {
ResponseTemplate::new(429).insert_header("Retry-After", "0")
} else {
ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": true}))
}
}
}
#[tokio::test]
async fn retries_transient_429_until_success() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/thing"))
.respond_with(FlakyResponder {
failures: std::sync::atomic::AtomicU32::new(2),
})
.expect(3)
.mount(&server)
.await;
let http = reqwest::Client::new();
let resp = send_with_retry(http.get(format!("{}/thing", server.uri())))
.await
.unwrap();
assert_eq!(resp.status(), 200);
}
#[tokio::test]
async fn persistent_429_becomes_throttled_after_max_attempts() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/thing"))
.respond_with(ResponseTemplate::new(429).insert_header("Retry-After", "0"))
.expect(4)
.mount(&server)
.await;
let http = reqwest::Client::new();
let err = send_with_retry(http.get(format!("{}/thing", server.uri())))
.await
.unwrap_err();
match err {
ClientError::Throttled { retry_after } => assert_eq!(retry_after, Some(0)),
other => panic!("expected Throttled, got {other:?}"),
}
}
#[tokio::test]
async fn non_transient_errors_pass_through_without_retry() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/thing"))
.respond_with(ResponseTemplate::new(400).set_body_string("bad"))
.expect(1)
.mount(&server)
.await;
let http = reqwest::Client::new();
let resp = send_with_retry(http.get(format!("{}/thing", server.uri())))
.await
.unwrap();
assert_eq!(resp.status(), 400);
}
}
#[cfg(test)]
mod continuation_tests {
use super::*;
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
const REFUSAL: &str = "refusing to follow a continuation link off graph.microsoft.com";
fn refused(r: Result<(), ClientError>) -> bool {
matches!(r, Err(ClientError::Graph { status: 400, ref message }) if message == REFUSAL)
}
#[test]
fn graph_links_are_followed_and_others_refused_in_production() {
let base = config::GRAPH_BASE;
assert!(
check_continuation(
"https://graph.microsoft.com/v1.0/me/messages?$skiptoken=x",
base
)
.is_ok()
);
for url in [
"https://evil.example.com/v1.0/me/messages",
"http://graph.microsoft.com/v1.0/me/messages",
"https://graph.microsoft.com.evil.example.com/v1.0",
"https://graph.microsoft.com:8443/v1.0",
"not a url",
] {
assert!(refused(check_continuation(url, base)), "{url}");
}
}
#[test]
fn a_test_base_url_allows_its_own_origin_only() {
let base = "http://127.0.0.1:4000/v1.0";
assert!(check_continuation("http://127.0.0.1:4000/v1.0/page-2", base).is_ok());
assert!(refused(check_continuation(
"http://127.0.0.1:4001/v1.0/page-2",
base
)));
}
#[tokio::test]
async fn off_graph_next_links_are_refused_without_a_request() {
let graph = MockServer::start().await;
let other = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": []
})))
.expect(0)
.mount(&other)
.await;
let http = reqwest::Client::new();
let base = format!("{}/v1.0", graph.uri());
let link = format!("{}/v1.0/page-2", other.uri());
let Err(err) = events::list_events_at(&http, &base, "tok", "a@b.se", &link).await else {
panic!("list_events_at followed an off-Graph link");
};
assert!(refused(Err(err)));
let Err(err) = mail::list_messages_at(&http, &base, "tok", "a@b.se", &link).await else {
panic!("list_messages_at followed an off-Graph link");
};
assert!(refused(Err(err)));
assert!(other.received_requests().await.unwrap().is_empty());
}
}