From 86ed1a355b6ef626a86d541f1ecfe84cace02814 Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Thu, 23 Apr 2026 14:00:30 +0100 Subject: [PATCH] store the posts (basic for now) in the database for those accounts that are tracked as follows Signed-off-by: Will Andrews --- Cargo.lock | 118 ++++++++++++++++++++++++++++ Cargo.toml | 4 +- migrations/20260422183446_posts.sql | 9 +++ src/tap.rs | 76 ++++++++++++++++-- 4 files changed, 200 insertions(+), 7 deletions(-) create mode 100644 migrations/20260422183446_posts.sql diff --git a/Cargo.lock b/Cargo.lock index 8407175..90d1826 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8,6 +8,15 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" +[[package]] +name = "android_system_properties" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "819e7219dbd41043ac279b19830f2efc897156490d7fd6ea916720117ee66311" +dependencies = [ + "libc", +] + [[package]] name = "anyhow" version = "1.0.102" @@ -326,7 +335,11 @@ version = "0.4.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c673075a2e0e5f4a1dde27ce9dee1ea4558c7ffe648f576438a20ca1d2acc4b0" dependencies = [ + "iana-time-zone", + "js-sys", "num-traits", + "wasm-bindgen", + "windows-link", ] [[package]] @@ -582,6 +595,15 @@ dependencies = [ "zeroize", ] +[[package]] +name = "deranged" +version = "0.5.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7cd812cc2bc1d69d4764bd80df88b4317eaef9e773c75226407d9bc0876b211c" +dependencies = [ + "powerfmt", +] + [[package]] name = "digest" version = "0.10.7" @@ -1260,6 +1282,30 @@ dependencies = [ "windows-registry", ] +[[package]] +name = "iana-time-zone" +version = "0.1.65" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e31bc9ad994ba00e440a8aa5c9ef0ec67d5cb5e5cb0cc7f8b744a35b389cc470" +dependencies = [ + "android_system_properties", + "core-foundation-sys", + "iana-time-zone-haiku", + "js-sys", + "log", + "wasm-bindgen", + "windows-core", +] + +[[package]] +name = "iana-time-zone-haiku" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f" +dependencies = [ + "cc", +] + [[package]] name = "icu_collections" version = "2.2.0" @@ -1630,10 +1676,12 @@ dependencies = [ "anyhow", "atproto-tap", "axum", + "chrono", "dotenv", "serde", "serde_json", "sqlx", + "time", "tokio", "tokio-stream", ] @@ -1680,6 +1728,12 @@ dependencies = [ "zeroize", ] +[[package]] +name = "num-conv" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c6673768db2d862beb9b39a78fdcb1a69439615d5794a1be50caa9bc92c81967" + [[package]] name = "num-integer" version = "0.1.46" @@ -1886,6 +1940,12 @@ dependencies = [ "zerovec", ] +[[package]] +name = "powerfmt" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" + [[package]] name = "ppv-lite86" version = "0.2.21" @@ -2546,6 +2606,7 @@ checksum = "ee6798b1838b6a0f69c007c133b8df5866302197e404e8b6ee8ed3e3a5e68dc6" dependencies = [ "base64", "bytes", + "chrono", "crc", "crossbeam-queue", "either", @@ -2622,6 +2683,7 @@ dependencies = [ "bitflags", "byteorder", "bytes", + "chrono", "crc", "digest 0.10.7", "dotenvy", @@ -2663,6 +2725,7 @@ dependencies = [ "base64", "bitflags", "byteorder", + "chrono", "crc", "dotenvy", "etcetera", @@ -2697,6 +2760,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2d12fe70b2c1b4401038055f90f151b78208de1f9f89a7dbfd41587a10c3eea" dependencies = [ "atoi", + "chrono", "flume", "futures-channel", "futures-core", @@ -2834,6 +2898,25 @@ dependencies = [ "syn", ] +[[package]] +name = "time" +version = "0.3.47" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "743bd48c283afc0388f9b8827b976905fb217ad9e647fae3a379a9283c4def2c" +dependencies = [ + "deranged", + "num-conv", + "powerfmt", + "serde_core", + "time-core", +] + +[[package]] +name = "time-core" +version = "0.1.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7694e1cfe791f8d31026952abf09c69ca6f6fa4e1a1229e18988f06a04a12dca" + [[package]] name = "tinystr" version = "0.8.3" @@ -3292,6 +3375,41 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72069c3113ab32ab29e5584db3c6ec55d416895e60715417b5b883a357c3e471" +[[package]] +name = "windows-core" +version = "0.62.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb" +dependencies = [ + "windows-implement", + "windows-interface", + "windows-link", + "windows-result", + "windows-strings", +] + +[[package]] +name = "windows-implement" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "windows-interface" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "windows-link" version = "0.2.1" diff --git a/Cargo.toml b/Cargo.toml index 1c67a6e..664523b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,9 +7,11 @@ edition = "2024" anyhow = "1.0.102" atproto-tap = "0.14.5" axum = "0.8.9" +chrono = "0.4.44" dotenv = "0.15.0" serde = { version = "1.0.228", features = ["derive"] } serde_json = "1.0.149" -sqlx = { version = "0.8.6", features = ["runtime-tokio-native-tls", "sqlite"] } +sqlx = { version = "0.8.6", features = ["chrono", "runtime-tokio-native-tls", "sqlite"] } +time = "0.3.47" tokio = { version = "1.52.1", features = ["full"] } tokio-stream = "0.1.18" diff --git a/migrations/20260422183446_posts.sql b/migrations/20260422183446_posts.sql new file mode 100644 index 0000000..6f03e1f --- /dev/null +++ b/migrations/20260422183446_posts.sql @@ -0,0 +1,9 @@ +-- Create a posts table +CREATE TABLE IF NOT EXISTS posts +( + 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 +); diff --git a/src/tap.rs b/src/tap.rs index 50f2178..a549d97 100644 --- a/src/tap.rs +++ b/src/tap.rs @@ -1,6 +1,7 @@ +use ::chrono::{DateTime, ParseError, TimeZone, Utc}; use atproto_tap::{RecordAction, RecordEvent, TapClient, TapEvent, connect_to}; -use serde::Deserialize; -use sqlx::Row; +use serde::{Deserialize, Serialize}; +use sqlx::{FromRow, Row, types::chrono}; use tokio_stream::StreamExt; #[derive(Deserialize)] @@ -78,7 +79,7 @@ pub async fn run_tap(users_did: String, pool: &sqlx::SqlitePool) { async fn handle_user_follow_event(record: &RecordEvent, pool: &sqlx::SqlitePool) { match record.action { RecordAction::Create => { - let follow: FollowRecord = record.parse_record().unwrap(); + let follow: FollowRecord = record.parse_record().unwrap(); // TODO - bad error handling there let result = sqlx::query("INSERT INTO follows (subject, rkey) VALUES ($1, $2)") .bind(&follow.subject) .bind(record.rkey.to_string()) @@ -111,10 +112,33 @@ async fn handle_user_follow_event(record: &RecordEvent, pool: &sqlx::SqlitePool) } } -async fn handle_post_event(record: &RecordEvent, _pool: &sqlx::SqlitePool) { +async fn handle_post_event(record: &RecordEvent, pool: &sqlx::SqlitePool) { match record.action { - RecordAction::Create => {} - RecordAction::Delete => {} + RecordAction::Create => { + let mut post = PostRecord { + created: Utc::now(), + indexed: Utc::now(), + author: record.did.clone().to_string(), + rkey: record.rkey.to_string(), + }; + + let tap_post: TapPost = record.parse_record().unwrap(); // TODO: better error handling here + + let created_at = DateTime::parse_from_rfc3339(&tap_post.created_at); + match created_at { + Ok(ca) => { + post.created = ca.to_utc(); + } + Err(e) => { + println!("parsing created at: {e}"); + } + } + + insert_post(post, pool).await; + } + RecordAction::Delete => { + delete_post(record.rkey.to_string(), pool).await; + } RecordAction::Update => {} } } @@ -141,3 +165,43 @@ async fn get_follow_by_rkey(rkey: &String, pool: &sqlx::SqlitePool) -> String { } } } +#[derive(Debug, Deserialize)] +struct TapPost { + #[serde(rename = "createdAt")] + created_at: String, +} + +#[derive(Debug, FromRow)] +struct PostRecord { + created: chrono::DateTime, + indexed: chrono::DateTime, + author: String, + rkey: 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)") + .bind(post.created) + .bind(post.indexed) + .bind(post.author) + .bind(post.rkey) + .execute(pool) + .await; + + if result.is_err() { + println!("Error inserting post 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) + .execute(pool) + .await; + + if result.is_err() { + println!("Error deleting post from the database: {result:?}"); + } +} -- 2.51.2