use std::path::PathBuf;
use std::sync::mpsc::{Receiver, RecvError, Sender, TryRecvError, channel};
use std::sync::{Arc, Mutex};
use std::thread::JoinHandle;
use anyhow::{Result, anyhow};
use rmut_core::config::Account;
use rmut_core::maildir::Flags;
use rmut_core::remote::{Progress, Remote};
pub enum Job {
FetchBodies(Vec<PathBuf>),
Sync {
flags: Vec<(PathBuf, Flags)>,
deletes: Vec<PathBuf>,
},
CopyToFolder {
paths: Vec<PathBuf>,
mailbox: String,
},
Append {
mailbox: Option<String>,
flags: Flags,
body: Vec<u8>,
},
CheckNew,
Folders,
Unseen(Vec<String>),
SearchBody(String),
Switch(String),
Manage(Manage),
}
#[derive(Clone)]
pub enum Manage {
Create(String),
Delete(String),
Rename(String, String),
Subscribe(String, bool),
}
impl Job {
pub fn what(&self) -> &'static str {
match self {
Job::FetchBodies(paths) => match paths.len() {
1 => "fetching the message",
_ => "fetching the messages",
},
Job::Sync { .. } => "syncing",
Job::CopyToFolder { .. } => "copying to the trash",
Job::Append { .. } => "saving to the server",
Job::CheckNew => "checking for new mail",
Job::Folders => "listing folders",
Job::Manage(_) => "managing folders",
Job::Unseen(_) => "counting unread",
Job::SearchBody(_) => "searching on the server",
Job::Switch(_) => "opening the folder",
}
}
}
pub enum Done {
Nothing,
Arrived(usize),
Folders(Vec<(String, usize)>),
Counts(Vec<usize>),
Uids(Vec<u32>),
Folder(String),
Switched(Box<Facts>),
}
#[derive(Clone)]
pub struct Facts {
pub spec: String,
pub account: Account,
pub mailbox: String,
pub cache: PathBuf,
pub pending_backfill: Vec<u32>,
}
impl Facts {
fn of(remote: &Remote) -> Facts {
Facts {
spec: remote.spec.clone(),
account: remote.account.clone(),
mailbox: remote.mailbox.clone(),
cache: remote.cache.clone(),
pending_backfill: remote.pending_backfill.clone(),
}
}
}
pub struct Imap {
pub facts: Facts,
jobs: Sender<Job>,
answers: Receiver<Result<Done>>,
thread: Option<JoinHandle<()>>,
progress: Arc<Mutex<Option<String>>>,
busy: Option<&'static str>,
cutoff: rmut_core::net::Cutoff,
}
impl Imap {
pub fn new(remote: Remote) -> Imap {
let facts = Facts::of(&remote);
let cutoff = remote.cutoff();
let (jobs, inbox) = channel::<Job>();
let (outbox, answers) = channel::<Result<Done>>();
let progress = Arc::new(Mutex::new(None));
let thread = std::thread::spawn({
let progress = progress.clone();
move || run(remote, inbox, outbox, progress)
});
Imap {
facts,
jobs,
answers,
thread: Some(thread),
progress,
busy: None,
cutoff,
}
}
pub fn abort(&self) {
if self.busy.is_some() {
self.cutoff.cut();
}
}
pub fn progress_sink(slot: &Arc<Mutex<Option<String>>>) -> Progress {
let slot = slot.clone();
Box::new(move |line: &str| {
if let Ok(mut slot) = slot.lock() {
*slot = Some(line.to_string());
}
})
}
pub fn take_progress(&mut self) -> Option<String> {
self.progress.lock().ok().and_then(|mut slot| slot.take())
}
pub fn busy(&self) -> Option<&'static str> {
self.busy
}
pub fn start(&mut self, job: Job) -> Result<()> {
let what = job.what();
self.jobs
.send(job)
.map_err(|_| anyhow!("the connection is gone"))?;
self.busy = Some(what);
Ok(())
}
pub fn collect(&mut self) -> Option<Result<Done>> {
match self.answers.try_recv() {
Ok(done) => {
self.busy = None;
Some(self.remember(done))
}
Err(TryRecvError::Empty) => None,
Err(TryRecvError::Disconnected) => {
self.busy = None;
Some(Err(anyhow!("the connection is gone")))
}
}
}
pub fn blocking(&mut self, job: Job) -> Result<Done> {
self.start(job)?;
self.wait()
}
pub fn wait(&mut self) -> Result<Done> {
let done = match self.answers.recv() {
Ok(done) => done,
Err(RecvError) => Err(anyhow!("the connection is gone")),
};
self.busy = None;
self.remember(done)
}
fn remember(&mut self, done: Result<Done>) -> Result<Done> {
if let Ok(Done::Switched(facts)) = &done {
self.facts = (**facts).clone();
}
done
}
pub fn take_backfill(&mut self) -> Vec<u32> {
std::mem::take(&mut self.facts.pending_backfill)
}
}
impl Drop for Imap {
fn drop(&mut self) {
let (jobs, _) = channel();
let _ = std::mem::replace(&mut self.jobs, jobs);
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
fn run(
mut remote: Remote,
jobs: Receiver<Job>,
answers: Sender<Result<Done>>,
progress: Arc<Mutex<Option<String>>>,
) {
remote.set_progress(Imap::progress_sink(&progress));
while let Ok(job) = jobs.recv() {
let what = job.what();
let done = do_job(&mut remote, job).map_err(|err| err.context(what));
if let Ok(mut slot) = progress.lock() {
*slot = None;
}
if answers.send(done).is_err() {
break; }
}
}
fn do_job(remote: &mut Remote, job: Job) -> Result<Done> {
match job {
Job::FetchBodies(paths) => {
for path in &paths {
remote.fetch_body(path)?;
}
Ok(Done::Nothing)
}
Job::Sync { flags, deletes } => {
for (path, flags) in &flags {
remote.push_flags(path, *flags)?;
}
if !deletes.is_empty() {
remote.delete(&deletes)?;
}
Ok(Done::Nothing)
}
Job::CopyToFolder { paths, mailbox } => {
remote.copy_to_folder(&paths, &mailbox)?;
Ok(Done::Nothing)
}
Job::Append {
mailbox,
flags,
body,
} => {
let folder = match &mailbox {
Some(mailbox) => remote.append_to(mailbox, flags, &body)?,
None => remote.append_sent(&body)?,
};
Ok(Done::Folder(folder))
}
Job::CheckNew => Ok(Done::Arrived(remote.check_new()?)),
Job::Folders => Ok(Done::Folders(remote.folders()?)),
Job::Unseen(folders) => Ok(Done::Counts(
folders.iter().map(|f| remote.unseen(f)).collect(),
)),
Job::SearchBody(term) => Ok(Done::Uids(remote.search_body(&term)?)),
Job::Switch(mailbox) => {
remote.switch(&mailbox)?;
Ok(Done::Switched(Box::new(Facts::of(remote))))
}
Job::Manage(action) => {
match action {
Manage::Create(name) => remote.create_folder(&name)?,
Manage::Delete(name) => remote.delete_folder(&name)?,
Manage::Rename(from, to) => remote.rename_folder(&from, &to)?,
Manage::Subscribe(name, on) => remote.subscribe_folder(&name, on)?,
}
Ok(Done::Nothing)
}
}
}