From 9610deac6b785addb58728f0e112f9cecf90758f Mon Sep 17 00:00:00 2001 From: afterlifepro Date: Wed, 28 Jan 2026 22:02:38 +0000 Subject: [PATCH] move from Arc>> to broadcast for channel system this creates a spmc architecture so we can create com.atproto.sync.subscribeRepos and implement jetstream etc --- src/ingest/handler.rs | 20 +++++++++++++++----- src/ingest/queue.rs | 26 ++++++++------------------ src/main.rs | 9 ++++++--- 3 files changed, 29 insertions(+), 26 deletions(-) diff --git a/src/ingest/handler.rs b/src/ingest/handler.rs index cf67c8e..eb1026a 100644 --- a/src/ingest/handler.rs +++ b/src/ingest/handler.rs @@ -1,4 +1,4 @@ -use std::{collections::VecDeque, sync::Arc}; +use std::sync::Arc; use futures_util::future; use ipld_core::ipld::Ipld; @@ -9,7 +9,7 @@ use jacquard::{ use jacquard_repo::{BlockStore, MemoryBlockStore}; use sqlx::{Pool, Postgres, query}; use thiserror::Error; -use tokio::{sync::Mutex, task::JoinHandle}; +use tokio::{sync::broadcast, task::JoinHandle}; use crate::{backfill::backfill, utils::ipld_json::ipld_to_json_value}; @@ -157,13 +157,23 @@ impl Ingest for Sync<'_> { } pub fn ingest( - queue: Arc>>>, + mut reciever: broadcast::Receiver>, conn: Arc>, ) -> JoinHandle<()> { tokio::spawn(async move { loop { - let Some(next) = queue.lock().await.pop_front() else { - continue; + let next = match reciever.recv().await { + Ok(val) => val, + Err(err) => match err { + broadcast::error::RecvError::Closed => { + eprintln!("Ingestion failed. Quitting"); + break; + } + broadcast::error::RecvError::Lagged(skipped) => { + eprintln!("Warning: lagging behind. Skipping {skipped} messages"); + continue; + } + }, }; match next { diff --git a/src/ingest/queue.rs b/src/ingest/queue.rs index 291f18a..a3aa224 100644 --- a/src/ingest/queue.rs +++ b/src/ingest/queue.rs @@ -1,23 +1,17 @@ use std::thread; use std::time::Duration; -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 tokio::sync::broadcast; +use tokio::task::{self, JoinHandle}; use crate::config; -pub async fn queue() -> ( - Arc>>>, - JoinHandle<()>, -) { - let queue = Arc::new(Mutex::new(VecDeque::new())); +pub async fn queue<'a>(tx: broadcast::Sender>) -> 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)) @@ -29,10 +23,8 @@ pub async fn queue() -> ( .expect("Could not subscribe to new events") .into_stream(); - let queue_clone = queue.clone(); - let handle = task::spawn(async move { - let queue = queue_clone; - + // let queue_clone = queue.clone(); + task::spawn(async move { loop { if let Some(msg) = messages.next().await { let msg = match msg { @@ -126,10 +118,8 @@ pub async fn queue() -> ( SubscribeReposMessage::Unknown(data) => SubscribeReposMessage::Unknown(data), }; - queue.lock().await.push_back(ev); + let _ = tx.send(ev); } } - }); - - (queue, handle) + }) } diff --git a/src/main.rs b/src/main.rs index d198b58..2931c97 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,3 +1,5 @@ +use tokio::sync::broadcast; + use crate::backfill::backfill; mod backfill; @@ -35,7 +37,8 @@ async fn main() -> Result<(), Error> { let conn = db::conn().await; println!("Database connected and initialized"); - let (queue, queue_handle) = ingest::queue().await; + let (tx, ingest_reciever) = broadcast::channel(16); + let broadcast_handle = ingest::queue(tx).await; println!("Starting backfill"); let timer = std::time::Instant::now(); @@ -55,10 +58,10 @@ async fn main() -> Result<(), Error> { }); println!("Handling new events. Ctrl+C to quit."); - let ingest_handle = ingest::ingest(queue, conn.clone()); + let ingest_handle = ingest::ingest(ingest_reciever, conn.clone()); let _ = rx.await; - queue_handle.abort(); + broadcast_handle.abort(); ingest_handle.abort(); Ok(()) -- 2.51.2