From 2c2e4e4d7dd28795589917374abc37ee13437b11 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Sat, 25 Apr 2026 15:09:25 +0300 Subject: [PATCH] fix author feed more --- src/main.rs | 207 ++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 170 insertions(+), 37 deletions(-) diff --git a/src/main.rs b/src/main.rs index 74cd1fe..2b16af7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -380,6 +380,32 @@ async fn get_post_thread( } } +fn profile_to_basic(full: serde_json::Value) -> serde_json::Value { + let mut basic = serde_json::json!({ + "did": full["did"], + "handle": full["handle"], + }); + if let Some(v) = full.get("displayName") { + basic["displayName"] = v.clone(); + } + if let Some(v) = full.get("avatar") { + basic["avatar"] = v.clone(); + } + if let Some(v) = full.get("viewer") { + basic["viewer"] = v.clone(); + } + if let Some(v) = full.get("labels") { + basic["labels"] = v.clone(); + } + if let Some(v) = full.get("associated") { + basic["associated"] = v.clone(); + } + if let Some(v) = full.get("createdAt") { + basic["createdAt"] = v.clone(); + } + basic +} + async fn get_post_view( app_state: &AppState, uri_str: &str, @@ -414,7 +440,10 @@ async fn get_post_view( let author_profile = get_profile_internal(&app_state, repo.did.as_str(), viewer_did).await?; - let mut viewer_state = serde_json::json!({}); + let mut viewer_state = serde_json::json!({ + "muted": false, + "blockedBy": false, + }); if let Some(viewer) = viewer_did { let like = app_state .hydrant @@ -448,20 +477,31 @@ async fn get_post_view( } } + let mut val_json = serde_json::to_value(&record.value).unwrap_or(serde_json::json!({})); + if val_json.get("$type").is_none() { + val_json["$type"] = serde_json::json!("app.bsky.feed.post"); + } + + let created_at = val_json + .get("createdAt") + .and_then(|v| v.as_str()) + .unwrap_or(""); + let mut post_view = serde_json::json!({ + "$type": "app.bsky.feed.defs#postView", "uri": uri_str, "cid": record.cid.to_string(), - "author": author_profile, - "record": record.value, + "author": profile_to_basic(author_profile), + "record": val_json, "replyCount": app_state.hydrant.backlinks.count(uri_str.to_string()).source("app.bsky.feed.post").run().await.unwrap_or(0), "repostCount": app_state.hydrant.backlinks.count(uri_str.to_string()).source("app.bsky.feed.repost").run().await.unwrap_or(0), "likeCount": app_state.hydrant.backlinks.count(uri_str.to_string()).source("app.bsky.feed.like").run().await.unwrap_or(0), - "indexedAt": chrono::Utc::now().to_rfc3339(), + "quoteCount": 0, + "indexedAt": if created_at.is_empty() { chrono::Utc::now().to_rfc3339() } else { created_at.to_string() }, "viewer": viewer_state, "labels": [], }); - let val_json = serde_json::to_value(&record.value).unwrap_or(serde_json::json!({})); if let Some(record_embed) = val_json.get("embed") { let t = record_embed .get("$type") @@ -512,6 +552,10 @@ async fn get_post_view( "labels": quoted_post["labels"], "indexedAt": quoted_post["indexedAt"], "embeds": quoted_post.get("embed").map(|e| vec![e]).unwrap_or_default(), + "likeCount": quoted_post["likeCount"], + "repostCount": quoted_post["repostCount"], + "replyCount": quoted_post["replyCount"], + "quoteCount": quoted_post["quoteCount"], } }); } @@ -685,7 +729,10 @@ async fn get_profile_internal( .map(|h| h.to_string()) .unwrap_or_else(|| did_str.to_string()); - let mut viewer_state = serde_json::json!({}); + let mut viewer_state = serde_json::json!({ + "muted": false, + "blockedBy": false, + }); if let Some(viewer) = viewer_did { if viewer != did_str { // following: does viewer follow did_str? @@ -780,8 +827,7 @@ async fn get_profile_internal( }); if let Some(rec) = profile_record { - let value = - serde_json::to_value(rec.value).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; + let value = serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); if let Some(obj) = profile.as_object_mut() { if let Some(display_name) = value.get("displayName") { obj.insert("displayName".to_string(), display_name.clone()); @@ -811,12 +857,14 @@ async fn get_profile_internal( serde_json::json!(app_state.cdn("banner", did_str, link)), ); } + if let Some(created_at) = value.get("createdAt") { + obj.insert("createdAt".to_string(), created_at.clone()); + } } } Ok(profile) } - #[derive(Deserialize)] struct GetAuthorFeedParams { actor: String, @@ -826,6 +874,8 @@ struct GetAuthorFeedParams { cursor: Option, #[serde(default)] filter: Option, + #[serde(default, rename = "includePins")] + include_pins: Option, } async fn get_author_feed( @@ -855,9 +905,17 @@ async fn get_author_feed( let limit = params.limit.unwrap_or(50).min(100); + // Note: Official AppView uses timestamp cursors. + // Our local index uses rkeys for list_records. + let rkey_cursor = if let Some(c) = ¶ms.cursor { + if c.len() < 20 { Some(c.as_str()) } else { None } + } else { + None + }; + // Get posts let posts_list = match repo - .list_records("app.bsky.feed.post", limit, true, params.cursor.as_deref()) + .list_records("app.bsky.feed.post", limit, false, rkey_cursor) .await { Ok(rl) => rl, @@ -869,12 +927,7 @@ async fn get_author_feed( // Get reposts let reposts_list = match repo - .list_records( - "app.bsky.feed.repost", - limit, - true, - params.cursor.as_deref(), - ) + .list_records("app.bsky.feed.repost", limit, false, rkey_cursor) .await { Ok(rl) => rl, @@ -894,6 +947,8 @@ async fn get_author_feed( } }; + let filter = params.filter.as_deref().unwrap_or("posts_with_replies"); + let mut all_items = Vec::new(); for rec in posts_list.records { let val_json = serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); @@ -904,11 +959,27 @@ async fn get_author_feed( .to_string(); // Filter logic - let filter = params.filter.as_deref().unwrap_or("posts_with_replies"); if filter == "posts_no_replies" && val_json.get("reply").is_some() { continue; } + if filter == "posts_and_author_threads" { + if let Some(reply) = val_json.get("reply") { + let root_uri_str = reply + .get("root") + .and_then(|r| r.get("uri")) + .and_then(|u| u.as_str()) + .unwrap_or(""); + if let Ok(root_uri) = jacquard_common::types::string::AtUri::new(root_uri_str) { + if root_uri.authority().as_str() != did.as_str() { + continue; + } + } else { + continue; + } + } + } + all_items.push((created_at, "post", rec)); } for rec in reposts_list.records { @@ -926,6 +997,25 @@ async fn get_author_feed( all_items.truncate(limit); let mut feed = Vec::new(); + + // Handle Pins + let mut pinned_uri_opt = None; + if params.include_pins.unwrap_or(false) && params.cursor.is_none() { + if let Ok(Some(profile_rec)) = repo.get_record("app.bsky.actor.profile", "self").await { + let prof_val = serde_json::to_value(profile_rec.value).unwrap_or(serde_json::json!({})); + if let Some(pinned_uri) = prof_val.get("pinnedPost").and_then(|v| v.as_str()) { + if let Ok(post_view) = + get_post_view(&app_state, pinned_uri, viewer_did.as_deref()).await + { + feed.push(serde_json::json!({ + "post": post_view, + })); + pinned_uri_opt = Some(pinned_uri.to_string()); + } + } + } + } + for (created_at, kind, rec) in &all_items { if *kind == "post" { let uri = format!( @@ -933,29 +1023,45 @@ async fn get_author_feed( did.as_str(), rec.rkey.as_str() ); + + // Skip if it was already included as a pin + if Some(&uri) == pinned_uri_opt.as_ref() { + continue; + } + let post = match get_post_view(&app_state, &uri, viewer_did.as_deref()).await { Ok(mut p) => { - p["author"] = author_profile.clone(); + p["author"] = profile_to_basic(author_profile.clone()); p } Err(_) => { + let mut val_json = + serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); + if val_json.get("$type").is_none() { + val_json["$type"] = serde_json::json!("app.bsky.feed.post"); + } serde_json::json!({ + "$type": "app.bsky.feed.defs#postView", "uri": uri, "cid": rec.cid.to_string(), - "author": author_profile.clone(), - "record": rec.value, + "author": profile_to_basic(author_profile.clone()), + "record": val_json, "replyCount": app_state.hydrant.backlinks.count(uri.clone()).source("app.bsky.feed.post").run().await.unwrap_or(0), "repostCount": app_state.hydrant.backlinks.count(uri.clone()).source("app.bsky.feed.repost").run().await.unwrap_or(0), "likeCount": app_state.hydrant.backlinks.count(uri.clone()).source("app.bsky.feed.like").run().await.unwrap_or(0), - "indexedAt": chrono::Utc::now().to_rfc3339(), + "quoteCount": 0, + "indexedAt": if created_at.is_empty() { chrono::Utc::now().to_rfc3339() } else { created_at.clone() }, + "viewer": { "muted": false, "blockedBy": false }, "labels": [], }) } }; - let val_json = serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); - let mut feed_item = serde_json::json!({ "post": post }); + let mut feed_item = serde_json::json!({ + "post": post, + }); + let val_json = serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); // Add reply context if it's a reply if let Some(reply) = val_json.get("reply") { if let Some(parent_uri) = reply @@ -998,7 +1104,7 @@ async fn get_author_feed( "post": post_view, "reason": { "$type": "app.bsky.feed.defs#reasonRepost", - "by": author_profile.clone(), + "by": profile_to_basic(author_profile.clone()), "indexedAt": created_at, } })); @@ -1007,7 +1113,9 @@ async fn get_author_feed( } } - let next_cursor = all_items.last().map(|i| i.0.clone()); + let next_cursor = all_items.last().map(|i| i.2.rkey.as_str().to_string()); + + feed.reverse(); Ok(Json(serde_json::json!({ "feed": feed, @@ -1951,11 +2059,13 @@ async fn get_timeline( Ok(record_list) => { for rec in record_list.records { let rkey = rec.rkey.as_str().to_string(); + let val_json = + serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); all_items.push(( did.clone(), "app.bsky.feed.post".to_string(), rkey, - rec.value, + val_json, )); } } @@ -1972,11 +2082,13 @@ async fn get_timeline( Ok(record_list) => { for rec in record_list.records { let rkey = rec.rkey.as_str().to_string(); + let val_json = + serde_json::to_value(&rec.value).unwrap_or(serde_json::json!({})); all_items.push(( did.clone(), "app.bsky.feed.repost".to_string(), rkey, - rec.value, + val_json, )); } } @@ -1986,10 +2098,22 @@ async fn get_timeline( } } - all_items.sort_by(|a, b| b.2.cmp(&a.2)); + // Sort by createdAt descending + all_items.sort_by(|a, b| { + let ca = a.3.get("createdAt").and_then(|v| v.as_str()).unwrap_or(""); + let cb = b.3.get("createdAt").and_then(|v| v.as_str()).unwrap_or(""); + cb.cmp(ca) + }); if let Some(c) = params.cursor { - all_items.retain(|item| item.2 < c); + all_items.retain(|item| { + let item_ca = item + .3 + .get("createdAt") + .and_then(|v| v.as_str()) + .unwrap_or(""); + item_ca < c.as_str() + }); } all_items.truncate(limit); @@ -1998,24 +2122,31 @@ async fn get_timeline( let mut next_cursor = None; let viewer_did = get_auth_did(&req); for (did, col, rkey, value) in all_items { - next_cursor = Some(rkey.clone()); + let created_at = value + .get("createdAt") + .and_then(|v| v.as_str()) + .unwrap_or("") + .to_string(); + next_cursor = Some(created_at.clone()); let uri = format!("at://{}/{}/{}", did.as_str(), col, rkey); - let val_json = serde_json::to_value(&value).unwrap_or(serde_json::json!({})); - if col == "app.bsky.feed.post" { match get_post_view(&app_state, &uri, viewer_did.as_deref()).await { - Ok(post_view) => feed.push(serde_json::json!({ "post": post_view })), + Ok(mut post_view) => { + post_view["author"] = profile_to_basic(post_view["author"].clone()); + feed.push(serde_json::json!({ "post": post_view })); + } Err(e) => tracing::warn!("failed to get post view for {uri}: {e}"), } } else if col == "app.bsky.feed.repost" { - if let Some(subject_uri) = val_json + if let Some(subject_uri) = value .get("subject") .and_then(|s| s.get("uri")) .and_then(|u| u.as_str()) { match get_post_view(&app_state, subject_uri, viewer_did.as_deref()).await { - Ok(post_view) => { + Ok(mut post_view) => { + post_view["author"] = profile_to_basic(post_view["author"].clone()); match get_profile_internal(&app_state, did.as_str(), viewer_did.as_deref()) .await { @@ -2024,8 +2155,8 @@ async fn get_timeline( "post": post_view, "reason": { "$type": "app.bsky.feed.defs#reasonRepost", - "by": reposter_profile, - "indexedAt": val_json.get("createdAt").and_then(|c| c.as_str()).map(|s| s.to_string()).unwrap_or_else(|| chrono::Utc::now().to_rfc3339()), + "by": profile_to_basic(reposter_profile), + "indexedAt": if created_at.is_empty() { chrono::Utc::now().to_rfc3339() } else { created_at }, } })); } @@ -2044,6 +2175,8 @@ async fn get_timeline( } } + feed.reverse(); + Ok(Json(serde_json::json!({ "feed": feed, "cursor": next_cursor -- 2.51.2