From 217d1124543f8f328f4665949dcf502bbab3d931 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 18 Jul 2026 18:26:40 +0300 Subject: [PATCH] [stream] exit stream thread on consumer drop, add stream_head --- src/control/indexer.rs | 39 ++++++++++++++++++++++++++++++++---- src/control/stream/engine.rs | 4 ++++ src/db/keyspaces.rs | 14 +++++++++++++ 3 files changed, 53 insertions(+), 4 deletions(-) diff --git a/src/control/indexer.rs b/src/control/indexer.rs index aa7c212..180a80e 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -6,7 +6,10 @@ use super::*; /// implements [`futures::Stream`] and can be used with `StreamExt::next`, /// `while let Some(item) = stream.next().await`, `forward`, etc. /// the stream terminates when the underlying channel closes (i.e. hydrant shuts down). -pub struct EventStream(mpsc::Receiver>); +pub struct EventStream { + receiver: mpsc::Receiver>, + wake: tokio::sync::broadcast::Sender, +} #[cfg(feature = "indexer_stream")] #[derive(Debug, Clone, thiserror::Error)] @@ -15,6 +18,15 @@ pub enum StreamError { ConsumerTooSlow { reason: String }, } +#[cfg(feature = "indexer_stream")] +impl Drop for EventStream { + fn drop(&mut self) { + // Wake the blocking stream thread so it can observe that its receiver + // has gone away instead of retaining the Hydrant database indefinitely. + let _ = self.wake.send(crate::types::BroadcastEvent::Persisted(0)); + } +} + #[cfg(feature = "indexer_stream")] impl StreamError { pub fn code(&self) -> &'static str { @@ -29,7 +41,7 @@ impl Stream for EventStream { type Item = Result; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - self.0.poll_recv(cx) + self.receiver.poll_recv(cx) } } @@ -62,6 +74,14 @@ impl BackfillHandle { #[cfg(feature = "indexer_stream")] impl Hydrant { + /// Return the highest committed durable indexer-stream event ID. + /// + /// Consumers can snapshot this before subscribing and report readiness only + /// after their committed cursor reaches the snapshot. + pub fn stream_head(&self) -> miette::Result> { + self.state.db.stream.event_head() + } + /// subscribe to the ordered event stream. /// /// returns an [`EventStream`] that implements [`futures::Stream`]. @@ -91,7 +111,10 @@ impl Hydrant { }) .expect("failed to spawn stream thread"); - EventStream(rx) + EventStream { + receiver: rx, + wake: self.state.db.stream.event_tx.clone(), + } } #[cfg(feature = "indexer_stream")] @@ -113,6 +136,14 @@ impl Hydrant { let trimmed = TrimmedDid::from(&did).into_static(); let rev = DbTid::new_from_bytes([0u8; 8]); let collection = CowStr::Borrowed("app.bsky.feed.post").into_static(); + let block = bytes::Bytes::from( + serde_ipld_dagcbor::to_vec(&serde_json::json!({ + "$type": "app.bsky.feed.post", + "text": "seeded event", + "createdAt": "2026-01-01T00:00:00Z" + })) + .expect("fixture record must encode as DAG-CBOR"), + ); for i in 0..count { let event_id = db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); @@ -124,7 +155,7 @@ impl Hydrant { collection: collection.clone(), rkey, action: DbAction::Create, - data: StoredData::Nothing, + data: StoredData::Block(block.clone()), }; let bytes = rmp_serde::to_vec(&evt).expect("msgpack serialization cannot fail"); total_bytes += bytes.len(); diff --git a/src/control/stream/engine.rs b/src/control/stream/engine.rs index 46e67e6..c8fcf34 100644 --- a/src/control/stream/engine.rs +++ b/src/control/stream/engine.rs @@ -27,6 +27,10 @@ pub(crate) fn run_ordered_stream( let mut replay_blocked_since: Option = None; loop { + if tx.is_closed() { + return; + } + if let Err(err) = drain_pending_broadcasts(&mut event_rx, &mut pending) { send_stream_error(&tx, err.into()); return; diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index f275b6d..2bac6ae 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -144,6 +144,20 @@ impl StreamDb { self.events.approximate_len() } + pub(crate) fn event_head(&self) -> Result> { + let Some(guard) = self.events.iter().next_back() else { + return Ok(None); + }; + let key = guard.key().into_diagnostic()?; + let id = u64::from_be_bytes( + key.as_ref() + .try_into() + .into_diagnostic() + .wrap_err("expected stream event ID to be 8 bytes")?, + ); + Ok(Some(id)) + } + /// resume event ids after the last stored event. pub(super) fn init(&self) -> Result<()> { let mut last_id = 0; -- 2.51.2