use crate::db::trees::Tree;
use crate::loro::AtomicLoroDoc;
use crate::{Db, Storelike};
use super::protocol;
#[derive(Debug, Clone, Copy)]
pub enum AuthBinding<'a> {
Unbound,
Origin(&'a str),
Origins(&'a [&'a str]),
}
#[derive(Debug, Clone, Copy)]
pub enum AuthChallenge<'a> {
None,
Issued(&'a str),
Required(&'a str),
}
fn url_origin(url: &str) -> Option<String> {
let parsed = url::Url::parse(url).ok()?;
if !matches!(parsed.scheme(), "http" | "https") {
return None;
}
let host = parsed.host_str()?.to_ascii_lowercase();
Some(match parsed.port() {
Some(p) => format!("{}://{}:{}", parsed.scheme(), host, p),
None => format!("{}://{}", parsed.scheme(), host),
})
}
pub async fn handle_auth_frame(
payload: &[u8],
store: &Db,
agent: &mut crate::agents::ForAgent,
binding: AuthBinding<'_>,
challenge: AuthChallenge<'_>,
) -> Vec<Vec<u8>> {
let refuse = |msg: String| {
vec![protocol::encode_error(
0,
protocol::error_code::AUTH_FAILED,
&msg,
)]
};
let Ok(json) = std::str::from_utf8(payload) else {
return refuse("Invalid UTF-8 in auth".into());
};
let auth = match serde_json::from_str::<crate::authentication::AuthValues>(json) {
Ok(a) => a,
Err(e) => return refuse(format!("Invalid auth JSON: {e}")),
};
let (_, carried_nonce) = protocol::split_challenge_fragment(&auth.requested_subject);
match (challenge, carried_nonce) {
(AuthChallenge::None, _) => {}
(AuthChallenge::Issued(_), None) => {}
(AuthChallenge::Required(_), None) => {
return refuse(
"Auth failed: this server requires the CHALLENGE nonce in requestedSubject".into(),
);
}
(AuthChallenge::Issued(issued) | AuthChallenge::Required(issued), Some(carried)) => {
if carried != issued {
return refuse(
"Auth failed: requestedSubject nonce does not answer this connection's CHALLENGE"
.into(),
);
}
}
}
let expected: Vec<&str> = match binding {
AuthBinding::Unbound => Vec::new(),
AuthBinding::Origin(one) => vec![one],
AuthBinding::Origins(many) => many.to_vec(),
};
let expected_origins: Vec<String> = expected.iter().filter_map(|o| url_origin(o)).collect();
if !expected_origins.is_empty() {
let signed_origin = url_origin(&auth.requested_subject);
let named_this_server = signed_origin
.as_deref()
.is_some_and(|signed| expected_origins.iter().any(|e| e == signed));
if !named_this_server {
return refuse(format!(
"Auth failed: requestedSubject {} does not name this server ({})",
auth.requested_subject,
expected_origins.join(" or ")
));
}
}
match crate::authentication::get_agent_from_auth_values_and_check(Some(auth), store).await {
Ok(a) => {
*agent = a;
vec![protocol::encode_auth_ok()]
}
Err(e) => refuse(format!("Auth failed: {e}")),
}
}
#[derive(Debug, Default, Clone)]
pub struct HandleOutput {
pub frames: Vec<Vec<u8>>,
pub subscribe: Option<String>,
pub unsubscribe: Option<String>,
}
pub async fn handle_frame(
frame: &[u8],
store: &Db,
agent: &mut crate::agents::ForAgent,
) -> Vec<Vec<u8>> {
handle_frame_full(frame, store, agent).await.frames
}
pub async fn handle_frame_full(
frame: &[u8],
store: &Db,
agent: &mut crate::agents::ForAgent,
) -> HandleOutput {
if frame.is_empty() {
return HandleOutput::default();
}
let tag = frame[0];
let payload = &frame[1..];
match tag {
protocol::tag::SUB => return handle_sub(payload, store, agent).await,
protocol::tag::UNSUB => return handle_unsub(payload),
_ => {}
}
let frames = match tag {
protocol::tag::AUTH => {
handle_auth_frame(
payload,
store,
agent,
AuthBinding::Unbound,
AuthChallenge::None,
)
.await
}
protocol::tag::GET => {
if let Some(decoded) = protocol::decode_get(payload) {
let subject =
crate::Subject::from_raw(decoded.subject, store.get_base_domain().as_deref());
match store.get_resource_extended(&subject, false, agent).await {
Ok(r) => {
let resource = r.to_single();
let snapshot = resource.materialized_state().unwrap_or_else(|| {
resource
.build_state_doc()
.map(|doc| doc.export_snapshot())
.unwrap_or_default()
});
if snapshot.is_empty() {
vec![protocol::encode_error(
decoded.request_id,
protocol::error_code::UNKNOWN,
"No state",
)]
} else {
let origin = store
.get_base_domain()
.unwrap_or_else(|| "http://localhost".to_string());
let subject_resolved = resource.get_subject().resolve(&origin);
let last_commit = resource
.get(crate::urls::LAST_COMMIT)
.ok()
.map(|v| v.to_string())
.filter(|s| !s.is_empty());
let mut flags = protocol::flags::SNAPSHOT;
if last_commit.is_some() {
flags |= protocol::flags::HAS_COMMIT_ID;
}
vec![protocol::encode_update(
flags,
decoded.request_id,
&subject_resolved,
last_commit.as_deref(),
&snapshot,
)]
}
}
Err(e) => {
vec![protocol::encode_error(
decoded.request_id,
protocol::error_code::UNKNOWN,
&e.to_string(),
)]
}
}
} else {
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid GET frame",
)]
}
}
protocol::tag::COMMIT => {
match protocol::decode_commit(payload) {
Some(decoded) => {
let request_id = decoded.request_id;
match apply_peer_commit(store, decoded.commit_json).await {
Ok(commit_json) => {
vec![protocol::encode_commit_ok(request_id, &commit_json)]
}
Err(e) => {
let msg = e.to_string();
vec![protocol::encode_error(
request_id,
protocol::classify_commit_error(&msg),
&msg,
)]
}
}
}
None => vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid COMMIT frame",
)],
}
}
protocol::tag::SYNC => match protocol::decode_sync(payload) {
Some(sync) if sync.probe => {
match drive_sync_hash_for(store, &sync.drive, agent).await {
Ok(server_hash) if server_hash == sync.drive_hash => {
vec![protocol::encode_sync_ok(&sync.drive)]
}
Ok(_) => vec![protocol::encode_sync_resend(&sync.drive)],
Err(reason) => vec![protocol::encode_error(
0,
protocol::error_code::UNAUTHORIZED_READ,
&format!("SYNC refused for {}: {reason}", sync.drive),
)],
}
}
Some(sync) => {
let filter = sync
.subjects
.as_ref()
.map(|s| s.iter().cloned().collect::<std::collections::HashSet<_>>());
handle_sync_vv_filtered(
&sync.drive,
&sync.drive_hash,
&sync.peers,
&sync.resources,
filter.as_ref(),
store,
agent,
)
.await
}
None => vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid SYNC frame",
)],
},
protocol::tag::SYNC_PUSH => {
if let Some(push) = protocol::decode_sync_push(payload) {
match import_sync_push(&push, store, agent, false).await {
Ok((_count, mut blob_requests)) => {
let mut responses = vec![protocol::encode_sync_ok(&push.drive)];
responses.append(&mut blob_requests);
responses
}
Err(rejected) => vec![rejected.to_error_frame()],
}
} else {
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid SYNC_PUSH frame",
)]
}
}
protocol::tag::BLOB_REQUEST => {
if let Some(hash) = protocol::decode_blob_request(payload) {
match store.kv.get(Tree::Blobs, &hash) {
Ok(Some(bytes)) => vec![protocol::encode_blob_response(&hash, &bytes)],
_ => vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Blob not found",
)],
}
} else {
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid BLOB_REQUEST frame",
)]
}
}
protocol::tag::BLOB_RESPONSE => {
if let Some(resp) = protocol::decode_blob_response(payload) {
match store.take_pending_blob_request(&resp.hash) {
Some(_drive) if blake3::hash(&resp.bytes).as_bytes() != &resp.hash => {
tracing::warn!(
"BLOB_RESPONSE: bytes do not hash to {}, dropped",
hex::encode(resp.hash)
);
vec![]
}
Some(drive) if store.sync_policy().admit_drive_write(&drive) => {
let _ = store.kv.insert(Tree::Blobs, &resp.hash, &resp.bytes);
vec![]
}
Some(drive) => {
tracing::warn!(
"BLOB_RESPONSE: drive {} not admitted by sync policy, dropping blob",
drive
);
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Drive not admitted for sync",
)]
}
None => {
tracing::warn!(
"BLOB_RESPONSE: no matching pending BLOB_REQUEST, dropping blob"
);
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Unsolicited blob response",
)]
}
}
} else {
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid BLOB_RESPONSE frame",
)]
}
}
_ => {
tracing::debug!("Unhandled frame tag: 0x{:02x}", tag);
vec![]
}
};
HandleOutput {
frames,
subscribe: None,
unsubscribe: None,
}
}
async fn handle_sub(payload: &[u8], store: &Db, agent: &crate::agents::ForAgent) -> HandleOutput {
let Ok(subject_str) = std::str::from_utf8(payload) else {
return HandleOutput {
frames: vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid SUB frame",
)],
..HandleOutput::default()
};
};
let subject = crate::Subject::from_raw(subject_str, store.get_base_domain().as_deref());
let refuse = |reason: &str| HandleOutput {
frames: vec![protocol::encode_error(
0,
protocol::error_code::UNAUTHORIZED_READ,
&format!("SUB refused for {subject}: {reason}"),
)],
..HandleOutput::default()
};
if !subject.is_local() {
tracing::warn!("can't subscribe to external resource: {subject}");
return HandleOutput::default();
}
let resource = match store.get_resource(&subject).await {
Ok(r) => r,
Err(_) => return refuse("not readable"),
};
if let Err(e) = crate::hierarchy::check_read(store, &resource, agent).await {
return refuse(&e.to_string());
}
HandleOutput {
frames: vec![],
subscribe: Some(subject_str.to_string()),
unsubscribe: None,
}
}
fn handle_unsub(payload: &[u8]) -> HandleOutput {
match std::str::from_utf8(payload) {
Ok(subject) => HandleOutput {
frames: vec![],
subscribe: None,
unsubscribe: Some(subject.to_string()),
},
Err(_) => HandleOutput::default(),
}
}
pub struct CommitIngestOpts {
pub source_id: Option<String>,
pub validate_loro_causality: bool,
pub enforce_subject_ownership: bool,
pub suppress_live_echo: bool,
pub response_origin: Option<String>,
}
impl CommitIngestOpts {
pub fn hub(source_id: Option<String>, response_origin: Option<String>) -> Self {
Self {
source_id,
validate_loro_causality: true,
enforce_subject_ownership: true,
suppress_live_echo: false,
response_origin,
}
}
pub fn peer() -> Self {
Self {
source_id: None,
validate_loro_causality: false,
enforce_subject_ownership: false,
suppress_live_echo: true,
response_origin: None,
}
}
}
pub async fn ingest_commit_json(
store: &Db,
commit_json: &str,
opts: &CommitIngestOpts,
) -> crate::errors::AtomicResult<String> {
let response = ingest_commit(store, commit_json, opts).await?;
let base_domain = store.get_base_domain();
let origin = opts.response_origin.as_deref().or(base_domain.as_deref());
let json = response.commit_resource.to_json_ad(origin)?;
Ok(json)
}
pub async fn ingest_commit(
store: &Db,
commit_json: &str,
opts: &CommitIngestOpts,
) -> crate::errors::AtomicResult<crate::commit::CommitResponse> {
if commit_json.contains("\"https://atomicdata.dev/properties/set\"")
|| commit_json.contains("\"https://atomicdata.dev/properties/push\"")
|| commit_json.contains("\"https://atomicdata.dev/properties/remove\"")
{
return Err(
"Commits with `set`, `push`, or `remove` fields are no longer accepted. Use `loroUpdate` instead."
.into(),
);
}
let incoming_commit_resource =
crate::parse::parse_json_ad_commit_resource(commit_json, store).await?;
let incoming_commit = crate::commit::Commit::from_resource(incoming_commit_resource)?;
if let Some(loro_bytes) = &incoming_commit.loro_update {
let doc = crate::loro::AtomicLoroDoc::new();
if doc.import_update(loro_bytes).is_ok() {
let props = doc.get_all_properties();
let prop_summary: Vec<String> = props
.keys()
.map(|k| k.rsplit('/').next().unwrap_or(k).to_string())
.collect();
tracing::info!(
subject = %incoming_commit.subject,
signer = %incoming_commit.signer,
properties = ?prop_summary,
loro_bytes = loro_bytes.len(),
"Incoming commit"
);
}
} else {
tracing::info!(
subject = %incoming_commit.subject,
destroy = ?incoming_commit.destroy,
"Incoming commit (no loroUpdate)"
);
}
if opts.enforce_subject_ownership {
let is_internal = incoming_commit.subject.is_internal();
let is_did = incoming_commit.subject.is_did();
let matches_base = if let Some(base) = store.get_base_domain() {
incoming_commit.subject.as_str().contains(&base)
} else {
false
};
let is_local_path =
!is_did && !is_internal && incoming_commit.subject.as_str().ends_with('/');
if !is_internal && !is_did && !matches_base && !is_local_path {
return Err(
"Subject of commit should be sent to other domain - this store can not own this resource."
.into(),
);
}
}
let signer = incoming_commit.signer.clone();
let signer_pure = signer.pure_id();
let is_self_creating_agent =
incoming_commit.subject.is_agent_did() && incoming_commit.subject == signer;
let needs_agent_resource = signer.is_agent_did()
&& !is_self_creating_agent
&& store.get_resource(&signer).await.is_err();
let commit_opts = crate::commit::CommitOpts {
validate_schema: true,
validate_signature: true,
validate_timestamp: true,
validate_rights: true,
validate_loro_causality: opts.validate_loro_causality,
validate_for_agent: Some(signer.to_string()),
update_index: true,
source_id: opts.source_id.clone(),
};
let response = if opts.suppress_live_echo {
super::ws_apply::set_importing(true);
let result = store.apply_commit(incoming_commit, &commit_opts).await;
super::ws_apply::set_importing(false);
result
} else {
store.apply_commit(incoming_commit, &commit_opts).await
}?;
if needs_agent_resource {
let mut new_agent = crate::Resource::new_instance(crate::urls::AGENT, store).await?;
new_agent.set_subject(signer_pure.clone());
if let Some(pk) = signer.as_str().strip_prefix("did:ad:agent:") {
new_agent
.set_string(crate::urls::PUBLIC_KEY.into(), pk, store)
.await?;
}
new_agent.save_locally(store).await?;
tracing::info!("Auto-created agent resource for {}", signer_pure);
}
Ok(response)
}
async fn apply_peer_commit(store: &Db, commit_json: &str) -> crate::errors::AtomicResult<String> {
ingest_commit_json(store, commit_json, &CommitIngestOpts::peer()).await
}
pub async fn collect_drive_subjects(
store: &Db,
drive_subject: &crate::Subject,
) -> std::collections::HashSet<String> {
let drive_str = drive_subject.pure_id();
let mut result = std::collections::HashSet::new();
result.insert(drive_str.clone());
if drive_subject.is_did() {
let mut queue = vec![drive_str];
while let Some(current) = queue.pop() {
let q = crate::storelike::Query {
property: Some(crate::urls::PARENT.into()),
value: Some(crate::Value::AtomicUrl(current.clone().into())),
filters: Vec::new(),
limit: None,
start_val: None,
end_val: None,
offset: 0,
sort_by: None,
sort_desc: false,
include_external: true,
include_nested: false,
for_agent: crate::agents::ForAgent::Sudo,
aggregation: None,
expression_filters: Vec::new(),
drive: None,
};
if let Ok(qr) = store.query(&q).await {
for child in qr.subjects {
let child_str = child.pure_id();
if result.insert(child_str.clone()) {
queue.push(child_str);
}
}
}
}
} else {
let drive_pure = drive_subject.pure_id();
for resource in store.all_resources(false) {
let subject = resource.get_subject();
if subject.pure_id().starts_with(&drive_pure) {
result.insert(subject.pure_id());
}
}
}
result
}
pub fn compute_drive_hash(
vvs: &std::collections::HashMap<String, std::collections::HashMap<String, i32>>,
) -> String {
let mut peer_set = std::collections::BTreeSet::new();
for vv in vvs.values() {
for peer_id in vv.keys() {
peer_set.insert(peer_id.clone());
}
}
let peers: Vec<String> = peer_set.into_iter().collect();
let peer_index: std::collections::HashMap<&str, usize> = peers
.iter()
.enumerate()
.map(|(i, p)| (p.as_str(), i))
.collect();
let mut entries: Vec<(String, Vec<i32>)> = vvs
.iter()
.map(|(subject, vv)| {
let mut counters = vec![0i32; peers.len()];
for (peer_id, &counter) in vv {
if let Some(&idx) = peer_index.get(peer_id.as_str()) {
counters[idx] = counter;
}
}
(subject.clone(), counters)
})
.collect();
entries.sort_by(|(a, _), (b, _)| a.cmp(b));
let hash_input: String = entries
.iter()
.map(|(s, c)| {
let counters = c
.iter()
.map(|n| n.to_string())
.collect::<Vec<_>>()
.join(",");
format!("{s}:{counters}")
})
.collect::<Vec<_>>()
.join("|");
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(hash_input.as_bytes());
hex::encode(hasher.finalize())
}
pub fn build_drive_vvs(
store: &Db,
drive_subjects: &std::collections::HashSet<String>,
) -> std::collections::HashMap<String, std::collections::HashMap<String, i32>> {
let mut vvs = std::collections::HashMap::new();
for subject_str in drive_subjects {
if let Ok(Some(snapshot_bytes)) = store.kv.get(Tree::LoroSnapshots, subject_str.as_bytes())
{
if let Ok(vv) = AtomicLoroDoc::vv_map_from_snapshot(&snapshot_bytes) {
vvs.insert(subject_str.clone(), vv);
}
}
}
vvs
}
pub async fn drive_items_for(
store: &Db,
drive: &str,
agent: &crate::agents::ForAgent,
) -> Result<Vec<crate::sync::rbsr::Item>, String> {
let drive_subject = crate::Subject::from_raw(drive, store.get_base_domain().as_deref());
let drive_resource = store
.get_resource(&drive_subject)
.await
.map_err(|_| "not readable".to_string())?;
crate::hierarchy::check_read(store, &drive_resource, agent)
.await
.map_err(|e| e.to_string())?;
let drive_subjects = collect_drive_subjects(store, &drive_subject).await;
let vvs = build_drive_vvs(store, &drive_subjects);
let mut items: Vec<crate::sync::rbsr::Item> = Vec::with_capacity(vvs.len());
for (subject, vv) in vvs {
let readable = match store
.get_resource(&crate::Subject::from_raw(
&subject,
store.get_base_domain().as_deref(),
))
.await
{
Ok(r) => crate::hierarchy::check_read(store, &r, agent).await.is_ok(),
Err(_) => false,
};
if readable {
items.push((subject, vv.into_iter().collect()));
}
}
items.sort_by(|a, b| a.0.cmp(&b.0));
Ok(items)
}
pub async fn drive_sync_hash_for(
store: &Db,
drive: &str,
agent: &crate::agents::ForAgent,
) -> Result<String, String> {
let items = drive_items_for(store, drive, agent).await?;
let vvs: std::collections::HashMap<String, std::collections::HashMap<String, i32>> = items
.into_iter()
.map(|(subject, vv)| (subject, vv.into_iter().collect()))
.collect();
Ok(compute_drive_hash(&vvs))
}
pub async fn handle_sync_vv(
drive: &str,
drive_hash: &str,
client_peers: &[String],
client_resources: &std::collections::HashMap<String, Vec<i32>>,
store: &Db,
agent: &crate::agents::ForAgent,
) -> Vec<Vec<u8>> {
handle_sync_vv_filtered(
drive,
drive_hash,
client_peers,
client_resources,
None,
store,
agent,
)
.await
}
pub async fn handle_sync_vv_filtered(
drive: &str,
drive_hash: &str,
client_peers: &[String],
client_resources: &std::collections::HashMap<String, Vec<i32>>,
subjects: Option<&std::collections::HashSet<String>>,
store: &Db,
agent: &crate::agents::ForAgent,
) -> Vec<Vec<u8>> {
let server_vvs = match subjects {
Some(set) => {
let mut vvs = std::collections::HashMap::new();
for subject in set {
if let Ok(Some(bytes)) = store.kv.get(Tree::LoroSnapshots, subject.as_bytes()) {
if let Ok(vv) = AtomicLoroDoc::vv_map_from_snapshot(&bytes) {
vvs.insert(subject.clone(), vv);
}
}
}
vvs
}
None => {
let drive_subject = crate::Subject::from_raw(drive, store.get_base_domain().as_deref());
let drive_subjects = collect_drive_subjects(store, &drive_subject).await;
build_drive_vvs(store, &drive_subjects)
}
};
if !drive_hash.is_empty() {
let server_hash = compute_drive_hash(&server_vvs);
if server_hash == drive_hash {
tracing::info!("SYNC: drive {} — hashes match, in sync", drive);
return vec![protocol::encode_sync_ok(drive)];
}
}
let mut client_vvs: std::collections::HashMap<String, std::collections::HashMap<String, i32>> =
std::collections::HashMap::new();
for (subject, counters) in client_resources {
let mut vv = std::collections::HashMap::new();
for (i, &counter) in counters.iter().enumerate() {
if counter != 0 {
if let Some(peer_id) = client_peers.get(i) {
vv.insert(peer_id.clone(), counter);
}
}
}
client_vvs.insert(subject.clone(), vv);
}
let mut pull: Vec<String> = Vec::new();
let mut pull_from: std::collections::HashMap<String, std::collections::HashMap<String, i32>> =
std::collections::HashMap::new();
let mut remove: Vec<String> = Vec::new();
let mut remove_commits: std::collections::HashMap<String, String> =
std::collections::HashMap::new();
let mut may_read_drive: Option<bool> = None;
let mut push_entries: Vec<(String, Vec<u8>)> = Vec::new();
for (subject, server_vv) in &server_vvs {
let resource = match store
.get_resource(&crate::Subject::from_raw(
subject,
store.get_base_domain().as_deref(),
))
.await
{
Ok(r) => {
if crate::hierarchy::check_read(store, &r, agent)
.await
.is_err()
{
continue;
}
r
}
Err(_) => continue,
};
if let Some(client_vv) = client_vvs.get(subject) {
let server_ahead = server_vv
.iter()
.any(|(p, &sc)| client_vv.get(p).copied().unwrap_or(0) < sc);
let client_ahead = client_vv
.iter()
.any(|(p, &cc)| server_vv.get(p).copied().unwrap_or(0) < cc);
if server_ahead {
if let Ok(Some(snapshot_bytes)) =
store.kv.get(Tree::LoroSnapshots, subject.as_bytes())
{
if let Ok(doc) = AtomicLoroDoc::from_snapshot(&snapshot_bytes) {
let client_loro_vv = AtomicLoroDoc::vv_from_map(client_vv);
let delta = doc.export_updates_since(&client_loro_vv);
if !delta.is_empty() {
push_entries.push((subject.clone(), delta));
}
}
}
}
if client_ahead {
pull.push(subject.clone());
pull_from.insert(subject.clone(), server_vv.clone());
}
if let Ok(blob_val) = resource.get(crate::urls::BLOB) {
let blob_did = blob_val.to_string();
if let Some(hash_hex) = crate::Subject::from_raw(&blob_did, None).blob_hash_hex() {
if let Ok(hash_bytes) = hex::decode(hash_hex) {
if hash_bytes.len() == 32 {
let mut hash = [0u8; 32];
hash.copy_from_slice(&hash_bytes);
if !store.kv.contains_key(Tree::Blobs, &hash).unwrap_or(false) {
if !pull.contains(subject) {
pull.push(subject.clone());
pull_from.insert(subject.clone(), server_vv.clone());
}
}
}
}
}
}
} else {
if let Ok(Some(snapshot_bytes)) = store.kv.get(Tree::LoroSnapshots, subject.as_bytes())
{
push_entries.push((subject.clone(), snapshot_bytes));
}
}
}
for subject in client_vvs.keys() {
if subjects.is_some_and(|set| !set.contains(subject)) {
continue;
}
if !server_vvs.contains_key(subject) {
if super::tombstones::is_tombstoned(store, subject) {
remove.push(subject.clone());
if let Some(json) = super::tombstones::destroy_envelope(store, subject) {
let allowed = match may_read_drive {
Some(v) => v,
None => {
let v = agent_may_read_drive(store, drive, agent).await;
may_read_drive = Some(v);
v
}
};
if allowed {
remove_commits.insert(subject.clone(), json);
}
}
} else {
pull.push(subject.clone());
pull_from
.entry(subject.clone())
.or_insert_with(std::collections::HashMap::new);
}
}
}
let push_subjects: Vec<String> = push_entries.iter().map(|(s, _)| s.clone()).collect();
tracing::info!(
"SYNC: drive {} — {} to push, {} to pull, {} to remove",
drive,
push_subjects.len(),
pull.len(),
remove.len(),
);
let mut frames = Vec::new();
frames.push(protocol::encode_sync_diff(
drive,
&pull,
&push_subjects,
&remove,
&pull_from,
&remove_commits,
));
if !push_entries.is_empty() {
let entries: Vec<(&str, &[u8])> = push_entries
.iter()
.map(|(s, b)| (s.as_str(), b.as_slice()))
.collect();
for chunk in protocol::encode_sync_push_chunks(drive, &entries) {
frames.push(chunk);
}
}
frames
}
pub(crate) async fn may_accept_drive_write(
store: &Db,
drive_resource: &crate::Resource,
for_agent: &crate::agents::ForAgent,
trust_owned: bool,
) -> bool {
if crate::hierarchy::check_write(store, drive_resource, for_agent)
.await
.is_ok()
{
return true;
}
if trust_owned {
if let Ok(own) = store.get_default_agent() {
let own_agent = crate::agents::ForAgent::from(own);
if crate::hierarchy::check_write(store, drive_resource, &own_agent)
.await
.is_ok()
{
return true;
}
}
}
false
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SyncPushRejected {
pub drive: String,
pub reason: String,
}
impl SyncPushRejected {
pub fn to_error_frame(&self) -> Vec<u8> {
protocol::encode_error(0, protocol::error_code::SYNC_REJECTED, &self.to_string())
}
}
impl std::fmt::Display for SyncPushRejected {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"SYNC_PUSH rejected for drive {}: {}",
self.drive, self.reason
)
}
}
async fn agent_may_read_drive(store: &Db, drive: &str, agent: &crate::agents::ForAgent) -> bool {
let drive_subject = crate::Subject::from_raw(drive, store.get_base_domain().as_deref());
match store.get_resource(&drive_subject).await {
Ok(resource) => crate::hierarchy::check_read(store, &resource, agent)
.await
.is_ok(),
Err(_) => false,
}
}
pub(crate) fn admit_unknown_drive(
store: &Db,
drive_subject: &str,
agent: &crate::agents::ForAgent,
) -> bool {
if matches!(agent, crate::agents::ForAgent::Public) {
return false;
}
let policy = store.sync_policy();
if policy.admit_drive_write(drive_subject) {
return true;
}
if policy.may_enroll_drive(drive_subject, agent) {
tracing::info!("enrolling new drive {} for {:?}", drive_subject, agent);
policy.enroll_drive(drive_subject);
true
} else {
false
}
}
pub async fn import_sync_push(
push: &protocol::DecodedSyncPush,
store: &Db,
for_agent: &crate::agents::ForAgent,
trust_owned: bool,
) -> Result<(usize, Vec<Vec<u8>>), SyncPushRejected> {
let drive_subject = crate::Subject::from_raw(&push.drive, store.get_base_domain().as_deref());
let policy = store.sync_policy();
if let Ok(drive_resource) = store.get_resource(&drive_subject).await {
if !may_accept_drive_write(store, &drive_resource, for_agent, trust_owned).await {
tracing::warn!(
"import_sync_push: agent {:?} has no write access to drive {} (trust_owned={})",
for_agent,
push.drive,
trust_owned
);
return Err(SyncPushRejected {
drive: push.drive.clone(),
reason: format!("agent {for_agent} has no write right on the drive"),
});
}
let decision = policy.admit_decision(&push.drive);
if !decision.is_admitted() {
tracing::warn!(
"import_sync_push: drive {} not admitted by sync policy ({:?}, agent {:?})",
push.drive,
decision,
for_agent
);
let reason = match decision {
super::policy::AdmitDecision::OverQuota => {
format!(
"drive {} is over its storage quota on this node",
push.drive
)
}
_ => policy.not_enrolled_message(&push.drive),
};
return Err(SyncPushRejected {
drive: push.drive.clone(),
reason,
});
}
} else if !admit_unknown_drive(store, &push.drive, for_agent) {
tracing::warn!(
"import_sync_push: refusing bootstrap of unknown drive {} for {:?}",
push.drive,
for_agent
);
let reason = if matches!(for_agent, crate::agents::ForAgent::Public) {
"unauthenticated agent cannot create a drive".to_string()
} else {
policy.not_enrolled_message(&push.drive)
};
return Err(SyncPushRejected {
drive: push.drive.clone(),
reason,
});
}
let mut count = 0;
let mut blob_requests = Vec::new();
let base_domain = store.get_base_domain();
let normalize = |s: &str| crate::Subject::from_raw(s, base_domain.as_deref()).pure_id();
let admitted_drive = normalize(&push.drive);
for entry in &push.entries {
if super::tombstones::is_tombstoned(store, &entry.subject) {
tracing::debug!(
"import_sync_push: skip {:?} (tombstoned locally)",
&entry.subject[..entry.subject.len().min(24)]
);
continue;
}
let snapshot_key =
crate::Subject::from_raw(&entry.subject, store.get_base_domain().as_deref()).pure_id();
let _subject_guard = store.subject_locks.lock(&snapshot_key).await;
let existing_resource = store
.get_resource(&crate::Subject::from_raw(
&snapshot_key,
base_domain.as_deref(),
))
.await
.ok();
if let Some(existing) = &existing_resource {
let stored_drive = existing
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| existing.get_subject().to_string());
if normalize(&stored_drive) != admitted_drive {
tracing::warn!(
"import_sync_push: {} belongs to drive {}, not to {} this push was admitted for; skipped",
entry.subject,
stored_drive,
push.drive
);
continue;
}
}
let doc = if let Ok(Some(existing)) =
store.kv.get(Tree::LoroSnapshots, snapshot_key.as_bytes())
{
match AtomicLoroDoc::from_snapshot(&existing) {
Ok(d) => {
if d.import_update(&entry.loro_bytes).is_err() {
tracing::warn!(
"import_sync_push: delta import failed for {}",
entry.subject
);
continue;
}
d
}
Err(_) => {
match AtomicLoroDoc::from_snapshot(&entry.loro_bytes) {
Ok(d) => d,
Err(_) => continue,
}
}
}
} else {
let doc = AtomicLoroDoc::new();
if doc.import_update(&entry.loro_bytes).is_err() {
match AtomicLoroDoc::from_snapshot(&entry.loro_bytes) {
Ok(d) => d,
Err(_) => {
tracing::warn!("import_sync_push: import failed for {}", entry.subject);
continue;
}
}
} else {
doc
}
};
let subject = crate::Subject::from_raw(&snapshot_key, store.get_base_domain().as_deref());
let mut resource = crate::Resource::new(subject.to_string());
if resource.apply_state_doc(doc).is_err() {
continue;
}
if existing_resource.is_none() {
let mut claimed = resource
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| 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 {
claimed = parent_res
.get(crate::urls::DRIVE_PROP)
.map(|v| v.to_string())
.unwrap_or_else(|_| parent_subject.to_string());
}
}
if normalize(&claimed) != admitted_drive {
tracing::warn!(
"import_sync_push: new resource {} resolves to drive {}, not to {} this push was admitted for; skipped",
entry.subject,
claimed,
push.drive
);
continue;
}
}
let has_strokes = resource
.get("https://atomicdata.dev/ontology/canvas/strokeData")
.is_ok();
tracing::info!(
" sync imported {}: {} props, has_strokes={}",
&entry.subject[..entry.subject.len().min(30)],
resource.get_propvals().len(),
has_strokes,
);
if store
.add_resource_opts(&resource, false, true, true)
.await
.is_err()
{
continue;
}
count += 1;
if let Ok(blob_val) = resource.get(crate::urls::BLOB) {
let blob_did = blob_val.to_string();
if let Some(hash_hex) = crate::Subject::from_raw(&blob_did, None).blob_hash_hex() {
if let Ok(hash_bytes) = hex::decode(hash_hex) {
if hash_bytes.len() == 32 {
let mut hash = [0u8; 32];
hash.copy_from_slice(&hash_bytes);
if !store.kv.contains_key(Tree::Blobs, &hash).unwrap_or(false) {
store.note_pending_blob_request(hash, push.drive.clone());
blob_requests.push(protocol::encode_blob_request(&hash));
}
}
}
}
}
}
tracing::info!(
"import_sync_push: imported {} resources for drive {}",
count,
push.drive
);
for entry in &push.entries {
tracing::info!(
" imported: {} ({} bytes)",
&entry.subject[..entry.subject.len().min(30)],
entry.loro_bytes.len()
);
}
Ok((count, blob_requests))
}
#[cfg(feature = "iroh")]
fn peer_is_paired(store: &Db, node_id: &str) -> bool {
crate::sync::peer::is_paired_peer(store, node_id)
}
#[cfg(not(feature = "iroh"))]
fn peer_is_paired(_store: &Db, _node_id: &str) -> bool {
false
}
pub async fn collect_readable_snapshots(
store: &Db,
agent: &crate::agents::ForAgent,
subjects: &[String],
paired_peer: Option<&str>,
) -> Vec<(String, Vec<u8>)> {
let own_agent = if paired_peer.is_some_and(|node| peer_is_paired(store, node)) {
store
.get_default_agent()
.ok()
.map(crate::agents::ForAgent::from)
} else {
None
};
let mut entries = Vec::new();
for subject in subjects {
let subj = crate::Subject::from_raw(subject, store.get_base_domain().as_deref());
match store.get_resource(&subj).await {
Ok(resource) => {
let mut readable = crate::hierarchy::check_read(store, &resource, agent)
.await
.is_ok();
if !readable {
if let Some(own) = own_agent.as_ref() {
readable = crate::hierarchy::check_read(store, &resource, own)
.await
.is_ok();
if readable {
tracing::debug!(
"[sync] serving {} to a paired replica",
&subject[..subject.len().min(30)]
);
}
}
}
if !readable {
tracing::warn!(
"[sync] refusing to serve {} to peer: no read access for {:?}",
&subject[..subject.len().min(30)],
agent
);
continue;
}
}
Err(_) => continue,
}
if let Ok(Some(snapshot)) = store
.kv
.get(crate::db::trees::Tree::LoroSnapshots, subject.as_bytes())
{
entries.push((subject.clone(), snapshot));
}
}
entries
}
#[cfg(test)]
mod bootstrap_and_sub_tests {
use super::*;
use crate::agents::ForAgent;
use crate::sync::policy::OwnerPolicy;
use crate::sync::protocol::{self, error_code, tag};
use std::sync::Arc;
fn empty_push(drive: &str) -> protocol::DecodedSyncPush {
let frame = protocol::encode_sync_push(drive, &[], true);
protocol::decode_sync_push(&frame[1..]).unwrap()
}
#[tokio::test]
async fn rejected_sync_entry_does_not_persist_snapshot() {
let db = Db::init_temp("rejected_sync_snapshot").await.unwrap();
let (alice, drive) = db.setup("Alice").await.unwrap();
let subject = "did:ad:unseen-victim-resource";
let doc = AtomicLoroDoc::new();
doc.set_property(
crate::urls::DRIVE_PROP,
&crate::Value::AtomicUrl("did:ad:other-drive".into()),
)
.unwrap();
doc.set_property(
crate::urls::NAME,
&crate::Value::String("Rejected data".into()),
)
.unwrap();
let frame = protocol::encode_sync_push(&drive, &[(subject, &doc.export_snapshot())], true);
let push = protocol::decode_sync_push(&frame[1..]).unwrap();
let (count, _) = import_sync_push(&push, &db, &ForAgent::from(alice.clone()), false)
.await
.unwrap();
assert_eq!(count, 0);
assert!(
db.kv
.get(Tree::LoroSnapshots, subject.as_bytes())
.unwrap()
.is_none(),
"rejected data must not seed a later legitimate merge"
);
assert!(db.get_resource(&subject.into()).await.is_err());
let valid = AtomicLoroDoc::new();
valid
.set_property(
crate::urls::DRIVE_PROP,
&crate::Value::AtomicUrl(drive.clone().into()),
)
.unwrap();
let frame =
protocol::encode_sync_push(&drive, &[(subject, &valid.export_snapshot())], true);
let push = protocol::decode_sync_push(&frame[1..]).unwrap();
let (count, _) = import_sync_push(&push, &db, &ForAgent::from(alice), false)
.await
.unwrap();
assert_eq!(count, 1);
let stored = db.get_resource(&subject.into()).await.unwrap();
assert!(
stored.get(crate::urls::NAME).is_err(),
"rejected properties must not reappear"
);
assert!(db
.kv
.get(Tree::LoroSnapshots, subject.as_bytes())
.unwrap()
.is_some());
}
#[tokio::test]
async fn public_cannot_bootstrap_a_missing_drive_even_on_open() {
let db = Db::init_temp("oq5_public_open").await.unwrap();
let drive = "did:ad:newdrivepublic";
let push = empty_push(drive);
let err = import_sync_push(&push, &db, &ForAgent::Public, false)
.await
.expect_err("Public must not create a drive on an open node");
assert!(
err.reason.contains("unauthenticated"),
"reason names the cause: {}",
err.reason
);
assert!(
db.get_resource(&drive.into()).await.is_err(),
"the refused push must not have stored the drive"
);
}
#[tokio::test]
async fn authenticated_first_sync_on_open_still_admits_a_new_drive() {
let db = Db::init_temp("oq5_auth_open").await.unwrap();
let (alice, _) = db.setup("Alice").await.unwrap();
let drive = "did:ad:newdrivealice";
let push = empty_push(drive);
import_sync_push(
&push,
&db,
&ForAgent::AgentSubject(alice.subject.clone()),
false,
)
.await
.expect("an authenticated agent on Open may bootstrap a drive");
}
#[tokio::test]
async fn owner_mode_refuses_a_stranger_bootstrapping_a_new_drive() {
let db = Db::init_temp("oq5_owner_stranger").await.unwrap();
let (owner, _) = db.setup("Owner").await.unwrap();
let stranger = db.create_agent(Some("Stranger")).await.unwrap();
db.set_sync_policy(Arc::new(OwnerPolicy::new(owner.subject.to_string())));
let drive = "did:ad:strangerdrive";
let push = empty_push(drive);
let err = import_sync_push(
&push,
&db,
&ForAgent::AgentSubject(stranger.subject.clone()),
false,
)
.await
.expect_err("a stranger must not dump a drive onto an owner-gated node");
assert!(
err.reason.contains("does not host new Drives"),
"the refusal speaks to the visitor: {}",
err.reason
);
assert!(
!db.sync_policy().admit_drive_write(drive),
"the refused drive must not have been enrolled"
);
}
#[tokio::test]
async fn owner_mode_enrolls_the_owners_new_drive_from_sync_push() {
let db = Db::init_temp("oq5_owner_self").await.unwrap();
let (owner, _) = db.setup("Owner").await.unwrap();
db.set_sync_policy(Arc::new(OwnerPolicy::new(owner.subject.to_string())));
let drive = "did:ad:ownerssecond";
let push = empty_push(drive);
import_sync_push(
&push,
&db,
&ForAgent::AgentSubject(owner.subject.clone()),
false,
)
.await
.expect("the owner may bootstrap a second drive");
assert!(
db.sync_policy().admit_drive_write(drive),
"the owner's new drive must be enrolled so later writes land"
);
}
#[tokio::test]
async fn sub_on_a_public_drive_is_a_session_command_not_an_error() {
let db = Db::init_temp("sub_public").await.unwrap();
let (_alice, drive) = db.setup("Alice").await.unwrap();
let mut resource = db.get_resource(&drive.as_str().into()).await.unwrap();
resource
.set_unsafe(
crate::urls::READ.into(),
crate::Value::ResourceArray(vec![crate::urls::PUBLIC_AGENT.into()]),
)
.unwrap();
db.add_resource_opts(&resource, false, true, true)
.await
.unwrap();
let frame = protocol::encode_sub(&drive);
let mut agent = ForAgent::Public;
let out = handle_frame_full(&frame, &db, &mut agent).await;
assert!(
out.frames.is_empty(),
"a granted SUB has no reply on the wire"
);
assert_eq!(out.subscribe.as_deref(), Some(drive.as_str()));
assert!(out.unsubscribe.is_none());
}
#[tokio::test]
async fn sub_without_read_right_is_refused_out_loud() {
let db = Db::init_temp("sub_denied").await.unwrap();
let (_alice, drive) = db.setup("Alice").await.unwrap();
let mallory = db.create_agent(Some("Mallory")).await.unwrap();
let frame = protocol::encode_sub(&drive);
let mut agent = ForAgent::AgentSubject(mallory.subject.clone());
let out = handle_frame_full(&frame, &db, &mut agent).await;
assert!(out.subscribe.is_none(), "must not ask the hub to register");
let err = out
.frames
.iter()
.find(|f| f.first() == Some(&tag::ERROR))
.expect("an unauthorized SUB is answered with ERROR");
assert_eq!(
u16::from_be_bytes([err[3], err[4]]),
error_code::UNAUTHORIZED_READ
);
let msg = String::from_utf8_lossy(&err[5..]);
assert!(msg.contains("SUB refused"), "{msg}");
assert!(msg.contains(&drive), "{msg}");
}
async fn signed_destroy_json(
db: &Db,
agent: &crate::agents::Agent,
subject: &crate::Subject,
) -> String {
let resource = db.get_resource(subject).await.unwrap();
let mut builder = crate::commit::CommitBuilder::new(subject.clone());
builder.destroy(true);
let commit = builder.sign(agent, db, &resource).await.unwrap();
commit
.into_resource(db)
.await
.unwrap()
.to_json_ad(None)
.unwrap()
}
async fn secret_child(db: &Db, drive: &str) -> String {
db.create_resource(
crate::urls::CLASS,
drive,
"Secret Doc",
Some(vec![
(
crate::urls::DESCRIPTION,
crate::Value::String("top secret".into()),
),
(
crate::urls::SHORTNAME,
crate::Value::Slug("secret-doc".into()),
),
]),
)
.await
.unwrap()
}
fn decode_diff(frames: Vec<Vec<u8>>) -> protocol::DecodedSyncDiff {
frames
.into_iter()
.find(|f| f.first() == Some(&tag::SYNC_DIFF))
.and_then(|f| protocol::decode_sync_diff(&f[1..]))
.expect("reconcile answers with a SYNC_DIFF")
}
#[tokio::test]
async fn sync_diff_envelope_is_only_handed_to_drive_readers() {
let db = Db::init_temp("envelope_read_gate").await.unwrap();
let (alice, drive) = db.setup("Alice").await.unwrap();
let child = secret_child(&db, &drive).await;
let subject = crate::Subject::from_raw(&child, None);
let json = signed_destroy_json(&db, &alice, &subject).await;
ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect("the owner's destroy applies");
assert!(super::super::tombstones::destroy_envelope(&db, &child).is_some());
let mut client_resources = std::collections::HashMap::new();
client_resources.insert(child.clone(), Vec::<i32>::new());
let public = decode_diff(
handle_sync_vv_filtered(
&drive,
"",
&[],
&client_resources,
None,
&db,
&ForAgent::Public,
)
.await,
);
assert_eq!(public.remove, vec![child.clone()]);
assert!(
public.remove_commits.is_empty(),
"an anonymous session must not receive the signed destroy envelope"
);
let owner = decode_diff(
handle_sync_vv_filtered(
&drive,
"",
&[],
&client_resources,
None,
&db,
&ForAgent::AgentSubject(alice.subject.clone()),
)
.await,
);
assert!(
owner.remove_commits.contains_key(&child),
"a drive reader receives the envelope"
);
}
#[tokio::test]
async fn replayed_destroy_is_refused_after_the_subject_is_recreated() {
let db = Db::init_temp("destroy_replay_stored").await.unwrap();
let (alice, drive) = db.setup("Alice").await.unwrap();
let subject = crate::Subject::from_raw(&drive, None);
let json = signed_destroy_json(&db, &alice, &subject).await;
ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect("the first destroy applies");
assert!(db.get_resource(&subject).await.is_err());
let again = db.ensure_private_drive().await.unwrap();
assert_eq!(again, drive, "repeat genesis recreates the same subject");
assert!(db.get_resource(&subject).await.is_ok());
let err = ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect_err("the same envelope must not destroy the recreated drive");
assert!(err.to_string().contains("already applied"), "{err}");
assert!(
db.get_resource(&subject).await.is_ok(),
"the recreated drive survives the replay"
);
assert!(!super::super::tombstones::is_tombstoned(&db, &drive));
}
#[tokio::test]
async fn destroy_predating_the_genesis_is_refused() {
let db = Db::init_temp("destroy_replay_genesis").await.unwrap();
let (alice, drive) = db.setup("Alice").await.unwrap();
let child = secret_child(&db, &drive).await;
let subject = crate::Subject::from_raw(&child, None);
let mut resource = db.get_resource(&subject).await.unwrap();
resource
.set_unsafe(
crate::urls::CREATED_AT.into(),
crate::Value::Timestamp(crate::utils::now() + 60_000),
)
.unwrap();
db.add_resource_opts(&resource, false, true, true)
.await
.unwrap();
let json = signed_destroy_json(&db, &alice, &subject).await;
let err = ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect_err("a destroy stamped before the genesis is a replay");
assert!(err.to_string().contains("predates"), "{err}");
assert!(db.get_resource(&subject).await.is_ok());
assert!(!super::super::tombstones::is_tombstoned(&db, &child));
}
async fn signed_rename_json(
db: &Db,
agent: &crate::agents::Agent,
subject: &crate::Subject,
) -> String {
let resource = db.get_resource(subject).await.unwrap();
let mut builder = crate::commit::CommitBuilder::new(subject.clone());
builder.set(
crate::urls::NAME.into(),
crate::Value::String(format!("renamed by {}", agent.subject)),
);
let commit = builder.sign(agent, db, &resource).await.unwrap();
commit
.into_resource(db)
.await
.unwrap()
.to_json_ad(None)
.unwrap()
}
#[tokio::test]
async fn signer_agent_is_only_created_for_an_accepted_commit() {
let db = Db::init_temp("agent_after_accept").await.unwrap();
let (alice, drive) = db.setup("Alice").await.unwrap();
let child = secret_child(&db, &drive).await;
let subject = crate::Subject::from_raw(&child, None);
let stranger = crate::agents::Agent::new(Some("Stranger")).unwrap();
let json = signed_rename_json(&db, &stranger, &subject).await;
ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect_err("a stranger may not edit Alice's document");
assert!(
!db.has_resource_locally(&stranger.subject.pure_id()),
"a refused commit must not store the signer's Agent resource"
);
let bob = crate::agents::Agent::new(Some("Bob")).unwrap();
let drive_subject = crate::Subject::from_raw(&drive, None);
let mut drive_res = db.get_resource(&drive_subject).await.unwrap();
drive_res
.set_unsafe(
crate::urls::WRITE.into(),
crate::Value::ResourceArray(vec![
alice.subject.to_string().into(),
bob.subject.to_string().into(),
]),
)
.unwrap();
db.add_resource_opts(&drive_res, false, true, true)
.await
.unwrap();
assert!(!db.has_resource_locally(&bob.subject.pure_id()));
let json = signed_rename_json(&db, &bob, &subject).await;
ingest_commit_json(&db, &json, &CommitIngestOpts::peer())
.await
.expect("a drive writer's commit applies");
let renamed = db.get_resource(&subject).await.unwrap();
assert!(
renamed
.get(crate::urls::NAME)
.unwrap()
.to_string()
.contains("Bob")
|| renamed
.get(crate::urls::NAME)
.unwrap()
.to_string()
.contains(&bob.subject.to_string()),
"Bob's accepted commit lands"
);
}
#[tokio::test]
async fn unsub_is_a_session_command() {
let db = Db::init_temp("unsub_cmd").await.unwrap();
let frame = protocol::encode_unsub("did:ad:whatever");
let mut agent = ForAgent::Public;
let out = handle_frame_full(&frame, &db, &mut agent).await;
assert!(out.frames.is_empty());
assert_eq!(out.unsubscribe.as_deref(), Some("did:ad:whatever"));
}
}