use std::{
collections::HashSet,
net::{Shutdown, TcpStream},
sync::Arc,
time::Duration,
};
use async_trait::async_trait;
use imap::{
types::{Flag, Uid, UnsolicitedResponse},
Client, ClientBuilder, Session,
};
use mail_parser::{Addr, HeaderValue, Message as EmailParser, MimeHeaders};
use native_tls::{TlsConnector, TlsStream};
use tokio::{
runtime::Handle,
sync::{mpsc, Notify},
};
use crate::{
message::{Kind, Message, Part},
transport::{Receiver, Transport},
};
const IDLE_REFRESH: Duration = Duration::from_secs(2 * 60);
const UNPROCESSABLE_KEYWORD: &str = "mailfred-unprocessable";
#[derive(Clone)]
pub struct Imap {
pub domain: String,
pub port: u16,
pub user: String,
pub password: String,
pub folder: String,
}
#[async_trait]
impl Transport for Imap {
const NAME: &'static str = "imap";
type Connection = ImapConnection;
type Error = imap::Error;
async fn connect(&self) -> imap::Result<ImapConnection> {
let (session, tcp) = tokio::task::block_in_place(move || -> imap::Result<_> {
let tcp = TcpStream::connect((self.domain.as_str(), self.port))?;
let tcp_handle = tcp.try_clone()?;
let tls = TlsConnector::builder()
.build()?
.connect(&self.domain, tcp)?;
let mut client = Client::new(tls);
client.read_greeting()?;
let mut session = client
.login(&self.user, &self.password)
.map_err(|(e, _)| e)?;
session.select(&self.folder)?;
Ok((session, tcp_handle))
})?;
let ready_to_recv = Arc::new(Notify::new());
let (tx, rx) = mpsc::channel(1);
tokio::task::spawn_blocking({
let ready_to_recv = ready_to_recv.clone();
move || {
if let Err(err) = listener(session, ready_to_recv, tx.clone()) {
tx.blocking_send(Err(err)).ok();
}
}
});
Ok(ImapConnection {
rx,
tcp,
ready_to_recv,
})
}
}
fn listener(
mut session: Session<TlsStream<TcpStream>>,
ready_to_recv: Arc<Notify>,
tx: mpsc::Sender<imap::Result<Message>>,
) -> imap::Result<()> {
let mut unprocessable = HashSet::new();
let unprocessable_flag = Flag::Custom(UNPROCESSABLE_KEYWORD.into());
loop {
let fetches = session.uid_fetch("1:*", "(UID FLAGS)")?;
let mut expunge_required = fetches
.iter()
.any(|fetch| fetch.flags().contains(&Flag::Deleted));
let pending = fetches
.iter()
.filter(|fetch| {
let flags = fetch.flags();
!flags.contains(&Flag::Deleted) && !flags.contains(&unprocessable_flag)
})
.filter_map(|fetch| fetch.uid)
.filter(|uid| !unprocessable.contains(uid))
.collect::<Vec<_>>();
for uid in &pending {
let uid = *uid;
let Some(msg) = fetch_email(&mut session, uid)? else {
log::warn!(
"imap: message with uid {} can not be read, marking it to \
not be processed again",
uid
);
let stored = session.uid_store(
uid.to_string(),
format!("+FLAGS ({})", UNPROCESSABLE_KEYWORD),
);
if let Err(err) = stored {
log::debug!("imap: the folder does not accept keywords: {}", err);
}
unprocessable.insert(uid);
continue;
};
let ready_to_recv = ready_to_recv.clone();
Handle::current().block_on(async move { ready_to_recv.notified().await });
if tx.blocking_send(Ok(msg)).is_err() {
return Ok(());
}
session.uid_store(uid.to_string(), "+FLAGS (\\Deleted)")?;
expunge_required = true;
}
if expunge_required {
session.expunge()?;
}
if pending.is_empty() {
session
.idle()
.timeout(IDLE_REFRESH)
.wait_while(|response| {
!matches!(
response,
UnsolicitedResponse::Exists(_) | UnsolicitedResponse::Recent(_)
)
})?;
}
}
}
fn fetch_email(
session: &mut Session<TlsStream<TcpStream>>,
uid: Uid,
) -> imap::Result<Option<Message>> {
let fetches = session.uid_fetch(uid.to_string(), "RFC822")?;
Ok(fetches
.iter()
.find_map(|fetch| fetch.body())
.and_then(read_email))
}
fn read_address(header: &HeaderValue) -> Option<String> {
fn first(addrs: &[Addr]) -> Option<String> {
addrs
.iter()
.find_map(|addr| addr.address.as_deref())
.map(Into::into)
}
match header {
HeaderValue::Address(addr) => addr.address.as_deref().map(Into::into),
HeaderValue::AddressList(addrs) => first(addrs),
HeaderValue::Group(group) => first(&group.addresses),
HeaderValue::GroupList(groups) => groups.iter().find_map(|group| first(&group.addresses)),
_ => None,
}
}
fn read_email(email_raw: &[u8]) -> Option<Message> {
let email = EmailParser::parse(email_raw)?;
let subject = email.subject().unwrap_or_default().into();
let from = read_address(email.from())?;
let mut body = Vec::default();
for part in email.text_bodies() {
body.push(Part {
kind: if part.is_text_html() {
Kind::Html
} else {
Kind::Text
},
content: part.contents().into(),
});
}
for part in email.attachments() {
if !part.is_empty() {
body.push(Part {
kind: Kind::Attachment(part.attachment_name().unwrap_or_default().into()),
content: part.contents().into(),
});
}
}
Some(Message {
address: from,
header: subject,
body,
})
}
impl Imap {
pub fn clear_folder(&self, folder: &str) -> imap::Result<()> {
let client = ClientBuilder::new(&self.domain, self.port).connect()?;
let mut session = client.login(&self.user, &self.password).map_err(|e| e.0)?;
session.select(folder)?;
session.store("1:*", "+FLAGS (\\Deleted)")?;
session.expunge()?;
Ok(())
}
}
pub struct ImapConnection {
rx: mpsc::Receiver<imap::Result<Message>>,
tcp: TcpStream,
ready_to_recv: Arc<Notify>,
}
#[async_trait]
impl Receiver for ImapConnection {
type Error = imap::Error;
async fn recv(&mut self) -> imap::Result<Message> {
self.ready_to_recv.notify_one();
match self.rx.recv().await {
Some(message) => message,
None => unreachable!(),
}
}
}
impl Drop for ImapConnection {
fn drop(&mut self) {
self.tcp.shutdown(Shutdown::Both).ok();
}
}
#[cfg(test)]
mod tests {
use super::*;
fn email(headers: &str) -> Option<Message> {
read_email(format!("{headers}\r\n\r\nbody\r\n").as_bytes())
}
fn remitter(headers: &str) -> Option<String> {
email(headers).map(|msg| msg.address)
}
#[test]
fn remitter_of_a_single_address() {
assert_eq!(remitter("From: a@b.com").as_deref(), Some("a@b.com"));
assert_eq!(remitter("From: Bob <a@b.com>").as_deref(), Some("a@b.com"));
}
#[test]
fn remitter_of_a_list_or_a_group() {
assert_eq!(
remitter("From: a@b.com, c@d.com").as_deref(),
Some("a@b.com")
);
assert_eq!(
remitter("From: Team: a@b.com, c@d.com;").as_deref(),
Some("a@b.com")
);
}
#[test]
fn message_without_a_usable_remitter_is_discarded() {
assert_eq!(remitter("Subject: no from header"), None);
assert_eq!(remitter("From: "), None);
assert_eq!(remitter("From: undisclosed-recipients:;"), None);
}
#[test]
fn subject_is_read_as_the_header() {
let msg = email("From: a@b.com\r\nSubject: Count").unwrap();
assert_eq!(msg.header, "Count");
let msg = email("From: a@b.com").unwrap();
assert_eq!(msg.header, "");
}
}