use std::{env, error::Error, io, sync::Arc};
use io_imap::{
client::{ImapClientAsync, ImapClientError},
codec::fragmentizer::Fragmentizer,
coroutine::*,
session::*,
types::{
mailbox::{ListMailbox, Mailbox},
response::Capability,
},
};
use io_sasl::{mechanism::Sasl, rfc4616::plain::SaslPlainCreds};
use rustls::{ClientConfig, pki_types::ServerName};
use rustls_platform_verifier::ConfigVerifierExt;
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::{TcpStream, UnixStream},
};
use tokio_rustls::{TlsConnector, client::TlsStream};
use url::Url;
const FRAGMENTIZER_MAX_MESSAGE_SIZE: u32 = 100 * 1024 * 1024;
const READ_BUFFER_SIZE: usize = 16 * 1024;
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
env_logger::init();
let url = Url::parse(&env::var("URL")?)?;
let starttls = env::var("STARTTLS").is_ok();
rustls::crypto::ring::default_provider()
.install_default()
.ok();
let config = ClientConfig::with_platform_verifier()?;
let connector = TlsConnector::from(Arc::new(config));
let sasl = env::var("LOGIN")
.ok()
.zip(env::var("PASSWORD").ok())
.map(|(authcid, passwd)| SaslPlainCreds {
authzid: None,
authcid,
passwd: passwd.into(),
});
let opts = ImapSessionOpenOptions {
starttls,
..Default::default()
};
let (mut client, capability) = ImapClientTokio::connect(&url, &connector, sasl, opts).await?;
for capability in &capability {
println!("{capability:?}");
}
let reference = Mailbox::try_from(String::new())?;
let pattern = ListMailbox::try_from("*")?;
let listing = tokio::spawn(async move { client.list(reference, pattern).await }).await??;
for (mailbox, _delimiter, _attributes) in listing {
println!("{mailbox:?}");
}
Ok(())
}
struct ImapClientTokio {
stream: TokioStream,
fragmentizer: Fragmentizer,
}
impl ImapClientTokio {
async fn connect(
url: &Url,
connector: &TlsConnector,
sasl: Option<impl Into<Sasl>>,
opts: ImapSessionOpenOptions,
) -> Result<(Self, Vec<Capability<'static>>), Box<dyn Error>> {
let transport = ImapSessionTransport::from_url(url)?;
let mut session = ImapSessionOpen::new(transport, sasl, opts);
let mut fragmentizer = Fragmentizer::new(FRAGMENTIZER_MAX_MESSAGE_SIZE);
let mut stream: Option<TokioStream> = None;
let mut buf = [0u8; READ_BUFFER_SIZE];
let mut arg: Option<&[u8]> = None;
let mut tls_host = String::new();
let missing = || String::from("IMAP session yielded I/O before connecting");
loop {
match session.resume(&mut fragmentizer, arg.take()) {
ImapCoroutineState::Complete(Err(err)) => return Err(err.into()),
ImapCoroutineState::Complete(Ok(data)) => {
let client = Self {
stream: stream.ok_or_else(missing)?,
fragmentizer,
};
return Ok((client, data.capability));
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsTcpConnect {
host,
port,
}) => {
let sock = TcpStream::connect((host.as_str(), port)).await?;
tls_host = host;
stream = Some(TokioStream::Tcp(sock));
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsTlsConnect {
host,
port,
}) => {
let name = ServerName::try_from(host.as_str())?.to_owned();
let sock = TcpStream::connect((host.as_str(), port)).await?;
stream = Some(TokioStream::Tls(Box::new(
connector.connect(name, sock).await?,
)));
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsUnixConnect(path)) => {
let sock = UnixStream::connect(path).await?;
stream = Some(TokioStream::Unix(sock));
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsTlsUpgrade) => {
let plain = stream.take().ok_or_else(missing)?;
stream = Some(plain.upgrade_tls(connector, &tls_host).await?);
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsRead) => {
let n = stream.as_mut().ok_or_else(missing)?.read(&mut buf).await?;
arg = Some(&buf[..n]);
}
ImapCoroutineState::Yielded(ImapSessionOpenYield::WantsWrite(bytes)) => {
stream
.as_mut()
.ok_or_else(missing)?
.write_all(&bytes)
.await?;
}
}
}
}
}
impl ImapClientAsync for ImapClientTokio {
#[allow(clippy::manual_async_fn)]
fn run<C, T, E>(
&mut self,
mut coroutine: C,
) -> impl Future<Output = Result<T, ImapClientError>> + Send
where
C: ImapCoroutine<Yield = ImapYield, Return = Result<T, E>> + Send,
T: Send,
E: Send,
ImapClientError: From<E>,
{
async move {
let mut buf = [0u8; READ_BUFFER_SIZE];
let mut arg: Option<&[u8]> = None;
loop {
match coroutine.resume(&mut self.fragmentizer, arg.take()) {
ImapCoroutineState::Complete(Ok(out)) => return Ok(out),
ImapCoroutineState::Complete(Err(err)) => return Err(err.into()),
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
let n = self.stream.read(&mut buf).await?;
if n == 0 {
let kind = io::ErrorKind::UnexpectedEof;
let err = io::Error::new(kind, "IMAP server closed the connection");
return Err(err.into());
}
arg = Some(&buf[..n]);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.stream.write_all(&bytes).await?;
}
}
}
}
}
}
enum TokioStream {
Tcp(TcpStream),
Tls(Box<TlsStream<TcpStream>>),
Unix(UnixStream),
}
impl TokioStream {
async fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
match self {
Self::Tcp(stream) => stream.read(buf).await,
Self::Tls(stream) => stream.read(buf).await,
Self::Unix(stream) => stream.read(buf).await,
}
}
async fn write_all(&mut self, bytes: &[u8]) -> io::Result<()> {
match self {
Self::Tcp(stream) => stream.write_all(bytes).await,
Self::Tls(stream) => stream.write_all(bytes).await,
Self::Unix(stream) => stream.write_all(bytes).await,
}
}
async fn upgrade_tls(
self,
connector: &TlsConnector,
host: &str,
) -> Result<Self, Box<dyn Error>> {
let Self::Tcp(sock) = self else {
let err = String::from("IMAP STARTTLS upgrade on an already-encrypted transport");
return Err(err.into());
};
let name = ServerName::try_from(host)?.to_owned();
let tls = connector.connect(name, sock).await?;
Ok(Self::Tls(Box::new(tls)))
}
}