use crate::db::trees::Tree;
use crate::loro::AtomicLoroDoc;
use crate::{Db, Storelike};
use super::protocol;
pub async fn handle_frame(
frame: &[u8],
store: &Db,
agent: &mut crate::agents::ForAgent,
) -> Vec<Vec<u8>> {
if frame.is_empty() {
return vec![];
}
let tag = frame[0];
let payload = &frame[1..];
match tag {
protocol::tag::AUTH => {
if let Ok(json) = std::str::from_utf8(payload) {
match serde_json::from_str::<crate::authentication::AuthValues>(json) {
Ok(auth) => {
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) => vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
&format!("Auth failed: {e}"),
)],
}
}
Err(e) => vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
&format!("Invalid auth JSON: {e}"),
)],
}
} else {
vec![protocol::encode_error(
0,
protocol::error_code::UNKNOWN,
"Invalid UTF-8 in auth",
)]
}
}
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 => {
if let Some(sync) = protocol::decode_sync(payload) {
handle_sync_vv(
&sync.drive,
&sync.drive_hash,
&sync.peers,
&sync.resources,
store,
agent,
)
.await
} else {
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) {
let (_count, mut blob_requests) =
import_sync_push(&push, store, agent, false).await;
let mut responses = vec![protocol::encode_sync_ok(&push.drive)];
responses.append(&mut blob_requests);
responses
} 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 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![]
}
}
}
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>,
}
pub async fn ingest_commit_json(
store: &Db,
commit_json: &str,
opts: &CommitIngestOpts,
) -> crate::errors::AtomicResult<String> {
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;
if signer.is_agent_did()
&& !is_self_creating_agent
&& store.get_resource(&signer).await.is_err()
{
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);
}
let commit_opts = crate::commit::CommitOpts {
validate_schema: true,
validate_signature: true,
validate_timestamp: true,
validate_rights: true,
validate_previous_commit: false,
validate_loro_causality: opts.validate_loro_causality,
validate_for_agent: Some(signer.to_string()),
update_index: true,
source_id: opts.source_id.clone(),
};
let base_domain = store.get_base_domain();
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?
};
let origin = opts.response_origin.as_deref().or(base_domain.as_deref());
let json = response.commit_resource.to_json_ad(origin)?;
Ok(json)
}
async fn apply_peer_commit(store: &Db, commit_json: &str) -> crate::errors::AtomicResult<String> {
ingest_commit_json(
store,
commit_json,
&CommitIngestOpts {
source_id: None,
validate_loro_causality: false,
enforce_subject_ownership: false,
suppress_live_echo: true,
response_origin: None,
},
)
.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(store: &Db, drive: &str) -> Vec<crate::sync::rbsr::Item> {
let drive_subject = crate::Subject::from_raw(drive, store.get_base_domain().as_deref());
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> = vvs
.into_iter()
.map(|(subject, vv)| (subject, vv.into_iter().collect()))
.collect();
items.sort_by(|a, b| a.0.cmp(&b.0));
items
}
pub async fn drive_sync_hash(store: &Db, drive: &str) -> String {
let drive_subject = crate::Subject::from_raw(drive, store.get_base_domain().as_deref());
let drive_subjects = collect_drive_subjects(store, &drive_subject).await;
let server_vvs = build_drive_vvs(store, &drive_subjects);
compute_drive_hash(&server_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_VV: 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 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());
} 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_VV: 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,
));
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
}
pub async fn import_sync_push(
push: &protocol::DecodedSyncPush,
store: &Db,
for_agent: &crate::agents::ForAgent,
trust_owned: bool,
) -> (usize, Vec<Vec<u8>>) {
let drive_subject = crate::Subject::from_raw(&push.drive, store.get_base_domain().as_deref());
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 (0, vec![]);
}
}
let decision = store.sync_policy().admit_decision(&push.drive);
if !decision.is_admitted() {
tracing::warn!(
"import_sync_push: drive {} not admitted by sync policy ({:?})",
push.drive,
decision
);
return (0, vec![]);
}
let mut count = 0;
let mut blob_requests = Vec::new();
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 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 snapshot = doc.export_snapshot();
if store
.kv
.insert(Tree::LoroSnapshots, snapshot_key.as_bytes(), &snapshot)
.is_err()
{
continue;
}
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;
}
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,
);
let _ = store.add_resource_opts(&resource, false, true, true).await;
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()
);
}
(count, blob_requests)
}
pub async fn collect_readable_snapshots(
store: &Db,
agent: &crate::agents::ForAgent,
subjects: &[String],
) -> Vec<(String, Vec<u8>)> {
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) => {
if crate::hierarchy::check_read(store, &resource, agent)
.await
.is_err()
{
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
}