use std::collections::HashSet;
use std::fs;
use std::io;
use std::path::PathBuf;
use std::sync::{
atomic::{AtomicBool, Ordering::SeqCst},
Arc,
};
use std::time::{Duration, SystemTime};
use log::{error, warn};
use super::defs::*;
use crate::account::mailbox_state::*;
use crate::account::model::*;
use crate::support::error::Error;
use crate::support::file_ops::IgnoreKinds;
use crate::support::threading;
const EXCESS_ROLLUP_THRESHOLD: usize = 4;
#[cfg(not(test))]
const OLD_ROLLUP_GRACE_PERIOD: Duration = Duration::from_secs(24 * 3600);
#[cfg(not(test))]
const EXCESS_ROLLUP_GRACE_PERIOD: Duration = Duration::from_secs(60);
#[cfg(test)]
const OLD_ROLLUP_GRACE_PERIOD: Duration = Duration::from_secs(2);
#[cfg(test)]
const EXCESS_ROLLUP_GRACE_PERIOD: Duration = Duration::from_secs(1);
impl StatelessMailbox {
pub fn select(self) -> Result<(StatefulMailbox, SelectResponse), Error> {
StatefulMailbox::select(self)
}
fn do_gc(&self, rollups: Vec<RollupInfo>) {
assert!(!self.read_only);
if let Err(err) = self.message_scheme().gc(
&self.common_paths.tmp,
&self.common_paths.garbage,
0,
) {
warn!(
"{} Error garbage collecting messages: {}",
self.log_prefix, err
);
return;
}
let expunge_before_cid = rollups
.iter()
.filter(|r| r.delete_transactions)
.map(|r| r.cid)
.max()
.unwrap_or(Cid(0));
if let Err(err) = self.change_scheme().gc(
&self.common_paths.tmp,
&self.common_paths.garbage,
expunge_before_cid.0,
) {
warn!(
"{} Error garbage collecting changes: {}",
self.log_prefix, err
);
} else {
for rollup in rollups {
if rollup.delete_rollup {
if let Err(err) =
fs::remove_file(&rollup.path).ignore_not_found()
{
warn!(
"{} Error removing {}: {}",
self.log_prefix,
rollup.path.display(),
err
);
}
}
}
}
}
}
impl StatefulMailbox {
fn select(s: StatelessMailbox) -> Result<(Self, SelectResponse), Error> {
let mut rollups = list_rollups(&s)?;
let state = rollups
.pop()
.and_then(|r| match s.read_state_file::<MailboxState>(&r.path) {
Ok(state) => Some(state),
Err(e) => {
error!(
"{} Error reading {}, starting from empty state: {}",
s.log_prefix,
r.path.display(),
e
);
None
}
})
.unwrap_or_else(MailboxState::new);
let mut this = Self {
recency_frontier: state.max_modseq().map(Modseq::uid),
s,
state,
fetch_loopbreaker: HashSet::new(),
client_known_flags: Vec::new(),
suggest_rollup: 0,
rollups_since_gc: 0,
gc_in_progress: Arc::new(AtomicBool::new(false)),
synchronous_gc: false,
};
this.state.init_transient();
this.poll()?;
if rollups.iter().any(|r| r.delete_rollup) {
this.start_gc(rollups, false);
}
this.state.flag_id_mut(Flag::Answered);
this.state.flag_id_mut(Flag::Deleted);
this.state.flag_id_mut(Flag::Draft);
this.state.flag_id_mut(Flag::Flagged);
this.state.flag_id_mut(Flag::Seen);
let select_response = SelectResponse {
flags: this.flags_response(),
exists: this.state.num_messages(),
recent: this.count_recent(),
unseen: this
.state
.seqnums_uids()
.find(|&(_, uid)| {
this.state
.flag_id(&Flag::Seen)
.map(|fid| !this.state.test_flag(fid, uid))
.unwrap_or(true)
})
.map(|(s, _)| s),
uidnext: this.state.next_uid().unwrap_or(Uid::MAX),
uidvalidity: this.s.uid_validity()?,
read_only: this.s.read_only,
max_modseq: this.state.report_max_modseq(),
};
Ok((this, select_response))
}
pub fn qresync(
&mut self,
request: QresyncRequest,
) -> Result<QresyncResponse, Error> {
if request.uid_validity != self.s.uid_validity()? {
return Ok(QresyncResponse::default());
}
let (seqnum_reference, uid_reference) =
request.mapping_reference.unwrap_or_default();
let known_uids = request.known_uids;
let response = self.state.qresync(
request.resync_from,
|&uid| known_uids.as_ref().map_or(true, |k| k.contains(uid)),
seqnum_reference.items(u32::MAX),
uid_reference.items(u32::MAX),
);
self.state.add_changed_flags_uids(&response.changed);
Ok(response)
}
fn start_gc(&self, rollups: Vec<RollupInfo>, force: bool) {
if self.s.read_only && !force {
return;
}
if Ok(false)
== self
.gc_in_progress
.compare_exchange(false, true, SeqCst, SeqCst)
{
return;
}
let s_clone = self.s.clone();
let gc_in_progress = Arc::clone(&self.gc_in_progress);
if self.synchronous_gc {
s_clone.do_gc(rollups);
gc_in_progress.store(false, SeqCst);
} else {
threading::run_in_background(move || {
s_clone.do_gc(rollups);
gc_in_progress.store(false, SeqCst);
});
}
}
pub fn schedule_gc(&self, force: bool) -> Result<(), Error> {
if self.s.read_only && !force {
return Ok(());
}
self.start_gc(list_rollups(&self.s)?, force);
Ok(())
}
pub(super) fn flags_response(&mut self) -> Vec<Flag> {
let ret: Vec<Flag> =
self.state.flags().map(|(_, f)| f.to_owned()).collect();
for flag in &ret {
if !self.client_known_flags.contains(flag) {
self.client_known_flags.push(flag.to_owned());
}
}
ret
}
}
pub(super) fn list_rollups(
s: &StatelessMailbox,
) -> Result<Vec<RollupInfo>, Error> {
match fs::read_dir(s.root.join("rollup")) {
Err(e) if io::ErrorKind::NotFound == e.kind() => Ok(vec![]),
Err(e) => Err(e.into()),
Ok(it) => {
let mut ret = Vec::new();
let now = SystemTime::now();
for entry in it {
let entry = entry?;
let modseq = match entry
.file_name()
.to_str()
.and_then(|n| u64::from_str_radix(n, 10).ok())
.and_then(Modseq::of)
{
Some(ms) => ms,
None => continue,
};
let md = match entry.metadata() {
Ok(md) => md,
Err(e) if io::ErrorKind::NotFound == e.kind() => continue,
Err(e) => return Err(e.into()),
};
ret.push(RollupInfo {
cid: modseq.cid(),
path: entry.path(),
age: md
.modified()
.ok()
.and_then(|modified| now.duration_since(modified).ok())
.unwrap_or(Duration::from_secs(0)),
delete_rollup: false,
delete_transactions: false,
});
}
classify_rollups(&mut ret);
Ok(ret)
}
}
}
fn classify_rollups(rollups: &mut [RollupInfo]) {
if rollups.is_empty() {
return;
}
rollups.sort_unstable();
let len = rollups.len();
for rollup in &mut rollups[..len - 1] {
if rollup.age >= OLD_ROLLUP_GRACE_PERIOD {
rollup.delete_rollup = true;
rollup.delete_transactions = true;
}
}
if len > EXCESS_ROLLUP_THRESHOLD {
for rollup in &mut rollups[..len - EXCESS_ROLLUP_THRESHOLD] {
if rollup.age >= EXCESS_ROLLUP_GRACE_PERIOD {
rollup.delete_rollup = true;
}
}
}
}
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub(super) struct RollupInfo {
cid: Cid,
age: Duration,
path: PathBuf,
delete_rollup: bool,
delete_transactions: bool,
}
#[cfg(test)]
mod test {
use super::super::test_prelude::*;
use super::*;
fn r(cid: u32, age_ms: u64) -> RollupInfo {
RollupInfo {
cid: Cid(cid),
path: PathBuf::new(),
age: Duration::from_millis(age_ms),
delete_rollup: false,
delete_transactions: false,
}
}
#[test]
fn classify_rollups_empty() {
classify_rollups(&mut []);
}
#[test]
fn classify_rollups_single_young() {
let mut rollups = [r(1234, 100)];
classify_rollups(&mut rollups);
assert_eq!([r(1234, 100)], rollups);
}
#[test]
fn classify_rollups_single_old() {
let mut rollups = [r(1234, 10_000_000)];
classify_rollups(&mut rollups);
assert_eq!([r(1234, 10_000_000)], rollups);
}
#[test]
fn classify_rollups_one_young_one_old() {
let mut rollups = [r(1000, 100), r(900, 10_000_000)];
classify_rollups(&mut rollups);
assert_eq!(
[
RollupInfo {
delete_rollup: true,
delete_transactions: true,
..r(900, 10_000_000)
},
r(1000, 100)
],
rollups
);
}
#[test]
fn classify_rollups_one_old_one_young() {
let mut rollups = [r(900, 10_000_000), r(1000, 100)];
classify_rollups(&mut rollups);
assert_eq!(
[
RollupInfo {
delete_rollup: true,
delete_transactions: true,
..r(900, 10_000_000)
},
r(1000, 100)
],
rollups
);
}
#[test]
fn classify_rollups_excess() {
let mut rollups = [
r(1, 5_000), r(2, 1_900), r(3, 1_800), r(4, 1_700), r(5, 1_600), r(6, 1_500), ];
classify_rollups(&mut rollups);
assert_eq!(
[
RollupInfo {
delete_rollup: true,
delete_transactions: true,
..r(1, 5_000)
},
RollupInfo {
delete_rollup: true,
..r(2, 1_900)
},
r(3, 1_800),
r(4, 1_700),
r(5, 1_600),
r(6, 1_500),
],
rollups
);
}
#[test]
fn resume_from_rollup() {
let setup = set_up();
let uid = simple_append(&setup.stateless);
{
let (mut mb1, _) = setup.stateless.clone().select().unwrap();
mb1.synchronous_gc = true;
mb1.store(&StoreRequest {
ids: &SeqRange::just(uid),
flags: &[Flag::Seen],
remove_listed: false,
remove_unlisted: false,
loud: false,
unchanged_since: None,
})
.unwrap();
for _ in 0..500 {
mb1.store(&StoreRequest {
ids: &SeqRange::just(uid),
flags: &[Flag::Flagged],
remove_listed: false,
remove_unlisted: false,
loud: false,
unchanged_since: None,
})
.unwrap();
mb1.store(&StoreRequest {
ids: &SeqRange::just(uid),
flags: &[Flag::Flagged],
remove_listed: true,
remove_unlisted: false,
loud: false,
unchanged_since: None,
})
.unwrap();
mb1.poll().unwrap();
std::thread::sleep(std::time::Duration::from_millis(16));
}
assert!(!mb1
.stateless()
.change_scheme()
.access_path_for_id(1)
.assume_exists()
.is_file());
}
let (mb2, _) = setup.stateless.clone().select().unwrap();
assert!(mb2.state.test_flag_o(&Flag::Seen, uid));
}
}