use 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>>, users_did: String, } impl Handler { pub fn new( pool: sqlx::SqlitePool, resolver: JacquardResolver, users_did: String, following: Arc>>, ) -> 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) -> 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; }