use vti_common::error::AppError;
use vti_common::identifier::validate_did;
use vti_common::store::KeyspaceHandle;
use super::{Record, RecordStatus, Room};
use crate::wire::EpochLink;
pub const ROOMS_PREFIX: &str = "rooms:";
pub const RECORDS_PREFIX: &str = "room_records:";
pub const EPOCH_LINKS_PREFIX: &str = "room_epoch_links:";
fn room_key(room_id: &str) -> String {
format!("{ROOMS_PREFIX}{room_id}")
}
fn record_key(room_id: &str, key: &str) -> String {
format!("{RECORDS_PREFIX}{room_id}:{key}")
}
fn record_prefix(room_id: &str) -> String {
format!("{RECORDS_PREFIX}{room_id}:")
}
fn epoch_link_key(room_id: &str, epoch: u32) -> String {
format!("{EPOCH_LINKS_PREFIX}{room_id}:{epoch:010}")
}
fn epoch_link_prefix(room_id: &str) -> String {
format!("{EPOCH_LINKS_PREFIX}{room_id}:")
}
pub async fn put_epoch_link(
links: &KeyspaceHandle,
room_id: &str,
link: &EpochLink,
) -> Result<(), AppError> {
if link.epoch < 2 {
return Err(AppError::Validation(format!(
"epoch {} has no predecessor to link to",
link.epoch
)));
}
let k = epoch_link_key(room_id, link.epoch);
if links.get_raw(k.clone()).await?.is_some() {
return Err(AppError::Conflict(format!(
"room `{room_id}` already has an epoch link at {}",
link.epoch
)));
}
links.insert(k, link).await
}
pub async fn list_epoch_links(
links: &KeyspaceHandle,
room_id: &str,
) -> Result<Vec<EpochLink>, AppError> {
let pairs = links.prefix_iter_raw(epoch_link_prefix(room_id)).await?;
let mut out = Vec::with_capacity(pairs.len());
for (_k, v) in pairs {
let link: EpochLink = serde_json::from_slice(&v)
.map_err(|e| AppError::Internal(format!("decode epoch link in `{room_id}`: {e}")))?;
out.push(link);
}
out.sort_by_key(|l| l.epoch);
Ok(out)
}
pub async fn prune_epoch_links_before(
links: &KeyspaceHandle,
room_id: &str,
epoch: u32,
) -> Result<usize, AppError> {
let pairs = links.prefix_iter_raw(epoch_link_prefix(room_id)).await?;
let mut dropped = 0;
for (k, v) in pairs {
let link: EpochLink = serde_json::from_slice(&v)
.map_err(|e| AppError::Internal(format!("decode epoch link in `{room_id}`: {e}")))?;
if link.epoch <= epoch {
let key = String::from_utf8(k)
.map_err(|e| AppError::Internal(format!("epoch link key is not utf-8: {e}")))?;
links.remove(key).await?;
dropped += 1;
}
}
Ok(dropped)
}
pub async fn create_room(rooms: &KeyspaceHandle, room: &Room) -> Result<(), AppError> {
validate_did("ownerDid", &room.owner_did)?;
if rooms.get_raw(room_key(&room.room_id)).await?.is_some() {
return Err(AppError::Conflict(format!(
"room `{}` is already registered here",
room.room_id
)));
}
rooms.insert(room_key(&room.room_id), room).await
}
fn epoch_lifetime_seconds() -> u64 {
u64::from(crate::lifecycle::DEFAULT_EPOCH_LIFETIME_DAYS) * 24 * 60 * 60
}
pub async fn get_room(rooms: &KeyspaceHandle, room_id: &str) -> Result<Room, AppError> {
let raw = rooms
.get_raw(room_key(room_id))
.await?
.ok_or_else(|| AppError::NotFound(format!("room `{room_id}` not found")))?;
serde_json::from_slice(&raw)
.map_err(|e| AppError::Internal(format!("decode room `{room_id}`: {e}")))
}
pub async fn advance_epoch(
rooms: &KeyspaceHandle,
room_id: &str,
new_epoch: u32,
now: u64,
) -> Result<Room, AppError> {
let mut room = get_room(rooms, room_id).await?;
if new_epoch != room.epoch + 1 {
return Err(AppError::Validation(format!(
"epoch must advance by exactly one: room `{room_id}` is at {}, got {new_epoch}",
room.epoch
)));
}
room.epoch = new_epoch;
room.epoch_expires_at = Some(now + epoch_lifetime_seconds());
room.updated_at = now;
rooms.insert(room_key(room_id), &room).await?;
Ok(room)
}
pub async fn set_owner(
rooms: &KeyspaceHandle,
room_id: &str,
new_owner_did: &str,
now: u64,
) -> Result<Room, AppError> {
validate_did("ownerDid", new_owner_did)?;
let mut room = get_room(rooms, room_id).await?;
room.owner_did = new_owner_did.to_string();
room.updated_at = now;
rooms.insert(room_key(room_id), &room).await?;
Ok(room)
}
pub async fn put_record(
rooms: &KeyspaceHandle,
records: &KeyspaceHandle,
room_id: &str,
mut record: Record,
expected_version: Option<u64>,
now: u64,
) -> Result<Record, AppError> {
let mut room = get_room(rooms, room_id).await?;
let existing = get_record(records, room_id, &record.key).await.ok();
if let Some(expected) = expected_version {
match (&existing, expected) {
(Some(_), 0) => {
return Err(AppError::Conflict(format!(
"record `{}` already exists (create-only requested)",
record.key
)));
}
(None, 0) => {}
(Some(current), n) if current.version != n => {
return Err(AppError::Conflict(format!(
"record `{}` is at version {}, not {n}",
record.key, current.version
)));
}
(None, n) => {
return Err(AppError::Conflict(format!(
"record `{}` does not exist, so it cannot be at version {n}",
record.key
)));
}
_ => {}
}
}
if room.visibility.stores_cleartext() {
if record.sealed.is_some() {
return Err(AppError::Validation(
"an open room stores cleartext records; `sealed` was supplied".into(),
));
}
if record.cleartext.is_none() {
return Err(AppError::Validation(
"an open room requires `cleartext`".into(),
));
}
} else {
if record.cleartext.is_some() {
return Err(AppError::Validation(format!(
"room `{room_id}` is {:?}; cleartext must not be stored here",
room.visibility
)));
}
if record.sealed.is_none() {
return Err(AppError::Validation(
"a sealed room requires `sealed` content".into(),
));
}
if record.epoch != Some(room.epoch) {
return Err(AppError::Validation(format!(
"record is sealed under epoch {:?}, room `{room_id}` is at {}",
record.epoch, room.epoch
)));
}
}
if matches!(room.visibility, super::Visibility::Private) && record.author.is_some() {
return Err(AppError::Validation(
"a private room does not record an author; authorship belongs inside the sealed body"
.into(),
));
}
record.version = room.next_version;
record.updated_at = now;
room.next_version += 1;
room.updated_at = now;
records
.insert(record_key(room_id, &record.key), &record)
.await?;
rooms.insert(room_key(room_id), &room).await?;
Ok(record)
}
pub async fn store_mirrored_record(
rooms: &KeyspaceHandle,
records: &KeyspaceHandle,
room_id: &str,
record: Record,
now: u64,
) -> Result<Record, AppError> {
let mut room = get_room(rooms, room_id).await?;
if !room.is_mirror() {
return Err(AppError::Validation(format!(
"room `{room_id}` is primaried here; a mirrored record has no meaning on a primary"
)));
}
if record.version <= room.watermark() && room.watermark() > 0 {
return Err(AppError::Conflict(format!(
"record `{}` arrived at version {}, at or below this mirror's watermark {}; a mirror does not go backwards",
record.key,
record.version,
room.watermark()
)));
}
room.next_version = record.version + 1;
room.updated_at = now;
records
.insert(record_key(room_id, &record.key), &record)
.await?;
rooms.insert(room_key(room_id), &room).await?;
Ok(record)
}
pub async fn get_record(
records: &KeyspaceHandle,
room_id: &str,
key: &str,
) -> Result<Record, AppError> {
let raw = records
.get_raw(record_key(room_id, key))
.await?
.ok_or_else(|| AppError::NotFound(format!("record `{key}` not found in `{room_id}`")))?;
serde_json::from_slice(&raw)
.map_err(|e| AppError::Internal(format!("decode record `{key}`: {e}")))
}
pub async fn list_rooms(rooms: &KeyspaceHandle) -> Result<Vec<Room>, AppError> {
let pairs = rooms.prefix_iter_raw(ROOMS_PREFIX.to_string()).await?;
let mut out = Vec::with_capacity(pairs.len());
for (k, v) in pairs {
let room: Room = serde_json::from_slice(&v).map_err(|e| {
let which = String::from_utf8_lossy(&k).to_string();
AppError::Internal(format!("decode room at `{which}`: {e}"))
})?;
out.push(room);
}
out.sort_by(|a, b| a.room_id.cmp(&b.room_id));
Ok(out)
}
pub async fn list_records(
records: &KeyspaceHandle,
room_id: &str,
key_prefix: Option<&str>,
since_version: Option<u64>,
) -> Result<Vec<Record>, AppError> {
let scan_prefix = match key_prefix {
Some(p) => format!("{}{p}", record_prefix(room_id)),
None => record_prefix(room_id),
};
let pairs = records.prefix_iter_raw(scan_prefix).await?;
let mut out = Vec::with_capacity(pairs.len());
for (_k, v) in pairs {
let record: Record = serde_json::from_slice(&v)
.map_err(|e| AppError::Internal(format!("decode record in `{room_id}`: {e}")))?;
if let Some(since) = since_version
&& record.version <= since
{
continue;
}
out.push(record);
}
out.sort_by_key(|r| r.version);
Ok(out)
}
#[derive(Debug, Clone, Default)]
pub struct Curation {
pub status: Option<RecordStatus>,
pub pinned: Option<bool>,
pub expected_version: Option<u64>,
}
pub async fn curate_record(
rooms: &KeyspaceHandle,
records: &KeyspaceHandle,
room_id: &str,
key: &str,
curation: Curation,
now: u64,
) -> Result<Record, AppError> {
let Curation {
status,
pinned,
expected_version,
} = curation;
let mut record = get_record(records, room_id, key).await?;
if let Some(expected) = expected_version
&& record.version != expected
{
return Err(AppError::Conflict(format!(
"record `{key}` is at version {}, not {expected}",
record.version
)));
}
if matches!(record.status, RecordStatus::Retracted)
&& matches!(status, Some(RecordStatus::Active))
{
return Err(AppError::Validation(format!(
"record `{key}` is retracted; its body is gone and a status change cannot bring \
it back"
)));
}
let mut room = get_room(rooms, room_id).await?;
if let Some(status) = status {
record.status = status;
if matches!(status, RecordStatus::Retracted) {
record.sealed = None;
record.nonce = None;
record.cleartext = None;
}
}
if let Some(pinned) = pinned {
record.pinned = pinned;
}
record.version = room.next_version;
record.updated_at = now;
room.next_version += 1;
room.updated_at = now;
records.insert(record_key(room_id, key), &record).await?;
rooms.insert(room_key(room_id), &room).await?;
Ok(record)
}
pub async fn purge_record(
records: &KeyspaceHandle,
room_id: &str,
key: &str,
) -> Result<(), AppError> {
let k = record_key(room_id, key);
if records.get_raw(k.clone()).await?.is_none() {
return Err(AppError::NotFound(format!(
"record `{key}` not found in `{room_id}`"
)));
}
records.remove(k).await
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Visibility;
use vti_common::config::StoreConfig;
use vti_common::store::Store;
async fn open() -> (tempfile::TempDir, KeyspaceHandle, KeyspaceHandle) {
let dir = tempfile::tempdir().unwrap();
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.unwrap();
let rooms = store.keyspace(crate::ROOMS_KEYSPACE).unwrap();
let records = store.keyspace(crate::ROOM_RECORDS_KEYSPACE).unwrap();
(dir, rooms, records)
}
async fn open_links() -> (tempfile::TempDir, KeyspaceHandle) {
let dir = tempfile::tempdir().unwrap();
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.unwrap();
let links = store.keyspace(crate::ROOM_EPOCH_LINKS_KEYSPACE).unwrap();
(dir, links)
}
fn link(epoch: u32) -> EpochLink {
EpochLink {
epoch,
wrapped: format!("d3JhcHBlZA{epoch}"),
nonce: "bm9uY2U".into(),
}
}
fn room(id: &str, visibility: Visibility) -> Room {
Room {
room_id: id.into(),
owner_did: "did:key:zOwner".into(),
visibility,
retention_policy: crate::RetentionPolicy::Chained,
epoch: 1,
next_version: 1,
retention_days: 90,
epoch_expires_at: None,
created_at: 0,
updated_at: 0,
mirror_of: None,
}
}
fn sealed_record(key: &str, epoch: u32) -> Record {
Record {
key: key.into(),
version: 0,
epoch: Some(epoch),
status: RecordStatus::Active,
pinned: false,
sealed: Some("c2VhbGVk".into()),
nonce: Some("bm9uY2U".into()),
cleartext: None,
author: None,
updated_at: 0,
}
}
fn open_record(key: &str) -> Record {
Record {
key: key.into(),
version: 0,
epoch: None,
status: RecordStatus::Active,
pinned: false,
sealed: None,
nonce: None,
cleartext: Some(serde_json::json!({ "body": "hello" })),
author: Some("did:key:zBob".into()),
updated_at: 0,
}
}
#[tokio::test]
async fn a_room_cannot_be_registered_twice() {
let (_d, rooms, _rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
let err = create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap_err();
assert!(
matches!(err, AppError::Conflict(_)),
"re-registering would reset the epoch and version counter: {err:?}"
);
}
#[tokio::test]
async fn versions_are_monotonic_across_records_in_a_room() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
let a = put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let b = put_record(&rooms, &rec, "r1", open_record("b"), None, 2)
.await
.unwrap();
let a2 = put_record(&rooms, &rec, "r1", open_record("a"), None, 3)
.await
.unwrap();
assert_eq!((a.version, b.version, a2.version), (1, 2, 3));
assert!(
a2.version > b.version,
"rewriting `a` must advance past `b`"
);
}
#[tokio::test]
async fn a_version_precondition_reports_what_it_lost_to() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let err = put_record(&rooms, &rec, "r1", open_record("a"), Some(99), 2)
.await
.unwrap_err();
let msg = format!("{err}");
assert!(
msg.contains("is at version 1"),
"the conflict must carry the current version so a caller need not re-read: {msg}"
);
}
#[tokio::test]
async fn create_only_refuses_an_existing_key() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), Some(0), 1)
.await
.unwrap();
let err = put_record(&rooms, &rec, "r1", open_record("a"), Some(0), 2)
.await
.unwrap_err();
assert!(matches!(err, AppError::Conflict(_)), "{err:?}");
}
#[tokio::test]
async fn a_sealed_room_refuses_cleartext_and_an_open_room_refuses_ciphertext() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("sealed", Visibility::Attributed))
.await
.unwrap();
create_room(&rooms, &room("plain", Visibility::Open))
.await
.unwrap();
let err = put_record(&rooms, &rec, "sealed", open_record("a"), None, 1)
.await
.unwrap_err();
assert!(
format!("{err}").contains("cleartext must not be stored"),
"{err}"
);
let err = put_record(&rooms, &rec, "plain", sealed_record("a", 1), None, 1)
.await
.unwrap_err();
assert!(format!("{err}").contains("stores cleartext"), "{err}");
}
#[tokio::test]
async fn a_private_room_refuses_a_recorded_author() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("p", Visibility::Private))
.await
.unwrap();
let mut r = sealed_record("a", 1);
r.author = Some("did:key:zBob".into());
let err = put_record(&rooms, &rec, "p", r, None, 1).await.unwrap_err();
assert!(
format!("{err}").contains("inside the sealed body"),
"on a private room the author must not reach this service: {err}"
);
}
#[tokio::test]
async fn a_record_sealed_under_a_stale_epoch_is_refused() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Attributed))
.await
.unwrap();
advance_epoch(&rooms, "r1", 2, 10).await.unwrap();
let err = put_record(&rooms, &rec, "r1", sealed_record("a", 1), None, 11)
.await
.unwrap_err();
assert!(format!("{err}").contains("room `r1` is at 2"), "{err}");
}
#[tokio::test]
async fn minting_an_epoch_renews_a_lapsed_room() {
use crate::lifecycle::Lifecycle;
let (_d, rooms, _rec) = open().await;
let mut r = room("r1", Visibility::Open);
r.epoch_expires_at = Some(1_000);
create_room(&rooms, &r).await.expect("create");
let now = 2_000;
assert_eq!(
get_room(&rooms, "r1").await.unwrap().lifecycle(now),
Lifecycle::Lapsed,
"expired an epoch ago"
);
let renewed = advance_epoch(&rooms, "r1", 2, now).await.expect("renew");
assert_eq!(renewed.lifecycle(now), Lifecycle::Live);
assert!(
renewed.epoch_expires_at.expect("renewal sets an expiry") > now,
"a renewal moves the clock forward, it does not merely clear it"
);
assert_eq!(
get_room(&rooms, "r1").await.unwrap().lifecycle(now),
Lifecycle::Live
);
}
#[tokio::test]
async fn an_epoch_must_advance_by_exactly_one() {
let (_d, rooms, _rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Attributed))
.await
.unwrap();
assert!(
advance_epoch(&rooms, "r1", 3, 1).await.is_err(),
"a gap is refused"
);
assert!(
advance_epoch(&rooms, "r1", 1, 1).await.is_err(),
"a repeat is refused"
);
assert_eq!(advance_epoch(&rooms, "r1", 2, 1).await.unwrap().epoch, 2);
}
#[tokio::test]
async fn one_room_never_lists_anothers_records() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("a", Visibility::Open))
.await
.unwrap();
create_room(&rooms, &room("ab", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "a", open_record("x"), None, 1)
.await
.unwrap();
put_record(&rooms, &rec, "ab", open_record("y"), None, 2)
.await
.unwrap();
let listed = list_records(&rec, "a", None, None).await.unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].key, "x");
}
#[tokio::test]
async fn a_watermark_listing_returns_tombstones() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let seen = put_record(&rooms, &rec, "r1", open_record("b"), None, 2)
.await
.unwrap()
.version;
let tomb = curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Retracted),
pinned: None,
expected_version: None,
},
3,
)
.await
.unwrap();
assert!(
tomb.sealed.is_none() && tomb.cleartext.is_none(),
"body is dropped"
);
let changed = list_records(&rec, "r1", Some(""), Some(seen))
.await
.unwrap();
assert_eq!(changed.len(), 1, "only what changed since the watermark");
assert_eq!(changed[0].key, "a");
assert!(
matches!(changed[0].status, RecordStatus::Retracted),
"the retraction must reach a puller, or the record resurrects"
);
}
#[tokio::test]
async fn deprecating_demotes_without_dropping_the_body() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let out = curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Deprecated),
pinned: None,
expected_version: None,
},
2,
)
.await
.unwrap();
assert_eq!(out.status, RecordStatus::Deprecated);
assert!(out.cleartext.is_some(), "a demotion keeps the body");
assert_eq!(out.version, 2, "and assigns a new version to converge on");
}
#[tokio::test]
async fn a_retracted_record_cannot_be_restored() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Retracted),
pinned: None,
expected_version: None,
},
2,
)
.await
.unwrap();
let err = curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Active),
pinned: None,
expected_version: None,
},
3,
)
.await
.unwrap_err();
assert!(format!("{err}").contains("cannot bring it back"), "{err}");
}
#[tokio::test]
async fn a_record_can_be_pinned_and_deprecated_at_once() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let out = curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Deprecated),
pinned: Some(true),
expected_version: None,
},
2,
)
.await
.unwrap();
assert!(out.pinned);
assert_eq!(out.status, RecordStatus::Deprecated);
}
#[tokio::test]
async fn curation_honours_a_version_precondition() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
let err = curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Deprecated),
pinned: None,
expected_version: Some(99),
},
2,
)
.await
.unwrap_err();
assert!(format!("{err}").contains("is at version 1"), "{err}");
}
#[tokio::test]
async fn purge_removes_the_tombstone_and_is_not_idempotent() {
let (_d, rooms, rec) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
put_record(&rooms, &rec, "r1", open_record("a"), None, 1)
.await
.unwrap();
curate_record(
&rooms,
&rec,
"r1",
"a",
Curation {
status: Some(RecordStatus::Retracted),
pinned: None,
expected_version: None,
},
2,
)
.await
.unwrap();
purge_record(&rec, "r1", "a").await.unwrap();
assert!(
list_records(&rec, "r1", None, None)
.await
.unwrap()
.is_empty()
);
assert!(
purge_record(&rec, "r1", "a").await.is_err(),
"purging an absent record reports not-found rather than succeeding quietly"
);
}
#[tokio::test]
async fn an_owner_must_be_a_did() {
let (_d, rooms, _records) = open().await;
create_room(&rooms, &room("did:key:zRoom", Visibility::Open))
.await
.expect("a real DID");
for bad in ["", "alice@example.com", "did:key:z Alice", "not-a-did"] {
set_owner(&rooms, "did:key:zRoom", bad, 1)
.await
.unwrap_err();
}
assert_eq!(
get_room(&rooms, "did:key:zRoom").await.unwrap().owner_did,
"did:key:zOwner",
"and a refused transfer left the room where it was"
);
let moved = set_owner(&rooms, "did:key:zRoom", "did:key:zBob", 99)
.await
.expect("a real DID");
assert_eq!(moved.owner_did, "did:key:zBob");
assert_eq!(moved.updated_at, 99);
}
fn mirror(id: &str, visibility: Visibility) -> Room {
Room {
mirror_of: Some("https://primary.example.org".into()),
..room(id, visibility)
}
}
#[tokio::test]
async fn a_mirrored_record_keeps_the_primarys_version() {
let (_d, rooms, records) = open().await;
create_room(&rooms, &mirror("r1", Visibility::Open))
.await
.unwrap();
let mut pulled = open_record("k1");
pulled.version = 42; let stored = store_mirrored_record(&rooms, &records, "r1", pulled, 10)
.await
.expect("a mirror stores what its primary assigned");
assert_eq!(stored.version, 42, "the primary's number, verbatim");
let room = get_room(&rooms, "r1").await.unwrap();
assert_eq!(room.watermark(), 42);
assert_eq!(room.next_version, 43);
}
#[tokio::test]
async fn a_mirror_does_not_go_backwards() {
let (_d, rooms, records) = open().await;
create_room(&rooms, &mirror("r1", Visibility::Open))
.await
.unwrap();
let mut first = open_record("k1");
first.version = 42;
store_mirrored_record(&rooms, &records, "r1", first, 10)
.await
.unwrap();
let mut replayed = open_record("k2");
replayed.version = 41;
let err = store_mirrored_record(&rooms, &records, "r1", replayed, 11)
.await
.expect_err("a version at or below the watermark is refused");
assert!(
matches!(err, AppError::Conflict(ref m) if m.contains("backwards")),
"got {err:?}"
);
}
#[tokio::test]
async fn a_primary_refuses_a_mirrored_record() {
let (_d, rooms, records) = open().await;
create_room(&rooms, &room("r1", Visibility::Open))
.await
.unwrap();
let mut pulled = open_record("k1");
pulled.version = 42;
let err = store_mirrored_record(&rooms, &records, "r1", pulled, 10)
.await
.expect_err("this host is the primary");
assert!(matches!(err, AppError::Validation(_)), "got {err:?}");
}
#[tokio::test]
async fn a_mirror_copies_a_record_sealed_under_another_epoch() {
let (_d, rooms, records) = open().await;
create_room(&rooms, &mirror("r1", Visibility::Attributed))
.await
.unwrap();
let mut pulled = sealed_record("k1", 7);
pulled.version = 3;
store_mirrored_record(&rooms, &records, "r1", pulled, 10)
.await
.expect("a mirror stores what the primary sealed");
}
#[tokio::test]
async fn links_round_trip_in_epoch_order() {
let (_d, links) = open_links().await;
for e in [4u32, 2, 3] {
put_epoch_link(&links, "r1", &link(e)).await.expect("put");
}
let got = list_epoch_links(&links, "r1").await.expect("list");
assert_eq!(
got.iter().map(|l| l.epoch).collect::<Vec<_>>(),
vec![2, 3, 4]
);
}
#[tokio::test]
async fn a_link_cannot_be_replaced() {
let (_d, links) = open_links().await;
put_epoch_link(&links, "r1", &link(2)).await.expect("put");
let err = put_epoch_link(&links, "r1", &link(2)).await.unwrap_err();
assert!(matches!(err, AppError::Conflict(_)), "got {err:?}");
}
#[tokio::test]
async fn the_first_epoch_has_no_link() {
let (_d, links) = open_links().await;
for e in [0u32, 1] {
assert!(matches!(
put_epoch_link(&links, "r1", &link(e)).await.unwrap_err(),
AppError::Validation(_)
));
}
}
#[tokio::test]
async fn a_scan_is_room_exact() {
let (_d, links) = open_links().await;
put_epoch_link(&links, "a", &link(2)).await.expect("put");
put_epoch_link(&links, "ab", &link(3)).await.expect("put");
let got = list_epoch_links(&links, "a").await.expect("list");
assert_eq!(got.len(), 1, "`a` must not pick up `ab`'s chain");
assert_eq!(got[0].epoch, 2);
}
#[tokio::test]
async fn pruning_severs_at_and_below_the_named_epoch() {
let (_d, links) = open_links().await;
for e in 2..=5u32 {
put_epoch_link(&links, "r1", &link(e)).await.expect("put");
}
let dropped = prune_epoch_links_before(&links, "r1", 3)
.await
.expect("prune");
assert_eq!(dropped, 2, "links 2 and 3");
let left = list_epoch_links(&links, "r1").await.expect("list");
assert_eq!(left.iter().map(|l| l.epoch).collect::<Vec<_>>(), vec![4, 5]);
}
#[tokio::test]
async fn pruning_is_idempotent() {
let (_d, links) = open_links().await;
put_epoch_link(&links, "r1", &link(2)).await.expect("put");
assert_eq!(
prune_epoch_links_before(&links, "r1", 2)
.await
.expect("first"),
1
);
assert_eq!(
prune_epoch_links_before(&links, "r1", 2)
.await
.expect("again"),
0
);
}
#[tokio::test]
async fn a_pruned_link_is_gone_and_the_epoch_is_still_spent() {
let (_d, links) = open_links().await;
put_epoch_link(&links, "r1", &link(2)).await.expect("put");
prune_epoch_links_before(&links, "r1", 2)
.await
.expect("prune");
assert!(
list_epoch_links(&links, "r1")
.await
.expect("list")
.is_empty()
);
}
#[tokio::test]
async fn listing_rooms_returns_every_room_in_identifier_order() {
let (_d, rooms, _r) = open().await;
for id in ["r3", "r1", "r2"] {
create_room(&rooms, &room(id, Visibility::Attributed))
.await
.expect("create");
}
let got = list_rooms(&rooms).await.expect("list");
assert_eq!(
got.iter().map(|r| r.room_id.as_str()).collect::<Vec<_>>(),
vec!["r1", "r2", "r3"],
"an operator's list is stable, or two reads disagree about nothing"
);
}
#[tokio::test]
async fn listing_rooms_is_empty_on_a_host_that_holds_none() {
let (_d, rooms, _r) = open().await;
assert!(list_rooms(&rooms).await.expect("list").is_empty());
}
}