use super::records; use ipld_core::cid::Cid; use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] #[serde(tag = "$type")] pub enum RecordTypes { #[serde(rename = "app.bsky.actor.profile")] AppBskyActorProfile(records::AppBskyActorProfile), #[serde(rename = "app.bsky.actor.status")] AppBskyActorStatus(records::AppBskyActorStatus), #[serde(rename = "app.bsky.feed.generator")] AppBskyFeedGenerator(records::AppBskyFeedGenerator), #[serde(rename = "app.bsky.feed.like")] AppBskyFeedLike(records::AppBskyFeedLike), #[serde(rename = "app.bsky.feed.post")] AppBskyFeedPost(records::AppBskyFeedPost), #[serde(rename = "app.bsky.feed.postgate")] AppBskyFeedPostgate(records::AppBskyFeedPostgate), #[serde(rename = "app.bsky.feed.repost")] AppBskyFeedRepost(records::AppBskyFeedRepost), #[serde(rename = "app.bsky.feed.threadgate")] AppBskyFeedThreadgate(records::AppBskyFeedThreadgate), #[serde(rename = "app.bsky.graph.block")] AppBskyGraphBlock(records::AppBskyGraphBlock), #[serde(rename = "app.bsky.graph.follow")] AppBskyGraphFollow(records::AppBskyGraphFollow), #[serde(rename = "app.bsky.graph.list")] AppBskyGraphList(records::AppBskyGraphList), #[serde(rename = "app.bsky.graph.listblock")] AppBskyGraphListBlock(records::AppBskyGraphListBlock), #[serde(rename = "app.bsky.graph.listitem")] AppBskyGraphListItem(records::AppBskyGraphListItem), #[serde(rename = "app.bsky.graph.starterpack")] AppBskyGraphStarterPack(records::AppBskyGraphStarterPack), #[serde(rename = "app.bsky.graph.verification")] AppBskyGraphVerification(records::AppBskyGraphVerification), #[serde(rename = "app.bsky.labeler.service")] AppBskyLabelerService(records::AppBskyLabelerService), #[serde(rename = "app.bsky.notification.declaration")] AppBskyNotificationDeclaration(records::AppBskyNotificationDeclaration), #[serde(rename = "chat.bsky.actor.declaration")] ChatBskyActorDeclaration(records::ChatBskyActorDeclaration), #[serde(rename = "community.lexicon.bookmarks.bookmark")] CommunityLexiconBookmark(lexica::community_lexicon::bookmarks::Bookmark) } #[derive(Debug, PartialOrd, PartialEq, Deserialize, Serialize)] pub enum CollectionType { BskyProfile, BskyStatus, BskyFeedGen, BskyFeedLike, BskyFeedPost, BskyFeedPostgate, BskyFeedRepost, BskyFeedThreadgate, BskyBlock, BskyFollow, BskyList, BskyListBlock, BskyListItem, BskyStarterPack, BskyVerification, BskyLabelerService, BskyNotificationDeclaration, ChatActorDecl, CommunityLexiconBookmark, Unsupported, } impl CollectionType { pub(crate) fn from_str(input: &str) -> CollectionType { match input { "app.bsky.actor.profile" => CollectionType::BskyProfile, "app.bsky.actor.status" => CollectionType::BskyStatus, "app.bsky.feed.generator" => CollectionType::BskyFeedGen, "app.bsky.feed.like" => CollectionType::BskyFeedLike, "app.bsky.feed.post" => CollectionType::BskyFeedPost, "app.bsky.feed.postgate" => CollectionType::BskyFeedPostgate, "app.bsky.feed.repost" => CollectionType::BskyFeedRepost, "app.bsky.feed.threadgate" => CollectionType::BskyFeedThreadgate, "app.bsky.graph.block" => CollectionType::BskyBlock, "app.bsky.graph.follow" => CollectionType::BskyFollow, "app.bsky.graph.list" => CollectionType::BskyList, "app.bsky.graph.listblock" => CollectionType::BskyListBlock, "app.bsky.graph.listitem" => CollectionType::BskyListItem, "app.bsky.graph.starterpack" => CollectionType::BskyStarterPack, "app.bsky.graph.verification" => CollectionType::BskyVerification, "app.bsky.labeler.service" => CollectionType::BskyLabelerService, "app.bsky.notification.declaration" => CollectionType::BskyNotificationDeclaration, "chat.bsky.actor.declaration" => CollectionType::ChatActorDecl, "community.lexicon.bookmarks.bookmark" => CollectionType::CommunityLexiconBookmark, _ => CollectionType::Unsupported, } } pub fn can_update(&self) -> bool { match self { CollectionType::BskyProfile => true, CollectionType::BskyStatus => true, CollectionType::BskyFeedGen => true, CollectionType::BskyFeedLike => false, CollectionType::BskyFeedPost => false, CollectionType::BskyFeedPostgate => true, CollectionType::BskyFeedRepost => false, CollectionType::BskyFeedThreadgate => true, CollectionType::BskyBlock => false, CollectionType::BskyFollow => false, CollectionType::BskyList => true, CollectionType::BskyListBlock => false, CollectionType::BskyListItem => false, CollectionType::ChatActorDecl => true, CollectionType::BskyStarterPack => true, CollectionType::BskyVerification => false, CollectionType::BskyLabelerService => true, CollectionType::BskyNotificationDeclaration => true, CollectionType::CommunityLexiconBookmark => true, CollectionType::Unsupported => false, } } } #[derive(Debug, Deserialize, Serialize)] pub struct BackfillItem { pub collection: CollectionType, pub inner: BackfillItemInner, pub at_uri: String, pub cid: Option, } #[derive(Debug, Deserialize, Serialize)] #[serde(tag = "action")] pub enum BackfillItemInner { Create(RecordTypes), Update(RecordTypes), Delete, } pub trait AggregateDeltaStore { async fn add_delta(&mut self, uri: &str, typ: parakeet_index::AggregateType, delta: i32); async fn incr(&mut self, uri: &str, typ: parakeet_index::AggregateType) { self.add_delta(uri, typ, 1).await } async fn decr(&mut self, uri: &str, typ: parakeet_index::AggregateType) { self.add_delta(uri, typ, -1).await } } impl AggregateDeltaStore for tokio::sync::mpsc::Sender { async fn add_delta(&mut self, uri: &str, typ: parakeet_index::AggregateType, delta: i32) { let res = self .send(parakeet_index::AggregateDeltaReq { typ: typ.into(), uri: uri.to_string(), delta, }) .await; if let Err(e) = res { tracing::error!("failed to send aggregate delta: {e}"); } } } impl AggregateDeltaStore for std::collections::HashMap<(String, i32), i32> { async fn add_delta(&mut self, uri: &str, typ: parakeet_index::AggregateType, delta: i32) { let key = (uri.to_string(), typ.into()); self.entry(key).and_modify(|v| *v += delta).or_insert(delta); } }