From 23d526f7e92fe95dfd5afac836b8e68cb749f4f9 Mon Sep 17 00:00:00 2001 From: Mia Date: Sun, 26 Jan 2025 16:19:15 +0000 Subject: [PATCH] start on the relay consumer and indexer --- consumer/src/indexer/mod.rs | 35 +++++++++++++++++++++++++++++++ consumer/src/main.rs | 41 ++++++++++++++++++++++++++++++++++++- 2 files changed, 75 insertions(+), 1 deletion(-) create mode 100644 consumer/src/indexer/mod.rs diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs new file mode 100644 index 00000000..182e3610 --- /dev/null +++ b/consumer/src/indexer/mod.rs @@ -0,0 +1,35 @@ +use crate::firehose::{AtpAccountEvent, AtpCommitEvent, AtpIdentityEvent, FirehoseEvent}; +use tokio::sync::mpsc::Receiver; + +pub async fn relay_indexer(mut rx: Receiver) -> eyre::Result<()> { + while let Some(event) = rx.recv().await { + let res = match event { + FirehoseEvent::Identity(identity) => index_identity(identity).await, + FirehoseEvent::Account(account) => index_account(account).await, + FirehoseEvent::Commit(commit) => index_commit(commit).await, + FirehoseEvent::Label(_) => { + // We handle all labels through direct connections to labelers + tracing::warn!("got #labels from the relay"); + Ok(()) + } + }; + + if let Err(e) = res { + tracing::error!("Indexing error: {e}"); + } + } + + Ok(()) +} + +async fn index_identity(identity: AtpIdentityEvent) -> eyre::Result<()> { + Ok(()) +} + +async fn index_account(account: AtpAccountEvent) -> eyre::Result<()> { + Ok(()) +} + +async fn index_commit(commit: AtpCommitEvent) -> eyre::Result<()> { + Ok(()) +} diff --git a/consumer/src/main.rs b/consumer/src/main.rs index 8eca4928..8cfc7845 100644 --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -1,11 +1,50 @@ +use tokio::sync::mpsc::Sender; + mod config; mod firehose; +mod indexer; #[tokio::main] -async fn main() -> eyre::Result<()> { +async fn main() -> eyre::Result<()> { tracing_subscriber::fmt::init(); let conf = config::load_config()?; + let (tx, rx) = tokio::sync::mpsc::channel::(64); + + let relay_firehose = firehose::FirehoseConsumer::new_relay(&conf.relay_source, None).await?; + + let firehose_handle = tokio::spawn(relay_consumer(relay_firehose, tx)); + let indexer_handle = tokio::spawn(indexer::relay_indexer(rx)); + + let (firehose_res, indexer_res) = tokio::try_join!{ + firehose_handle, + indexer_handle, + }?; + + firehose_res.and(indexer_res) +} + +async fn relay_consumer( + mut consumer: firehose::FirehoseConsumer, + tx: Sender, +) -> eyre::Result<()> { + loop { + let event = consumer.drive().await?; + match event { + firehose::FirehoseOutput::Close => break, + firehose::FirehoseOutput::Continue => continue, + firehose::FirehoseOutput::Error(err) => { + tracing::error!("Firehose sent an error, exiting: {err:?}"); + break; + } + firehose::FirehoseOutput::Event(event) => { + if let Err(e) = tx.send(event).await { + tracing::error!("Error sending event: {e}"); + } + } + } + } + Ok(()) } -- 2.51.2