use anyhow::{Result, anyhow, bail};
use chrono::DateTime;
use io_pimdir::PimdirItem;
use io_replica::{
client::ReplicaStorage,
collection::ReplicaCollectionId,
coroutine::{ReplicaArg, ReplicaCoroutine, ReplicaCoroutineState, ReplicaYield},
mutate::{ReplicaMutate, ReplicaMutation},
object::ReplicaObject,
placement::{ReplicaFlags, ReplicaHandle, ReplicaLinkId, ReplicaMeta, ReplicaPlacement},
};
use log::warn;
use mail_parser::{Address as MailParserAddress, HeaderValue, MessageParser};
use serde::Deserialize;
use crate::{
email::{
address::Address,
envelope::{Envelope, normalize_message_id, parse_message_ids},
flag::{Flag, FlagOp, IanaFlag},
mailbox::Mailbox,
search::{eval, query::SearchEmailsQuery},
},
pimdir::{client::PimdirClient, hash::content_hash},
};
const MAIL_KIND: &str = "message/rfc822";
const SCAN_BATCH: usize = 500;
#[derive(Default, Deserialize)]
struct MetaView {
#[serde(default)]
message_id: Option<String>,
#[serde(default)]
in_reply_to: Vec<String>,
#[serde(default)]
subject: String,
#[serde(default)]
from: Option<String>,
#[serde(default)]
to: Option<String>,
#[serde(default)]
date: Option<String>,
#[serde(default)]
size: Option<u64>,
}
impl PimdirClient {
pub fn list_mailboxes(&mut self, with_counts: bool) -> Result<Vec<Mailbox>> {
let mut mailboxes = Vec::new();
for collection in self.store.list_collections()? {
if !collection.kind.is_empty() && collection.kind != MAIL_KIND {
continue;
}
let total = if with_counts {
Some(self.store.count_items(&collection.id)?)
} else {
None
};
mailboxes.push(Mailbox {
id: collection.id,
name: collection.name,
total,
unread: None,
});
}
mailboxes.sort_by(|a, b| a.name.cmp(&b.name));
Ok(mailboxes)
}
pub fn list_envelopes(
&mut self,
mailbox: &str,
page: Option<u32>,
page_size: Option<u32>,
_with_attachment: bool,
) -> Result<Vec<Envelope>> {
let mut envelopes: Vec<Envelope> = self
.scan_items(mailbox)?
.iter()
.map(envelope_from_item)
.collect();
envelopes.sort_by_key(|envelope| std::cmp::Reverse(envelope.date));
Ok(paginate(envelopes, page, page_size))
}
pub fn search_envelopes(
&mut self,
mailbox: &str,
query: Option<&SearchEmailsQuery>,
page: Option<u32>,
page_size: Option<u32>,
_with_attachment: bool,
) -> Result<Vec<Envelope>> {
let filter = query.and_then(|q| q.filter.as_ref());
let mut hits: Vec<Envelope> = self
.scan_items(mailbox)?
.iter()
.map(envelope_from_item)
.filter(|envelope| match filter {
Some(filter) => eval::matches_filter(envelope, &[], filter),
None => true,
})
.collect();
eval::sort_envelopes(&mut hits, query.and_then(|q| q.sort.as_deref()));
Ok(paginate(hits, page, page_size))
}
pub fn get_message(&mut self, mailbox: &str, id: &str, seen: bool) -> Result<Vec<u8>> {
let Some(item) = self.store.get_item(mailbox, parse_id(id)?)? else {
bail!("Message `{id}` not found in `{mailbox}`");
};
let Some(hash) = item.object else {
bail!(
"Message `{id}` in `{mailbox}` is not downloaded yet (body not fetched); \
run a sync to hydrate it"
);
};
let bytes = self
.blobs
.get(&hash)?
.ok_or_else(|| anyhow!("Body blob missing for `{id}` in `{mailbox}`"))?;
if seen {
let seen_flag = Flag::from_iana(IanaFlag::Seen);
if let Err(err) = self.store_flags(mailbox, &[id], &[seen_flag], FlagOp::Add) {
warn!("could not stage \\Seen on `{id}` in `{mailbox}`: {err:#}");
}
}
Ok(bytes)
}
pub fn store_flags(
&mut self,
mailbox: &str,
ids: &[&str],
flags: &[Flag],
op: FlagOp,
) -> Result<()> {
for id in ids {
let placement = self.synced_placement(mailbox, id)?;
let next = apply_flag_op(&placement.flags, flags, op);
self.run_mutation(
mailbox,
ReplicaMutation::SetFlags {
handle: placement.handle,
flags: next,
},
)?;
}
Ok(())
}
pub fn add_message(&mut self, mailbox: &str, flags: &[Flag], raw: Vec<u8>) -> Result<String> {
let (link_id, meta) = derive_link_and_meta(&raw);
let object = ReplicaObject {
hash: content_hash(&raw),
size: raw.len(),
};
let handle = ReplicaHandle(format!("local:{}", link_id.0));
let link = link_id.0.clone();
self.run_mutation(
mailbox,
ReplicaMutation::Add {
handle,
link_id,
flags: to_replica_flags(flags),
object,
body: raw,
meta: Some(meta),
},
)?;
let seq = self
.store
.seq_for_link(mailbox, &link)?
.ok_or_else(|| anyhow!("Added message `{link}` in `{mailbox}` has no public id"))?;
Ok(seq.to_string())
}
pub fn copy_messages(&mut self, from: &str, to: &str, ids: &[&str]) -> Result<usize> {
for id in ids {
let placement = self.synced_placement(from, id)?;
let placeholder = ReplicaHandle(format!("copy:{}:{}", to, placement.handle.0));
self.run_mutation(
from,
ReplicaMutation::Copy {
handle: placement.handle,
target: ReplicaCollectionId(to.to_string()),
placeholder,
},
)?;
}
Ok(ids.len())
}
pub fn move_messages(&mut self, from: &str, to: &str, ids: &[&str]) -> Result<usize> {
for id in ids {
let placement = self.synced_placement(from, id)?;
let placeholder = ReplicaHandle(format!("move:{}:{}", to, placement.handle.0));
self.run_mutation(
from,
ReplicaMutation::Move {
handle: placement.handle,
target: ReplicaCollectionId(to.to_string()),
placeholder,
},
)?;
}
Ok(ids.len())
}
pub fn delete_messages(&mut self, mailbox: &str, ids: &[&str]) -> Result<()> {
for id in ids {
let placement = self.synced_placement(mailbox, id)?;
self.run_mutation(mailbox, ReplicaMutation::Remove(placement.handle))?;
}
Ok(())
}
fn scan_items(&self, mailbox: &str) -> Result<Vec<PimdirItem>> {
let mut all = Vec::new();
let mut cursor: Option<String> = None;
loop {
let page = self
.store
.list_items(mailbox, cursor.as_deref(), SCAN_BATCH)?;
let n = page.len();
if let Some(last) = page.last() {
cursor = Some(last.link_id.0.clone());
}
all.extend(page);
if n < SCAN_BATCH {
break;
}
}
Ok(all)
}
fn synced_placement(&self, collection: &str, id: &str) -> Result<ReplicaPlacement> {
let seq = parse_id(id)?;
let link_id = self
.store
.get_item(collection, seq)?
.map(|item| item.link_id.0)
.ok_or_else(|| anyhow!("Message `{id}` not found in `{collection}`"))?;
let loaded = self
.store
.load(&ReplicaCollectionId(collection.to_string()))?;
let placement = loaded
.placements
.into_iter()
.find(|p| p.link_id.as_ref().map(|l| l.0.as_str()) == Some(link_id.as_str()))
.ok_or_else(|| anyhow!("Message `{id}` not found in `{collection}`"))?;
if placement.base.is_none() {
bail!(
"`{collection}` was not synced as source `{}`, so `{id}` cannot be edited \
here; set `pimdir.source` to the sync source and sync first",
self.source
);
}
Ok(placement)
}
fn run_mutation(&mut self, collection: &str, mutation: ReplicaMutation) -> Result<()> {
let mut coroutine = ReplicaMutate::new(collection.to_string(), mutation);
let mut arg: Option<ReplicaArg> = None;
loop {
match coroutine.resume(arg.take()) {
ReplicaCoroutineState::Yielded(ReplicaYield::WantsLoad(collection)) => {
let loaded = self.store.load(&collection)?;
arg = Some(ReplicaArg::Load(loaded));
}
ReplicaCoroutineState::Yielded(ReplicaYield::WantsWrite(ops)) => {
self.store.write(ops)?;
arg = Some(ReplicaArg::Write);
}
ReplicaCoroutineState::Yielded(_) => {
bail!("pimdir mutate asked for an unexpected step");
}
ReplicaCoroutineState::Complete(result) => {
return result.map_err(|err| anyhow!("pimdir mutate failed: {err}"));
}
}
}
}
}
fn envelope_from_item(item: &PimdirItem) -> Envelope {
let view: MetaView = item
.meta
.as_ref()
.and_then(|meta| serde_json::from_str(&meta.0).ok())
.unwrap_or_default();
let flags = item
.flags
.0
.iter()
.map(|raw| Flag::from_raw(raw.clone()))
.collect();
let from = view
.from
.map(|email| vec![Address { name: None, email }])
.unwrap_or_default();
let to = view
.to
.map(|email| vec![Address { name: None, email }])
.unwrap_or_default();
let date = view
.date
.as_deref()
.and_then(|d| DateTime::parse_from_rfc3339(d).ok());
Envelope {
id: item.seq.to_string(),
message_id: view.message_id,
in_reply_to: view.in_reply_to,
flags,
subject: view.subject,
from,
to,
date,
size: view.size.unwrap_or(0),
has_attachment: None,
}
}
fn derive_link_and_meta(raw: &[u8]) -> (ReplicaLinkId, ReplicaMeta) {
let parsed = MessageParser::default().parse(raw);
let message_id = parsed.as_ref().and_then(|m| m.message_id());
let in_reply_to: Vec<String> = parsed
.as_ref()
.map(|m| match m.in_reply_to() {
HeaderValue::TextList(ids) => ids
.iter()
.filter_map(|id| normalize_message_id(id))
.collect(),
HeaderValue::Text(id) => parse_message_ids(id),
_ => Vec::new(),
})
.unwrap_or_default();
let subject = parsed.as_ref().and_then(|m| m.subject()).unwrap_or("");
let from = parsed
.as_ref()
.and_then(|m| m.from())
.and_then(first_address);
let to = parsed.as_ref().and_then(|m| m.to()).and_then(first_address);
let date = parsed
.as_ref()
.and_then(|m| m.date())
.map(|d| d.to_rfc3339());
let link_id = match message_id {
Some(mid) if !mid.trim().is_empty() => {
ReplicaLinkId(format!("mid:{}", mid.trim().trim_matches(['<', '>'])))
}
_ => ReplicaLinkId(format!(
"alt:{subject}|{}|{}",
date.as_deref().unwrap_or(""),
from.as_deref().unwrap_or("")
)),
};
#[derive(serde::Serialize)]
struct MetaWrite<'a> {
v: u8,
#[serde(skip_serializing_if = "Option::is_none")]
message_id: Option<&'a str>,
#[serde(skip_serializing_if = "<[String]>::is_empty")]
in_reply_to: &'a [String],
subject: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
from: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
to: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
date: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
size: Option<u64>,
}
let summary = MetaWrite {
v: 1,
message_id,
in_reply_to: &in_reply_to,
subject,
from: from.as_deref(),
to: to.as_deref(),
date: date.as_deref(),
size: (!raw.is_empty()).then_some(raw.len() as u64),
};
let meta = ReplicaMeta(serde_json::to_string(&summary).unwrap_or_default());
(link_id, meta)
}
fn first_address(addrs: &MailParserAddress<'_>) -> Option<String> {
addrs
.clone()
.into_list()
.into_iter()
.find_map(|a| a.address.map(|s| s.into_owned()))
}
fn apply_flag_op(current: &ReplicaFlags, flags: &[Flag], op: FlagOp) -> ReplicaFlags {
let incoming = flags.iter().map(|f| f.raw().to_string());
match op {
FlagOp::Set => ReplicaFlags(incoming.collect()),
FlagOp::Add => {
let mut set = current.0.clone();
set.extend(incoming);
ReplicaFlags(set)
}
FlagOp::Remove => {
let mut set = current.0.clone();
for flag in incoming {
set.remove(&flag);
}
ReplicaFlags(set)
}
}
}
fn to_replica_flags(flags: &[Flag]) -> ReplicaFlags {
ReplicaFlags(flags.iter().map(|f| f.raw().to_string()).collect())
}
fn parse_id(id: &str) -> Result<i64> {
id.parse::<i64>()
.map_err(|_| anyhow!("Invalid message id `{id}` (expected a number)"))
}
fn paginate<T>(items: Vec<T>, page: Option<u32>, page_size: Option<u32>) -> Vec<T> {
let Some(size) = page_size else {
return items;
};
if size == 0 {
return Vec::new();
}
let page = page.unwrap_or(1).max(1);
let skip = ((page - 1) as usize).saturating_mul(size as usize);
if skip >= items.len() {
return Vec::new();
}
items.into_iter().skip(skip).take(size as usize).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn envelope_is_built_from_meta_without_a_body() {
let item = PimdirItem {
seq: 42,
link_id: ReplicaLinkId("mid:x@y".into()),
flags: ReplicaFlags(["\\Seen".to_string()].into_iter().collect()),
meta: Some(ReplicaMeta(
r#"{"v":1,"message_id":"x@y","subject":"Hi","from":"a@x.org","to":"b@x.org","size":99}"#
.into(),
)),
object: None,
level: io_replica::placement::ReplicaLevel::Meta,
};
let envelope = envelope_from_item(&item);
assert_eq!(envelope.id, "42");
assert_eq!(envelope.subject, "Hi");
assert_eq!(envelope.from[0].email, "a@x.org");
assert_eq!(envelope.to[0].email, "b@x.org");
assert_eq!(envelope.size, 99);
assert!(envelope.flags.iter().any(|f| f.raw() == "\\Seen"));
}
#[test]
fn add_derives_a_mid_link_and_v1_meta() {
let raw = b"Message-ID: <new@host>\r\nSubject: Compose\r\nFrom: a@x.org\r\n\r\nbody";
let (link, meta) = derive_link_and_meta(raw);
assert_eq!(link.0, "mid:new@host");
assert!(meta.0.contains("\"v\":1"));
assert!(meta.0.contains("Compose"));
}
#[test]
fn flag_ops_add_set_and_remove() {
let base = ReplicaFlags(["\\Seen".to_string()].into_iter().collect());
let flagged = [Flag::from_raw("\\Flagged")];
let added = apply_flag_op(&base, &flagged, FlagOp::Add);
assert!(added.contains("\\Seen") && added.contains("\\Flagged"));
let set = apply_flag_op(&base, &flagged, FlagOp::Set);
assert!(!set.contains("\\Seen") && set.contains("\\Flagged"));
let seen = [Flag::from_raw("\\Seen")];
let removed = apply_flag_op(&base, &seen, FlagOp::Remove);
assert!(!removed.contains("\\Seen"));
}
}