use std::fs;
use std::io::Write;
use std::path::Path;
use serde::{de::DeserializeOwned, Serialize};
use tempfile::NamedTempFile;
use super::defs::*;
use crate::account::mailbox_state::*;
use crate::crypt::data_stream;
use crate::support::compression::{Compression, FinishWrite};
use crate::support::error::Error;
impl StatelessMailbox {
pub(super) fn read_state_file<T: DeserializeOwned>(
&self,
src: &Path,
) -> Result<T, Error> {
let file = fs::File::open(src)?;
let stream = data_stream::Reader::new(file, |k| {
let mut ks = self.key_store.lock().unwrap();
ks.get_private_key(k)
})?;
let compression = stream.metadata.compression;
let stream = compression.decompressor(stream)?;
serde_cbor::from_reader(stream).map_err(|e| e.into())
}
pub(super) fn write_state_file(
&self,
data: &impl Serialize,
) -> Result<NamedTempFile, Error> {
let mut buffer_file = NamedTempFile::new_in(&self.common_paths.tmp)?;
{
let compression = Compression::DEFAULT_FOR_STATE;
let mut crypt_writer = {
let mut ks = self.key_store.lock().unwrap();
let (key_name, pub_key) = ks.get_default_public_key()?;
data_stream::Writer::new(
&mut buffer_file,
pub_key,
key_name.to_owned(),
compression,
)?
};
{
let mut compressor =
compression.compressor(&mut crypt_writer)?;
serde_cbor::to_writer(&mut compressor, data)?;
compressor.finish()?;
}
crypt_writer.flush()?;
}
buffer_file.as_file_mut().sync_all()?;
Ok(buffer_file)
}
}
impl StatefulMailbox {
pub(super) fn change_transaction<R>(
&mut self,
mut f: impl FnMut(&Self, &mut StateTransaction) -> Result<R, Error>,
) -> Result<R, Error> {
self.poll_for_new_changes()?;
for _ in 0..1000 {
let (cid, mut tx) = self.state.start_tx()?;
let res = f(self, &mut tx)?;
if tx.is_empty() {
return Ok(res);
}
let buffer_file = self.s.write_state_file(&tx)?;
if self.s.change_scheme().emplace(buffer_file.path(), cid.0)? {
self.state.commit(cid, tx);
self.see_cid(cid);
self.s.notify_all_best_effort();
return Ok(res);
}
self.poll_for_new_changes()?;
}
Err(Error::GaveUpInsertion)
}
}
#[cfg(test)]
mod test {
use chrono::prelude::*;
use super::super::test_prelude::*;
use super::*;
use crate::account::model::*;
#[test]
fn write_and_read_state_files() {
let setup = set_up();
assert_eq!(
Some(Cid(1)),
setup
.stateless
.set_flags_blind(vec![(Uid::u(1), vec![(true, Flag::Flagged)])])
.unwrap()
);
let tx: StateTransaction = setup
.stateless
.read_state_file(
&setup
.stateless
.change_scheme()
.access_path_for_id(1)
.assume_exists(),
)
.unwrap();
assert!(!tx.is_empty());
}
#[test]
fn append_with_flags() {
let setup = set_up();
let (mut mb1, _) = setup.stateless.select().unwrap();
let uid = mb1
.stateless()
.append(
FixedOffset::east(0)
.from_utc_datetime(&Utc::now().naive_local()),
vec![Flag::Flagged, Flag::Keyword("foo".to_owned())],
&mut "foobar".as_bytes(),
)
.unwrap();
mb1.poll().unwrap();
assert!(mb1.state.test_flag_o(&Flag::Flagged, uid));
assert!(mb1.state.test_flag_o(&Flag::Keyword("foo".to_owned()), uid));
}
#[test]
fn misordered_blind_change_accepted() {
let setup = set_up();
let (mut mb1, _) = setup.stateless.select().unwrap();
let uid1 = simple_append(mb1.stateless());
let uid2 = simple_append(mb1.stateless());
mb1.poll().unwrap();
mb1.store(&StoreRequest {
ids: &SeqRange::just(uid2),
flags: &[Flag::Flagged],
remove_listed: false,
remove_unlisted: false,
loud: false,
unchanged_since: None,
})
.unwrap();
mb1.poll().unwrap();
mb1.s
.set_flags_blind(vec![(uid1, vec![(true, Flag::Seen)])])
.unwrap();
let poll = mb1.poll().unwrap();
assert_eq!(Some(Modseq::new(uid2, Cid(2))), poll.max_modseq);
assert_eq!(
Modseq::new(uid2, Cid(2)),
mb1.state.message_status(uid1).unwrap().last_modified()
);
}
}