use std::sync::Arc;
use fraiseql_functions::IngestError;
use futures::{StreamExt, future::BoxFuture};
use tokio::{
io::{AsyncRead, AsyncWrite},
net::TcpStream,
};
use tokio_rustls::{
TlsConnector,
rustls::{ClientConfig, RootCertStore, crypto::ring, pki_types::ServerName},
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchedMessage {
pub uid: u32,
pub raw: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchBatch {
pub uid_validity: u32,
pub messages: Vec<FetchedMessage>,
}
pub trait MailboxFetcher: Send + Sync {
fn fetch(
&self,
stored: Option<super::cursor::Cursor>,
batch_size: u32,
) -> BoxFuture<'_, Result<FetchBatch, IngestError>>;
}
pub struct ImapMailboxFetcher {
host: String,
port: u16,
username: String,
password: String,
mailbox: String,
connector: TlsConnector,
}
impl ImapMailboxFetcher {
pub fn new(
host: impl Into<String>,
port: u16,
username: impl Into<String>,
password: impl Into<String>,
mailbox: impl Into<String>,
) -> Result<Self, IngestError> {
Ok(Self {
host: host.into(),
port,
username: username.into(),
password: password.into(),
mailbox: mailbox.into(),
connector: tls_connector()?,
})
}
async fn fetch_batch(
&self,
stored: Option<super::cursor::Cursor>,
batch_size: u32,
) -> Result<FetchBatch, IngestError> {
let tcp = TcpStream::connect((self.host.as_str(), self.port))
.await
.map_err(|error| IngestError::new(format!("imap connect {}: {error}", self.host)))?;
let server_name = ServerName::try_from(self.host.clone())
.map_err(|error| IngestError::new(format!("imap server name: {error}")))?;
let tls = self
.connector
.connect(server_name, tcp)
.await
.map_err(|error| IngestError::new(format!("imap TLS handshake: {error}")))?;
let client = async_imap::Client::new(tls);
let mut session = client
.login(&self.username, &self.password)
.await
.map_err(|(error, _client)| IngestError::new(format!("imap login: {error}")))?;
let batch = self.fetch_selected(&mut session, stored, batch_size).await;
drop(session);
batch
}
async fn fetch_selected<T>(
&self,
session: &mut async_imap::Session<T>,
stored: Option<super::cursor::Cursor>,
batch_size: u32,
) -> Result<FetchBatch, IngestError>
where
T: AsyncRead + AsyncWrite + Unpin + Send + std::fmt::Debug,
{
let mailbox = session
.select(&self.mailbox)
.await
.map_err(|error| IngestError::new(format!("imap SELECT {}: {error}", self.mailbox)))?;
let uid_validity = mailbox.uid_validity.ok_or_else(|| {
IngestError::new("imap SELECT did not report UIDVALIDITY; cannot cursor safely")
})?;
let fetch_start = super::cursor::fetch_start(stored, uid_validity);
let query = format!("{fetch_start}:*");
let mut stream = session
.uid_fetch(query, "(UID BODY.PEEK[])")
.await
.map_err(|error| IngestError::new(format!("imap UID FETCH: {error}")))?;
let mut messages = Vec::new();
while let Some(item) = stream.next().await {
let fetch =
item.map_err(|error| IngestError::new(format!("imap FETCH item: {error}")))?;
let (Some(uid), Some(body)) = (fetch.uid, fetch.body()) else {
continue;
};
messages.push(FetchedMessage {
uid,
raw: body.to_vec(),
});
if messages.len() >= batch_size as usize {
break;
}
}
drop(stream);
Ok(FetchBatch {
uid_validity,
messages,
})
}
}
impl MailboxFetcher for ImapMailboxFetcher {
fn fetch(
&self,
stored: Option<super::cursor::Cursor>,
batch_size: u32,
) -> BoxFuture<'_, Result<FetchBatch, IngestError>> {
Box::pin(self.fetch_batch(stored, batch_size))
}
}
fn tls_connector() -> Result<TlsConnector, IngestError> {
let mut roots = RootCertStore::empty();
roots.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
let config = ClientConfig::builder_with_provider(Arc::new(ring::default_provider()))
.with_safe_default_protocol_versions()
.map_err(|error| IngestError::new(format!("imap rustls config: {error}")))?
.with_root_certificates(roots)
.with_no_client_auth();
Ok(TlsConnector::from(Arc::new(config)))
}
#[cfg(test)]
mod tests;