diff --git a/docker-compose.yaml b/docker-compose.yaml index b6ba165..20dffd2 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -15,6 +15,8 @@ services: USERS_DID: ${USERS_DID} TAP_URL: ${TAP_URL} DATABASE_URL: ${DATABASE_URL} + depends_on: + - tap tap: container_name: tap diff --git a/readme.md b/readme.md index ec68b25..1cac005 100644 --- a/readme.md +++ b/readme.md @@ -18,10 +18,10 @@ It is not (at the moment) supposed to be a complex appview, just a way for a use - [x] Run as part of a docker-compose file - [x] Configure to be in `Dynamically Configured` mode - [x] Configure the filters to be for `app.bsky.*` (maybe limit this just to be the ones needed) -- [ ] - App start up configuration for user +- [x] - App start up configuration for user - [x] Get users did from config - - [ ] Fetch and store that users follows - - [ ] For each user call the `/repos/add` endpoint on tap to add the follow to be tracked + - [x] Fetch and store that users follows + - [x] For each user call the `/repos/add` endpoint on tap to add the follow to be tracked - [x] Call `/repos/add` for the user of the appview - [ ] - Handle the events for tracked users - [ ] Ignore anything other than the post lexicon types for the follows (work out what lexicon types need to be used for the user using the appview) diff --git a/src/tap.rs b/src/tap.rs index 3d6dc7c..50f2178 100644 --- a/src/tap.rs +++ b/src/tap.rs @@ -1,6 +1,6 @@ use atproto_tap::{RecordAction, RecordEvent, TapClient, TapEvent, connect_to}; -use axum::response::sse::Event; use serde::Deserialize; +use sqlx::Row; use tokio_stream::StreamExt; #[derive(Deserialize)] @@ -13,11 +13,35 @@ pub async fn add_repo(did: &String) -> anyhow::Result<()> { let client = TapClient::new(tap_url.as_str(), Some("password".to_string())); // Add repositories to track - client.add_repos(&[did.as_str()]).await?; + let res = client.add_repos(&[did.as_str()]).await; + match res { + Err(e) => { + println!("failed to add repo with tap: {e}"); + } + Ok(()) => { + println!("added repo {did}"); + } + } Ok(()) } +pub async fn remove_repo(did: &String) { + let tap_url = std::env::var("TAP_URL").unwrap_or("localhost:2480".to_string()); + let client = TapClient::new(tap_url.as_str(), Some("password".to_string())); + + let res = client.remove_repos(&[did.as_str()]).await; + + match res { + Err(e) => { + println!("failed to remove repo with tap: {e}"); + } + Ok(()) => { + println!("removed repo {did}"); + } + } +} + pub async fn run_tap(users_did: String, pool: &sqlx::SqlitePool) { let tap_url = std::env::var("TAP_URL").unwrap_or("localhost:2480".to_string()); let mut stream = connect_to(tap_url.as_str()); @@ -56,7 +80,7 @@ async fn handle_user_follow_event(record: &RecordEvent, pool: &sqlx::SqlitePool) RecordAction::Create => { let follow: FollowRecord = record.parse_record().unwrap(); let result = sqlx::query("INSERT INTO follows (subject, rkey) VALUES ($1, $2)") - .bind(follow.subject) + .bind(&follow.subject) .bind(record.rkey.to_string()) .execute(pool) .await; @@ -64,8 +88,16 @@ async fn handle_user_follow_event(record: &RecordEvent, pool: &sqlx::SqlitePool) if result.is_err() { println!("Error inserting follow into the database: {result:?}"); } + + // track the follow anyway + _ = add_repo(&follow.subject).await; } RecordAction::Delete => { + let follow_subject = get_follow_by_rkey(&record.rkey.to_string(), pool).await; + if follow_subject != "" { + _ = remove_repo(&follow_subject).await; + } + let result = sqlx::query("DELETE FROM follows WHERE rkey = $1") .bind(record.rkey.to_string()) .execute(pool) @@ -79,10 +111,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::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) + .fetch_one(pool) + .await; + + match result { + Ok(row) => { + if row.len() == 0 { + println!("did not find subject in follows with rkey: {rkey}"); + return "".to_string(); + } + let subject = row.get::(0); + + return subject; + } + Err(e) => { + println!("error getting follow {e}"); + return "".to_string(); + } + } +}