diff --git a/consumer/src/indexer/db.rs b/consumer/src/indexer/db.rs index 13a3652d..db4640d4 100644 --- a/consumer/src/indexer/db.rs +++ b/consumer/src/indexer/db.rs @@ -693,3 +693,51 @@ pub async fn delete_chat_decl(conn: &mut AsyncPgConnection, did: &str) -> QueryR .execute(conn) .await } + +pub async fn upsert_starterpack( + conn: &mut AsyncPgConnection, + did: &str, + cid: Cid, + at_uri: &str, + rec: records::AppBskyGraphStarterPack, +) -> QueryResult { + let record = serde_json::to_value(&rec).unwrap(); + + let feeds = rec + .feeds + .map(|v| v.into_iter().map(|item| item.uri).collect()); + + let description_facets = rec + .description_facets + .and_then(|v| serde_json::to_value(v).ok()); + + + let data = models::NewStarterPack { + at_uri, + cid: cid.to_string(), + owner: did, + record, + name: &rec.name, + description: rec.description, + description_facets, + list: &rec.list, + feeds, + created_at: rec.created_at.naive_utc(), + indexed_at: Utc::now().naive_utc(), + }; + + diesel::insert_into(schema::starterpacks::table) + .values(&data) + .on_conflict(schema::starterpacks::at_uri) + .do_update() + .set(&data) + .execute(conn) + .await +} + +pub async fn delete_starterpack(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { + diesel::delete(schema::starterpacks::table) + .filter(schema::starterpacks::at_uri.eq(at_uri)) + .execute(conn) + .await +} diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index a08e971d..0c201c41 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -420,6 +420,9 @@ pub async fn index_op( db::insert_list_item(conn, at_uri, record).await?; } + RecordTypes::AppBskyGraphStarterPack(record) => { + db::upsert_starterpack(conn, repo, cid, at_uri, record).await?; + }, RecordTypes::ChatBskyActorDeclaration(record) => { if rkey == "self" { db::upsert_chat_decl(conn, repo, record).await?; @@ -458,6 +461,7 @@ pub async fn index_op_delete( } CollectionType::BskyListBlock => db::delete_list_block(conn, at_uri).await?, CollectionType::BskyListItem => db::delete_list_item(conn, at_uri).await?, + CollectionType::BskyStarterPack => db::delete_starterpack(conn, at_uri).await?, CollectionType::ChatActorDecl => db::delete_chat_decl(conn, at_uri).await?, _ => unreachable!(), }; diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs index c0fc3d39..89680cb6 100644 --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -293,6 +293,24 @@ pub struct AppBskyGraphListBlock { pub created_at: DateTime, } +#[derive(Debug, Deserialize, Serialize)] +#[serde(tag = "$type")] +#[serde(rename = "app.bsky.graph.starterpack")] +#[serde(rename_all = "camelCase")] +pub struct AppBskyGraphStarterPack { + pub name: String, + pub description: Option, + pub description_facets: Option>, + pub list: String, + pub feeds: Option>, + pub created_at: DateTime, +} + +#[derive(Debug, Deserialize, Serialize)] +pub struct StarterPackFeedItem { + pub uri: String, +} + #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct ChatBskyActorDeclaration { diff --git a/consumer/src/indexer/types.rs b/consumer/src/indexer/types.rs index 046aec7f..8e810576 100644 --- a/consumer/src/indexer/types.rs +++ b/consumer/src/indexer/types.rs @@ -29,6 +29,8 @@ pub enum RecordTypes { AppBskyGraphListBlock(records::AppBskyGraphListBlock), #[serde(rename = "app.bsky.graph.listitem")] AppBskyGraphListItem(records::AppBskyGraphListItem), + #[serde(rename = "app.bsky.graph.starterpack")] + AppBskyGraphStarterPack(records::AppBskyGraphStarterPack), #[serde(rename = "chat.bsky.actor.declaration")] ChatBskyActorDeclaration(records::ChatBskyActorDeclaration), } @@ -47,6 +49,7 @@ pub enum CollectionType { BskyList, BskyListBlock, BskyListItem, + BskyStarterPack, ChatActorDecl, Unsupported, } @@ -66,6 +69,7 @@ impl CollectionType { "app.bsky.graph.list" => CollectionType::BskyList, "app.bsky.graph.listblock" => CollectionType::BskyListBlock, "app.bsky.graph.listitem" => CollectionType::BskyListItem, + "app.bsky.graph.starterpack" => CollectionType::BskyStarterPack, "chat.bsky.actor.declaration" => CollectionType::ChatActorDecl, _ => CollectionType::Unsupported, } @@ -86,6 +90,7 @@ impl CollectionType { CollectionType::BskyListBlock => false, CollectionType::BskyListItem => false, CollectionType::ChatActorDecl => true, + CollectionType::BskyStarterPack => true, CollectionType::Unsupported => false, } } diff --git a/lexica/src/app_bsky/graph.rs b/lexica/src/app_bsky/graph.rs index 04f223d4..48e2b84b 100644 --- a/lexica/src/app_bsky/graph.rs +++ b/lexica/src/app_bsky/graph.rs @@ -1,4 +1,5 @@ -use crate::app_bsky::actor::ProfileView; +use crate::app_bsky::actor::{ProfileView, ProfileViewBasic}; +use crate::app_bsky::feed::GeneratorView; use crate::app_bsky::richtext::FacetMain; use chrono::prelude::*; use serde::{Deserialize, Serialize}; @@ -48,7 +49,7 @@ pub struct ListView { pub indexed_at: NaiveDateTime, } -#[derive(Debug, Serialize)] +#[derive(Clone, Debug, Serialize)] pub struct ListItemView { pub uri: String, pub subject: ProfileView, @@ -79,3 +80,42 @@ impl FromStr for ListPurpose { } } } + +#[derive(Clone, Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct StarterPackView { + pub uri: String, + pub cid: String, + pub record: serde_json::Value, + pub creator: ProfileViewBasic, + + #[serde(skip_serializing_if = "Option::is_none")] + pub list: Option, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub list_items_sample: Vec, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub feeds: Vec, + + pub list_item_count: i64, + pub joined_week_count: i64, + pub joined_all_time_count: i64, + + // pub labels: Vec<()>, + pub indexed_at: NaiveDateTime, +} + +#[derive(Clone, Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct StarterPackViewBasic { + pub uri: String, + pub cid: String, + pub record: serde_json::Value, + pub creator: ProfileViewBasic, + + pub list_item_count: i64, + pub joined_week_count: i64, + pub joined_all_time_count: i64, + + // pub labels: Vec<()>, + pub indexed_at: NaiveDateTime, +} diff --git a/migrations/2025-04-09-175735_starterpacks/down.sql b/migrations/2025-04-09-175735_starterpacks/down.sql new file mode 100644 index 00000000..cf7d9ff5 --- /dev/null +++ b/migrations/2025-04-09-175735_starterpacks/down.sql @@ -0,0 +1 @@ +drop table starterpacks; \ No newline at end of file diff --git a/migrations/2025-04-09-175735_starterpacks/up.sql b/migrations/2025-04-09-175735_starterpacks/up.sql new file mode 100644 index 00000000..c188c56a --- /dev/null +++ b/migrations/2025-04-09-175735_starterpacks/up.sql @@ -0,0 +1,18 @@ +create table starterpacks +( + at_uri text primary key, + owner text not null references actors (did), + cid text not null, + record jsonb not null, + + name text not null, + description text, + description_facets jsonb, + list text not null, + feeds text[], + + created_at timestamptz not null default now(), + indexed_at timestamp not null default now() +); + +create index starterpacks_owner_index on starterpacks using hash (owner); diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index 37afd9f6..f1744c47 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -554,3 +554,42 @@ pub struct NewChatDecl<'a> { pub did: &'a str, pub allow_incoming: String, } + +#[derive(Clone, Debug, Queryable, Selectable, Identifiable)] +#[diesel(table_name = crate::schema::starterpacks)] +#[diesel(primary_key(at_uri))] +#[diesel(check_for_backend(diesel::pg::Pg))] +pub struct StaterPack { + pub at_uri: String, + pub cid: String, + pub owner: String, + pub record: serde_json::Value, + + pub name: String, + pub description: Option, + pub description_facets: Option, + pub list: String, + pub feeds: Option>>, + + pub created_at: NaiveDateTime, + pub indexed_at: NaiveDateTime, +} + +#[derive(Insertable, AsChangeset)] +#[diesel(table_name = crate::schema::starterpacks)] +#[diesel(check_for_backend(diesel::pg::Pg))] +pub struct NewStarterPack<'a> { + pub at_uri: &'a str, + pub cid: String, + pub owner: &'a str, + pub record: serde_json::Value, + + pub name: &'a str, + pub description: Option, + pub description_facets: Option, + pub list: &'a str, + pub feeds: Option>, + + pub created_at: NaiveDateTime, + pub indexed_at: NaiveDateTime, +} diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index 19e12c42..d108f490 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -233,6 +233,22 @@ diesel::table! { } } +diesel::table! { + starterpacks (at_uri) { + at_uri -> Text, + owner -> Text, + cid -> Text, + record -> Jsonb, + name -> Text, + description -> Nullable, + description_facets -> Nullable, + list -> Text, + feeds -> Nullable>>, + created_at -> Timestamptz, + indexed_at -> Timestamp, + } +} + diesel::table! { threadgates (at_uri) { at_uri -> Text, @@ -264,6 +280,7 @@ diesel::joinable!(postgates -> posts (post_uri)); diesel::joinable!(posts -> actors (did)); diesel::joinable!(profiles -> actors (did)); diesel::joinable!(reposts -> actors (did)); +diesel::joinable!(starterpacks -> actors (owner)); diesel::joinable!(threadgates -> posts (post_uri)); diesel::allow_tables_to_appear_in_same_query!( @@ -287,5 +304,6 @@ diesel::allow_tables_to_appear_in_same_query!( posts, profiles, reposts, + starterpacks, threadgates, ); diff --git a/parakeet/src/hydration/mod.rs b/parakeet/src/hydration/mod.rs index 0382b4a1..51f5a72e 100644 --- a/parakeet/src/hydration/mod.rs +++ b/parakeet/src/hydration/mod.rs @@ -5,8 +5,10 @@ pub mod feedgen; pub mod list; pub mod posts; pub mod profile; +pub mod starter_packs; pub use feedgen::*; pub use list::*; pub use posts::*; -pub use profile::*; \ No newline at end of file +pub use profile::*; +pub use starter_packs::*; diff --git a/parakeet/src/hydration/starter_packs.rs b/parakeet/src/hydration/starter_packs.rs new file mode 100644 index 00000000..2728eaae --- /dev/null +++ b/parakeet/src/hydration/starter_packs.rs @@ -0,0 +1,145 @@ +use crate::hydration::{ + hydrate_feedgens, hydrate_list_basic, hydrate_lists_basic, hydrate_profile_basic, + hydrate_profiles_basic, +}; +use crate::loaders::Dataloaders; +use lexica::app_bsky::actor::ProfileViewBasic; +use lexica::app_bsky::feed::GeneratorView; +use lexica::app_bsky::graph::{ListViewBasic, StarterPackView, StarterPackViewBasic}; +use parakeet_db::models; +use std::collections::HashMap; + +fn build_basic( + starter_pack: models::StaterPack, + creator: ProfileViewBasic, + list_item_count: i64, +) -> StarterPackViewBasic { + StarterPackViewBasic { + uri: starter_pack.at_uri, + cid: starter_pack.cid, + record: starter_pack.record, + creator, + list_item_count, + joined_week_count: 0, + joined_all_time_count: 0, + indexed_at: starter_pack.indexed_at, + } +} + +fn build_spview( + starter_pack: models::StaterPack, + creator: ProfileViewBasic, + list: Option, + feeds: Vec, +) -> StarterPackView { + let list_item_count = list.as_ref().map(|list| list.list_item_count).unwrap_or(0); + + StarterPackView { + uri: starter_pack.at_uri, + cid: starter_pack.cid, + record: starter_pack.record, + creator, + list, + list_items_sample: vec![], // TODO: we should do this, but the app seems to be okay without? + feeds, + list_item_count, + joined_week_count: 0, + joined_all_time_count: 0, + indexed_at: starter_pack.indexed_at, + } +} + +pub async fn hydrate_starterpack_basic( + loaders: &Dataloaders, + pack: String, +) -> Option { + let sp = loaders.starterpacks.load(pack).await?; + let creator = hydrate_profile_basic(loaders, sp.owner.clone()).await?; + let (_, list_item_count) = loaders.list.load(sp.list.clone()).await?; + + Some(build_basic(sp, creator, list_item_count)) +} + +pub async fn hydrate_starterpacks_basic( + loaders: &Dataloaders, + packs: Vec, +) -> HashMap { + let packs = loaders.starterpacks.load_many(packs).await; + + let (creators, lists) = packs + .values() + .map(|pack| (pack.owner.clone(), pack.list.clone())) + .unzip(); + + let creators = hydrate_profiles_basic(loaders, creators).await; + let lists = loaders.list.load_many(lists).await; + + packs + .into_iter() + .filter_map(|(at_uri, pack)| { + let creator = creators.get(&pack.owner).cloned()?; + let list_item_count = lists.get(&pack.list).map(|(_, v)| *v).unwrap_or(0); + + Some((at_uri, build_basic(pack, creator, list_item_count))) + }) + .collect() +} + +pub async fn hydrate_starterpack(loaders: &Dataloaders, pack: String) -> Option { + let sp = loaders.starterpacks.load(pack).await?; + + let creator = hydrate_profile_basic(loaders, sp.owner.clone()).await?; + let list = hydrate_list_basic(loaders, sp.list.clone()).await; + + let feeds = sp + .feeds + .clone() + .unwrap_or_default() + .into_iter() + .flatten() + .collect(); + let feeds = hydrate_feedgens(loaders, feeds) + .await + .into_values() + .collect(); + + Some(build_spview(sp, creator, list, feeds)) +} + +pub async fn hydrate_starterpacks( + loaders: &Dataloaders, + packs: Vec, +) -> HashMap { + let packs = loaders.starterpacks.load_many(packs).await; + + let (creators, lists) = packs + .values() + .map(|pack| (pack.owner.clone(), pack.list.clone())) + .unzip(); + let feeds = packs + .values() + .filter_map(|pack| pack.feeds.clone()) + .flat_map(|feeds| feeds.into_iter().flatten()) + .collect(); + + let creators = hydrate_profiles_basic(loaders, creators).await; + let lists = hydrate_lists_basic(loaders, lists).await; + let feeds = hydrate_feedgens(loaders, feeds).await; + + packs + .into_iter() + .filter_map(|(at_uri, pack)| { + let creator = creators.get(&pack.owner).cloned()?; + let list = lists.get(&pack.list).cloned(); + let feeds = pack.feeds.as_ref().map(|v| { + v.iter() + .flatten() + .filter_map(|feed| feeds.get(feed).cloned()) + .collect() + }); + let feeds = feeds.unwrap_or_default(); + + Some((at_uri, build_spview(pack, creator, list, feeds))) + }) + .collect() +} diff --git a/parakeet/src/loaders.rs b/parakeet/src/loaders.rs index 14e0294f..3fe5cd07 100644 --- a/parakeet/src/loaders.rs +++ b/parakeet/src/loaders.rs @@ -16,6 +16,7 @@ pub struct Dataloaders { pub list: Loader, pub posts: Loader, pub profile: Loader, + pub starterpacks: Loader, } impl Dataloaders { @@ -29,6 +30,7 @@ impl Dataloaders { list: Loader::new(ListLoader(pool.clone())), posts: Loader::new(PostLoader(pool.clone())), profile: Loader::new(ProfileLoader(pool.clone())), + starterpacks: Loader::new(StarterPackLoader(pool.clone())), } } } @@ -274,3 +276,28 @@ impl BatchFn for EmbedLoader { )) } } + +pub struct StarterPackLoader(Pool); +type StarterPackLoaderRet = models::StaterPack; +impl BatchFn for StarterPackLoader { + async fn load(&mut self, keys: &[String]) -> HashMap { + let mut conn = self.0.get().await.unwrap(); + + let res = schema::starterpacks::table + .select(models::StaterPack::as_select()) + .filter(schema::starterpacks::at_uri.eq_any(keys)) + .load(&mut conn) + .await; + + match res { + Ok(res) => HashMap::from_iter( + res.into_iter() + .map(|starterpack| (starterpack.at_uri.clone(), (starterpack))), + ), + Err(e) => { + tracing::error!("starterpack load failed: {e}"); + HashMap::new() + } + } + } +} diff --git a/parakeet/src/xrpc/app_bsky/graph/mod.rs b/parakeet/src/xrpc/app_bsky/graph/mod.rs index c4005390..3b553cc6 100644 --- a/parakeet/src/xrpc/app_bsky/graph/mod.rs +++ b/parakeet/src/xrpc/app_bsky/graph/mod.rs @@ -1,2 +1,3 @@ pub mod lists; pub mod relations; +pub mod starter_packs; diff --git a/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs b/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs new file mode 100644 index 00000000..29776ea1 --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs @@ -0,0 +1,111 @@ +use crate::xrpc::error::{Error, XrpcResult}; +use crate::xrpc::{datetime_cursor, get_actor_did, ActorWithCursorQuery}; +use crate::{hydration, GlobalState}; +use axum::extract::{Query, State}; +use axum::Json; +use axum_extra::extract::Query as ExtraQuery; +use diesel::prelude::*; +use diesel_async::RunQueryDsl; +use lexica::app_bsky::graph::{StarterPackView, StarterPackViewBasic}; +use parakeet_db::schema; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct StarterPacksRes { + pub starter_packs: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + cursor: Option, +} + +pub async fn get_actor_starter_packs( + State(state): State, + Query(query): Query, +) -> XrpcResult> { + let mut conn = state.pool.get().await?; + + let subj_did = get_actor_did(&state.dataloaders, query.actor).await?; + + let limit = query.limit.unwrap_or(50).clamp(1, 100); + + let mut sp_query = schema::starterpacks::table + .select(( + schema::starterpacks::created_at, + schema::starterpacks::at_uri, + )) + .filter(schema::starterpacks::owner.eq(subj_did)) + .into_boxed(); + + if let Some(cursor) = datetime_cursor(query.cursor.as_ref()) { + sp_query = sp_query.filter(schema::starterpacks::created_at.lt(cursor)); + } + + let results = sp_query + .order(schema::starterpacks::created_at.desc()) + .limit(limit as i64) + .load::<(chrono::DateTime, String)>(&mut conn) + .await?; + + let cursor = results + .last() + .map(|(last, _)| last.timestamp_millis().to_string()); + + let uris = results.iter().map(|(_, uri)| uri.clone()).collect(); + + let mut starter_packs = hydration::hydrate_starterpacks_basic(&state.dataloaders, uris).await; + + let starter_packs = results + .into_iter() + .filter_map(|(_, uri)| starter_packs.remove(&uri)) + .collect(); + + Ok(Json(StarterPacksRes { + starter_packs, + cursor, + })) +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct GetStarterPackQuery { + pub starter_pack: String, +} + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct GetStarterPackRes { + pub starter_pack: StarterPackView, +} + +pub async fn get_starter_pack( + State(state): State, + Query(query): Query, +) -> XrpcResult> { + let Some(starter_pack) = + hydration::hydrate_starterpack(&state.dataloaders, query.starter_pack).await + else { + return Err(Error::not_found()); + }; + + Ok(Json(GetStarterPackRes { starter_pack })) +} + +#[derive(Debug, Deserialize)] +pub struct GetStarterPacksQuery { + pub uris: Vec, +} + +pub async fn get_starter_packs( + State(state): State, + ExtraQuery(query): ExtraQuery, +) -> XrpcResult> { + let starter_packs = hydration::hydrate_starterpacks_basic(&state.dataloaders, query.uris) + .await + .into_values() + .collect(); + + Ok(Json(StarterPacksRes { + starter_packs, + cursor: None, + })) +} diff --git a/parakeet/src/xrpc/app_bsky/mod.rs b/parakeet/src/xrpc/app_bsky/mod.rs index c5d6da9b..aa962ec3 100644 --- a/parakeet/src/xrpc/app_bsky/mod.rs +++ b/parakeet/src/xrpc/app_bsky/mod.rs @@ -19,8 +19,11 @@ pub fn routes() -> Router { .route("/app.bsky.feed.getRepostedBy", get(feed::posts::get_reposted_by)) .route("/app.bsky.feed.getFeedGenerator", get(feed::feedgen::get_feed_generator)) .route("/app.bsky.feed.getFeedGenerators", get(feed::feedgen::get_feed_generators)) + .route("/app.bsky.graph.getActorStarterPacks", get(graph::starter_packs::get_actor_starter_packs)) .route("/app.bsky.graph.getFollowers", get(graph::relations::get_followers)) .route("/app.bsky.graph.getFollows", get(graph::relations::get_follows)) .route("/app.bsky.graph.getList", get(graph::lists::get_list)) .route("/app.bsky.graph.getLists", get(graph::lists::get_lists)) + .route("/app.bsky.graph.getStarterPack", get(graph::starter_packs::get_starter_pack)) + .route("/app.bsky.graph.getStarterPacks", get(graph::starter_packs::get_starter_packs)) }