use std::{io::Write, path::Path};
use io_pimdir::{
PimdirBlobs, PimdirError, PimdirProducer, PimdirReader, PimdirSourceStore, PimdirStore,
codec::PimdirAction, hash::PimdirHashAlgo, sql,
};
use io_replica::{
change::{ReplicaDropReason, ReplicaWriteOp},
client::ReplicaStorage,
collection::ReplicaCollectionId,
object::{ReplicaHash, ReplicaObject},
placement::{
ReplicaBase, ReplicaFlags, ReplicaHandle, ReplicaLevel, ReplicaLinkId, ReplicaMeta,
ReplicaPlacement, ReplicaStatus,
},
storage::ReplicaLoadScope,
};
const NOW: &str = "2026-08-07T00:00:00Z";
fn inbox() -> ReplicaCollectionId {
ReplicaCollectionId("INBOX".into())
}
fn placement(handle: &str, link: &str, hash: &str, flags: &[&str]) -> ReplicaPlacement {
let flags = ReplicaFlags::from_iter(flags.iter().copied());
ReplicaPlacement {
sort_key: Default::default(),
collection: inbox(),
handle: ReplicaHandle(handle.into()),
link_id: Some(ReplicaLinkId(link.into())),
object: Some(ReplicaHash(hash.into())),
level: ReplicaLevel::Full,
meta: None,
flags: flags.clone(),
status: ReplicaStatus::Clean,
conflict_revision: None,
conflict_object: None,
base: Some(ReplicaBase {
flags,
revision: None,
object: Some(ReplicaHash(hash.into())),
}),
origin: None,
}
}
fn store_object(hash: &str, body: &[u8]) -> ReplicaWriteOp {
ReplicaWriteOp::StoreObject {
object: ReplicaObject {
hash: ReplicaHash(hash.into()),
size: body.len(),
},
body: Some(body.to_vec()),
}
}
fn blob_exists(dir: &Path, hash: &str) -> bool {
dir.join("objects")
.join(&hash[0..2])
.join(&hash[2..4])
.join(hash)
.exists()
}
fn seeded(dir: &Path) -> (PimdirSourceStore, i64) {
let mut store = PimdirStore::open(dir).unwrap().for_source("local");
store
.write(vec![
store_object("cafebabe", b"abc"),
ReplicaWriteOp::UpsertPlacement(placement("1", "mid:a", "cafebabe", &["\\Seen"])),
])
.unwrap();
let seq = store.list_items("INBOX", None, 10).unwrap()[0].seq;
(store, seq)
}
fn write_blob(dir: &Path, hash: &str, body: &[u8]) -> u64 {
let blobs = PimdirBlobs::open(dir, PimdirHashAlgo::default());
let mut writer = blobs.writer().unwrap();
writer.write_all(body).unwrap();
writer.commit(&ReplicaHash(hash.into())).unwrap()
}
#[test]
fn a_queued_add_round_trips_into_a_staged_item() {
let dir = tempfile::tempdir().unwrap();
let (mut store, _) = seeded(dir.path());
let size = write_blob(dir.path(), "beef0000", b"new body");
let mut producer = PimdirProducer::open(dir.path(), "smtp").unwrap();
let add = PimdirAction::Add {
link_id: Some(ReplicaLinkId("mid:new".into())),
flags: ReplicaFlags::from_iter(["\\Draft"]),
object: Some(ReplicaHash("beef0000".into())),
meta: Some(ReplicaMeta("{\"subject\":\"hi\",\"v\":1}".into())),
handle: Some(ReplicaHandle("draft-1".into())),
};
producer.enqueue("INBOX", &add, Some(size), NOW).unwrap();
let pending = producer.pending_actions("INBOX").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].action, add);
assert_eq!(pending[0].producer, "smtp");
assert_eq!(store.queued_collections().unwrap(), ["INBOX"]);
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
assert!(store.pending_actions("INBOX").unwrap().is_empty());
assert!(store.queued_collections().unwrap().is_empty());
let items = store.list_items("INBOX", None, 10).unwrap();
let item = items.iter().find(|i| i.link_id.0 == "mid:new").unwrap();
assert_eq!(item.object, Some(ReplicaHash("beef0000".into())));
assert_eq!(item.level, ReplicaLevel::Full);
assert!(item.flags.contains("\\Draft"));
assert_eq!(
item.meta.as_ref().unwrap().0,
"{\"subject\":\"hi\",\"v\":1}"
);
let projected = store
.load(&inbox(), &ReplicaLoadScope::All)
.unwrap()
.placements;
let created = projected
.iter()
.find(|p| p.link_id.as_ref().is_some_and(|l| l.0 == "mid:new"))
.unwrap();
assert_ne!(created.status, ReplicaStatus::Clean, "a pending push");
assert!(created.base.is_none(), "no prior sync: an append push");
}
#[test]
fn a_duplicate_add_parks_instead_of_clobbering() {
let dir = tempfile::tempdir().unwrap();
let (mut store, _) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
let add = PimdirAction::Add {
link_id: Some(ReplicaLinkId("mid:a".into())),
flags: ReplicaFlags::default(),
object: None,
meta: None,
handle: None,
};
producer.enqueue("INBOX", &add, None, NOW).unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (0, 1));
let parked = store.parked_actions().unwrap();
assert_eq!(parked.len(), 1);
assert!(parked[0].error.contains("already present"), "{parked:?}");
assert_eq!(parked[0].attempts, 1);
assert!(
store.list_items("INBOX", None, 10).unwrap()[0]
.flags
.contains("\\Seen")
);
}
#[test]
fn set_flags_is_absolute_and_reapplies_idempotently() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
let set = PimdirAction::SetFlags {
seq,
flags: ReplicaFlags::from_iter(["\\Flagged"]),
};
producer.enqueue("INBOX", &set, None, NOW).unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
let item = store.get_item("INBOX", seq).unwrap().unwrap();
assert_eq!(item.flags, ReplicaFlags::from_iter(["\\Flagged"]));
producer.enqueue("INBOX", &set, None, NOW).unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
let item = store.get_item("INBOX", seq).unwrap().unwrap();
assert_eq!(item.flags, ReplicaFlags::from_iter(["\\Flagged"]));
}
#[test]
fn remove_hides_the_item_and_an_absent_remove_succeeds() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
producer
.enqueue("INBOX", &PimdirAction::Remove { seq }, None, NOW)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
assert!(store.get_item("INBOX", seq).unwrap().is_none());
let projected = store
.load(&inbox(), &ReplicaLoadScope::All)
.unwrap()
.placements;
assert_eq!(projected[0].status, ReplicaStatus::Tombstone);
producer
.enqueue("INBOX", &PimdirAction::Remove { seq }, None, NOW)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
assert!(store.parked_actions().unwrap().is_empty());
}
#[test]
fn copy_fills_the_target_and_move_also_empties_the_source() {
let dir = tempfile::tempdir().unwrap();
let mut store = PimdirStore::open(dir.path()).unwrap().for_source("local");
store
.write(vec![
store_object("cafebabe", b"a"),
store_object("cafebabf", b"b"),
ReplicaWriteOp::UpsertPlacement(placement("1", "mid:a", "cafebabe", &[])),
ReplicaWriteOp::UpsertPlacement(placement("2", "mid:b", "cafebabf", &[])),
])
.unwrap();
let seq_of = |store: &PimdirStore, link: &str| {
store
.list_items("INBOX", None, 10)
.unwrap()
.into_iter()
.find(|i| i.link_id.0 == link)
.unwrap()
.seq
};
let seq_a = seq_of(&store, "mid:a");
let seq_b = seq_of(&store, "mid:b");
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::Copy {
seq: seq_a,
to: ReplicaCollectionId("Backup".into()),
},
None,
NOW,
)
.unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::Move {
seq: seq_b,
to: ReplicaCollectionId("Archive".into()),
},
None,
NOW,
)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (2, 0));
assert_eq!(store.count_items("INBOX").unwrap(), 1, "mid:b moved away");
let backup = store.list_items("Backup", None, 10).unwrap();
assert_eq!(backup.len(), 1);
assert_eq!(backup[0].link_id.0, "mid:a");
assert_eq!(backup[0].seq, seq_a, "a message keeps one public id");
let archive = store.list_items("Archive", None, 10).unwrap();
assert_eq!(archive.len(), 1);
assert_eq!(archive[0].link_id.0, "mid:b");
assert_eq!(archive[0].seq, seq_b);
assert!(store.get_item("INBOX", seq_b).unwrap().is_none());
}
#[test]
fn update_repoints_the_body() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let size = write_blob(dir.path(), "beef0000", b"edited body");
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::Update {
seq,
object: ReplicaHash("beef0000".into()),
meta: None,
},
Some(size),
NOW,
)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
let item = store.get_item("INBOX", seq).unwrap().unwrap();
assert_eq!(item.object, Some(ReplicaHash("beef0000".into())));
assert!(blob_exists(dir.path(), "cafebabe"));
assert!(blob_exists(dir.path(), "beef0000"));
let projected = store
.load(&inbox(), &ReplicaLoadScope::All)
.unwrap()
.placements;
assert_eq!(projected[0].status, ReplicaStatus::Dirty);
}
#[test]
fn a_parked_action_does_not_block_later_actions() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
producer
.enqueue("INBOX", &PimdirAction::Remove { seq: 9999 }, None, NOW)
.unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::SetFlags {
seq: 9999,
flags: ReplicaFlags::default(),
},
None,
NOW,
)
.unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::SetFlags {
seq,
flags: ReplicaFlags::from_iter(["\\Answered"]),
},
None,
NOW,
)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (2, 1));
let item = store.get_item("INBOX", seq).unwrap().unwrap();
assert!(item.flags.contains("\\Answered"));
let parked = store.parked_actions().unwrap();
assert_eq!(parked.len(), 1);
assert_eq!(parked[0].collection, "INBOX");
assert_eq!(parked[0].action, "set-flags");
assert!(parked[0].error.contains("unknown seq"), "{parked:?}");
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (0, 0));
assert_eq!(store.parked_actions().unwrap().len(), 1, "never deleted");
}
#[test]
fn gc_never_sweeps_a_queued_body() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seeded_seq) = seeded(dir.path());
let size = write_blob(dir.path(), "beef0000", b"queued body");
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::Add {
link_id: Some(ReplicaLinkId("mid:new".into())),
flags: ReplicaFlags::default(),
object: Some(ReplicaHash("beef0000".into())),
meta: None,
handle: Some(ReplicaHandle("draft-1".into())),
},
Some(size),
NOW,
)
.unwrap();
assert!(matches!(
store.collect_garbage(),
Err(PimdirError::Staging(_))
));
drop(producer);
store
.write(vec![ReplicaWriteOp::DropPlacement {
collection: inbox(),
handle: ReplicaHandle("1".into()),
reason: ReplicaDropReason::Deleted,
}])
.unwrap();
assert!(store.purge(&inbox(), seeded_seq).unwrap());
let collected = store.collect_garbage().unwrap();
assert_eq!((collected.objects, collected.blobs), (1, 1));
assert!(!blob_exists(dir.path(), "cafebabe"), "the orphan is taken");
assert!(
blob_exists(dir.path(), "beef0000"),
"the queued body is not"
);
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked), (1, 0));
assert!(blob_exists(dir.path(), "beef0000"), "now the item pins it");
let seq = store.seq_for_link("INBOX", "mid:new").unwrap().unwrap();
store
.write(vec![ReplicaWriteOp::DropPlacement {
collection: inbox(),
handle: ReplicaHandle("draft-1".into()),
reason: ReplicaDropReason::Deleted,
}])
.unwrap();
assert!(store.purge(&inbox(), seq).unwrap());
assert_eq!(store.collect_garbage().unwrap().blobs, 1);
assert!(!blob_exists(dir.path(), "beef0000"), "no refcount leak");
}
#[test]
fn an_unknown_kind_is_skipped_and_never_blocks_the_queue() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "himalaya").unwrap();
let size = write_blob(dir.path(), "beef0000", b"a message to send");
let submit = PimdirAction::Unknown {
kind: "submit".into(),
payload: "{\"v\":1,\"object\":\"beef0000\",\"to\":[\"a@b.c\"]}".into(),
object_hash: Some(ReplicaHash("beef0000".into())),
};
let id = producer.enqueue("INBOX", &submit, Some(size), NOW).unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::SetFlags {
seq,
flags: ReplicaFlags::from_iter(["\\Answered"]),
},
None,
NOW,
)
.unwrap();
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked, report.skipped), (1, 0, 1));
assert!(
store
.get_item("INBOX", seq)
.unwrap()
.unwrap()
.flags
.contains("\\Answered")
);
assert!(store.parked_actions().unwrap().is_empty());
let pending = store.pending_actions("INBOX").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, id);
assert_eq!(pending[0].action, submit);
assert!(blob_exists(dir.path(), "beef0000"));
let report = store.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.parked, report.skipped), (0, 0, 1));
assert_eq!(store.pending_actions("INBOX").unwrap()[0].attempts, 0);
}
#[test]
fn an_acknowledged_action_releases_its_queued_body() {
let dir = tempfile::tempdir().unwrap();
let (mut store, _) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "himalaya").unwrap();
let size = write_blob(dir.path(), "beef0000", b"sent already");
let id = producer
.enqueue(
"INBOX",
&PimdirAction::Unknown {
kind: "submit".into(),
payload: "{\"v\":1,\"object\":\"beef0000\"}".into(),
object_hash: Some(ReplicaHash("beef0000".into())),
},
Some(size),
NOW,
)
.unwrap();
assert!(blob_exists(dir.path(), "beef0000"));
assert!(store.drop_action(id).unwrap());
assert!(store.pending_actions("INBOX").unwrap().is_empty());
assert!(blob_exists(dir.path(), "beef0000"), "no write collects");
drop(producer);
assert_eq!(store.collect_garbage().unwrap().blobs, 1);
assert!(
!blob_exists(dir.path(), "beef0000"),
"the pin went with the row"
);
assert!(!store.drop_action(id).unwrap());
}
#[test]
fn a_failed_action_retries_until_it_is_parked() {
let dir = tempfile::tempdir().unwrap();
let (mut store, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "test").unwrap();
let id = producer
.enqueue("INBOX", &PimdirAction::Remove { seq }, None, NOW)
.unwrap();
store.fail_action(id, None).unwrap();
store.fail_action(id, None).unwrap();
let pending = store.pending_actions("INBOX").unwrap();
assert_eq!(pending.len(), 1, "a transient failure stays pending");
assert_eq!(pending[0].attempts, 2);
store.fail_action(id, Some("the channel is gone")).unwrap();
assert!(store.pending_actions("INBOX").unwrap().is_empty());
let parked = store.parked_actions().unwrap();
assert_eq!(parked.len(), 1);
assert_eq!(parked[0].id, id);
assert_eq!(parked[0].attempts, 3, "the parking attempt counts too");
assert_eq!(parked[0].error, "the channel is gone");
assert!(store.drop_action(id).unwrap());
assert!(store.parked_actions().unwrap().is_empty());
store.fail_action(id, Some("gone")).unwrap();
}
#[test]
fn a_generation_bump_is_visible_to_a_reader() {
let dir = tempfile::tempdir().unwrap();
let mut store = PimdirStore::open(dir.path()).unwrap().for_source("local");
store.ensure_collection("INBOX", "message/rfc822").unwrap();
assert_eq!(store.generation("INBOX").unwrap(), Some(1));
store
.write(vec![
store_object("cafebabe", b"abc"),
ReplicaWriteOp::UpsertPlacement(placement("1", "mid:a", "cafebabe", &[])),
])
.unwrap();
assert_eq!(store.generation("INBOX").unwrap(), Some(1));
let mut rekeyed = placement("101", "mid:a", "cafebabe", &[]);
rekeyed.status = ReplicaStatus::Clean;
let generation = store
.write_rekeyed(
"INBOX",
vec![
ReplicaWriteOp::DropPlacement {
collection: inbox(),
handle: ReplicaHandle("1".into()),
reason: ReplicaDropReason::Superseded,
},
ReplicaWriteOp::UpsertPlacement(rekeyed),
],
)
.unwrap();
assert_eq!(generation, 2);
let carried = store.load(&inbox(), &ReplicaLoadScope::All).unwrap();
assert_eq!(carried.placements.len(), 1);
assert_eq!(carried.placements[0].handle, ReplicaHandle("101".into()));
let reader = PimdirStore::open(dir.path()).unwrap();
assert_eq!(reader.generation("INBOX").unwrap(), Some(2));
let collections = reader.list_collections().unwrap();
assert_eq!(collections[0].generation, 2);
assert_eq!(reader.generation("nope").unwrap(), None);
}
#[test]
fn a_store_stamped_with_a_higher_version_is_refused() {
let dir = tempfile::tempdir().unwrap();
{
let _store = PimdirStore::open(dir.path()).unwrap();
let conn = rusqlite::Connection::open(dir.path().join("pimdir.db")).unwrap();
let user_version: i64 = conn
.pragma_query_value(None, "user_version", |r| r.get(0))
.unwrap();
let meta_version: i64 = conn
.query_row("SELECT version FROM store_meta WHERE id = 1", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!((user_version, meta_version), (sql::VERSION, sql::VERSION));
}
{
let conn = rusqlite::Connection::open(dir.path().join("pimdir.db")).unwrap();
conn.pragma_update(None, "user_version", sql::VERSION + 1)
.unwrap();
}
for result in [
PimdirStore::open(dir.path()).map(drop),
PimdirProducer::open(dir.path(), "test").map(drop),
PimdirReader::open(dir.path()).map(drop),
] {
match result {
Err(PimdirError::Version { found }) => assert_eq!(found, sql::VERSION + 1),
Err(other) => panic!("expected a version refusal, got {other:?}"),
Ok(()) => panic!("expected a version refusal, got an open store"),
}
}
let fresh = tempfile::tempdir().unwrap();
rusqlite::Connection::open(fresh.path().join("pimdir.db")).unwrap();
match PimdirProducer::open(fresh.path(), "test") {
Err(PimdirError::Uncreated) => {}
Err(other) => panic!("expected an uncreated-store refusal, got {other:?}"),
Ok(_) => panic!("expected an uncreated-store refusal, got a producer"),
}
}
#[test]
fn an_action_the_draining_source_cannot_place_is_skipped_not_parked() {
let dir = tempfile::tempdir().unwrap();
let (mut owner, seq) = seeded(dir.path());
let mut producer = PimdirProducer::open(dir.path(), "frontend").unwrap();
producer
.enqueue(
"INBOX",
&PimdirAction::SetFlags {
seq,
flags: ReplicaFlags::default(),
},
None,
NOW,
)
.unwrap();
let mut stranger = PimdirStore::open(dir.path()).unwrap().for_source("other");
let report = stranger.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.skipped, report.parked), (0, 1, 0));
assert!(stranger.parked_actions().unwrap().is_empty());
let report = owner.drain_collection("INBOX").unwrap();
assert_eq!((report.applied, report.skipped, report.parked), (1, 0, 0));
let item = owner.get_item("INBOX", seq).unwrap().unwrap();
assert_eq!(item.flags, ReplicaFlags::default());
}