From 28cec531213aa0ab405d3fb20b66df5e326e0f59 Mon Sep 17 00:00:00 2001 From: Thomas Karpiniec Date: Sun, 25 Jan 2026 02:19:01 +1100 Subject: [PATCH] Improve record slice handling to avoid needing self_cell --- Cargo.lock | 7 --- tapped/Cargo.toml | 1 - tapped/src/channel.rs | 37 ++++---------- tapped/src/types.rs | 109 ++++++++++++++++++------------------------ 4 files changed, 55 insertions(+), 99 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7ec7ec8..20ad7c6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2224,12 +2224,6 @@ dependencies = [ "libc", ] -[[package]] -name = "self_cell" -version = "1.2.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b12e76d157a900eb52e81bc6e9f3069344290341720e9178cde2407113ac8d89" - [[package]] name = "semver" version = "1.0.27" @@ -2602,7 +2596,6 @@ dependencies = [ "futures-util", "libc", "reqwest", - "self_cell", "serde", "serde_json", "thiserror 2.0.18", diff --git a/tapped/Cargo.toml b/tapped/Cargo.toml index 6f01483..b822649 100644 --- a/tapped/Cargo.toml +++ b/tapped/Cargo.toml @@ -22,7 +22,6 @@ futures-util = "0.3" tracing = "0.1" base64 = "0.22" libc = "0.2" -self_cell = "1.2.2" [dev-dependencies] tokio = { version = "1", features = ["macros", "rt-multi-thread"] } diff --git a/tapped/src/channel.rs b/tapped/src/channel.rs index a5a5a76..6775d08 100644 --- a/tapped/src/channel.rs +++ b/tapped/src/channel.rs @@ -7,7 +7,7 @@ use tokio_tungstenite::{connect_async, tungstenite::Message}; use tungstenite::protocol::frame::Utf8Bytes; use url::Url; -use crate::types::{RawEvent, RawEventOwned, RawRecordEventOwned, Record, UnparsedRecord}; +use crate::types::RawEvent; use crate::{Error, Event, Result}; type WsStream = @@ -124,35 +124,16 @@ impl EventReceiver { loop { match self.event_rx.recv().await { Some(event_with_ack) => { - let mut raw_event_owned = None::; - let mut raw_record_owned = None::; - let inner_record = Record::new(event_with_ack.event, |json| { - match serde_json::from_str::(json) { - Ok(mut raw) => { - let inner = if let Some(mut rec) = raw.record.take() { - let inner = rec.record.take(); - raw_record_owned = Some(rec.owned); - inner - } else { - None - }; - raw_event_owned = Some(raw.owned); - UnparsedRecord(inner) - } - Err(e) => { - tracing::warn!("Failed to parse event: {}", e); - UnparsedRecord(None) - } + let json = event_with_ack.event; + let raw = match serde_json::from_str::(json.as_str()) { + Ok(raw) => raw, + Err(e) => { + tracing::warn!("Failed to parse event: {}", e); + continue; } - }); - let inner_record = if inner_record.borrow_dependent().0.is_some() { - Some(inner_record) - } else { - None }; - if let Some(event) = - raw_event_owned.and_then(|r| r.into_event(raw_record_owned, inner_record)) - { + + if let Some(event) = raw.into_event(json.clone()) { let id = event.id(); break Ok(ReceivedEvent { event, diff --git a/tapped/src/types.rs b/tapped/src/types.rs index 5c4a500..d0a917d 100644 --- a/tapped/src/types.rs +++ b/tapped/src/types.rs @@ -7,17 +7,6 @@ use tungstenite::protocol::frame::Utf8Bytes; use crate::Error; -self_cell::self_cell!( - pub(crate) struct Record { - owner: Utf8Bytes, - - #[covariant] - dependent: UnparsedRecord, - } -); - -pub(crate) struct UnparsedRecord<'a>(pub(crate) Option<&'a RawValue>); - /// Action performed on a record. #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "lowercase")] @@ -73,17 +62,18 @@ pub struct RecordEvent { pub action: RecordAction, /// CID of the record (None on delete). pub cid: Option, - /// The record data (None on delete). - pub(crate) record: Option, + // Inner record JSON pointing into the outer JSON + json: Option, + record_offset: usize, + record_len: usize, } impl RecordEvent { /// Get the record's content as a reference to a JSON string pub fn record_as_str(&self) -> Option<&str> { - self.record + self.json .as_ref() - .and_then(|r| r.borrow_dependent().0) - .map(|rv| rv.get()) + .map(|j| &j.as_str()[self.record_offset..self.record_offset + self.record_len]) } /// Parse the record's content to a compatible struct @@ -217,32 +207,19 @@ pub struct Service { // Internal deserialisation structures for parsing tap's JSON format -#[derive(Deserialize, Clone)] +#[derive(Deserialize)] #[serde(bound(deserialize = "'de: 'a"))] pub(crate) struct RawEvent<'a> { - #[serde(flatten)] - pub(crate) owned: RawEventOwned, - pub(crate) record: Option>, -} - -#[derive(Deserialize, Clone)] -pub(crate) struct RawEventOwned { - pub(crate) id: u64, + pub id: u64, #[serde(rename = "type")] - pub(crate) type_: String, - pub(crate) identity: Option, + pub type_: String, + pub identity: Option, + pub record: Option>, } -#[derive(Deserialize, Clone)] +#[derive(Deserialize)] #[serde(bound(deserialize = "'de: 'a"))] pub(crate) struct RawRecordEvent<'a> { - #[serde(flatten)] - pub owned: RawRecordEventOwned, - pub record: Option<&'a RawValue>, -} - -#[derive(Deserialize, Clone)] -pub(crate) struct RawRecordEventOwned { pub live: bool, pub did: String, pub rev: String, @@ -250,6 +227,7 @@ pub(crate) struct RawRecordEventOwned { pub rkey: String, pub action: RecordAction, pub cid: Option, + pub record: Option<&'a RawValue>, } #[derive(Deserialize, Clone)] @@ -261,16 +239,20 @@ pub(crate) struct RawIdentityEvent { pub status: AccountStatus, } -impl RawEventOwned { +impl RawEvent<'_> { /// Convert to the public Event type. - pub fn into_event( - self, - raw_record: Option, - inner_record: Option, - ) -> Option { + pub fn into_event(self, json: Utf8Bytes) -> Option { match self.type_.as_str() { "record" => { - let r = raw_record?; + let r = self.record?; + let (json, record_offset, record_len) = if let Some(rv) = r.record.as_ref() { + let json_str = json.as_str(); + let rv_str = rv.get(); + let offset = rv_str.as_ptr() as usize - json_str.as_ptr() as usize; + (Some(json), offset, rv_str.len()) + } else { + (None, 0, 0) + }; Some(Event::Record(RecordEvent { id: self.id, live: r.live, @@ -280,7 +262,9 @@ impl RawEventOwned { rkey: r.rkey, action: r.action, cid: r.cid, - record: inner_record, + json, + record_offset, + record_len, })) } "identity" => { @@ -477,14 +461,12 @@ mod tests { }) .to_string(); - let raw: RawEvent = serde_json::from_str(&json).unwrap(); - assert_eq!(raw.owned.id, 12345); - assert_eq!(raw.owned.type_, "record"); + let json: Utf8Bytes = json.into(); + let raw: RawEvent = serde_json::from_str(json.as_str()).unwrap(); + assert_eq!(raw.id, 12345); + assert_eq!(raw.type_, "record"); - let event = raw - .owned - .into_event(raw.record.map(|r| r.owned), None) - .unwrap(); + let event = raw.into_event(json.clone()).unwrap(); match event { Event::Record(r) => { assert_eq!(r.id, 12345); @@ -499,7 +481,7 @@ mod tests { #[test] fn raw_identity_event_deserialize() { - let json = json!({ + let json: Utf8Bytes = json!({ "id": 99999, "type": "identity", "identity": { @@ -509,10 +491,11 @@ mod tests { "status": "active" } }) - .to_string(); + .to_string() + .into(); - let raw: RawEvent = serde_json::from_str(&json).unwrap(); - let event = raw.owned.into_event(None, None).unwrap(); + let raw: RawEvent = serde_json::from_str(json.as_str()).unwrap(); + let event = raw.into_event(json.clone()).unwrap(); match event { Event::Identity(i) => { @@ -528,7 +511,7 @@ mod tests { #[test] fn raw_delete_event_no_record() { - let json = json!({ + let json: Utf8Bytes = json!({ "id": 55555, "type": "record", "record": { @@ -542,19 +525,17 @@ mod tests { "record": null } }) - .to_string(); + .to_string() + .into(); - let raw: RawEvent = serde_json::from_str(&json).unwrap(); - let event = raw - .owned - .into_event(raw.record.map(|r| r.owned), None) - .unwrap(); + let raw: RawEvent = serde_json::from_str(json.as_str()).unwrap(); + let event = raw.into_event(json.clone()).unwrap(); match event { Event::Record(r) => { assert_eq!(r.action, RecordAction::Delete); assert!(r.cid.is_none()); - assert!(r.record.is_none()); + assert!(r.record_as_str().is_none()); } _ => panic!("Expected Record event"), } @@ -571,7 +552,9 @@ mod tests { rkey: "key".to_string(), action: RecordAction::Create, cid: None, - record: None, + json: None, + record_offset: 0, + record_len: 0, }); assert_eq!(record_event.id(), 123); -- 2.51.2