use std::fmt::Write as _;
use std::fs;
use std::os::unix::fs::PermissionsExt;
use std::path::Path;
use std::time::Duration;
use chrono::prelude::*;
use rusqlite::OptionalExtension as _;
use super::types::*;
use crate::support::{error::Error, log_prefix::LogPrefix};
pub struct Connection {
cxn: rusqlite::Connection,
}
static MIGRATIONS: &[&str] = &[include_str!("deliverydb.v1.sql")];
impl Connection {
pub fn new(log_prefix: &LogPrefix, path: &Path) -> Result<Self, Error> {
let mut cxn = rusqlite::Connection::open_with_flags(
path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_WRITE
| rusqlite::OpenFlags::SQLITE_OPEN_CREATE
| rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)?;
let _ = fs::set_permissions(
path,
fs::Permissions::from_mode(0o660),
);
cxn.pragma_update(None, "foreign_keys", true)?;
cxn.pragma_update(None, "journal_mode", "PERSIST")?;
cxn.pragma_update(None, "journal_size_limit", 64 * 1024)?;
cxn.busy_timeout(Duration::from_secs(10))?;
super::db_migrations::apply_migrations(
log_prefix, &mut cxn, "delivery", MIGRATIONS,
)?;
Ok(Self { cxn })
}
pub fn queue_delivery(&mut self, delivery: &Delivery) -> Result<(), Error> {
let mut flags = String::new();
for flag in &delivery.flags {
if !flags.is_empty() {
flags.push(' ');
}
let _ = write!(flags, "{}", flag);
}
self.cxn.execute(
"INSERT INTO `delivery` (`path`, `mailbox`, `flags`, `savedate`) \
VALUES (?, ?, ?, ?)",
(&delivery.path, &delivery.mailbox, &flags, delivery.savedate),
)?;
Ok(())
}
pub fn pop_delivery(&mut self) -> Result<Option<Delivery>, Error> {
self.cxn
.prepare_cached(
"UPDATE `delivery` \
SET `delivered` = ?
WHERE ROWID = ( \
SELECT MIN(ROWID) FROM `delivery` WHERE `delivered` IS NULL \
) RETURNING *",
)?
.query_row((UnixTimestamp::now(),), from_row)
.optional()
.map_err(Into::into)
}
pub fn is_delivery(&mut self, path: &str) -> Result<bool, Error> {
self.cxn
.prepare_cached("SELECT 1 FROM `delivery` WHERE `path` = ?")?
.exists((path,))
.map_err(Into::into)
}
pub fn clear_old_deliveries(&mut self) -> Result<(), Error> {
self.cxn.execute(
"DELETE FROM `delivery` WHERE `delivered` < ?",
(UnixTimestamp(Utc::now() - chrono::Duration::hours(1)),),
)?;
Ok(())
}
}
#[cfg(test)]
mod test {
use super::*;
use crate::account::model::Flag;
use tempfile::TempDir;
#[test]
fn test_delivery() {
let tmpdir = TempDir::new().unwrap();
let mut cxn = Connection::new(
&LogPrefix::new("test".to_owned()),
&tmpdir.path().join("delivery.sqlite"),
)
.unwrap();
let delivery1 = Delivery {
path: "foo/bar".to_owned(),
mailbox: "INBOX".to_owned(),
flags: vec![Flag::Flagged, Flag::Keyword("foo".to_owned())],
savedate: UnixTimestamp(DateTime::from_timestamp(42, 0).unwrap()),
};
let delivery2 = Delivery {
path: "baz/quux".to_owned(),
mailbox: "Spam".to_owned(),
flags: vec![],
savedate: UnixTimestamp(DateTime::from_timestamp(54, 0).unwrap()),
};
cxn.queue_delivery(&delivery1).unwrap();
cxn.queue_delivery(&delivery2).unwrap();
assert!(cxn.is_delivery("foo/bar").unwrap());
assert!(cxn.is_delivery("baz/quux").unwrap());
assert!(!cxn.is_delivery("nonexistent").unwrap());
cxn.clear_old_deliveries().unwrap();
assert!(cxn.is_delivery("foo/bar").unwrap());
assert!(cxn.is_delivery("baz/quux").unwrap());
let mut popped = Vec::new();
popped.push(cxn.pop_delivery().unwrap().unwrap());
popped.push(cxn.pop_delivery().unwrap().unwrap());
assert_eq!(None, cxn.pop_delivery().unwrap());
assert!(popped.contains(&delivery1));
assert!(popped.contains(&delivery2));
assert!(cxn.is_delivery("foo/bar").unwrap());
assert!(cxn.is_delivery("baz/quux").unwrap());
cxn.clear_old_deliveries().unwrap();
assert!(cxn.is_delivery("foo/bar").unwrap());
assert!(cxn.is_delivery("baz/quux").unwrap());
}
}