use std::sync::atomic::{AtomicBool, Ordering};
use crate::{
commit::{Commit, CommitOpts},
db::Db,
errors::AtomicResult,
parse::parse_json_ad_commit_resource,
Storelike,
};
static IMPORTING: AtomicBool = AtomicBool::new(false);
pub fn is_importing() -> bool {
IMPORTING.load(Ordering::Relaxed)
}
pub(crate) fn set_importing(v: bool) {
IMPORTING.store(v, Ordering::Relaxed);
}
pub async fn apply_state_update(store: &Db, subject: &str, state_bytes: &[u8]) -> AtomicResult<()> {
set_importing(true);
let result = async {
if let Some(resolved) = resolve_update(store, subject, state_bytes).await {
persist_update(store, subject, resolved).await?;
}
Ok(())
}
.await;
set_importing(false);
result
}
pub struct ResolvedUpdate {
snapshot: Vec<u8>,
resource: crate::Resource,
pub drive_subject: String,
}
pub async fn resolve_update(
store: &Db,
subject: &str,
state_bytes: &[u8],
) -> Option<ResolvedUpdate> {
if state_bytes.is_empty() {
return None;
}
let snapshot_key =
crate::Subject::from_raw(subject, store.get_base_domain().as_deref()).pure_id();
let doc = if let Ok(Some(existing)) = store.kv.get(
crate::db::trees::Tree::LoroSnapshots,
snapshot_key.as_bytes(),
) {
match crate::loro::AtomicLoroDoc::from_snapshot(&existing) {
Ok(d) => {
if let Err(e) = d.import_update(state_bytes) {
tracing::warn!(
"[ws_apply] import_update failed for {}: {e}",
&subject[..subject.len().min(20)]
);
}
d
}
Err(_) => crate::loro::AtomicLoroDoc::from_snapshot(state_bytes).ok()?,
}
} else {
match crate::loro::AtomicLoroDoc::from_snapshot(state_bytes) {
Ok(d) => d,
Err(_) => {
let d = crate::loro::AtomicLoroDoc::new();
if d.import_update(state_bytes).is_err() {
return None;
}
d
}
}
};
let snapshot = doc.export_snapshot();
let subj = crate::Subject::from_raw(subject, store.get_base_domain().as_deref());
let existing = store.get_resource(&subj).await.ok();
let drive_subject = if let Some(existing) = &existing {
Some(
existing
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| existing.get_subject().to_string()),
)
} else {
None
};
let mut resource = existing.unwrap_or_else(|| crate::Resource::new(subject.to_string()));
if resource.apply_state_doc(doc).is_err() {
return None;
}
let drive_subject = match drive_subject {
Some(d) => d,
None => {
let mut resolved = resource.get_subject().to_string();
if let Ok(parent_val) = resource.get(crate::urls::PARENT) {
let parent_subject = crate::Subject::from(parent_val.to_string());
if let Ok(parent_res) = store.get_resource(&parent_subject).await {
resolved = parent_res
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| parent_subject.to_string());
}
}
resolved
}
};
Some(ResolvedUpdate {
snapshot,
resource,
drive_subject,
})
}
pub async fn persist_update(
store: &Db,
subject: &str,
resolved: ResolvedUpdate,
) -> AtomicResult<()> {
let snapshot_key =
crate::Subject::from_raw(subject, store.get_base_domain().as_deref()).pure_id();
let _subject_guard = store.subject_locks.lock(&snapshot_key).await;
let doc = match store.kv.get(
crate::db::trees::Tree::LoroSnapshots,
snapshot_key.as_bytes(),
) {
Ok(Some(current)) => match crate::loro::AtomicLoroDoc::from_snapshot(¤t) {
Ok(doc) => {
if let Err(e) = doc.import_update(&resolved.snapshot) {
tracing::warn!("[ws_apply] re-merge failed for {snapshot_key}: {e}");
}
Some(doc)
}
Err(_) => None,
},
_ => None,
};
let mut resource = resolved.resource;
let snapshot = match doc {
Some(doc) => {
let snapshot = doc.export_snapshot();
let _ = resource.apply_state_doc(doc);
snapshot
}
None => resolved.snapshot,
};
let _ = store.kv.insert(
crate::db::trees::Tree::LoroSnapshots,
snapshot_key.as_bytes(),
&snapshot,
);
let _ = store.add_resource_opts(&resource, false, true, true).await;
Ok(())
}
pub async fn apply_destroy(store: &Db, subject: &str) -> AtomicResult<()> {
if subject.is_empty() {
return Ok(());
}
set_importing(true);
let result = apply_destroy_unchecked(store, subject).await;
set_importing(false);
result
}
async fn apply_destroy_unchecked(store: &Db, subject: &str) -> AtomicResult<()> {
let subj = crate::Subject::from_raw(subject, store.get_base_domain().as_deref());
let existed = store.get_resource(&subj).await.is_ok();
if !existed && !crate::sync::tombstones::is_tombstoned(store, subject) {
tracing::warn!(
"[ws_apply] ignoring DESTROY for locally-unknown subject {} (F10: not recording a phantom tombstone)",
&subject[..subject.len().min(20)]
);
return Ok(());
}
let _ = store.remove_resource(&subj).await;
crate::sync::tombstones::record_tombstone(store, subject);
tracing::info!("[ws_apply] deleted {}", &subject[..subject.len().min(20)]);
Ok(())
}
pub async fn resolve_destroy_drive(store: &Db, subject: &str) -> Option<String> {
let subj = crate::Subject::from_raw(subject, store.get_base_domain().as_deref());
let resource = store.get_resource(&subj).await.ok()?;
Some(
resource
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| resource.get_subject().to_string()),
)
}
pub async fn apply_destroy_checked(store: &Db, subject: &str) -> AtomicResult<()> {
if subject.is_empty() {
return Ok(());
}
set_importing(true);
let result = apply_destroy_unchecked(store, subject).await;
set_importing(false);
result
}
pub async fn apply_commit_json(store: &Db, body: &str) -> AtomicResult<()> {
set_importing(true);
let result = async {
let resource = parse_json_ad_commit_resource(body, store).await?;
let commit = Commit::from_resource(resource)?;
let opts = CommitOpts {
validate_signature: true,
validate_timestamp: false,
validate_previous_commit: false,
validate_rights: false,
update_index: true,
..CommitOpts::no_validations_no_index()
};
store.apply_commit(commit, &opts).await?;
Ok::<(), crate::AtomicError>(())
}
.await;
set_importing(false);
result
}
#[cfg(test)]
mod resolve_update_drive_spoof_tests {
use super::*;
use crate::loro::AtomicLoroDoc;
use crate::values::Value;
#[tokio::test]
async fn existing_resource_ignores_spoofed_drive_in_payload() {
let db = Db::init_temp("resolve_update_spoof_existing")
.await
.unwrap();
let (_alice, real_drive) = db.setup("Alice").await.unwrap();
let doc_subject = db
.create_resource(
"https://atomicdata.dev/classes/Folder",
&real_drive,
"Alice's doc",
None,
)
.await
.unwrap();
let stored = db.get_resource(&doc_subject.as_str().into()).await.unwrap();
assert_eq!(
stored.get(crate::urls::DRIVE_PROP).unwrap().to_string(),
real_drive
);
let snapshot_key =
crate::Subject::from_raw(&doc_subject, db.get_base_domain().as_deref()).pure_id();
let real_snapshot = db
.kv
.get(
crate::db::trees::Tree::LoroSnapshots,
snapshot_key.as_bytes(),
)
.unwrap()
.expect("resource should have a stored Loro snapshot");
let spoofed_drive = "https://attacker.example/not-your-drive";
let malicious = AtomicLoroDoc::from_snapshot(&real_snapshot).unwrap();
malicious
.set_property(
crate::urls::DRIVE_PROP,
&Value::AtomicUrl(spoofed_drive.to_string().into()),
)
.unwrap();
let malicious_bytes = malicious.export_snapshot();
let resolved = resolve_update(&db, &doc_subject, &malicious_bytes)
.await
.expect("a well-formed snapshot should still resolve");
assert_eq!(
resolved.drive_subject, real_drive,
"the existing resource's real drive must win over a spoofed payload assertion"
);
assert_ne!(resolved.drive_subject, spoofed_drive);
}
#[tokio::test]
async fn new_subject_with_no_resolvable_parent_falls_back_to_own_subject() {
let db = Db::init_temp("resolve_update_new_subject_no_parent")
.await
.unwrap();
let _ = db.setup("Alice").await.unwrap();
let new_subject = "https://example.test/brand-new-resource";
let spoofed_drive = "https://attacker.example/not-your-drive";
let malicious = AtomicLoroDoc::new();
malicious
.set_property(
crate::urls::DRIVE_PROP,
&Value::AtomicUrl(spoofed_drive.to_string().into()),
)
.unwrap();
let malicious_bytes = malicious.export_snapshot();
let resolved = resolve_update(&db, new_subject, &malicious_bytes)
.await
.expect("a well-formed snapshot should still resolve");
assert_eq!(
resolved.drive_subject, new_subject,
"no parent to borrow a drive from — must fall back to its own subject, not the payload's claimed drive"
);
assert_ne!(resolved.drive_subject, spoofed_drive);
}
}
#[cfg(test)]
mod destroy_phantom_tombstone_tests {
use super::*;
use crate::sync::tombstones;
#[tokio::test]
async fn destroy_of_unknown_subject_does_not_record_tombstone() {
let db = Db::init_temp("ws_apply_f10_unknown_subject").await.unwrap();
let _ = db.setup("Alice").await.unwrap();
let unknown_subject = "https://example.test/never-existed-here";
assert!(!tombstones::is_tombstoned(&db, unknown_subject));
apply_destroy(&db, unknown_subject).await.unwrap();
assert!(
!tombstones::is_tombstoned(&db, unknown_subject),
"F10: DESTROY for a locally-unknown subject must not record a phantom tombstone"
);
}
#[tokio::test]
async fn destroy_of_known_subject_still_records_tombstone() {
let db = Db::init_temp("ws_apply_f10_known_subject").await.unwrap();
let (_alice, drive) = db.setup("Alice").await.unwrap();
let subject = db
.create_resource(
"https://atomicdata.dev/classes/Folder",
&drive,
"Alice's doc",
None,
)
.await
.unwrap();
apply_destroy(&db, &subject).await.unwrap();
assert!(
tombstones::is_tombstoned(&db, &subject),
"a DESTROY for a subject we actually knew about must still record a tombstone"
);
}
#[tokio::test]
async fn destroy_of_already_tombstoned_subject_stays_tombstoned() {
let db = Db::init_temp("ws_apply_f10_already_tombstoned")
.await
.unwrap();
let _ = db.setup("Alice").await.unwrap();
let subject = "https://example.test/already-gone";
tombstones::record_tombstone(&db, subject);
assert!(tombstones::is_tombstoned(&db, subject));
apply_destroy(&db, subject).await.unwrap();
assert!(tombstones::is_tombstoned(&db, subject));
}
}