From 016d3c80222ae4e3ab6998840a3df0f874078963 Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Wed, 17 Dec 2025 14:58:30 -0500 Subject: [PATCH] whoops, uncleanup --- .../src/xrpc/app_bsky/feed/posts/feeds.rs | 173 ++++++++++++ .../src/xrpc/app_bsky/feed/posts/helpers.rs | 88 ++++++ .../src/xrpc/app_bsky/feed/posts/queries.rs | 265 ++++++++++++++++++ .../src/xrpc/app_bsky/feed/posts/threads.rs | 170 +++++++++++ 4 files changed, 696 insertions(+) create mode 100644 parakeet/src/xrpc/app_bsky/feed/posts/feeds.rs create mode 100644 parakeet/src/xrpc/app_bsky/feed/posts/helpers.rs create mode 100644 parakeet/src/xrpc/app_bsky/feed/posts/queries.rs create mode 100644 parakeet/src/xrpc/app_bsky/feed/posts/threads.rs diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/feeds.rs b/parakeet/src/xrpc/app_bsky/feed/posts/feeds.rs new file mode 100644 index 00000000..60f27554 --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/posts/feeds.rs @@ -0,0 +1,173 @@ +use axum::extract::{Query, State}; +use axum::Json; +use lexica::app_bsky::feed::{FeedViewPost, GeneratorView}; +use serde::{Deserialize, Serialize}; + +use crate::xrpc::datetime_cursor; +use crate::common::errors::XrpcResult; +use crate::common::auth::{AtpAcceptLabelers, AtpAuth}; +use crate::GlobalState; + +#[derive(Debug, Deserialize)] +pub struct GetFeedQuery { + pub feed: String, + pub limit: Option, + pub cursor: Option, +} + +#[derive(Debug, Serialize)] +pub struct GetFeedRes { + #[serde(skip_serializing_if = "Option::is_none")] + pub cursor: Option, + pub feed: Vec, +} + +pub async fn get_feed( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + let _limit = query.limit.unwrap_or(50).clamp(1, 100); + + // TODO: Implement feed generation logic + // For now, return empty feed + Ok(Json(GetFeedRes { + cursor: None, + feed: Vec::new(), + })) +} + +pub async fn get_feed_skeleton( + State(_state): State, + auth: AtpAuth, + Query(_query): Query, +) -> XrpcResult> { + let _viewer_did = auth.0.clone(); + + // TODO: Implement skeleton feed logic + Ok(Json(GetFeedSkeletonRes { + cursor: None, + feed: Vec::new(), + })) +} + +#[derive(Debug, Serialize)] +pub struct GetFeedSkeletonRes { + #[serde(skip_serializing_if = "Option::is_none")] + pub cursor: Option, + pub feed: Vec, +} + +#[derive(Debug, Deserialize)] +pub struct GetFeedGeneratorQuery { + pub feed: String, +} + +#[derive(Debug, Serialize)] +pub struct GetFeedGeneratorRes { + pub view: GeneratorView, + pub is_online: bool, + pub is_valid: bool, +} + +pub async fn get_feed_generator( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + // Get feed generator using FeedGeneratorEntity + let view = state.feedgen_entity + .get_by_uri(&query.feed, viewer_did.as_deref()) + .await? + .ok_or_else(|| crate::common::errors::Error::not_found())?; + + Ok(Json(GetFeedGeneratorRes { + view, + is_online: true, + is_valid: true, + })) +} + +#[derive(Debug, Deserialize)] +pub struct GetFeedGeneratorsQuery { + pub feeds: Vec, +} + +#[derive(Debug, Serialize)] +pub struct GetFeedGeneratorsRes { + pub feeds: Vec, +} + +pub async fn get_feed_generators( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + // Limit the number of feeds + let uris = query.feeds.into_iter().take(25).collect::>(); + + // Get feed generators using FeedGeneratorEntity + let feeds_map = state.feedgen_entity + .get_by_uris(uris.clone(), viewer_did.as_deref()) + .await?; + + // Preserve order + let feeds = uris + .into_iter() + .filter_map(|uri| feeds_map.get(&uri).cloned()) + .collect(); + + Ok(Json(GetFeedGeneratorsRes { feeds })) +} + +// Stub implementations for missing feed functions + +#[derive(Debug, Deserialize)] +pub struct GetAuthorFeedQuery { + pub actor: String, + pub limit: Option, + pub cursor: Option, + #[serde(default)] + pub filter: String, +} + +pub async fn get_author_feed( + State(_state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + _maybe_auth: Option, + Query(_query): Query, +) -> XrpcResult> { + // TODO: Implement author feed logic + Ok(Json(GetFeedRes { + cursor: None, + feed: Vec::new(), + })) +} + +#[derive(Debug, Deserialize)] +pub struct GetListFeedQuery { + pub list: String, + pub limit: Option, + pub cursor: Option, +} + +pub async fn get_list_feed( + State(_state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + _maybe_auth: Option, + Query(_query): Query, +) -> XrpcResult> { + // TODO: Implement list feed logic + Ok(Json(GetFeedRes { + cursor: None, + feed: Vec::new(), + })) +} \ No newline at end of file diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/helpers.rs b/parakeet/src/xrpc/app_bsky/feed/posts/helpers.rs new file mode 100644 index 00000000..6f4f6aaa --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/posts/helpers.rs @@ -0,0 +1,88 @@ +use axum_extra::headers::authorization::Bearer; +use axum_extra::headers::Authorization; +use axum_extra::TypedHeader; +use chrono::NaiveDateTime; +use diesel_async::AsyncPgConnection; +use lexica::app_bsky::feed::{ + BlockedAuthor, FeedSkeletonResponse, PostView, ThreadViewPost, ThreadViewPostType, +}; +use reqwest::Url; +use std::collections::HashMap; + +use crate::common::errors::{Error, XrpcResult}; + +#[expect(dead_code)] +const FEEDGEN_SERVICE_ID: &str = "#bsky_fg"; + +pub(super) async fn get_feed_skeleton( + client: &reqwest::Client, + feed: &str, + service: &str, + maybe_tok: Option<&TypedHeader>>, + limit: Option, + cursor: Option, +) -> XrpcResult { + let mut params = vec![("feed", feed.to_owned())]; + + if let Some(cursor) = cursor { + params.push(("cursor", cursor)); + } + if let Some(limit) = limit { + params.push(("limit", limit.to_string())); + } + let url = Url::parse_with_params( + &format!("{service}/xrpc/app.bsky.feed.getFeedSkeleton"), + params, + ) + .unwrap(); + + let mut req = client.get(url); + if let Some(auth) = maybe_tok { + req = req.bearer_auth(auth.token()); + } + + match req.send().await { + Ok(skeleton) => match skeleton.json().await { + Ok(skeleton) => Ok(skeleton), + Err(err) => { + tracing::error!("Failed to parse feed skeleton: {err}"); + Err(Error::server_error(Some("Failed to fetch feed skeleton"))) + } + }, + Err(err) => { + tracing::error!("Failed to fetch feed skeleton: {err}"); + Err(Error::server_error(Some("Failed to fetch feed skeleton"))) + } + } +} + +pub(super) async fn get_skeleton_repost_data( + _conn: &mut AsyncPgConnection, + _reposts: Vec, +) -> HashMap { + // TODO: Implement proper repost tracking with timestamps + // For now, return empty map + HashMap::new() +} + +pub(super) fn postview_to_tvpt( + post: PostView, + parent: Option, + replies: Vec, +) -> ThreadViewPostType { + match &post.author.viewer { + Some(v) if v.blocked_by || v.blocking.is_some() => ThreadViewPostType::Blocked { + uri: post.uri.clone(), + blocked: true, + author: Box::new(BlockedAuthor { + did: post.author.did, + viewer: post.author.viewer, + }), + }, + _ => ThreadViewPostType::Post(Box::new(ThreadViewPost { + post, + parent, + replies, + })), + } +} diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs b/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs new file mode 100644 index 00000000..42119e6e --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs @@ -0,0 +1,265 @@ +use axum::extract::{Query, State}; +use axum::Json; +use axum_extra::extract::Query as ExtraQuery; +use lexica::app_bsky::feed::PostView; +use serde::{Deserialize, Serialize}; + +use crate::common::errors::XrpcResult; +use crate::common::auth::{AtpAcceptLabelers, AtpAuth}; +use crate::GlobalState; + +#[derive(Debug, Deserialize)] +pub struct PostsQuery { + pub uris: Vec, +} + +#[derive(Debug, Serialize)] +pub struct PostsRes { + pub posts: Vec, +} + +pub async fn get_posts( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + ExtraQuery(query): ExtraQuery, +) -> XrpcResult> { + // Get viewer DID if authenticated + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + // Limit the number of URIs + let uris = query.uris.into_iter().take(25).collect::>(); + + // Get posts using PostEntity + let posts_map = state.post_entity.get_by_uris(uris.clone(), viewer_did.as_deref()).await?; + + // Preserve the order of the original URIs + let posts: Vec = uris + .into_iter() + .filter_map(|uri| posts_map.get(&uri).cloned()) + .collect(); + + Ok(Json(PostsRes { posts })) +} + +#[derive(Debug, Deserialize)] +pub struct PostQuery { + pub uri: String, +} + +#[derive(Debug, Serialize)] +pub struct PostRes { + pub uri: String, + pub cid: String, + pub value: serde_json::Value, +} + +pub async fn get_post( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + // Get viewer DID if authenticated + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + // Get the post using PostEntity + let post_view = state.post_entity + .get_by_uri(&query.uri, viewer_did.as_deref()) + .await? + .ok_or_else(|| crate::common::errors::Error::not_found())?; + + // Extract the record value from the PostView + // This is a simplified response - the actual endpoint returns the raw record + let value = serde_json::json!({ + "$type": "app.bsky.feed.post", + "text": post_view.record.get("text").unwrap_or(&serde_json::Value::String("".to_string())), + "createdAt": post_view.indexed_at.to_rfc3339(), + }); + + Ok(Json(PostRes { + uri: post_view.uri, + cid: post_view.cid, + value, + })) +} + +#[derive(Debug, Deserialize)] +pub struct GetQuotesQuery { + pub uri: String, + pub cid: Option, + pub limit: Option, + pub cursor: Option, +} + +#[derive(Debug, Serialize)] +pub struct GetQuotesRes { + pub uri: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub cid: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cursor: Option, + pub posts: Vec, +} + +pub async fn get_quotes( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let mut conn = state.pool.get().await?; + let _viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + let limit = query.limit.unwrap_or(50).clamp(1, 100); + + // Parse post URI to extract actor_id and rkey + let parts: Vec<&str> = query.uri.trim_start_matches("at://").split('/').collect(); + if parts.len() < 3 { + return Ok(Json(GetQuotesRes { + uri: query.uri, + cid: query.cid, + cursor: None, + posts: vec![], + })); + } + + let embed_did = parts[0]; + let embed_rkey_str = parts[2]; + + // Resolve DID to actor_id + let embed_actor_id = state.profile_entity + .resolve_identifier(embed_did) + .await + .map_err(|_| crate::common::errors::Error::not_found())?; + + // Decode rkey + let embed_rkey = parakeet_db::tid_util::decode_tid(embed_rkey_str) + .map_err(|_| crate::common::errors::Error::invalid_request(Some("Invalid rkey".to_string())))?; + + // Parse cursor + let (cursor_actor_id, cursor_rkey) = if let Some(ref cursor) = query.cursor { + let parts: Vec<&str> = cursor.split(':').collect(); + if parts.len() == 2 { + ( + parts[0].parse::().ok(), + parts[1].parse::().ok() + ) + } else { + (None, None) + } + } else { + (None, None) + }; + + // Get quotes from database using PostEntity + let cursor = cursor_actor_id.zip(cursor_rkey); + let results = state.post_entity.get_quotes(embed_actor_id, embed_rkey, cursor, limit).await + .map_err(|e| crate::common::errors::Error::server_error(Some(&e.to_string())))?; + + // Convert to PostViews + let mut posts = Vec::new(); + for (actor_id, rkey) in &results { + // Build URI + let author_did = state.profile_entity.get_did_by_id(*actor_id).await + .unwrap_or_else(|_| format!("did:plc:unknown{}", actor_id)); + let uri = format!( + "at://{}/app.bsky.feed.post/{}", + author_did, + parakeet_db::tid_util::encode_tid(*rkey) + ); + + // Get post view + if let Ok(Some(post_view)) = state.post_entity.get_by_uri(&uri, _viewer_did.as_deref()).await { + posts.push(post_view); + } + } + + // Build cursor from last result + let cursor = results.last().map(|(actor_id, rkey)| { + format!("{}:{}", actor_id, rkey) + }); + + Ok(Json(GetQuotesRes { + uri: query.uri, + cid: query.cid, + cursor, + posts, + })) +} + +#[derive(Debug, Deserialize)] +pub struct GetRepostedByQuery { + pub uri: String, + pub cid: Option, + pub limit: Option, + pub cursor: Option, +} + +#[derive(Debug, Serialize)] +pub struct GetRepostedByRes { + pub uri: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub cid: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cursor: Option, + pub reposted_by: Vec, +} + +pub async fn get_reposted_by( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let mut conn = state.pool.get().await?; + let _viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + let limit = query.limit.unwrap_or(50).clamp(1, 100); + + // Parse post URI to extract actor_id and rkey + let parts: Vec<&str> = query.uri.trim_start_matches("at://").split('/').collect(); + if parts.len() < 3 { + return Ok(Json(GetRepostedByRes { + uri: query.uri, + cid: query.cid, + cursor: None, + reposted_by: vec![], + })); + } + + let post_did = parts[0]; + let post_rkey_str = parts[2]; + + // Resolve DID to actor_id + let post_actor_id = state.profile_entity + .resolve_identifier(post_did) + .await + .map_err(|_| crate::common::errors::Error::not_found())?; + + // Decode rkey + let post_rkey = parakeet_db::tid_util::decode_tid(post_rkey_str) + .map_err(|_| crate::common::errors::Error::invalid_request(Some("Invalid rkey".to_string())))?; + + // Parse cursor + let cursor_rkey = query.cursor.as_ref() + .and_then(|c| c.parse::().ok()); + + // Get reposted by from database using PostEntity + let results = state.post_entity.get_reposted_by(post_actor_id, post_rkey, cursor_rkey, limit).await + .map_err(|e| crate::common::errors::Error::server_error(Some(&e.to_string())))?; + + // Convert to ProfileViews + let actor_ids: Vec = results.iter().map(|(actor_id, _)| *actor_id).collect(); + let reposted_by = state.profile_entity.get_profile_views(&actor_ids).await; + + // Build cursor from last result + let cursor = results.last().map(|(_, rkey)| rkey.to_string()); + + Ok(Json(GetRepostedByRes { + uri: query.uri, + cid: query.cid, + cursor, + reposted_by, + })) +} \ No newline at end of file diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs b/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs new file mode 100644 index 00000000..50a92682 --- /dev/null +++ b/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs @@ -0,0 +1,170 @@ +use axum::extract::{Query, State}; +use axum::Json; +use lexica::app_bsky::feed::{BlockedAuthor, PostView, ThreadViewPost, ThreadViewPostType, ThreadgateView}; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; + +use crate::entities::post::ThreadItem; +use crate::common::errors::{Error, XrpcResult}; +use crate::common::auth::{AtpAcceptLabelers, AtpAuth}; +use crate::GlobalState; + +use super::helpers::postview_to_tvpt; + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct GetPostThreadQuery { + pub uri: String, + pub depth: Option, + pub parent_height: Option, +} + +#[derive(Debug, Serialize)] +pub struct GetPostThreadRes { + pub thread: ThreadViewPostType, + #[serde(skip_serializing_if = "Option::is_none")] + pub threadgate: Option, +} + +pub async fn get_post_thread( + State(state): State, + AtpAcceptLabelers(_labelers): AtpAcceptLabelers, + maybe_auth: Option, + Query(query): Query, +) -> XrpcResult> { + let mut conn = state.pool.get().await?; + let viewer_did = maybe_auth.as_ref().map(|auth| auth.0.clone()); + + let uri = query.uri.clone(); // TODO: Normalize URI if needed + let depth = query.depth.unwrap_or(6).clamp(0, 1000); + let parent_height = query.parent_height.unwrap_or(80).clamp(0, 1000); + + // Get the root post using PostEntity + let root = state.post_entity.get_by_uri(&uri, viewer_did.as_deref()).await? + .ok_or_else(Error::not_found)?; + + let threadgate = root.threadgate.clone(); + + // Check if author is blocked + if let Some(viewer) = &root.author.viewer { + if viewer.blocked_by || viewer.blocking.is_some() { + return Ok(Json(GetPostThreadRes { + thread: ThreadViewPostType::Blocked { + uri, + blocked: true, + author: Box::new(BlockedAuthor { + did: root.author.did.clone(), + viewer: root.author.viewer.clone(), + }), + }, + threadgate, + })); + } + } + + // Extract anchor URI parts + let parts: Vec<&str> = uri.trim_start_matches("at://").split('/').collect(); + if parts.len() < 3 { + return Err(Error::invalid_request(Some("Invalid URI format".to_string()))); + } + let anchor_did = parts[0]; + let anchor_rkey_base32 = parts[2]; + + // Get actor_id using ProfileEntity + let anchor_actor_id = state.profile_entity.resolve_identifier(anchor_did).await + .map_err(|_| Error::actor_not_found(anchor_did))?; + + let anchor_rkey = parakeet_db::tid_util::decode_tid(anchor_rkey_base32) + .map_err(|_| Error::invalid_request(Some("Invalid rkey".to_string())))?; + + // Get root info for parent query filtering + let root_uri = root.record.get("reply") + .and_then(|reply| reply.get("root")) + .and_then(|root| root.get("uri")) + .and_then(|uri| uri.as_str()) + .map(|s| s.to_string()); + + let root_info = if let Some(ref root_uri_str) = root_uri { + let root_parts: Vec<&str> = root_uri_str.trim_start_matches("at://").split('/').collect(); + if root_parts.len() >= 3 { + let root_did = root_parts[0]; + let root_rkey_base32 = root_parts[2]; + + let root_actor_id = state.profile_entity.resolve_identifier(root_did).await + .map_err(|_| Error::actor_not_found(root_did))?; + + let root_rkey = parakeet_db::tid_util::decode_tid(root_rkey_base32).ok(); + root_rkey.map(|rkey| (root_actor_id, rkey)) + } else { + None + } + } else { + Some((anchor_actor_id, anchor_rkey)) + }; + + // Query thread data using PostEntity + let parents = if let Some((root_actor_id, root_rkey)) = root_info { + state.post_entity.get_thread_parents( + anchor_actor_id, + anchor_rkey, + parent_height as i32, + root_actor_id, + root_rkey + ).await + .map_err(|e| Error::server_error(Some(&e.to_string())))? + } else { + Vec::new() + }; + + let children = state.post_entity.get_thread_children(anchor_actor_id, anchor_rkey, depth as i32).await + .map_err(|e| Error::server_error(Some(&e.to_string())))?; + + // Build all URIs to hydrate + let mut all_uris = Vec::new(); + + // Build URIs from actor_id/rkey for parents + for item in &parents { + let did = state.profile_entity.get_did_by_id(item.actor_id).await + .unwrap_or_else(|_| format!("did:plc:unknown{}", item.actor_id)); + let rkey_str = parakeet_db::tid_util::encode_tid(item.rkey); + all_uris.push(format!("at://{}/app.bsky.feed.post/{}", did, rkey_str)); + } + + // Build URIs from actor_id/rkey for children + for item in &children { + let did = state.profile_entity.get_did_by_id(item.actor_id).await + .unwrap_or_else(|_| format!("did:plc:unknown{}", item.actor_id)); + let rkey_str = parakeet_db::tid_util::encode_tid(item.rkey); + all_uris.push(format!("at://{}/app.bsky.feed.post/{}", did, rkey_str)); + } + + // Get all posts using PostEntity + let posts_map = state.post_entity.get_by_uris(all_uris, viewer_did.as_deref()).await + .unwrap_or_default(); + + // Build thread structure + let thread = build_thread_structure(root, parents, children, posts_map); + + Ok(Json(GetPostThreadRes { + thread, + threadgate, + })) +} + +fn build_thread_structure( + root: PostView, + parents: Vec, + children: Vec, + posts_map: HashMap, +) -> ThreadViewPostType { + // Convert root to ThreadViewPost + let root_tvp = ThreadViewPost { + post: root, + parent: None, + replies: Vec::new(), // Will be populated + }; + + // TODO: Build full thread structure with parents and children + // For now, return simplified version + ThreadViewPostType::Post(Box::new(root_tvp)) +} \ No newline at end of file -- 2.51.2