use core::{fmt, num::NonZeroU32};
use alloc::{string::String, string::ToString, vec, vec::Vec};
use imap_codec::{
CommandCodec, ResponseCodec,
encode::Encoder,
fragmentizer::{FragmentInfo, Fragmentizer},
imap_types::{
command::{Command, CommandBody},
core::TagGenerator,
fetch::{MacroOrMessageDataItemNames, MessageDataItemName},
response::{Response, Status, StatusKind},
sequence::{SeqOrUid, SequenceSet},
},
};
use log::{debug, trace};
use thiserror::Error;
use crate::coroutine::*;
#[derive(Clone, Debug, Error)]
pub enum ImapMessageFetchStreamError {
#[error("IMAP FETCH failed: NO {0}")]
No(String),
#[error("IMAP FETCH failed: BAD {0}")]
Bad(String),
#[error("IMAP FETCH failed: BYE {0}")]
Bye(String),
#[error("IMAP FETCH failed: server did not return a tagged response")]
MissingTagged,
#[error("IMAP FETCH failed: stream ended before the declared body length")]
ShortBody,
#[error("IMAP FETCH failed: unexpected literal in response trailer")]
UnexpectedLiteral,
}
#[derive(Debug)]
pub enum ImapMessageFetchStreamYield {
WantsRead,
WantsWrite(Vec<u8>),
BodyChunk(Vec<u8>),
WantsStream {
len: u32,
},
}
pub struct ImapMessageFetchStream {
state: State,
command: Option<Vec<u8>>,
pending: Vec<u8>,
remaining: u32,
stream_pending: bool,
codec: ResponseCodec,
}
impl ImapMessageFetchStream {
pub fn new(id: NonZeroU32, uid: bool) -> Self {
let command = Command {
tag: TagGenerator::new().generate(),
body: CommandBody::Fetch {
sequence_set: SequenceSet::from(SeqOrUid::from(id)),
macro_or_item_names: MacroOrMessageDataItemNames::MessageDataItemNames(vec![
MessageDataItemName::BodyExt {
section: None,
partial: None,
peek: true,
},
]),
uid,
modifiers: Vec::new(),
},
};
trace!("send IMAP command {command:?}");
let command = CommandCodec::new().encode(&command).dump();
Self {
state: State::SendCommand,
command: Some(command),
pending: Vec::new(),
remaining: 0,
stream_pending: false,
codec: ResponseCodec::new(),
}
}
}
impl ImapCoroutine for ImapMessageFetchStream {
type Yield = ImapMessageFetchStreamYield;
type Return = Result<(), ImapMessageFetchStreamError>;
fn resume(
&mut self,
fragmentizer: &mut Fragmentizer,
mut arg: Option<&[u8]>,
) -> ImapCoroutineState<Self::Yield, Self::Return> {
loop {
match self.state {
State::SendCommand => {
let command = self.command.take().expect("command sent once");
self.state = State::Header;
debug!("{}", self.state);
return ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::WantsWrite(
command,
));
}
State::Header => {
if let Some(bytes) = arg.take() {
if bytes.is_empty() {
let err = ImapMessageFetchStreamError::MissingTagged;
return ImapCoroutineState::Complete(Err(err));
}
self.pending.extend_from_slice(bytes);
}
loop {
let Some(nl) = self.pending.iter().position(|&b| b == b'\n') else {
return ImapCoroutineState::Yielded(
ImapMessageFetchStreamYield::WantsRead,
);
};
let line: Vec<u8> = self.pending.drain(..=nl).collect();
fragmentizer.enqueue_bytes(&line);
match fragmentizer.progress() {
Some(FragmentInfo::Line {
announcement: Some(announcement),
..
}) => {
self.remaining = announcement.length;
self.state = State::Stream;
debug!("{}", self.state);
break;
}
Some(FragmentInfo::Line {
announcement: None, ..
}) => {
if let Some(result) = self.decode_terminal(fragmentizer) {
return result;
}
}
_ => {}
}
}
}
State::Stream => {
if self.remaining == 0 {
fragmentizer.skip_message();
self.state = State::Trailer;
debug!("{}", self.state);
continue;
}
if !self.pending.is_empty() {
let take = (self.remaining as usize).min(self.pending.len());
let chunk: Vec<u8> = self.pending.drain(..take).collect();
self.remaining -= take as u32;
return ImapCoroutineState::Yielded(
ImapMessageFetchStreamYield::BodyChunk(chunk),
);
}
if self.stream_pending {
self.stream_pending = false;
if matches!(arg.take(), Some(&[])) {
let err = ImapMessageFetchStreamError::ShortBody;
return ImapCoroutineState::Complete(Err(err));
}
self.remaining = 0;
continue;
}
self.stream_pending = true;
return ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::WantsStream {
len: self.remaining,
});
}
State::Trailer => {
if let Some(bytes) = arg.take() {
if bytes.is_empty() {
let err = ImapMessageFetchStreamError::MissingTagged;
return ImapCoroutineState::Complete(Err(err));
}
self.pending.extend_from_slice(bytes);
}
loop {
let Some(nl) = self.pending.iter().position(|&b| b == b'\n') else {
return ImapCoroutineState::Yielded(
ImapMessageFetchStreamYield::WantsRead,
);
};
let line: Vec<u8> = self.pending.drain(..=nl).collect();
fragmentizer.enqueue_bytes(&line);
match fragmentizer.progress() {
Some(FragmentInfo::Line {
announcement: Some(_),
..
}) => {
let err = ImapMessageFetchStreamError::UnexpectedLiteral;
return ImapCoroutineState::Complete(Err(err));
}
Some(FragmentInfo::Line {
announcement: None, ..
}) => {
if let Some(result) = self.decode_terminal(fragmentizer) {
return result;
}
}
_ => {}
}
}
}
}
}
}
}
impl ImapMessageFetchStream {
fn decode_terminal(
&self,
fragmentizer: &Fragmentizer,
) -> Option<
ImapCoroutineState<ImapMessageFetchStreamYield, Result<(), ImapMessageFetchStreamError>>,
> {
match fragmentizer.decode_message(&self.codec) {
Ok(Response::Status(Status::Tagged(tagged))) => {
let text = tagged.body.text.to_string();
let result = match tagged.body.kind {
StatusKind::Ok => Ok(()),
StatusKind::No => Err(ImapMessageFetchStreamError::No(text)),
StatusKind::Bad => Err(ImapMessageFetchStreamError::Bad(text)),
};
Some(ImapCoroutineState::Complete(result))
}
Ok(Response::Status(Status::Bye(bye))) => {
let err = ImapMessageFetchStreamError::Bye(bye.text.to_string());
Some(ImapCoroutineState::Complete(Err(err)))
}
_ => None,
}
}
}
#[derive(Clone, Copy)]
enum State {
SendCommand,
Header,
Stream,
Trailer,
}
impl fmt::Display for State {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::SendCommand => f.write_str("send fetch command"),
Self::Header => f.write_str("parse fetch header"),
Self::Stream => f.write_str("stream body"),
Self::Trailer => f.write_str("parse fetch trailer"),
}
}
}
#[cfg(test)]
mod tests {
use core::str;
use alloc::{borrow::ToOwned, format};
use crate::rfc3501::fetch_stream::*;
#[test]
fn streams_body_in_one_read() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(1).unwrap(), true);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let cmd = expect_wants_write(&mut cor, &mut frag, None);
let line = str::from_utf8(&cmd).expect("utf8 command");
let tag = first_word(line).to_owned();
assert!(line.contains("UID FETCH 1 BODY.PEEK[]"));
expect_wants_read(&mut cor, &mut frag, None);
let reply = format!("* 1 FETCH (BODY[] {{5}}\r\nhello)\r\n{tag} OK FETCH completed\r\n");
let chunk = expect_body_chunk(&mut cor, &mut frag, Some(reply.as_bytes()));
assert_eq!(chunk, b"hello");
expect_complete_ok(&mut cor, &mut frag, None);
}
#[test]
fn streams_body_via_wants_stream() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(9).unwrap(), false);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let cmd = expect_wants_write(&mut cor, &mut frag, None);
let line = str::from_utf8(&cmd).expect("utf8 command");
let tag = first_word(line).to_owned();
assert!(line.contains("FETCH 9 BODY.PEEK[]"));
assert!(!line.contains("UID"));
expect_wants_read(&mut cor, &mut frag, None);
let len = expect_wants_stream(&mut cor, &mut frag, Some(b"* 9 FETCH (BODY[] {12}\r\n"));
assert_eq!(len, 12);
expect_wants_read(&mut cor, &mut frag, None);
let reply = format!(")\r\n{tag} OK FETCH completed\r\n");
expect_complete_ok(&mut cor, &mut frag, Some(reply.as_bytes()));
}
#[test]
fn partial_body_in_header_read_chunks_then_streams() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(1).unwrap(), true);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let cmd = expect_wants_write(&mut cor, &mut frag, None);
let tag = first_word(str::from_utf8(&cmd).expect("utf8 command")).to_owned();
expect_wants_read(&mut cor, &mut frag, None);
let chunk = expect_body_chunk(&mut cor, &mut frag, Some(b"* 1 FETCH (BODY[] {5}\r\nhel"));
assert_eq!(chunk, b"hel");
let len = expect_wants_stream(&mut cor, &mut frag, None);
assert_eq!(len, 2);
expect_wants_read(&mut cor, &mut frag, None);
let reply = format!(")\r\n{tag} OK done\r\n");
expect_complete_ok(&mut cor, &mut frag, Some(reply.as_bytes()));
}
#[test]
fn missing_message_returns_ok_without_body() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(7).unwrap(), true);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let cmd = expect_wants_write(&mut cor, &mut frag, None);
let tag = first_word(str::from_utf8(&cmd).expect("utf8 command")).to_owned();
expect_wants_read(&mut cor, &mut frag, None);
let reply = format!("{tag} OK FETCH completed\r\n");
expect_complete_ok(&mut cor, &mut frag, Some(reply.as_bytes()));
}
#[test]
fn tagged_no_returns_no_error() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(7).unwrap(), true);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let cmd = expect_wants_write(&mut cor, &mut frag, None);
let tag = first_word(str::from_utf8(&cmd).expect("utf8 command")).to_owned();
expect_wants_read(&mut cor, &mut frag, None);
let reply = format!("{tag} NO mailbox not selected\r\n");
let err = expect_complete_err(&mut cor, &mut frag, Some(reply.as_bytes()));
let ImapMessageFetchStreamError::No(text) = err else {
panic!("expected ImapMessageFetchStreamError::No, got {err:?}");
};
assert_eq!(text, "mailbox not selected");
}
#[test]
fn short_stream_returns_short_body() {
let mut cor = ImapMessageFetchStream::new(NonZeroU32::new(1).unwrap(), true);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let _ = expect_wants_write(&mut cor, &mut frag, None);
expect_wants_read(&mut cor, &mut frag, None);
let _ = expect_wants_stream(&mut cor, &mut frag, Some(b"* 1 FETCH (BODY[] {12}\r\n"));
let err = expect_complete_err(&mut cor, &mut frag, Some(&[]));
assert!(matches!(err, ImapMessageFetchStreamError::ShortBody));
}
fn expect_wants_write(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> Vec<u8> {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::WantsWrite(bytes)) => bytes,
state => panic!("expected WantsWrite, got {state:?}"),
}
}
fn expect_wants_read(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::WantsRead) => {}
state => panic!("expected WantsRead, got {state:?}"),
}
}
fn expect_body_chunk(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> Vec<u8> {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::BodyChunk(bytes)) => bytes,
state => panic!("expected BodyChunk, got {state:?}"),
}
}
fn expect_wants_stream(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> u32 {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageFetchStreamYield::WantsStream { len }) => len,
state => panic!("expected WantsStream, got {state:?}"),
}
}
fn expect_complete_ok(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) {
match cor.resume(frag, arg) {
ImapCoroutineState::Complete(Ok(())) => {}
state => panic!("expected Complete(Ok), got {state:?}"),
}
}
fn expect_complete_err(
cor: &mut ImapMessageFetchStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> ImapMessageFetchStreamError {
match cor.resume(frag, arg) {
ImapCoroutineState::Complete(Err(err)) => err,
state => panic!("expected Complete(Err), got {state:?}"),
}
}
fn first_word(line: &str) -> &str {
line.split_whitespace()
.next()
.expect("first whitespace-separated token")
}
}