From ffd155a92158aacfc9dc40981ee63a93ff77c04c Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Fri, 17 Apr 2026 00:41:40 +0300 Subject: [PATCH] [all] add indexer_stream feature to decouple the events stream from the indexer --- Cargo.toml | 3 +- README.md | 1 + src/api/debug.rs | 47 ++++++++++++++++------- src/api/mod.rs | 4 +- src/backfill/mod.rs | 86 ++++++++++++++++++++++++------------------ src/control/indexer.rs | 3 ++ src/control/mod.rs | 19 +++++++--- src/control/stream.rs | 6 +-- src/db/ephemeral.rs | 22 ++++++----- src/db/keys/indexer.rs | 2 + src/db/keys/mod.rs | 1 + src/db/mod.rs | 34 +++++++++-------- src/ingest/indexer.rs | 39 +++++++++++++------ src/lib.rs | 12 ++++-- src/ops.rs | 81 +++++++++++++++++++++++---------------- src/types.rs | 11 +++++- 16 files changed, 236 insertions(+), 135 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 971b3e1..ab86c89 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,11 +4,12 @@ version = "0.1.0" edition = "2024" [features] -default = ["indexer"] +default = ["indexer", "indexer_stream"] __persist_sync_all = [] backlinks = [] relay = [] indexer = [] +indexer_stream = ["indexer"] [dependencies] tokio = { version = "1.0", features = ["full"] } diff --git a/README.md b/README.md index 26d185d..1d23e26 100644 --- a/README.md +++ b/README.md @@ -250,6 +250,7 @@ directory, it will also be loaded automatically. | feature | default | description | | :--- | :--- | :--- | | `indexer` | yes | makes hydrant act as an indexer. incompatible with the relay feature. | +| `indexer_stream` | yes | enables the event stream for the indexer. requires indexer feature. | | `relay` | no | makes hydrant act as a relay. incompatible with the indexer feature. | | `backlinks` | no | enables the backlinks indexer and XRPC endpoints (`blue.microcosm.links.*`). requires indexer feature. | diff --git a/src/api/debug.rs b/src/api/debug.rs index f792ef3..6ea20bb 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -1,6 +1,8 @@ use crate::api::AppState; use crate::db::keys; -use crate::types::{RepoState, ResyncState, StoredEvent}; +#[cfg(feature = "indexer_stream")] +use crate::types::StoredEvent; +use crate::types::{RepoState, ResyncState}; use axum::routing::{get, post}; use axum::{ Json, @@ -32,7 +34,10 @@ pub fn router() -> axum::Router> { let r = axum::Router::new() .route("/debug/get", get(handle_debug_get)) .route("/debug/iter", get(handle_debug_iter)) - .route("/debug/compact", post(handle_debug_compact)) + .route("/debug/compact", post(handle_debug_compact)); + + #[cfg(any(feature = "indexer_stream", feature = "relay"))] + let r = r .route( "/debug/ephemeral_ttl_tick", post(handle_debug_ephemeral_ttl_tick), @@ -100,6 +105,7 @@ fn deserialize_value(partition: &str, value: &[u8]) -> Value { return serde_json::to_value(state).unwrap_or(Value::Null); } } + #[cfg(feature = "indexer_stream")] "events" => { if let Ok(event) = rmp_serde::from_slice::(value) { return serde_json::to_value(event).unwrap_or(Value::Null); @@ -279,7 +285,7 @@ fn get_keyspace_by_name(db: &crate::db::Db, name: &str) -> Result Ok(db.pending.clone()), #[cfg(feature = "indexer")] "resync" => Ok(db.resync.clone()), - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] "events" => Ok(db.events.clone()), #[cfg(feature = "indexer")] "records" => Ok(db.records.clone()), @@ -312,15 +318,16 @@ pub async fn handle_debug_compact( Ok(StatusCode::OK) } +#[cfg(any(feature = "indexer_stream", feature = "relay"))] pub async fn handle_debug_ephemeral_ttl_tick( State(state): State>, ) -> Result { - tokio::task::spawn_blocking(move || { - #[cfg(feature = "indexer")] - let res = crate::db::ephemeral::ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl); + tokio::task::spawn_blocking(move || -> miette::Result<()> { + #[cfg(feature = "indexer_stream")] + crate::db::ephemeral::ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; #[cfg(feature = "relay")] - let res = crate::db::ephemeral::relay_events_ttl_tick(&state.db, &state.ephemeral_ttl); - res + crate::db::ephemeral::relay_events_ttl_tick(&state.db, &state.ephemeral_ttl)?; + Ok(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)? @@ -330,6 +337,7 @@ pub async fn handle_debug_ephemeral_ttl_tick( } #[derive(Deserialize)] +#[cfg(any(feature = "indexer_stream", feature = "relay"))] pub struct DebugSeedWatermarkRequest { /// unix timestamp (seconds) to write the watermark at pub ts: u64, @@ -340,20 +348,31 @@ pub struct DebugSeedWatermarkRequest { /// writes an event watermark entry directly to the cursors keyspace, using identical /// key/value encoding to the real TTL worker. used in tests to plant a past watermark /// so the real `ephemeral_ttl_tick` code path is exercised without waiting 3600 seconds. +#[cfg(any(feature = "indexer_stream", feature = "relay"))] pub async fn handle_debug_seed_watermark( State(state): State>, Query(req): Query, ) -> Result { - tokio::task::spawn_blocking(move || { - #[cfg(feature = "indexer")] - let key = crate::db::keys::event_watermark_key(req.ts); + tokio::task::spawn_blocking(move || -> Result<(), StatusCode> { + #[cfg(feature = "indexer_stream")] + state + .db + .cursors + .insert( + crate::db::keys::event_watermark_key(req.ts), + req.event_id.to_be_bytes(), + ) + .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; #[cfg(feature = "relay")] - let key = crate::db::keys::relay_event_watermark_key(req.ts); state .db .cursors - .insert(key, req.event_id.to_be_bytes()) - .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR) + .insert( + crate::db::keys::relay_event_watermark_key(req.ts), + req.event_id.to_be_bytes(), + ) + .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; + Ok(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)??; diff --git a/src/api/mod.rs b/src/api/mod.rs index be67521..027ece6 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -15,7 +15,7 @@ mod ingestion; mod pds; mod repos; mod stats; -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] mod stream; mod xrpc; @@ -39,7 +39,7 @@ pub async fn serve(hydrant: Hydrant, port: u16) -> miette::Result<()> { .route("/health", get(async || "OK")) .route("/_health", get(async || "OK")) .route("/stats", get(stats::get_stats)); - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] let app = app.nest("/stream", stream::router()); let app = app .merge(xrpc::router(blocks_available)) diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 97dec01..af99d9b 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -4,18 +4,15 @@ use crate::filter::FilterMode; use crate::ops; use crate::resolver::ResolverError; use crate::state::AppState; -use crate::types::{ - AccountEvt, BroadcastEvent, Commit, GaugeState, RepoState, RepoStatus, ResyncErrorKind, - ResyncState, StoredData, StoredEvent, -}; +use crate::types::{Commit, GaugeState, RepoState, RepoStatus, ResyncErrorKind, ResyncState}; use fjall::Slice; use jacquard_api::com_atproto::sync::get_repo::{GetRepo, GetRepoError}; +use jacquard_common::IntoStatic; use jacquard_common::error::{ClientError, ClientErrorKind}; use jacquard_common::types::cid::Cid; use jacquard_common::types::did::Did; use jacquard_common::xrpc::{XrpcError, XrpcExt}; -use jacquard_common::{CowStr, IntoStatic}; use jacquard_repo::mst::Mst; use jacquard_repo::{BlockStore, MemoryBlockStore}; use miette::{Diagnostic, IntoDiagnostic, Result}; @@ -23,11 +20,17 @@ use reqwest::StatusCode; use smol_str::{SmolStr, ToSmolStr}; use std::collections::HashMap; use std::sync::Arc; -use std::sync::atomic::Ordering; use std::time::{Duration, Instant}; + use thiserror::Error; use tokio::sync::Semaphore; use tracing::{Instrument, debug, error, info, trace, warn}; +#[cfg(feature = "indexer_stream")] +use { + crate::types::{AccountEvt, BroadcastEvent, StoredData, StoredEvent}, + jacquard_common::CowStr, + std::sync::atomic::Ordering, +}; pub mod manager; @@ -412,6 +415,7 @@ async fn process_did<'i>( ); state.update_from_doc(doc); + #[cfg(feature = "indexer_stream")] let emit_identity = |status: &RepoStatus, active: bool| { let status = match status { RepoStatus::Deactivated => "deactivated", @@ -463,6 +467,7 @@ async fn process_did<'i>( if let Some(status) = inactive_status { warn!(?status, "repo is inactive, stopping backfill"); + #[cfg(feature = "indexer_stream")] emit_identity(&status, false); let resync_state = ResyncState::Gone { @@ -491,6 +496,7 @@ async fn process_did<'i>( }; // emit identity event so any consumers know, but only if something changed + #[cfg(feature = "indexer_stream")] if state.active != previous_state.active || state.status != previous_state.status || previous_state.pds.is_none() @@ -555,6 +561,7 @@ async fn process_did<'i>( let result = { let app_state = app_state.clone(); let did = did.clone(); + #[cfg(feature = "indexer_stream")] let rev = root_commit.rev; tokio::task::spawn_blocking(move || { @@ -663,24 +670,27 @@ async fn process_did<'i>( *collection_counts.entry(path.0.clone()).or_default() += 1; } - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev, - collection: CowStr::Borrowed(collection), - rkey, - action, - data: if ephemeral { - StoredData::Block(val) - } else if only_index_links { - StoredData::Nothing - } else { - StoredData::Ptr(cid_obj.to_ipld().expect("valid cid")) - }, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + #[cfg(feature = "indexer_stream")] + { + let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let evt = StoredEvent { + live: false, + did: TrimmedDid::from(&did), + rev, + collection: CowStr::Borrowed(collection), + rkey, + action, + data: if ephemeral { + StoredData::Block(val) + } else if only_index_links { + StoredData::Nothing + } else { + StoredData::Ptr(cid_obj.to_ipld().expect("valid cid")) + }, + }; + let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; + batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + } count += 1; } @@ -705,18 +715,21 @@ async fn process_did<'i>( &rkey.to_smolstr(), )?; - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev, - collection: CowStr::Borrowed(&collection), - rkey, - action: DbAction::Delete, - data: StoredData::Nothing, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + #[cfg(feature = "indexer_stream")] + { + let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let evt = StoredEvent { + live: false, + did: TrimmedDid::from(&did), + rev, + collection: CowStr::Borrowed(&collection), + rkey, + action: DbAction::Delete, + data: StoredData::Nothing, + }; + let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; + batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + } delta -= 1; count += 1; @@ -812,6 +825,7 @@ async fn process_did<'i>( "committed backfill batch" ); + #[cfg(feature = "indexer_stream")] let _ = db.event_tx.send(BroadcastEvent::Persisted( db.next_event_id.load(Ordering::SeqCst) - 1, )); diff --git a/src/control/indexer.rs b/src/control/indexer.rs index fc083fb..d6bb725 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -1,5 +1,6 @@ use super::*; +#[cfg(feature = "indexer_stream")] /// a stream of [`Event`]s. returned by [`Hydrant::subscribe`]. /// /// implements [`futures::Stream`] and can be used with `StreamExt::next`, @@ -7,6 +8,7 @@ use super::*; /// the stream terminates when the underlying channel closes (i.e. hydrant shuts down). pub struct EventStream(mpsc::Receiver); +#[cfg(feature = "indexer_stream")] impl Stream for EventStream { type Item = Event; @@ -42,6 +44,7 @@ impl BackfillHandle { } } +#[cfg(feature = "indexer_stream")] impl Hydrant { /// subscribe to the ordered event stream. /// diff --git a/src/control/mod.rs b/src/control/mod.rs index 4883d0a..4e78ee4 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -48,10 +48,10 @@ use crate::filter::FilterMode; use crate::ingest::indexer::FirehoseWorker; use crate::pds_meta::{PdsMeta, PdsMetaHandle}; use crate::state::AppState; +#[cfg(feature = "indexer_stream")] use crate::types::MarshallableEvt; - use firehose::FirehoseShared; -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] use stream::event_stream_thread; #[cfg(feature = "relay")] use stream::relay_stream_thread; @@ -77,6 +77,7 @@ pub struct Host { /// - `"account"`: a repo's active/inactive status changed. carries an [`AccountEvt`]. ephemeral, not replayable. /// /// the `id` field is a monotonically increasing sequence number usable as a cursor for [`Hydrant::subscribe`]. +#[cfg(feature = "indexer_stream")] pub type Event = MarshallableEvt<'static>; /// the top-level handle to a hydrant instance. @@ -294,7 +295,7 @@ impl Hydrant { } // 7. ephemeral GC thread (not used in relay mode) - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] if config.ephemeral { let state = state.clone(); std::thread::Builder::new() @@ -344,10 +345,11 @@ impl Hydrant { }); // 9. events/sec stats ticker + #[cfg(any(feature = "indexer_stream", feature = "relay"))] tokio::spawn({ let state = state.clone(); let get_id = |state: &AppState| { - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] let id = state.db.next_event_id.load(Ordering::Relaxed); #[cfg(feature = "relay")] let id = state.db.next_relay_seq.load(Ordering::Relaxed); @@ -681,6 +683,10 @@ impl Hydrant { count_keys.push("resync"); } + #[cfg_attr( + not(any(feature = "indexer_stream", feature = "relay")), + allow(unused_mut) + )] let mut counts: BTreeMap<&'static str, u64> = futures::future::join_all(count_keys.into_iter().map(|name| { let state = state.clone(); @@ -690,7 +696,7 @@ impl Hydrant { .into_iter() .collect(); - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] counts.insert("events", state.db.events.approximate_len() as u64); #[cfg(feature = "relay")] @@ -714,8 +720,9 @@ impl Hydrant { s.insert("pending", state.db.pending.disk_space()); s.insert("resync", state.db.resync.disk_space()); s.insert("resync_buffer", state.db.resync_buffer.disk_space()); - s.insert("events", state.db.events.disk_space()); } + #[cfg(feature = "indexer_stream")] + s.insert("events", state.db.events.disk_space()); #[cfg(feature = "relay")] s.insert("relay_events", state.db.relay_events.disk_space()); diff --git a/src/control/stream.rs b/src/control/stream.rs index 88e4b74..f556f1a 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -7,7 +7,7 @@ use crate::db::keys; use crate::state::AppState; use std::sync::atomic::Ordering; -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] use { super::Event, crate::db, @@ -20,7 +20,7 @@ use { sha2::{Digest, Sha256}, }; -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] pub(super) fn event_stream_thread( state: Arc, tx: mpsc::Sender, @@ -156,7 +156,7 @@ pub(super) fn relay_stream_thread( } } -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] fn stored_to_event(state: &AppState, id: u64, stored: StoredEvent<'_>) -> Option { let StoredEvent { live, diff --git a/src/db/ephemeral.rs b/src/db/ephemeral.rs index bb53abd..3e2a230 100644 --- a/src/db/ephemeral.rs +++ b/src/db/ephemeral.rs @@ -1,12 +1,15 @@ -use crate::db::{Db, keys}; -use fjall::Keyspace; -use miette::{IntoDiagnostic, WrapErr}; -use std::sync::Arc; -use std::sync::atomic::Ordering; -use std::time::Duration; -use tracing::{debug, error, info}; +#[cfg(any(feature = "indexer_stream", feature = "relay"))] +use { + crate::db::{Db, keys}, + fjall::Keyspace, + miette::{IntoDiagnostic, WrapErr}, + std::sync::Arc, + std::sync::atomic::Ordering, + std::time::Duration, + tracing::{debug, error, info}, +}; -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] pub fn ephemeral_ttl_worker(state: Arc) { info!("ephemeral TTL worker started"); loop { @@ -28,7 +31,7 @@ pub fn relay_events_ttl_worker(state: Arc) { } } -#[cfg(feature = "indexer")] +#[cfg(feature = "indexer_stream")] pub fn ephemeral_ttl_tick(db: &Db, ttl: &Duration) -> miette::Result<()> { let current_seq = db.next_event_id.load(Ordering::SeqCst); ttl_tick_inner( @@ -54,6 +57,7 @@ pub fn relay_events_ttl_tick(db: &Db, ttl: &Duration) -> miette::Result<()> { ) } +#[cfg(any(feature = "indexer_stream", feature = "relay"))] fn ttl_tick_inner( db: &Db, ttl: &Duration, diff --git a/src/db/keys/indexer.rs b/src/db/keys/indexer.rs index 532b4cb..506fd42 100644 --- a/src/db/keys/indexer.rs +++ b/src/db/keys/indexer.rs @@ -4,12 +4,14 @@ use smol_str::SmolStr; use super::SEP; use crate::db::types::{DbRkey, DbTid, TrimmedDid}; +#[cfg(feature = "indexer_stream")] pub const EVENT_WATERMARK_PREFIX: &[u8] = b"ewm|"; pub fn pending_key(id: u64) -> [u8; 8] { id.to_be_bytes() } +#[cfg(feature = "indexer_stream")] pub fn event_watermark_key(timestamp_secs: u64) -> Vec { let mut key = Vec::with_capacity(EVENT_WATERMARK_PREFIX.len() + 8); key.extend_from_slice(EVENT_WATERMARK_PREFIX); diff --git a/src/db/keys/mod.rs b/src/db/keys/mod.rs index 2cae544..aaf13dc 100644 --- a/src/db/keys/mod.rs +++ b/src/db/keys/mod.rs @@ -51,6 +51,7 @@ pub fn relay_event_watermark_key(timestamp_secs: u64) -> Vec { } // key format: {SEQ} +#[cfg(any(feature = "indexer_stream", feature = "relay"))] pub fn event_key(seq: u64) -> [u8; 8] { seq.to_be_bytes() } diff --git a/src/db/mod.rs b/src/db/mod.rs index 0016186..3fa7416 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -1,10 +1,11 @@ use crate::config::Compression; use crate::db::compaction::DropPrefixFilterFactory; -#[cfg(feature = "indexer")] +use crate::types::{RepoMetadata, RepoState}; + +#[cfg(feature = "indexer_stream")] use crate::types::BroadcastEvent; #[cfg(feature = "relay")] use crate::types::RelayBroadcast; -use crate::types::{RepoMetadata, RepoState}; use fjall::config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}; use fjall::{ @@ -18,8 +19,6 @@ use smol_str::SmolStr; use std::cell::RefCell; use std::collections::HashSet; use std::sync::Arc; -use std::sync::atomic::AtomicU64; - use url::Url; pub mod compaction; @@ -30,9 +29,11 @@ pub mod migration; pub mod pds_meta; pub mod types; -use tokio::sync::broadcast; use tracing::error; +#[cfg(any(feature = "indexer_stream", feature = "relay"))] +use {std::sync::atomic::AtomicU64, tokio::sync::broadcast}; + fn default_opts() -> KeyspaceCreateOptions { KeyspaceCreateOptions::default() } @@ -56,13 +57,13 @@ pub struct Db { pub resync: Keyspace, #[cfg(feature = "indexer")] pub resync_buffer: Keyspace, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] pub events: Keyspace, #[cfg(feature = "backlinks")] pub backlinks: Keyspace, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] pub(crate) event_tx: broadcast::Sender, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] pub next_event_id: Arc, #[cfg(feature = "relay")] pub(crate) relay_events: Keyspace, @@ -323,7 +324,7 @@ impl Db { .data_block_compression_policy(CompressionPolicy::disabled()) .data_block_restart_interval_policy(RestartIntervalPolicy::all(16)), )?; - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] let events = open_ks( "events", opts() @@ -433,7 +434,7 @@ impl Db { // when adding new keyspaces, make sure to add them to the /stats endpoint // and also update any relevant /debug/* endpoints - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] let (event_tx, _) = broadcast::channel(10000); #[cfg(feature = "relay")] @@ -455,17 +456,17 @@ impl Db { resync, #[cfg(feature = "indexer")] resync_buffer, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] events, counts, filter, crawler, #[cfg(feature = "backlinks")] backlinks, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] event_tx, counts_map: HashMap::new(), - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] next_event_id: Arc::new(AtomicU64::new(0)), #[cfg(feature = "relay")] relay_events, @@ -493,7 +494,7 @@ impl Db { .store(last_relay_seq + 1, std::sync::atomic::Ordering::Relaxed); } - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] { let mut last_id = 0; if let Some(guard) = this.events.iter().next_back() { @@ -529,7 +530,7 @@ impl Db { let ks = match ks_name { #[cfg(feature = "indexer")] "blocks" => &self.blocks, - #[cfg(feature = "indexer")] + #[cfg(feature = "indexer_stream")] "events" => &self.events, "repos" => &self.repos, #[cfg(feature = "backlinks")] @@ -647,8 +648,9 @@ impl Db { tasks.push(compact(self.pending.clone())); tasks.push(compact(self.resync.clone())); tasks.push(compact(self.resync_buffer.clone())); - tasks.push(compact(self.events.clone())); } + #[cfg(feature = "indexer_stream")] + tasks.push(compact(self.events.clone())); #[cfg(feature = "relay")] tasks.push(compact(self.relay_events.clone())); diff --git a/src/ingest/indexer.rs b/src/ingest/indexer.rs index 53844ad..ae7d841 100644 --- a/src/ingest/indexer.rs +++ b/src/ingest/indexer.rs @@ -4,16 +4,14 @@ use crate::ingest::stream::{Account, Commit, Identity}; use crate::ingest::validation; use crate::resolver::{NoSigningKeyError, ResolverError}; use crate::state::AppState; -use crate::types::{ - AccountEvt, BroadcastEvent, GaugeState, IdentityEvt, RepoMetadata, RepoState, RepoStatus, -}; +use crate::types::{GaugeState, RepoMetadata, RepoState, RepoStatus}; use crate::{ops, util}; use fjall::OwnedWriteBatch; use jacquard_common::IntoStatic; -use jacquard_common::cowstr::ToCowStr; use jacquard_common::types::did::Did; + use jacquard_repo::error::CommitError; use miette::{Diagnostic, IntoDiagnostic, Result}; use std::sync::Arc; @@ -22,6 +20,11 @@ use thiserror::Error; use tokio::runtime::Handle as TokioHandle; use tokio::sync::mpsc; use tracing::{debug, error, info, warn}; +#[cfg(feature = "indexer_stream")] +use { + crate::types::{AccountEvt, BroadcastEvent, IdentityEvt}, + jacquard_common::cowstr::ToCowStr, +}; #[derive(Debug)] pub struct IndexerCommitData { @@ -124,6 +127,7 @@ struct WorkerContext<'a> { batch: OwnedWriteBatch, added_blocks: &'a mut i64, records_delta: &'a mut i64, + #[cfg(feature = "indexer_stream")] broadcast_events: &'a mut Vec, } @@ -194,10 +198,12 @@ impl FirehoseWorker { let _guard = handle.enter(); debug!(shard = id, "shard started"); + #[cfg(feature = "indexer_stream")] let mut broadcast_events = Vec::new(); while let Some(msg) = rx.blocking_recv() { let batch = state.db.inner.batch(); + #[cfg(feature = "indexer_stream")] broadcast_events.clear(); let mut added_blocks = 0; @@ -208,6 +214,7 @@ impl FirehoseWorker { batch, added_blocks: &mut added_blocks, records_delta: &mut records_delta, + #[cfg(feature = "indexer_stream")] broadcast_events: &mut broadcast_events, }; @@ -393,6 +400,7 @@ impl FirehoseWorker { if records_delta != 0 { state.db.update_count("records", records_delta); } + #[cfg(feature = "indexer_stream")] for evt in broadcast_events.drain(..) { let _ = state.db.event_tx.send(evt); } @@ -473,26 +481,31 @@ impl FirehoseWorker { let repo_state = res.repo_state; *ctx.added_blocks += res.blocks_count; *ctx.records_delta += res.records_delta; + #[cfg(feature = "indexer_stream")] ctx.broadcast_events .push(BroadcastEvent::Persisted(db.next_event_id.load(SeqCst) - 1)); Ok(RepoProcessResult::Ok(repo_state)) } + #[cfg_attr(not(feature = "indexer_stream"), allow(unused_variables))] fn handle_identity<'s>( ctx: &mut WorkerContext, repo_state: RepoState<'s>, identity: &Identity<'_>, changed: bool, ) -> Result, IngestError> { - let db = &ctx.state.db; - let did = &identity.did; - if changed { - let evt = IdentityEvt { - did: did.clone().into_static(), - handle: repo_state.handle.clone().map(IntoStatic::into_static), - }; - ctx.broadcast_events.push(ops::make_identity_event(db, evt)); + #[cfg(feature = "indexer_stream")] + { + let db = &ctx.state.db; + let did = &identity.did; + if changed { + let evt = IdentityEvt { + did: did.clone().into_static(), + handle: repo_state.handle.clone().map(IntoStatic::into_static), + }; + ctx.broadcast_events.push(ops::make_identity_event(db, evt)); + } } Ok(RepoProcessResult::Ok(repo_state)) @@ -508,6 +521,7 @@ impl FirehoseWorker { let db = &ctx.state.db; let did = &account.did; let is_inactive = !account.active; + #[cfg(feature = "indexer_stream")] let evt = AccountEvt { did: did.clone().into_static(), active: account.active, @@ -537,6 +551,7 @@ impl FirehoseWorker { } } + #[cfg(feature = "indexer_stream")] if changed { ctx.broadcast_events.push(ops::make_account_event(db, evt)); } diff --git a/src/lib.rs b/src/lib.rs index f6b237b..29492b1 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -14,10 +14,16 @@ pub mod deps { pub use smol_str; } -#[cfg(all(feature = "relay", feature = "indexer"))] +#[cfg(all( + feature = "relay", + any(feature = "indexer", feature = "indexer_stream", feature = "backlinks") +))] compile_error!("can't be relay and indexer at the same time"); -#[cfg(all(feature = "relay", feature = "backlinks"))] -compile_error!("can't index backlinks while running as a relay"); +#[cfg(all( + not(feature = "indexer"), + any(feature = "indexer_stream", feature = "backlinks") +))] +compile_error!("indexer dependent features (stream, backlinks) without indexer can't be enabled"); pub(crate) mod api; #[cfg(feature = "indexer")] diff --git a/src/ops.rs b/src/ops.rs index 0d893d2..bbbb689 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -1,13 +1,11 @@ use fjall::OwnedWriteBatch; use fjall::Slice; -use jacquard_common::CowStr; #[cfg(feature = "backlinks")] use jacquard_common::Data; use jacquard_common::types::did::Did; use miette::{Context, IntoDiagnostic, Result}; use std::collections::HashMap; -use std::sync::atomic::Ordering; use tracing::debug; use crate::db::types::{DbAction, DbRkey, DbTid, TrimmedDid}; @@ -16,10 +14,15 @@ use crate::filter::FilterConfig; use crate::ingest::stream::Commit; use crate::ingest::validation::ValidatedCommit; use crate::state::AppState; -use crate::types::StoredData; -use crate::types::{ - AccountEvt, BroadcastEvent, IdentityEvt, MarshallableEvt, RepoState, RepoStatus, ResyncState, - StoredEvent, +use crate::types::{RepoState, RepoStatus, ResyncState}; + +#[cfg(feature = "indexer_stream")] +use { + crate::types::{ + AccountEvt, BroadcastEvent, IdentityEvt, MarshallableEvt, StoredData, StoredEvent, + }, + jacquard_common::CowStr, + std::sync::atomic::Ordering, }; pub fn persist_to_resync_buffer(db: &Db, did: &Did, commit: &Commit) -> Result<()> { @@ -36,6 +39,7 @@ pub fn persist_to_resync_buffer(db: &Db, did: &Did, commit: &Commit) -> Result<( // emitting identity is ephemeral // we dont replay these, consumers can just fetch identity themselves if they need it +#[cfg(feature = "indexer_stream")] pub fn make_identity_event(db: &Db, evt: IdentityEvt<'static>) -> BroadcastEvent { let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); let marshallable = MarshallableEvt { @@ -48,6 +52,7 @@ pub fn make_identity_event(db: &Db, evt: IdentityEvt<'static>) -> BroadcastEvent BroadcastEvent::Ephemeral(Box::new(marshallable)) } +#[cfg(feature = "indexer_stream")] pub fn make_account_event(db: &Db, evt: AccountEvt<'static>) -> BroadcastEvent { let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); let marshallable = MarshallableEvt { @@ -237,10 +242,11 @@ pub fn apply_commit<'s>( let rkey = DbRkey::new(rkey); let db_key = keys::record_key(did, collection, &rkey); + #[cfg(feature = "indexer_stream")] let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); let action = DbAction::try_from(op.action.as_str())?; - let block = match action { + let block: Option = match action { DbAction::Create | DbAction::Update => { let Some(cid) = &op.cid else { continue; @@ -281,10 +287,17 @@ pub fn apply_commit<'s>( )?; } None - } else if action == DbAction::Create || action == DbAction::Update { - Some(bytes.clone()) } else { - unreachable!("we tested if we are in create or update action") + // in ephemeral mode, capture bytes inline for event emission + #[cfg(feature = "indexer_stream")] + { + Some(bytes.clone()) + } + #[cfg(not(feature = "indexer_stream"))] + { + let _ = bytes; + None + } } } DbAction::Delete => { @@ -309,28 +322,32 @@ pub fn apply_commit<'s>( } }; - let evt = StoredEvent { - live: true, - did: TrimmedDid::from(did), - rev: DbTid::from(&commit.rev), - collection: CowStr::Borrowed(collection), - rkey, - action, - data: block - .map(StoredData::Block) - .or_else(|| { - (!only_index_links).then(|| { - op.cid - .as_ref() - .map(|c| c.to_ipld().expect("valid cid")) - .map(StoredData::Ptr) - })? - }) - .unwrap_or(StoredData::Nothing), - }; - - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&db.events, keys::event_key(event_id), bytes); + #[cfg(feature = "indexer_stream")] + { + let evt = StoredEvent { + live: true, + did: TrimmedDid::from(did), + rev: DbTid::from(&commit.rev), + collection: CowStr::Borrowed(collection), + rkey, + action, + data: block + .map(StoredData::Block) + .or_else(|| { + (!only_index_links).then(|| { + op.cid + .as_ref() + .map(|c| c.to_ipld().expect("valid cid")) + .map(StoredData::Ptr) + })? + }) + .unwrap_or(StoredData::Nothing), + }; + let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; + batch.insert(&db.events, keys::event_key(event_id), bytes); + } + #[cfg(not(feature = "indexer_stream"))] + drop(block); } // update counts diff --git a/src/types.rs b/src/types.rs index fe003fb..1d72b47 100644 --- a/src/types.rs +++ b/src/types.rs @@ -11,7 +11,10 @@ use serde::{Deserialize, Serialize, Serializer}; use serde_json::Value; use smol_str::{SmolStr, ToSmolStr}; -use crate::db::types::{DbAction, DbRkey, DbTid, DidKey, TrimmedDid}; +use crate::db::types::{DbTid, DidKey}; + +#[cfg(feature = "indexer_stream")] +use crate::db::types::{DbAction, DbRkey, TrimmedDid}; use crate::resolver::MiniDoc; pub(crate) mod v2 { @@ -200,6 +203,7 @@ mod indexer { } } + #[cfg(feature = "indexer_stream")] #[derive(Clone, Debug)] pub(crate) enum BroadcastEvent { #[allow(dead_code)] @@ -367,6 +371,7 @@ pub struct AccountEvt<'i> { pub status: Option>, } +#[cfg(feature = "indexer_stream")] #[derive(Serialize, Deserialize, Clone)] pub(crate) enum StoredData { Nothing, @@ -375,18 +380,21 @@ pub(crate) enum StoredData { Block(Bytes), } +#[cfg(feature = "indexer_stream")] impl StoredData { pub fn is_nothing(&self) -> bool { matches!(self, StoredData::Nothing) } } +#[cfg(feature = "indexer_stream")] impl Default for StoredData { fn default() -> Self { Self::Nothing } } +#[cfg(feature = "indexer_stream")] impl Debug for StoredData { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { @@ -397,6 +405,7 @@ impl Debug for StoredData { } } +#[cfg(feature = "indexer_stream")] #[derive(Debug, Serialize, Deserialize, Clone)] #[serde(bound(deserialize = "'i: 'de"))] pub(crate) struct StoredEvent<'i> { -- 2.51.2