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::resources::Resource;
use crate::storelike::Storelike;
use crate::Subject;
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()
}
fn read_lane_state(store: &Db, drive_pseudonym: &str, device_pubkey: &str) -> Vec<String> {
store
.kv
.get(
crate::db::trees::Tree::PluginMeta,
&lane_state_key(drive_pseudonym, device_pubkey),
)
.ok()
.flatten()
.and_then(|bytes| serde_json::from_slice(&bytes).ok())
.unwrap_or_default()
}
fn write_lane_state(store: &Db, drive_pseudonym: &str, device_pubkey: &str, subjects: &[String]) {
if let Ok(bytes) = serde_json::to_vec(subjects) {
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 subjects: Vec<String> =
serde_json::from_slice(&bytes).map_err(|e| format!("malformed pending lane state: {e}"))?;
write_lane_state(store, drive_pseudonym, device_pubkey, &subjects);
let _ = store
.kv
.remove(crate::db::trees::Tree::PluginMeta, &pending_key);
Ok(())
}
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)
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BackupSummary {
pub object_key: String,
pub resources: usize,
pub tombstones: usize,
pub sealed_bytes: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct RestoreSummary {
pub packs_read: usize,
pub resources_restored: usize,
pub tombstones_applied: usize,
}
pub async fn export_vault_delta(
store: &Db,
drive: &Subject,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
drive_pseudonym: &str,
device_pubkey: &str,
segment: u32,
) -> AtomicResult<Option<BackupSummary>> {
let subjects = crate::sync::engine::collect_drive_subjects(store, drive).await;
let mut entries = Vec::new();
for subject_str in &subjects {
let subject = Subject::from_raw(subject_str, None);
let Ok(resource) = store.get_resource(&subject).await else {
continue;
};
let doc = resource.build_state_doc()?;
let update = doc.export_updates_since(&Default::default());
if update.is_empty() {
continue;
}
entries.push(PackEntry {
subject: subject_str.clone(),
update,
});
}
entries.sort_by(|a, b| a.subject.cmp(&b.subject));
let present: std::collections::HashSet<&String> = entries.iter().map(|e| &e.subject).collect();
let tombstones: Vec<String> = read_lane_state(store, drive_pseudonym, device_pubkey)
.into_iter()
.filter(|subject| {
!present.contains(subject) && crate::sync::tombstones::is_tombstoned(store, subject)
})
.collect();
let exported: Vec<String> = entries.iter().map(|e| e.subject.clone()).collect();
let pack = Pack::new(entries, tombstones);
if pack.is_empty() {
return Ok(None);
}
let resources = pack.entries.len();
let tombstones = pack.tombstones.len();
let sealed = envelope::seal(key, ObjectKind::Pack, &pack.encode()?)?;
let object_key = segment_key(drive_pseudonym, device_pubkey, segment);
vault.put(&object_key, &sealed)?;
if let Ok(bytes) = serde_json::to_vec(&exported) {
let _ = store.kv.insert(
crate::db::trees::Tree::PluginMeta,
&pending_lane_state_key(drive_pseudonym, device_pubkey, segment),
&bytes,
);
}
Ok(Some(BackupSummary {
object_key,
resources,
tombstones,
sealed_bytes: sealed.len(),
}))
}
pub async fn import_vault_batch(
store: &Db,
key: &DriveVaultKey,
vault: &dyn VaultObjectStore,
prefix: &str,
) -> AtomicResult<RestoreSummary> {
let mut summary = RestoreSummary::default();
for object_key in vault.list(prefix)? {
let sealed = vault.get(&object_key)?;
let (header, plaintext) = envelope::open(key, &sealed)?;
if header.kind != ObjectKind::Pack {
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;
}
}
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 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 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 =
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.expect("a populated drive must produce a pack");
assert_eq!(backup.resources, before.len());
let restored = Db::init_temp("vault_round_trip_restored").await.unwrap();
let result = import_vault_batch(&restored, &key, &vault, &lane_prefix(PSEUDONYM, DEVICE))
.await
.unwrap();
assert_eq!(result.packs_read, 1);
assert_eq!(result.resources_restored, before.len());
let after = drive_contents(&restored, &drive_subject).await;
assert_eq!(after, before, "restored drive must match the original");
}
#[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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.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 =
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 2)
.await
.unwrap()
.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();
import_vault_batch(&restored, &key, &vault, &lane_prefix(PSEUDONYM, DEVICE))
.await
.unwrap();
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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.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));
import_vault_batch(&restored, &key, &vault, &lane_prefix(PSEUDONYM, DEVICE))
.await
.unwrap();
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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
source
.remove_resource(&Subject::from_raw(
&doomed,
source.get_base_domain().as_deref(),
))
.await
.unwrap();
let second =
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 2)
.await
.unwrap()
.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);
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
export_vault_delta(
&source,
&drive_subject,
&key,
&vault,
PSEUDONYM,
&other_device,
1,
)
.await
.unwrap()
.unwrap();
let restored = Db::init_temp("vault_multi_lane_restored").await.unwrap();
let summary = import_vault_batch(&restored, &key, &vault, &drive_prefix(PSEUDONYM))
.await
.unwrap();
assert_eq!(
summary.packs_read, 2,
"the drive prefix must reach both lanes"
);
let one_lane = Db::init_temp("vault_multi_lane_one").await.unwrap();
let partial = import_vault_batch(&one_lane, &key, &vault, &lane_prefix(PSEUDONYM, DEVICE))
.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();
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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
commit_lane_state(&source, PSEUDONYM, DEVICE, 1).unwrap();
let second =
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 2)
.await
.unwrap()
.expect("unchanged drive still exports its oplog in Phase 1");
assert_eq!(second.tombstones, 0);
}
#[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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
let restored = Db::init_temp("vault_idempotent_restored").await.unwrap();
let prefix = lane_prefix(PSEUDONYM, DEVICE);
import_vault_batch(&restored, &key, &vault, &prefix)
.await
.unwrap();
let once = drive_contents(&restored, &drive_subject).await;
import_vault_batch(&restored, &key, &vault, &prefix)
.await
.unwrap();
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();
export_vault_delta(
&source,
&drive_subject,
&key(),
&vault,
PSEUDONYM,
DEVICE,
1,
)
.await
.unwrap()
.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,
&lane_prefix(PSEUDONYM, DEVICE),
)
.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 = export_vault_delta(&store, &unknown, &key(), &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap();
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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
let target = Db::init_temp("vault_history_roundtrip_target")
.await
.unwrap();
import_vault_batch(&target, &key, &vault, &drive_prefix(PSEUDONYM))
.await
.unwrap();
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_segment_alone_restores_the_drive() {
let source = Db::init_temp("vault_latest_segment").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();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 1)
.await
.unwrap()
.unwrap();
commit_lane_state(&source, PSEUDONYM, DEVICE, 1).unwrap();
source.remove_resource(&doomed).await.unwrap();
export_vault_delta(&source, &drive_subject, &key, &vault, PSEUDONYM, DEVICE, 2)
.await
.unwrap()
.unwrap();
commit_lane_state(&source, PSEUDONYM, DEVICE, 2).unwrap();
let keys = vault.list(&drive_prefix(PSEUDONYM)).unwrap();
assert_eq!(keys.len(), 2, "the fixture needs two segments");
let newest = keys.iter().max().unwrap().clone();
let latest_only = MemoryVaultStore::new();
latest_only
.put(&newest, &vault.get(&newest).unwrap())
.unwrap();
let target = Db::init_temp("vault_latest_segment_target").await.unwrap();
import_vault_batch(&target, &key, &latest_only, &drive_prefix(PSEUDONYM))
.await
.unwrap();
let restored = target.get_resource(&kept).await.unwrap();
assert_eq!(
crate::history::versions(&restored).unwrap().len(),
5,
"the newest segment 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 segment must stay deleted — recovering it is what an OLDER segment would be kept for"
);
}
}