use core::fmt;
use alloc::{format, string::String, string::ToString, vec::Vec};
use imap_codec::{
CommandCodec,
encode::{Encoder, Fragment},
fragmentizer::Fragmentizer,
imap_types::{
command::{Command, CommandBody},
core::{Literal, TagGenerator},
extensions::binary::LiteralOrLiteral8,
mailbox::Mailbox,
response::{Code, Data, StatusKind, Tagged},
},
};
use log::{debug, trace};
use thiserror::Error;
use crate::{
coroutine::*,
imap_try,
rfc3501::{
append::{ImapMessageAppendOptions, ImapMessageAppendOutput},
mailbox::encode_inplace,
},
send::*,
};
#[derive(Clone, Debug, Error)]
pub enum ImapMessageAppendStreamError {
#[error("IMAP APPEND failed: NO {0}")]
No(String),
#[error("IMAP APPEND failed: BAD {0}")]
Bad(String),
#[error("IMAP APPEND failed: BYE {0}")]
Bye(String),
#[error("IMAP APPEND failed: server did not return a tagged response")]
MissingTagged,
#[error("IMAP APPEND failed: message source delivered fewer octets than declared")]
ShortMessage,
#[error("IMAP APPEND failed: {0}")]
Send(#[from] ImapSendError),
}
#[derive(Debug)]
pub enum ImapMessageAppendStreamYield {
WantsRead,
WantsWrite(Vec<u8>),
WantsStream,
}
impl From<ImapYield> for ImapMessageAppendStreamYield {
fn from(yielded: ImapYield) -> Self {
match yielded {
ImapYield::WantsRead => Self::WantsRead,
ImapYield::WantsWrite(bytes) => Self::WantsWrite(bytes),
}
}
}
pub struct ImapMessageAppendStream {
state: State,
header: Option<Vec<u8>>,
crlf: Option<Vec<u8>>,
command: Command<'static>,
non_sync: bool,
stream_pending: bool,
}
impl ImapMessageAppendStream {
pub fn new(mut mailbox: Mailbox<'static>, len: u32, opts: ImapMessageAppendOptions) -> Self {
encode_inplace(&mut mailbox);
let command = Command {
tag: TagGenerator::new().generate(),
body: CommandBody::Append {
mailbox,
flags: opts.flags,
date: opts.date,
message: LiteralOrLiteral8::Literal(Literal::unvalidated_non_sync(Vec::new())),
},
};
trace!("send IMAP command {command:?}");
let fragments: Vec<Fragment> = CommandCodec::new().encode(&command).collect();
let last = fragments
.iter()
.rposition(|fragment| matches!(fragment, Fragment::Literal { .. }))
.expect("APPEND always encodes a message literal");
let mut header = Vec::new();
let mut crlf = Vec::new();
for (index, fragment) in fragments.into_iter().enumerate() {
match fragment {
Fragment::Line { data } if index < last => header.extend(data),
Fragment::Line { data } => crlf.extend(data),
Fragment::Literal { data, .. } if index < last => header.extend(data),
Fragment::Literal { .. } => {}
}
}
const EMPTY_LITERAL: &[u8] = b"{0+}\r\n";
debug_assert!(header.ends_with(EMPTY_LITERAL));
header.truncate(header.len() - EMPTY_LITERAL.len());
if opts.non_sync {
header.extend_from_slice(format!("{{{len}+}}\r\n").as_bytes());
} else {
header.extend_from_slice(format!("{{{len}}}\r\n").as_bytes());
}
Self {
state: State::WriteHeader,
header: Some(header),
crlf: Some(crlf),
command,
non_sync: opts.non_sync,
stream_pending: false,
}
}
}
impl ImapCoroutine for ImapMessageAppendStream {
type Yield = ImapMessageAppendStreamYield;
type Return = Result<ImapMessageAppendOutput, ImapMessageAppendStreamError>;
fn resume(
&mut self,
fragmentizer: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> ImapCoroutineState<Self::Yield, Self::Return> {
loop {
match &mut self.state {
State::WriteHeader => {
let header = self.header.take().expect("header written once");
self.state = if self.non_sync {
State::Stream
} else {
State::Continuation(ImapSend::receive(self.command.clone()))
};
debug!("{}", self.state);
return ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsWrite(
header,
));
}
State::Continuation(recv) => {
let out = imap_try!(recv, fragmentizer, arg);
if let Some(bye) = out.bye {
let err = ImapMessageAppendStreamError::Bye(bye.text.to_string());
return ImapCoroutineState::Complete(Err(err));
}
if let Some(Tagged { body, .. }) = out.tagged {
let err = match body.kind {
StatusKind::No => {
ImapMessageAppendStreamError::No(body.text.to_string())
}
_ => ImapMessageAppendStreamError::Bad(body.text.to_string()),
};
return ImapCoroutineState::Complete(Err(err));
}
self.state = State::Stream;
debug!("{}", self.state);
}
State::Stream => {
if self.stream_pending {
self.stream_pending = false;
if matches!(arg, Some(&[])) {
let err = ImapMessageAppendStreamError::ShortMessage;
return ImapCoroutineState::Complete(Err(err));
}
self.state = State::WriteCrlf;
debug!("{}", self.state);
continue;
}
self.stream_pending = true;
return ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsStream);
}
State::WriteCrlf => {
let crlf = self.crlf.take().expect("crlf written once");
self.state = State::Recv(ImapSend::receive(self.command.clone()));
debug!("{}", self.state);
return ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsWrite(
crlf,
));
}
State::Recv(recv) => {
let out = imap_try!(recv, fragmentizer, arg);
if let Some(bye) = out.bye {
let err = ImapMessageAppendStreamError::Bye(bye.text.to_string());
return ImapCoroutineState::Complete(Err(err));
}
let Some(Tagged { body, .. }) = out.tagged else {
let err = ImapMessageAppendStreamError::MissingTagged;
return ImapCoroutineState::Complete(Err(err));
};
let mut exists = None;
for data in out.data {
if let Data::Exists(seq) = data {
exists = Some(seq);
}
}
return match body.kind {
StatusKind::Ok => {
let appenduid =
if let Some(Code::AppendUid { uid_validity, uid }) = body.code {
Some((uid_validity.get(), uid.get()))
} else {
None
};
ImapCoroutineState::Complete(Ok((exists, appenduid)))
}
StatusKind::No => {
let err = ImapMessageAppendStreamError::No(body.text.to_string());
ImapCoroutineState::Complete(Err(err))
}
StatusKind::Bad => {
let err = ImapMessageAppendStreamError::Bad(body.text.to_string());
ImapCoroutineState::Complete(Err(err))
}
};
}
}
}
}
}
enum State {
WriteHeader,
Continuation(ImapSend<CommandCodec>),
Stream,
WriteCrlf,
Recv(ImapSend<CommandCodec>),
}
impl fmt::Display for State {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::WriteHeader => f.write_str("write append header"),
Self::Continuation(_) => f.write_str("await continuation"),
Self::Stream => f.write_str("stream message"),
Self::WriteCrlf => f.write_str("write append crlf"),
Self::Recv(_) => f.write_str("receive append response"),
}
}
}
#[cfg(test)]
mod tests {
use core::str;
use alloc::borrow::ToOwned;
use crate::rfc3501::append_stream::*;
#[test]
fn sync_success_with_appenduid_returns_pair() {
let mut append = ImapMessageAppendStream::new(
"INBOX".try_into().expect("valid mailbox"),
15,
ImapMessageAppendOptions::default(),
);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let header = expect_wants_write(&mut append, &mut frag, None);
let line = str::from_utf8(&header).expect("utf8 header");
let tag = first_word(line).to_owned();
assert!(line.contains("APPEND INBOX"));
assert!(line.ends_with("{15}\r\n"));
expect_wants_read(&mut append, &mut frag, None);
expect_wants_stream(
&mut append,
&mut frag,
Some(b"+ Ready for literal data\r\n"),
);
let crlf = expect_wants_write(&mut append, &mut frag, None);
assert_eq!(crlf, b"\r\n");
expect_wants_read(&mut append, &mut frag, None);
let reply =
format!("* 42 EXISTS\r\n{tag} OK [APPENDUID 1700000000 7] APPEND completed\r\n");
let (exists, appenduid) =
expect_complete_ok(&mut append, &mut frag, Some(reply.as_bytes()));
assert_eq!(Some(42), exists);
assert_eq!(Some((1700000000, 7)), appenduid);
}
#[test]
fn non_sync_streams_without_continuation() {
let mut append = ImapMessageAppendStream::new(
"INBOX".try_into().expect("valid mailbox"),
15,
ImapMessageAppendOptions {
non_sync: true,
..Default::default()
},
);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let header = expect_wants_write(&mut append, &mut frag, None);
let line = str::from_utf8(&header).expect("utf8 header");
let tag = first_word(line).to_owned();
assert!(line.ends_with("{15+}\r\n"));
expect_wants_stream(&mut append, &mut frag, None);
let crlf = expect_wants_write(&mut append, &mut frag, None);
assert_eq!(crlf, b"\r\n");
expect_wants_read(&mut append, &mut frag, None);
let reply = format!("{tag} OK APPEND completed\r\n");
expect_complete_ok(&mut append, &mut frag, Some(reply.as_bytes()));
}
#[test]
fn continuation_no_returns_no_error() {
let mut append = ImapMessageAppendStream::new(
"INBOX".try_into().expect("valid mailbox"),
15,
ImapMessageAppendOptions::default(),
);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let header = expect_wants_write(&mut append, &mut frag, None);
let tag = first_word(str::from_utf8(&header).expect("utf8 header")).to_owned();
expect_wants_read(&mut append, &mut frag, None);
let reply = format!("{tag} NO over quota\r\n");
let err = expect_complete_err(&mut append, &mut frag, Some(reply.as_bytes()));
let ImapMessageAppendStreamError::No(text) = err else {
panic!("expected ImapMessageAppendStreamError::No, got {err:?}");
};
assert_eq!(text, "over quota");
}
#[test]
fn short_stream_returns_short_message_error() {
let mut append = ImapMessageAppendStream::new(
"INBOX".try_into().expect("valid mailbox"),
15,
ImapMessageAppendOptions::default(),
);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let _ = expect_wants_write(&mut append, &mut frag, None);
expect_wants_read(&mut append, &mut frag, None);
expect_wants_stream(&mut append, &mut frag, Some(b"+ go\r\n"));
let err = expect_complete_err(&mut append, &mut frag, Some(&[]));
assert!(matches!(err, ImapMessageAppendStreamError::ShortMessage));
}
#[test]
fn bye_returns_bye_error() {
let mut append = ImapMessageAppendStream::new(
"INBOX".try_into().expect("valid mailbox"),
15,
ImapMessageAppendOptions::default(),
);
let mut frag = Fragmentizer::new(50 * 1024 * 1024);
let _ = expect_wants_write(&mut append, &mut frag, None);
expect_wants_read(&mut append, &mut frag, None);
expect_wants_stream(&mut append, &mut frag, Some(b"+ go\r\n"));
let _ = expect_wants_write(&mut append, &mut frag, None);
expect_wants_read(&mut append, &mut frag, None);
let err = expect_complete_err(&mut append, &mut frag, Some(b"* BYE shutting down\r\n"));
let ImapMessageAppendStreamError::Bye(text) = err else {
panic!("expected ImapMessageAppendStreamError::Bye, got {err:?}");
};
assert_eq!(text, "shutting down");
}
fn expect_wants_write(
cor: &mut ImapMessageAppendStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> Vec<u8> {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsWrite(bytes)) => bytes,
state => panic!("expected WantsWrite, got {state:?}"),
}
}
fn expect_wants_read(
cor: &mut ImapMessageAppendStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsRead) => {}
state => panic!("expected WantsRead, got {state:?}"),
}
}
fn expect_wants_stream(
cor: &mut ImapMessageAppendStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) {
match cor.resume(frag, arg) {
ImapCoroutineState::Yielded(ImapMessageAppendStreamYield::WantsStream) => {}
state => panic!("expected WantsStream, got {state:?}"),
}
}
fn expect_complete_ok(
cor: &mut ImapMessageAppendStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> ImapMessageAppendOutput {
match cor.resume(frag, arg) {
ImapCoroutineState::Complete(Ok(value)) => value,
state => panic!("expected Complete(Ok), got {state:?}"),
}
}
fn expect_complete_err(
cor: &mut ImapMessageAppendStream,
frag: &mut Fragmentizer,
arg: Option<&[u8]>,
) -> ImapMessageAppendStreamError {
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")
}
}