From df5b836758cbd73fee3914b85cd07538ac60a59e Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Sun, 15 Mar 2026 13:50:43 +0300 Subject: [PATCH] [ingest] log relay url --- src/crawler/mod.rs | 6 +-- src/db/keys.rs | 10 ++-- src/db/mod.rs | 10 ++-- src/ingest/firehose.rs | 112 ++++++++++++++++++++--------------------- src/ingest/mod.rs | 4 +- src/ingest/worker.rs | 11 ++-- src/main.rs | 4 +- src/state.rs | 5 +- src/util.rs | 7 --- 9 files changed, 81 insertions(+), 88 deletions(-) diff --git a/src/crawler/mod.rs b/src/crawler/mod.rs index d3308d0..85df0e9 100644 --- a/src/crawler/mod.rs +++ b/src/crawler/mod.rs @@ -3,7 +3,7 @@ use crate::db::keys::crawler_cursor_key; use crate::db::{Db, keys, ser_repo_state}; use crate::state::AppState; use crate::types::RepoState; -use crate::util::{ErrorForStatus, RetryOutcome, RetryWithBackoff, parse_retry_after, relay_id}; +use crate::util::{ErrorForStatus, RetryOutcome, RetryWithBackoff, parse_retry_after}; use chrono::{DateTime, TimeDelta, Utc}; use futures::FutureExt; use jacquard_api::com_atproto::repo::describe_repo::DescribeRepoOutput; @@ -214,7 +214,7 @@ impl Crawler { } async fn get_cursor(&self, relay_host: &Url) -> Result { - let key = crawler_cursor_key(&relay_id(relay_host)); + let key = crawler_cursor_key(relay_host); let cursor_bytes = Db::get(self.state.db.cursors.clone(), &key).await?; let cursor: Cursor = cursor_bytes .as_deref() @@ -603,7 +603,7 @@ impl Crawler { } batch.insert( &db.cursors, - crawler_cursor_key(&relay_id(relay_host)), + crawler_cursor_key(relay_host), rmp_serde::to_vec(&cursor) .into_diagnostic() .wrap_err("cant serialize cursor")?, diff --git a/src/db/keys.rs b/src/db/keys.rs index 9adc642..029e6fc 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -2,7 +2,7 @@ use jacquard_common::types::string::Did; use smol_str::SmolStr; use crate::db::types::{DbRkey, DbTid, TrimmedDid}; -use crate::util::RelayId; +use url::Url; /// separator used for composite keys pub const SEP: u8 = b'|'; @@ -163,14 +163,14 @@ pub fn crawler_retry_parse_key(key: &[u8]) -> miette::Result> { TrimmedDid::try_from(&key[CRAWLER_RETRY_PREFIX.len()..]) } -pub fn crawler_cursor_key(relay_id: &RelayId) -> Vec { +pub fn crawler_cursor_key(relay: &Url) -> Vec { let mut key = b"crawler_cursor|".to_vec(); - key.extend_from_slice(relay_id); + key.extend_from_slice(relay.as_str().as_bytes()); key } -pub fn firehose_cursor_key(relay_id: &RelayId) -> Vec { +pub fn firehose_cursor_key(relay: &Url) -> Vec { let mut key = b"firehose_cursor|".to_vec(); - key.extend_from_slice(relay_id); + key.extend_from_slice(relay.as_str().as_bytes()); key } diff --git a/src/db/mod.rs b/src/db/mod.rs index a4a5359..39b4497 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -13,7 +13,7 @@ use smol_str::SmolStr; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64}; -use crate::util::RelayId; +use url::Url; pub mod compaction; pub mod filter; @@ -517,14 +517,14 @@ impl Db { } } -pub fn set_firehose_cursor(db: &Db, relay_id: &RelayId, cursor: i64) -> Result<()> { +pub fn set_firehose_cursor(db: &Db, relay: &Url, cursor: i64) -> Result<()> { db.cursors - .insert(keys::firehose_cursor_key(relay_id), cursor.to_be_bytes()) + .insert(keys::firehose_cursor_key(relay), cursor.to_be_bytes()) .into_diagnostic() } -pub async fn get_firehose_cursor(db: &Db, relay_id: &RelayId) -> Result> { - let per_relay_key = keys::firehose_cursor_key(relay_id); +pub async fn get_firehose_cursor(db: &Db, relay: &Url) -> Result> { + let per_relay_key = keys::firehose_cursor_key(relay); if let Some(v) = Db::get(db.cursors.clone(), per_relay_key).await? { return Ok(Some(i64::from_be_bytes( v.as_ref() diff --git a/src/ingest/firehose.rs b/src/ingest/firehose.rs index 16598b1..dd5f61e 100644 --- a/src/ingest/firehose.rs +++ b/src/ingest/firehose.rs @@ -1,22 +1,20 @@ -use crate::db; +use crate::db::{self, deser_repo_state}; use crate::filter::{FilterHandle, FilterMode}; use crate::ingest::stream::{FirehoseStream, SubscribeReposMessage, decode_frame}; use crate::ingest::{BufferTx, IngestMessage}; use crate::state::AppState; -use crate::util::RelayId; use jacquard_common::IntoStatic; use jacquard_common::types::did::Did; use miette::{IntoDiagnostic, Result}; use std::sync::Arc; use std::time::Duration; -use tracing::{debug, error, info, trace}; +use tracing::{Span, debug, error, info, trace}; use url::Url; pub struct FirehoseIngestor { state: Arc, buffer_tx: BufferTx, relay_host: Url, - relay_id: RelayId, filter: FilterHandle, _verify_signatures: bool, } @@ -29,57 +27,55 @@ impl FirehoseIngestor { filter: FilterHandle, verify_signatures: bool, ) -> Self { - let relay_id = crate::util::relay_id(&relay_host); Self { state, buffer_tx, relay_host, - relay_id, filter, _verify_signatures: verify_signatures, } } + #[tracing::instrument(skip(self), fields(relay = %self.relay_host))] pub async fn run(self) -> Result<()> { loop { - let start_cursor = db::get_firehose_cursor(&self.state.db, &self.relay_id).await?; + let start_cursor = db::get_firehose_cursor(&self.state.db, &self.relay_host).await?; match start_cursor { - Some(c) => info!(relay = %self.relay_host, cursor = %c, "resuming from cursor"), - None => info!(relay = %self.relay_host, "no cursor found, live tailing"), + Some(c) => info!(cursor = %c, "resuming from cursor"), + None => info!("no cursor found, live tailing"), } - let mut stream = match FirehoseStream::connect(self.relay_host.clone(), start_cursor) - .await - { - Ok(s) => s, - Err(e) => { - error!(relay = %self.relay_host, err = %e, "failed to connect to firehose, retrying in 5s"); - tokio::time::sleep(Duration::from_secs(5)).await; - continue; - } - }; + let mut stream = + match FirehoseStream::connect(self.relay_host.clone(), start_cursor).await { + Ok(s) => s, + Err(e) => { + error!(err = %e, "failed to connect to firehose, retrying in 5s"); + tokio::time::sleep(Duration::from_secs(5)).await; + continue; + } + }; - info!(relay = %self.relay_host, "firehose connected"); + info!("firehose connected"); while let Some(bytes_res) = stream.next().await { let bytes = match bytes_res { Ok(b) => b, Err(e) => { - error!(relay = %self.relay_host, err = %e, "firehose stream error"); + error!(err = %e, "firehose stream error"); break; } }; match decode_frame(&bytes) { Ok(msg) => self.handle_message(msg).await, Err(e) => { - error!(relay = %self.relay_host, err = %e, "firehose stream error"); + error!(err = %e, "firehose stream error"); break; } } } - error!(relay = %self.relay_host, "firehose disconnected, reconnecting in 5s..."); + error!("firehose disconnected, reconnecting in 5s..."); tokio::time::sleep(Duration::from_secs(5)).await; } } @@ -105,7 +101,7 @@ impl FirehoseIngestor { trace!(did = %did, "forwarding message to ingest buffer"); if let Err(e) = self.buffer_tx.send(IngestMessage::Firehose { - relay_id: self.relay_id.clone(), + relay: self.relay_host.clone(), msg: msg.into_static(), }) { error!(err = %e, "failed to send message to buffer processor"); @@ -114,43 +110,47 @@ impl FirehoseIngestor { async fn should_process(&self, did: &Did<'_>) -> Result { let filter = self.filter.load(); + let state = self.state.clone(); + let did = did.clone().into_static(); + let span = Span::current(); - let excl_key = crate::db::filter::exclude_key(did.as_str())?; - if self - .state - .db - .filter - .contains_key(&excl_key) - .into_diagnostic()? - { - return Ok(false); - } + tokio::task::spawn_blocking(move || { + let _entered = span.entered(); + let _entered = tracing::info_span!("should_process", repo = %did).entered(); - match filter.mode { - FilterMode::Full => Ok(true), - FilterMode::Filter => { - let repo_key = crate::db::keys::repo_key(did); - if let Some(state_bytes) = self.state.db.repos.get(&repo_key).into_diagnostic()? { - let repo_state: crate::types::RepoState = - rmp_serde::from_slice(&state_bytes).into_diagnostic()?; - - if repo_state.tracked { - trace!(did = %did, "tracked repo, processing"); - return Ok(true); - } else { - debug!(did = %did, "known but explicitly untracked, skipping"); - return Ok(false); + let excl_key = crate::db::filter::exclude_key(did.as_str())?; + if state.db.filter.contains_key(&excl_key).into_diagnostic()? { + return Ok(false); + } + + match filter.mode { + FilterMode::Full => Ok(true), + FilterMode::Filter => { + let repo_key = crate::db::keys::repo_key(&did); + if let Some(bytes) = state.db.repos.get(&repo_key).into_diagnostic()? { + let repo_state = deser_repo_state(&bytes)?; + + if repo_state.tracked { + trace!(did = %did, "tracked repo, processing"); + return Ok(true); + } else { + debug!(did = %did, "known but explicitly untracked, skipping"); + return Ok(false); + } } - } - if !filter.signals.is_empty() { - trace!(did = %did, "unknown — passing to worker for signal check"); - Ok(true) - } else { - trace!(did = %did, "unknown and no signals configured, skipping"); - Ok(false) + if !filter.signals.is_empty() { + trace!(did = %did, "unknown — passing to worker for signal check"); + Ok(true) + } else { + trace!(did = %did, "unknown and no signals configured, skipping"); + Ok(false) + } } } - } + }) + .await + .into_diagnostic() + .flatten() } } diff --git a/src/ingest/mod.rs b/src/ingest/mod.rs index cd9a8eb..a5abbf4 100644 --- a/src/ingest/mod.rs +++ b/src/ingest/mod.rs @@ -7,12 +7,12 @@ pub mod worker; use jacquard_common::types::did::Did; use crate::ingest::stream::SubscribeReposMessage; -use crate::util::RelayId; +use url::Url; #[derive(Debug)] pub enum IngestMessage { Firehose { - relay_id: RelayId, + relay: Url, msg: SubscribeReposMessage<'static>, }, BackfillFinished(Did<'static>), diff --git a/src/ingest/worker.rs b/src/ingest/worker.rs index 387c06b..b885ce3 100644 --- a/src/ingest/worker.rs +++ b/src/ingest/worker.rs @@ -236,7 +236,8 @@ impl FirehoseWorker { } } } - IngestMessage::Firehose { relay_id, msg } => { + IngestMessage::Firehose { relay, msg } => { + let _span = tracing::info_span!("firehose", relay = %relay).entered(); let (did, seq) = match &msg { SubscribeReposMessage::Commit(c) => (&c.repo, c.seq), SubscribeReposMessage::Identity(i) => (&i.did, i.seq), @@ -252,7 +253,7 @@ impl FirehoseWorker { db::check_poisoned_report(r); } error!(did = %did, err = %e, "error in check_repo_state"); - if let Some((_, cursor)) = state.relay_cursors.get(&relay_id) { + if let Some(cursor) = state.relay_cursors.get(&relay) { cursor.store(seq, std::sync::atomic::Ordering::SeqCst); } continue; @@ -306,8 +307,8 @@ impl FirehoseWorker { did = %did, err = %e, "failed to transition inactive repo to synced" ); - if let Some((_, cursor)) = - state.relay_cursors.get(&relay_id) + if let Some(cursor) = + state.relay_cursors.get(&relay) { cursor.store( seq, @@ -358,7 +359,7 @@ impl FirehoseWorker { } } - if let Some((_, cursor)) = state.relay_cursors.get(&relay_id) { + if let Some(cursor) = state.relay_cursors.get(&relay) { cursor.store(seq, std::sync::atomic::Ordering::SeqCst); } } diff --git a/src/main.rs b/src/main.rs index 43cc5f0..79b8778 100644 --- a/src/main.rs +++ b/src/main.rs @@ -177,10 +177,10 @@ async fn main() -> miette::Result<()> { std::thread::sleep(persist_interval); // persist firehose cursors - for (relay_id, (relay, cursor)) in &state.relay_cursors { + for (relay, cursor) in &state.relay_cursors { let seq = cursor.load(Ordering::SeqCst); if seq > 0 { - if let Err(e) = db::set_firehose_cursor(&state.db, relay_id, seq) { + if let Err(e) = db::set_firehose_cursor(&state.db, relay, seq) { error!(relay = %relay, err = %e, "failed to save cursor"); db::check_poisoned_report(&e); } diff --git a/src/state.rs b/src/state.rs index 969b409..df7e19f 100644 --- a/src/state.rs +++ b/src/state.rs @@ -10,14 +10,13 @@ use crate::{ db::Db, filter::{FilterHandle, new_handle}, resolver::Resolver, - util::{RelayId, relay_id}, }; pub struct AppState { pub db: Db, pub resolver: Resolver, pub filter: FilterHandle, - pub relay_cursors: HashMap, + pub relay_cursors: HashMap, pub backfill_notify: Notify, } @@ -31,7 +30,7 @@ impl AppState { let relay_cursors = config .relays .iter() - .map(|url| (relay_id(url), (url.clone(), AtomicI64::new(0)))) + .map(|url| (url.clone(), AtomicI64::new(0))) .collect(); Ok(Self { diff --git a/src/util.rs b/src/util.rs index f85d453..cb802fa 100644 --- a/src/util.rs +++ b/src/util.rs @@ -3,13 +3,6 @@ use std::time::Duration; use rand::RngExt; use reqwest::StatusCode; use serde::{Deserialize, Deserializer, Serializer}; -use url::Url; - -pub type RelayId = Vec; - -pub fn relay_id(url: &Url) -> RelayId { - url.as_str().as_bytes().to_vec() -} /// outcome of [`RetryWithBackoff::retry`] when the operation does not succeed. pub enum RetryOutcome { -- 2.51.2