From d19a50fd9f0194131aa8d427bd5e3c2b5b179077 Mon Sep 17 00:00:00 2001 From: Mia Date: Wed, 3 Sep 2025 18:43:04 +0100 Subject: [PATCH] feat(lexica): move StrongRef and Blob into lexica --- Cargo.lock | 1 + consumer/src/backfill/mod.rs | 15 +++---------- consumer/src/db/copy.rs | 5 +++-- consumer/src/db/record.rs | 10 ++++----- consumer/src/indexer/records.rs | 23 +------------------ consumer/src/utils.rs | 39 +++++---------------------------- lexica/Cargo.toml | 1 + lexica/src/lib.rs | 28 ++++++++++++++++++++++- lexica/src/utils.rs | 31 ++++++++++++++++++++++++++ 9 files changed, 77 insertions(+), 76 deletions(-) create mode 100644 lexica/src/utils.rs diff --git a/Cargo.lock b/Cargo.lock index c25506ea..93ed98b2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2276,6 +2276,7 @@ name = "lexica" version = "0.1.0" dependencies = [ "chrono", + "cid", "serde", "serde_json", ] diff --git a/consumer/src/backfill/mod.rs b/consumer/src/backfill/mod.rs index fe73ebb1..fd4a0d2b 100644 --- a/consumer/src/backfill/mod.rs +++ b/consumer/src/backfill/mod.rs @@ -6,6 +6,7 @@ use chrono::prelude::*; use deadpool_postgres::{Object, Pool, Transaction}; use did_resolver::Resolver; use ipld_core::cid::Cid; +use lexica::StrongRef; use metrics::counter; use parakeet_db::types::{ActorStatus, ActorSyncState}; use redis::aio::MultiplexedConnection; @@ -267,19 +268,9 @@ async fn check_pds_repo_status( #[derive(Debug, Default)] struct CopyStore { - likes: Vec<( - String, - records::StrongRef, - Option, - DateTime, - )>, + likes: Vec<(String, StrongRef, Option, DateTime)>, posts: Vec<(String, Cid, records::AppBskyFeedPost)>, - reposts: Vec<( - String, - records::StrongRef, - Option, - DateTime, - )>, + reposts: Vec<(String, StrongRef, Option, DateTime)>, blocks: Vec<(String, String, DateTime)>, follows: Vec<(String, String, DateTime)>, list_items: Vec<(String, records::AppBskyGraphListItem)>, diff --git a/consumer/src/db/copy.rs b/consumer/src/db/copy.rs index 1b622991..17db6f0f 100644 --- a/consumer/src/db/copy.rs +++ b/consumer/src/db/copy.rs @@ -7,6 +7,7 @@ use futures::pin_mut; use ipld_core::cid::Cid; use tokio_postgres::binary_copy::BinaryCopyInWriter; use tokio_postgres::types::Type; +use lexica::StrongRef; // StrongRefs are used in both likes and reposts const STRONGREF_TYPES: &[Type] = &[ @@ -20,8 +21,8 @@ const STRONGREF_TYPES: &[Type] = &[ ]; type StrongRefRow = ( String, - records::StrongRef, - Option, + StrongRef, + Option, DateTime, ); diff --git a/consumer/src/db/record.rs b/consumer/src/db/record.rs index d739b0c4..6d9b632b 100644 --- a/consumer/src/db/record.rs +++ b/consumer/src/db/record.rs @@ -371,7 +371,7 @@ async fn post_embed_image_insert( let stmt = conn.prepare("INSERT INTO post_embed_images (post_uri, seq, cid, mime_type, alt, width, height) VALUES ($1, $2, $3, $4, $5, $6, $7)").await?; for (idx, image) in embed.images.iter().enumerate() { - let cid = image.image.r#ref.to_string(); + let cid = image.image.cid.to_string(); let width = image.aspect_ratio.as_ref().map(|v| v.width); let height = image.aspect_ratio.as_ref().map(|v| v.height); @@ -398,7 +398,7 @@ async fn post_embed_video_insert( post: &str, embed: AppBskyEmbedVideo, ) -> PgExecResult { - let cid = embed.video.r#ref.to_string(); + let cid = embed.video.cid.to_string(); let width = embed.aspect_ratio.as_ref().map(|v| v.width); let height = embed.aspect_ratio.as_ref().map(|v| v.height); @@ -411,7 +411,7 @@ async fn post_embed_video_insert( let stmt = conn.prepare_cached("INSERT INTO post_embed_video_captions (post_uri, cid, mime_type, language) VALUES ($1, $2, $3, $4)").await?; for caption in captions { - let cid = caption.file.r#ref.to_string(); + let cid = caption.file.cid.to_string(); conn.execute( &stmt, &[&post, &cid, &caption.file.mime_type, &caption.lang], @@ -429,7 +429,7 @@ async fn post_embed_external_insert( embed: AppBskyEmbedExternal, ) -> PgExecResult { let thumb_mime = embed.external.thumb.as_ref().map(|v| v.mime_type.clone()); - let thumb_cid = embed.external.thumb.as_ref().map(|v| v.r#ref.to_string()); + let thumb_cid = embed.external.thumb.as_ref().map(|v| v.cid.to_string()); conn.execute( "INSERT INTO post_embed_ext (post_uri, uri, title, description, thumb_mime_type, thumb_cid) VALUES ($1, $2, $3, $4, $5, $6)", @@ -634,7 +634,7 @@ pub async fn status_upsert( let record = serde_json::to_value(&rec).unwrap(); let thumb = rec.embed.as_ref().and_then(|v| v.external.thumb.clone()); let thumb_mime = thumb.as_ref().map(|v| v.mime_type.clone()); - let thumb_cid = thumb.as_ref().map(|v| v.r#ref.to_string()); + let thumb_cid = thumb.as_ref().map(|v| v.cid.to_string()); conn.execute( include_str!("sql/status_upsert.sql"), diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs index d3729831..7d8f983a 100644 --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -1,36 +1,15 @@ use crate::utils; use chrono::{DateTime, Utc}; -use ipld_core::cid::Cid; use lexica::app_bsky::actor::{ChatAllowIncoming, ProfileAllowSubscriptions, Status}; use lexica::app_bsky::embed::AspectRatio; use lexica::app_bsky::labeler::LabelerPolicy; use lexica::app_bsky::richtext::FacetMain; use lexica::com_atproto::label::SelfLabels; use lexica::com_atproto::moderation::{ReasonType, SubjectType}; +use lexica::{Blob, StrongRef}; use serde::{Deserialize, Serialize}; use serde_with::serde_as; -#[derive(Clone, Debug, Deserialize, Serialize)] -pub struct StrongRef { - #[serde( - deserialize_with = "utils::cid_from_string", - serialize_with = "utils::cid_as_str" - )] - pub cid: Cid, - pub uri: String, -} - -#[derive(Clone, Debug, Deserialize, Serialize)] -#[serde(tag = "$type")] -#[serde(rename = "blob")] -#[serde(rename_all = "camelCase")] -pub struct Blob { - pub mime_type: String, - #[serde(serialize_with = "utils::cid_as_link")] - pub r#ref: Cid, - pub size: i32, -} - #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] #[serde_as] diff --git a/consumer/src/utils.rs b/consumer/src/utils.rs index 88825ba2..cd74eea6 100644 --- a/consumer/src/utils.rs +++ b/consumer/src/utils.rs @@ -1,5 +1,5 @@ -use ipld_core::cid::Cid; -use serde::{Deserialize, Deserializer, Serialize, Serializer}; +use serde::{Deserialize, Deserializer}; +use lexica::{Blob, StrongRef}; // see https://deer.social/profile/did:plc:63y3oh7iakdueqhlj6trojbq/post/3ltuv4skhqs2h pub fn safe_string<'de, D: Deserializer<'de>>(deserializer: D) -> Result { @@ -8,41 +8,12 @@ pub fn safe_string<'de, D: Deserializer<'de>>(deserializer: D) -> Result>(deserializer: D) -> Result { - let str = String::deserialize(deserializer)?; - - Cid::try_from(str).map_err(serde::de::Error::custom) -} - -pub fn cid_as_str(inp: &Cid, serializer: S) -> Result -where - S: Serializer, -{ - inp.to_string().serialize(serializer) -} - -#[derive(Debug, Deserialize, Serialize)] -pub struct LinkRef { - #[serde(rename = "$link")] - link: String, -} - -pub fn cid_as_link(inp: &Cid, serializer: S) -> Result -where - S: Serializer, -{ - LinkRef { - link: inp.to_string(), - } - .serialize(serializer) -} - -pub fn blob_ref(blob: Option) -> Option { - blob.map(|blob| blob.r#ref.to_string()) +pub fn blob_ref(blob: Option) -> Option { + blob.map(|blob| blob.cid.to_string()) } pub fn strongref_to_parts( - strongref: Option<&crate::indexer::records::StrongRef>, + strongref: Option<&StrongRef>, ) -> (Option, Option) { strongref .map(|sr| (sr.uri.clone(), sr.cid.to_string())) diff --git a/lexica/Cargo.toml b/lexica/Cargo.toml index 1eed6c36..41e7edef 100644 --- a/lexica/Cargo.toml +++ b/lexica/Cargo.toml @@ -5,5 +5,6 @@ edition = "2021" [dependencies] chrono = { version = "0.4.39", features = ["serde"] } +cid = { version = "0.11", features = ["serde"] } serde = { version = "1.0.216", features = ["derive"] } serde_json = "1.0.134" diff --git a/lexica/src/lib.rs b/lexica/src/lib.rs index 5242c85a..7c1a9188 100644 --- a/lexica/src/lib.rs +++ b/lexica/src/lib.rs @@ -1,10 +1,36 @@ -use serde::Serialize; +use cid::Cid; +use serde::{Deserialize, Serialize}; + +pub use utils::LinkRef; pub mod app_bsky; pub mod com_atproto; +mod utils; #[derive(Clone, Debug, Serialize)] pub struct JsonBytes { #[serde(rename = "$bytes")] pub bytes: String, } + +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct StrongRef { + #[serde( + deserialize_with = "utils::cid_from_string", + serialize_with = "utils::cid_as_str" + )] + pub cid: Cid, + pub uri: String, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(tag = "$type")] +#[serde(rename = "blob")] +#[serde(rename_all = "camelCase")] +pub struct Blob { + pub mime_type: String, + #[serde(rename = "ref")] + #[serde(serialize_with = "utils::cid_as_link")] + pub cid: Cid, + pub size: i32, +} diff --git a/lexica/src/utils.rs b/lexica/src/utils.rs new file mode 100644 index 00000000..a78932b8 --- /dev/null +++ b/lexica/src/utils.rs @@ -0,0 +1,31 @@ +use cid::Cid; +use serde::{Deserialize, Deserializer, Serialize, Serializer}; + +pub fn cid_from_string<'de, D: Deserializer<'de>>(deserializer: D) -> Result { + let str = String::deserialize(deserializer)?; + + Cid::try_from(str).map_err(serde::de::Error::custom) +} + +pub fn cid_as_str(inp: &Cid, serializer: S) -> Result +where + S: Serializer, +{ + inp.to_string().serialize(serializer) +} + +#[derive(Debug, Deserialize, Serialize)] +pub struct LinkRef { + #[serde(rename = "$link")] + link: String, +} + +pub fn cid_as_link(inp: &Cid, serializer: S) -> Result +where + S: Serializer, +{ + LinkRef { + link: inp.to_string(), + } + .serialize(serializer) +} \ No newline at end of file -- 2.51.2