Something went wrong. Try again.
A personal app view to see Bsky posts of your followers (for when their app view goes down)
Something went wrong. Try again.
11 kB · 343 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344use crate::{external_data, records};use ::chrono::{DateTime, Utc};use async_trait::async_trait;use atproto_jetstream::{ CancellationToken, Consumer, ConsumerTaskConfig, JetstreamEventCommit, JetstreamEventDelete,};use atproto_jetstream::{EventHandler, JetstreamEvent};use jacquard::api::app_bsky::feed::post::Post;use jacquard::api::app_bsky::feed::repost::Repost;use jacquard::api::app_bsky::graph::follow::Follow;use jacquard_identity::JacquardResolver;use serde::Deserialize;use std::sync::Arc;use std::time::{Duration, SystemTime, UNIX_EPOCH};use std::vec;use tokio::sync::RwLock;
use crate::store::{self};
pub async fn run_posts_jetstream( users_did: &String, jetstream_host: &String, resolver: JacquardResolver, pool: &sqlx::SqlitePool,) { let mut dids = store::get_follows(pool).await; dids.push(users_did.to_string());
let mut config = ConsumerTaskConfig { dids: Vec::new(), user_agent: "my-appview/1.0".to_string(), compression: false, zstd_dictionary_location: String::new(), jetstream_hostname: jetstream_host.to_string(), collections: vec![ "app.bsky.feed.post".to_string(), "app.bsky.feed.repost".to_string(), "app.bsky.graph.follow".to_string(), ], max_message_size_bytes: None, cursor: None, require_hello: false, };
let time_back = Duration::from_secs(3 * 60 * 60); if let Some(start) = SystemTime::now().checked_sub(time_back) { match start.duration_since(UNIX_EPOCH) { Ok(t) => { config.cursor = Some(t.as_millis() as i64); } Err(e) => println!("parsing cursor time: {e}"), } }
let cancellation_token = CancellationToken::new();
let consumer = Consumer::new(config.clone());
let following = Arc::new(RwLock::new(dids));
let handler = Handler::new(pool.clone(), resolver, users_did.to_string(), following);
consumer .register_handler(std::sync::Arc::new(handler)) .await .unwrap();
consumer .run_background(cancellation_token.clone()) .await .unwrap();}
#[derive(Clone)]pub struct Handler { pool: sqlx::SqlitePool, resolver: JacquardResolver, following: Arc<RwLock<Vec<String>>>, users_did: String,}
impl Handler { pub fn new( pool: sqlx::SqlitePool, resolver: JacquardResolver, users_did: String, following: Arc<RwLock<Vec<String>>>, ) -> Self { Self { pool: pool, users_did: users_did, following: following, resolver: resolver, } }}
#[async_trait]impl EventHandler for Handler { async fn handle_event(&self, event: Arc<JetstreamEvent>) -> anyhow::Result<()> { let mut follows = self.following.write().await;
match event.as_ref() { JetstreamEvent::Commit { did, time_us: _, kind: _, commit, } => { if !follows.contains(did) { return Ok(()); }
let did_doc = external_data::get_did_doc(self.resolver.clone(), &self.pool, did).await; match did_doc { Ok(doc) => { let _ = external_data::get_profile_from_pds( self.pool.clone(), did.clone(), doc.pds_endpoint, ) .await .map_err(|err| println!("caching did doc {err}")); } Err(e) => println!("caching did doc {e}"), }
match commit.collection.as_str() { "app.bsky.feed.post" => { handle_new_post_event(&self.pool, did.to_string(), commit).await } "app.bsky.feed.repost" => { handle_new_repost_event( self.resolver.clone(), &self.pool, did.to_string(), commit, ) .await }
"app.bsky.graph.follow" => { if did.to_string() != self.users_did { return Ok(()); } let subject = handle_new_follow_event(&self.pool, commit).await; follows.push(subject); } _ => { println!("it was something else {}", commit.collection); return Ok(()); } } } JetstreamEvent::Delete { did, time_us: _, kind: _, commit, } => { if !follows.contains(did) { return Ok(()); } match commit.collection.as_str() { "app.bsky.feed.post" => handle_delete_post_event(&self.pool, commit).await, "app.bsky.feed.repost" => { handle_delete_repost_event(&self.pool, did.to_string(), commit).await } "app.bsky.graph.follow" => { if did.to_string() != self.users_did { return Ok(()); }
// TODO: delete posts for that user
let subject = handle_delete_follow_event(&self.pool, commit).await; if let Some(index) = follows.iter().position(|value| *value == subject) { follows.remove(index); } } _ => { println!("it was something else {}", commit.collection); } } } _ => {} }
Ok(()) }
fn handler_id(&self) -> &str { "handler" }}
async fn handle_new_post_event( pool: &sqlx::SqlitePool, did: String, commit: &JetstreamEventCommit,) { match Post::deserialize(commit.record.clone()) { Ok(rec) => { let mut post = store::PostRecord { created: Utc::now().timestamp_millis(), indexed: Utc::now(), author: did.clone(), rkey: commit.rkey.to_string(), cid: commit.cid.to_string(), subject: "".to_string(), raw_blob: Vec::new(), subject_raw_blob: Vec::new(), };
match serde_json::to_vec(&commit.record) { Ok(blob) => { post.raw_blob = blob; } Err(e) => { println!("parsing commit record: {e}") } }
let created_at = DateTime::parse_from_rfc3339(&rec.created_at.as_str()); match created_at { Ok(ca) => { post.created = ca.to_utc().timestamp_millis(); post.indexed = ca.to_utc(); } Err(e) => { println!("parsing created at: {e}"); } }
store::insert_post(pool, post).await; } Err(err) => println!("parsing post record: {}", err), }}
async fn handle_delete_post_event(pool: &sqlx::SqlitePool, commit: &JetstreamEventDelete) { store::delete_post(pool, commit.rkey.clone()).await;}
async fn handle_new_repost_event( resolver: JacquardResolver, pool: &sqlx::SqlitePool, did: String, commit: &JetstreamEventCommit,) { match Repost::deserialize(commit.record.clone()) { Ok(rec) => { let mut repost = store::PostRecord { created: Utc::now().timestamp_millis(), indexed: Utc::now(), author: did, rkey: commit.rkey.to_string(), cid: commit.cid.to_string(), subject: rec.subject.uri.to_string(), raw_blob: Vec::new(), subject_raw_blob: Vec::new(), };
match serde_json::to_vec(&commit.record) { Ok(blob) => { repost.raw_blob = blob; } Err(e) => { println!("parsing commit record: {e}") } }
let did = rec.subject.uri.authority(); let rkey = match rec.subject.uri.rkey() { Some(r) => r.0.to_string(), _ => "".to_string(), };
match external_data::get_did_doc(resolver, &pool, &did.to_string()).await { Ok(doc) => { match records::get_post_from_pds(doc.did, doc.pds_endpoint, rkey).await { Ok(raw) => match serde_json::to_vec(&raw.value) { Ok(blob) => repost.subject_raw_blob = blob, Err(e) => { println!("parsing raw post: {e}") } }, Err(e) => println!("getting subject raw blob {e}"), } } Err(e) => println!("caching did doc {e}"), }
let created_at = DateTime::parse_from_rfc3339(&rec.created_at.as_str()); match created_at { Ok(ca) => { repost.created = ca.to_utc().timestamp_millis(); } Err(e) => { println!("parsing created at: {e}"); } }
store::insert_repost(pool, repost).await; } Err(err) => println!("parsing repost record: {}", err), }}
async fn handle_delete_repost_event( pool: &sqlx::SqlitePool, did: String, commit: &JetstreamEventDelete,) { store::delete_repost(pool, did, commit.rkey.clone()).await;}
async fn handle_new_follow_event(pool: &sqlx::SqlitePool, commit: &JetstreamEventCommit) -> String { match Follow::deserialize(commit.record.clone()) { Ok(rec) => { let follow_record = store::FollowRecord { subject: rec.subject.to_string(), };
store::insert_follow_record(pool, &follow_record, commit.rkey.to_string()).await; return follow_record.subject; } Err(err) => println!("parsing post record: {}", err), }
return "".to_string();}
async fn handle_delete_follow_event( pool: &sqlx::SqlitePool, commit: &JetstreamEventDelete,) -> String { let did = store::get_follow_by_rkey(pool, &commit.rkey.to_string()).await; store::delete_follow_record(pool, commit.rkey.to_string()).await;
return did;}