use pidge_core::{FlagStatus, Message, MessageFrom};
use serde::Deserialize;
use crate::error::ClientError;
#[derive(Debug, Deserialize)]
struct GraphMessage {
id: String,
subject: Option<String>,
from: Option<GraphFromWrapper>,
#[serde(rename = "receivedDateTime")]
received_date_time: chrono::DateTime<chrono::Utc>,
#[serde(rename = "isRead")]
is_read: Option<bool>,
#[serde(rename = "bodyPreview")]
body_preview: Option<String>,
body: Option<GraphBody>,
#[serde(rename = "hasAttachments")]
has_attachments: Option<bool>,
#[serde(rename = "conversationId", default)]
conversation_id: Option<String>,
#[serde(default)]
flag: Option<GraphFlag>,
#[serde(rename = "toRecipients", default)]
to_recipients: Vec<GraphFromWrapper>,
#[serde(rename = "ccRecipients", default)]
cc_recipients: Vec<GraphFromWrapper>,
#[serde(rename = "@odata.type", default)]
odata_type: Option<String>,
}
fn is_event_message(odata_type: Option<&str>) -> bool {
odata_type == Some("#microsoft.graph.eventMessageRequest")
}
#[derive(Debug, Deserialize)]
struct GraphFlag {
#[serde(rename = "flagStatus", default)]
flag_status: Option<String>,
}
fn recipient_from(addr: GraphFromAddress) -> MessageFrom {
MessageFrom {
name: clean(addr.name),
address: clean(addr.address),
}
}
pub(crate) fn clean(text: Option<String>) -> String {
text.as_deref()
.map(pidge_core::render::strip_controls)
.unwrap_or_default()
}
pub(crate) fn clean_string<'de, D: serde::Deserializer<'de>>(d: D) -> Result<String, D::Error> {
let s = <String as serde::Deserialize>::deserialize(d)?;
Ok(pidge_core::render::strip_controls(&s))
}
fn unwrap_recipients(rs: Vec<GraphFromWrapper>) -> Vec<MessageFrom> {
rs.into_iter()
.map(|w| recipient_from(w.email_address))
.collect()
}
fn flag_status_from(g: Option<GraphFlag>) -> FlagStatus {
match g.and_then(|f| f.flag_status).as_deref() {
Some("flagged") => FlagStatus::Flagged,
Some("complete") => FlagStatus::Complete,
_ => FlagStatus::NotFlagged,
}
}
#[derive(Debug, Deserialize)]
struct GraphFromWrapper {
#[serde(rename = "emailAddress")]
email_address: GraphFromAddress,
}
#[derive(Debug, Deserialize)]
struct GraphFromAddress {
name: Option<String>,
address: Option<String>,
}
#[derive(Debug, Deserialize)]
struct GraphList {
value: Vec<GraphMessage>,
#[serde(rename = "@odata.nextLink", default)]
next_link: Option<String>,
}
pub struct InboxPage {
pub messages: Vec<Message>,
pub has_more: bool,
pub next_link: Option<String>,
}
#[derive(Debug, Deserialize)]
struct GraphFullMessage {
id: String,
#[serde(rename = "conversationId", default)]
conversation_id: Option<String>,
subject: Option<String>,
from: Option<GraphFromWrapper>,
#[serde(rename = "toRecipients", default)]
to_recipients: Vec<GraphFromWrapper>,
#[serde(rename = "ccRecipients", default)]
cc_recipients: Vec<GraphFromWrapper>,
#[serde(rename = "bccRecipients", default)]
bcc_recipients: Vec<GraphFromWrapper>,
#[serde(rename = "receivedDateTime")]
received_date_time: chrono::DateTime<chrono::Utc>,
#[serde(rename = "sentDateTime")]
sent_date_time: chrono::DateTime<chrono::Utc>,
#[serde(rename = "isRead")]
is_read: Option<bool>,
body: GraphBody,
#[serde(rename = "hasAttachments")]
has_attachments: Option<bool>,
#[serde(default)]
flag: Option<GraphFlag>,
#[serde(rename = "@odata.type", default)]
odata_type: Option<String>,
#[serde(rename = "isDraft", default)]
is_draft: Option<bool>,
}
#[derive(Debug, Deserialize)]
struct GraphBody {
#[serde(rename = "contentType")]
content_type: String,
content: String,
}
#[derive(Debug, Deserialize)]
struct GraphAttachmentList {
value: Vec<GraphAttachment>,
}
#[derive(Debug, Deserialize)]
struct GraphAttachment {
id: String,
name: Option<String>,
#[serde(rename = "contentType")]
content_type: Option<String>,
size: Option<u64>,
#[serde(rename = "isInline")]
is_inline: Option<bool>,
#[serde(rename = "contentId")]
content_id: Option<String>,
#[serde(rename = "@odata.type", default)]
odata_type: Option<String>,
#[serde(rename = "contentBytes", default)]
content_bytes: Option<String>,
}
pub async fn list_inbox(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
limit: usize,
skip: usize,
unread_only: bool,
) -> Result<InboxPage, ClientError> {
list_folder(
http,
base_url,
access_token,
account,
"inbox",
limit,
skip,
unread_only,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn list_folder_messages(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
folder_id: &str,
limit: usize,
skip: usize,
unread_only: bool,
) -> Result<InboxPage, ClientError> {
list_folder(
http,
base_url,
access_token,
account,
folder_id,
limit,
skip,
unread_only,
)
.await
}
pub async fn list_drafts(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
limit: usize,
skip: usize,
) -> Result<InboxPage, ClientError> {
list_folder(
http,
base_url,
access_token,
account,
"drafts",
limit,
skip,
false,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn list_folder(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
folder: &str,
limit: usize,
skip: usize,
unread_only: bool,
) -> Result<InboxPage, ClientError> {
let url = format!("{base_url}/me/mailFolders/{folder}/messages");
let mut req = http.get(&url).bearer_auth(access_token).query(&[
(
"$select",
"id,subject,from,receivedDateTime,isRead,bodyPreview,body,hasAttachments,flag,conversationId,toRecipients,ccRecipients",
),
("$orderby", "receivedDateTime desc"),
("$top", &limit.to_string()),
]);
if skip > 0 {
req = req.query(&[("$skip", &skip.to_string())]);
}
if unread_only {
req = req.query(&[("$filter", "isRead eq false")]);
}
let resp = super::send_with_retry(req).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphList = resp.json().await?;
Ok(InboxPage {
has_more: list.next_link.is_some(),
next_link: list.next_link,
messages: list
.value
.into_iter()
.map(|g| to_message(g, account))
.collect(),
})
}
pub(crate) fn message_from_delta_value(
value: serde_json::Value,
account: &str,
) -> Option<pidge_core::Message> {
let g: GraphMessage = serde_json::from_value(value).ok()?;
Some(to_message(g, account))
}
pub async fn list_conversation(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
conversation_id: &str,
) -> Result<Vec<pidge_core::Message>, ClientError> {
let url = format!("{base_url}/me/messages");
let filter = format!(
"conversationId eq '{}'",
conversation_id.replace('\'', "''")
);
let req = http
.get(&url)
.bearer_auth(access_token)
.header("Prefer", "outlook.body-content-type=\"text\"")
.query(&[
(
"$select",
"id,subject,from,receivedDateTime,isRead,bodyPreview,body,hasAttachments,flag,conversationId,toRecipients,ccRecipients",
),
("$filter", filter.as_str()),
("$top", "100"),
]);
let resp = super::send_with_retry(req).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphList = resp.json().await?;
let mut messages: Vec<pidge_core::Message> = list
.value
.into_iter()
.map(|g| to_message(g, account))
.collect();
messages.sort_by_key(|m| m.received_at);
Ok(messages)
}
pub async fn list_messages_at(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
url: &str,
) -> Result<InboxPage, ClientError> {
super::check_continuation(url, base_url)?;
let req = http.get(url).bearer_auth(access_token);
let resp = super::send_with_retry(req).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphList = resp.json().await?;
Ok(InboxPage {
has_more: list.next_link.is_some(),
next_link: list.next_link,
messages: list
.value
.into_iter()
.map(|g| to_message(g, account))
.collect(),
})
}
pub async fn search_messages(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
query: &str,
limit: usize,
) -> Result<InboxPage, ClientError> {
search_in(http, base_url, access_token, account, None, query, limit).await
}
pub async fn search_folder_messages(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
folder: &str,
query: &str,
limit: usize,
) -> Result<InboxPage, ClientError> {
search_in(
http,
base_url,
access_token,
account,
Some(folder),
query,
limit,
)
.await
}
async fn search_in(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
folder: Option<&str>,
query: &str,
limit: usize,
) -> Result<InboxPage, ClientError> {
let escaped = query.replace('\\', "\\\\").replace('"', "\\\"");
let quoted = format!("\"{escaped}\"");
let url = match folder {
Some(f) => format!("{base_url}/me/mailFolders/{f}/messages"),
None => format!("{base_url}/me/messages"),
};
let resp = super::send_with_retry(
http.get(&url)
.bearer_auth(access_token)
.header("Prefer", "outlook.body-content-type=\"text\"")
.query(&[
(
"$select",
"id,subject,from,receivedDateTime,isRead,bodyPreview,body,hasAttachments,flag,conversationId,toRecipients,ccRecipients",
),
("$top", &limit.to_string()),
("$search", "ed),
]),
)
.await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphList = resp.json().await?;
Ok(InboxPage {
has_more: list.next_link.is_some(),
next_link: list.next_link,
messages: list
.value
.into_iter()
.map(|g| to_message(g, account))
.collect(),
})
}
fn to_message(g: GraphMessage, account: &str) -> Message {
let (body, body_content_type) = match g.body {
Some(b) => {
let kind = if b.content_type.eq_ignore_ascii_case("html") {
pidge_core::BodyContentType::Html
} else {
pidge_core::BodyContentType::Text
};
(clean(Some(b.content)), kind)
}
None => (String::new(), pidge_core::BodyContentType::Text),
};
Message {
account: account.to_string(),
id: g.id,
conversation_id: g.conversation_id.unwrap_or_default(),
from: MessageFrom {
name: clean(g.from.as_ref().and_then(|f| f.email_address.name.clone())),
address: clean(
g.from
.as_ref()
.and_then(|f| f.email_address.address.clone()),
),
},
subject: clean(g.subject),
received_at: g.received_date_time,
is_read: g.is_read.unwrap_or(true),
preview: clean(g.body_preview),
flag_status: flag_status_from(g.flag),
has_attachments: g.has_attachments.unwrap_or(false),
body,
body_content_type,
to: unwrap_recipients(g.to_recipients),
cc: unwrap_recipients(g.cc_recipients),
is_invite: is_event_message(g.odata_type.as_deref()),
}
}
pub async fn get_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
account: &str,
message_id: &str,
) -> Result<pidge_core::FullMessage, ClientError> {
let url = format!(
"{base_url}/me/messages/{message_id}\
?$select=id,subject,from,toRecipients,ccRecipients,bccRecipients,\
receivedDateTime,sentDateTime,isRead,body,hasAttachments,flag,conversationId,isDraft"
);
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let g: GraphFullMessage = resp.json().await?;
let is_invite = is_event_message(g.odata_type.as_deref());
let event_id = if is_invite {
fetch_event_id(http, base_url, access_token, message_id).await
} else {
None
};
let content_type = match g.body.content_type.to_lowercase().as_str() {
"html" => pidge_core::BodyContentType::Html,
_ => pidge_core::BodyContentType::Text,
};
Ok(pidge_core::FullMessage {
account: account.to_string(),
id: g.id,
conversation_id: g.conversation_id.unwrap_or_default(),
from: g
.from
.map(|w| recipient_from(w.email_address))
.unwrap_or_else(|| pidge_core::MessageFrom {
name: String::new(),
address: String::new(),
}),
to: unwrap_recipients(g.to_recipients),
cc: unwrap_recipients(g.cc_recipients),
bcc: unwrap_recipients(g.bcc_recipients),
subject: clean(g.subject),
received_at: g.received_date_time,
sent_at: g.sent_date_time,
is_read: g.is_read.unwrap_or(true),
body_content_type: content_type,
body_content: clean(Some(g.body.content)),
has_attachments: g.has_attachments.unwrap_or(false),
flag_status: flag_status_from(g.flag),
is_invite,
event_id,
is_draft: g.is_draft.unwrap_or(false),
})
}
async fn fetch_event_id(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Option<String> {
#[derive(Deserialize)]
struct WithEvent {
event: Option<EventRef>,
}
#[derive(Deserialize)]
struct EventRef {
id: String,
}
let url = format!("{base_url}/me/messages/{message_id}");
let req = http.get(&url).bearer_auth(access_token).query(&[
("$select", "id"),
("$expand", "microsoft.graph.eventMessage/event($select=id)"),
]);
let resp = super::send_with_retry(req).await.ok()?;
if !resp.status().is_success() {
return None;
}
let body: WithEvent = resp.json().await.ok()?;
body.event.map(|e| e.id)
}
pub async fn fetch_message_headers(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<Vec<(String, String)>, ClientError> {
let url = format!("{base_url}/me/messages/{message_id}?$select=internetMessageHeaders");
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let body: GraphHeadersResponse = resp.json().await?;
Ok(body
.internet_message_headers
.unwrap_or_default()
.into_iter()
.map(|h| (h.name, h.value))
.collect())
}
#[derive(serde::Deserialize)]
struct GraphHeadersResponse {
#[serde(rename = "internetMessageHeaders", default)]
internet_message_headers: Option<Vec<GraphHeader>>,
}
#[derive(serde::Deserialize)]
struct GraphHeader {
name: String,
value: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum OneClickPolicy {
PublicOnly,
AllowLoopback,
}
pub async fn unsubscribe_one_click(url: &str) -> Result<(), ClientError> {
post_one_click(url, OneClickPolicy::PublicOnly).await
}
pub(crate) async fn post_one_click(url: &str, policy: OneClickPolicy) -> Result<(), ClientError> {
let rejected = || ClientError::UnsubscribeRejected;
let loopback_ok = policy == OneClickPolicy::AllowLoopback;
let parsed = url::Url::parse(url).map_err(|_| rejected())?;
match parsed.scheme() {
"https" => {}
"http" if loopback_ok => {}
_ => return Err(rejected()),
}
let (host, is_name) = match parsed.host() {
Some(url::Host::Domain(name)) => (name.to_string(), true),
Some(url::Host::Ipv4(ip)) if loopback_ok && ip.is_loopback() => (ip.to_string(), false),
_ => return Err(rejected()),
};
let port = parsed.port_or_known_default().ok_or_else(rejected)?;
let addrs: Vec<std::net::SocketAddr> = tokio::net::lookup_host((host.as_str(), port))
.await
.map_err(|_| rejected())?
.collect();
if addrs.is_empty()
|| addrs
.iter()
.any(|a| !one_click_address_allowed(a.ip(), loopback_ok))
{
return Err(rejected());
}
let mut builder = reqwest::Client::builder()
.user_agent(format!("pidge/{}", env!("CARGO_PKG_VERSION")))
.timeout(std::time::Duration::from_secs(10))
.redirect(reqwest::redirect::Policy::none());
if is_name {
builder = builder.resolve_to_addrs(&host, &addrs);
}
let client = builder.build().map_err(|_| rejected())?;
let resp = client
.post(parsed)
.header("Content-Type", "application/x-www-form-urlencoded")
.body("List-Unsubscribe=One-Click")
.send()
.await
.map_err(|_| rejected())?;
if !resp.status().is_success() {
return Err(rejected());
}
Ok(())
}
fn one_click_address_allowed(ip: std::net::IpAddr, loopback_ok: bool) -> bool {
use std::net::IpAddr;
let ip = match ip {
IpAddr::V6(v6) => v6.to_ipv4_mapped().map_or(IpAddr::V6(v6), IpAddr::V4),
v4 => v4,
};
if ip.is_loopback() {
return loopback_ok;
}
match ip {
IpAddr::V4(v4) => {
let [a, b, ..] = v4.octets();
!(v4.is_private()
|| v4.is_link_local()
|| v4.is_unspecified()
|| v4.is_broadcast()
|| v4.is_multicast()
|| a == 0
|| (a == 100 && (b & 0xc0) == 64))
}
IpAddr::V6(v6) => {
let first = v6.segments()[0];
!(v6.is_unspecified()
|| v6.is_multicast()
|| (first & 0xfe00) == 0xfc00
|| (first & 0xffc0) == 0xfe80)
}
}
}
pub async fn list_attachments(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<Vec<pidge_core::Attachment>, ClientError> {
let url = format!(
"{base_url}/me/messages/{message_id}/attachments\
?$select=id,name,contentType,size,isInline"
);
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphAttachmentList = resp.json().await?;
Ok(list
.value
.into_iter()
.filter(|a| {
a.odata_type
.as_deref()
.map(|t| t == "#microsoft.graph.fileAttachment")
.unwrap_or(true)
})
.map(|a| pidge_core::Attachment {
id: a.id,
name: clean(a.name),
content_type: clean(a.content_type),
size_bytes: a.size.unwrap_or(0),
is_inline: a.is_inline.unwrap_or(false),
content_id: a.content_id,
})
.collect())
}
pub async fn get_attachment_bytes(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
attachment_id: &str,
) -> Result<Vec<u8>, ClientError> {
use base64::Engine;
use base64::engine::general_purpose::STANDARD;
let url = format!("{base_url}/me/messages/{message_id}/attachments/{attachment_id}");
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let g: GraphAttachment = resp.json().await?;
let b64 = g.content_bytes.ok_or_else(|| ClientError::Graph {
status: 200,
message: "attachment response missing contentBytes".to_string(),
})?;
STANDARD.decode(&b64).map_err(|e| ClientError::Graph {
status: 200,
message: format!("attachment base64 decode: {e}"),
})
}
pub async fn mark_read(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<(), ClientError> {
patch_message(
http,
base_url,
access_token,
message_id,
&serde_json::json!({ "isRead": true }),
)
.await
}
pub async fn mark_unread(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<(), ClientError> {
patch_message(
http,
base_url,
access_token,
message_id,
&serde_json::json!({ "isRead": false }),
)
.await
}
pub async fn set_flag(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
flagged: bool,
) -> Result<(), ClientError> {
let status = if flagged { "flagged" } else { "notFlagged" };
patch_message(
http,
base_url,
access_token,
message_id,
&serde_json::json!({ "flag": { "flagStatus": status } }),
)
.await
}
#[derive(serde::Deserialize)]
struct GraphCategories {
#[serde(default)]
categories: Vec<String>,
}
pub async fn get_categories(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<Vec<String>, ClientError> {
let url = format!("{base_url}/me/messages/{message_id}?$select=categories");
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let body: GraphCategories = resp.json().await?;
Ok(body.categories)
}
pub async fn set_categories(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
categories: &[String],
) -> Result<(), ClientError> {
patch_message(
http,
base_url,
access_token,
message_id,
&serde_json::json!({ "categories": categories }),
)
.await
}
async fn patch_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
body: &serde_json::Value,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}");
let resp =
super::send_with_retry(http.patch(&url).bearer_auth(access_token).json(body)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn send_mail(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message: &Outgoing,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/sendMail");
let body = serde_json::json!({
"message": message.to_graph_json(),
"saveToSentItems": true,
});
post_no_body(http, &url, access_token, &body).await
}
fn text_to_html(text: &str) -> String {
fn escape(s: &str) -> String {
s.replace('&', "&")
.replace('<', "<")
.replace('>', ">")
.replace('"', """)
}
let normalized = text.replace("\r\n", "\n").replace('\r', "\n");
normalized
.split("\n\n")
.filter(|p| !p.trim().is_empty())
.map(|paragraph| {
let lines: Vec<String> = paragraph
.trim_matches('\n')
.split('\n')
.map(escape)
.collect();
format!("<p>{}</p>", lines.join("<br>"))
})
.collect::<Vec<_>>()
.join("")
}
#[derive(Debug, Deserialize)]
struct GraphBodyOnly {
body: GraphBody,
}
async fn prepend_html_to_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
html_to_prepend: &str,
) -> Result<(), ClientError> {
let get_url = format!("{base_url}/me/messages/{message_id}?$select=body");
let resp = super::send_with_retry(http.get(&get_url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let existing: GraphBodyOnly = resp.json().await?;
let existing_content = existing.body.content;
let new_content = match find_body_tag_end(&existing_content) {
Some(pos) => {
let mut s = String::with_capacity(existing_content.len() + html_to_prepend.len());
s.push_str(&existing_content[..pos]);
s.push_str(html_to_prepend);
s.push_str(&existing_content[pos..]);
s
}
None => format!("{html_to_prepend}{existing_content}"),
};
let patch_url = format!("{base_url}/me/messages/{message_id}");
let patch_body = serde_json::json!({
"body": {
"contentType": "HTML",
"content": new_content,
}
});
let resp = super::send_with_retry(
http.patch(&patch_url)
.bearer_auth(access_token)
.json(&patch_body),
)
.await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
fn find_body_tag_end(html: &str) -> Option<usize> {
let lc = html.to_ascii_lowercase();
let start = lc.find("<body")?;
let after = &html[start..];
let close_rel = after.find('>')?;
Some(start + close_rel + 1)
}
pub async fn reply_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
comment: &str,
) -> Result<(), ClientError> {
let draft_id = create_reply_draft(http, base_url, access_token, message_id, comment).await?;
send_draft(http, base_url, access_token, &draft_id).await
}
pub async fn reply_all_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
comment: &str,
) -> Result<(), ClientError> {
let draft_id =
create_reply_all_draft(http, base_url, access_token, message_id, comment).await?;
send_draft(http, base_url, access_token, &draft_id).await
}
pub async fn forward_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
to: &[String],
comment: &str,
) -> Result<(), ClientError> {
let draft_id =
create_forward_draft(http, base_url, access_token, message_id, to, comment).await?;
send_draft(http, base_url, access_token, &draft_id).await
}
async fn post_no_body(
http: &reqwest::Client,
url: &str,
access_token: &str,
body: &serde_json::Value,
) -> Result<(), ClientError> {
let resp = super::send_with_retry(http.post(url).bearer_auth(access_token).json(body)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn create_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message: &Outgoing,
) -> Result<String, ClientError> {
let url = format!("{base_url}/me/messages");
let body = message.to_graph_json();
let resp =
super::send_with_retry(http.post(&url).bearer_auth(access_token).json(&body)).await?;
parse_id_from_response(resp).await
}
pub async fn create_reply_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
comment: &str,
) -> Result<String, ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/createReply");
let resp = super::send_with_retry(
http.post(&url)
.bearer_auth(access_token)
.json(&serde_json::json!({})),
)
.await?;
let draft_id = parse_id_from_response(resp).await?;
if !comment.is_empty() {
prepend_html_to_draft(
http,
base_url,
access_token,
&draft_id,
&text_to_html(comment),
)
.await?;
}
Ok(draft_id)
}
pub async fn create_reply_all_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
comment: &str,
) -> Result<String, ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/createReplyAll");
let resp = super::send_with_retry(
http.post(&url)
.bearer_auth(access_token)
.json(&serde_json::json!({})),
)
.await?;
let draft_id = parse_id_from_response(resp).await?;
if !comment.is_empty() {
prepend_html_to_draft(
http,
base_url,
access_token,
&draft_id,
&text_to_html(comment),
)
.await?;
}
Ok(draft_id)
}
pub async fn create_forward_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
to: &[String],
comment: &str,
) -> Result<String, ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/createForward");
let resp = super::send_with_retry(http.post(&url).bearer_auth(access_token).json(
&serde_json::json!({
"toRecipients": to.iter().map(|addr| serde_json::json!({
"emailAddress": { "address": addr }
})).collect::<Vec<_>>(),
}),
))
.await?;
let draft_id = parse_id_from_response(resp).await?;
if !comment.is_empty() {
prepend_html_to_draft(
http,
base_url,
access_token,
&draft_id,
&text_to_html(comment),
)
.await?;
}
Ok(draft_id)
}
pub async fn send_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/send");
let resp = super::send_with_retry(
http.post(&url)
.bearer_auth(access_token)
.header(reqwest::header::CONTENT_LENGTH, 0)
.body(reqwest::Body::from("")),
)
.await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn update_draft(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
message: &Outgoing,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}");
let body = message.to_graph_json();
let resp =
super::send_with_retry(http.patch(&url).bearer_auth(access_token).json(&body)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn update_draft_recipients(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
to: Option<&[String]>,
cc: Option<&[String]>,
bcc: Option<&[String]>,
) -> Result<(), ClientError> {
let mut body = serde_json::Map::new();
for (field, list) in [
("toRecipients", to),
("ccRecipients", cc),
("bccRecipients", bcc),
] {
if let Some(list) = list {
let addresses: Vec<_> = list
.iter()
.map(|addr| serde_json::json!({ "emailAddress": { "address": addr } }))
.collect();
body.insert(field.to_string(), addresses.into());
}
}
let url = format!("{base_url}/me/messages/{message_id}");
let resp =
super::send_with_retry(http.patch(&url).bearer_auth(access_token).json(&body)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn delete_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}");
let resp = super::send_with_retry(http.delete(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
pub async fn add_attachment(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
name: &str,
content_type: &str,
bytes: &[u8],
) -> Result<String, ClientError> {
use base64::Engine;
use base64::engine::general_purpose::STANDARD;
let url = format!("{base_url}/me/messages/{message_id}/attachments");
let body = serde_json::json!({
"@odata.type": "#microsoft.graph.fileAttachment",
"name": name,
"contentType": content_type,
"contentBytes": STANDARD.encode(bytes),
});
let resp =
super::send_with_retry(http.post(&url).bearer_auth(access_token).json(&body)).await?;
parse_id_from_response(resp).await
}
pub async fn delete_attachment(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
attachment_id: &str,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/attachments/{attachment_id}");
let resp = super::send_with_retry(http.delete(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
async fn parse_id_from_response(resp: reqwest::Response) -> Result<String, ClientError> {
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let v: serde_json::Value = resp.json().await?;
v["id"]
.as_str()
.map(|s| s.to_string())
.ok_or_else(|| ClientError::Graph {
status: 200,
message: "draft response missing 'id'".to_string(),
})
}
#[derive(Debug, Clone)]
pub struct Outgoing {
pub subject: String,
pub body_text: String,
pub to: Vec<String>,
pub cc: Vec<String>,
pub bcc: Vec<String>,
}
impl Outgoing {
fn to_graph_json(&self) -> serde_json::Value {
fn addresses(list: &[String]) -> Vec<serde_json::Value> {
list.iter()
.map(|addr| serde_json::json!({ "emailAddress": { "address": addr } }))
.collect()
}
serde_json::json!({
"subject": self.subject,
"body": {
"contentType": "Text",
"content": self.body_text,
},
"toRecipients": addresses(&self.to),
"ccRecipients": addresses(&self.cc),
"bccRecipients": addresses(&self.bcc),
})
}
}
pub async fn move_message(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
message_id: &str,
destination: &str,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/messages/{message_id}/move");
let resp = super::send_with_retry(
http.post(&url)
.bearer_auth(access_token)
.json(&serde_json::json!({ "destinationId": destination })),
)
.await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
#[derive(Debug, Clone, Deserialize)]
pub struct MailFolder {
pub id: String,
#[serde(rename = "displayName", deserialize_with = "clean_string")]
pub display_name: String,
#[serde(rename = "totalItemCount", default)]
pub total_item_count: Option<u64>,
#[serde(rename = "unreadItemCount", default)]
pub unread_item_count: Option<u64>,
#[serde(rename = "childFolderCount", default)]
pub child_folder_count: Option<u64>,
}
#[derive(Debug, Deserialize)]
struct GraphFolderList {
value: Vec<MailFolder>,
#[serde(rename = "@odata.nextLink", default)]
next_link: Option<String>,
}
const FOLDER_SELECT: &str = "id,displayName,totalItemCount,unreadItemCount,childFolderCount";
async fn fetch_folder_pages(
http: &reqwest::Client,
access_token: &str,
mut url: String,
) -> Result<Vec<MailFolder>, ClientError> {
let mut folders: Vec<MailFolder> = Vec::new();
loop {
let resp = super::send_with_retry(http.get(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let list: GraphFolderList = resp.json().await?;
folders.extend(list.value);
match list.next_link {
Some(next) => url = next,
None => break,
}
}
Ok(folders)
}
pub async fn list_mail_folders(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
) -> Result<Vec<MailFolder>, ClientError> {
let url = format!("{base_url}/me/mailFolders?$select={FOLDER_SELECT}&$top=100");
fetch_folder_pages(http, access_token, url).await
}
pub async fn list_child_folders(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
parent_id: &str,
) -> Result<Vec<MailFolder>, ClientError> {
let url = format!(
"{base_url}/me/mailFolders/{parent_id}/childFolders?$select={FOLDER_SELECT}&$top=100"
);
fetch_folder_pages(http, access_token, url).await
}
async fn post_folder(
http: &reqwest::Client,
url: &str,
access_token: &str,
display_name: &str,
) -> Result<MailFolder, ClientError> {
let resp = super::send_with_retry(
http.post(url)
.bearer_auth(access_token)
.json(&serde_json::json!({ "displayName": display_name })),
)
.await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
let folder: MailFolder = resp.json().await?;
Ok(folder)
}
pub async fn create_mail_folder(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
display_name: &str,
) -> Result<MailFolder, ClientError> {
let url = format!("{base_url}/me/mailFolders");
post_folder(http, &url, access_token, display_name).await
}
pub async fn create_child_folder(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
parent_id: &str,
display_name: &str,
) -> Result<MailFolder, ClientError> {
let url = format!("{base_url}/me/mailFolders/{parent_id}/childFolders");
post_folder(http, &url, access_token, display_name).await
}
pub async fn delete_mail_folder(
http: &reqwest::Client,
base_url: &str,
access_token: &str,
folder_id: &str,
) -> Result<(), ClientError> {
let url = format!("{base_url}/me/mailFolders/{folder_id}");
let resp = super::send_with_retry(http.delete(&url).bearer_auth(access_token)).await?;
let status = resp.status();
if !status.is_success() {
let text = resp.text().await.unwrap_or_default();
return Err(ClientError::Graph {
status: status.as_u16(),
message: text,
});
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn third_party_text_loses_its_control_characters_at_the_graph_boundary() {
assert_eq!(
clean(Some("Re: \x1b[2K\x1b[1Ahidden\u{9b}x".into())),
"Re: [2K[1Ahiddenx"
);
assert_eq!(
clean(Some("keeps\nlines\tand tabs".into())),
"keeps\nlines\tand tabs"
);
assert_eq!(clean(None), "");
let folder: MailFolder = serde_json::from_value(serde_json::json!({
"id": "F1", "displayName": "Inbox\u{1b}]0;pwned\u{7}"
}))
.unwrap();
assert_eq!(folder.display_name, "Inbox]0;pwned");
let from = recipient_from(GraphFromAddress {
name: Some("Eve\u{1b}[31m".into()),
address: Some("eve@example.com".into()),
});
assert_eq!(from.name, "Eve[31m");
}
use wiremock::matchers::{body_string, header, method, path, path_regex, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
#[test]
fn event_message_rows_are_invites_and_plain_messages_are_not() {
let row = |odata_type: Option<&str>| {
let mut v = serde_json::json!({
"id": "m1",
"receivedDateTime": "2026-09-23T08:00:00Z",
});
if let Some(t) = odata_type {
v["@odata.type"] = t.into();
}
message_from_delta_value(v, "a@example.com").unwrap()
};
assert!(row(Some("#microsoft.graph.eventMessageRequest")).is_invite);
assert!(!row(Some("#microsoft.graph.eventMessage")).is_invite);
assert!(!row(Some("#microsoft.graph.eventMessageResponse")).is_invite);
assert!(!row(Some("#microsoft.graph.message")).is_invite);
assert!(!row(None).is_invite);
}
#[tokio::test]
async fn list_inbox_parses_graph_response() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/inbox/messages"))
.and(header("authorization", "Bearer AT"))
.and(query_param("$top", "5"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{
"id": "AAAA",
"subject": "Hello",
"from": {
"emailAddress": {
"name": "Maria",
"address": "maria@mklab.se"
}
},
"receivedDateTime": "2026-05-13T22:00:00Z",
"isRead": false,
"bodyPreview": "Hi there"
}
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let page = list_inbox(&http, &server.uri(), "AT", "u@e.com", 5, 0, false)
.await
.unwrap();
assert_eq!(page.messages.len(), 1);
assert_eq!(page.messages[0].subject, "Hello");
assert_eq!(page.messages[0].from.address, "maria@mklab.se");
assert!(!page.messages[0].is_read);
assert_eq!(page.messages[0].account, "u@e.com");
assert!(!page.has_more);
}
#[tokio::test]
async fn list_inbox_adds_filter_when_unread_only() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/inbox/messages"))
.and(query_param("$filter", "isRead eq false"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": []
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let page = list_inbox(&http, &server.uri(), "AT", "u@e.com", 5, 0, true)
.await
.unwrap();
assert!(page.messages.is_empty());
}
#[tokio::test]
async fn list_inbox_passes_skip_when_nonzero() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/inbox/messages"))
.and(query_param("$skip", "25"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [],
"@odata.nextLink": "https://graph.example/next"
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let page = list_inbox(&http, &server.uri(), "AT", "u@e.com", 25, 25, false)
.await
.unwrap();
assert!(page.has_more);
}
#[tokio::test]
async fn search_messages_passes_search_query() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/messages"))
.and(query_param("$search", "\"alice budget\""))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{
"id": "S1",
"subject": "Q4 budget review",
"from": { "emailAddress": { "name": "Alice", "address": "alice@example.com" } },
"receivedDateTime": "2026-05-13T22:00:00Z",
"isRead": true,
"bodyPreview": "Numbers attached"
}
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let msgs = search_messages(&http, &server.uri(), "AT", "u@e.com", "alice budget", 25)
.await
.unwrap()
.messages;
assert_eq!(msgs.len(), 1);
assert_eq!(msgs[0].subject, "Q4 budget review");
}
#[tokio::test]
async fn search_folder_messages_searches_within_the_folder() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/sentitems/messages"))
.and(query_param("$search", "\"budget\""))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [{ "id": "S1", "receivedDateTime": "2026-05-13T22:00:00Z" }]
})))
.expect(1)
.mount(&server)
.await;
let http = reqwest::Client::new();
let msgs = search_folder_messages(
&http,
&server.uri(),
"AT",
"u@e.com",
"sentitems",
"budget",
25,
)
.await
.unwrap()
.messages;
assert_eq!(msgs[0].id, "S1");
}
#[tokio::test]
async fn search_escapes_backslashes_before_quotes() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/messages"))
.and(query_param("$search", r#""a\\b \"c\"""#))
.respond_with(
ResponseTemplate::new(200).set_body_json(serde_json::json!({ "value": [] })),
)
.expect(1)
.mount(&server)
.await;
let http = reqwest::Client::new();
search_messages(&http, &server.uri(), "AT", "u@e.com", r#"a\b "c""#, 5)
.await
.unwrap();
}
#[tokio::test]
async fn get_message_fetches_the_event_id_for_a_meeting_request_only() {
let server = MockServer::start().await;
let message = |id: &str, odata_type: &str| {
serde_json::json!({
"@odata.type": odata_type,
"id": id,
"receivedDateTime": "2026-09-23T08:00:00Z",
"sentDateTime": "2026-09-23T08:00:00Z",
"body": { "contentType": "text", "content": "" },
})
};
Mock::given(method("GET"))
.and(path("/me/messages/INV"))
.and(query_param(
"$expand",
"microsoft.graph.eventMessage/event($select=id)",
))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"id": "INV", "event": { "id": "EV1" }
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/messages/INV"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(message("INV", "#microsoft.graph.eventMessageRequest")),
)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/messages/PLAIN"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(message("PLAIN", "#microsoft.graph.message")),
)
.expect(1)
.mount(&server)
.await;
let http = reqwest::Client::new();
let inv = get_message(&http, &server.uri(), "AT", "u@e.com", "INV")
.await
.unwrap();
assert!(inv.is_invite);
assert_eq!(inv.event_id.as_deref(), Some("EV1"));
let plain = get_message(&http, &server.uri(), "AT", "u@e.com", "PLAIN")
.await
.unwrap();
assert!(!plain.is_invite);
assert_eq!(plain.event_id, None);
}
#[tokio::test]
async fn mark_unread_patches_isread_false() {
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path_regex("/me/messages/[A-Za-z0-9]+"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({})))
.mount(&server)
.await;
let http = reqwest::Client::new();
mark_unread(&http, &server.uri(), "AT", "MSG")
.await
.unwrap();
}
#[tokio::test]
async fn set_flag_patches_flag_status() {
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path_regex("/me/messages/[A-Za-z0-9]+"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({})))
.mount(&server)
.await;
let http = reqwest::Client::new();
set_flag(&http, &server.uri(), "AT", "MSG", true)
.await
.unwrap();
set_flag(&http, &server.uri(), "AT", "MSG", false)
.await
.unwrap();
}
#[tokio::test]
async fn send_mail_wraps_outgoing_in_message_envelope() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/me/sendMail"))
.and(body_partial_json(serde_json::json!({
"saveToSentItems": true,
"message": {
"subject": "Hello",
"body": { "contentType": "Text", "content": "Hi there" },
"toRecipients": [{ "emailAddress": { "address": "alice@example.com" } }]
}
})))
.respond_with(ResponseTemplate::new(202))
.mount(&server)
.await;
let http = reqwest::Client::new();
let msg = Outgoing {
subject: "Hello".into(),
body_text: "Hi there".into(),
to: vec!["alice@example.com".into()],
cc: vec![],
bcc: vec![],
};
send_mail(&http, &server.uri(), "AT", &msg).await.unwrap();
}
#[tokio::test]
async fn reply_message_creates_draft_patches_body_and_sends() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path_regex("/me/messages/MSG/createReply"))
.respond_with(
ResponseTemplate::new(201).set_body_json(serde_json::json!({ "id": "DRAFT" })),
)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/messages/DRAFT"))
.and(query_param("$select", "body"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"body": { "contentType": "HTML", "content": "<html><body><div></div></body></html>" }
})))
.mount(&server)
.await;
Mock::given(method("PATCH"))
.and(path("/me/messages/DRAFT"))
.and(body_partial_json(serde_json::json!({
"body": { "contentType": "HTML" }
})))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/me/messages/DRAFT/send"))
.respond_with(ResponseTemplate::new(202))
.mount(&server)
.await;
let http = reqwest::Client::new();
reply_message(&http, &server.uri(), "AT", "MSG", "Thanks!")
.await
.unwrap();
}
#[tokio::test]
async fn forward_message_creates_draft_with_recipients_and_sends() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path_regex("/me/messages/MSG/createForward"))
.and(body_partial_json(serde_json::json!({
"toRecipients": [{ "emailAddress": { "address": "bob@example.com" } }]
})))
.respond_with(
ResponseTemplate::new(201).set_body_json(serde_json::json!({ "id": "DRAFT" })),
)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/messages/DRAFT"))
.and(query_param("$select", "body"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"body": { "contentType": "HTML", "content": "<html><body></body></html>" }
})))
.mount(&server)
.await;
Mock::given(method("PATCH"))
.and(path("/me/messages/DRAFT"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/me/messages/DRAFT/send"))
.respond_with(ResponseTemplate::new(202))
.mount(&server)
.await;
let http = reqwest::Client::new();
forward_message(
&http,
&server.uri(),
"AT",
"MSG",
&["bob@example.com".into()],
"FYI",
)
.await
.unwrap();
}
#[test]
fn text_to_html_escapes_and_breaks_paragraphs() {
let out = text_to_html("Hej Edward,\n\nLine 1\nLine 2\n\n<script>x</script>");
assert!(out.contains("<p>Hej Edward,</p>"));
assert!(out.contains("<p>Line 1<br>Line 2</p>"));
assert!(out.contains("<script>x</script>"));
assert!(!out.contains("<script>"));
}
#[test]
fn text_to_html_normalizes_crlf() {
let out = text_to_html("A\r\nB\r\n\r\nC");
assert!(out.contains("<p>A<br>B</p>"));
assert!(out.contains("<p>C</p>"));
}
#[test]
fn find_body_tag_end_handles_attributes_and_case() {
let html = "<html><BODY class=\"x\">content</BODY></html>";
let pos = find_body_tag_end(html).unwrap();
assert_eq!(&html[pos..pos + 7], "content");
}
#[tokio::test]
async fn list_mail_folders_parses_and_pages() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders"))
.and(query_param("$top", "100"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{ "id": "F1", "displayName": "Biljetter", "totalItemCount": 3, "unreadItemCount": 0 }
],
"@odata.nextLink": format!("{}/me/mailFolders?page=2", server.uri())
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/mailFolders"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{ "id": "F2", "displayName": "Kvitton", "totalItemCount": 7, "unreadItemCount": 2 }
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let folders = list_mail_folders(&http, &server.uri(), "AT").await.unwrap();
assert_eq!(folders.len(), 2);
assert_eq!(folders[0].display_name, "Biljetter");
assert_eq!(folders[0].total_item_count, Some(3));
assert_eq!(folders[1].display_name, "Kvitton");
assert_eq!(folders[1].unread_item_count, Some(2));
}
#[tokio::test]
async fn create_mail_folder_posts_display_name() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/me/mailFolders"))
.and(body_partial_json(
serde_json::json!({ "displayName": "Biljetter" }),
))
.respond_with(
ResponseTemplate::new(201)
.set_body_json(serde_json::json!({ "id": "NEWF", "displayName": "Biljetter" })),
)
.mount(&server)
.await;
let http = reqwest::Client::new();
let folder = create_mail_folder(&http, &server.uri(), "AT", "Biljetter")
.await
.unwrap();
assert_eq!(folder.id, "NEWF");
assert_eq!(folder.display_name, "Biljetter");
}
#[tokio::test]
async fn list_child_folders_hits_child_endpoint() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/PARENT/childFolders"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{ "id": "C1", "displayName": "MKLab", "totalItemCount": 5, "unreadItemCount": 0, "childFolderCount": 0 }
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let children = list_child_folders(&http, &server.uri(), "AT", "PARENT")
.await
.unwrap();
assert_eq!(children.len(), 1);
assert_eq!(children[0].display_name, "MKLab");
}
#[tokio::test]
async fn create_child_folder_posts_to_parent() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/me/mailFolders/PARENT/childFolders"))
.and(body_partial_json(
serde_json::json!({ "displayName": "MKLab" }),
))
.respond_with(
ResponseTemplate::new(201)
.set_body_json(serde_json::json!({ "id": "C9", "displayName": "MKLab" })),
)
.mount(&server)
.await;
let http = reqwest::Client::new();
let folder = create_child_folder(&http, &server.uri(), "AT", "PARENT", "MKLab")
.await
.unwrap();
assert_eq!(folder.id, "C9");
}
#[tokio::test]
async fn delete_mail_folder_hits_delete_endpoint() {
let server = MockServer::start().await;
Mock::given(method("DELETE"))
.and(path("/me/mailFolders/F1"))
.respond_with(ResponseTemplate::new(204))
.mount(&server)
.await;
let http = reqwest::Client::new();
delete_mail_folder(&http, &server.uri(), "AT", "F1")
.await
.unwrap();
}
#[tokio::test]
async fn get_categories_parses_field() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_regex(r"^/me/messages/.+$"))
.and(query_param("$select", "categories"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"categories": ["Receipts", "Urgent"]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let cats = get_categories(&http, &server.uri(), "AT", "MSG")
.await
.unwrap();
assert_eq!(cats, vec!["Receipts", "Urgent"]);
}
#[tokio::test]
async fn set_categories_patches_array() {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path_regex(r"^/me/messages/.+$"))
.and(body_partial_json(
serde_json::json!({ "categories": ["receipt", "ticket"] }),
))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let http = reqwest::Client::new();
set_categories(
&http,
&server.uri(),
"AT",
"MSG",
&["receipt".to_string(), "ticket".to_string()],
)
.await
.unwrap();
}
#[tokio::test]
async fn move_message_posts_destination() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path_regex("/me/messages/[A-Za-z0-9]+/move"))
.respond_with(
ResponseTemplate::new(201).set_body_json(serde_json::json!({"id": "NEW"})),
)
.mount(&server)
.await;
let http = reqwest::Client::new();
move_message(&http, &server.uri(), "AT", "MSG", "archive")
.await
.unwrap();
}
#[tokio::test]
async fn get_message_parses_graph_response() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_regex("/me/messages/[A-Za-z0-9]+"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"id": "AAA",
"subject": "Hello",
"from": { "emailAddress": { "name": "Maria", "address": "maria@mklab.se" } },
"toRecipients": [
{ "emailAddress": { "name": "Kristofer", "address": "kristofer@mklab.se" } }
],
"ccRecipients": [],
"bccRecipients": [],
"receivedDateTime": "2026-05-14T22:00:00Z",
"sentDateTime": "2026-05-14T21:59:30Z",
"isRead": false,
"body": { "contentType": "html", "content": "<p>Hi</p>" },
"hasAttachments": true
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let m = get_message(&http, &server.uri(), "AT", "u@e.com", "AAA")
.await
.unwrap();
assert_eq!(m.id, "AAA");
assert_eq!(m.subject, "Hello");
assert_eq!(m.from.name, "Maria");
assert_eq!(m.to.len(), 1);
assert_eq!(m.to[0].address, "kristofer@mklab.se");
assert!(matches!(
m.body_content_type,
pidge_core::BodyContentType::Html
));
assert_eq!(m.body_content, "<p>Hi</p>");
assert!(m.has_attachments);
}
#[tokio::test]
async fn get_message_selects_and_maps_the_draft_flag() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/me/messages/D1"))
.respond_with(move |req: &wiremock::Request| {
let select = req
.url
.query_pairs()
.find(|(k, _)| k == "$select")
.map(|(_, v)| v.into_owned())
.unwrap_or_default();
assert!(select.split(',').any(|f| f == "isDraft"), "{select}");
ResponseTemplate::new(200).set_body_json(serde_json::json!({
"id": "D1",
"receivedDateTime": "2026-05-14T22:00:00Z",
"sentDateTime": "2026-05-14T21:59:30Z",
"body": { "contentType": "text", "content": "" },
"isDraft": true
}))
})
.mount(&server)
.await;
let http = reqwest::Client::new();
let m = get_message(&http, &server.uri(), "AT", "u@e.com", "D1")
.await
.unwrap();
assert!(m.is_draft);
}
#[tokio::test]
async fn update_draft_recipients_patches_only_the_given_lists() {
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path("/me/messages/D1"))
.and(body_string(
serde_json::json!({
"ccRecipients": [{ "emailAddress": { "address": "cc@example.com" } }]
})
.to_string(),
))
.respond_with(ResponseTemplate::new(200))
.expect(1)
.mount(&server)
.await;
let http = reqwest::Client::new();
update_draft_recipients(
&http,
&server.uri(),
"AT",
"D1",
None,
Some(&["cc@example.com".to_string()]),
None,
)
.await
.unwrap();
}
#[tokio::test]
async fn list_attachments_filters_file_attachments() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_regex("/me/messages/[A-Za-z0-9]+/attachments"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [
{
"@odata.type": "#microsoft.graph.fileAttachment",
"id": "att-1",
"name": "report.pdf",
"contentType": "application/pdf",
"size": 12345,
"isInline": false
},
{
"@odata.type": "#microsoft.graph.itemAttachment",
"id": "att-2",
"name": "an-email.eml",
"contentType": "message/rfc822",
"size": 7777,
"isInline": false
}
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let atts = list_attachments(&http, &server.uri(), "AT", "MSG")
.await
.unwrap();
assert_eq!(atts.len(), 1);
assert_eq!(atts[0].name, "report.pdf");
assert_eq!(atts[0].size_bytes, 12345);
}
#[tokio::test]
async fn get_attachment_bytes_decodes_base64() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_regex(
"/me/messages/[A-Za-z0-9]+/attachments/[A-Za-z0-9-]+",
))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"id": "att-1",
"name": "report.pdf",
"contentType": "application/pdf",
"size": 5,
"isInline": false,
"contentBytes": "aGVsbG8="
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let bytes = get_attachment_bytes(&http, &server.uri(), "AT", "MSG", "att-1")
.await
.unwrap();
assert_eq!(bytes, b"hello");
}
#[tokio::test]
async fn mark_read_patches_isread_true() {
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path_regex("/me/messages/[A-Za-z0-9]+"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({})))
.mount(&server)
.await;
let http = reqwest::Client::new();
mark_read(&http, &server.uri(), "AT", "MSG").await.unwrap();
}
#[tokio::test]
async fn fetch_message_headers_parses_array() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_regex(r"^/me/messages/.+$"))
.and(query_param("$select", "internetMessageHeaders"))
.and(header("authorization", "Bearer AT"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"internetMessageHeaders": [
{ "name": "List-Unsubscribe", "value": "<mailto:u@x>, <https://x/u>" },
{ "name": "List-Unsubscribe-Post", "value": "List-Unsubscribe=One-Click" }
]
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let headers = fetch_message_headers(&http, &server.uri(), "AT", "MSGID")
.await
.unwrap();
assert_eq!(headers.len(), 2);
assert_eq!(headers[0].0, "List-Unsubscribe");
assert_eq!(headers[0].1, "<mailto:u@x>, <https://x/u>");
assert_eq!(headers[1].0, "List-Unsubscribe-Post");
}
#[tokio::test]
async fn unsubscribe_one_click_posts_exact_form_body() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/u"))
.and(body_string("List-Unsubscribe=One-Click"))
.respond_with(ResponseTemplate::new(200))
.expect(1)
.mount(&server)
.await;
post_one_click(
&format!("{}/u", server.uri()),
OneClickPolicy::AllowLoopback,
)
.await
.unwrap();
}
#[tokio::test]
async fn unsubscribe_one_click_reports_a_500_without_its_status_or_body() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/u"))
.respond_with(ResponseTemplate::new(500).set_body_string("nope"))
.mount(&server)
.await;
let err = post_one_click(
&format!("{}/u", server.uri()),
OneClickPolicy::AllowLoopback,
)
.await
.unwrap_err();
assert!(matches!(err, ClientError::UnsubscribeRejected), "{err:?}");
assert!(!err.to_string().contains("500"));
assert!(!err.to_string().contains("nope"));
}
#[tokio::test]
async fn unsubscribe_one_click_refuses_http_without_a_request() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.expect(0)
.mount(&server)
.await;
let port = server.address().port();
for url in [
format!("http://localhost:{port}/u"),
format!("{}/u", server.uri()),
] {
let err = unsubscribe_one_click(&url).await.unwrap_err();
assert!(matches!(err, ClientError::UnsubscribeRejected), "{url}");
}
}
#[tokio::test]
async fn unsubscribe_one_click_refuses_ip_literals_and_internal_hosts() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.expect(0)
.mount(&server)
.await;
let port = server.address().port();
for url in [
format!("https://127.0.0.1:{port}/u"),
"https://10.0.0.1/u".to_string(),
"https://[::1]/u".to_string(),
"https://93.184.215.14/u".to_string(),
format!("https://localhost:{port}/u"),
"ftp://example.com/u".to_string(),
"not a url".to_string(),
] {
let err = unsubscribe_one_click(&url).await.unwrap_err();
assert!(matches!(err, ClientError::UnsubscribeRejected), "{url}");
}
let err = post_one_click("http://10.0.0.1/u", OneClickPolicy::AllowLoopback)
.await
.unwrap_err();
assert!(matches!(err, ClientError::UnsubscribeRejected));
}
#[tokio::test]
async fn unsubscribe_one_click_does_not_follow_a_redirect() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/u"))
.respond_with(ResponseTemplate::new(302).insert_header("Location", "/internal"))
.expect(1)
.mount(&server)
.await;
Mock::given(path("/internal"))
.respond_with(ResponseTemplate::new(200))
.expect(0)
.mount(&server)
.await;
let err = post_one_click(
&format!("{}/u", server.uri()),
OneClickPolicy::AllowLoopback,
)
.await
.unwrap_err();
assert!(matches!(err, ClientError::UnsubscribeRejected), "{err:?}");
}
#[test]
fn one_click_address_policy() {
use std::net::IpAddr;
let allowed = |s: &str, lo| one_click_address_allowed(s.parse::<IpAddr>().unwrap(), lo);
for bad in [
"10.1.2.3",
"172.16.0.1",
"172.31.255.255",
"192.168.1.1",
"169.254.169.254",
"0.0.0.0",
"100.64.0.1",
"::",
"fe80::1",
"fc00::1",
"fd12:3456::1",
"::ffff:10.0.0.1",
"::ffff:127.0.0.1",
] {
assert!(!allowed(bad, false), "{bad}");
}
assert!(!allowed("127.0.0.1", false));
assert!(!allowed("::1", false));
assert!(allowed("127.0.0.1", true));
assert!(!allowed("10.0.0.1", true));
for good in ["93.184.215.14", "172.32.0.1", "2606:2800:220:1::1"] {
assert!(allowed(good, false), "{good}");
}
}
}
#[cfg(test)]
mod cursor_paging_tests {
use super::*;
use wiremock::matchers::{method, path, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
fn msg(id: &str, received: &str) -> serde_json::Value {
serde_json::json!({
"id": id,
"subject": format!("s-{id}"),
"from": {"emailAddress": {"name": "N", "address": "n@x.se"}},
"receivedDateTime": received,
"isRead": true,
"bodyPreview": "p",
"hasAttachments": false
})
}
#[tokio::test]
async fn next_link_is_surfaced_and_followable() {
let server = MockServer::start().await;
let page2_url = format!("{}/me/mailFolders/inbox/messages?page=2", server.uri());
Mock::given(method("GET"))
.and(path("/me/mailFolders/inbox/messages"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [msg("m3", "2026-07-01T10:00:00Z")]
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/me/mailFolders/inbox/messages"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"value": [msg("m1", "2026-07-03T10:00:00Z"), msg("m2", "2026-07-02T10:00:00Z")],
"@odata.nextLink": page2_url
})))
.mount(&server)
.await;
let http = reqwest::Client::new();
let page1 = list_inbox(&http, &server.uri(), "tok", "a@b.se", 2, 0, false)
.await
.unwrap();
assert_eq!(page1.messages.len(), 2);
let next = page1.next_link.expect("first page links onward");
let page2 = list_messages_at(&http, &server.uri(), "tok", "a@b.se", &next)
.await
.unwrap();
assert_eq!(page2.messages.len(), 1);
assert_eq!(page2.messages[0].id, "m3");
assert!(page2.next_link.is_none(), "final page has no continuation");
}
}