diff --git a/src/jetstream.rs b/src/jetstream.rs index 91db6c8..e9315e5 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -1,9 +1,9 @@ +use ::chrono::{DateTime, Utc}; use async_trait::async_trait; use atproto_jetstream::{EventHandler, JetstreamEvent}; -use serde::Serialize; use std::sync::Arc; -use crate::jetstream; +use crate::store; pub struct MyEventHandler { pub pool: sqlx::SqlitePool, @@ -20,15 +20,22 @@ impl EventHandler for MyEventHandler { time_us, kind, commit, - } => { - println!("hello: {}", commit.record) - } - JetstreamEvent::Delete { - did, - time_us, - kind, - commit, - } => {} + } => match commit.collection.as_str() { + "app.bsky.feed.post" => { + let post = store::PostRecord { + created: Utc::now(), + indexed: Utc::now(), + author: did.to_string(), + rkey: commit.rkey.to_string(), + cid: commit.cid.to_string(), + }; + + store::insert_post(post, &self.pool).await; + } + _ => { + println!("it was something else {}", commit.collection); + } + }, JetstreamEvent::Delete { did, time_us, @@ -49,12 +56,6 @@ impl EventHandler for MyEventHandler { } => {} } - // 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(()) }