From 58c6eed75b494da531f4d6a55099c79fdf71e083 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 10 Aug 2026 19:59:03 +0300 Subject: [PATCH] [db,stream] decode rkeys by msgpack marker --- src/control/stream/indexer.rs | 20 +++++++ src/db/types.rs | 106 +++++++++++++++++++++++++++++++++- 2 files changed, 125 insertions(+), 1 deletion(-) diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index b0c5816..a3cd354 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -429,6 +429,26 @@ mod tests { assert!(rec.record.is_some()); } + #[test] + fn persisted_eight_byte_string_rkey_replays_verbatim() { + let (_tmp, state) = test_state(); + let did = jacquard_common::types::string::Did::new(DID).unwrap(); + let stored = StoredEvent { + live: false, + did: TrimmedDid::from(&did), + rev: rev("3kzbif5moe22m"), + collection: CowStr::Borrowed(COL), + rkey: DbRkey::Str("stinkpot".into()), + action: crate::db::types::DbAction::Create, + data: StoredData::Block(bytes::Bytes::from_static(b"\xa0")), + }; + let serialized = rmp_serde::to_vec(&stored).unwrap(); + let deserialized: StoredEvent = rmp_serde::from_slice(&serialized).unwrap(); + + let event = stored_to_event(&state, 153_752, deserialized, None).unwrap(); + assert_eq!(event.record.unwrap().rkey.as_ref(), "stinkpot"); + } + #[test] fn redacted_head_suppresses_queued_inline_body_but_keeps_cid_metadata() { let (_tmp, state) = test_state(); diff --git a/src/db/types.rs b/src/db/types.rs index 5bd1819..edadb93 100644 --- a/src/db/types.rs +++ b/src/db/types.rs @@ -297,7 +297,7 @@ impl Display for DbAction { } } -#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize)] #[serde(untagged)] #[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub enum DbRkey { @@ -305,6 +305,79 @@ pub enum DbRkey { Str(SmolStr), } +impl<'de> Deserialize<'de> for DbRkey { + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + struct DbRkeyVisitor; + + fn tid_from_bytes(value: &[u8]) -> Result + where + E: serde::de::Error, + { + let bytes = value + .try_into() + .map_err(|_| E::invalid_length(value.len(), &"exactly 8 bytes"))?; + Ok(DbRkey::Tid(DbTid::new_from_bytes(bytes))) + } + + impl<'de> serde::de::Visitor<'de> for DbRkeyVisitor { + type Value = DbRkey; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("an rkey string or an 8-byte tid") + } + + fn visit_str(self, value: &str) -> Result + where + E: serde::de::Error, + { + Ok(DbRkey::Str(SmolStr::from(value))) + } + + fn visit_string(self, value: String) -> Result + where + E: serde::de::Error, + { + Ok(DbRkey::Str(SmolStr::from(value))) + } + + fn visit_bytes(self, value: &[u8]) -> Result + where + E: serde::de::Error, + { + tid_from_bytes(value) + } + + fn visit_byte_buf(self, value: Vec) -> Result + where + E: serde::de::Error, + { + tid_from_bytes(&value) + } + + fn visit_seq(self, mut seq: A) -> Result + where + A: serde::de::SeqAccess<'de>, + { + let mut bytes = [0; 8]; + for (index, byte) in bytes.iter_mut().enumerate() { + *byte = seq + .next_element()? + .ok_or_else(|| serde::de::Error::invalid_length(index, &self))?; + } + if seq.next_element::()?.is_some() { + return Err(serde::de::Error::invalid_length(9, &self)); + } + Ok(DbRkey::Tid(DbTid::new_from_bytes(bytes))) + } + } + + deserializer.deserialize_any(DbRkeyVisitor) + } +} + #[cfg_attr(not(feature = "indexer"), allow(dead_code))] impl DbRkey { pub fn new(s: &str) -> Self { @@ -432,4 +505,35 @@ mod tests { panic!("expected str"); } } + + #[test] + fn dbrkey_msgpack_string_marker_preserves_eight_byte_rkey() { + let encoded = rmp_serde::to_vec(&DbRkey::Str("stinkpot".into())).unwrap(); + assert_eq!(encoded, b"\xa8stinkpot"); + + let decoded: DbRkey = rmp_serde::from_slice(&encoded).unwrap(); + assert_eq!(decoded, DbRkey::Str("stinkpot".into())); + assert_eq!(decoded.to_smolstr(), "stinkpot"); + } + + #[test] + fn dbrkey_msgpack_binary_marker_preserves_tid() { + let expected = DbRkey::new("3jzfcijpj2z2a"); + let encoded = rmp_serde::to_vec(&expected).unwrap(); + assert_eq!(encoded.get(..2), Some([0xc4, 8].as_slice())); + + let decoded: DbRkey = rmp_serde::from_slice(&encoded).unwrap(); + assert_eq!(decoded, expected); + assert_eq!(decoded.to_smolstr(), "3jzfcijpj2z2a"); + } + + #[test] + fn dbrkey_json_tid_sequence_stays_compatible() { + let expected = DbRkey::new("3jzfcijpj2z2a"); + let encoded = serde_json::to_vec(&expected).unwrap(); + assert_eq!(encoded.first(), Some(&b'[')); + + let decoded: DbRkey = serde_json::from_slice(&encoded).unwrap(); + assert_eq!(decoded, expected); + } } -- 2.51.2