use std::cmp::Ordering;
use std::collections::HashMap;
use std::collections::VecDeque;
use std::mem;
use chrono::prelude::*;
use serde::{Deserialize, Serialize};
use super::model::*;
use crate::support::error::Error;
use crate::support::small_bitset::SmallBitset;
#[cfg(not(test))]
const MAX_RECENT_EXPUNGEMENTS: usize = 1024;
#[cfg(test)]
const MAX_RECENT_EXPUNGEMENTS: usize = 4;
#[derive(Serialize, Deserialize, Debug, Clone, Default)]
#[serde(default)]
pub struct MailboxState {
flags: Vec<Flag>,
extant_messages: Vec<Uid>,
message_status: HashMap<Uid, MessageStatus>,
max_modseq: Option<Modseq>,
max_tx_modseq: Option<Modseq>,
recent_expungements: VecDeque<(Modseq, Uid)>,
soft_expungements: Vec<SoftExpungement>,
#[serde(skip)]
unapplied_create: usize,
#[serde(skip)]
unapplied_expunge: Vec<Uid>,
#[serde(skip)]
report_max_modseq: Option<Modseq>,
#[serde(skip)]
changed_flags_uids: Vec<Uid>,
}
#[derive(Serialize, Deserialize, Debug, Clone, Copy)]
struct SoftExpungement {
#[serde(with = "chrono::serde::ts_seconds")]
deadline: DateTime<Utc>,
uid: Uid,
}
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct MessageStatus {
#[serde(rename = "f")]
flags: SmallBitset,
#[serde(rename = "m")]
last_modified: Modseq,
#[serde(skip)]
recent: bool,
}
impl MailboxState {
pub fn new() -> Self {
Self::default()
}
pub fn init_transient(&mut self) {
self.report_max_modseq = self.max_modseq;
}
pub fn seen(&mut self, uid: Uid) {
let this_uid = uid.0.get();
let last_uid = self.max_modseq.map(|m| m.uid().0.get()).unwrap_or(0);
if this_uid > last_uid {
for uid in last_uid + 1..=this_uid {
let uid = Uid::of(uid).unwrap();
self.extant_messages.push(uid);
self.unapplied_create += 1;
self.message_status.insert(
uid,
MessageStatus {
flags: SmallBitset::new(),
last_modified: Modseq::new(uid, Cid::GENESIS),
recent: false,
},
);
}
self.max_modseq = Some(
self.max_modseq
.map(|m| m.with_uid(uid))
.unwrap_or_else(|| Modseq::new(uid, Cid::GENESIS)),
);
}
}
pub fn start_tx(&self) -> Result<(Cid, StateTransaction), Error> {
let m = self
.max_modseq
.expect("start_tx with no messages")
.next()
.ok_or(Error::MailboxFull)?;
Ok((
m.cid(),
StateTransaction {
max_uid: m.uid(),
ops: Vec::new(),
},
))
}
pub fn commit(&mut self, cid: Cid, tx: StateTransaction) {
assert_eq!(self.next_cid().unwrap(), cid);
self.seen(tx.max_uid);
let nominal_modseq = Modseq::new(tx.max_uid, cid);
let canonical_modseq = self
.max_tx_modseq
.map(|m| m.combine(nominal_modseq))
.unwrap_or(nominal_modseq);
self.max_tx_modseq = Some(canonical_modseq);
self.max_modseq = Some(
self.max_modseq
.map_or(canonical_modseq, |m| m.combine(canonical_modseq)),
);
let m = canonical_modseq;
for op in tx.ops {
use self::StateMutation::*;
match op {
AddFlag(uid, flag) => self.set_flag(uid, m, flag, true),
RmFlag(uid, flag) => self.set_flag(uid, m, flag, false),
Expunge(deadline, uid) => self.expunge(deadline, uid, m),
}
}
}
pub fn flush(&mut self) -> FlushResponse {
let mut unapplied_create = self.unapplied_create;
let mut expunged_pairs: Vec<(Seqnum, Uid)> = Vec::new();
let mut stillborn: Vec<Uid> = Vec::new();
if !self.unapplied_expunge.is_empty() {
let max_index = self.num_messages();
self.unapplied_expunge.sort_unstable();
self.unapplied_expunge.dedup();
for uid in &self.unapplied_expunge {
self.message_status.remove(uid);
}
let mut expunged = self.unapplied_expunge.drain(..).peekable();
let mut index = 0;
self.extant_messages.retain(|&uid| loop {
let next_expunged =
expunged.peek().copied().unwrap_or(Uid::MAX);
match next_expunged.cmp(&uid) {
Ordering::Less => {
expunged.next();
}
Ordering::Equal => {
if index < max_index {
expunged_pairs
.push((Seqnum::from_index(index), uid));
} else {
unapplied_create -= 1;
stillborn.push(uid);
}
index += 1;
return false;
}
Ordering::Greater => {
index += 1;
return true;
}
}
});
}
let first_new = self.extant_messages.len() - unapplied_create;
let new_message_pairs: Vec<(Seqnum, Uid)> = self.extant_messages
[first_new..]
.iter()
.copied()
.enumerate()
.map(|(ix, uid)| (Seqnum::from_index(ix + first_new), uid))
.collect();
self.unapplied_create = 0;
self.report_max_modseq = self.max_modseq;
FlushResponse {
new: new_message_pairs,
expunged: expunged_pairs,
stillborn,
max_modseq: self.max_modseq,
}
}
pub fn has_pending_expunge(&self) -> bool {
!self.unapplied_expunge.is_empty()
}
pub fn take_changed_flags_uids(&mut self) -> Vec<Uid> {
self.changed_flags_uids.sort_unstable();
self.changed_flags_uids.dedup();
mem::replace(&mut self.changed_flags_uids, Vec::new())
}
pub fn add_changed_flags_uid(&mut self, uid: Uid) {
self.changed_flags_uids.push(uid);
}
pub fn add_changed_flags_uids(&mut self, uids: &[Uid]) {
self.changed_flags_uids.extend_from_slice(uids);
}
pub fn num_messages(&self) -> usize {
self.extant_messages.len() - self.unapplied_create
}
pub fn seqnum_to_uid(&self, seqnum: Seqnum) -> Result<Uid, Error> {
self.extant_messages[..self.num_messages()]
.get(seqnum.to_index())
.copied()
.ok_or(Error::NxMessage)
}
pub fn seqnum_range_to_uid(
&self,
seqnums: &SeqRange<Seqnum>,
silent: bool,
) -> Result<SeqRange<Uid>, Error> {
let mut ret = SeqRange::new();
for seqnum in seqnums.items(u32::MAX) {
match self.seqnum_to_uid(seqnum) {
Ok(uid) => ret.append(uid),
Err(_) if silent => (),
Err(e) => return Err(e),
}
}
Ok(ret)
}
pub fn uid_to_seqnum(&self, uid: Uid) -> Result<Seqnum, Error> {
if self.max_modseq.map(|m| uid > m.uid()).unwrap_or(true) {
return Err(Error::NxMessage);
}
self.extant_messages[..self.num_messages()]
.binary_search(&uid)
.map(Seqnum::from_index)
.map_err(|_| {
if self.extant_messages[self.num_messages()..]
.binary_search(&uid)
.is_ok()
{
Error::UnaddressableMessage
} else {
Error::ExpungedMessage
}
})
}
pub fn uid_range_to_seqnum(
&self,
uids: &SeqRange<Uid>,
silent: bool,
) -> Result<SeqRange<Seqnum>, Error> {
let mut ret = SeqRange::new();
for uid in uids.items(u32::MAX) {
match self.uid_to_seqnum(uid) {
Ok(seqnum) => ret.append(seqnum),
Err(_) if silent => (),
Err(e) => return Err(e),
}
}
Ok(ret)
}
pub fn is_assigned_uid(&self, uid: Uid) -> bool {
self.extant_messages[..self.num_messages()]
.binary_search(&uid)
.is_ok()
}
pub fn missing_uid_error(&self, uid: Uid) -> Error {
if self.max_modseq.map(|m| uid > m.uid()).unwrap_or(true) {
Error::NxMessage
} else {
Error::ExpungedMessage
}
}
pub fn max_uid_val(&self) -> u32 {
self.max_modseq.map(|m| m.uid().0.get()).unwrap_or(0)
}
pub fn flag_id_mut(&mut self, flag: Flag) -> FlagId {
FlagId(self.flag_ix_mut(flag))
}
pub fn flag_id(&self, flag: &Flag) -> Option<FlagId> {
self.flag_ix(flag).map(FlagId)
}
pub fn flag(&self, flag_id: FlagId) -> Option<&Flag> {
self.flags.get(flag_id.0)
}
pub fn flags<'a>(
&'a self,
) -> impl Iterator<Item = (FlagId, &'a Flag)> + 'a {
self.flags.iter().enumerate().map(|(ix, f)| (FlagId(ix), f))
}
pub fn test_flag(&self, flag: FlagId, message: Uid) -> bool {
if let Some(status) = self.message_status.get(&message) {
status.flags.contains(flag.0)
} else {
false
}
}
#[cfg(test)]
pub fn test_flag_o(&self, flag: &Flag, message: Uid) -> bool {
self.flag_id(flag)
.map(|f| self.test_flag(f, message))
.unwrap_or(false)
}
pub fn set_recent(&mut self, uid: Uid) {
if let Some(status) = self.message_status.get_mut(&uid) {
status.recent = true;
}
}
pub fn is_recent(&self, uid: Uid) -> bool {
self.message_status
.get(&uid)
.map(|m| m.recent)
.unwrap_or(false)
}
pub fn message_status(&self, uid: Uid) -> Option<&MessageStatus> {
self.message_status.get(&uid)
}
pub fn uids<'a>(&'a self) -> impl Iterator<Item = Uid> + 'a {
self.extant_messages[..self.num_messages()].iter().copied()
}
pub fn seqnums_uids<'a>(
&'a self,
) -> impl Iterator<Item = (Seqnum, Uid)> + 'a {
self.uids()
.enumerate()
.map(|(ix, uid)| (Seqnum::from_index(ix), uid))
}
pub fn max_modseq(&self) -> Option<Modseq> {
self.max_modseq
}
pub fn report_max_modseq(&self) -> Option<Modseq> {
self.report_max_modseq
}
pub fn next_uid(&self) -> Option<Uid> {
self.max_modseq
.map(|m| m.uid().next())
.unwrap_or(Some(Uid::MIN))
}
pub fn next_cid(&self) -> Option<Cid> {
self.max_modseq
.map(|m| m.cid().next())
.unwrap_or(Some(Cid::MIN))
}
pub fn max_uid(&self) -> Option<Uid> {
self.max_modseq.map(Modseq::uid)
}
pub fn max_seqnum(&self) -> Option<Seqnum> {
Seqnum::of(self.num_messages() as u32)
}
pub fn qresync(
&self,
resync_from: Option<Modseq>,
mut filter: impl FnMut(&Uid) -> bool,
seqnum_reference: impl IntoIterator<Item = Seqnum>,
uid_reference: impl IntoIterator<Item = Uid>,
) -> QresyncResponse {
let max_modseq = match self.max_modseq {
Some(m) => m,
None => return QresyncResponse::default(),
};
let expunged = if self
.recent_expungements
.front()
.copied()
.and_then(|(m, _)| resync_from.map(|r| r >= m))
.unwrap_or(false)
{
let resync_from = resync_from.unwrap();
let mut uids = self
.recent_expungements
.iter()
.copied()
.filter(|&(m, _)| m > resync_from)
.map(|(_, uid)| uid)
.filter(&mut filter)
.collect::<Vec<_>>();
uids.sort_unstable();
uids.dedup();
let mut s = SeqRange::new();
for uid in uids {
s.append(uid);
}
s
} else {
let expunge_start = seqnum_reference
.into_iter()
.zip(uid_reference.into_iter())
.take_while(|&(seqnum, uid)| {
Some(uid) == self.seqnum_to_uid(seqnum).ok()
})
.last()
.map(|(_, uid)| uid)
.unwrap_or(Uid::MIN);
let mut s = SeqRange::new();
for uid in (expunge_start.0.get()..=max_modseq.uid().0.get())
.map(|uid| Uid::of(uid).unwrap())
.filter(|uid| !self.message_status.contains_key(uid))
.filter(&mut filter)
{
s.append(uid);
}
s
};
let changed = self
.message_status
.iter()
.filter(|&(_, status)| {
resync_from
.map(|r| status.last_modified > r)
.unwrap_or(true)
})
.map(|(&uid, _)| uid)
.filter(filter)
.collect::<Vec<_>>();
QresyncResponse { expunged, changed }
}
pub fn uids_expunged_since<'a>(
&'a self,
since: Modseq,
) -> Option<impl Iterator<Item = Uid> + 'a> {
if self
.recent_expungements
.front()
.copied()
.map(|(m, _)| since >= m)
.unwrap_or(false)
{
Some(
self.recent_expungements
.iter()
.copied()
.filter(move |&(m, _)| m > since)
.map(|(_, u)| u),
)
} else {
None
}
}
pub fn drain_soft_expunged(&mut self, before: DateTime<Utc>) -> Vec<Uid> {
let ret = self
.soft_expungements
.iter()
.filter(|s| s.deadline < before)
.map(|s| s.uid)
.collect::<Vec<_>>();
self.soft_expungements.retain(|s| s.deadline >= before);
ret
}
pub fn silent_expunge(&mut self, uid: Uid) {
self.unapplied_expunge.push(uid);
}
fn set_flag(
&mut self,
uid: Uid,
canonical_modseq: Modseq,
flag: Flag,
val: bool,
) {
let flag = self.flag_ix_mut(flag);
if let Some(status) = self.message_status.get_mut(&uid) {
if val {
status.flags.insert(flag);
} else {
status.flags.remove(flag);
}
status.last_modified = canonical_modseq;
}
self.changed_flags_uids.push(uid);
}
fn expunge(
&mut self,
deadline: DateTime<Utc>,
uid: Uid,
canonical_modseq: Modseq,
) {
self.soft_expungements
.push(SoftExpungement { deadline, uid });
self.recent_expungements.push_back((canonical_modseq, uid));
while self.recent_expungements.len() > MAX_RECENT_EXPUNGEMENTS {
self.recent_expungements.pop_front();
}
if let Some(status) = self.message_status.get_mut(&uid) {
status.last_modified = canonical_modseq;
self.unapplied_expunge.push(uid);
if let Some(back) = self.recent_expungements.back() {
assert!(canonical_modseq >= back.0);
}
}
}
fn flag_ix_mut(&mut self, flag: Flag) -> usize {
if let Some(ix) = self.flag_ix(&flag) {
ix
} else {
let ix = self.flags.len();
self.flags.push(flag);
ix
}
}
fn flag_ix(&self, flag: &Flag) -> Option<usize> {
self.flags
.iter()
.enumerate()
.find(|&(_, f)| f == flag)
.map(|(ix, _)| ix)
}
}
impl MessageStatus {
pub fn is_recent(&self) -> bool {
self.recent
}
pub fn flags<'a>(&'a self) -> impl Iterator<Item = FlagId> + 'a {
self.flags.iter().map(FlagId)
}
pub fn test_flag(&self, flag: FlagId) -> bool {
self.flags.contains(flag.0)
}
pub fn last_modified(&self) -> Modseq {
self.last_modified
}
}
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct StateTransaction {
max_uid: Uid,
ops: Vec<StateMutation>,
}
#[derive(Serialize, Deserialize, Debug, Clone)]
enum StateMutation {
AddFlag(Uid, Flag),
RmFlag(Uid, Flag),
Expunge(
#[serde(with = "chrono::serde::ts_seconds")] DateTime<Utc>,
Uid,
),
}
impl StateTransaction {
pub fn new_unordered(uid: Uid) -> Self {
StateTransaction {
max_uid: uid,
ops: vec![],
}
}
pub fn add_flag(&mut self, uid: Uid, flag: Flag) {
assert!(uid <= self.max_uid);
self.ops.push(StateMutation::AddFlag(uid, flag));
}
pub fn rm_flag(&mut self, uid: Uid, flag: Flag) {
assert!(uid <= self.max_uid);
self.ops.push(StateMutation::RmFlag(uid, flag));
}
pub fn expunge(&mut self, deadline: DateTime<Utc>, uid: Uid) {
assert!(uid <= self.max_uid);
self.ops.push(StateMutation::Expunge(deadline, uid));
}
pub fn is_empty(&self) -> bool {
self.ops.is_empty()
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct FlushResponse {
pub new: Vec<(Seqnum, Uid)>,
pub expunged: Vec<(Seqnum, Uid)>,
pub stillborn: Vec<Uid>,
pub max_modseq: Option<Modseq>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct FlagId(usize);
#[cfg(test)]
mod test {
use std::collections::BTreeSet;
use proptest::prelude::*;
use super::*;
#[test]
fn seqnum_mapping() {
let mut state = MailboxState::new();
let flush = state.flush();
assert!(flush.new.is_empty());
assert!(flush.expunged.is_empty());
assert_eq!(None, flush.max_modseq);
state.seen(Uid::u(3));
let flush = state.flush();
assert_eq!(
vec![
(Seqnum::u(1), Uid::u(1)),
(Seqnum::u(2), Uid::u(2)),
(Seqnum::u(3), Uid::u(3))
],
flush.new
);
assert!(flush.expunged.is_empty());
assert_eq!(
Some(Modseq::new(Uid::u(3), Cid::GENESIS)),
flush.max_modseq
);
assert_eq!(Some(Uid::u(1)), state.seqnum_to_uid(Seqnum::u(1)).ok());
assert_eq!(Some(Uid::u(2)), state.seqnum_to_uid(Seqnum::u(2)).ok());
assert_eq!(Some(Uid::u(3)), state.seqnum_to_uid(Seqnum::u(3)).ok());
assert_eq!(None, state.seqnum_to_uid(Seqnum::u(4)).ok());
assert_eq!(Some(Seqnum::u(1)), state.uid_to_seqnum(Uid::u(1)).ok());
assert_eq!(Some(Seqnum::u(2)), state.uid_to_seqnum(Uid::u(2)).ok());
assert_eq!(Some(Seqnum::u(3)), state.uid_to_seqnum(Uid::u(3)).ok());
assert_eq!(None, state.uid_to_seqnum(Uid::u(4)).ok());
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(2));
state.commit(cid, tx);
assert_eq!(3, state.num_messages());
state.seen(Uid::u(5));
assert_eq!(3, state.num_messages());
assert_eq!(Some(Uid::u(1)), state.seqnum_to_uid(Seqnum::u(1)).ok());
assert_eq!(Some(Uid::u(2)), state.seqnum_to_uid(Seqnum::u(2)).ok());
assert_eq!(Some(Uid::u(3)), state.seqnum_to_uid(Seqnum::u(3)).ok());
assert_eq!(None, state.seqnum_to_uid(Seqnum::u(4)).ok());
assert_eq!(Some(Seqnum::u(1)), state.uid_to_seqnum(Uid::u(1)).ok());
assert_eq!(Some(Seqnum::u(2)), state.uid_to_seqnum(Uid::u(2)).ok());
assert_eq!(Some(Seqnum::u(3)), state.uid_to_seqnum(Uid::u(3)).ok());
assert_eq!(None, state.uid_to_seqnum(Uid::u(4)).ok());
assert_eq!(
Some(Modseq::new(Uid::u(3), Cid::GENESIS)),
state.report_max_modseq()
);
let flush = state.flush();
assert_eq!(
vec![(Seqnum::u(3), Uid::u(4)), (Seqnum::u(4), Uid::u(5))],
flush.new
);
assert_eq!(vec![(Seqnum::u(2), Uid::u(2))], flush.expunged);
assert_eq!(Some(Modseq::new(Uid::u(5), cid)), flush.max_modseq);
assert_eq!(4, state.num_messages());
assert_eq!(Some(Uid::u(1)), state.seqnum_to_uid(Seqnum::u(1)).ok());
assert_eq!(Some(Uid::u(3)), state.seqnum_to_uid(Seqnum::u(2)).ok());
assert_eq!(Some(Uid::u(4)), state.seqnum_to_uid(Seqnum::u(3)).ok());
assert_eq!(Some(Uid::u(5)), state.seqnum_to_uid(Seqnum::u(4)).ok());
assert_eq!(None, state.seqnum_to_uid(Seqnum::u(5)).ok());
assert_eq!(Some(Seqnum::u(1)), state.uid_to_seqnum(Uid::u(1)).ok());
assert_eq!(None, state.uid_to_seqnum(Uid::u(2)).ok());
assert_eq!(Some(Seqnum::u(2)), state.uid_to_seqnum(Uid::u(3)).ok());
assert_eq!(Some(Seqnum::u(3)), state.uid_to_seqnum(Uid::u(4)).ok());
assert_eq!(Some(Seqnum::u(4)), state.uid_to_seqnum(Uid::u(5)).ok());
assert_eq!(None, state.uid_to_seqnum(Uid::u(6)).ok());
assert_eq!(
Some(Modseq::new(Uid::u(5), cid)),
state.report_max_modseq()
);
}
#[test]
fn expunged_new_messages() {
let mut state = MailboxState::new();
state.seen(Uid::u(5));
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(1));
tx.expunge(Utc::now(), Uid::u(3));
tx.expunge(Utc::now(), Uid::u(5));
state.commit(cid, tx);
let flush = state.flush();
assert_eq!(
vec![(Seqnum::u(1), Uid::u(2)), (Seqnum::u(2), Uid::u(4))],
flush.new
);
assert!(flush.expunged.is_empty());
assert_eq!(vec![Uid::u(1), Uid::u(3), Uid::u(5)], flush.stillborn);
state.seen(Uid::u(7));
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(4));
tx.expunge(Utc::now(), Uid::u(6));
state.commit(cid, tx);
let flush = state.flush();
assert_eq!(vec![(Seqnum::u(2), Uid::u(7))], flush.new);
assert_eq!(vec![(Seqnum::u(2), Uid::u(4))], flush.expunged);
assert_eq!(vec![Uid::u(6)], flush.stillborn);
}
#[test]
fn foreign_commit_updates_max_uid() {
let mut state = MailboxState::new();
state.seen(Uid::u(2));
state.flush();
let mut state2 = MailboxState::new();
state2.seen(Uid::u(4));
let (cid, mut tx) = state2.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(4));
state.commit(cid, tx);
let flush = state.flush();
assert_eq!(vec![(Seqnum::u(3), Uid::u(3))], flush.new);
assert_eq!(vec![Uid::u(4)], flush.stillborn);
}
#[test]
fn seqnum_to_uid_error_states() {
let mut state = MailboxState::new();
state.seen(Uid::u(3));
state.flush();
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(2));
state.commit(cid, tx);
state.seen(Uid::u(4));
assert!(matches!(
state.seqnum_to_uid(Seqnum::u(4)),
Err(Error::NxMessage)
));
assert!(matches!(
state.seqnum_to_uid(Seqnum::u(5)),
Err(Error::NxMessage)
));
}
#[test]
fn uid_to_seqnum_error_states() {
let mut state = MailboxState::new();
state.seen(Uid::u(3));
state.flush();
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), Uid::u(2));
state.commit(cid, tx);
state.flush();
state.seen(Uid::u(4));
assert!(matches!(
state.uid_to_seqnum(Uid::u(2)),
Err(Error::ExpungedMessage)
));
assert!(matches!(
state.uid_to_seqnum(Uid::u(4)),
Err(Error::UnaddressableMessage)
));
assert!(matches!(
state.uid_to_seqnum(Uid::u(5)),
Err(Error::NxMessage)
));
}
#[test]
fn different_tx_uid_orders_produce_consistent_results() {
let mut state1 = MailboxState::new();
let mut state2 = MailboxState::new();
state1.seen(Uid::u(1));
state2.seen(Uid::u(1));
state1.flush();
state2.flush();
state1.seen(Uid::u(2));
state1.flush();
let (cid, mut tx) = state2.start_tx().unwrap();
tx.add_flag(Uid::u(1), Flag::Deleted);
state2.commit(cid, tx.clone());
state2.flush();
state1.commit(cid, tx.clone());
state1.flush();
assert_eq!(
state1.message_status(Uid::u(1)).unwrap().last_modified,
state2.message_status(Uid::u(1)).unwrap().last_modified,
);
}
#[test]
fn seqnum_mapping_not_confused_by_intermediate_transactions() {
let mut state = MailboxState::new();
state.seen(Uid::u(1));
let (cid, mut tx) = state.start_tx().unwrap();
tx.add_flag(Uid::u(1), Flag::Deleted);
state.seen(Uid::u(2));
state.commit(cid, tx);
state.flush();
assert_eq!(Seqnum::u(1), state.uid_to_seqnum(Uid::u(1)).unwrap());
assert_eq!(Seqnum::u(2), state.uid_to_seqnum(Uid::u(2)).unwrap());
assert_eq!(Uid::u(1), state.seqnum_to_uid(Seqnum::u(1)).unwrap());
assert_eq!(Uid::u(2), state.seqnum_to_uid(Seqnum::u(2)).unwrap());
}
proptest! {
#[test]
fn qresync_always_finds_all_expungements(
expunge_before in prop::collection::vec(1u32..1000u32, 0..100),
expunge_after in prop::collection::vec(1u32..1000u32, 0..100),
) {
let mut state = MailboxState::new();
let mut all_expunged = BTreeSet::<Uid>::new();
for uid in expunge_before {
let uid = Uid::u(uid);
state.seen(uid.next().unwrap());
all_expunged.insert(uid);
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), uid);
state.commit(cid, tx);
}
state.flush();
let mut client_expunged = all_expunged.clone();
let resync_point = state.report_max_modseq();
let resync_checkpoints = state.seqnums_uids()
.filter(|&(seq, _)| 0 == seq.0.get() % 10)
.collect::<Vec<_>>();
for uid in expunge_after {
let uid = Uid::u(uid);
state.seen(uid.next().unwrap());
all_expunged.insert(uid);
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(Utc::now(), uid);
state.commit(cid, tx);
}
state.flush();
let qresync = state.qresync(
resync_point,
|_| true,
resync_checkpoints.iter().copied().map(|(s, _)| s),
resync_checkpoints.iter().copied().map(|(_, u)| u));
for uid in qresync.expunged.items(u32::MAX) {
client_expunged.insert(uid);
}
prop_assert_eq!(all_expunged, client_expunged);
}
}
#[test]
fn flag_operations() {
let mut state = MailboxState::new();
state.seen(Uid::u(10));
state.flush();
let flagged_flag = state.flag_id_mut(Flag::Flagged);
let seen_flag = state.flag_id_mut(Flag::Seen);
let kw_flag = state.flag_id_mut(Flag::Keyword("NotJunk".to_owned()));
assert!(!state.test_flag(flagged_flag, Uid::u(1)));
assert!(!state.test_flag(seen_flag, Uid::u(1)));
assert!(!state.test_flag(kw_flag, Uid::u(1)));
assert!(!state.test_flag(flagged_flag, Uid::u(11)));
assert!(!state.test_flag(seen_flag, Uid::u(11)));
assert!(!state.test_flag(kw_flag, Uid::u(11)));
let (cid, mut tx) = state.start_tx().unwrap();
tx.add_flag(Uid::u(1), Flag::Seen);
tx.add_flag(Uid::u(2), Flag::Flagged);
state.commit(cid, tx);
assert!(state.test_flag(flagged_flag, Uid::u(2)));
let (cid, mut tx) = state.start_tx().unwrap();
tx.add_flag(Uid::u(3), Flag::Keyword("NotJunk".to_owned()));
tx.rm_flag(Uid::u(2), Flag::Flagged);
state.commit(cid, tx);
assert!(!state.test_flag(flagged_flag, Uid::u(1)));
assert!(!state.test_flag(flagged_flag, Uid::u(2)));
assert!(!state.test_flag(flagged_flag, Uid::u(3)));
assert!(state.test_flag(seen_flag, Uid::u(1)));
assert!(!state.test_flag(seen_flag, Uid::u(2)));
assert!(!state.test_flag(seen_flag, Uid::u(3)));
assert!(!state.test_flag(kw_flag, Uid::u(1)));
assert!(!state.test_flag(kw_flag, Uid::u(2)));
assert!(state.test_flag(kw_flag, Uid::u(3)));
assert_eq!(
vec![Uid::u(1), Uid::u(2), Uid::u(3)],
state.take_changed_flags_uids()
);
assert!(state.take_changed_flags_uids().is_empty());
}
#[test]
fn test_drain_soft_expunged() {
let mut state = MailboxState::new();
state.seen(Uid::u(10));
state.flush();
let date = Date::<Utc>::from_utc(NaiveDate::from_ymd(2020, 1, 1), Utc);
let (cid, mut tx) = state.start_tx().unwrap();
tx.expunge(date.and_hms(1, 1, 1), Uid::u(2));
tx.expunge(date.and_hms(2, 2, 2), Uid::u(3));
tx.expunge(date.and_hms(3, 3, 3), Uid::u(1));
state.commit(cid, tx);
let mut se = state.drain_soft_expunged(date.and_hms(2, 59, 59));
se.sort();
assert_eq!(vec![Uid::u(2), Uid::u(3)], se);
assert!(state
.drain_soft_expunged(date.and_hms(2, 59, 59))
.is_empty());
assert!(state.drain_soft_expunged(date.and_hms(3, 3, 3)).is_empty());
assert_eq!(
vec![Uid::u(1)],
state.drain_soft_expunged(date.and_hms(3, 3, 4))
);
}
}