use core::{
mem,
num::{NonZeroU32, NonZeroU64},
sync::atomic::{AtomicBool, Ordering},
time::Duration,
};
use alloc::{
collections::{BTreeMap, VecDeque},
string::String,
sync::Arc,
vec,
vec::Vec,
};
use imap_codec::{
fragmentizer::Fragmentizer,
imap_types::{
command::SelectParameter,
core::{Atom, Vec1},
extensions::enable::CapabilityEnable,
fetch::{MacroOrMessageDataItemNames, MessageDataItem, MessageDataItemName},
flag::{Flag, FlagFetch},
mailbox::Mailbox,
response::Capability,
sequence::SequenceSet,
},
};
use log::{debug, trace};
use thiserror::Error;
use crate::{
coroutine::*,
rfc2177::idle::{ImapIdle, ImapIdleError, ImapIdleOptions, ImapIdleYield},
rfc3501::{
examine::{ImapMailboxExamine, ImapMailboxExamineError, ImapMailboxExamineOptions},
fetch::{ImapMessageFetch, ImapMessageFetchError, ImapMessageFetchOptions},
select::ImapMailboxSelectData,
},
rfc5161::enable::{ImapExtensionEnable, ImapExtensionEnableError},
};
#[derive(Clone, Debug)]
pub enum ImapMailboxWatchEvent {
EnvelopeAdded {
uid: NonZeroU32,
items: Vec<MessageDataItem<'static>>,
},
FlagsAdded {
uid: NonZeroU32,
flags: Vec<Flag<'static>>,
},
FlagsRemoved {
uid: NonZeroU32,
flags: Vec<Flag<'static>>,
},
EnvelopeRemoved {
uid: NonZeroU32,
},
}
#[derive(Debug, Error)]
pub enum ImapMailboxWatchError {
#[error("IMAP server did not return UIDVALIDITY in EXAMINE response")]
MissingUidValidity,
#[error("IMAP server did not return HIGHESTMODSEQ in EXAMINE response")]
MissingHighestModSeq,
#[error("IMAP mailbox UIDVALIDITY changed from {known} to {seen}")]
UidValidityChanged {
known: NonZeroU32,
seen: NonZeroU32,
},
#[error("Invalid `1:*` sequence set: {0}")]
InvalidSequenceSet(String),
#[error("IMAP EXAMINE error")]
Examine(#[from] ImapMailboxExamineError),
#[error("IMAP FETCH error")]
Fetch(#[from] ImapMessageFetchError),
#[error("IMAP IDLE error")]
Idle(#[from] ImapIdleError),
#[error("IMAP ENABLE error")]
Enable(#[from] ImapExtensionEnableError),
}
#[derive(Clone, Copy, Debug, Default)]
pub struct ImapMailboxWatchOptions {
pub idle_timeout: Option<Duration>,
pub poll: bool,
}
#[derive(Debug)]
pub enum ImapMailboxWatchYield {
WantsRead,
WantsWait,
WantsWrite(Vec<u8>),
Event(ImapMailboxWatchEvent),
}
enum State {
EnableQresync(ImapExtensionEnable),
ExamineInitial(ImapMailboxExamine),
FetchBaseline(ImapMessageFetch),
BeginIdle,
Idle(ImapIdle),
Waiting,
ExamineQresync(ImapMailboxExamine),
ExamineResync(ImapMailboxExamine),
FetchResync(ImapMessageFetch),
EmitDeltas,
Terminal,
}
pub struct ImapMailboxWatch {
state: State,
opts: ImapMailboxWatchOptions,
qresync: bool,
shutdown: Arc<AtomicBool>,
idle_done: Arc<AtomicBool>,
idle_saw_data: bool,
mailbox: Mailbox<'static>,
uid_validity: Option<NonZeroU32>,
highest_mod_seq: u64,
shadow: BTreeMap<NonZeroU32, Vec<Flag<'static>>>,
pending: VecDeque<ImapMailboxWatchEvent>,
}
impl ImapMailboxWatch {
pub fn new(
capability: &[Capability<'static>],
mailbox: Mailbox<'static>,
shutdown: Arc<AtomicBool>,
opts: ImapMailboxWatchOptions,
) -> Self {
let qresync = capability.contains(&Capability::QResync);
let state = if qresync {
let condstore = CapabilityEnable::CondStore;
let qresync = CapabilityEnable::from(
Atom::try_from("QRESYNC").expect("`QRESYNC` is a syntactically valid IMAP atom"),
);
let capabilities =
Vec1::try_from(vec![condstore, qresync]).expect("two capabilities is non-empty");
State::EnableQresync(ImapExtensionEnable::new(capabilities))
} else {
debug!("qresync unsupported, watching the whole mailbox");
State::ExamineInitial(ImapMailboxExamine::new(
mailbox.clone(),
ImapMailboxExamineOptions::default(),
))
};
Self {
state,
opts,
qresync,
shutdown,
idle_done: Arc::new(AtomicBool::new(false)),
idle_saw_data: false,
mailbox,
uid_validity: None,
highest_mod_seq: 0,
shadow: BTreeMap::new(),
pending: VecDeque::new(),
}
}
fn resync(&self) -> State {
if !self.qresync {
let examine =
ImapMailboxExamine::new(self.mailbox.clone(), ImapMailboxExamineOptions::default());
return State::ExamineResync(examine);
}
let uid_validity = self.uid_validity.unwrap();
let modseq = NonZeroU64::new(self.highest_mod_seq)
.unwrap_or_else(|| NonZeroU64::new(1).expect("1 is non-zero"));
let parameters = vec![SelectParameter::QResync {
uid_validity,
mod_sequence_value: modseq,
known_uids: None,
seq_match_data: None,
}];
let examine = ImapMailboxExamine::new(
self.mailbox.clone(),
ImapMailboxExamineOptions { parameters },
);
State::ExamineQresync(examine)
}
fn check_uid_validity(
&self,
data: &ImapMailboxSelectData,
) -> Result<(), ImapMailboxWatchError> {
let (Some(known), Some(seen)) = (self.uid_validity, data.uid_validity) else {
return Ok(());
};
if known != seen {
return Err(ImapMailboxWatchError::UidValidityChanged { known, seen });
}
Ok(())
}
fn push_flag_deltas(
&mut self,
uid: NonZeroU32,
old_flags: &[Flag<'static>],
new_flags: &[Flag<'static>],
) {
let added: Vec<Flag<'static>> = new_flags
.iter()
.filter(|f| !old_flags.contains(f))
.cloned()
.collect();
let removed: Vec<Flag<'static>> = old_flags
.iter()
.filter(|f| !new_flags.contains(f))
.cloned()
.collect();
if !added.is_empty() {
self.pending
.push_back(ImapMailboxWatchEvent::FlagsAdded { uid, flags: added });
}
if !removed.is_empty() {
self.pending.push_back(ImapMailboxWatchEvent::FlagsRemoved {
uid,
flags: removed,
});
}
}
fn compute_snapshot_deltas(
&mut self,
snapshot: BTreeMap<NonZeroU32, Vec<MessageDataItem<'static>>>,
) {
let vanished: Vec<NonZeroU32> = self
.shadow
.keys()
.filter(|uid| !snapshot.contains_key(uid))
.copied()
.collect();
for uid in vanished {
self.shadow.remove(&uid);
self.pending
.push_back(ImapMailboxWatchEvent::EnvelopeRemoved { uid });
}
for (uid, items) in snapshot {
let (_uid, new_flags) = extract_uid_flags(&items);
match self.shadow.insert(uid, new_flags.clone()) {
None => {
self.pending
.push_back(ImapMailboxWatchEvent::EnvelopeAdded { uid, items });
}
Some(old_flags) => self.push_flag_deltas(uid, &old_flags, &new_flags),
}
}
}
fn compute_deltas(&mut self, data: &ImapMailboxSelectData) {
for uid in &data.vanished_earlier {
if self.shadow.remove(uid).is_some() {
self.pending
.push_back(ImapMailboxWatchEvent::EnvelopeRemoved { uid: *uid });
}
}
for fetch in &data.changed {
let items_vec: Vec<MessageDataItem<'static>> =
fetch.items.clone().into_inner().into_iter().collect();
let (uid_opt, new_flags) = extract_uid_flags(&items_vec);
let Some(uid) = uid_opt else {
continue;
};
match self.shadow.insert(uid, new_flags.clone()) {
None => {
self.pending
.push_back(ImapMailboxWatchEvent::EnvelopeAdded {
uid,
items: items_vec,
});
}
Some(old_flags) => self.push_flag_deltas(uid, &old_flags, &new_flags),
}
}
}
}
fn fetch_uid_flags() -> Result<ImapMessageFetch, ImapMailboxWatchError> {
let sequence_set: SequenceSet = "1:*"
.try_into()
.map_err(|_| ImapMailboxWatchError::InvalidSequenceSet("1:*".into()))?;
let item_names = MacroOrMessageDataItemNames::MessageDataItemNames(vec![
MessageDataItemName::Uid,
MessageDataItemName::Flags,
]);
Ok(ImapMessageFetch::new(
sequence_set,
item_names,
ImapMessageFetchOptions::default(),
))
}
impl ImapCoroutine for ImapMailboxWatch {
type Yield = ImapMailboxWatchYield;
type Return = Result<(), ImapMailboxWatchError>;
fn resume(
&mut self,
fragmentizer: &mut Fragmentizer,
mut arg: Option<&[u8]>,
) -> ImapCoroutineState<Self::Yield, Self::Return> {
if self.shutdown.load(Ordering::SeqCst) {
self.idle_done.store(true, Ordering::SeqCst);
}
loop {
let state = mem::replace(&mut self.state, State::Terminal);
match state {
State::EnableQresync(mut enable) => match enable.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::EnableQresync(enable);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::EnableQresync(enable);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(enabled)) => {
debug!("enabled qresync");
trace!("{enabled:?}");
let parameters = vec![SelectParameter::CondStore];
let examine = ImapMailboxExamine::new(
self.mailbox.clone(),
ImapMailboxExamineOptions { parameters },
);
self.state = State::ExamineInitial(examine);
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
},
State::ExamineInitial(mut examine) => {
match examine.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::ExamineInitial(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::ExamineInitial(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(data)) => {
let Some(uid_validity) = data.uid_validity else {
return ImapCoroutineState::Complete(Err(
ImapMailboxWatchError::MissingUidValidity,
));
};
self.uid_validity = Some(uid_validity);
trace!("uid_validity: {uid_validity}");
if self.qresync {
let Some(highest_mod_seq) = data.highest_mod_seq else {
return ImapCoroutineState::Complete(Err(
ImapMailboxWatchError::MissingHighestModSeq,
));
};
self.highest_mod_seq = highest_mod_seq;
debug!("examined mailbox with condstore");
trace!("highest_mod_seq: {highest_mod_seq}");
} else {
debug!("examined mailbox");
}
let fetch = match fetch_uid_flags() {
Ok(fetch) => fetch,
Err(err) => return ImapCoroutineState::Complete(Err(err)),
};
self.state = State::FetchBaseline(fetch);
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
}
}
State::FetchBaseline(mut fetch) => match fetch.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::FetchBaseline(fetch);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::FetchBaseline(fetch);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(data)) => {
for (_seq, items) in data {
let items_vec = items.into_inner();
if let (Some(uid), flags) = extract_uid_flags(&items_vec) {
self.shadow.insert(uid, flags);
}
}
debug!("seeded baseline shadow");
trace!("uids: {}", self.shadow.len());
self.state = State::BeginIdle;
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
},
State::BeginIdle => {
if self.shutdown.load(Ordering::SeqCst) {
return ImapCoroutineState::Complete(Ok(()));
}
if self.opts.poll {
self.state = State::Waiting;
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWait);
}
self.idle_done.store(false, Ordering::SeqCst);
self.idle_saw_data = false;
let opts = ImapIdleOptions {
timeout: self.opts.idle_timeout,
};
let idle = ImapIdle::new(self.idle_done.clone(), opts);
self.state = State::Idle(idle);
}
State::Waiting => {
if self.shutdown.load(Ordering::SeqCst) {
return ImapCoroutineState::Complete(Ok(()));
}
self.state = self.resync();
}
State::Idle(mut idle) => match idle.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapIdleYield::Event(_)) => {
debug!("idle saw untagged data");
self.idle_saw_data = true;
self.idle_done.store(true, Ordering::SeqCst);
self.state = State::Idle(idle);
}
ImapCoroutineState::Yielded(ImapIdleYield::WantsRead) => {
self.state = State::Idle(idle);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapIdleYield::WantsWrite(bytes)) => {
self.state = State::Idle(idle);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(())) => {
if self.shutdown.load(Ordering::SeqCst) {
return ImapCoroutineState::Complete(Ok(()));
}
if self.idle_saw_data {
self.state = self.resync();
} else {
debug!("idle timed out with no data, restarting");
self.state = State::BeginIdle;
}
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
},
State::ExamineQresync(mut examine) => {
match examine.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::ExamineQresync(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::ExamineQresync(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(data)) => {
if let Err(err) = self.check_uid_validity(&data) {
return ImapCoroutineState::Complete(Err(err));
}
self.compute_deltas(&data);
if let Some(new_modseq) = data.highest_mod_seq {
self.highest_mod_seq = new_modseq;
}
self.state = State::EmitDeltas;
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
}
}
State::ExamineResync(mut examine) => {
match examine.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::ExamineResync(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::ExamineResync(examine);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(data)) => {
if let Err(err) = self.check_uid_validity(&data) {
return ImapCoroutineState::Complete(Err(err));
}
let fetch = match fetch_uid_flags() {
Ok(fetch) => fetch,
Err(err) => return ImapCoroutineState::Complete(Err(err)),
};
self.state = State::FetchResync(fetch);
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
}
}
State::FetchResync(mut fetch) => match fetch.resume(fragmentizer, arg.take()) {
ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
self.state = State::FetchResync(fetch);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
}
ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
self.state = State::FetchResync(fetch);
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
bytes,
));
}
ImapCoroutineState::Complete(Ok(data)) => {
let mut snapshot = BTreeMap::new();
for (_seq, items) in data {
let items_vec = items.into_inner();
if let (Some(uid), _flags) = extract_uid_flags(&items_vec) {
snapshot.insert(uid, items_vec);
}
}
debug!("re-read the whole mailbox");
trace!("uids: {}", snapshot.len());
self.compute_snapshot_deltas(snapshot);
self.state = State::EmitDeltas;
}
ImapCoroutineState::Complete(Err(err)) => {
return ImapCoroutineState::Complete(Err(err.into()));
}
},
State::EmitDeltas => {
if let Some(event) = self.pending.pop_front() {
self.state = State::EmitDeltas;
return ImapCoroutineState::Yielded(ImapMailboxWatchYield::Event(event));
}
self.state = State::BeginIdle;
}
State::Terminal => {
self.state = State::Terminal;
return ImapCoroutineState::Complete(Ok(()));
}
}
}
}
}
fn extract_uid_flags(
items: &[MessageDataItem<'static>],
) -> (Option<NonZeroU32>, Vec<Flag<'static>>) {
let mut uid = None;
let mut flags = Vec::new();
for item in items {
match item {
MessageDataItem::Uid(u) => uid = Some(*u),
MessageDataItem::Flags(fs) => {
flags = fs
.iter()
.filter_map(|f| match f {
FlagFetch::Flag(flag) => Some(flag.clone()),
_ => None,
})
.collect();
}
_ => {}
}
}
(uid, flags)
}
#[cfg(test)]
mod tests {
use core::str;
use alloc::{borrow::ToOwned, format, string::ToString};
use crate::watch::*;
const UID_VALIDITY: u32 = 1700;
type Step<'a> = (&'a str, &'a [&'a str]);
fn watcher(capability: &[Capability<'static>]) -> (ImapMailboxWatch, Fragmentizer) {
watcher_with(capability, ImapMailboxWatchOptions::default())
}
fn watcher_with(
capability: &[Capability<'static>],
opts: ImapMailboxWatchOptions,
) -> (ImapMailboxWatch, Fragmentizer) {
let watch = ImapMailboxWatch::new(
capability,
"INBOX".try_into().expect("valid mailbox"),
Arc::new(AtomicBool::new(false)),
opts,
);
(watch, Fragmentizer::new(50 * 1024 * 1024))
}
fn first_word(line: &str) -> &str {
line.split_whitespace()
.next()
.expect("first whitespace-separated token")
}
fn drive(
cor: &mut ImapMailboxWatch,
frag: &mut Fragmentizer,
steps: &[Step],
) -> Result<Vec<ImapMailboxWatchEvent>, ImapMailboxWatchError> {
let mut events = Vec::new();
let mut replies: VecDeque<String> = VecDeque::new();
let mut tag = String::new();
let mut next = 0;
let mut arg: Option<Vec<u8>> = None;
loop {
match cor.resume(frag, arg.take().as_deref()) {
ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(bytes)) => {
let line = str::from_utf8(&bytes).expect("utf8 command").to_string();
let Some((expected, scripted)) = steps.get(next) else {
return Ok(events);
};
assert!(line.contains(expected), "expected {expected}, wrote {line}");
next += 1;
if !line.starts_with("DONE") {
tag = first_word(&line).to_owned();
}
replies.extend(scripted.iter().map(|reply| reply.replace("{tag}", &tag)));
}
ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead) => {
let reply = replies
.pop_front()
.expect("the script owes a reply to every read");
arg = Some(reply.into_bytes());
}
ImapCoroutineState::Yielded(ImapMailboxWatchYield::Event(event)) => {
events.push(event);
}
ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWait) => {}
ImapCoroutineState::Complete(Ok(())) => panic!("the watch stopped early"),
ImapCoroutineState::Complete(Err(err)) => return Err(err),
}
}
}
fn first_command(cor: &mut ImapMailboxWatch, frag: &mut Fragmentizer) -> String {
match cor.resume(frag, None) {
ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(bytes)) => {
str::from_utf8(&bytes).expect("utf8 command").to_string()
}
state => panic!("expected WantsWrite, got {state:?}"),
}
}
fn examined(uid_validity: u32) -> String {
format!(
"* 2 EXISTS\r\n\
* OK [UIDVALIDITY {uid_validity}] uid validity\r\n\
{{tag}} OK [READ-ONLY] EXAMINE completed\r\n",
)
}
fn fetched(messages: &[(u32, &str)]) -> String {
let mut reply = String::new();
for (seq, (uid, flags)) in messages.iter().enumerate() {
reply.push_str(&format!(
"* {} FETCH (UID {uid} FLAGS ({flags}))\r\n",
seq + 1
));
}
reply.push_str("{tag} OK FETCH completed\r\n");
reply
}
#[test]
fn qresync_capability_enables_it_first() {
let (mut watch, mut frag) = watcher(&[Capability::QResync]);
let line = first_command(&mut watch, &mut frag);
assert!(line.contains("ENABLE"), "wrote {line}");
}
#[test]
fn missing_qresync_examines_straight_away() {
let (mut watch, mut frag) = watcher(&[]);
let line = first_command(&mut watch, &mut frag);
assert!(line.contains("EXAMINE INBOX"), "wrote {line}");
assert!(!line.contains("ENABLE"), "wrote {line}");
assert!(!line.contains("CONDSTORE"), "wrote {line}");
}
#[test]
fn fallback_diffs_the_whole_mailbox_on_every_wake() {
let (mut watch, mut frag) = watcher(&[]);
let baseline = fetched(&[(1, ""), (2, "")]);
let resynced = fetched(&[(2, "\\Seen"), (3, "")]);
let steps: &[Step] = &[
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&baseline]),
("IDLE", &["+ idling\r\n", "* 3 EXISTS\r\n"]),
("DONE", &["{tag} OK IDLE terminated\r\n"]),
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&resynced]),
];
let events = drive(&mut watch, &mut frag, steps).expect("watch running");
assert_eq!(3, events.len(), "got {events:?}");
let ImapMailboxWatchEvent::EnvelopeRemoved { uid } = &events[0] else {
panic!("expected EnvelopeRemoved, got {:?}", events[0]);
};
assert_eq!(1, uid.get());
let ImapMailboxWatchEvent::FlagsAdded { uid, flags } = &events[1] else {
panic!("expected FlagsAdded, got {:?}", events[1]);
};
assert_eq!(2, uid.get());
assert_eq!(&vec![Flag::Seen], flags);
let ImapMailboxWatchEvent::EnvelopeAdded { uid, .. } = &events[2] else {
panic!("expected EnvelopeAdded, got {:?}", events[2]);
};
assert_eq!(3, uid.get());
}
#[test]
fn fallback_reports_nothing_when_the_mailbox_is_unchanged() {
let (mut watch, mut frag) = watcher(&[]);
let snapshot = fetched(&[(1, "\\Seen")]);
let steps: &[Step] = &[
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&snapshot]),
("IDLE", &["+ idling\r\n", "* 1 EXISTS\r\n"]),
("DONE", &["{tag} OK IDLE terminated\r\n"]),
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&snapshot]),
];
let events = drive(&mut watch, &mut frag, steps).expect("watch running");
assert!(events.is_empty(), "got {events:?}");
}
#[test]
fn a_polling_watch_re_reads_instead_of_idling() {
let opts = ImapMailboxWatchOptions {
poll: true,
..Default::default()
};
let (mut watch, mut frag) = watcher_with(&[], opts);
let baseline = fetched(&[(1, "")]);
let resynced = fetched(&[(1, "\\Seen")]);
let steps: &[Step] = &[
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&baseline]),
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&resynced]),
];
let events = drive(&mut watch, &mut frag, steps).expect("watch running");
assert_eq!(1, events.len(), "got {events:?}");
let ImapMailboxWatchEvent::FlagsAdded { uid, flags } = &events[0] else {
panic!("expected FlagsAdded, got {:?}", events[0]);
};
assert_eq!(1, uid.get());
assert_eq!(&vec![Flag::Seen], flags);
}
#[test]
fn a_recreated_mailbox_ends_the_watch() {
let (mut watch, mut frag) = watcher(&[]);
let baseline = fetched(&[(1, "")]);
let steps: &[Step] = &[
("EXAMINE INBOX", &[&examined(UID_VALIDITY)]),
("FETCH 1:* (UID FLAGS)", &[&baseline]),
("IDLE", &["+ idling\r\n", "* 1 EXPUNGE\r\n"]),
("DONE", &["{tag} OK IDLE terminated\r\n"]),
("EXAMINE INBOX", &[&examined(UID_VALIDITY + 1)]),
];
let err = drive(&mut watch, &mut frag, steps).expect_err("uid validity changed");
let ImapMailboxWatchError::UidValidityChanged { known, seen } = err else {
panic!("expected UidValidityChanged, got {err:?}");
};
assert_eq!(UID_VALIDITY, known.get());
assert_eq!(UID_VALIDITY + 1, seen.get());
}
}