use super::dek::DriveVaultKey;
use super::envelope::{self, ObjectKind};
use super::pack::{Pack, PackEntry};
use super::store::VaultObjectStore;
use crate::db::Db;
use crate::errors::AtomicResult;
use crate::loro::AtomicLoroDoc;
use crate::resources::Resource;
use crate::storelike::Storelike;
use crate::Subject;
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, HashSet};
pub type VersionVectorMap = BTreeMap<String, i32>;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct LaneState {
#[serde(default)]
pub cursors: BTreeMap<String, VersionVectorMap>,
#[serde(default)]
pub segments_since_checkpoint: u32,
#[serde(default)]
pub bytes_since_checkpoint: u64,
#[serde(default)]
pub last_checkpoint_bytes: u64,
#[serde(default)]
pub imported_lanes: BTreeMap<String, u32>,
}
impl LaneState {
fn decode(bytes: &[u8]) -> Self {
if let Ok(state) = serde_json::from_slice::<LaneState>(bytes) {
return state;
}
if let Ok(subjects) = serde_json::from_slice::<Vec<String>>(bytes) {
return Self {
cursors: subjects
.into_iter()
.map(|subject| (subject, VersionVectorMap::new()))
.collect(),
..Default::default()
};
}
Self::default()
}
fn known_subjects(&self) -> impl Iterator<Item = &String> {
self.cursors.keys()
}
}
fn lane_state_key(drive_pseudonym: &str, device_pubkey: &str) -> Vec<u8> {
format!("vault-lane:{drive_pseudonym}:{device_pubkey}").into_bytes()
}
fn pending_lane_state_key(drive_pseudonym: &str, device_pubkey: &str, segment: u32) -> Vec<u8> {
format!("vault-lane-pending:{drive_pseudonym}:{device_pubkey}:{segment}").into_bytes()
}
pub fn read_lane_state(store: &Db, drive_pseudonym: &str, device_pubkey: &str) -> LaneState {
store
.kv
.get(
crate::db::trees::Tree::PluginMeta,
&lane_state_key(drive_pseudonym, device_pubkey),
)
.ok()
.flatten()
.map(|bytes| LaneState::decode(&bytes))
.unwrap_or_default()
}
fn write_lane_state(store: &Db, drive_pseudonym: &str, device_pubkey: &str, state: &LaneState) {
if let Ok(bytes) = serde_json::to_vec(state) {
let _ = store.kv.insert(
crate::db::trees::Tree::PluginMeta,
&lane_state_key(drive_pseudonym, device_pubkey),
&bytes,
);
}
}
pub fn commit_lane_state(
store: &Db,
drive_pseudonym: &str,
device_pubkey: &str,
segment: u32,
) -> AtomicResult<()> {
let pending_key = pending_lane_state_key(drive_pseudonym, device_pubkey, segment);
let Some(bytes) = store
.kv
.get(crate::db::trees::Tree::PluginMeta, &pending_key)
.ok()
.flatten()
else {
return Ok(());
};
let state: LaneState =
serde_json::from_slice(&bytes).map_err(|e| format!("malformed pending lane state: {e}"))?;
write_lane_state(store, drive_pseudonym, device_pubkey, &state);
let _ = store
.kv
.remove(crate::db::trees::Tree::PluginMeta, &pending_key);
Ok(())
}
pub fn record_imported_lane(
store: &Db,
drive_pseudonym: &str,
device_pubkey: &str,
lane: &str,
segment: u32,
) {
let mut state = read_lane_state(store, drive_pseudonym, device_pubkey);
let slot = state.imported_lanes.entry(lane.to_string()).or_insert(0);
if segment <= *slot {
return;
}
*slot = segment;
write_lane_state(store, drive_pseudonym, device_pubkey, &state);
}
pub fn drive_prefix(drive_pseudonym: &str) -> String {
format!("vault/{drive_pseudonym}/")
}
pub fn lane_prefix(drive_pseudonym: &str, device_pubkey: &str) -> String {
format!("vault/{drive_pseudonym}/lanes/{device_pubkey}/")
}
pub fn segment_key(drive_pseudonym: &str, device_pubkey: &str, segment: u32) -> String {
format!(
"{}seg-{segment:06}.pack",
lane_prefix(drive_pseudonym, device_pubkey)
)
}
pub fn checkpoint_prefix(drive_pseudonym: &str) -> String {
format!("vault/{drive_pseudonym}/checkpoints/")
}
pub fn checkpoint_key(drive_pseudonym: &str, checkpoint_n: u64) -> String {
format!(
"{}ckpt-{checkpoint_n:06}.loro",
checkpoint_prefix(drive_pseudonym)
)
}
pub fn parse_segment_key(object_key: &str) -> Option<(String, u32)> {
let rest = object_key.split("/lanes/").nth(1)?;
let (device, file) = rest.split_once('/')?;
let number = file.strip_prefix("seg-")?.strip_suffix(".pack")?;
Some((device.to_string(), number.parse().ok()?))
}
pub fn parse_checkpoint_key(object_key: &str) -> Option<u64> {
let file = object_key.split("/checkpoints/").nth(1)?;
file.strip_prefix("ckpt-")?
.strip_suffix(".loro")?
.parse()
.ok()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SegmentKind {
Pack,
Checkpoint,
}
#[derive(Debug, Clone, Copy)]
pub struct CheckpointPolicy {
pub max_segments: u32,
pub bytes_ratio: f64,
}
impl Default for CheckpointPolicy {
fn default() -> Self {
Self {
max_segments: 64,
bytes_ratio: 1.0,
}
}
}
impl CheckpointPolicy {
fn wants_checkpoint(&self, state: &LaneState, drive_has_checkpoint: bool) -> bool {
if !drive_has_checkpoint {
return true;
}
if state.segments_since_checkpoint >= self.max_segments {
return true;
}
state.last_checkpoint_bytes > 0
&& state.bytes_since_checkpoint as f64
>= state.last_checkpoint_bytes as f64 * self.bytes_ratio
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BackupSummary {
pub kind: SegmentKind,
pub object_key: String,
pub resources: usize,
pub unchanged: usize,
pub tombstones: usize,
pub sealed_bytes: usize,
pub coverage: BTreeMap<String, u32>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct RestoreSummary {
pub packs_read: usize,
pub resources_restored: usize,
pub tombstones_applied: usize,
pub objects_skipped: usize,
pub objects_unreadable: usize,
}
async fn drive_agent_subjects(store: &Db, drive: &Subject) -> Vec<String> {
let Ok(resource) = store.get_resource(drive).await else {
return Vec::new();
};
[crate::urls::READ, crate::urls::WRITE]
.iter()
.filter_map(|right| resource.get(right).ok())
.filter_map(|value| value.to_subjects(None).ok())
.flatten()
.filter(|subject| Subject::from_raw(subject, None).is_agent_did())
.collect()
}
fn cheap_version_vector(store: &Db, subject: &str) -> Option<VersionVectorMap> {
let snapshot = store
.kv
.get(crate::db::trees::Tree::LoroSnapshots, subject.as_bytes())
.ok()
.flatten()?;
let vv = AtomicLoroDoc::vv_map_from_snapshot(&snapshot).ok()?;
Some(vv.into_iter().collect())
}
enum Contribution {
Update(Vec<u8>, VersionVectorMap),
Unchanged,
Unreadable,
}
async fn contribution(
store: &Db,
subject_str: &str,
cursor: Option<&VersionVectorMap>,
full: bool,
) -> Contribution {
if !full {
if let Some(cursor) = cursor {
if cheap_version_vector(store, subject_str).as_ref() == Some(cursor) {
return Contribution::Unchanged;
}
}
}
let subject = Subject::from_raw(subject_str, None);
let Ok(resource) = store.get_resource(&subject).await else {
return Contribution::Unreadable;
};
let Ok(doc) = resource.build_state_doc() else {
return Contribution::Unreadable;
};
let since = match (full, cursor) {
(false, Some(cursor)) => AtomicLoroDoc::vv_from_map(&cursor.clone().into_iter().collect()),
_ => Default::default(),
};
let update = doc.export_updates_since(&since);
let reached: VersionVectorMap = doc.oplog_vv_map().into_iter().collect();
if update.is_empty() {
return Contribution::Unchanged;
}
Contribution::Update(update, reached)
}
#[allow(clippy::too_many_arguments)]
pub async fn export_vault_segment(
store: &Db,
drive: &Subject,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
drive_pseudonym: &str,
device_pubkey: &str,
segment: u32,
checkpoint_n: u64,
drive_has_checkpoint: bool,
observed_lanes: &BTreeMap<String, u32>,
policy: CheckpointPolicy,
) -> AtomicResult<Option<BackupSummary>> {
let mut state = read_lane_state(store, drive_pseudonym, device_pubkey);
let full = policy.wants_checkpoint(&state, drive_has_checkpoint);
let mut subjects: HashSet<String> =
crate::sync::engine::collect_drive_subjects(store, drive).await;
subjects.extend(drive_agent_subjects(store, drive).await);
let mut entries = Vec::new();
let mut cursors: BTreeMap<String, VersionVectorMap> = BTreeMap::new();
let mut unchanged = 0usize;
for subject_str in &subjects {
match contribution(store, subject_str, state.cursors.get(subject_str), full).await {
Contribution::Update(update, reached) => {
entries.push(PackEntry {
subject: subject_str.clone(),
update,
});
cursors.insert(subject_str.clone(), reached);
}
Contribution::Unchanged => {
unchanged += 1;
if let Some(existing) = state.cursors.get(subject_str) {
cursors.insert(subject_str.clone(), existing.clone());
} else if let Some(current) = cheap_version_vector(store, subject_str) {
cursors.insert(subject_str.clone(), current);
}
}
Contribution::Unreadable => {
if let Some(existing) = state.cursors.get(subject_str) {
cursors.insert(subject_str.clone(), existing.clone());
}
}
}
}
entries.sort_by(|a, b| a.subject.cmp(&b.subject));
let mut tombstones = Vec::new();
let mut carried = Vec::new();
for subject in state.known_subjects() {
if cursors.contains_key(subject) {
continue;
}
if crate::sync::tombstones::is_tombstoned(store, subject) {
tombstones.push(subject.clone());
} else {
carried.push(subject.clone());
}
}
for subject in carried {
if let Some(cursor) = state.cursors.get(&subject) {
cursors.insert(subject, cursor.clone());
}
}
let (pack, kind, object_key, coverage) = if full {
let mut coverage = state.imported_lanes.clone();
if segment > 1 {
coverage.insert(device_pubkey.to_string(), segment - 1);
}
let observed = observed_lanes.clone();
(
Pack::checkpoint(entries, tombstones, coverage.clone(), observed),
SegmentKind::Checkpoint,
checkpoint_key(drive_pseudonym, checkpoint_n),
coverage,
)
} else {
(
Pack::new(entries, tombstones),
SegmentKind::Pack,
segment_key(drive_pseudonym, device_pubkey, segment),
BTreeMap::new(),
)
};
if pack.is_empty() {
return Ok(None);
}
let resources = pack.entries.len();
let tombstone_count = pack.tombstones.len();
let object_kind = match kind {
SegmentKind::Pack => ObjectKind::Pack,
SegmentKind::Checkpoint => ObjectKind::Checkpoint,
};
let sealed = envelope::seal(key, object_kind, &pack.encode()?)?;
vault.put(&object_key, &sealed)?;
state.cursors = cursors;
match kind {
SegmentKind::Checkpoint => {
state.segments_since_checkpoint = 0;
state.bytes_since_checkpoint = 0;
state.last_checkpoint_bytes = sealed.len() as u64;
}
SegmentKind::Pack => {
state.segments_since_checkpoint = state.segments_since_checkpoint.saturating_add(1);
state.bytes_since_checkpoint = state
.bytes_since_checkpoint
.saturating_add(sealed.len() as u64);
}
}
if let Ok(bytes) = serde_json::to_vec(&state) {
let _ = store.kv.insert(
crate::db::trees::Tree::PluginMeta,
&pending_lane_state_key(drive_pseudonym, device_pubkey, segment),
&bytes,
);
}
Ok(Some(BackupSummary {
kind,
object_key,
resources,
unchanged,
tombstones: tombstone_count,
sealed_bytes: sealed.len(),
coverage,
}))
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct RestorePlan {
pub order: Vec<String>,
pub skipped: Vec<String>,
}
pub fn plan_restore(
keys: &[String],
anchor: Option<&str>,
coverage: &BTreeMap<String, u32>,
observed: &BTreeMap<String, u32>,
) -> RestorePlan {
let Some(anchor) = anchor else {
let mut order = keys.to_vec();
order.sort();
return RestorePlan {
order,
skipped: Vec::new(),
};
};
let mut before = Vec::new();
let mut after = Vec::new();
let mut skipped = Vec::new();
for object_key in keys {
if object_key == anchor {
continue;
}
match parse_segment_key(object_key) {
Some((lane, segment)) => {
if coverage.get(&lane).is_some_and(|&c| segment <= c) {
skipped.push(object_key.clone());
} else if observed.get(&lane).is_some_and(|&o| segment <= o) {
before.push(object_key.clone());
} else {
after.push(object_key.clone());
}
}
None => before.push(object_key.clone()),
}
}
before.sort();
after.sort();
skipped.sort();
let mut order = before;
order.push(anchor.to_string());
order.extend(after);
RestorePlan { order, skipped }
}
fn newest_checkpoint(keys: &[String]) -> Option<String> {
keys.iter()
.filter_map(|key| parse_checkpoint_key(key).map(|n| (n, key)))
.max_by_key(|(n, _)| *n)
.map(|(_, key)| key.clone())
}
pub async fn import_vault_batch(
store: &Db,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
prefix: &str,
drive_pseudonym: Option<(&str, &str)>,
) -> AtomicResult<RestoreSummary> {
let mut summary = RestoreSummary::default();
let keys = vault.list(prefix)?;
let anchor = newest_checkpoint(&keys);
let (mut coverage, mut observed) = (BTreeMap::new(), BTreeMap::new());
let mut anchor_key = None;
if let Some(candidate) = anchor {
match vault
.get(&candidate)
.and_then(|sealed| envelope::open(key, &sealed))
{
Ok((_, plaintext)) => {
let pack = Pack::decode(&plaintext)?;
coverage = pack.coverage.clone();
observed = pack.observed.clone();
anchor_key = Some(candidate);
}
Err(err) => {
tracing::warn!(
"Cloud Vault restore: checkpoint {candidate} could not be opened ({err}); \
replaying every object instead"
);
}
}
}
let plan = plan_restore(&keys, anchor_key.as_deref(), &coverage, &observed);
summary.objects_skipped = plan.skipped.len();
for object_key in &plan.order {
let sealed = vault.get(object_key)?;
let Ok((header, plaintext)) = envelope::open(key, &sealed) else {
tracing::warn!("Cloud Vault restore: {object_key} could not be opened; skipping it");
summary.objects_unreadable += 1;
continue;
};
if !matches!(header.kind, ObjectKind::Pack | ObjectKind::Checkpoint) {
continue;
}
let pack = Pack::decode(&plaintext)?;
summary.packs_read += 1;
for entry in pack.entries {
let subject = Subject::from_raw(&entry.subject, None);
let mut resource = match store.get_resource(&subject).await {
Ok(existing) => existing,
Err(_) => Resource::new(entry.subject.clone()),
};
let doc = resource.build_state_doc()?;
doc.import_update(&entry.update)?;
resource.apply_state_doc(doc)?;
store
.add_resource_opts(&resource, false, true, true)
.await?;
crate::sync::tombstones::clear_tombstone(store, &entry.subject);
summary.resources_restored += 1;
}
for subject_str in pack.tombstones {
let subject = Subject::from_raw(&subject_str, None);
match store.remove_resource(&subject).await {
Ok(()) => {}
Err(_) => crate::sync::tombstones::record_tombstone(store, &subject_str),
}
summary.tombstones_applied += 1;
}
if let (Some((pseudonym, device)), Some((lane, segment))) =
(drive_pseudonym, parse_segment_key(object_key))
{
record_imported_lane(store, pseudonym, device, &lane, segment);
}
}
if summary.packs_read == 0 && summary.objects_unreadable > 0 {
return Err(format!(
"could not decrypt any of the {} vault objects for this drive — wrong drive key",
summary.objects_unreadable
)
.into());
}
Ok(summary)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::vault::store::MemoryVaultStore;
const PSEUDONYM: &str = "testpseudonym";
const DEVICE: &str = "testdevice";
fn key() -> DriveVaultKey {
DriveVaultKey::from_bytes([5u8; 32], 1)
}
#[test]
fn segment_keys_sort_in_replay_order() {
let mut keys = vec![
segment_key(PSEUDONYM, DEVICE, 10),
segment_key(PSEUDONYM, DEVICE, 2),
segment_key(PSEUDONYM, DEVICE, 1),
];
keys.sort();
assert_eq!(
keys,
vec![
segment_key(PSEUDONYM, DEVICE, 1),
segment_key(PSEUDONYM, DEVICE, 2),
segment_key(PSEUDONYM, DEVICE, 10),
],
"lexical order must match segment order, or replay applies packs backwards"
);
}
#[test]
fn segment_key_matches_the_published_format() {
let device = "0303030303030303030303030303030303030303030303030303030303030303";
let key = segment_key("testpseudonym", device, 1);
assert_eq!(
key,
"vault/testpseudonym/lanes/0303030303030303030303030303030303030303030303030303030303030303/seg-000001.pack",
"atomic-saas build_object_key must produce this byte-for-byte"
);
assert!(
key < segment_key("testpseudonym", device, 10),
"lexical order must match segment order"
);
}
#[test]
fn segment_keys_live_under_the_device_lane() {
let k = segment_key(PSEUDONYM, DEVICE, 1);
assert_eq!(k, "vault/testpseudonym/lanes/testdevice/seg-000001.pack");
assert!(k.starts_with(&lane_prefix(PSEUDONYM, DEVICE)));
}
#[test]
fn a_vault_with_no_checkpoint_replays_everything() {
let keys = vec![
segment_key(PSEUDONYM, DEVICE, 2),
segment_key(PSEUDONYM, DEVICE, 1),
];
let plan = plan_restore(&keys, None, &BTreeMap::new(), &BTreeMap::new());
assert_eq!(
plan.order,
vec![
segment_key(PSEUDONYM, DEVICE, 1),
segment_key(PSEUDONYM, DEVICE, 2),
]
);
assert!(plan.skipped.is_empty());
}
#[test]
fn a_checkpoint_orders_the_segments_around_itself() {
let other = "ff".repeat(32);
let anchor = checkpoint_key(PSEUDONYM, 7);
let keys = vec![
checkpoint_key(PSEUDONYM, 3),
anchor.clone(),
segment_key(PSEUDONYM, DEVICE, 1),
segment_key(PSEUDONYM, DEVICE, 4),
segment_key(PSEUDONYM, DEVICE, 9),
segment_key(PSEUDONYM, &other, 2),
];
let plan = plan_restore(
&keys,
Some(&anchor),
&BTreeMap::from([(DEVICE.to_string(), 1u32)]),
&BTreeMap::from([(DEVICE.to_string(), 4u32)]),
);
assert_eq!(
plan.skipped,
vec![segment_key(PSEUDONYM, DEVICE, 1)],
"a covered segment is not read"
);
assert_eq!(
plan.order,
vec![
checkpoint_key(PSEUDONYM, 3),
segment_key(PSEUDONYM, DEVICE, 4),
anchor.clone(),
segment_key(PSEUDONYM, &other, 2),
segment_key(PSEUDONYM, DEVICE, 9),
]
);
}
#[tokio::test]
async fn an_unopenable_checkpoint_falls_back_to_replaying_everything() {
let source = Db::init_temp("vault_bad_anchor_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
rename(&source, &drive, "renamed").await;
backup_delta(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let mangled = MemoryVaultStore::new();
for object_key in vault.list(&drive_prefix(PSEUDONYM)).unwrap() {
mangled
.put(&object_key, &vault.get(&object_key).unwrap())
.unwrap();
}
mangled
.put(&checkpoint_key(PSEUDONYM, 9), b"not a vault object")
.unwrap();
let target = Db::init_temp("vault_bad_anchor_target").await.unwrap();
let summary = restore(&target, &key, &mangled).await;
assert_eq!(
summary.objects_skipped, 0,
"an anchor that could not be opened must not decide what to skip"
);
assert_eq!(
summary.packs_read, 2,
"both real objects still replay: {summary:?}"
);
}
#[test]
fn sealed_packs_do_not_reveal_subjects() {
let pack = Pack::new(
vec![PackEntry {
subject: "did:ad:drive/secret-resource".into(),
update: vec![1, 2, 3],
}],
vec![],
);
let sealed = envelope::seal(&key(), ObjectKind::Pack, &pack.encode().unwrap()).unwrap();
let needle = b"secret-resource";
assert!(
!sealed.windows(needle.len()).any(|w| w == needle),
"subject leaked into the sealed pack"
);
}
const FOLDER: &str = "https://atomicdata.dev/classes/Folder";
async fn backup_with(
store: &Db,
drive: &Subject,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
device: &str,
policy: CheckpointPolicy,
confirm: bool,
) -> Option<BackupSummary> {
let keys = vault.list(&drive_prefix(PSEUDONYM)).unwrap();
let mut observed: BTreeMap<String, u32> = BTreeMap::new();
let mut highest_checkpoint = 0u64;
for object_key in &keys {
if let Some((lane, segment)) = parse_segment_key(object_key) {
let slot = observed.entry(lane).or_insert(0);
*slot = (*slot).max(segment);
}
if let Some(n) = parse_checkpoint_key(object_key) {
highest_checkpoint = highest_checkpoint.max(n);
}
}
let segment = observed.get(device).copied().unwrap_or(0) + 1;
let summary = export_vault_segment(
store,
drive,
key,
vault,
PSEUDONYM,
device,
segment,
highest_checkpoint + 1,
highest_checkpoint > 0,
&observed,
policy,
)
.await
.unwrap();
if confirm && summary.is_some() {
commit_lane_state(store, PSEUDONYM, device, segment).unwrap();
}
summary
}
async fn rename(store: &Db, subject_str: &str, name: &str) {
let subject = Subject::from_raw(subject_str, store.get_base_domain().as_deref());
let mut resource = store.get_resource(&subject).await.unwrap();
let doc = resource.build_state_doc().unwrap();
doc.set_property(crate::urls::NAME, &crate::Value::String(name.into()))
.unwrap();
doc.commit_with_message(&format!("rename to {name}"));
resource.apply_state_doc(doc).unwrap();
store
.add_resource_opts(&resource, false, true, true)
.await
.unwrap();
}
async fn backup(
store: &Db,
drive: &Subject,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
device: &str,
) -> Option<BackupSummary> {
backup_with(
store,
drive,
key,
vault,
device,
CheckpointPolicy::default(),
true,
)
.await
}
async fn backup_delta(
store: &Db,
drive: &Subject,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
device: &str,
) -> Option<BackupSummary> {
backup_with(
store,
drive,
key,
vault,
device,
CheckpointPolicy {
max_segments: u32::MAX,
bytes_ratio: f64::INFINITY,
},
true,
)
.await
}
async fn restore(
store: &Db,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
) -> RestoreSummary {
import_vault_batch(store, key, vault, &drive_prefix(PSEUDONYM), None)
.await
.unwrap()
}
async fn drive_contents(store: &Db, drive: &Subject) -> Vec<(String, String)> {
let mut out = Vec::new();
for subject_str in crate::sync::engine::collect_drive_subjects(store, drive).await {
let subject = Subject::from_raw(&subject_str, None);
if let Ok(resource) = store.get_resource(&subject).await {
let mut props: Vec<String> = resource
.get_propvals()
.iter()
.filter(|(k, _)| k.as_str() != crate::urls::LORO_UPDATE)
.map(|(k, v)| format!("{k}={v}"))
.collect();
props.sort();
out.push((subject_str, props.join("|")));
}
}
out.sort();
out
}
#[tokio::test]
async fn a_wiped_store_is_restored_from_the_vault() {
let source = Db::init_temp("vault_round_trip_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
for i in 0..3 {
source
.create_resource(FOLDER, &drive, &format!("note-{i}"), None)
.await
.unwrap();
}
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let before = drive_contents(&source, &drive_subject).await;
assert!(
before.len() >= 4,
"expected a drive plus children: {before:?}"
);
let key = key();
let vault = MemoryVaultStore::new();
let backup = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("a populated drive must produce a pack");
assert_eq!(backup.resources, before.len() + 1);
let restored = Db::init_temp("vault_round_trip_restored").await.unwrap();
let result = restore(&restored, &key, &vault).await;
assert_eq!(result.packs_read, 1);
assert_eq!(result.resources_restored, before.len() + 1);
let after = drive_contents(&restored, &drive_subject).await;
assert_eq!(after, before, "restored drive must match the original");
}
#[tokio::test]
async fn the_drive_owner_is_restored_with_the_drive() {
let source = Db::init_temp("vault_agent_source").await.unwrap();
let (agent, drive) = source.setup("alice").await.unwrap();
let agent_subject = agent.subject.clone();
let mut profile = source.get_resource(&agent_subject).await.unwrap();
profile
.set_unsafe(
crate::urls::NAME.into(),
crate::Value::String("Alice Returning".into()),
)
.unwrap();
source
.add_resource_opts(&profile, false, true, true)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("a populated drive must produce a pack");
let restored = Db::init_temp("vault_agent_restored").await.unwrap();
restore(&restored, &key, &vault).await;
let name = restored
.get_resource(&agent_subject)
.await
.expect("the drive's agent must come back with the drive")
.get(crate::urls::NAME)
.expect("with its profile")
.to_string();
assert_eq!(name, "Alice Returning");
}
#[tokio::test]
async fn a_deleted_resource_does_not_come_back_after_restore() {
let source = Db::init_temp("vault_delete_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let keep = source
.create_resource(FOLDER, &drive, "keep", None)
.await
.unwrap();
let doomed = source
.create_resource(FOLDER, &drive, "doomed", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("first pack");
commit_lane_state(&source, PSEUDONYM, DEVICE, 1).unwrap();
let doomed_subject = Subject::from_raw(&doomed, source.get_base_domain().as_deref());
source.remove_resource(&doomed_subject).await.unwrap();
let second = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("second pack carries the deletion");
assert_eq!(
second.tombstones, 1,
"the deletion must reach the vault, not just the local store"
);
let restored = Db::init_temp("vault_delete_restored").await.unwrap();
restore(&restored, &key, &vault).await;
let subjects = crate::sync::engine::collect_drive_subjects(&restored, &drive_subject).await;
assert!(
subjects.contains(&Subject::from_raw(&keep, None).pure_id()),
"the kept resource must survive: {subjects:?}"
);
assert!(
!subjects.contains(&doomed_subject.pure_id()),
"the deleted resource must not be resurrected: {subjects:?}"
);
}
#[tokio::test]
async fn restoring_a_locally_destroyed_subject_clears_its_tombstone() {
let source = Db::init_temp("vault_clear_tombstone_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let note = source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let restored = Db::init_temp("vault_clear_tombstone_restored")
.await
.unwrap();
crate::sync::tombstones::record_tombstone(&restored, ¬e);
assert!(crate::sync::tombstones::is_tombstoned(&restored, ¬e));
restore(&restored, &key, &vault).await;
assert!(
!crate::sync::tombstones::is_tombstoned(&restored, ¬e),
"a subject the backup re-created must not stay tombstoned"
);
}
#[tokio::test]
async fn an_uncommitted_export_does_not_advance_the_lane() {
let source = Db::init_temp("vault_uncommitted_export").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let doomed = source
.create_resource(FOLDER, &drive, "doomed", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup_with(
&source,
&drive_subject,
&key,
&vault,
DEVICE,
CheckpointPolicy::default(),
false,
)
.await
.unwrap();
source
.remove_resource(&Subject::from_raw(
&doomed,
source.get_base_domain().as_deref(),
))
.await
.unwrap();
let second = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("still exports the surviving resources");
assert_eq!(
second.tombstones, 0,
"nothing was ever backed up, so nothing can be reported deleted"
);
}
#[tokio::test]
async fn restore_spans_every_device_lane() {
let source = Db::init_temp("vault_multi_lane_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
let other_device = "ff".repeat(32);
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
source
.create_resource(FOLDER, &drive, "second-note", None)
.await
.unwrap();
backup(&source, &drive_subject, &key, &vault, &other_device)
.await
.unwrap();
let restored = Db::init_temp("vault_multi_lane_restored").await.unwrap();
let summary = restore(&restored, &key, &vault).await;
assert_eq!(
summary.packs_read, 2,
"the drive prefix must reach the anchor and the second lane"
);
let one_lane = Db::init_temp("vault_multi_lane_one").await.unwrap();
let partial = import_vault_batch(
&one_lane,
&key,
&vault,
&lane_prefix(PSEUDONYM, &other_device),
None,
)
.await
.unwrap();
assert_eq!(partial.packs_read, 1);
}
#[tokio::test]
async fn a_vanished_but_untombstoned_subject_is_not_reported_deleted() {
let source = Db::init_temp("vault_no_false_tombstone").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let note = source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
rename(&source, ¬e, "renamed").await;
let second = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("an edited drive still exports");
assert_eq!(second.tombstones, 0);
}
#[tokio::test]
async fn a_subject_that_vanishes_without_a_tombstone_stays_known() {
let source = Db::init_temp("vault_vanished_stays_known").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let doomed_str = source
.create_resource(FOLDER, &drive, "doomed", None)
.await
.unwrap();
let keeper = source
.create_resource(FOLDER, &drive, "keeper", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let doomed = Subject::from_raw(&doomed_str, source.get_base_domain().as_deref());
source
.kv
.remove(crate::db::trees::Tree::LoroSnapshots, doomed_str.as_bytes())
.unwrap();
source.remove_resource(&doomed).await.ok();
crate::sync::tombstones::clear_tombstone(&source, &doomed_str);
rename(&source, &keeper, "touched").await;
let second = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("the edit still ships");
assert_eq!(
second.tombstones, 0,
"no tombstone exists, so none may be claimed"
);
assert!(
read_lane_state(&source, PSEUDONYM, DEVICE)
.cursors
.contains_key(&doomed_str),
"the lane must still remember it backed this subject up, or it can \
never claim the deletion when one finally arrives"
);
}
#[tokio::test]
async fn an_unchanged_drive_costs_nothing() {
let source = Db::init_temp("vault_unchanged_drive").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
for i in 0..5 {
source
.create_resource(FOLDER, &drive, &format!("folder-{i}"), None)
.await
.unwrap();
}
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
let first = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("the first pass anchors the vault");
assert_eq!(first.unchanged, 0, "nothing was cached yet");
let objects_after_first = vault.list(&drive_prefix(PSEUDONYM)).unwrap().len();
for _ in 0..4 {
assert!(
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.is_none(),
"an unchanged drive must not produce an object"
);
}
assert_eq!(
vault.list(&drive_prefix(PSEUDONYM)).unwrap().len(),
objects_after_first,
"five passes over an unchanged drive must cost exactly one object"
);
}
#[tokio::test]
async fn one_edit_costs_one_edit() {
let source = Db::init_temp("vault_small_delta").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let mut subjects = Vec::new();
for i in 0..40 {
subjects.push(
source
.create_resource(FOLDER, &drive, &format!("folder-{i}"), None)
.await
.unwrap(),
);
}
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
let anchor = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("the first pass anchors the vault");
assert!(anchor.resources >= 40);
rename(&source, &subjects[0], "touched").await;
let delta = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("one edit must still be backed up");
assert_eq!(delta.kind, SegmentKind::Pack);
assert_eq!(
delta.resources, 1,
"only the edited resource belongs in the pack"
);
assert_eq!(
delta.unchanged,
anchor.resources - 1,
"every other resource must be recognised as unchanged from its version vector alone"
);
assert!(
delta.sealed_bytes * 4 < anchor.sealed_bytes,
"a one-resource delta ({} B) must be far smaller than a {}-resource anchor ({} B)",
delta.sealed_bytes,
anchor.resources,
anchor.sealed_bytes
);
}
#[tokio::test]
async fn a_delta_chain_restores_the_drive() {
let source = Db::init_temp("vault_delta_chain").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let note = source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let subject = Subject::from_raw(¬e, source.get_base_domain().as_deref());
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
for name in ["one", "two", "three"] {
rename(&source, ¬e, name).await;
let pass = backup_delta(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("each edit ships");
assert_eq!(pass.kind, SegmentKind::Pack);
}
let expected_versions =
crate::history::versions(&source.get_resource(&subject).await.unwrap())
.unwrap()
.len();
let target = Db::init_temp("vault_delta_chain_target").await.unwrap();
restore(&target, &key, &vault).await;
let restored = target.get_resource(&subject).await.unwrap();
assert_eq!(
restored.get(crate::urls::NAME).unwrap().to_string(),
"three",
"the chain must replay to the latest state"
);
assert_eq!(
crate::history::versions(&restored).unwrap().len(),
expected_versions,
"and carry every version the source had — a chain must lose no history \
a single self-sufficient segment would have kept"
);
}
#[tokio::test]
async fn a_broken_delta_chain_loses_the_edits_in_the_missing_link() {
let source = Db::init_temp("vault_broken_chain").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let note = source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let subject = Subject::from_raw(¬e, source.get_base_domain().as_deref());
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
rename(&source, ¬e, "one").await;
backup_delta(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
rename(&source, ¬e, "two").await;
backup_delta(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let mutilated = MemoryVaultStore::new();
let dropped = segment_key(PSEUDONYM, DEVICE, 1);
for object_key in vault.list(&drive_prefix(PSEUDONYM)).unwrap() {
if object_key != dropped {
mutilated
.put(&object_key, &vault.get(&object_key).unwrap())
.unwrap();
}
}
let target = Db::init_temp("vault_broken_chain_target").await.unwrap();
restore(&target, &key, &mutilated).await;
let restored = target.get_resource(&subject).await.unwrap();
assert_ne!(
restored.get(crate::urls::NAME).unwrap().to_string(),
"two",
"a delta whose predecessor is gone cannot be applied — if this ever \
passes, deltas became self-sufficient and the pruning rule can relax"
);
}
#[tokio::test]
async fn importing_the_same_pack_twice_is_idempotent() {
let source = Db::init_temp("vault_idempotent_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let before = drive_contents(&source, &drive_subject).await;
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let restored = Db::init_temp("vault_idempotent_restored").await.unwrap();
restore(&restored, &key, &vault).await;
let once = drive_contents(&restored, &drive_subject).await;
restore(&restored, &key, &vault).await;
let twice = drive_contents(&restored, &drive_subject).await;
assert_eq!(once, before);
assert_eq!(twice, once, "a second import must not change the result");
}
#[tokio::test]
async fn a_restore_without_the_right_key_fails() {
let source = Db::init_temp("vault_wrong_key_source").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
source
.create_resource(FOLDER, &drive, "note", None)
.await
.unwrap();
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key(), &vault, DEVICE)
.await
.unwrap();
let restored = Db::init_temp("vault_wrong_key_restored").await.unwrap();
let attacker = DriveVaultKey::from_bytes([0xFF; 32], 1);
let err = import_vault_batch(&restored, &attacker, &vault, &drive_prefix(PSEUDONYM), None)
.await
.unwrap_err()
.to_string();
assert!(err.contains("decrypt"), "{err}");
}
#[tokio::test]
async fn an_empty_drive_produces_no_object() {
let store = Db::init_temp("vault_empty_drive").await.unwrap();
let unknown = Subject::from_raw("did:ad:nonexistentdrive", None);
let vault = MemoryVaultStore::new();
let out = backup(&store, &unknown, &key(), &vault, DEVICE).await;
assert!(out.is_none(), "nothing to back up should mean no object");
assert!(vault.is_empty());
}
#[test]
fn a_pack_round_trips_through_a_store() {
let vault = MemoryVaultStore::new();
let pack = Pack::new(
vec![PackEntry {
subject: "s".into(),
update: vec![9, 9, 9],
}],
vec!["gone".into()],
);
let sealed = envelope::seal(&key(), ObjectKind::Pack, &pack.encode().unwrap()).unwrap();
let object_key = segment_key(PSEUDONYM, DEVICE, 1);
vault.put(&object_key, &sealed).unwrap();
let fetched = vault.get(&object_key).unwrap();
let (_, plaintext) = envelope::open(&key(), &fetched).unwrap();
assert_eq!(Pack::decode(&plaintext).unwrap(), pack);
}
#[tokio::test]
async fn a_restore_keeps_edit_history() {
let source = Db::init_temp("vault_history_roundtrip").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let subject_str = source
.create_resource(FOLDER, &drive, "history", None)
.await
.unwrap();
let subject = Subject::from_raw(&subject_str, source.get_base_domain().as_deref());
for name in ["one", "two", "three"] {
let mut resource = source.get_resource(&subject).await.unwrap();
let doc = resource.build_state_doc().unwrap();
doc.set_property(crate::urls::NAME, &crate::Value::String(name.into()))
.unwrap();
doc.commit_with_message(&format!("rename to {name}"));
resource.apply_state_doc(doc).unwrap();
source
.add_resource_opts(&resource, false, true, true)
.await
.unwrap();
}
let before = crate::history::versions(&source.get_resource(&subject).await.unwrap())
.unwrap()
.len();
assert!(
before > 1,
"the fixture must produce a history worth preserving, got {before}"
);
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.unwrap();
let target = Db::init_temp("vault_history_roundtrip_target")
.await
.unwrap();
restore(&target, &key, &vault).await;
let restored = target.get_resource(&subject).await.unwrap();
let after = crate::history::versions(&restored).unwrap().len();
assert_eq!(
after, before,
"a restore must bring back every version, not just the latest state"
);
}
#[tokio::test]
async fn the_newest_checkpoint_alone_restores_the_drive() {
let source = Db::init_temp("vault_latest_checkpoint").await.unwrap();
let (_agent, drive) = source.setup("alice").await.unwrap();
let kept_str = source
.create_resource(FOLDER, &drive, "kept", None)
.await
.unwrap();
let kept = Subject::from_raw(&kept_str, source.get_base_domain().as_deref());
for name in ["one", "two", "three"] {
let mut resource = source.get_resource(&kept).await.unwrap();
let doc = resource.build_state_doc().unwrap();
doc.set_property(crate::urls::NAME, &crate::Value::String(name.into()))
.unwrap();
doc.commit_with_message(&format!("rename to {name}"));
resource.apply_state_doc(doc).unwrap();
source
.add_resource_opts(&resource, false, true, true)
.await
.unwrap();
}
let doomed_str = source
.create_resource(FOLDER, &drive, "doomed", None)
.await
.unwrap();
let doomed = Subject::from_raw(&doomed_str, source.get_base_domain().as_deref());
let drive_subject = Subject::from_raw(&drive, source.get_base_domain().as_deref());
let key = key();
let vault = MemoryVaultStore::new();
let first = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("a populated drive must produce an object");
assert_eq!(first.kind, SegmentKind::Checkpoint);
source.remove_resource(&doomed).await.unwrap();
let second = backup(&source, &drive_subject, &key, &vault, DEVICE)
.await
.expect("a deletion must be shipped");
assert_eq!(second.kind, SegmentKind::Pack);
assert_eq!(second.tombstones, 1);
source
.create_resource(FOLDER, &drive, "after", None)
.await
.unwrap();
let third = backup_with(
&source,
&drive_subject,
&key,
&vault,
DEVICE,
CheckpointPolicy {
max_segments: 1,
bytes_ratio: 1.0,
},
true,
)
.await
.expect("the policy asked for a fresh anchor");
assert_eq!(third.kind, SegmentKind::Checkpoint);
assert_eq!(
third.coverage.get(DEVICE),
Some(&1),
"a checkpoint consumes no segment number, so an anchor taken when \
this lane's next segment is 2 subsumes everything through 1"
);
let newest = checkpoint_key(PSEUDONYM, 2);
let latest_only = MemoryVaultStore::new();
latest_only
.put(&newest, &vault.get(&newest).unwrap())
.unwrap();
let target = Db::init_temp("vault_latest_checkpoint_target")
.await
.unwrap();
restore(&target, &key, &latest_only).await;
let restored = target.get_resource(&kept).await.unwrap();
assert_eq!(
crate::history::versions(&restored).unwrap().len(),
5,
"the newest checkpoint must carry the whole version history on its own"
);
assert_eq!(
restored.get(crate::urls::NAME).unwrap().to_string(),
"three",
"and the current state"
);
assert!(
target.get_resource(&doomed).await.is_err(),
"a resource deleted before this checkpoint must stay deleted — recovering it is what an OLDER object would be kept for"
);
}
}