From 2d23aa5b7d27490cf19a1e52b3b87287b2eb589a Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Sat, 7 Feb 2026 02:12:01 +0300 Subject: [PATCH] [db,api] make block keys use cid raw bytes, deduplicate listRecords impl --- Cargo.lock | 1 + Cargo.toml | 1 + src/api/debug.rs | 18 +++-- src/api/stream.rs | 7 +- src/api/xrpc.rs | 172 +++++++++++++++++++------------------------- src/backfill/mod.rs | 68 +++++++++--------- src/db/keys.rs | 15 ++-- src/db/mod.rs | 47 ++++++++++-- src/ops.rs | 44 +++++++----- 9 files changed, 198 insertions(+), 175 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e55787d..8e77e8a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1621,6 +1621,7 @@ dependencies = [ "async-stream", "axum", "chrono", + "cid", "data-encoding", "fjall", "futures", diff --git a/Cargo.toml b/Cargo.toml index 75e85fa..ba2a483 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -41,3 +41,4 @@ mimalloc = { version = "0.1", features = ["v3"] } hex = "0.4" scc = "3" data-encoding = "2.10.0" +cid = "0.11.1" diff --git a/src/api/debug.rs b/src/api/debug.rs index ef9097c..0843d05 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -42,14 +42,14 @@ pub async fn handle_debug_count( .map_err(|_| StatusCode::BAD_REQUEST)?; let db = &state.db; - let ks = db.records.clone(); + let ks = db + .record_partition(&req.collection) + .map_err(|_| StatusCode::NOT_FOUND)?; - // {did_prefix}\x00{collection}\x00 + // {TrimmedDid}\x00 let mut prefix = Vec::new(); TrimmedDid::from(&did).write_to_vec(&mut prefix); prefix.push(keys::SEP); - prefix.extend_from_slice(req.collection.as_bytes()); - prefix.push(keys::SEP); let count = tokio::task::spawn_blocking(move || { let start_key = prefix.clone(); @@ -185,7 +185,6 @@ pub async fn handle_debug_iter( let items = tokio::task::spawn_blocking(move || { let limit = req.limit.unwrap_or(50).min(1000); - // Helper closure to avoid generic type complexity let collect = |iter: &mut dyn Iterator| { let mut items = Vec::new(); for guard in iter.take(limit) { @@ -247,13 +246,18 @@ pub async fn handle_debug_iter( fn get_keyspace_by_name(db: &crate::db::Db, name: &str) -> Result { match name { "repos" => Ok(db.repos.clone()), - "records" => Ok(db.records.clone()), "blocks" => Ok(db.blocks.clone()), "cursors" => Ok(db.cursors.clone()), "pending" => Ok(db.pending.clone()), "resync" => Ok(db.resync.clone()), "events" => Ok(db.events.clone()), "counts" => Ok(db.counts.clone()), - _ => Err(StatusCode::BAD_REQUEST), + _ => { + if let Some(col) = name.strip_prefix(crate::db::RECORDS_PARTITION_PREFIX) { + db.record_partition(col).map_err(|_| StatusCode::NOT_FOUND) + } else { + Err(StatusCode::BAD_REQUEST) + } + } } } diff --git a/src/api/stream.rs b/src/api/stream.rs index 9de482c..9e66e78 100644 --- a/src/api/stream.rs +++ b/src/api/stream.rs @@ -97,11 +97,8 @@ async fn handle_socket(mut socket: WebSocket, state: Arc, query: Strea let marshallable = { let mut record_val = None; - if let Some(cid_struct) = &cid { - let cid_str = cid_struct.to_string(); - if let Ok(Some(block_bytes)) = - db.blocks.get(keys::block_key(&cid_str)) - { + if let Some(cid) = &cid { + if let Ok(Some(block_bytes)) = db.blocks.get(&cid.to_bytes()) { if let Ok(raw_data) = serde_ipld_dagcbor::from_slice::(&block_bytes) { diff --git a/src/api/xrpc.rs b/src/api/xrpc.rs index 5292feb..f68803f 100644 --- a/src/api/xrpc.rs +++ b/src/api/xrpc.rs @@ -3,6 +3,7 @@ use crate::db::types::TrimmedDid; use crate::db::{self, Db, keys}; use axum::{Json, Router, extract::State, http::StatusCode}; use futures::TryFutureExt; +use jacquard::cowstr::ToCowStr; use jacquard::types::ident::AtIdentifier; use jacquard::{ IntoStatic, @@ -21,6 +22,7 @@ use jacquard_common::{ }, xrpc::{GenericXrpcError, XrpcError}, }; +use miette::IntoDiagnostic; use serde::{Deserialize, Serialize}; use smol_str::ToSmolStr; use std::{fmt::Display, sync::Arc}; @@ -76,17 +78,19 @@ pub async fn handle_get_record( .await .map_err(|e| bad_request(GetRecord::NSID, e))?; - let db_key = keys::record_key(&did, req.collection.as_str(), req.rkey.0.as_str()); + let partition = db + .record_partition(req.collection.as_str()) + .map_err(|e| internal_error(GetRecord::NSID, e))?; + + let db_key = keys::record_key(&did, req.rkey.0.as_str()); // 2 args - let cid_bytes = Db::get(db.records.clone(), db_key) + let cid_bytes = Db::get(partition, db_key) .await .map_err(|e| internal_error(GetRecord::NSID, e))?; if let Some(cid_bytes) = cid_bytes { - let cid_str = - std::str::from_utf8(&cid_bytes).map_err(|e| internal_error(GetRecord::NSID, e))?; - - let block_bytes = Db::get(db.blocks.clone(), keys::block_key(cid_str)) + // lookup block using binary cid + let block_bytes = Db::get(db.blocks.clone(), &cid_bytes) .await .map_err(|e| internal_error(GetRecord::NSID, e))? .ok_or_else(|| internal_error(GetRecord::NSID, "not found"))?; @@ -94,7 +98,9 @@ pub async fn handle_get_record( let value: Data = serde_ipld_dagcbor::from_slice(&block_bytes) .map_err(|e| internal_error(GetRecord::NSID, e))?; - let cid = Cid::new(cid_str.as_bytes()).unwrap().into_static(); + let cid = Cid::new(&cid_bytes) + .map_err(|e| internal_error(GetRecord::NSID, e))? + .into_static(); Ok(Json(GetRecordOutput { uri: AtUri::from_parts_owned( @@ -103,7 +109,7 @@ pub async fn handle_get_record( req.rkey.0.as_str(), ) .unwrap(), - cid: Some(cid), + cid: Some(Cid::Str(cid.to_cowstr()).into_static()), value: value.into_static(), extra_data: Default::default(), })) @@ -126,17 +132,16 @@ pub async fn handle_list_records( .await .map_err(|e| bad_request(ListRecords::NSID, e))?; - let prefix = format!( - "{}{}{}{}", - TrimmedDid::from(&did), - keys::SEP as char, - req.collection.as_str(), - keys::SEP as char - ); + let ks = db + .record_partition(req.collection.as_str()) + .map_err(|e| internal_error(ListRecords::NSID, e))?; + + let mut prefix = Vec::new(); + TrimmedDid::from(&did).write_to_vec(&mut prefix); + prefix.push(keys::SEP); let limit = req.limit.unwrap_or(50).min(100) as usize; let reverse = req.reverse.unwrap_or(false); - let ks = db.records.clone(); let blocks_ks = db.blocks.clone(); let did_str = smol_str::SmolStr::from(did.as_str()); @@ -146,107 +151,78 @@ pub async fn handle_list_records( let mut results = Vec::new(); let mut cursor = None; - let mut end_prefix = prefix.clone().into_bytes(); - if let Some(last) = end_prefix.last_mut() { - *last += 1; - } + let iter: Box> = if !reverse { + let mut end_prefix = prefix.clone(); + if let Some(last) = end_prefix.last_mut() { + *last += 1; + } - if !reverse { let end_key = if let Some(cursor) = &req.cursor { - format!("{}{}", prefix, cursor).into_bytes() + let mut k = prefix.clone(); + k.extend_from_slice(cursor.as_bytes()); + k } else { - end_prefix.clone() + end_prefix }; - for item in ks.range(prefix.as_bytes()..end_key.as_slice()).rev() { - let (key, cid_bytes) = item.into_inner().ok()?; - - if !key.starts_with(prefix.as_bytes()) { - break; - } - if results.len() >= limit { - let key_str = String::from_utf8_lossy(&key); - if let Some(last_part) = key_str.split(keys::SEP as char).last() { - cursor = Some(smol_str::SmolStr::from(last_part)); - } - break; - } - - let key_str = String::from_utf8_lossy(&key); - let parts: Vec<&str> = key_str.split(keys::SEP as char).collect(); - if parts.len() == 3 { - let rkey = parts[2]; - let cid_str = std::str::from_utf8(&cid_bytes).ok()?; - - if let Ok(Some(block_bytes)) = blocks_ks.get(keys::block_key(cid_str)) { - let val: Data = - serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); - let cid = Cid::new(cid_str.as_bytes()).unwrap().into_static(); - results.push(RepoRecord { - uri: AtUri::from_parts_owned( - did_str.as_str(), - collection_str.as_str(), - rkey, - ) - .unwrap(), - cid, - value: val.into_static(), - extra_data: Default::default(), - }); - } - } - } + Box::new(ks.range(prefix.as_slice()..end_key.as_slice()).rev()) } else { let start_key = if let Some(cursor) = &req.cursor { - format!("{}{}\0", prefix, cursor).into_bytes() + let mut k = prefix.clone(); + k.extend_from_slice(cursor.as_bytes()); + k.push(0); + k } else { - prefix.clone().into_bytes() + prefix.clone() }; - for item in ks.range(start_key.as_slice()..) { - let (key, cid_bytes) = item.into_inner().ok()?; + Box::new(ks.range(start_key.as_slice()..)) + }; - if !key.starts_with(prefix.as_bytes()) { - break; - } - if results.len() >= limit { - let key_str = String::from_utf8_lossy(&key); - if let Some(last_part) = key_str.split(keys::SEP as char).last() { - cursor = Some(smol_str::SmolStr::from(last_part)); - } - break; - } + for item in iter { + let (key, cid_bytes) = item.into_inner().into_diagnostic()?; + if !key.starts_with(prefix.as_slice()) { + break; + } + if results.len() >= limit { let key_str = String::from_utf8_lossy(&key); - let parts: Vec<&str> = key_str.split(keys::SEP as char).collect(); - if parts.len() == 3 { - let rkey = parts[2]; - let cid_str = std::str::from_utf8(&cid_bytes).ok()?; - - if let Ok(Some(block_bytes)) = blocks_ks.get(keys::block_key(cid_str)) { - let val: Data = - serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); - let cid = Cid::new(cid_str.as_bytes()).unwrap().into_static(); - results.push(RepoRecord { - uri: AtUri::from_parts_owned( - did_str.as_str(), - collection_str.as_str(), - rkey, - ) - .unwrap(), - cid, - value: val.into_static(), - extra_data: Default::default(), - }); - } + if let Some(last_part) = key_str.split(keys::SEP as char).last() { + cursor = Some(smol_str::SmolStr::from(last_part)); + } + break; + } + + // key: {TrimmedDid}|{RKey} + let key_str = String::from_utf8_lossy(&key); + let parts: Vec<&str> = key_str.split(keys::SEP as char).collect(); + if parts.len() == 2 { + let rkey = parts[1]; + // look up using binary cid bytes from the record + if let Ok(Some(block_bytes)) = blocks_ks.get(&cid_bytes) { + let val: Data = + serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); + let cid = + Cid::Str(Cid::new(&cid_bytes).into_diagnostic()?.to_cowstr()).into_static(); + results.push(RepoRecord { + uri: AtUri::from_parts_owned( + did_str.as_str(), + collection_str.as_str(), + rkey, + ) + .into_diagnostic()?, + cid, + value: val.into_static(), + extra_data: Default::default(), + }); } } } - Some((results, cursor)) + Result::<_, miette::Report>::Ok((results, cursor)) }) .await .map_err(|e| internal_error(ListRecords::NSID, e))? - .ok_or_else(|| internal_error(ListRecords::NSID, "not found"))?; + .map_err(|e| internal_error(ListRecords::NSID, e))?; Ok(Json(ListRecordsOutput { records: results, diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 4be2053..162cd29 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -11,7 +11,7 @@ use jacquard::{CowStr, IntoStatic, prelude::*}; use jacquard_common::xrpc::XrpcError; use jacquard_repo::mst::Mst; use jacquard_repo::{BlockStore, MemoryBlockStore}; -use miette::{Context, IntoDiagnostic, Result}; +use miette::{IntoDiagnostic, Result}; use smol_str::{SmolStr, ToSmolStr}; use std::collections::HashMap; use std::sync::Arc; @@ -369,25 +369,32 @@ impl BackfillWorker { let mut batch = app_state.db.inner.batch(); let store = mst.storage(); - // pre-load existing record CIDs for this DID to detect duplicates/updates let prefix = keys::record_prefix(&did); - let prefix_len = prefix.len(); let mut existing_cids: HashMap<(SmolStr, SmolStr), SmolStr> = HashMap::new(); - for guard in app_state.db.records.prefix(&prefix) { - let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - // extract path (collection/rkey) from key by skipping the DID prefix - let mut path_split = key[prefix_len..].split(|b| *b == keys::SEP); - let collection = - std::str::from_utf8(path_split.next().wrap_err("collection not found")?) - .into_diagnostic()? - .to_smolstr(); - let rkey = std::str::from_utf8(path_split.next().wrap_err("rkey not found")?) - .into_diagnostic()? - .to_smolstr(); - let cid = std::str::from_utf8(&cid_bytes) - .into_diagnostic()? - .to_smolstr(); - existing_cids.insert((collection, rkey), cid); + + let mut partitions = Vec::new(); + app_state.db.record_partitions.iter_sync(|col, ks| { + partitions.push((col.clone(), ks.clone())); + true + }); + + for (col_name, ks) in partitions { + for guard in ks.prefix(&prefix) { + let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; + // key: {DID}|{RKey} + let key_str = std::str::from_utf8(&key).into_diagnostic()?; + let parts: Vec<&str> = key_str.split(keys::SEP as char).collect(); + if parts.len() == 2 { + let rkey = parts[1].to_smolstr(); + let cid = if let Ok(c) = cid::Cid::read_bytes(cid_bytes.as_ref()) { + c.to_string().to_smolstr() + } else { + continue; + }; + + existing_cids.insert((col_name.as_str().into(), rkey), cid); + } + } } for (key, cid) in leaves { @@ -398,12 +405,13 @@ impl BackfillWorker { if let Some(val) = val_bytes { let (collection, rkey) = ops::parse_path(&key)?; let path = (collection.to_smolstr(), rkey.to_smolstr()); - let cid = Cid::ipld(cid); + let cid_obj = Cid::ipld(cid); + let partition = app_state.db.record_partition(collection)?; // check if this record already exists with same CID let (action, is_new) = if let Some(existing_cid) = existing_cids.remove(&path) { - if existing_cid == cid.as_str() { + if existing_cid == cid_obj.as_str() { debug!("skip {did}/{collection}/{rkey} ({cid})"); continue; // skip unchanged record } @@ -413,14 +421,11 @@ impl BackfillWorker { }; debug!("{action} {did}/{collection}/{rkey} ({cid})"); - let db_key = keys::record_key(&did, &collection, &rkey); + // Key is just did|rkey + let db_key = keys::record_key(&did, &rkey); - batch.insert( - &app_state.db.blocks, - keys::block_key(cid.as_str()), - val.as_ref(), - ); - batch.insert(&app_state.db.records, db_key, cid.as_str().as_bytes()); + batch.insert(&app_state.db.blocks, cid.to_bytes(), val.as_ref()); + batch.insert(&partition, db_key, cid.to_bytes()); added_blocks += 1; if is_new { @@ -436,7 +441,7 @@ impl BackfillWorker { collection: CowStr::Borrowed(collection), rkey: DbRkey::new(rkey), action: DbAction::from(action), - cid: Some(cid.to_ipld().expect("valid cid")), + cid: Some(cid_obj.to_ipld().expect("valid cid")), }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); @@ -448,10 +453,9 @@ impl BackfillWorker { // remove any remaining existing records (they weren't in the new MST) for ((collection, rkey), cid) in existing_cids { debug!("remove {did}/{collection}/{rkey} ({cid})"); - batch.remove( - &app_state.db.records, - keys::record_key(&did, &collection, &rkey), - ); + let partition = app_state.db.record_partition(collection.as_str())?; + + batch.remove(&partition, keys::record_key(&did, &rkey)); let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { diff --git a/src/db/keys.rs b/src/db/keys.rs index b44d1ac..d7c03c0 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -14,19 +14,17 @@ pub fn repo_key<'a>(did: &'a Did) -> Vec { vec } -// key format: {DID}\x00{Collection}\x00{RKey} (DID trimmed) -pub fn record_key(did: &Did, collection: &str, rkey: &str) -> Vec { +// key format: {DID}\x00{RKey} (DID trimmed) +pub fn record_key(did: &Did, rkey: &str) -> Vec { let repo = TrimmedDid::from(did); - let mut key = Vec::with_capacity(repo.len() + collection.len() + rkey.len() + 2); + let mut key = Vec::with_capacity(repo.len() + rkey.len() + 1); repo.write_to_vec(&mut key); key.push(SEP); - key.extend_from_slice(collection.as_bytes()); - key.push(SEP); key.extend_from_slice(rkey.as_bytes()); key } -// prefix format: {DID}\x00 (DID trimmed) - for scanning all records of a DID +// prefix format: {DID}\x00 (DID trimmed) - for scanning all records of a DID within a collection pub fn record_prefix(did: &Did) -> Vec { let repo = TrimmedDid::from(did); let mut prefix = Vec::with_capacity(repo.len() + 1); @@ -42,11 +40,6 @@ pub fn event_key(seq: u64) -> [u8; 8] { seq.to_be_bytes() } -// key format: {CID} -pub fn block_key(cid: &str) -> &[u8] { - cid.as_bytes() -} - pub const COUNT_KS_PREFIX: &[u8] = &[b'k', SEP]; // count keys for the counts keyspace diff --git a/src/db/mod.rs b/src/db/mod.rs index b874d13..d920972 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -4,7 +4,7 @@ use jacquard::IntoStatic; use jacquard_common::types::string::Did; use miette::{Context, IntoDiagnostic, Result}; use scc::HashMap; -use smol_str::SmolStr; +use smol_str::{SmolStr, format_smolstr}; use std::path::Path; use std::sync::Arc; @@ -15,10 +15,20 @@ use std::sync::atomic::AtomicU64; use tokio::sync::broadcast; use tracing::error; +pub const RECORDS_PARTITION_PREFIX: &str = "r:"; + +fn default_opts() -> KeyspaceCreateOptions { + KeyspaceCreateOptions::default() +} + +fn record_partition_opts() -> KeyspaceCreateOptions { + default_opts().max_memtable_size(32 * 2_u64.pow(20)) +} + pub struct Db { pub inner: Arc, pub repos: Keyspace, - pub records: Keyspace, + pub record_partitions: HashMap, pub blocks: Keyspace, pub cursors: Keyspace, pub pending: Keyspace, @@ -48,13 +58,12 @@ impl Db { .into_diagnostic()?; let db = Arc::new(db); - let opts = KeyspaceCreateOptions::default; + let opts = default_opts; let open_ks = |name: &str, opts: KeyspaceCreateOptions| { db.keyspace(name, move || opts).into_diagnostic() }; let repos = open_ks("repos", opts().expect_point_read_hits(true))?; - let records = open_ks("records", opts().max_memtable_size(32 * 2_u64.pow(20)))?; let blocks = open_ks( "blocks", opts() @@ -68,6 +77,19 @@ impl Db { let events = open_ks("events", opts())?; let counts = open_ks("counts", opts().expect_point_read_hits(true))?; + let record_partitions = HashMap::new(); + { + let names = db.list_keyspace_names(); + for name in names { + let name_str: &str = name.as_ref(); + if let Some(collection) = name_str.strip_prefix(RECORDS_PARTITION_PREFIX) { + let popts = record_partition_opts(); + let ks = db.keyspace(name_str, move || popts).into_diagnostic()?; + let _ = record_partitions.insert_sync(collection.to_string(), ks); + } + } + } + let mut last_id = 0; if let Some(guard) = events.iter().next_back() { let k = guard.key().into_diagnostic()?; @@ -97,7 +119,7 @@ impl Db { Ok(Self { inner: db, repos, - records, + record_partitions, blocks, cursors, pending, @@ -110,6 +132,21 @@ impl Db { }) } + pub fn record_partition(&self, collection: &str) -> Result { + use scc::hash_map::Entry; + match self.record_partitions.entry_sync(collection.to_string()) { + Entry::Occupied(o) => Ok(o.get().clone()), + Entry::Vacant(v) => { + let name = format_smolstr!("{}{}", RECORDS_PARTITION_PREFIX, collection); + let ks = self + .inner + .keyspace(&name, record_partition_opts) + .into_diagnostic()?; + Ok(v.insert_entry(ks).get().clone()) + } + } + } + pub fn persist(&self) -> Result<()> { self.inner.persist(PersistMode::SyncAll).into_diagnostic()?; Ok(()) diff --git a/src/ops.rs b/src/ops.rs index 981c042..b1314da 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -8,7 +8,7 @@ use crate::types::{ use fjall::OwnedWriteBatch; use jacquard::CowStr; use jacquard::IntoStatic; -use jacquard::cowstr::ToCowStr; + use jacquard::types::cid::Cid; use jacquard_api::com_atproto::sync::subscribe_repos::Commit; use jacquard_common::types::crypto::PublicKey; @@ -67,12 +67,19 @@ pub fn delete_repo<'batch>( batch.remove(&db.pending, &repo_key); batch.remove(&db.resync, &repo_key); - // 2. delete from records (prefix: repo_key + SEP) - let mut records_prefix = repo_key.clone(); - records_prefix.push(keys::SEP); - for guard in db.records.prefix(&records_prefix) { - let k = guard.key().into_diagnostic()?; - batch.remove(&db.records, k); + // 2. delete from records (all partitions) + let mut partitions = Vec::new(); + db.record_partitions.iter_sync(|_, v| { + partitions.push(v.clone()); + true + }); + + let records_prefix = keys::record_prefix(did); + for ks in partitions { + for guard in ks.prefix(&records_prefix) { + let k = guard.key().into_diagnostic()?; + batch.remove(&ks, k); + } } // 3. reset collection counts @@ -218,11 +225,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( // store all blocks in the CAS for (cid, bytes) in &parsed.blocks { - batch.insert( - &db.blocks, - keys::block_key(&cid.to_cowstr()), - bytes.to_vec(), - ); + batch.insert(&db.blocks, cid.to_bytes(), bytes.to_vec()); } // 2. iterate ops and update records index @@ -231,7 +234,8 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( for op in &commit.ops { let (collection, rkey) = parse_path(&op.path)?; - let db_key = keys::record_key(did, collection, rkey); + let partition = db.record_partition(collection)?; + let db_key = keys::record_key(did, rkey); // removed collection arg let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); @@ -240,8 +244,14 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( let Some(cid) = &op.cid else { continue; }; - let s = smol_str::SmolStr::from(cid.as_str()); - batch.insert(&db.records, db_key, s.as_bytes().to_vec()); + batch.insert( + &partition, + db_key.clone(), + cid.to_ipld() + .into_diagnostic() + .wrap_err("expected valid cid from relay")? + .to_bytes(), + ); // accumulate counts if op.action.as_str() == "create" { @@ -250,7 +260,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( } } "delete" => { - batch.remove(&db.records, db_key); + batch.remove(&partition, db_key); // accumulate counts records_delta -= 1; @@ -268,7 +278,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( collection: CowStr::Borrowed(collection), rkey: DbRkey::new(rkey), action: DbAction::from(op.action.as_str()), - cid: op.cid.as_ref().map(|c| c.0.to_ipld().expect("valid cid")), + cid: op.cid.as_ref().map(|c| c.to_ipld().expect("valid cid")), }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; -- 2.51.2