diff --git a/src/jetstream.rs b/src/jetstream.rs index 93c4074..91db6c8 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -1,13 +1,60 @@ use async_trait::async_trait; use atproto_jetstream::{EventHandler, JetstreamEvent}; +use serde::Serialize; use std::sync::Arc; -pub struct MyEventHandler; +use crate::jetstream; + +pub struct MyEventHandler { + pub pool: sqlx::SqlitePool, +} #[async_trait] impl EventHandler for MyEventHandler { async fn handle_event(&self, event: Arc) -> anyhow::Result<()> { println!("Received event: {:?}", event); + + match event.as_ref() { + JetstreamEvent::Commit { + did, + time_us, + kind, + commit, + } => { + println!("hello: {}", commit.record) + } + JetstreamEvent::Delete { + did, + time_us, + kind, + commit, + } => {} + JetstreamEvent::Delete { + did, + time_us, + kind, + commit, + } => {} + JetstreamEvent::Identity { + did, + time_us, + kind, + identity, + } => {} + JetstreamEvent::Account { + did, + time_us, + kind, + account, + } => {} + } + + // let result = sqlx::query("INSERT INTO follows (subject, rkey) VALUES ($1, $2)") + // .bind(&follow.subject) + // .bind(record.rkey.to_string()) + // .execute(&self.pool) + // .await; + Ok(()) } diff --git a/src/main.rs b/src/main.rs index c4f908b..50e82af 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,9 +4,7 @@ mod store; mod tap; mod xrpc; -use atproto_jetstream::{ - CancellationToken, Consumer, ConsumerTaskConfig, EventHandler, JetstreamEvent, -}; +use atproto_jetstream::{CancellationToken, Consumer, ConsumerTaskConfig}; use sqlx::migrate::MigrateDatabase; #[tokio::main] @@ -45,8 +43,11 @@ async fn main() -> anyhow::Result<()> { }; let consumer = Consumer::new(config); + + let handler = jetstream::MyEventHandler { pool: pool.clone() }; + consumer - .register_handler(std::sync::Arc::new(jetstream::MyEventHandler)) + .register_handler(std::sync::Arc::new(handler)) .await?; let cancellation_token = CancellationToken::new();