From 0c6b2445a61d45632246b42ab5d727937acfc812 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 11 Jul 2026 18:27:54 +0300 Subject: [PATCH] [ingest] extract per-mode event sinks from the shared worker RelayWorker/WorkerContext/handlers no longer mix cfg(indexer)/cfg(relay) blocks inline. a per-mode EventSink (indexer: forward IndexerMessages to the hook; relay: re-sequence + emit frames and jetstream events; none: drop) owns the pending queues, batch commit staging, post-commit flush, and cursor advance, selected once in ingest/relay/sink.rs. queue_emit moves into the relay sink; the relay jetstream path now reuses the CAR already parsed during validation instead of re-parsing commit blocks. composition root (run.rs) picks the SinkSeed with one cfg pair. verified live: full-network indexer ingest (1147 repos, 635k events, zero worker errors) with /stream replay + live tail via bun ws, and relay-mode subscribeRepos serving well-formed #commit CBOR frames. --- src/control/hydrant/run.rs | 8 +- src/ingest/relay.rs | 1 + src/ingest/relay/context.rs | 45 +---- src/ingest/relay/handlers.rs | 250 +++----------------------- src/ingest/relay/sink.rs | 22 +++ src/ingest/relay/sink/indexer.rs | 170 ++++++++++++++++++ src/ingest/relay/sink/none.rs | 102 +++++++++++ src/ingest/relay/sink/relay.rs | 293 +++++++++++++++++++++++++++++++ src/ingest/relay/worker.rs | 87 ++------- 9 files changed, 644 insertions(+), 334 deletions(-) create mode 100644 src/ingest/relay/sink.rs create mode 100644 src/ingest/relay/sink/indexer.rs create mode 100644 src/ingest/relay/sink/none.rs create mode 100644 src/ingest/relay/sink/relay.rs diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index 9a89402..254abbe 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -49,11 +49,15 @@ impl Hydrant { let (indexer_tx, firehose_worker) = FirehoseWorker::new(state.clone(), config.firehose_workers); + #[cfg(feature = "indexer")] + let sink_seed = crate::ingest::relay::sink::SinkSeed::new(indexer_tx.clone()); + #[cfg(not(feature = "indexer"))] + let sink_seed = crate::ingest::relay::sink::SinkSeed::default(); + // raw firehose events from pds/relay to RelayWorker. let (buffer_tx, relay_worker) = crate::ingest::relay::RelayWorker::new( state.clone(), - #[cfg(feature = "indexer")] - indexer_tx.clone(), + sink_seed, matches!(config.verify_signatures, SignatureVerification::Full), config.firehose_workers, crate::ingest::validation::ValidationOptions { diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index 08804bb..6a3a39f 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -14,6 +14,7 @@ use crate::types::{RepoState, RepoStatus}; pub mod context; pub mod handlers; +pub(crate) mod sink; pub mod worker; pub(crate) use context::WorkerContext; diff --git a/src/ingest/relay/context.rs b/src/ingest/relay/context.rs index aaad979..f381c36 100644 --- a/src/ingest/relay/context.rs +++ b/src/ingest/relay/context.rs @@ -29,10 +29,6 @@ use crate::ingest::stream::SubscribeReposMessage; use crate::ingest::validation::{ CommitValidationError, SyncValidationError, ValidatedCommit, ValidatedSync, ValidationContext, }; -#[cfg(feature = "relay")] -use crate::types::RelayBroadcast; -#[cfg(feature = "relay")] -use std::sync::atomic::Ordering; use super::{commit_validation_outcome, sync_validation_outcome}; use crate::ingest::firehose_stats::StatsInstant; @@ -44,17 +40,7 @@ pub(crate) struct WorkerContext<'a> { pub(crate) stats: Arc, pub(crate) batch: OwnedWriteBatch, pub(crate) count_deltas: CountDeltas, - #[cfg(feature = "relay")] - pub(crate) pending_broadcasts: Vec, - #[cfg(all(feature = "relay", feature = "jetstream"))] - pub(crate) pending_jetstream_events: Vec<( - crate::types::StoredJetstreamEvent<'static>, - Option, - )>, - #[cfg(feature = "indexer")] - pub(crate) pending_hook_messages: Vec, - #[cfg(feature = "indexer")] - pub(crate) hook: crate::ingest::indexer::IndexerTx, + pub(crate) sink: super::sink::EventSink, pub(crate) http: reqwest::Client, pub(crate) error_counts: HashMap>, pub(crate) wrong_host_authority: HashMap>, @@ -466,13 +452,7 @@ impl WorkerContext<'_> { crate::db::ser_repo_state(&repo_state)?, ); - #[cfg(feature = "indexer")] - { - self.pending_hook_messages - .push(crate::ingest::indexer::IndexerMessage::NewRepo( - did.clone().into_static(), - )); - } + self.sink.new_repo(did.clone().into_static()); self.count_deltas.add_repos(1); @@ -483,25 +463,4 @@ impl WorkerContext<'_> { Ok(Some(repo_state)) } - - #[cfg(feature = "relay")] - pub(crate) fn queue_emit( - &mut self, - make_frame: impl FnOnce(i64) -> Result, - ) -> Result { - let started = StatsInstant::now(); - let result = (|| { - let db = &self.state.db; - let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); - let frame = make_frame(seq as i64)?; - self.batch - .insert(&db.relay_events, keys::relay_event_key(seq), frame.as_ref()); - self.pending_broadcasts - .push(RelayBroadcast::Ephemeral(seq, frame)); - self.pending_broadcasts.push(RelayBroadcast::Persisted(seq)); - Ok(seq) - })(); - self.stats.record_queue_emit(started.elapsed()); - result - } } diff --git a/src/ingest/relay/handlers.rs b/src/ingest/relay/handlers.rs index 7e9a085..c66128d 100644 --- a/src/ingest/relay/handlers.rs +++ b/src/ingest/relay/handlers.rs @@ -1,7 +1,5 @@ use miette::Result; use smol_str::SmolStr; -#[cfg(all(feature = "relay", feature = "jetstream"))] -use smol_str::ToSmolStr; use tokio::runtime::Handle; use tracing::{debug, warn}; use url::Url; @@ -9,18 +7,8 @@ use url::Url; use crate::db::{self, keys}; use crate::types::{RepoState, RepoStatus}; -#[cfg(all(feature = "relay", feature = "jetstream"))] -use crate::db::types::TrimmedDid; -#[cfg(feature = "relay")] -use crate::ingest::stream::encode_frame; use crate::ingest::stream::{Account, AccountStatus, Commit, Identity, Sync}; use crate::ingest::validation::ValidatedCommit; -#[cfg(feature = "indexer")] -use jacquard_common::IntoStatic; -#[cfg(all(feature = "relay", feature = "jetstream"))] -use crate::types::StoredJetstreamEvent; -#[cfg(all(feature = "relay", feature = "jetstream"))] -use jacquard_common::CowStr; use super::{RelayWorker, WorkerContext}; @@ -28,8 +16,8 @@ impl RelayWorker { pub(crate) fn handle_commit( ctx: &mut WorkerContext, repo_state: &mut RepoState, - #[allow(unused_variables)] firehose: &Url, - #[allow(unused_mut)] mut commit: Commit<'static>, + firehose: &Url, + commit: Commit<'static>, ) -> Result<()> { if !repo_state.active { return Ok(()); @@ -47,9 +35,6 @@ impl RelayWorker { .. } = validated; - #[cfg(not(feature = "indexer"))] - let _ = parsed_blocks; - if chain_break.is_broken() { // chain breaks are not grounds for blocking when acting as a relay debug!(broken = ?chain_break, "chain break, forwarding anyway"); @@ -57,72 +42,14 @@ impl RelayWorker { let repo_key = keys::repo_key(&commit.repo); - #[cfg(feature = "indexer")] - { - ctx.pending_hook_messages - .push(crate::ingest::indexer::IndexerMessage::Event(Box::new( - crate::ingest::indexer::IndexerEvent { - seq: commit.seq, - firehose: firehose.clone(), - data: crate::ingest::indexer::IndexerEventData::Commit( - crate::ingest::indexer::IndexerCommitData { - commit, - chain_break: chain_break.is_broken(), - parsed_blocks, - }, - ), - }, - ))); - } - #[cfg(feature = "relay")] - { - #[cfg(feature = "jetstream")] - let jetstream_ops = commit - .ops - .iter() - .enumerate() - .filter_map(|(idx, op)| { - matches!(op.action.as_str(), "create" | "update" | "delete") - .then(|| split_collection(&op.path).map(|col| (idx as u32, col))) - .flatten() - }) - .collect::>(); - #[cfg(feature = "jetstream")] - let jetstream_did = TrimmedDid::from(&commit.repo).into_static(); - - let _relay_seq = ctx.queue_emit(|seq| { - commit.seq = seq; - encode_frame("#commit", &commit) - })?; - #[cfg(feature = "jetstream")] - { - // skip car parse and record materialization when no jetstream subscribers are - // connected — the stream thread re-reads from relay_events as a fallback. - let parsed_car = (ctx.state.db.jetstream_tx.receiver_count() > 0) - .then(|| { - tokio::runtime::Handle::current() - .block_on(jacquard_repo::car::reader::parse_car_bytes( - commit.blocks.as_ref(), - )) - .ok() - }) - .flatten(); - for (op_index, collection) in &jetstream_ops { - let ephemeral = parsed_car.as_ref().and_then(|car| { - build_relay_commit_ephemeral(&commit, *op_index, collection, car) - }); - ctx.pending_jetstream_events.push(( - StoredJetstreamEvent::RelayCommit { - did: jetstream_did.clone(), - collection: collection.clone(), - relay_seq: _relay_seq, - op_index: *op_index, - }, - ephemeral, - )); - } - } - } + ctx.sink.commit( + ctx.state, + &mut ctx.batch, + firehose, + commit, + chain_break.is_broken(), + parsed_blocks, + )?; repo_state.root = Some(commit_obj.into()); repo_state.touch(); @@ -138,8 +65,8 @@ impl RelayWorker { pub(crate) fn handle_sync( ctx: &mut WorkerContext, repo_state: &mut RepoState, - #[allow(unused_variables)] firehose: &Url, - #[allow(unused_mut)] mut sync: Sync<'static>, + firehose: &Url, + sync: Sync<'static>, ) -> Result<()> { if !repo_state.active { return Ok(()); @@ -153,26 +80,7 @@ impl RelayWorker { let repo_key = keys::repo_key(&sync.did); - #[cfg(feature = "indexer")] - { - ctx.pending_hook_messages - .push(crate::ingest::indexer::IndexerMessage::Event(Box::new( - crate::ingest::indexer::IndexerEvent { - seq: sync.seq, - firehose: firehose.clone(), - data: crate::ingest::indexer::IndexerEventData::Sync( - sync.did.into_static(), - ), - }, - ))); - } - #[cfg(feature = "relay")] - { - ctx.queue_emit(|seq| { - sync.seq = seq; - encode_frame("#sync", &sync) - })?; - } + ctx.sink.sync(ctx.state, &mut ctx.batch, firehose, sync)?; repo_state.root = Some(validated.commit_obj.into()); repo_state.touch(); @@ -188,7 +96,7 @@ impl RelayWorker { pub(crate) fn handle_identity( ctx: &mut WorkerContext, repo_state: &mut RepoState, - #[allow(unused_variables)] firehose: &Url, + firehose: &Url, mut identity: Identity<'static>, is_pds: bool, ) -> Result<()> { @@ -200,12 +108,7 @@ impl RelayWorker { repo_state.advance_identity_time(event_ms); let was_active = repo_state.active; let was_pds_host = Self::pds_host(repo_state.pds.as_deref()); - - #[cfg(feature = "indexer")] - let (was_handle, was_signing_key) = ( - repo_state.handle.clone().map(IntoStatic::into_static), - repo_state.signing_key.clone().map(IntoStatic::into_static), - ); + let snapshot = ctx.sink.identity_snapshot(repo_state); // refresh did doc if a pds sent this event // or if there is no handle specified @@ -239,36 +142,8 @@ impl RelayWorker { let repo_key = keys::repo_key(&identity.did); - #[cfg(feature = "indexer")] - { - let changed = - repo_state.handle != was_handle || repo_state.signing_key != was_signing_key; - ctx.pending_hook_messages - .push(crate::ingest::indexer::IndexerMessage::Event(Box::new( - crate::ingest::indexer::IndexerEvent { - seq: identity.seq, - firehose: firehose.clone(), - data: crate::ingest::indexer::IndexerEventData::Identity( - crate::ingest::indexer::IndexerIdentityData { identity, changed }, - ), - }, - ))); - } - #[cfg(feature = "relay")] - { - let _relay_seq = ctx.queue_emit(|seq| { - identity.seq = seq; - encode_frame("#identity", &identity) - })?; - #[cfg(feature = "jetstream")] - ctx.pending_jetstream_events.push(( - StoredJetstreamEvent::RelayIdentity { - did: TrimmedDid::from(&identity.did).into_static(), - relay_seq: _relay_seq, - }, - None, - )); - } + ctx.sink + .identity(ctx.state, &mut ctx.batch, firehose, identity, repo_state, snapshot)?; ctx.batch.insert( &ctx.state.db.repos, @@ -282,8 +157,8 @@ impl RelayWorker { pub(crate) fn handle_account( ctx: &mut WorkerContext, repo_state: &mut RepoState, - #[allow(unused_variables)] firehose: &Url, - #[allow(unused_mut)] mut account: Account<'static>, + firehose: &Url, + account: Account<'static>, _is_pds: bool, ) -> Result<()> { let event_ms = account.time.0.timestamp_millis(); @@ -297,8 +172,7 @@ impl RelayWorker { // always capture was_active for count tracking, not just in indexer mode let was_active = repo_state.active; let was_pds_host = Self::pds_host(repo_state.pds.as_deref()); - #[cfg(feature = "indexer")] - let was_status = repo_state.status.clone(); + let snapshot = ctx.sink.account_snapshot(repo_state); repo_state.active = account.active; if !account.active { @@ -334,39 +208,15 @@ impl RelayWorker { let repo_key = keys::repo_key(&account.did); - #[cfg(feature = "indexer")] - { - let changed = repo_state.active != was_active || repo_state.status != was_status; - ctx.pending_hook_messages - .push(crate::ingest::indexer::IndexerMessage::Event(Box::new( - crate::ingest::indexer::IndexerEvent { - seq: account.seq, - firehose: firehose.clone(), - data: crate::ingest::indexer::IndexerEventData::Account( - crate::ingest::indexer::IndexerAccountData { - account, - was_active, - changed, - }, - ), - }, - ))); - } - #[cfg(feature = "relay")] - { - let _relay_seq = ctx.queue_emit(|seq| { - account.seq = seq; - encode_frame("#account", &account) - })?; - #[cfg(feature = "jetstream")] - ctx.pending_jetstream_events.push(( - StoredJetstreamEvent::RelayAccount { - did: TrimmedDid::from(&account.did).into_static(), - relay_seq: _relay_seq, - }, - None, - )); - } + ctx.sink.account( + ctx.state, + &mut ctx.batch, + firehose, + account, + repo_state, + snapshot, + was_active, + )?; repo_state.touch(); ctx.batch.insert( @@ -412,45 +262,3 @@ impl RelayWorker { } } } - -#[cfg(all(feature = "relay", feature = "jetstream"))] -fn split_collection(path: &str) -> Option> { - path.split_once('/') - .map(|(collection, _)| CowStr::Owned(collection.to_smolstr())) -} - -#[cfg(all(feature = "relay", feature = "jetstream"))] -fn build_relay_commit_ephemeral( - commit: &crate::ingest::stream::Commit, - op_index: u32, - collection: &jacquard_common::CowStr<'static>, - car: &jacquard_repo::car::reader::ParsedCar, -) -> Option { - let op = commit.ops.get(op_index as usize)?; - let (_, rkey) = op.path.split_once('/')?; - let action = op.action.as_str(); - - let (record, cid) = if matches!(action, "create" | "update") { - let cid_link = op.cid.as_ref()?; - let cid_ipld = cid_link.to_ipld().ok()?; - let block = car.blocks.get(&cid_ipld)?; - let val = serde_ipld_dagcbor::from_slice::(block).ok()?; - let record = serde_json::value::to_raw_value(&val) - .ok() - .map(std::sync::Arc::from); - (record, Some(cid_link.to_string())) - } else { - (None, None) - }; - - Some(crate::jetstream::JetstreamEphemeral { - did: commit.repo.as_str().to_string(), - rev: commit.rev.as_str().to_string(), - operation: action.to_string(), - collection: collection.as_str().to_string(), - rkey: rkey.to_string(), - record, - cid, - live: true, - }) -} diff --git a/src/ingest/relay/sink.rs b/src/ingest/relay/sink.rs new file mode 100644 index 0000000..6c198bf --- /dev/null +++ b/src/ingest/relay/sink.rs @@ -0,0 +1,22 @@ +//! per-mode event sinks for the shared firehose worker. +//! +//! the worker validates and persists repo state the same way in every mode; +//! what differs is where accepted events go: the indexer forwards them to the +//! `FirehoseWorker` hook, the relay re-sequences and re-emits them (plus +//! optional jetstream events), and a bare build drops them. each mode is a +//! parallel module with an identical interface, selected once here. + +#[cfg(feature = "indexer")] +mod indexer; +#[cfg(feature = "indexer")] +pub(crate) use indexer::*; + +#[cfg(feature = "relay")] +mod relay; +#[cfg(feature = "relay")] +pub(crate) use relay::*; + +#[cfg(not(any(feature = "indexer", feature = "relay")))] +mod none; +#[cfg(not(any(feature = "indexer", feature = "relay")))] +pub(crate) use none::*; diff --git a/src/ingest/relay/sink/indexer.rs b/src/ingest/relay/sink/indexer.rs new file mode 100644 index 0000000..5a6123f --- /dev/null +++ b/src/ingest/relay/sink/indexer.rs @@ -0,0 +1,170 @@ +use fjall::OwnedWriteBatch; +use jacquard_common::IntoStatic; +use jacquard_common::types::string::Handle; +use miette::{IntoDiagnostic, Result}; +use url::Url; + +use crate::ingest::indexer::{ + IndexerAccountData, IndexerCommitData, IndexerEvent, IndexerEventData, IndexerIdentityData, + IndexerMessage, IndexerTx, +}; +use crate::ingest::stream::{Account, Commit, Identity, Sync}; +use crate::state::AppState; +use crate::db::types::DidKey; +use crate::types::{RepoState, RepoStatus}; + +/// per-worker seed from which each shard builds its sink. +#[derive(Clone)] +pub(crate) struct SinkSeed { + hook: IndexerTx, +} + +impl SinkSeed { + pub(crate) fn new(hook: IndexerTx) -> Self { + Self { hook } + } +} + +/// forwards validated events to the indexer hook after batch commit. +pub(crate) struct EventSink { + hook: IndexerTx, + pending: Vec, +} + +pub(crate) struct Staged; + +pub(crate) struct IdentitySnapshot { + handle: Option>, + signing_key: Option>, +} + +pub(crate) struct AccountSnapshot { + status: RepoStatus, +} + +impl EventSink { + pub(crate) fn new( + seed: &SinkSeed, + _stats: std::sync::Arc, + ) -> Self { + Self { + hook: seed.hook.clone(), + pending: Vec::with_capacity(2), + } + } + + pub(crate) fn new_repo(&mut self, did: jacquard_common::types::did::Did<'static>) { + self.pending.push(IndexerMessage::NewRepo(did)); + } + + pub(crate) fn identity_snapshot(&self, repo_state: &RepoState) -> IdentitySnapshot { + IdentitySnapshot { + handle: repo_state.handle.clone().map(IntoStatic::into_static), + signing_key: repo_state.signing_key.clone().map(IntoStatic::into_static), + } + } + + pub(crate) fn account_snapshot(&self, repo_state: &RepoState) -> AccountSnapshot { + AccountSnapshot { + status: repo_state.status.clone(), + } + } + + pub(crate) fn commit( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + firehose: &Url, + commit: Commit<'static>, + chain_break: bool, + parsed_blocks: jacquard_repo::car::reader::ParsedCar, + ) -> Result<()> { + self.pending.push(IndexerMessage::Event(Box::new(IndexerEvent { + seq: commit.seq, + firehose: firehose.clone(), + data: IndexerEventData::Commit(IndexerCommitData { + commit, + chain_break, + parsed_blocks, + }), + }))); + Ok(()) + } + + pub(crate) fn sync( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + firehose: &Url, + sync: Sync<'static>, + ) -> Result<()> { + self.pending.push(IndexerMessage::Event(Box::new(IndexerEvent { + seq: sync.seq, + firehose: firehose.clone(), + data: IndexerEventData::Sync(sync.did.into_static()), + }))); + Ok(()) + } + + pub(crate) fn identity( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + firehose: &Url, + identity: Identity<'static>, + repo_state: &RepoState, + snapshot: IdentitySnapshot, + ) -> Result<()> { + let changed = repo_state.handle != snapshot.handle + || repo_state.signing_key != snapshot.signing_key; + self.pending.push(IndexerMessage::Event(Box::new(IndexerEvent { + seq: identity.seq, + firehose: firehose.clone(), + data: IndexerEventData::Identity(IndexerIdentityData { identity, changed }), + }))); + Ok(()) + } + + #[allow(clippy::too_many_arguments)] + pub(crate) fn account( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + firehose: &Url, + account: Account<'static>, + repo_state: &RepoState, + snapshot: AccountSnapshot, + was_active: bool, + ) -> Result<()> { + let changed = repo_state.active != was_active || repo_state.status != snapshot.status; + self.pending.push(IndexerMessage::Event(Box::new(IndexerEvent { + seq: account.seq, + firehose: firehose.clone(), + data: IndexerEventData::Account(IndexerAccountData { + account, + was_active, + changed, + }), + }))); + Ok(()) + } + + pub(crate) fn commit_batch( + &mut self, + _state: &AppState, + batch: OwnedWriteBatch, + ) -> Result { + batch.commit().into_diagnostic()?; + Ok(Staged) + } + + pub(crate) fn flush(&mut self, _state: &AppState, _staged: Staged) { + for msg in self.pending.drain(..) { + let _ = self.hook.blocking_send(msg); + } + } + + pub(crate) fn advance_cursor(&self, _state: &AppState, _firehose: &Url, _seq: i64) { + // in indexer mode the FirehoseWorker advances the cursor after processing + } +} diff --git a/src/ingest/relay/sink/none.rs b/src/ingest/relay/sink/none.rs new file mode 100644 index 0000000..73507a9 --- /dev/null +++ b/src/ingest/relay/sink/none.rs @@ -0,0 +1,102 @@ +//! sink for builds with neither indexer nor relay: validated events update +//! repo state but are not forwarded anywhere. + +use fjall::OwnedWriteBatch; +use miette::{IntoDiagnostic, Result}; +use url::Url; + +use crate::ingest::stream::{Account, Commit, Identity, Sync}; +use crate::state::AppState; +use crate::types::RepoState; + +/// per-worker seed from which each shard builds its sink. +#[derive(Clone, Default)] +pub(crate) struct SinkSeed; + +pub(crate) struct EventSink; + +pub(crate) struct Staged; + +pub(crate) struct IdentitySnapshot; + +pub(crate) struct AccountSnapshot; + +impl EventSink { + pub(crate) fn new( + _seed: &SinkSeed, + _stats: std::sync::Arc, + ) -> Self { + Self + } + + pub(crate) fn new_repo(&mut self, _did: jacquard_common::types::did::Did<'static>) {} + + pub(crate) fn identity_snapshot(&self, _repo_state: &RepoState) -> IdentitySnapshot { + IdentitySnapshot + } + + pub(crate) fn account_snapshot(&self, _repo_state: &RepoState) -> AccountSnapshot { + AccountSnapshot + } + + pub(crate) fn commit( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + _firehose: &Url, + _commit: Commit<'static>, + _chain_break: bool, + _parsed_blocks: jacquard_repo::car::reader::ParsedCar, + ) -> Result<()> { + Ok(()) + } + + pub(crate) fn sync( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + _firehose: &Url, + _sync: Sync<'static>, + ) -> Result<()> { + Ok(()) + } + + pub(crate) fn identity( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + _firehose: &Url, + _identity: Identity<'static>, + _repo_state: &RepoState, + _snapshot: IdentitySnapshot, + ) -> Result<()> { + Ok(()) + } + + #[allow(clippy::too_many_arguments)] + pub(crate) fn account( + &mut self, + _state: &AppState, + _batch: &mut OwnedWriteBatch, + _firehose: &Url, + _account: Account<'static>, + _repo_state: &RepoState, + _snapshot: AccountSnapshot, + _was_active: bool, + ) -> Result<()> { + Ok(()) + } + + pub(crate) fn commit_batch( + &mut self, + _state: &AppState, + batch: OwnedWriteBatch, + ) -> Result { + batch.commit().into_diagnostic()?; + Ok(Staged) + } + + pub(crate) fn flush(&mut self, _state: &AppState, _staged: Staged) {} + + pub(crate) fn advance_cursor(&self, _state: &AppState, _firehose: &Url, _seq: i64) {} +} diff --git a/src/ingest/relay/sink/relay.rs b/src/ingest/relay/sink/relay.rs new file mode 100644 index 0000000..b4f91b0 --- /dev/null +++ b/src/ingest/relay/sink/relay.rs @@ -0,0 +1,293 @@ +use fjall::OwnedWriteBatch; +use miette::{IntoDiagnostic, Result}; +#[cfg(feature = "jetstream")] +use smol_str::ToSmolStr; +use std::sync::atomic::Ordering; +use url::Url; + +#[cfg(feature = "jetstream")] +use crate::db::types::TrimmedDid; +use crate::db::keys; +use crate::ingest::stream::{Account, Commit, Identity, Sync, encode_frame}; +use crate::state::AppState; +#[cfg(feature = "jetstream")] +use crate::types::StoredJetstreamEvent; +use crate::types::{RelayBroadcast, RepoState}; +#[cfg(feature = "jetstream")] +use jacquard_common::CowStr; + +/// per-worker seed from which each shard builds its sink. +#[derive(Clone, Default)] +pub(crate) struct SinkSeed; + +/// re-sequences validated events and re-emits them on the relay stream, +/// plus jetstream events when that feature is enabled. +pub(crate) struct EventSink { + stats: std::sync::Arc, + broadcasts: Vec, + #[cfg(feature = "jetstream")] + jetstream_events: Vec<( + StoredJetstreamEvent<'static>, + Option, + )>, +} + +pub(crate) struct Staged { + #[cfg(feature = "jetstream")] + jetstream_broadcasts: Vec, +} + +pub(crate) struct IdentitySnapshot; + +pub(crate) struct AccountSnapshot; + +impl EventSink { + pub(crate) fn new( + _seed: &SinkSeed, + stats: std::sync::Arc, + ) -> Self { + Self { + stats, + broadcasts: Vec::with_capacity(2), + #[cfg(feature = "jetstream")] + jetstream_events: Vec::with_capacity(2), + } + } + + pub(crate) fn new_repo(&mut self, _did: jacquard_common::types::did::Did<'static>) {} + + pub(crate) fn identity_snapshot(&self, _repo_state: &RepoState) -> IdentitySnapshot { + IdentitySnapshot + } + + pub(crate) fn account_snapshot(&self, _repo_state: &RepoState) -> AccountSnapshot { + AccountSnapshot + } + + /// assigns the next relay sequence number, persists the re-encoded frame, + /// and queues its broadcast. + fn queue_emit( + &mut self, + state: &AppState, + batch: &mut OwnedWriteBatch, + make_frame: impl FnOnce(i64) -> Result, + ) -> Result { + let started = crate::ingest::firehose_stats::StatsInstant::now(); + let result = (|| { + let db = &state.db; + let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); + let frame = make_frame(seq as i64)?; + batch.insert(&db.relay_events, keys::relay_event_key(seq), frame.as_ref()); + self.broadcasts.push(RelayBroadcast::Ephemeral(seq, frame)); + self.broadcasts.push(RelayBroadcast::Persisted(seq)); + Ok(seq) + })(); + self.stats.record_queue_emit(started.elapsed()); + result + } + + pub(crate) fn commit( + &mut self, + state: &AppState, + batch: &mut OwnedWriteBatch, + _firehose: &Url, + mut commit: Commit<'static>, + _chain_break: bool, + parsed_blocks: jacquard_repo::car::reader::ParsedCar, + ) -> Result<()> { + let _ = &parsed_blocks; + #[cfg(feature = "jetstream")] + let jetstream_ops = commit + .ops + .iter() + .enumerate() + .filter_map(|(idx, op)| { + matches!(op.action.as_str(), "create" | "update" | "delete") + .then(|| split_collection(&op.path).map(|col| (idx as u32, col))) + .flatten() + }) + .collect::>(); + #[cfg(feature = "jetstream")] + let jetstream_did = TrimmedDid::from(&commit.repo).into_static(); + + let _relay_seq = self.queue_emit(state, batch, |seq| { + commit.seq = seq; + encode_frame("#commit", &commit) + })?; + #[cfg(feature = "jetstream")] + { + // skip car parse and record materialization when no jetstream subscribers are + // connected — the stream thread re-reads from relay_events as a fallback. + let has_subscribers = state.db.jetstream_tx.receiver_count() > 0; + for (op_index, collection) in &jetstream_ops { + let ephemeral = has_subscribers + .then(|| build_relay_commit_ephemeral(&commit, *op_index, collection, &parsed_blocks)) + .flatten(); + self.jetstream_events.push(( + StoredJetstreamEvent::RelayCommit { + did: jetstream_did.clone(), + collection: collection.clone(), + relay_seq: _relay_seq, + op_index: *op_index, + }, + ephemeral, + )); + } + } + Ok(()) + } + + pub(crate) fn sync( + &mut self, + state: &AppState, + batch: &mut OwnedWriteBatch, + _firehose: &Url, + mut sync: Sync<'static>, + ) -> Result<()> { + self.queue_emit(state, batch, |seq| { + sync.seq = seq; + encode_frame("#sync", &sync) + })?; + Ok(()) + } + + pub(crate) fn identity( + &mut self, + state: &AppState, + batch: &mut OwnedWriteBatch, + _firehose: &Url, + mut identity: Identity<'static>, + _repo_state: &RepoState, + _snapshot: IdentitySnapshot, + ) -> Result<()> { + let _relay_seq = self.queue_emit(state, batch, |seq| { + identity.seq = seq; + encode_frame("#identity", &identity) + })?; + #[cfg(feature = "jetstream")] + self.jetstream_events.push(( + StoredJetstreamEvent::RelayIdentity { + did: TrimmedDid::from(&identity.did).into_static(), + relay_seq: _relay_seq, + }, + None, + )); + Ok(()) + } + + #[allow(clippy::too_many_arguments)] + pub(crate) fn account( + &mut self, + state: &AppState, + batch: &mut OwnedWriteBatch, + _firehose: &Url, + mut account: Account<'static>, + _repo_state: &RepoState, + _snapshot: AccountSnapshot, + _was_active: bool, + ) -> Result<()> { + let _relay_seq = self.queue_emit(state, batch, |seq| { + account.seq = seq; + encode_frame("#account", &account) + })?; + #[cfg(feature = "jetstream")] + self.jetstream_events.push(( + StoredJetstreamEvent::RelayAccount { + did: TrimmedDid::from(&account.did).into_static(), + relay_seq: _relay_seq, + }, + None, + )); + Ok(()) + } + + pub(crate) fn commit_batch( + &mut self, + state: &AppState, + batch: OwnedWriteBatch, + ) -> Result { + #[cfg(feature = "jetstream")] + { + let mut batch = batch; + let mut jetstream_broadcasts = Vec::new(); + let _lock = state.db.jetstream_lock.lock(); + for (event, ephemeral) in self.jetstream_events.drain(..) { + jetstream_broadcasts.push(crate::jetstream::stage_event( + &mut batch, &state.db, event, ephemeral, + )?); + } + batch.commit().into_diagnostic()?; + Ok(Staged { + jetstream_broadcasts, + }) + } + #[cfg(not(feature = "jetstream"))] + { + let _ = state; + batch.commit().into_diagnostic()?; + Ok(Staged {}) + } + } + + pub(crate) fn flush(&mut self, state: &AppState, staged: Staged) { + for broadcast in self.broadcasts.drain(..) { + let _ = state.db.relay_broadcast_tx.send(broadcast); + } + #[cfg(feature = "jetstream")] + for broadcast in staged.jetstream_broadcasts { + let _ = state.db.jetstream_tx.send(broadcast); + } + #[cfg(not(feature = "jetstream"))] + let _ = staged; + } + + /// the relay is the terminal consumer, so it advances the firehose cursor + /// itself once a message is fully processed. + pub(crate) fn advance_cursor(&self, state: &AppState, firehose: &Url, seq: i64) { + state + .firehose_cursors + .peek_with(firehose, |_, c| c.store(seq, Ordering::SeqCst)); + } +} + +#[cfg(feature = "jetstream")] +fn split_collection(path: &str) -> Option> { + path.split_once('/') + .map(|(collection, _)| CowStr::Owned(collection.to_smolstr())) +} + +#[cfg(feature = "jetstream")] +fn build_relay_commit_ephemeral( + commit: &Commit, + op_index: u32, + collection: &CowStr<'static>, + car: &jacquard_repo::car::reader::ParsedCar, +) -> Option { + let op = commit.ops.get(op_index as usize)?; + let (_, rkey) = op.path.split_once('/')?; + let action = op.action.as_str(); + + let (record, cid) = if matches!(action, "create" | "update") { + let cid_link = op.cid.as_ref()?; + let cid_ipld = cid_link.to_ipld().ok()?; + let block = car.blocks.get(&cid_ipld)?; + let val = serde_ipld_dagcbor::from_slice::(block).ok()?; + let record = serde_json::value::to_raw_value(&val) + .ok() + .map(std::sync::Arc::from); + (record, Some(cid_link.to_string())) + } else { + (None, None) + }; + + Some(crate::jetstream::JetstreamEphemeral { + did: commit.repo.as_str().to_string(), + rev: commit.rev.as_str().to_string(), + operation: action.to_string(), + collection: collection.as_str().to_string(), + rkey: rkey.to_string(), + record, + cid, + live: true, + }) +} diff --git a/src/ingest/relay/worker.rs b/src/ingest/relay/worker.rs index 6d60602..e2c93e0 100644 --- a/src/ingest/relay/worker.rs +++ b/src/ingest/relay/worker.rs @@ -1,6 +1,5 @@ use std::sync::Arc; -#[cfg(feature = "relay")] -use std::sync::atomic::Ordering; + use crate::ingest::firehose_stats::StatsInstant; use miette::{IntoDiagnostic, Result}; @@ -14,9 +13,8 @@ use crate::ingest::validation::ValidationOptions; use crate::ingest::{BufferRx, BufferTx, IngestMessage}; use crate::state::AppState; -use super::{AuthorityOutcome, WorkerContext}; - -use super::relay_message_kind; +use super::sink::{EventSink, SinkSeed}; +use super::{AuthorityOutcome, WorkerContext, relay_message_kind}; pub struct WorkerMessage { pub(crate) is_pds: bool, @@ -27,8 +25,7 @@ pub struct WorkerMessage { pub struct RelayWorker { pub(crate) state: Arc, pub(crate) rxs: Vec, - #[cfg(feature = "indexer")] - pub(crate) hook: crate::ingest::indexer::IndexerTx, + pub(crate) seed: SinkSeed, pub(crate) verify_signatures: bool, pub(crate) num_shards: usize, pub(crate) validation_opts: Arc, @@ -38,7 +35,7 @@ pub struct RelayWorker { impl RelayWorker { pub fn new( state: Arc, - #[cfg(feature = "indexer")] hook: crate::ingest::indexer::IndexerTx, + seed: SinkSeed, verify_signatures: bool, num_shards: usize, validation_opts: ValidationOptions, @@ -49,8 +46,7 @@ impl RelayWorker { Self { state, rxs, - #[cfg(feature = "indexer")] - hook, + seed, verify_signatures, num_shards, validation_opts: Arc::new(validation_opts), @@ -64,8 +60,7 @@ impl RelayWorker { for (i, rx) in self.rxs.into_iter().enumerate() { let state = Arc::clone(&self.state); - #[cfg(feature = "indexer")] - let hook = self.hook.clone(); + let seed = self.seed.clone(); let verify = self.verify_signatures; let h = handle.clone(); let opts = self.validation_opts.clone(); @@ -80,8 +75,7 @@ impl RelayWorker { i, rx, state, - #[cfg(feature = "indexer")] - hook, + seed, verify, h, opts, @@ -108,7 +102,7 @@ impl RelayWorker { id: usize, mut rx: BufferRx, state: Arc, - #[cfg(feature = "indexer")] hook: crate::ingest::indexer::IndexerTx, + seed: SinkSeed, verify_signatures: bool, handle: Handle, validation_opts: Arc, @@ -129,14 +123,7 @@ impl RelayWorker { stats: shard_stats.clone(), batch: state.db.inner.batch(), count_deltas: CountDeltas::default(), - #[cfg(feature = "relay")] - pending_broadcasts: Vec::with_capacity(2), - #[cfg(all(feature = "relay", feature = "jetstream"))] - pending_jetstream_events: Vec::with_capacity(2), - #[cfg(feature = "indexer")] - pending_hook_messages: Vec::with_capacity(2), - #[cfg(feature = "indexer")] - hook, + sink: EventSink::new(&seed, shard_stats.clone()), http, error_counts: Default::default(), wrong_host_authority: Default::default(), @@ -193,66 +180,30 @@ impl RelayWorker { .stage_count_deltas(&mut batch, &ctx.count_deltas); let stage_counts = stage_counts_started.elapsed(); - #[cfg(all(feature = "relay", feature = "jetstream"))] - let mut jetstream_broadcasts = Vec::new(); - let stage_and_commit_started = StatsInstant::now(); - #[cfg(all(feature = "relay", feature = "jetstream"))] - let res = { - let _lock = ctx.state.db.jetstream_lock.lock(); - let mut stage_res = Ok(()); - for (event, ephemeral) in ctx.pending_jetstream_events.drain(..) { - match crate::jetstream::stage_event(&mut batch, &ctx.state.db, event, ephemeral) - { - Ok(broadcast) => jetstream_broadcasts.push(broadcast), - Err(e) => { - stage_res = Err(e); - break; - } - } + let staged = match ctx.sink.commit_batch(&state, batch) { + Ok(staged) => staged, + Err(e) => { + shard_stats.record_commit_error(); + error!(shard = id, err = %e, "relay shard: failed to commit batch"); + drop(reservation); + continue; } - stage_res.and_then(|_| batch.commit().into_diagnostic()) }; - - #[cfg(not(all(feature = "relay", feature = "jetstream")))] - let res = batch.commit(); let stage_and_commit = stage_and_commit_started.elapsed(); - - if let Err(e) = res { - shard_stats.record_commit_error(); - error!(shard = id, err = %e, "relay shard: failed to commit batch"); - drop(reservation); - continue; - } let apply_counts_started = StatsInstant::now(); ctx.state.db.apply_count_deltas(&ctx.count_deltas); drop(reservation); let apply_counts = apply_counts_started.elapsed(); let broadcast_started = StatsInstant::now(); - #[cfg(feature = "relay")] - for broadcast in ctx.pending_broadcasts.drain(..) { - let _ = state.db.relay_broadcast_tx.send(broadcast); - } - #[cfg(all(feature = "relay", feature = "jetstream"))] - for broadcast in jetstream_broadcasts { - let _ = state.db.jetstream_tx.send(broadcast); - } - #[cfg(feature = "indexer")] - for msg in ctx.pending_hook_messages.drain(..) { - let _ = ctx.hook.blocking_send(msg); - } + ctx.sink.flush(&state, staged); let broadcast = broadcast_started.elapsed(); // advance cursor for this firehose only if we are the terminal consumer (relay mode) // in events mode, FirehoseWorker will advance the cursor after processing let cursor_started = StatsInstant::now(); - #[cfg(feature = "relay")] - { - ctx.state - .firehose_cursors - .peek_with(&firehose, |_, c| c.store(seq, Ordering::SeqCst)); - } + ctx.sink.advance_cursor(&state, &firehose, seq); shard_stats.record_processed(crate::ingest::firehose_stats::RelayShardTimings { process_message, stage_counts, -- 2.51.2