From c93a1abc6cd03cdccdc412974f622424fe3f0253 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Fri, 24 Apr 2026 04:01:48 +0300 Subject: [PATCH] [ingest,stream] dont read from db, use live event --- src/control/stream.rs | 149 +++++++++++++++++++++++++++++------------- src/ingest/indexer.rs | 15 ++++- src/ingest/relay.rs | 2 + src/ops.rs | 106 +++++++++++++++++++++--------- src/types.rs | 21 +++++- tests/run_all.nu | 48 +++++++++++--- 6 files changed, 249 insertions(+), 92 deletions(-) diff --git a/src/control/stream.rs b/src/control/stream.rs index f556f1a..4dc975c 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -30,15 +30,16 @@ pub(super) fn event_stream_thread( let mut event_rx = db.event_tx.subscribe(); let ks = db.events.clone(); let mut current_id = match cursor { - Some(c) => c.saturating_sub(1), - None => db.next_event_id.load(Ordering::SeqCst).saturating_sub(1), + Some(c) => c.checked_sub(1), + None => db.next_event_id.load(Ordering::SeqCst).checked_sub(1), }; + let mut needs_catch_up = cursor.is_some(); loop { - // catch up from db - loop { - let mut found = false; - for item in ks.range(keys::event_key(current_id + 1)..) { + if needs_catch_up { + // catch up from db (record events only; ids are sparse due to ephemeral events) + let start = current_id.map(|id| id.saturating_add(1)).unwrap_or(0); + for item in ks.range(keys::event_key(start)..) { let (k, v) = match item.into_inner() { Ok(kv) => kv, Err(e) => { @@ -54,7 +55,7 @@ pub(super) fn event_stream_thread( continue; } }; - current_id = id; + current_id = Some(id); let stored: StoredEvent = match rmp_serde::from_slice(&v) { Ok(e) => e, @@ -64,29 +65,48 @@ pub(super) fn event_stream_thread( } }; - let Some(evt) = stored_to_event(&state, id, stored) else { + let Some(out_evt) = stored_to_event(&state, id, stored, None) else { continue; }; - if tx.blocking_send(evt).is_err() { + if tx.blocking_send(out_evt).is_err() { return; // receiver dropped } - found = true; - } - if !found { - break; } + needs_catch_up = false; } // wait for live events match event_rx.blocking_recv() { - Ok(BroadcastEvent::Persisted(_)) => {} // re-run catch-up + Ok(BroadcastEvent::Persisted(_)) => needs_catch_up = true, + Ok(BroadcastEvent::LiveRecord(evt)) => { + let expected = current_id.map(|id| id.saturating_add(1)).unwrap_or(0); + if needs_catch_up || evt.id != expected { + needs_catch_up = true; + continue; + } + + let stored = evt.stored.clone(); + let Some(out_evt) = + stored_to_event(&state, evt.id, stored, evt.inline_block.clone()) + else { + needs_catch_up = true; + continue; + }; + let out_id = out_evt.id; + if tx.blocking_send(out_evt).is_err() { + return; + } + current_id = Some(out_id); + } Ok(BroadcastEvent::Ephemeral(evt)) => { + let evt_id = evt.id; if tx.blocking_send(*evt).is_err() { return; } + current_id = Some(current_id.unwrap_or(0).max(evt_id)); } - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {} + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => needs_catch_up = true, Err(tokio::sync::broadcast::error::RecvError::Closed) => break, } } @@ -105,17 +125,24 @@ pub(super) fn relay_stream_thread( let ks = state.db.relay_events.clone(); let mut current_seq = match cursor { Some(c) => c.saturating_sub(1), - None => state - .db - .next_relay_seq - .load(Ordering::Relaxed) - .saturating_sub(1), + None => ks + .iter() + .next_back() + .and_then(|guard| { + guard + .key() + .ok() + .and_then(|k| k.as_ref().try_into().ok()) + .map(u64::from_be_bytes) + }) + .unwrap_or(0), }; + let mut head_seq = current_seq; + let mut needs_catch_up = true; loop { - // catch up from db: send all stored frames from current_seq+1 onward - loop { - let mut found = false; + if needs_catch_up { + // catch up from db: send all stored frames from current_seq+1 onward for item in ks.range(crate::db::keys::relay_event_key(current_seq + 1)..) { let (k, v) = match item.into_inner() { Ok(kv) => kv, @@ -131,33 +158,51 @@ pub(super) fn relay_stream_thread( continue; } }; - current_seq = seq; + if seq != current_seq + 1 { + break; + } if tx.blocking_send(bytes::Bytes::copy_from_slice(&v)).is_err() { return; // subscriber dropped } - found = true; - } - if !found { - break; + current_seq = seq; + if current_seq >= head_seq { + break; + } } + needs_catch_up = false; } // wait for live events match relay_rx.blocking_recv() { - Ok(RelayBroadcast::Persisted(_)) => {} // re-run catch-up - Ok(RelayBroadcast::Ephemeral(frame)) => { + Ok(RelayBroadcast::Persisted(seq)) => { + head_seq = head_seq.max(seq); + needs_catch_up = current_seq < head_seq; + } + Ok(RelayBroadcast::Ephemeral(seq, frame)) => { + head_seq = head_seq.max(seq); + if seq != current_seq + 1 { + // out-of-order or gap: fall back to db catch-up to preserve ordering. + needs_catch_up = true; + continue; + } if tx.blocking_send(frame).is_err() { return; } + current_seq = seq; } - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {} // re-run catch-up + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => needs_catch_up = true, Err(tokio::sync::broadcast::error::RecvError::Closed) => break, } } } #[cfg(feature = "indexer_stream")] -fn stored_to_event(state: &AppState, id: u64, stored: StoredEvent<'_>) -> Option { +fn stored_to_event( + state: &AppState, + id: u64, + stored: StoredEvent<'_>, + inline_block: Option, +) -> Option { let StoredEvent { live, did, @@ -170,26 +215,38 @@ fn stored_to_event(state: &AppState, id: u64, stored: StoredEvent<'_>) -> Option let record = match data { StoredData::Ptr(cid) => { - let block = state - .db - .blocks - .get(&keys::block_key(collection.as_str(), &cid.to_bytes())); - match block { - Ok(Some(bytes)) => match serde_ipld_dagcbor::from_slice::(&bytes) { + if let Some(bytes) = inline_block { + match serde_ipld_dagcbor::from_slice::(&bytes) { Ok(val) => Some((cid, serde_json::to_value(val).ok()?)), Err(e) => { error!(err = %e, "cant parse block"); return None; } - }, - Ok(None) => { - error!("block not found, this is a bug"); - return None; } - Err(e) => { - error!(err = %e, "cant get block"); - db::check_poisoned(&e); - return None; + } else { + let block = state + .db + .blocks + .get(&keys::block_key(collection.as_str(), &cid.to_bytes())); + match block { + Ok(Some(bytes)) => { + match serde_ipld_dagcbor::from_slice::(bytes.as_ref()) { + Ok(val) => Some((cid, serde_json::to_value(val).ok()?)), + Err(e) => { + error!(err = %e, "cant parse block"); + return None; + } + } + } + Ok(None) => { + error!("block not found, this is a bug"); + return None; + } + Err(e) => { + error!(err = %e, "cant get block"); + db::check_poisoned(&e); + return None; + } } } } diff --git a/src/ingest/indexer.rs b/src/ingest/indexer.rs index 1954552..c8b6ed0 100644 --- a/src/ingest/indexer.rs +++ b/src/ingest/indexer.rs @@ -490,8 +490,19 @@ impl FirehoseWorker { *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)); + { + use std::sync::Arc; + + let mut live_events = res.live_events; + for evt in live_events.drain(..) { + ctx.broadcast_events + .push(BroadcastEvent::LiveRecord(Arc::new(evt))); + } + if let Some(last_id) = res.last_event_id { + ctx.broadcast_events + .push(BroadcastEvent::Persisted(last_id)); + } + } Ok(RepoProcessResult::Ok(repo_state)) } diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index 0981d1f..488577c 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -921,6 +921,8 @@ impl WorkerContext<'_> { 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(()) } diff --git a/src/ops.rs b/src/ops.rs index 3db7e26..d3c7f80 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -19,9 +19,10 @@ use crate::types::{RepoState, RepoStatus, ResyncState}; #[cfg(feature = "indexer_stream")] use { crate::types::{ - AccountEvt, BroadcastEvent, IdentityEvt, MarshallableEvt, StoredData, StoredEvent, + AccountEvt, BroadcastEvent, IdentityEvt, LiveRecordEvent, MarshallableEvt, StoredData, + StoredEvent, }, - jacquard_common::CowStr, + jacquard_common::{CowStr, IntoStatic}, std::sync::atomic::Ordering, }; @@ -207,6 +208,10 @@ pub struct ApplyCommitResults<'s> { pub repo_state: RepoState<'s>, pub records_delta: i64, pub blocks_count: i64, + #[cfg(feature = "indexer_stream")] + pub live_events: Vec, + #[cfg(feature = "indexer_stream")] + pub last_event_id: Option, } pub fn apply_commit<'s>( @@ -231,6 +236,14 @@ pub fn apply_commit<'s>( let mut records_delta = 0; let mut blocks_count = 0; let mut collection_deltas: HashMap<&str, i64> = HashMap::new(); + let rev = DbTid::from(&commit.rev); + + #[cfg(feature = "indexer_stream")] + let should_broadcast_live = db.event_tx.receiver_count() > 0; + #[cfg(feature = "indexer_stream")] + let mut live_events = Vec::new(); + #[cfg(feature = "indexer_stream")] + let mut last_event_id = None; for op in &commit.ops { let (collection, rkey) = parse_path(&op.path)?; @@ -242,11 +255,16 @@ pub fn apply_commit<'s>( let rkey = DbRkey::new(rkey); let db_key = keys::record_key(did, collection, &rkey); + let action = DbAction::try_from(op.action.as_str())?; + + #[cfg(feature = "indexer_stream")] + let mut cid_for_event: Option = None; + #[cfg(feature = "indexer_stream")] + let mut block_inline_for_event: Option = None; #[cfg(feature = "indexer_stream")] - let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + let mut inline_block: Option = None; - let action = DbAction::try_from(op.action.as_str())?; - let block: Option = match action { + match action { DbAction::Create | DbAction::Update => { let Some(cid) = &op.cid else { continue; @@ -255,6 +273,10 @@ pub fn apply_commit<'s>( .to_ipld() .into_diagnostic() .wrap_err("expected valid cid from relay")?; + #[cfg(feature = "indexer_stream")] + { + cid_for_event = Some(cid_ipld.clone()); + } let Some(bytes) = parsed.blocks.get(&cid_ipld) else { return Err(miette::miette!( @@ -286,17 +308,16 @@ pub fn apply_commit<'s>( &value, )?; } - None - } else { - // in ephemeral mode, capture bytes inline for event emission #[cfg(feature = "indexer_stream")] - { - Some(bytes.clone()) + if should_broadcast_live && !only_index_links { + // inline record bytes for live tailing so we don't have to load from blocks. + inline_block = Some(bytes.clone()); } - #[cfg(not(feature = "indexer_stream"))] + } else { + #[cfg(feature = "indexer_stream")] { - let _ = bytes; - None + // in ephemeral mode, the event payload is the only place we persist the record. + block_inline_for_event = Some(bytes.clone()); } } } @@ -317,37 +338,54 @@ pub fn apply_commit<'s>( &rkey.to_smolstr(), )?; } - - None } }; #[cfg(feature = "indexer_stream")] { + let data = block_inline_for_event + .clone() + .map(StoredData::Block) + .or_else(|| { + (!only_index_links) + .then(|| cid_for_event.clone().map(StoredData::Ptr)) + .flatten() + }) + .unwrap_or(StoredData::Nothing); + + let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + last_event_id = Some(event_id); + let did_trimmed = TrimmedDid::from(did); + let collection = CowStr::Borrowed(collection); + let evt = StoredEvent { live: true, - did: TrimmedDid::from(did), - rev: DbTid::from(&commit.rev), - collection: CowStr::Borrowed(collection), - rkey, + did: did_trimmed.clone(), + rev, + collection: collection.clone(), + rkey: rkey.clone(), 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), + data: data.clone(), }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; batch.insert(&db.events, keys::event_key(event_id), bytes); + + if should_broadcast_live { + live_events.push(LiveRecordEvent { + id: event_id, + stored: StoredEvent { + live: evt.live, + did: did_trimmed.into_static(), + rev: evt.rev, + collection: collection.into_static(), + rkey, + action: evt.action, + data, + }, + inline_block, + }); + } } - #[cfg(not(feature = "indexer_stream"))] - drop(block); } // update counts @@ -361,6 +399,10 @@ pub fn apply_commit<'s>( repo_state, records_delta, blocks_count, + #[cfg(feature = "indexer_stream")] + live_events, + #[cfg(feature = "indexer_stream")] + last_event_id, }) } diff --git a/src/types.rs b/src/types.rs index 1d72b47..3baae5c 100644 --- a/src/types.rs +++ b/src/types.rs @@ -206,8 +206,11 @@ mod indexer { #[cfg(feature = "indexer_stream")] #[derive(Clone, Debug)] pub(crate) enum BroadcastEvent { - #[allow(dead_code)] - Persisted(u64), + Persisted(#[allow(dead_code)] u64), + /// a durable record event with optional inline block bytes for live tailing. + /// + /// used to avoid re-reading `events`/`blocks` from the database when tailing. + LiveRecord(std::sync::Arc), Ephemeral(Box>), } @@ -423,10 +426,22 @@ pub(crate) struct StoredEvent<'i> { pub data: StoredData, } +/// a durable record event that is also emitted on the in-memory stream after commit. +/// +/// `inline_block` is only used for live tailing to avoid loading the record block from `blocks`. +/// cursor replay continues to read from the database. +#[cfg(feature = "indexer_stream")] +#[derive(Debug, Clone)] +pub(crate) struct LiveRecordEvent { + pub id: u64, + pub stored: StoredEvent<'static>, + pub inline_block: Option, +} + #[cfg(feature = "relay")] #[derive(Clone)] pub(crate) enum RelayBroadcast { Persisted(#[allow(dead_code)] u64), #[allow(dead_code)] - Ephemeral(bytes::Bytes), + Ephemeral(u64, bytes::Bytes), } diff --git a/tests/run_all.nu b/tests/run_all.nu index 352ba50..d4df5b4 100644 --- a/tests/run_all.nu +++ b/tests/run_all.nu @@ -21,22 +21,36 @@ def get_free_ports [count: int] { } def run-test [] { - let result = (with-env { + let name = $in.name + let base_env = { HYDRANT_API_PORT: $in.api HYDRANT_DEBUG_PORT: $in.debug HYDRANT_TEST_MOCK_PORT: $in.mock - HYDRANT_BINARY: "target/x86_64-unknown-linux-gnu/debug/hydrant" - } { - ^nu $"tests/($in.name).nu" | complete + } + let binary = $in.binary? + let env_vars = if $binary == null { $base_env } else { $base_env | insert HYDRANT_BINARY $binary } + + let result = (with-env $env_vars { + ^nu $"tests/($name).nu" | complete }) { - name: $in.name + name: $name success: ($result.exit_code == 0) output: $result.stdout stderr: $result.stderr } } +def test-needs-relay-binary [name: string] { + # tests that build a relay-only binary must run last (and serially) to avoid racing on + # `target/` artifacts while other tests are executing. + try { + open --raw $"tests/($name).nu" | str contains "build-hydrant-relay" + } catch { + false + } +} + def main [--only: list = [], --skip-creds] { print "building hydrant..." # build default features @@ -64,29 +78,45 @@ def main [--only: list = [], --skip-creds] { } let ports = get_free_ports (($tests | length) * 3) + let relay_tests = $tests | where {|t| test-needs-relay-binary $t } + mut assigned = [] for test in ($tests | enumerate) { let p = {($test | get index) * 3 + $in} + let name = ($test | get item) + let binary = if ($relay_tests | any {$in == $name}) { null } else { "target/x86_64-unknown-linux-gnu/debug/hydrant" } let entry = { - name: ($test | get item), + name: $name, api: ($ports | get (0 | do $p)), debug: ($ports | get (1 | do $p)), - mock: ($ports | get (2 | do $p)) + mock: ($ports | get (2 | do $p)), + binary: $binary, } $assigned = ($assigned | append $entry) } + let relay_assigned = $assigned | where {|t| $t.binary == null } + let parallel_assigned = $assigned | where {|t| $t.binary != null } + let groups = { "authenticated_stream": "event_dependent", "count_tracking": "event_dependent", "signal_filter": "event_dependent", } - let grouped = $assigned | group-by {|t| $groups | get -o $t.name | default $t.name} + let grouped = $parallel_assigned | group-by { + let name = $in.name + $groups | get -o $name | default $name + } print $"running ($assigned | length) tests...\n" + if not ($relay_assigned | is-empty) { + print $"note: relay-binary tests will run last and not in parallel: (($relay_assigned | get name) | str join ', ')\n" + } let run_group = {each {timeit -o {run-test} | {time: $in.time, ...$in.output}}}; - let results = $grouped | values | par-each {do $run_group} | flatten + let parallel_results = $grouped | values | par-each {do $run_group} | flatten + let relay_results = if ($relay_assigned | is-empty) { [] } else { $relay_assigned | do $run_group } + let results = $parallel_results | append $relay_results print "\n=== results ===\n" for r in $results { -- 2.51.2