use std::fs;
use std::path::PathBuf;
use std::sync::Arc;
use chrono::prelude::*;
use log::error;
use super::super::storage;
use super::defs::*;
use crate::{
account::{key_store::KeyStore, model::*},
support::{error::Error, log_prefix::LogPrefix, small_bitset::SmallBitset},
};
pub struct DeliveryAccount {
deliverydb: storage::DeliveryDb,
key_store: KeyStore,
message_store: storage::MessageStore,
common_paths: Arc<CommonPaths>,
log_prefix: LogPrefix,
}
impl DeliveryAccount {
pub fn new(log_prefix: LogPrefix, root: PathBuf) -> Result<Self, Error> {
let common_paths = Arc::new(CommonPaths {
tmp: root.join("tmp"),
garbage: root.join("garbage"),
});
let key_store = KeyStore::new(
log_prefix.clone(),
root.join("keys"),
common_paths.tmp.clone(),
None,
);
let deliverydb_path = root.join("delivery.sqlite");
let deliverydb =
storage::DeliveryDb::new(&log_prefix, &deliverydb_path)?;
let message_store = storage::MessageStore::new(root.join("messages"));
Ok(Self {
deliverydb,
key_store,
message_store,
common_paths,
log_prefix,
})
}
pub fn buffer_message(
&mut self,
data: impl std::io::Read,
) -> Result<BufferedMessage, Error> {
super::messages::buffer_message(
&mut self.key_store,
&self.common_paths,
Utc::now().into(),
data,
)
}
pub fn deliver(
&mut self,
mailbox: &str,
flags: &[Flag],
data: impl std::io::Read,
) -> Result<(), Error> {
let buffered = self.buffer_message(data)?;
self.deliver_buffered(mailbox, flags, &buffered)
}
pub fn deliver_buffered(
&mut self,
mailbox: &str,
flags: &[Flag],
message: &BufferedMessage,
) -> Result<(), Error> {
let canonical_path = fs::File::open(&message.0)
.and_then(storage::MessageStore::canonical_path)?;
self.message_store.insert(&message.0, &canonical_path)?;
self.deliverydb.queue_delivery(&storage::Delivery {
path: canonical_path
.into_os_string()
.into_string()
.expect("all canonical paths are UTF-8"),
mailbox: mailbox.to_owned(),
flags: flags.to_owned(),
savedate: storage::UnixTimestamp::now(),
})?;
if let Err(e) = self.deliverydb.clear_old_deliveries() {
error!("{} Failed to clear old deliveries: {e:?}", self.log_prefix,);
}
Ok(())
}
}
impl Account {
pub fn drain_deliveries(&mut self) {
loop {
let delivery = match self.deliverydb.pop_delivery() {
Ok(None) => return,
Ok(Some(d)) => d,
Err(e) => {
error!("{} Failed to pop delivery: {e:?}", self.log_prefix);
return;
},
};
let inbox_id = match self.metadb.find_mailbox("INBOX") {
Ok(id) => id,
Err(e) => {
error!(
"{} Failed to look up INBOX for delivery: {e:?}",
self.log_prefix,
);
return;
},
};
let dst_id = match self.metadb.find_mailbox(&delivery.mailbox) {
Ok(id) => id,
Err(e) => {
error!(
"{} Delivering message to INBOX instead of '{}': \
{e:?}",
self.log_prefix, delivery.mailbox,
);
inbox_id
},
};
let mut flags = SmallBitset::new();
for flag in delivery.flags {
if let Ok(flag_id) = self.metadb.intern_flag(&flag) {
flags.insert(flag_id.0);
}
}
let r = self.metadb.intern_and_append_mailbox_messages(
dst_id,
&mut [(delivery.path.as_str(), Some(&flags))].into_iter(),
);
if let Err(e) = r {
error!(
"{} Failed to deliver message to '{}', \
retrying with INBOX: {e:?}",
self.log_prefix, delivery.mailbox,
);
let r = self.metadb.intern_and_append_mailbox_messages(
inbox_id,
&mut [(delivery.path.as_str(), Some(&flags))].into_iter(),
);
if let Err(e) = r {
error!(
"{} Failed to deliver message to INBOX: {e:?}",
self.log_prefix,
);
return;
}
}
}
}
}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn deliver_success() {
let mut fixture = TestFixture::new();
let mut delivery = DeliveryAccount::new(
LogPrefix::new("delivery".to_owned()),
fixture.root.path().to_owned(),
)
.unwrap();
let buf1 = delivery.buffer_message(b"foobar" as &[u8]).unwrap();
let buf2 = delivery.buffer_message(b"barfoo" as &[u8]).unwrap();
delivery.deliver_buffered("INBOX", &[], &buf1).unwrap();
delivery
.deliver_buffered(
"iNbOx",
&[Flag::Flagged, Flag::Keyword("foo".to_owned())],
&buf2,
)
.unwrap();
delivery.deliver_buffered("Archive", &[], &buf1).unwrap();
let (mb, _) = fixture.select("INBOX", false, None).unwrap();
assert_eq!(2, mb.select_response().unwrap().exists);
assert!(mb.test_flag_o(&Flag::Flagged, Uid::u(2)));
assert!(mb.test_flag_o(&Flag::Keyword("foo".to_owned()), Uid::u(2)));
let (mb, _) = fixture.select("Archive", false, None).unwrap();
assert_eq!(1, mb.select_response().unwrap().exists);
}
#[test]
fn deliver_bad_destination() {
let mut fixture = TestFixture::new();
let mut delivery = DeliveryAccount::new(
LogPrefix::new("delivery".to_owned()),
fixture.root.path().to_owned(),
)
.unwrap();
fixture.create("noselect/foo");
fixture.delete("noselect").unwrap();
let buf1 = delivery.buffer_message(b"foobar" as &[u8]).unwrap();
delivery.deliver_buffered("", &[], &buf1).unwrap();
delivery.deliver_buffered("noselect", &[], &buf1).unwrap();
delivery
.deliver_buffered("nonexistent", &[], &buf1)
.unwrap();
let (mb, _) = fixture.select("INBOX", false, None).unwrap();
assert_eq!(3, mb.select_response().unwrap().exists);
}
}