From dd7f414a5e4b269b4b0196d82c65fbd38cdcfb61 Mon Sep 17 00:00:00 2001 From: Will Date: Mon, 17 Aug 2026 21:25:33 +0100 Subject: [PATCH] handle follow creation events Signed-off-by: Will --- src/jetstream.rs | 29 +++++++++++++++++++++++++++++ src/main.rs | 5 ++++- 2 files changed, 33 insertions(+), 1 deletion(-) diff --git a/src/jetstream.rs b/src/jetstream.rs index e9315e5..5968d8a 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -1,12 +1,20 @@ use ::chrono::{DateTime, Utc}; use async_trait::async_trait; use atproto_jetstream::{EventHandler, JetstreamEvent}; +use serde::Deserialize; + use std::sync::Arc; use crate::store; pub struct MyEventHandler { pub pool: sqlx::SqlitePool, + pub users_did: String, +} + +#[derive(Deserialize)] +struct FollowRecord { + subject: String, } #[async_trait] @@ -32,6 +40,16 @@ impl EventHandler for MyEventHandler { store::insert_post(post, &self.pool).await; } + "app.bsky.graph.follow" => { + if did.to_string() == self.users_did { + match serde_json::from_value(commit.record.clone()) { + Ok(foo) => { + store_follow_record(&self.pool, foo, commit.rkey.to_string()).await + } + Err(err) => println!("parsing follow record {}", err), + } + } + } _ => { println!("it was something else {}", commit.collection); } @@ -63,3 +81,14 @@ impl EventHandler for MyEventHandler { "my-handler" } } + +async fn store_follow_record(pool: &sqlx::SqlitePool, record: FollowRecord, rkey: String) { + let result = sqlx::query("INSERT INTO follows (subject, rkey) VALUES ($1, $2)") + .bind(record.subject) + .bind(rkey) + .execute(pool) + .await; + if result.is_err() { + println!("Error inserting follow into the database: {result:?}"); + } +} diff --git a/src/main.rs b/src/main.rs index 50e82af..13a5a6d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -44,7 +44,10 @@ async fn main() -> anyhow::Result<()> { let consumer = Consumer::new(config); - let handler = jetstream::MyEventHandler { pool: pool.clone() }; + let handler = jetstream::MyEventHandler { + pool: pool.clone(), + users_did: users_did.clone(), + }; consumer .register_handler(std::sync::Arc::new(handler)) -- 2.51.2