diff --git a/src/ingest/ingest.rs b/src/ingest/ingest.rs new file mode 100644 index 0000000..247b1d2 --- /dev/null +++ b/src/ingest/ingest.rs @@ -0,0 +1,196 @@ +use std::{collections::VecDeque, sync::Arc}; + +use futures_util::future; +use ipld_core::ipld::Ipld; +use jacquard::{ + api::com_atproto::sync::subscribe_repos::{Commit, SubscribeReposMessage, Sync}, + types::string::Handle, +}; +use jacquard_repo::{BlockStore, MemoryBlockStore}; +use sqlx::{Pool, Postgres, query}; +use thiserror::Error; +use tokio::{sync::Mutex, task::JoinHandle}; + +use crate::{backfill::backfill, utils::ipld_json::ipld_to_json_value}; + +trait Ingest { + type Error; + async fn ingest(&self, conn: Arc>) -> Result<(), Self::Error>; +} + +#[derive(Debug, Error)] +enum CommitError { + #[error("Error parsing #commit event: {}", .0)] + ParseCarBytes(#[from] jacquard_repo::RepoError), +} + +impl Ingest for Commit<'_> { + type Error = CommitError; + async fn ingest(&self, conn: Arc>) -> Result<(), Self::Error> { + let car = jacquard_repo::car::parse_car_bytes(&self.blocks).await?; + let storage = Arc::new(MemoryBlockStore::new_from_blocks(car.blocks)); + + let ops = future::join_all(self.ops.clone().into_iter().map(|op| async { + // get block data by cid, or None if errors/not found + if let Some(cid) = &op.cid { + if let Ok(cid) = cid.0.to_ipld() + && let Ok(contents) = storage.get(&cid).await + && let Some(contents) = contents + && let Ok(val) = serde_ipld_dagcbor::from_slice::(&contents) + { + (op, ipld_to_json_value(&val).ok()) + } else { + (op, None) + } + } else { + (op, None) + } + })) + .await; + + future::join_all(ops.into_iter().map(|(op, val)| async { + let mut path = op.path.split("/"); + let Some(collection) = path.next() else { + eprintln!("Invalid path ({})", op.path.as_str()); + return; + }; + let Some(rkey) = path.next() else { + eprintln!("Invalid path ({})", op.path.as_str()); + return; + }; + // assert the path is only collection/rkey + if path.next().is_some() { + eprintln!("Invalid path ({})", op.path.as_str()); + return; + }; + match op.action.clone().as_str() { + "create" => { + let Some(cid) = op.cid.map(|x| x.0.to_string()) else { + eprintln!("Missing cid for create {}/{}", collection, rkey); + return; + }; + let Some(val) = val else { + eprintln!("Missing value for create {}/{}/{}", collection, rkey, cid); + return; + }; + if let Err(err) = query!( + "INSERT INTO records (collection, rkey, cid, record) + VALUES ($1, $2, $3, $4)", + collection, + rkey, + cid, + val + ) + .execute(&*conn) + .await + { + eprintln!("Error creating {}/{}/{}\n{}", collection, rkey, cid, err); + } else { + println!("wrote {}/{}/{} successfully.", collection, rkey, cid); + }; + } + "update" => { + let Some(cid) = op.cid.map(|x| x.0.to_string()) else { + eprintln!("Missing cid for update {}/{}", collection, rkey); + return; + }; + let Some(val) = val else { + eprintln!("Missing value for update {}/{}/{}", collection, rkey, cid); + return; + }; + if let Err(err) = query!( + "UPDATE records SET + collection = $1, + rkey = $2, + cid = $3, + record = $4 + WHERE + collection = $1 + and rkey = $2", + collection, + rkey, + cid, + val + ) + .execute(&*conn) + .await + { + eprintln!("Error updating {}/{}/{}\n{}", collection, rkey, cid, err); + } else { + println!("updated {}/{}/{} successfully.", collection, rkey, cid); + }; + } + "delete" => { + if let Err(err) = query!( + "DELETE FROM records WHERE + collection = $1 + and rkey = $2", + collection, + rkey, + ) + .execute(&*conn) + .await + { + eprintln!("Error deleting {}/{}\n{}", collection, rkey, err); + } else { + println!("deleted {}/{} successfully.", collection, rkey); + }; + } + _ => { + println!("missing #{} {:#?} {:#?}", op.action.as_str(), op, val) + } + } + })) + .await; + + Ok(()) + } +} + +impl Ingest for Sync<'_> { + type Error = crate::backfill::Error; + async fn ingest(&self, conn: Arc>) -> Result<(), Self::Error> { + backfill(conn, None).await + } +} + +pub fn ingest( + queue: Arc>>>, + conn: Arc>, +) -> JoinHandle<()> { + tokio::spawn(async move { + loop { + let Some(next) = queue.lock().await.pop_front() else { + continue; + }; + + match next { + SubscribeReposMessage::Commit(commit) => { + commit.ingest(conn.clone()).await.unwrap_or_else(|err| { + eprintln!("error handling #commit({}): {:?}", commit.clone().rev, err) + }) + } + SubscribeReposMessage::Sync(sync) => { + sync.ingest(conn.clone()).await.unwrap_or_else(|err| { + eprintln!("error handling #sync({}): {:?}", sync.clone().rev, err) + }) + } + SubscribeReposMessage::Identity(identity) => println!( + "ignoring #identity({}) event. has user migrated?", + identity.handle.unwrap_or(Handle::raw("handle.invalid")) + ), + SubscribeReposMessage::Account(account) => println!( + "ignoring #account({} {}) event. has user deactivated?", + account.active, + account.status.unwrap_or("unknown".into()) + ), + SubscribeReposMessage::Info(info) => { + println!("ignoring #info({}) event", info.name) + } + SubscribeReposMessage::Unknown(_) => { + println!("ignoring unknown event. is meview outdated?") + } + }; + } + }) +} diff --git a/src/ingest/mod.rs b/src/ingest/mod.rs new file mode 100644 index 0000000..effd7a9 --- /dev/null +++ b/src/ingest/mod.rs @@ -0,0 +1,4 @@ +pub mod ingest; +pub mod queue; +pub use self::ingest::ingest; +pub use self::queue::queue; diff --git a/src/ingest/queue.rs b/src/ingest/queue.rs new file mode 100644 index 0000000..cba95b4 --- /dev/null +++ b/src/ingest/queue.rs @@ -0,0 +1,94 @@ +use std::{collections::VecDeque, sync::Arc}; + +use futures_util::stream::StreamExt; +use jacquard::api::com_atproto::sync::subscribe_repos::{SubscribeRepos, SubscribeReposMessage}; +use jacquard::url::Url; +use jacquard::{common::xrpc::TungsteniteSubscriptionClient, xrpc::SubscriptionClient}; +use tokio::{ + sync::Mutex, + task::{self, JoinHandle}, +}; + +use crate::config; + +pub async fn queue() -> ( + Arc>>>, + JoinHandle<()>, +) { + let queue = Arc::new(Mutex::new(VecDeque::new())); + + // USER_SUBSCRIBE_URL is formatted as a domain + let uri = Url::parse(&format!("wss://{}/", config::USER_SUBSCRIBE_URL)) + .expect("Env var USER_SUBSCRIBE_URL should be formated as a domain."); + let client = TungsteniteSubscriptionClient::from_base_uri(uri); + let (_sink, mut messages) = client + .subscribe(&SubscribeRepos::new().build()) + .await + .expect("Could not subscribe to new events") + .into_stream(); + + let queue_clone = queue.clone(); + let handle = task::spawn(async move { + let queue = queue_clone; + + loop { + if let Some(msg) = messages.next().await { + let msg = match msg { + Ok(val) => val, + Err(err) => { + eprintln!("Warning: Websocket error: {} ({:?})", err, err.source()); + continue; + } + }; + + // filter messages by user did + // note that #identity #account #info and #unknown will probably be ignored + let ev = match msg.clone() { + SubscribeReposMessage::Commit(commit) => { + if commit.repo != *config::USER_DID { + continue; + } else { + SubscribeReposMessage::Commit(commit) + } + } + SubscribeReposMessage::Sync(sync) => { + if sync.did != *config::USER_DID { + continue; + } else { + SubscribeReposMessage::Sync(sync) + } + } + + SubscribeReposMessage::Identity(identity) => { + if identity.did != *config::USER_DID { + continue; + } else { + eprintln!( + "Warning: Recieved #identity event. Configuration may be out of date" + ); + SubscribeReposMessage::Identity(identity) + } + } + SubscribeReposMessage::Account(account) => { + if account.did != *config::USER_DID { + continue; + } else { + eprintln!( + "Warning: Recieved #account event. Account active: `{}`. Account status: `{}`", + account.active, + account.status.clone().unwrap_or("Unknown".into()) + ); + SubscribeReposMessage::Account(account) + } + } + SubscribeReposMessage::Info(info) => SubscribeReposMessage::Info(info), + SubscribeReposMessage::Unknown(data) => SubscribeReposMessage::Unknown(data), + }; + + queue.lock().await.push_back(ev); + } + } + }); + + (queue, handle) +} diff --git a/src/main.rs b/src/main.rs index 204f1f0..d198b58 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,10 +1,9 @@ -use sqlx::{Pool, Postgres}; - use crate::backfill::backfill; mod backfill; mod config; mod db; +mod ingest; mod utils; #[derive(Debug)] @@ -33,9 +32,11 @@ async fn main() -> Result<(), Error> { config::USER_SUBSCRIBE_URL.get().await, config::DATABASE_URL.get().await ); - let conn: Pool = db::conn().await; + let conn = db::conn().await; println!("Database connected and initialized"); + let (queue, queue_handle) = ingest::queue().await; + println!("Starting backfill"); let timer = std::time::Instant::now(); backfill(conn.clone(), Some(timer)) @@ -45,5 +46,20 @@ async fn main() -> Result<(), Error> { println!("Backfill complete. Took {:?}", timer.elapsed()); println!("Completed sucessfully!"); + + // Set up Ctrl-C handler + let (tx, rx) = tokio::sync::oneshot::channel(); + tokio::spawn(async move { + tokio::signal::ctrl_c().await.ok(); + let _ = tx.send(()); + }); + + println!("Handling new events. Ctrl+C to quit."); + let ingest_handle = ingest::ingest(queue, conn.clone()); + + let _ = rx.await; + queue_handle.abort(); + ingest_handle.abort(); + Ok(()) }