From ac5ea372e79cbd19b19751888d6cb46b7edc30e7 Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Thu, 23 Apr 2026 21:03:42 +0100 Subject: [PATCH] track and store reposts Signed-off-by: Will Andrews --- docker-compose.yaml | 2 +- migrations/20260423190103_reposts.sql | 10 +++ src/tap.rs | 89 ++++++++++++++++++++++++++- 3 files changed, 98 insertions(+), 3 deletions(-) create mode 100644 migrations/20260423190103_reposts.sql diff --git a/docker-compose.yaml b/docker-compose.yaml index 20dffd2..e136a67 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -26,4 +26,4 @@ services: volumes: - ./data:/data environment: - TAP_COLLECTION_FILTERS: "app.bsky.feed.post,app.bsky.graph.*" + TAP_COLLECTION_FILTERS: "app.bsky.feed.post,app.bsky.feed.repost,app.bsky.graph.follow" diff --git a/migrations/20260423190103_reposts.sql b/migrations/20260423190103_reposts.sql new file mode 100644 index 0000000..d535dce --- /dev/null +++ b/migrations/20260423190103_reposts.sql @@ -0,0 +1,10 @@ +-- Create a reposts table +CREATE TABLE IF NOT EXISTS reposts +( + id INTEGER PRIMARY KEY NOT NULL, + created TIMESTAMP NOT NULL, + indexed TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + author TEXT NOT NULL, + rkey TEXT NOT NULL, + subject TEXT NOT NULL +); diff --git a/src/tap.rs b/src/tap.rs index a549d97..d02fe8a 100644 --- a/src/tap.rs +++ b/src/tap.rs @@ -1,6 +1,6 @@ -use ::chrono::{DateTime, ParseError, TimeZone, Utc}; +use ::chrono::{DateTime, Utc}; use atproto_tap::{RecordAction, RecordEvent, TapClient, TapEvent, connect_to}; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; use sqlx::{FromRow, Row, types::chrono}; use tokio_stream::StreamExt; @@ -62,6 +62,9 @@ pub async fn run_tap(users_did: String, pool: &sqlx::SqlitePool) { "app.bsky.feed.post" => { handle_post_event(record, pool).await; } + "app.bsky.feed.repost" => { + handle_repost_event(record, pool).await; + } _ => {} } @@ -143,6 +146,38 @@ async fn handle_post_event(record: &RecordEvent, pool: &sqlx::SqlitePool) { } } +async fn handle_repost_event(record: &RecordEvent, pool: &sqlx::SqlitePool) { + match record.action { + RecordAction::Create => { + let tap_post: TapRepost = record.parse_record().unwrap(); // TODO: better error handling here + + let mut repost = RepostRecord { + created: Utc::now(), + indexed: Utc::now(), + author: record.did.clone().to_string(), + rkey: record.rkey.to_string(), + subject: tap_post.subject.uri, + }; + + let created_at = DateTime::parse_from_rfc3339(&tap_post.created_at); + match created_at { + Ok(ca) => { + repost.created = ca.to_utc(); + } + Err(e) => { + println!("parsing created at: {e}"); + } + } + + insert_repost(repost, pool).await; + } + RecordAction::Delete => { + delete_repost(record.rkey.to_string(), pool).await; + } + RecordAction::Update => {} + } +} + async fn get_follow_by_rkey(rkey: &String, pool: &sqlx::SqlitePool) -> String { let result = sqlx::query("SELECT subject FROM follows where rkey = ? LIMIT 1") .bind(rkey) @@ -171,6 +206,18 @@ struct TapPost { created_at: String, } +#[derive(Debug, Deserialize)] +struct TapRepost { + #[serde(rename = "createdAt")] + created_at: String, + subject: Subject, +} + +#[derive(Debug, Deserialize)] +struct Subject { + uri: String, +} + #[derive(Debug, FromRow)] struct PostRecord { created: chrono::DateTime, @@ -180,6 +227,16 @@ struct PostRecord { // TODO: other fields like reply to etc } +#[derive(Debug, FromRow)] +struct RepostRecord { + created: chrono::DateTime, + indexed: chrono::DateTime, + author: String, + rkey: String, + subject: String, + // TODO: other fields like reply to etc +} + async fn insert_post(post: PostRecord, pool: &sqlx::SqlitePool) { let result = sqlx::query("INSERT INTO posts (created, indexed, author, rkey) VALUES ($1, $2, $3, $4)") @@ -195,6 +252,23 @@ async fn insert_post(post: PostRecord, pool: &sqlx::SqlitePool) { } } +async fn insert_repost(repost: RepostRecord, pool: &sqlx::SqlitePool) { + let result = sqlx::query( + "INSERT INTO reposts (created, indexed, author, rkey, subject) VALUES ($1, $2, $3, $4, $5)", + ) + .bind(repost.created) + .bind(repost.indexed) + .bind(repost.author) + .bind(repost.rkey) + .bind(repost.subject) + .execute(pool) + .await; + + if result.is_err() { + println!("Error inserting repost into the database: {result:?}"); + } +} + async fn delete_post(rkey: String, pool: &sqlx::SqlitePool) { let result = sqlx::query("DELETE FROM posts WHERE rkey = $1") .bind(rkey) @@ -205,3 +279,14 @@ async fn delete_post(rkey: String, pool: &sqlx::SqlitePool) { println!("Error deleting post from the database: {result:?}"); } } + +async fn delete_repost(rkey: String, pool: &sqlx::SqlitePool) { + let result = sqlx::query("DELETE FROM reposts WHERE rkey = $1") + .bind(rkey) + .execute(pool) + .await; + + if result.is_err() { + println!("Error deleting repost from the database: {result:?}"); + } +} -- 2.51.2