|
|
|
@ -30,12 +30,7 @@ use std::sync::Arc;
|
|
|
|
|
#[derive(Debug, Clone)]
|
|
|
|
|
pub struct ImapOp {
|
|
|
|
|
uid: usize,
|
|
|
|
|
bytes: Option<String>,
|
|
|
|
|
headers: Option<String>,
|
|
|
|
|
body: Option<String>,
|
|
|
|
|
mailbox_path: String,
|
|
|
|
|
mailbox_hash: MailboxHash,
|
|
|
|
|
flags: Arc<FutureMutex<Option<Flag>>>,
|
|
|
|
|
connection: Arc<FutureMutex<ImapConnection>>,
|
|
|
|
|
uid_store: Arc<UIDStore>,
|
|
|
|
|
}
|
|
|
|
@ -43,7 +38,6 @@ pub struct ImapOp {
|
|
|
|
|
impl ImapOp {
|
|
|
|
|
pub fn new(
|
|
|
|
|
uid: usize,
|
|
|
|
|
mailbox_path: String,
|
|
|
|
|
mailbox_hash: MailboxHash,
|
|
|
|
|
connection: Arc<FutureMutex<ImapConnection>>,
|
|
|
|
|
uid_store: Arc<UIDStore>,
|
|
|
|
@ -51,65 +45,63 @@ impl ImapOp {
|
|
|
|
|
ImapOp {
|
|
|
|
|
uid,
|
|
|
|
|
connection,
|
|
|
|
|
bytes: None,
|
|
|
|
|
headers: None,
|
|
|
|
|
body: None,
|
|
|
|
|
mailbox_path,
|
|
|
|
|
mailbox_hash,
|
|
|
|
|
flags: Arc::new(FutureMutex::new(None)),
|
|
|
|
|
uid_store,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl BackendOp for ImapOp {
|
|
|
|
|
fn as_bytes(&mut self) -> Result<&[u8]> {
|
|
|
|
|
if self.bytes.is_none() {
|
|
|
|
|
let mut bytes_cache = self.uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(self.uid).or_default();
|
|
|
|
|
if cache.bytes.is_some() {
|
|
|
|
|
self.bytes = cache.bytes.clone();
|
|
|
|
|
} else {
|
|
|
|
|
drop(cache);
|
|
|
|
|
drop(bytes_cache);
|
|
|
|
|
let ret: Result<()> = futures::executor::block_on(async {
|
|
|
|
|
let mut response = String::with_capacity(8 * 1024);
|
|
|
|
|
{
|
|
|
|
|
let mut conn = self.connection.lock().await;
|
|
|
|
|
conn.examine_mailbox(self.mailbox_hash, &mut response, false)
|
|
|
|
|
.await?;
|
|
|
|
|
conn.send_command(
|
|
|
|
|
format!("UID FETCH {} (FLAGS RFC822)", self.uid).as_bytes(),
|
|
|
|
|
)
|
|
|
|
|
fn as_bytes(&mut self) -> ResultFuture<Vec<u8>> {
|
|
|
|
|
let mut response = String::with_capacity(8 * 1024);
|
|
|
|
|
let connection = self.connection.clone();
|
|
|
|
|
let mailbox_hash = self.mailbox_hash;
|
|
|
|
|
let uid = self.uid;
|
|
|
|
|
let uid_store = self.uid_store.clone();
|
|
|
|
|
Ok(Box::pin(async move {
|
|
|
|
|
let exists_in_cache = {
|
|
|
|
|
let mut bytes_cache = uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(uid).or_default();
|
|
|
|
|
cache.bytes.is_some()
|
|
|
|
|
};
|
|
|
|
|
if !exists_in_cache {
|
|
|
|
|
let mut response = String::with_capacity(8 * 1024);
|
|
|
|
|
{
|
|
|
|
|
let mut conn = connection.lock().await;
|
|
|
|
|
conn.examine_mailbox(mailbox_hash, &mut response, false)
|
|
|
|
|
.await?;
|
|
|
|
|
conn.read_response(&mut response, RequiredResponses::FETCH_REQUIRED)
|
|
|
|
|
.await?;
|
|
|
|
|
}
|
|
|
|
|
debug!(
|
|
|
|
|
"fetch response is {} bytes and {} lines",
|
|
|
|
|
response.len(),
|
|
|
|
|
response.lines().collect::<Vec<&str>>().len()
|
|
|
|
|
);
|
|
|
|
|
let UidFetchResponse {
|
|
|
|
|
uid, flags, body, ..
|
|
|
|
|
} = protocol_parser::uid_fetch_response(&response)?.1;
|
|
|
|
|
assert_eq!(uid, self.uid);
|
|
|
|
|
assert!(body.is_some());
|
|
|
|
|
let mut bytes_cache = self.uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(self.uid).or_default();
|
|
|
|
|
if let Some((flags, _)) = flags {
|
|
|
|
|
//self.flags.set(Some(flags));
|
|
|
|
|
cache.flags = Some(flags);
|
|
|
|
|
}
|
|
|
|
|
cache.bytes =
|
|
|
|
|
Some(unsafe { std::str::from_utf8_unchecked(body.unwrap()).to_string() });
|
|
|
|
|
self.bytes = cache.bytes.clone();
|
|
|
|
|
Ok(())
|
|
|
|
|
});
|
|
|
|
|
ret?;
|
|
|
|
|
conn.send_command(format!("UID FETCH {} (FLAGS RFC822)", uid).as_bytes())
|
|
|
|
|
.await?;
|
|
|
|
|
conn.read_response(&mut response, RequiredResponses::FETCH_REQUIRED)
|
|
|
|
|
.await?;
|
|
|
|
|
}
|
|
|
|
|
debug!(
|
|
|
|
|
"fetch response is {} bytes and {} lines",
|
|
|
|
|
response.len(),
|
|
|
|
|
response.lines().collect::<Vec<&str>>().len()
|
|
|
|
|
);
|
|
|
|
|
let UidFetchResponse {
|
|
|
|
|
uid: _uid,
|
|
|
|
|
flags: _flags,
|
|
|
|
|
body,
|
|
|
|
|
..
|
|
|
|
|
} = protocol_parser::uid_fetch_response(&response)?.1;
|
|
|
|
|
assert_eq!(_uid, uid);
|
|
|
|
|
assert!(body.is_some());
|
|
|
|
|
let mut bytes_cache = uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(uid).or_default();
|
|
|
|
|
if let Some((_flags, _)) = _flags {
|
|
|
|
|
//flags.lock().await.set(Some(_flags));
|
|
|
|
|
cache.flags = Some(_flags);
|
|
|
|
|
}
|
|
|
|
|
cache.bytes =
|
|
|
|
|
Some(unsafe { std::str::from_utf8_unchecked(body.unwrap()).to_string() });
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Ok(self.bytes.as_ref().unwrap().as_bytes())
|
|
|
|
|
let mut bytes_cache = uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(uid).or_default();
|
|
|
|
|
let ret = cache.bytes.clone().unwrap().into_bytes();
|
|
|
|
|
Ok(ret)
|
|
|
|
|
}))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn fetch_flags(&self) -> ResultFuture<Flag> {
|
|
|
|
@ -118,12 +110,8 @@ impl BackendOp for ImapOp {
|
|
|
|
|
let mailbox_hash = self.mailbox_hash;
|
|
|
|
|
let uid = self.uid;
|
|
|
|
|
let uid_store = self.uid_store.clone();
|
|
|
|
|
let flags = self.flags.clone();
|
|
|
|
|
|
|
|
|
|
Ok(Box::pin(async move {
|
|
|
|
|
if let Some(val) = *flags.lock().await {
|
|
|
|
|
return Ok(val);
|
|
|
|
|
}
|
|
|
|
|
let exists_in_cache = {
|
|
|
|
|
let mut bytes_cache = uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(uid).or_default();
|
|
|
|
@ -158,7 +146,6 @@ impl BackendOp for ImapOp {
|
|
|
|
|
}
|
|
|
|
|
let (_uid, (_flags, _)) = v[0];
|
|
|
|
|
assert_eq!(uid, uid);
|
|
|
|
|
*flags.lock().await = Some(_flags);
|
|
|
|
|
let mut bytes_cache = uid_store.byte_cache.lock()?;
|
|
|
|
|
let cache = bytes_cache.entry(uid).or_default();
|
|
|
|
|
cache.flags = Some(_flags);
|
|
|
|
@ -170,8 +157,6 @@ impl BackendOp for ImapOp {
|
|
|
|
|
let val = cache.flags;
|
|
|
|
|
val
|
|
|
|
|
};
|
|
|
|
|
let mut f = flags.lock().await;
|
|
|
|
|
*f = val;
|
|
|
|
|
Ok(val.unwrap())
|
|
|
|
|
}
|
|
|
|
|
}))
|
|
|
|
|