diff --git a/consumer/src/indexer/db.rs b/consumer/src/indexer/db.rs index f4643e2a..1db59cef 100644 --- a/consumer/src/indexer/db.rs +++ b/consumer/src/indexer/db.rs @@ -257,3 +257,45 @@ pub async fn delete_list_item(conn: &mut AsyncPgConnection, at_uri: &str) -> Que .execute(conn) .await } + +pub async fn upsert_feedgen( + conn: &mut AsyncPgConnection, + repo: &str, + cid: Cid, + at_uri: &str, + rec: records::AppBskyFeedGenerator, +) -> QueryResult { + let description_facets = rec + .description_facets + .and_then(|v| serde_json::to_value(v).ok()); + + let data = models::UpsertFeedGen { + at_uri, + cid: &cid.to_string(), + owner: repo, + service_did: &rec.did, + content_mode: rec.content_mode, + name: &rec.display_name, + description: rec.description, + description_facets, + avatar_cid: blob_ref(rec.avatar), + accepts_interactions: rec.accepts_interactions, + created_at: rec.created_at.naive_utc(), + indexed_at: Utc::now().naive_utc(), + }; + + diesel::insert_into(schema::feedgens::table) + .values(&data) + .on_conflict(schema::feedgens::at_uri) + .do_update() + .set(&data) + .execute(conn) + .await +} + +pub async fn delete_feedgen(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { + diesel::delete(schema::feedgens::table) + .filter(schema::feedgens::at_uri.eq(at_uri)) + .execute(conn) + .await +} diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index 30deb143..dd3829b4 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -221,6 +221,9 @@ async fn process_op( db::upsert_profile(conn, repo, cid, record).await?; } + RecordTypes::AppBskyFeedGenerator(record) => { + db::upsert_feedgen(conn, repo, cid, &full_path, record).await?; + } RecordTypes::AppBskyGraphBlock(record) => { db::insert_block(conn, repo, &full_path, record).await?; } @@ -252,6 +255,7 @@ async fn process_op( } else if op.action == "delete" { match collection { CollectionType::BskyBlock => db::delete_block(conn, &full_path).await?, + CollectionType::BskyFeedGen => db::delete_feedgen(conn, &full_path).await?, CollectionType::BskyFollow => { if let Some(subject) = db::delete_follow(conn, &full_path).await? { db::update_follow_stats(conn, &[repo, &subject]).await? diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs index ef5e4ed1..227c7324 100644 --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -37,6 +37,20 @@ pub struct AppBskyActorProfile { pub created_at: Option>, } +#[derive(Debug, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AppBskyFeedGenerator { + pub did: String, + pub display_name: String, + pub description: Option, + pub description_facets: Option>, + pub avatar: Option, + pub accepts_interactions: Option, + // pub labels: Option>, + pub content_mode: Option, + pub created_at: DateTime, +} + #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct AppBskyGraphBlock { diff --git a/consumer/src/indexer/types.rs b/consumer/src/indexer/types.rs index a8b5e5a0..69dff210 100644 --- a/consumer/src/indexer/types.rs +++ b/consumer/src/indexer/types.rs @@ -6,6 +6,8 @@ use serde::{Deserialize, Serialize}; pub enum RecordTypes { #[serde(rename = "app.bsky.actor.profile")] AppBskyActorProfile(records::AppBskyActorProfile), + #[serde(rename = "app.bsky.feed.generator")] + AppBskyFeedGenerator(records::AppBskyFeedGenerator), #[serde(rename = "app.bsky.graph.block")] AppBskyGraphBlock(records::AppBskyGraphBlock), #[serde(rename = "app.bsky.graph.follow")] @@ -21,6 +23,7 @@ pub enum RecordTypes { #[derive(Debug, PartialOrd, PartialEq)] pub enum CollectionType { BskyProfile, + BskyFeedGen, BskyBlock, BskyFollow, BskyList, @@ -33,6 +36,7 @@ impl CollectionType { pub(crate) fn from_str(input: &str) -> CollectionType { match input { "app.bsky.actor.profile" => CollectionType::BskyProfile, + "app.bsky.feed.generator" => CollectionType::BskyFeedGen, "app.bsky.graph.block" => CollectionType::BskyBlock, "app.bsky.graph.follow" => CollectionType::BskyFollow, "app.bsky.graph.list" => CollectionType::BskyList, @@ -45,6 +49,7 @@ impl CollectionType { pub fn can_update(&self) -> bool { match self { CollectionType::BskyProfile => true, + CollectionType::BskyFeedGen => true, CollectionType::BskyBlock => false, CollectionType::BskyFollow => false, CollectionType::BskyList => true, diff --git a/lexica/src/app_bsky/feed.rs b/lexica/src/app_bsky/feed.rs new file mode 100644 index 00000000..d4aaae12 --- /dev/null +++ b/lexica/src/app_bsky/feed.rs @@ -0,0 +1,53 @@ +use crate::app_bsky::actor::ProfileView; +use crate::app_bsky::richtext::FacetMain; +use chrono::prelude::*; +use serde::Serialize; +use std::str::FromStr; + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct GeneratorView { + pub uri: String, + pub cid: String, + pub did: String, + pub creator: ProfileView, + pub display_name: String, + + #[serde(skip_serializing_if = "Option::is_none")] + pub description: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub description_facets: Option>, + + #[serde(skip_serializing_if = "Option::is_none")] + pub avatar: Option, + pub like_count: i64, + + pub accepts_interactions: bool, + // pub labels: Vec<()>, + // #[serde(skip_serializing_if = "Option::is_none")] + // pub viewer: Option<()>, + #[serde(skip_serializing_if = "Option::is_none")] + pub content_mode: Option, + + pub indexed_at: NaiveDateTime, +} + +#[derive(Debug, Serialize)] +pub enum GeneratorContentMode { + #[serde(rename = "app.bsky.feed.defs#contentModeUnspecified")] + Unspecified, + #[serde(rename = "app.bsky.feed.defs#contentModeVideo")] + Video, +} + +impl FromStr for GeneratorContentMode { + type Err = (); + + fn from_str(s: &str) -> Result { + match s { + "app.bsky.feed.defs#contentModeUnspecified" => Ok(Self::Unspecified), + "app.bsky.feed.defs#contentModeVideo" => Ok(Self::Video), + _ => Err(()), + } + } +} diff --git a/lexica/src/app_bsky/mod.rs b/lexica/src/app_bsky/mod.rs index c7cde70b..bfbde235 100644 --- a/lexica/src/app_bsky/mod.rs +++ b/lexica/src/app_bsky/mod.rs @@ -1,3 +1,4 @@ pub mod actor; +pub mod feed; pub mod graph; pub mod richtext; diff --git a/migrations/2025-02-08-162441_feedgen/down.sql b/migrations/2025-02-08-162441_feedgen/down.sql new file mode 100644 index 00000000..4d246885 --- /dev/null +++ b/migrations/2025-02-08-162441_feedgen/down.sql @@ -0,0 +1 @@ +drop table feedgens; \ No newline at end of file diff --git a/migrations/2025-02-08-162441_feedgen/up.sql b/migrations/2025-02-08-162441_feedgen/up.sql new file mode 100644 index 00000000..b91eecb3 --- /dev/null +++ b/migrations/2025-02-08-162441_feedgen/up.sql @@ -0,0 +1,20 @@ +create table feedgens +( + at_uri text primary key, + cid text not null, + owner text not null references actors (did), + + service_did text not null, + content_mode text, + name text not null, + + description text, + description_facets jsonb, + avatar_cid text, + accepts_interactions bool not null default false, + + created_at timestamptz not null default now(), + indexed_at timestamp not null default now() +); + +create index feedgens_owner_index on feedgens using hash (owner); \ No newline at end of file diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index 7be1e82b..d3afee53 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -193,4 +193,49 @@ pub struct NewListItem<'a> { pub created_at: NaiveDateTime, pub indexed_at: NaiveDateTime, -} \ No newline at end of file +} + +#[derive(Clone, Debug, Queryable, Selectable, Identifiable)] +#[diesel(table_name = crate::schema::feedgens)] +#[diesel(primary_key(at_uri))] +#[diesel(check_for_backend(diesel::pg::Pg))] +pub struct FeedGen { + pub at_uri: String, + pub cid: String, + pub owner: String, + + pub service_did: String, + pub content_mode: Option, + pub name: String, + + pub description: Option, + pub description_facets: Option, + pub avatar_cid: Option, + pub accepts_interactions: bool, + + pub created_at: NaiveDateTime, + pub indexed_at: NaiveDateTime, +} + +#[derive(Insertable, AsChangeset)] +#[diesel(table_name = crate::schema::feedgens)] +#[diesel(check_for_backend(diesel::pg::Pg))] +#[diesel(treat_none_as_null = true)] +pub struct UpsertFeedGen<'a> { + pub at_uri: &'a str, + pub cid: &'a str, + pub owner: &'a str, + + pub service_did: &'a str, + pub content_mode: Option, + pub name: &'a str, + + pub description: Option, + pub description_facets: Option, + pub avatar_cid: Option, + #[diesel(treat_none_as_null = false)] + pub accepts_interactions: Option, + + pub created_at: NaiveDateTime, + pub indexed_at: NaiveDateTime, +} diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index 2f231e5b..cef8a22e 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -21,6 +21,23 @@ diesel::table! { } } +diesel::table! { + feedgens (at_uri) { + at_uri -> Text, + cid -> Text, + owner -> Text, + service_did -> Text, + content_mode -> Nullable, + name -> Text, + description -> Nullable, + description_facets -> Nullable, + avatar_cid -> Nullable, + accepts_interactions -> Bool, + created_at -> Timestamptz, + indexed_at -> Timestamp, + } +} + diesel::table! { follow_stats (did) { did -> Text, @@ -91,6 +108,7 @@ diesel::table! { } diesel::joinable!(blocks -> actors (did)); +diesel::joinable!(feedgens -> actors (owner)); diesel::joinable!(follows -> actors (did)); diesel::joinable!(list_blocks -> actors (did)); diesel::joinable!(list_blocks -> lists (list_uri)); @@ -101,6 +119,7 @@ diesel::joinable!(profiles -> actors (did)); diesel::allow_tables_to_appear_in_same_query!( actors, blocks, + feedgens, follow_stats, follows, list_blocks, diff --git a/parakeet/src/hydration/feedgen.rs b/parakeet/src/hydration/feedgen.rs new file mode 100644 index 00000000..fcc98c17 --- /dev/null +++ b/parakeet/src/hydration/feedgen.rs @@ -0,0 +1,63 @@ +use crate::hydration::profile::{hydrate_profile, hydrate_profiles}; +use crate::loaders::Dataloaders; +use lexica::app_bsky::actor::ProfileView; +use lexica::app_bsky::feed::{GeneratorContentMode, GeneratorView}; +use parakeet_db::models; +use std::collections::HashMap; +use std::str::FromStr; + +fn build_feedgen(feedgen: models::FeedGen, creator: ProfileView) -> GeneratorView { + let content_mode = feedgen + .content_mode + .and_then(|cm| GeneratorContentMode::from_str(&cm).ok()); + + let description_facets = feedgen + .description_facets + .and_then(|v| serde_json::from_value(v).ok()); + + GeneratorView { + uri: feedgen.at_uri, + cid: feedgen.cid, + did: feedgen.service_did, + creator, + display_name: feedgen.name, + description: feedgen.description, + description_facets, + avatar: feedgen + .avatar_cid + .map(|v| format!("https://localhost/feedgen/{v}")), + like_count: 0, + accepts_interactions: feedgen.accepts_interactions, + content_mode, + indexed_at: feedgen.indexed_at, + } +} + +pub async fn hydrate_feedgen(loaders: &Dataloaders, feedgen: String) -> Option { + let feedgen = loaders.feedgen.load(feedgen).await?; + let profile = hydrate_profile(loaders, feedgen.owner.clone()).await?; + + Some(build_feedgen(feedgen, profile)) +} + +pub async fn hydrate_feedgens( + loaders: &Dataloaders, + feedgens: Vec, +) -> HashMap { + let feedgens = loaders.feedgen.load_many(feedgens).await; + + let creators = feedgens + .values() + .map(|feedgen| feedgen.owner.clone()) + .collect(); + let creators = hydrate_profiles(loaders, creators).await; + + feedgens + .into_iter() + .filter_map(|(uri, feedgen)| { + let creator = creators.get(&feedgen.owner)?; + + Some((uri, build_feedgen(feedgen, creator.to_owned()))) + }) + .collect() +} diff --git a/parakeet/src/hydration/mod.rs b/parakeet/src/hydration/mod.rs index 7c31afd1..6173fbf5 100644 --- a/parakeet/src/hydration/mod.rs +++ b/parakeet/src/hydration/mod.rs @@ -1,4 +1,5 @@ #![allow(unused)] +pub mod feedgen; pub mod list; pub mod profile; diff --git a/parakeet/src/loaders.rs b/parakeet/src/loaders.rs index e0722e25..7cc16a52 100644 --- a/parakeet/src/loaders.rs +++ b/parakeet/src/loaders.rs @@ -7,6 +7,7 @@ use parakeet_db::{models, schema}; use std::collections::HashMap; pub struct Dataloaders { + pub feedgen: Loader, pub handle: Loader, pub list: Loader, pub profile: Loader, @@ -17,6 +18,7 @@ impl Dataloaders { // we should build a redis/valkey backend at some point in the future. pub fn new(pool: Pool) -> Dataloaders { Dataloaders { + feedgen: Loader::new(FeedGenLoader(pool.clone())), handle: Loader::new(HandleLoader(pool.clone())), list: Loader::new(ListLoader(pool.clone())), profile: Loader::new(ProfileLoader(pool.clone())), @@ -119,3 +121,28 @@ impl BatchFn for ListLoader { } } } + +pub struct FeedGenLoader(Pool); +type FeedGenLoaderRet = models::FeedGen; //todo: when we have likes, we'll need the count here +impl BatchFn for FeedGenLoader { + async fn load(&mut self, keys: &[String]) -> HashMap { + let mut conn = self.0.get().await.unwrap(); + + let res = schema::feedgens::table + .select(models::FeedGen::as_select()) + .filter(schema::feedgens::at_uri.eq_any(keys)) + .load(&mut conn) + .await; + + match res { + Ok(res) => HashMap::from_iter( + res.into_iter() + .map(|feedgen| (feedgen.at_uri.clone(), feedgen)), + ), + Err(e) => { + tracing::error!("feedgen load failed: {e}"); + HashMap::new() + } + } + } +} diff --git a/parakeet/src/xrpc/app_bsky/feed/feedgen.rs b/parakeet/src/xrpc/app_bsky/feed/feedgen.rs new file mode 100644 index 00000000..565dfcc6 --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/feedgen.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::feed::GeneratorView; +use parakeet_db::schema; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize)] +pub struct GetActorFeedRes { + #[serde(skip_serializing_if = "Option::is_none")] + cursor: Option, + feeds: Vec, +} + +pub async fn get_actor_feeds( + State(state): State, + Query(query): Query, +) -> XrpcResult> { + let mut conn = state.pool.get().await?; + + let did = get_actor_did(&state.dataloaders, query.actor).await?; + + let limit = query.limit.unwrap_or(50).clamp(1, 100); + + let mut feeds_query = schema::feedgens::table + .select((schema::feedgens::created_at, schema::feedgens::at_uri)) + .filter(schema::feedgens::owner.eq(did)) + .into_boxed(); + + if let Some(cursor) = datetime_cursor(query.cursor.as_ref()) { + feeds_query = feeds_query.filter(schema::feedgens::created_at.lt(cursor)); + } + + let results = feeds_query + .order(schema::feedgens::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 at_uris = results.iter().map(|(_, uri)| uri.clone()).collect(); + + let mut feeds = hydration::feedgen::hydrate_feedgens(&state.dataloaders, at_uris).await; + + let feeds = results + .into_iter() + .filter_map(|(_, uri)| feeds.remove(&uri)) + .collect(); + + Ok(Json(GetActorFeedRes { cursor, feeds })) +} + +#[derive(Debug, Deserialize)] +pub struct GetFeedGeneratorQuery { + pub feed: String, +} + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct GetFeedGeneratorRes { + pub view: GeneratorView, + pub is_online: bool, + pub is_valid: bool, +} + +pub async fn get_feed_generator( + State(state): State, + Query(query): Query, +) -> XrpcResult> { + let Some(view) = hydration::feedgen::hydrate_feedgen(&state.dataloaders, query.feed).await + else { + return Err(Error::not_found()); + }; + + // todo: make the two flags work + Ok(Json(GetFeedGeneratorRes { + view, + is_online: true, + is_valid: true, + })) +} + +#[derive(Debug, Deserialize)] +pub struct FeedsQuery { + pub feeds: Vec, +} + +#[derive(Debug, Serialize)] +pub struct GetFeedGeneratorsRes { + pub feeds: Vec, +} + +pub async fn get_feed_generators( + State(state): State, + ExtraQuery(query): ExtraQuery, +) -> XrpcResult> { + let feeds = hydration::feedgen::hydrate_feedgens(&state.dataloaders, query.feeds) + .await + .into_values() + .collect(); + + Ok(Json(GetFeedGeneratorsRes { feeds })) +} diff --git a/parakeet/src/xrpc/app_bsky/feed/mod.rs b/parakeet/src/xrpc/app_bsky/feed/mod.rs new file mode 100644 index 00000000..1300d81f --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/mod.rs @@ -0,0 +1 @@ +pub mod feedgen; diff --git a/parakeet/src/xrpc/app_bsky/mod.rs b/parakeet/src/xrpc/app_bsky/mod.rs index adcd2a84..fe7d12ea 100644 --- a/parakeet/src/xrpc/app_bsky/mod.rs +++ b/parakeet/src/xrpc/app_bsky/mod.rs @@ -2,12 +2,16 @@ use axum::routing::get; use axum::Router; mod actor; +mod feed; mod graph; pub fn routes() -> Router { Router::new() .route("/app.bsky.actor.getProfile", get(actor::get_profile)) .route("/app.bsky.actor.getProfiles", get(actor::get_profiles)) + .route("/app.bsky.feed.getActorFeeds", get(feed::feedgen::get_actor_feeds)) + .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.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))