diff --git a/docs/api/README.md b/docs/api/README.md index 8ee044c..42aac59 100644 --- a/docs/api/README.md +++ b/docs/api/README.md @@ -6,7 +6,7 @@ hydrant's REST API is split into public endpoints (safe to expose) and managemen ## public -- `GET /stream`: subscribe to the event stream. query params: `cursor` (optional, start from a specific event ID). +- `GET /stream`: subscribe to the event stream. query params: `cursor` (optional, start from a specific event ID). slow consumers may receive a `{"type":"error","error":"ConsumerTooSlow",...}` frame before the connection closes. - `GET /stats`: get stats about the database (counts of repos, records, events; sizes of keyspaces on disk). - `GET /health` / `GET /_health`: health check. diff --git a/docs/configuration.md b/docs/configuration.md index 9072233..ae583e2 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -66,6 +66,10 @@ hydrant is configured via environment variables, all prefixed with `HYDRANT_` (e | variable | default | description | | :--- | :--- | :--- | | `CACHE_SIZE` | `256` | size of the database cache in MB | +| `STREAM_REPLAY_CHUNK_SIZE` | `64` | number of persisted events read per `/stream` replay batch | +| `STREAM_REPLAY_CHUNK_PAUSE` | `2ms` | pause between replay batches so database flush/compaction work can run | +| `STREAM_PENDING_EVENT_LIMIT` | `4096` | maximum live events buffered per stream subscriber while it catches up | +| `STREAM_SEND_TIMEOUT` | `30sec` | maximum time a stream subscriber may block delivery before being disconnected | ## rate limiting (relay mode) diff --git a/examples/statusphere.rs b/examples/statusphere.rs index 5c9bcb5..505e140 100644 --- a/examples/statusphere.rs +++ b/examples/statusphere.rs @@ -113,7 +113,14 @@ async fn handle_stream(index: Arc, repos: ReposControl, mut stream: .map(|h| h.to_string()) .unwrap_or_else(|| did.to_string()) }; - while let Some(event) = stream.next().await { + while let Some(item) = stream.next().await { + let event = match item { + Ok(event) => event, + Err(err) => { + tracing::warn!(err = %err, "hydrant stream closed"); + break; + } + }; if let Some(rec) = event.record { let did = rec.did.as_str().to_owned(); match rec.action.as_str() { diff --git a/src/api/stream.rs b/src/api/stream.rs index d8c54ea..df7db6d 100644 --- a/src/api/stream.rs +++ b/src/api/stream.rs @@ -28,13 +28,53 @@ pub async fn handle_stream( } async fn handle_socket(mut socket: WebSocket, hydrant: Hydrant, query: StreamQuery) { + let send_timeout = hydrant.stream_send_timeout(); let mut stream = hydrant.subscribe(query.cursor); - while let Some(evt) = stream.next().await { + while let Some(item) = stream.next().await { + let evt = match item { + Ok(evt) => evt, + Err(err) => { + let json = serde_json::json!({ + "type": "error", + "error": err.code(), + "message": err.to_string(), + }); + let _ = tokio::time::timeout( + send_timeout, + socket.send(Message::text(json.to_string())), + ) + .await; + let _ = + tokio::time::timeout(std::time::Duration::from_secs(1), socket.close()).await; + break; + } + }; + match serde_json::to_string(&evt) { Ok(json) => { - if socket.send(Message::text(json)).await.is_err() { - break; + match tokio::time::timeout(send_timeout, socket.send(Message::text(json))).await { + Ok(Ok(())) => {} + Ok(Err(_)) => break, + Err(_) => { + let err = serde_json::json!({ + "type": "error", + "error": "ConsumerTooSlow", + "message": format!( + "stream socket send blocked for at least {} seconds", + send_timeout.as_secs() + ), + }); + let _ = tokio::time::timeout( + std::time::Duration::from_secs(1), + socket.send(Message::text(err.to_string())), + ) + .await; + let _ = + tokio::time::timeout(std::time::Duration::from_secs(1), socket.close()) + .await; + break; + } } } Err(e) => { diff --git a/src/config.rs b/src/config.rs index 195c974..ce79bb6 100644 --- a/src/config.rs +++ b/src/config.rs @@ -453,6 +453,19 @@ pub struct Config { /// in-memory write buffer (memtable) size for the records keyspace in MB. /// set via `HYDRANT_DB_RECORDS_MEMTABLE_SIZE_MB`. pub db_records_memtable_size_mb: u64, + + /// maximum number of persisted events read from the database per replay batch. + /// set via `HYDRANT_STREAM_REPLAY_CHUNK_SIZE`. + pub stream_replay_chunk_size: usize, + /// pause between replay batches, giving database maintenance work a chance to run. + /// set via `HYDRANT_STREAM_REPLAY_CHUNK_PAUSE` (humantime duration, e.g. `2ms`). + pub stream_replay_chunk_pause: Duration, + /// maximum number of live in-memory stream events buffered per subscriber while it catches up. + /// set via `HYDRANT_STREAM_PENDING_EVENT_LIMIT`. + pub stream_pending_event_limit: usize, + /// maximum time a subscriber may block stream delivery before being disconnected. + /// set via `HYDRANT_STREAM_SEND_TIMEOUT` (humantime duration, e.g. `30sec`). + pub stream_send_timeout: Duration, } impl Default for Config { @@ -526,6 +539,10 @@ impl Default for Config { db_repos_memtable_size_mb: BASE_MEMTABLE_MB / 2, db_events_memtable_size_mb: BASE_MEMTABLE_MB, db_records_memtable_size_mb: BASE_MEMTABLE_MB / 3 * 2, + stream_replay_chunk_size: 64, + stream_replay_chunk_pause: Duration::from_millis(2), + stream_pending_event_limit: 4096, + stream_send_timeout: Duration::from_secs(30), } } } @@ -650,6 +667,20 @@ impl Config { "DB_REPOS_MEMTABLE_SIZE_MB", defaults.db_repos_memtable_size_mb ); + let stream_replay_chunk_size = cfg!( + "STREAM_REPLAY_CHUNK_SIZE", + defaults.stream_replay_chunk_size + ); + let stream_replay_chunk_pause = cfg!( + "STREAM_REPLAY_CHUNK_PAUSE", + defaults.stream_replay_chunk_pause, + sec + ); + let stream_pending_event_limit = cfg!( + "STREAM_PENDING_EVENT_LIMIT", + defaults.stream_pending_event_limit + ); + let stream_send_timeout = cfg!("STREAM_SEND_TIMEOUT", defaults.stream_send_timeout, sec); let crawler_max_pending_repos = cfg!( "CRAWLER_MAX_PENDING_REPOS", @@ -825,6 +856,10 @@ impl Config { db_repos_memtable_size_mb, db_events_memtable_size_mb, db_records_memtable_size_mb, + stream_replay_chunk_size, + stream_replay_chunk_pause, + stream_pending_event_limit, + stream_send_timeout, }) } } @@ -902,6 +937,18 @@ impl fmt::Display for Config { "db records memtable", format_args!("{} mb", self.db_records_memtable_size_mb) )?; + config_line!(f, "stream replay chunk", self.stream_replay_chunk_size)?; + config_line!( + f, + "stream replay pause", + format_args!("{}ms", self.stream_replay_chunk_pause.as_millis()) + )?; + config_line!(f, "stream pending limit", self.stream_pending_event_limit)?; + config_line!( + f, + "stream send timeout", + format_args!("{}sec", self.stream_send_timeout.as_secs()) + )?; config_line!(f, "crawler max pending", self.crawler_max_pending_repos)?; config_line!( f, diff --git a/src/control/indexer.rs b/src/control/indexer.rs index d6bb725..a48aa66 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -4,13 +4,29 @@ use super::*; /// a stream of [`Event`]s. returned by [`Hydrant::subscribe`]. /// /// implements [`futures::Stream`] and can be used with `StreamExt::next`, -/// `while let Some(evt) = stream.next().await`, `forward`, etc. +/// `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(mpsc::Receiver>); + +#[cfg(feature = "indexer_stream")] +#[derive(Debug, Clone, thiserror::Error)] +pub enum StreamError { + #[error("stream consumer too slow: {reason}")] + ConsumerTooSlow { reason: String }, +} + +#[cfg(feature = "indexer_stream")] +impl StreamError { + pub fn code(&self) -> &'static str { + match self { + Self::ConsumerTooSlow { .. } => "ConsumerTooSlow", + } + } +} #[cfg(feature = "indexer_stream")] impl Stream for EventStream { - type Item = Event; + type Item = Result; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.0.poll_recv(cx) @@ -59,20 +75,26 @@ impl Hydrant { /// a specific repository. /// /// multiple concurrent subscribers each receive a full independent copy of the stream. - /// the stream ends when the `EventStream` is dropped. + /// the stream ends when the `EventStream` is dropped. slow consumers receive + /// [`StreamError::ConsumerTooSlow`] before the stream terminates when possible. pub fn subscribe(&self, cursor: Option) -> EventStream { let (tx, rx) = mpsc::channel(500); let state = self.state.clone(); let runtime = tokio::runtime::Handle::current(); + let opts = stream::EventStreamOptions::from_config(&self.config); std::thread::Builder::new() .name("hydrant-stream".into()) .spawn(move || { let _g = runtime.enter(); - event_stream_thread(state, tx, cursor); + event_stream_thread(state, tx, cursor, opts); }) .expect("failed to spawn stream thread"); EventStream(rx) } + + pub(crate) fn stream_send_timeout(&self) -> std::time::Duration { + self.config.stream_send_timeout + } } diff --git a/src/control/stream.rs b/src/control/stream.rs index 4dc975c..9764ed2 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -1,7 +1,9 @@ +use std::collections::VecDeque; use std::sync::Arc; +use std::time::{Duration, Instant}; -use tokio::sync::mpsc; -use tracing::error; +use tokio::sync::{broadcast, mpsc}; +use tracing::{error, warn}; use crate::db::keys; use crate::state::AppState; @@ -9,7 +11,8 @@ use std::sync::atomic::Ordering; #[cfg(feature = "indexer_stream")] use { - super::Event, + super::{Event, StreamError}, + crate::config::Config, crate::db, crate::types::{BroadcastEvent, MarshallableEvt, RecordEvt, StoredData, StoredEvent}, jacquard_common::types::cid::{ATP_CID_HASH, IpldCid}, @@ -18,13 +21,91 @@ use { jacquard_common::{CowStr, IntoStatic, RawData}, jacquard_repo::DAG_CBOR_CID_CODEC, sha2::{Digest, Sha256}, + tokio::sync::mpsc::error::TrySendError, }; +#[cfg(feature = "indexer_stream")] +const STREAM_SEND_RETRY_PAUSE: Duration = Duration::from_millis(10); + +#[cfg(feature = "indexer_stream")] +#[derive(Debug, Clone, Copy)] +pub(crate) struct EventStreamOptions { + replay_chunk_size: usize, + replay_chunk_pause: Duration, + pending_event_limit: usize, + send_timeout: Duration, +} + +#[cfg(feature = "indexer_stream")] +impl EventStreamOptions { + pub(crate) fn from_config(config: &Config) -> Self { + Self { + replay_chunk_size: config.stream_replay_chunk_size.max(1), + replay_chunk_pause: config.stream_replay_chunk_pause, + pending_event_limit: config.stream_pending_event_limit.max(1), + send_timeout: config.stream_send_timeout, + } + } +} + +#[cfg(feature = "indexer_stream")] +struct PendingEvents { + queue: VecDeque, + persisted_head: Option, + limit: usize, +} + +#[cfg(feature = "indexer_stream")] +impl PendingEvents { + fn new(limit: usize) -> Self { + Self { + queue: VecDeque::new(), + persisted_head: None, + limit, + } + } + + fn push(&mut self, event: BroadcastEvent) -> Result<(), StreamError> { + match event { + BroadcastEvent::Persisted(id) => { + self.persisted_head = Some(self.persisted_head.unwrap_or(0).max(id)); + Ok(()) + } + event => { + if self.queue.len() >= self.limit { + return Err(StreamError::ConsumerTooSlow { + reason: format!( + "pending stream event buffer exceeded {} events", + self.limit + ), + }); + } + self.queue.push_back(event); + Ok(()) + } + } + } + + fn next_event_id(&self) -> Option { + self.queue.front().map(broadcast_event_id) + } + + fn pop_front(&mut self) -> Option { + self.queue.pop_front() + } + + fn take_persisted_after(&mut self, current_id: Option) -> Option { + let head = self.persisted_head.take()?; + is_after_current(head, current_id).then_some(head) + } +} + #[cfg(feature = "indexer_stream")] pub(super) fn event_stream_thread( state: Arc, - tx: mpsc::Sender, + tx: mpsc::Sender>, cursor: Option, + opts: EventStreamOptions, ) { let db = &state.db; let mut event_rx = db.event_tx.subscribe(); @@ -33,85 +114,322 @@ pub(super) fn event_stream_thread( 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(); + let mut catch_up_target = cursor + .and_then(|_| db.next_event_id.load(Ordering::SeqCst).checked_sub(1)) + .filter(|target| is_after_current(*target, current_id)); + let mut replay_gap_target = None; + let mut pending = PendingEvents::new(opts.pending_event_limit); loop { - 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) => { - error!(err = %e, "failed to read event from db"); - break; - } - }; + if let Err(err) = drain_pending_broadcasts(&mut event_rx, &mut pending) { + send_stream_error(&tx, err); + return; + } - let id = match k.as_ref().try_into().map(u64::from_be_bytes) { - Ok(id) => id, - Err(_) => { - error!("failed to parse event id"); - continue; - } - }; - current_id = Some(id); + drop_delivered_pending(&mut pending, current_id); - let stored: StoredEvent = match rmp_serde::from_slice(&v) { - Ok(e) => e, - Err(e) => { - error!(err = %e, "failed to deserialize stored event"); - continue; - } - }; + if let Some(target) = replay_gap_target { + advance_replay_gap(&mut current_id, target, &pending); + if current_id.is_some_and(|id| id >= target) { + replay_gap_target = None; + } + } - let Some(out_evt) = stored_to_event(&state, id, stored, None) else { - continue; - }; + if catch_up_target.is_none() { + catch_up_target = pending.take_persisted_after(current_id); + } + + let pending_next_id = pending.next_event_id(); + let pending_is_ready = pending_next_id.is_some_and(|id| id == next_expected_id(current_id)); - if tx.blocking_send(out_evt).is_err() { - return; // receiver dropped + if let Some(target) = catch_up_target.filter(|_| !pending_is_ready) { + let effective_target = pending_next_id + .and_then(|id| id.checked_sub(1)) + .map(|before_pending| before_pending.min(target)) + .unwrap_or(target); + let chunk = read_replay_chunk( + &state, + &ks, + current_id, + effective_target, + opts.replay_chunk_size, + ); + current_id = chunk.last_seen_id.or(current_id); + + for event in chunk.events { + match send_stream_event(&tx, event, &mut event_rx, &mut pending, opts) { + Ok(SendOutcome::Sent) => {} + Ok(SendOutcome::ReceiverDropped) => return, + Err(err) => { + send_stream_error(&tx, err); + return; + } } } - needs_catch_up = false; - } - // wait for live events - match event_rx.blocking_recv() { - 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; + if chunk.exhausted || current_id.is_some_and(|id| id >= effective_target) { + if effective_target == target { + catch_up_target = None; } + replay_gap_target = Some(effective_target); + } else if !opts.replay_chunk_pause.is_zero() { + std::thread::sleep(opts.replay_chunk_pause); + } - 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() { + continue; + } + + if let Some(event) = pending.pop_front() { + let id = broadcast_event_id(&event); + if !is_after_current(id, current_id) { + continue; + } + + let expected = next_expected_id(current_id); + if id != expected { + catch_up_target = id.checked_sub(1); + pending.queue.push_front(event); + continue; + } + + let Some(out_event) = broadcast_to_event(&state, event) else { + catch_up_target = Some(id); + continue; + }; + + match send_stream_event(&tx, out_event, &mut event_rx, &mut pending, opts) { + Ok(SendOutcome::Sent) => { + current_id = Some(id); + } + Ok(SendOutcome::ReceiverDropped) => return, + Err(err) => { + send_stream_error(&tx, err); return; } - current_id = Some(out_id); } - Ok(BroadcastEvent::Ephemeral(evt)) => { - let evt_id = evt.id; - if tx.blocking_send(*evt).is_err() { + continue; + } + + match event_rx.blocking_recv() { + Ok(BroadcastEvent::Persisted(id)) => { + if is_after_current(id, current_id) { + catch_up_target = Some(id); + } + } + Ok(event) => { + if let Err(err) = pending.push(event) { + send_stream_error(&tx, err); return; } - current_id = Some(current_id.unwrap_or(0).max(evt_id)); } - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => needs_catch_up = true, + Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { + let err = StreamError::ConsumerTooSlow { + reason: format!("subscriber lagged past {skipped} broadcast events"), + }; + warn!(%err, "closing slow stream subscriber"); + send_stream_error(&tx, err); + return; + } Err(tokio::sync::broadcast::error::RecvError::Closed) => break, } } } +#[cfg(feature = "indexer_stream")] +struct ReplayChunk { + events: Vec, + last_seen_id: Option, + exhausted: bool, +} + +#[cfg(feature = "indexer_stream")] +fn read_replay_chunk( + state: &AppState, + ks: &fjall::Keyspace, + current_id: Option, + target: u64, + chunk_size: usize, +) -> ReplayChunk { + let start = current_id.map(|id| id.saturating_add(1)).unwrap_or(0); + if start > target { + return ReplayChunk { + events: Vec::new(), + last_seen_id: current_id, + exhausted: true, + }; + } + + let mut events = Vec::with_capacity(chunk_size); + let mut last_seen_id = current_id; + let mut exhausted = false; + let max_scanned = chunk_size.saturating_mul(4).max(chunk_size); + let mut scanned = 0usize; + let mut iter = ks.range(keys::event_key(start)..=keys::event_key(target)); + + while events.len() < chunk_size && scanned < max_scanned { + let Some(item) = iter.next() else { + exhausted = true; + break; + }; + scanned += 1; + + let (k, v) = match item.into_inner() { + Ok(kv) => kv, + Err(e) => { + error!(err = %e, "failed to read event from db"); + exhausted = true; + break; + } + }; + + let id = match k.as_ref().try_into().map(u64::from_be_bytes) { + Ok(id) => id, + Err(_) => { + error!("failed to parse event id"); + continue; + } + }; + last_seen_id = Some(id); + + let stored: StoredEvent = match rmp_serde::from_slice(&v) { + Ok(e) => e, + Err(e) => { + error!(err = %e, "failed to deserialize stored event"); + continue; + } + }; + + let Some(out_evt) = stored_to_event(state, id, stored, None) else { + continue; + }; + + events.push(out_evt); + } + + ReplayChunk { + events, + last_seen_id, + exhausted, + } +} + +#[cfg(feature = "indexer_stream")] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum SendOutcome { + Sent, + ReceiverDropped, +} + +#[cfg(feature = "indexer_stream")] +fn send_stream_event( + tx: &mpsc::Sender>, + event: Event, + event_rx: &mut broadcast::Receiver, + pending: &mut PendingEvents, + opts: EventStreamOptions, +) -> Result { + let mut item = Ok(event); + let started = Instant::now(); + + loop { + match tx.try_send(item) { + Ok(()) => return Ok(SendOutcome::Sent), + Err(TrySendError::Closed(_)) => return Ok(SendOutcome::ReceiverDropped), + Err(TrySendError::Full(returned)) => { + item = returned; + drain_pending_broadcasts(event_rx, pending)?; + if started.elapsed() >= opts.send_timeout { + return Err(StreamError::ConsumerTooSlow { + reason: format!( + "stream delivery blocked for at least {} seconds", + opts.send_timeout.as_secs() + ), + }); + } + std::thread::sleep(STREAM_SEND_RETRY_PAUSE); + } + } + } +} + +#[cfg(feature = "indexer_stream")] +fn send_stream_error(tx: &mpsc::Sender>, err: StreamError) { + warn!(%err, "closing stream subscriber"); + let _ = tx.blocking_send(Err(err)); +} + +#[cfg(feature = "indexer_stream")] +fn drain_pending_broadcasts( + event_rx: &mut broadcast::Receiver, + pending: &mut PendingEvents, +) -> Result<(), StreamError> { + loop { + match event_rx.try_recv() { + Ok(event) => pending.push(event)?, + Err(broadcast::error::TryRecvError::Empty) => return Ok(()), + Err(broadcast::error::TryRecvError::Closed) => return Ok(()), + Err(broadcast::error::TryRecvError::Lagged(skipped)) => { + return Err(StreamError::ConsumerTooSlow { + reason: format!("subscriber lagged past {skipped} broadcast events"), + }); + } + } + } +} + +#[cfg(feature = "indexer_stream")] +fn drop_delivered_pending(pending: &mut PendingEvents, current_id: Option) { + while pending + .next_event_id() + .is_some_and(|id| !is_after_current(id, current_id)) + { + pending.pop_front(); + } +} + +#[cfg(feature = "indexer_stream")] +fn advance_replay_gap(current_id: &mut Option, target: u64, pending: &PendingEvents) { + let next_pending_id = pending.next_event_id().filter(|id| *id <= target); + let advance_to = next_pending_id + .and_then(|id| id.checked_sub(1)) + .unwrap_or(target); + + if is_after_current(advance_to, *current_id) { + *current_id = Some(advance_to); + } +} + +#[cfg(feature = "indexer_stream")] +fn broadcast_event_id(event: &BroadcastEvent) -> u64 { + match event { + BroadcastEvent::Persisted(id) => *id, + BroadcastEvent::LiveRecord(evt) => evt.id, + BroadcastEvent::Ephemeral(evt) => evt.id, + } +} + +#[cfg(feature = "indexer_stream")] +fn broadcast_to_event(state: &AppState, event: BroadcastEvent) -> Option { + match event { + BroadcastEvent::Persisted(_) => None, + BroadcastEvent::LiveRecord(evt) => { + let stored = evt.stored.clone(); + stored_to_event(state, evt.id, stored, evt.inline_block.clone()) + } + BroadcastEvent::Ephemeral(evt) => Some(*evt), + } +} + +#[cfg(feature = "indexer_stream")] +fn is_after_current(id: u64, current_id: Option) -> bool { + current_id.is_none_or(|current| id > current) +} + +#[cfg(feature = "indexer_stream")] +fn next_expected_id(current_id: Option) -> u64 { + current_id.map(|id| id.saturating_add(1)).unwrap_or(0) +} + #[cfg(feature = "relay")] pub(super) fn relay_stream_thread( state: Arc, diff --git a/src/db/ephemeral.rs b/src/db/ephemeral.rs index 1d0cd67..bca7062 100644 --- a/src/db/ephemeral.rs +++ b/src/db/ephemeral.rs @@ -11,6 +11,10 @@ use { #[cfg(any(feature = "indexer_stream", feature = "relay"))] const AUTO_COMPACT_PRUNED_SEQ_INTERVAL: u64 = 250_000; +#[cfg(any(feature = "indexer_stream", feature = "relay"))] +const TTL_PRUNE_BATCH_SIZE: usize = 10_000; +#[cfg(any(feature = "indexer_stream", feature = "relay"))] +const TTL_PRUNE_BATCH_PAUSE: Duration = Duration::from_millis(10); #[cfg(feature = "indexer_stream")] static LAST_EVENTS_COMPACTED_SEQ: AtomicU64 = AtomicU64::new(0); #[cfg(feature = "relay")] @@ -128,17 +132,46 @@ fn ttl_tick_inner( .map(u64::from_be_bytes) .unwrap_or(0); - let start_key_events = keys::event_key(last_pruned_seq); - let cutoff_key_events = keys::event_key(cutoff_seq); - let mut batch = db.inner.batch(); let mut pruned = 0usize; + let mut next_prune_seq = last_pruned_seq; - for guard in events_ks.range(start_key_events..cutoff_key_events) { - let k = guard.key().into_diagnostic()?; - batch.remove(events_ks, k); - pruned += 1; + loop { + let start_key_events = keys::event_key(next_prune_seq); + let cutoff_key_events = keys::event_key(cutoff_seq); + let mut keys_to_remove = Vec::with_capacity(TTL_PRUNE_BATCH_SIZE); + let mut last_removed_seq = None; + + for guard in events_ks.range(start_key_events..cutoff_key_events) { + let k = guard.key().into_diagnostic()?; + last_removed_seq = Some(read_event_seq(&k)?); + keys_to_remove.push(k); + + if keys_to_remove.len() >= TTL_PRUNE_BATCH_SIZE { + break; + } + } + + let Some(last_removed_seq) = last_removed_seq else { + break; + }; + + let mut batch = db.inner.batch(); + for key in keys_to_remove { + batch.remove(events_ks, key); + pruned += 1; + } + batch.insert( + &db.cursors, + pruned_key.clone(), + last_removed_seq.to_be_bytes(), + ); + batch.commit().into_diagnostic()?; + + next_prune_seq = last_removed_seq.saturating_add(1); + std::thread::sleep(TTL_PRUNE_BATCH_PAUSE); } + let mut batch = db.inner.batch(); batch.insert(&db.cursors, pruned_key, cutoff_seq.to_be_bytes()); // clean up consumed watermark entries (everything up to and including cutoff_ts) @@ -186,3 +219,11 @@ fn ttl_tick_inner( Ok(()) } + +#[cfg(any(feature = "indexer_stream", feature = "relay"))] +fn read_event_seq(key: &[u8]) -> miette::Result { + key.try_into() + .into_diagnostic() + .wrap_err("event key must be 8 bytes") + .map(u64::from_be_bytes) +}