diff --git a/src/api/xrpc.rs b/src/api/xrpc.rs index 3900bec..354cc86 100644 --- a/src/api/xrpc.rs +++ b/src/api/xrpc.rs @@ -180,25 +180,13 @@ pub async fn handle_list_records( if !key.starts_with(prefix.as_slice()) { break; } + + let rkey = keys::parse_rkey(&key[prefix.len()..])?; 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)); - } + cursor = Some(rkey); break; } - // manual deserialization of the key suffix since it is raw bytes, not msgpack - let suffix = &key[prefix.len()..]; - let rkey = if suffix.len() == 8 { - let mut bytes = [0u8; 8]; - bytes.copy_from_slice(suffix); - DbRkey::Tid(crate::db::types::DbTid::new_from_bytes(bytes)) - } else { - let s = String::from_utf8_lossy(suffix); - DbRkey::Str(smol_str::SmolStr::from(s.as_ref())) - }; - // 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); @@ -225,7 +213,7 @@ pub async fn handle_list_records( Ok(Json(ListRecordsOutput { records: results, - cursor: cursor.map(|c| c.into()), + cursor: cursor.map(|c| jacquard::CowStr::Owned(c.to_smolstr())), extra_data: Default::default(), })) } diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 1d32eef..ecfa9fe 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -506,14 +506,11 @@ async fn process_did<'i>( for (col_name, ks) in partitions { for guard in ks.prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - let rkey: DbRkey = - rmp_serde::from_slice(&key[prefix.len()..]).into_diagnostic()?; - let cid = if let Ok(c) = cid::Cid::read_bytes(cid_bytes.as_ref()) { - c.to_string().to_smolstr() - } else { - error!("invalid cid for {did}: {cid_bytes:?}"); - continue; - }; + let rkey = keys::parse_rkey(&key[prefix.len()..]) + .map_err(|e| miette::miette!("invalid rkey '{key:?}' for {did}: {e}"))?; + let cid = cid::Cid::read_bytes(cid_bytes.as_ref()) + .map_err(|e| miette::miette!("invalid cid '{cid_bytes:?}' for {did}: {e}"))? + .to_smolstr(); existing_cids.insert((col_name.as_str().into(), rkey), cid); } diff --git a/src/db/keys.rs b/src/db/keys.rs index ea6e6f9..916efef 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -1,4 +1,5 @@ use jacquard_common::types::string::Did; +use smol_str::SmolStr; use crate::db::types::{DbRkey, DbTid, TrimmedDid}; @@ -7,14 +8,14 @@ pub const SEP: u8 = b'|'; pub const CURSOR_KEY: &[u8] = b"firehose_cursor"; -// Key format: {DID} (trimmed) +// Key format: {DID} pub fn repo_key<'a>(did: &'a Did) -> Vec { let mut vec = Vec::with_capacity(32); TrimmedDid::from(did).write_to_vec(&mut vec); vec } -// prefix format: {DID}\x00 (DID trimmed) - for scanning all records of a DID within a collection +// prefix format: {DID}\x00 pub fn record_prefix(did: &Did) -> Vec { let repo = TrimmedDid::from(did); let mut prefix = Vec::with_capacity(repo.len() + 1); @@ -23,16 +24,48 @@ pub fn record_prefix(did: &Did) -> Vec { prefix } -// key format: {DID}\x00{rkey} (DID trimmed) +// key format: {DID}\x00{rkey} pub fn record_key(did: &Did, rkey: &DbRkey) -> Vec { let repo = TrimmedDid::from(did); let mut key = Vec::with_capacity(repo.len() + rkey.len() + 1); repo.write_to_vec(&mut key); key.push(SEP); - key.extend_from_slice(rkey.as_bytes()); + write_rkey(&mut key, rkey); key } +pub fn write_rkey(buf: &mut Vec, rkey: &DbRkey) { + match rkey { + DbRkey::Tid(tid) => { + buf.push(b't'); + buf.extend_from_slice(tid.as_bytes()); + } + DbRkey::Str(s) => { + buf.push(b's'); + buf.extend_from_slice(s.as_bytes()); + } + } +} + +pub fn parse_rkey(raw: &[u8]) -> miette::Result { + let Some(kind) = raw.first() else { + miette::bail!("record key is empty"); + }; + let rkey = match kind { + b't' => { + DbRkey::Tid(DbTid::new_from_bytes(raw[1..].try_into().map_err(|e| { + miette::miette!("record key '{raw:?}' is invalid: {e}") + })?)) + } + b's' => DbRkey::Str(SmolStr::new( + std::str::from_utf8(&raw[1..]) + .map_err(|e| miette::miette!("record key '{raw:?}' is invalid: {e}"))?, + )), + _ => miette::bail!("invalid record key kind: {}", *kind as char), + }; + Ok(rkey) +} + // key format: {SEQ} pub fn event_key(seq: u64) -> [u8; 8] { seq.to_be_bytes()