use std::{
cmp::Reverse,
collections::{BTreeSet, HashMap},
fmt,
io::{self, Write},
mem,
sync::{
Mutex,
atomic::{AtomicBool, Ordering},
},
thread,
};
use anyhow::{Context, Result, bail};
use crossbeam_queue::SegQueue;
use io_pimdir::{
change::{PimdirChange, PimdirChangeKind},
client::blobs::{PimdirBlobWriter, PimdirBlobs},
collection::{PimdirCheckpoint, PimdirCollectionId},
hash::PimdirHasher,
object::PimdirHash,
placement::{PimdirFlags, PimdirHandle},
remote::{
PimdirFetchedBody, PimdirFetchedItem, PimdirPushOutcome, PimdirPushResult, PimdirRemote,
PimdirRemoteItem, PimdirRemoteSnapshot, PimdirTier,
},
summary::PimdirDerivation,
};
use log::warn;
#[cfg(feature = "dav")]
use crate::dav::client::is_duplicate_uid;
use crate::{
client::{Client, Pool},
item::{
flag::{Flag, FlagOp},
summary::ItemSummary,
},
kind::{Kind, LinkId},
};
#[cfg(not(feature = "dav"))]
fn is_duplicate_uid(_err: &anyhow::Error) -> bool {
false
}
pub struct PimRemote<'a> {
kind: Kind,
pool: &'a mut Pool,
blob: PimdirBlobs,
namespace: String,
on_body: Option<&'a (dyn Fn() + Sync)>,
sizes: HashMap<String, u64>,
held: HeldHandles,
refused: Vec<RefusedCreate>,
rejected: Vec<RejectedPush>,
}
#[derive(Default)]
struct HeldHandles(HashMap<String, BTreeSet<String>>);
impl HeldHandles {
fn remember(
&mut self,
collection: &str,
items: &[PimdirRemoteItem],
vanished: &[PimdirHandle],
) {
let held = self.0.entry(collection.to_string()).or_default();
held.extend(items.iter().map(|item| item.handle.0.clone()));
for handle in vanished {
held.remove(handle.as_str());
}
}
fn claim(&mut self, collection: &str, handle: &str) -> bool {
self.0
.entry(collection.to_string())
.or_default()
.insert(handle.to_string())
}
}
pub struct RefusedCreate {
pub collection: String,
pub uid: String,
pub handle: String,
}
pub struct RejectedPush {
pub collection: String,
pub handle: String,
pub action: &'static str,
pub reason: String,
}
impl<'a> PimRemote<'a> {
pub fn new(pool: &'a mut Pool, blob: PimdirBlobs, namespace: impl Into<String>) -> Self {
let kind = resolve_kind(pool);
Self {
kind,
pool,
blob,
namespace: namespace.into(),
on_body: None,
sizes: HashMap::new(),
held: HeldHandles::default(),
refused: Vec::new(),
rejected: Vec::new(),
}
}
pub fn with_progress(
pool: &'a mut Pool,
blob: PimdirBlobs,
namespace: impl Into<String>,
on_body: &'a (dyn Fn() + Sync),
sizes: HashMap<String, u64>,
) -> Self {
let kind = resolve_kind(pool);
Self {
kind,
pool,
blob,
namespace: namespace.into(),
on_body: Some(on_body),
sizes,
held: HeldHandles::default(),
refused: Vec::new(),
rejected: Vec::new(),
}
}
fn wire_name<'n>(&self, collection: &'n str) -> &'n str {
wire_name(&self.namespace, collection)
}
pub fn take_refused(&mut self) -> Vec<RefusedCreate> {
mem::take(&mut self.refused)
}
pub fn take_rejected(&mut self) -> Vec<RejectedPush> {
mem::take(&mut self.rejected)
}
fn reject(
&mut self,
collection: &str,
handle: PimdirHandle,
action: &'static str,
reason: impl fmt::Display,
) -> PimdirPushResult {
let reason = reason.to_string();
warn!("{action} {} in {collection} rejected: {reason}", handle.0);
self.rejected.push(RejectedPush {
collection: collection.to_string(),
handle: handle.0.clone(),
action,
reason,
});
rejected_bare(handle)
}
}
pub(crate) fn wire_name<'n>(namespace: &str, collection: &'n str) -> &'n str {
collection
.strip_prefix(namespace)
.and_then(|rest| rest.strip_prefix('/'))
.unwrap_or(collection)
}
pub(crate) fn resolve_kind(pool: &mut Pool) -> Kind {
let media_type = pool.primary().media_type();
Kind::from_media_type(media_type).unwrap_or_else(|| {
warn!("unknown media type {media_type}, deriving link ids as mail");
Kind::Mail
})
}
pub(crate) const BATCH_SIZE: usize = 64;
fn to_offline_flags<'f>(flags: impl IntoIterator<Item = &'f Flag>) -> PimdirFlags {
flags.into_iter().map(|f| f.raw()).collect()
}
fn to_item_flags(flags: &PimdirFlags) -> Vec<Flag> {
let Some(flags) = flags.known() else {
return Vec::new();
};
flags.iter().map(|s| Flag::from_raw(s.clone())).collect()
}
impl PimdirRemote for PimRemote<'_> {
type Error = anyhow::Error;
fn enumerate(
&mut self,
collection: &PimdirCollectionId,
cursor: Option<PimdirCheckpoint>,
) -> Result<PimdirRemoteSnapshot, Self::Error> {
let collection = self.wire_name(collection.as_str());
let cursor = cursor.as_ref().map(|c| c.0.as_slice());
let enumeration = self
.pool
.primary()
.enumerate(collection, cursor)
.with_context(|| format!("Enumerate {collection} error"))?;
let items: Vec<PimdirRemoteItem> = enumeration
.items
.into_iter()
.map(|entry| PimdirRemoteItem {
handle: PimdirHandle::from(entry.id),
flags: to_offline_flags(&entry.flags),
revision: entry.revision,
})
.collect();
let vanished: Vec<PimdirHandle> = enumeration
.vanished
.into_iter()
.map(PimdirHandle::from)
.collect();
self.held.remember(collection, &items, &vanished);
Ok(PimdirRemoteSnapshot {
items,
vanished,
complete: enumeration.complete,
checkpoint: PimdirCheckpoint(enumeration.checkpoint),
})
}
fn fetch(
&mut self,
collection: &PimdirCollectionId,
handles: Vec<PimdirHandle>,
tier: PimdirTier,
) -> Result<Vec<PimdirFetchedItem>, Self::Error> {
let collection = self.wire_name(collection.as_str());
match tier {
PimdirTier::Meta => self.fetch_meta(collection, handles),
PimdirTier::Full => self.fetch_full(collection, handles),
}
}
fn push(
&mut self,
collection: &PimdirCollectionId,
changes: Vec<PimdirChange>,
) -> Result<Vec<PimdirPushResult>, Self::Error> {
let collection = self.wire_name(collection.as_str()).to_string();
let mut results = Vec::with_capacity(changes.len());
for change in changes {
let result = match change.kind {
PimdirChangeKind::SetFlags { handle, flags } => {
let email_flags = to_item_flags(&flags);
let stored = self.pool.primary().store_flags(
&collection,
&[handle.as_str()],
&email_flags,
FlagOp::Set,
);
match stored {
Ok(()) => accepted(handle, None),
Err(err) => {
self.reject(&collection, handle, "set flags", format!("{err:#}"))
}
}
}
PimdirChangeKind::Remove {
handle,
to,
link_id: _,
if_match,
} => match to {
Some(target) => {
let dest = wire_name(&self.namespace, target.as_str()).to_string();
let moved =
self.pool
.primary()
.move_items(&collection, &dest, &[handle.as_str()]);
match moved {
Ok(()) => accepted(handle, None),
Err(err) => {
self.reject(&collection, handle, "move", format!("{err:#}"))
}
}
}
None => {
let deleted = self.pool.primary().delete_item(
&collection,
handle.as_str(),
if_match.as_deref(),
);
match deleted {
Ok(()) => accepted(handle, None),
Err(err) => {
self.reject(&collection, handle, "delete", format!("{err:#}"))
}
}
}
},
PimdirChangeKind::Add {
handle,
link_id,
flags,
object,
..
} => {
let link = link_id
.as_ref()
.map(|link| self.kind.split_link_id(link))
.unwrap_or_default();
self.append(&collection, handle, &flags, object, link)
}
PimdirChangeKind::Update {
handle,
object,
if_match,
} => self.update(&collection, handle, object, if_match.as_deref()),
};
results.push(result);
}
Ok(results)
}
}
pub type FetchKey = (String, String);
pub struct CachedFetchRemote<'a> {
cache: &'a HashMap<FetchKey, PimdirFetchedItem>,
fallback: PimRemote<'a>,
}
impl<'a> CachedFetchRemote<'a> {
pub fn new(cache: &'a HashMap<FetchKey, PimdirFetchedItem>, fallback: PimRemote<'a>) -> Self {
Self { cache, fallback }
}
}
impl PimdirRemote for CachedFetchRemote<'_> {
type Error = anyhow::Error;
fn enumerate(
&mut self,
collection: &PimdirCollectionId,
cursor: Option<PimdirCheckpoint>,
) -> Result<PimdirRemoteSnapshot, Self::Error> {
self.fallback.enumerate(collection, cursor)
}
fn fetch(
&mut self,
collection: &PimdirCollectionId,
handles: Vec<PimdirHandle>,
tier: PimdirTier,
) -> Result<Vec<PimdirFetchedItem>, Self::Error> {
let coll = collection.as_str();
let mut items = Vec::with_capacity(handles.len());
let mut misses = Vec::new();
for handle in handles {
match self.cache.get(&(coll.to_string(), handle.0.clone())) {
Some(item) => items.push(item.clone()),
None => misses.push(handle),
}
}
if !misses.is_empty() {
items.extend(self.fallback.fetch(collection, misses, tier)?);
}
Ok(items)
}
fn push(
&mut self,
collection: &PimdirCollectionId,
changes: Vec<PimdirChange>,
) -> Result<Vec<PimdirPushResult>, Self::Error> {
self.fallback.push(collection, changes)
}
}
impl PimRemote<'_> {
fn fetch_meta(
&mut self,
collection: &str,
handles: Vec<PimdirHandle>,
) -> Result<Vec<PimdirFetchedItem>> {
let ids: Vec<&str> = handles.iter().map(|h| h.as_str()).collect();
let envelopes = self
.pool
.primary()
.fetch_summaries(collection, &ids)
.with_context(|| format!("Fetch envelopes {collection} error"))?;
let by_id: HashMap<&str, &ItemSummary> =
envelopes.iter().map(|e| (e.id.as_str(), e)).collect();
let mut items = Vec::with_capacity(handles.len());
for handle in handles {
let Some(env) = by_id.get(handle.as_str()) else {
continue;
};
let Some(PimdirDerivation {
link_id,
summary,
sort_key,
}) = self.kind.parse_summary(env)
else {
continue;
};
items.push(PimdirFetchedItem {
handle,
link_id,
summary,
sort_key,
body: None,
revision: None,
});
}
Ok(items)
}
fn fetch_full(
&mut self,
collection: &str,
mut handles: Vec<PimdirHandle>,
) -> Result<Vec<PimdirFetchedItem>> {
if handles.is_empty() {
return Ok(Vec::new());
}
if self.sizes.is_empty() {
handles.sort_by_key(|h| h.as_str().parse::<u64>().unwrap_or(u64::MAX));
} else {
handles.sort_by_key(|h| Reverse(self.sizes.get(h.as_str()).copied().unwrap_or(0)));
}
let total = handles.len();
let target = self.pool.max().min(total);
let batches: Vec<Vec<PimdirHandle>> = handles
.chunks(BATCH_SIZE)
.map(<[PimdirHandle]>::to_vec)
.collect();
if target <= 1 {
let blob = self.blob.clone();
let mut items = Vec::with_capacity(total);
for batch in &batches {
items.extend(hydrate_batch(
self.kind,
self.pool.primary(),
collection,
batch,
&blob,
self.on_body,
)?);
}
return Ok(items);
}
self.fetch_full_pooled(collection, batches, target)
}
fn fetch_full_pooled(
&mut self,
collection: &str,
batches: Vec<Vec<PimdirHandle>>,
target: usize,
) -> Result<Vec<PimdirFetchedItem>> {
let queue: SegQueue<Vec<PimdirHandle>> = SegQueue::new();
for batch in batches {
queue.push(batch);
}
let results: Mutex<Vec<PimdirFetchedItem>> = Mutex::new(Vec::new());
let failure: Mutex<Option<anyhow::Error>> = Mutex::new(None);
let stop = AtomicBool::new(false);
let kind = self.kind;
let blob = self.blob.clone();
let clients = self.pool.workers(target)?;
let queue_ref = &queue;
let results_ref = &results;
let failure_ref = &failure;
let stop_ref = &stop;
let blob_ref = &blob;
let on_body = self.on_body;
thread::scope(|scope| {
for client in clients.iter_mut() {
scope.spawn(move || {
while !stop_ref.load(Ordering::Relaxed) {
let Some(batch) = queue_ref.pop() else {
break;
};
match hydrate_batch(kind, client, collection, &batch, blob_ref, on_body) {
Ok(mut items) => results_ref.lock().unwrap().append(&mut items),
Err(err) => {
*failure_ref.lock().unwrap() = Some(err);
stop_ref.store(true, Ordering::Relaxed);
break;
}
}
}
});
}
});
if let Some(err) = failure.into_inner().unwrap() {
return Err(err);
}
Ok(results.into_inner().unwrap())
}
}
fn fetch_one_full(
kind: Kind,
client: &mut Client,
collection: &str,
handle: PimdirHandle,
blob: &PimdirBlobs,
) -> Result<PimdirFetchedItem> {
let writer = blob.writer().context("Open blob writer error")?;
let mut sink = HydrateSink::new(writer, blob.hasher());
let revision = client
.get_item_stream(collection, handle.as_str(), &mut sink)
.with_context(|| format!("Stream item {} in {collection} error", handle.as_str()))?;
let (hash, size, header) = sink
.finish()
.with_context(|| format!("Commit body {} in {collection} error", handle.as_str()))?;
if size == 0 {
bail!(
"Server returned an empty body for {} in {collection}",
handle.as_str(),
);
}
let PimdirDerivation {
link_id,
summary,
sort_key,
} = kind.parse_body(&header, size as u64);
Ok(PimdirFetchedItem {
handle,
link_id,
summary,
sort_key,
body: Some(PimdirFetchedBody::Persisted { hash, size }),
revision,
})
}
pub(crate) fn hydrate_batch(
kind: Kind,
client: &mut Client,
collection: &str,
handles: &[PimdirHandle],
blob: &PimdirBlobs,
on_body: Option<&(dyn Fn() + Sync)>,
) -> Result<Vec<PimdirFetchedItem>> {
let ids: Vec<&str> = handles.iter().map(|h| h.as_str()).collect();
let mut items: Vec<PimdirFetchedItem> = Vec::with_capacity(handles.len());
let batched = client.fetch_bodies(
collection,
&ids,
|_id| {
blob.writer()
.map(|writer| HydrateSink::new(writer, blob.hasher()))
},
|id, revision, sink: HydrateSink| {
let (hash, size, header) = sink.finish().map_err(io::Error::other)?;
let PimdirDerivation {
link_id,
summary,
sort_key,
} = kind.parse_body(&header, size as u64);
items.push(PimdirFetchedItem {
handle: PimdirHandle::from(id),
link_id,
summary,
sort_key,
body: Some(PimdirFetchedBody::Persisted { hash, size }),
revision: revision.map(str::to_string),
});
if let Some(cb) = on_body {
cb();
}
Ok(())
},
);
match batched {
Ok(()) => {
let fetched: BTreeSet<String> = items
.iter()
.map(|item| item.handle.as_str().to_owned())
.collect();
let missing: Vec<PimdirHandle> = handles
.iter()
.filter(|handle| !fetched.contains(handle.as_str()))
.cloned()
.collect();
if !missing.is_empty() {
warn!(
"batched fetch {collection} returned {} of {} bodies; \
fetching the rest one by one",
items.len(),
handles.len(),
);
let blob = blob.clone();
for handle in missing {
items.push(fetch_one_full(kind, client, collection, handle, &blob)?);
if let Some(cb) = on_body {
cb();
}
}
}
Ok(items)
}
Err(err) => {
warn!("batched fetch {collection} failed ({err:#}); falling back to per-item");
items.clear();
let blob = blob.clone();
for handle in handles {
items.push(fetch_one_full(
kind,
client,
collection,
handle.clone(),
&blob,
)?);
if let Some(cb) = on_body {
cb();
}
}
Ok(items)
}
}
}
impl PimRemote<'_> {
fn append(
&mut self,
collection: &str,
handle: PimdirHandle,
flags: &PimdirFlags,
object: Option<PimdirHash>,
link: LinkId<'_>,
) -> PimdirPushResult {
let Some(hash) = object else {
return self.reject(collection, handle, "append", "no body was stored for it");
};
let reader = match self.blob.reader(&hash) {
Ok(Some(file)) => file,
Ok(None) => {
let reason = format!("its body {} is missing from the blob tree", hash.as_str());
return self.reject(collection, handle, "append", reason);
}
Err(err) => {
return self.reject(collection, handle, "append", format!("{err:#}"));
}
};
let len = match reader.metadata() {
Ok(meta) => meta.len() as usize,
Err(err) => {
return self.reject(collection, handle, "append", format!("{err:#}"));
}
};
let item_flags = to_item_flags(flags);
let written =
self.pool
.primary()
.add_item_stream(collection, &item_flags, reader, len, link);
let written = match written {
Ok(written) => written,
Err(err) => {
if !is_duplicate_uid(&err) {
return self.reject(collection, handle, "append", format!("{err:#}"));
}
let uid = link.hint.unwrap_or(handle.as_str());
warn!("append to {collection} refused: it already holds UID {uid}");
self.refused.push(RefusedCreate {
collection: collection.to_string(),
uid: uid.to_string(),
handle: handle.0.clone(),
});
return rejected_bare(handle);
}
};
let assigned = PimdirHandle::from(written.id);
if !self.held.claim(collection, assigned.as_str()) {
let reason = format!(
"the server answered with {}, which it already holds",
assigned.as_str(),
);
return self.reject(collection, handle, "append", reason);
}
PimdirPushResult {
handle,
outcome: PimdirPushOutcome::Accepted,
assigned: Some(assigned),
revision: written.revision,
}
}
fn update(
&mut self,
collection: &str,
handle: PimdirHandle,
object: PimdirHash,
if_match: Option<&str>,
) -> PimdirPushResult {
let reader = match self.blob.reader(&object) {
Ok(Some(file)) => file,
Ok(None) => {
let reason = format!("its body {} is missing from the blob tree", object.as_str());
return self.reject(collection, handle, "update", reason);
}
Err(err) => {
return self.reject(collection, handle, "update", format!("{err:#}"));
}
};
let len = match reader.metadata() {
Ok(meta) => meta.len() as usize,
Err(err) => {
return self.reject(collection, handle, "update", format!("{err:#}"));
}
};
let updated = self.pool.primary().update_item_stream(
collection,
handle.as_str(),
reader,
len,
if_match,
);
match updated {
Ok(revision) => PimdirPushResult {
handle,
outcome: PimdirPushOutcome::Accepted,
assigned: None,
revision,
},
Err(err) => self.reject(collection, handle, "update", format!("{err:#}")),
}
}
}
fn accepted(handle: PimdirHandle, assigned: Option<PimdirHandle>) -> PimdirPushResult {
PimdirPushResult {
handle,
outcome: PimdirPushOutcome::Accepted,
assigned,
revision: None,
}
}
fn rejected_bare(handle: PimdirHandle) -> PimdirPushResult {
PimdirPushResult {
handle,
outcome: PimdirPushOutcome::Rejected,
assigned: None,
revision: None,
}
}
const HEADER_CAP: usize = 256 * 1024;
struct HydrateSink {
writer: PimdirBlobWriter,
hasher: PimdirHasher,
header: Vec<u8>,
header_done: bool,
}
impl HydrateSink {
fn new(writer: PimdirBlobWriter, hasher: PimdirHasher) -> Self {
Self {
writer,
hasher,
header: Vec::new(),
header_done: false,
}
}
fn finish(self) -> Result<(PimdirHash, usize, Vec<u8>)> {
let hash = self.hasher.finish();
let size = self.writer.commit(&hash)? as usize;
Ok((hash, size, self.header))
}
}
impl Write for HydrateSink {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.writer.write_all(buf)?;
self.hasher.update(buf);
if !self.header_done {
let from = self.header.len().saturating_sub(3);
self.header.extend_from_slice(buf);
if let Some(end) = header_boundary(&self.header[from..]) {
self.header.truncate(from + end);
self.header_done = true;
} else if self.header.len() >= HEADER_CAP {
self.header.truncate(HEADER_CAP);
self.header_done = true;
}
}
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
fn header_boundary(buf: &[u8]) -> Option<usize> {
buf.windows(4)
.position(|w| w == b"\r\n\r\n")
.map(|i| i + 4)
.or_else(|| buf.windows(2).position(|w| w == b"\n\n").map(|i| i + 2))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_hub_id_becomes_the_name_its_server_knows() {
assert_eq!(
wire_name("imap", "imap/Archives.Charlie"),
"Archives.Charlie"
);
assert_eq!(wire_name("mail", "mail/INBOX"), "INBOX");
assert_eq!(wire_name("mail", "mail/Archive/2026"), "Archive/2026");
assert_eq!(wire_name("mail", "mailbox/INBOX"), "mailbox/INBOX");
assert_eq!(wire_name("cards", "mail/INBOX"), "mail/INBOX");
}
fn member(handle: &str) -> PimdirRemoteItem {
PimdirRemoteItem {
handle: PimdirHandle(handle.into()),
flags: PimdirFlags::default(),
revision: None,
}
}
#[test]
fn a_create_answered_with_a_handle_the_side_holds_is_refused() {
let mut held = HeldHandles::default();
held.remember("agenda", &[member("event-1.ics")], &[]);
assert!(
!held.claim("agenda", "event-1.ics"),
"the server answered with a resource it already had",
);
assert!(
held.claim("agenda", "event-2.ics"),
"a genuinely new member is bound",
);
assert!(
!held.claim("agenda", "event-2.ics"),
"and is held from then on, so a second create onto it is refused too",
);
assert!(
held.claim("contacts", "event-1.ics"),
"handles are per collection, two collections never colliding",
);
}
#[test]
fn a_vanished_handle_stops_being_held() {
let mut held = HeldHandles::default();
held.remember("agenda", &[member("event-1.ics")], &[]);
held.remember("agenda", &[], &[PimdirHandle("event-1.ics".into())]);
assert!(held.claim("agenda", "event-1.ics"));
}
}