From 22e350224ea78bc3df07b8fbfde956a99b25271e Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Tue, 29 Sep 2026 16:44:50 +0300 Subject: [PATCH] [stream] share one ordered engine across /stream, /subscribe and subscribeRepos with positions assigned at commit, broadcasts arrive in position order and every unbroadcast row is announced by a later marker, so a subscriber only has to take broadcasts in order and read the db for a marker's rows. all three endpoints now run the same engine over a position type (event id, jetstream key, relay seq); each only says how to read its rows and render a live broadcast. the gap-guessing in the old engine and the separate jetstream loop are gone. a subscriber that lags the broadcast channel now catches up from the db instead of being disconnected. ConsumerTooSlow is only sent when the output stays full for the send timeout. /subscribe now broadcasts a historical marker after backfilled commits. subscribers that ask for wantedEventTypes=historical explicitly get backfills that finish while they're connected from the live tail. the default is unchanged: cursor replay returns both, the tail stays live, and default subscribers skip the marker reads entirely. a lagging subscriber's catch-up keeps the live filter: backfills committed past the head it connected at stay out of a default tail. a db read that fails ends the stream, and the subscriber resumes from its last event, instead of the replay counting as exhausted and skipping the rows after the failure. --- AGENTS.md | 1 + docs/api/jetstream.md | 6 +- docs/api/stream.md | 2 + src/config.rs | 3 +- src/control/indexer.rs | 9 +- src/control/jetstream.rs | 52 +- src/control/relay.rs | 5 +- src/control/stream.rs | 3 + src/control/stream/engine.rs | 889 ++++++++++++++++++++--------- src/control/stream/indexer.rs | 167 +++--- src/control/stream/interleaving.rs | 282 +++++++++ src/control/stream/jetstream.rs | 627 ++++---------------- src/control/stream/relay.rs | 122 ++-- src/control/stream/types.rs | 164 +++--- src/db/keyspaces.rs | 23 +- src/db/mod.rs | 2 + src/db/outbox/jetstream.rs | 41 +- src/types/event/jetstream.rs | 116 ++-- 18 files changed, 1374 insertions(+), 1140 deletions(-) create mode 100644 src/control/stream/interleaving.rs diff --git a/AGENTS.md b/AGENTS.md index b0fa2be..77fa83a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -100,6 +100,7 @@ Modes (`indexer`, `relay`) are compile-time choices and mutually exclusive (enfo - **Cursors**: Store cursors as big-endian bytes (`u64`/`i64`). - **Compression**: Configurable via `HYDRANT_DATA_COMPRESSION` (`lz4`, `zstd`, `none`). Per-keyspace zstd dictionaries can be trained via `POST /db/train` and are stored as `dict_{keyspace}.bin` in the database directory. - **Keyspaces**: Use the `keys.rs` module to maintain consistent composite key formats. +- **Stream positions**: `/stream` event ids, jetstream keys, and relay seqs are only assigned inside `db::sequencer::Sequencer::commit`. Stage stream events in the transaction's `outbox` (`txn.outbox.stream`, `.jetstream`, `.relay`) and let `Txn::commit` give them positions and broadcast them; never write event rows or broadcast stream events directly. This keeps position order equal to commit order, which the shared stream engine (`control::stream::engine`) relies on. - **Schema evolution**: Treat versioned wire types in `src/types.rs` as frozen snapshots once shipped. If a stored type changes shape, add a new versioned type, export the newest version for live code, and add a forward migration that explicitly deserializes the previous version and writes the new one. Do not modify older migration input/output types in place. For example, if `RepoState` changes after `v7`, add `v8::RepoState` and a `v7 -> v8` migration instead of mutating `v7`. ### Rate limiting & XRPC Client integration diff --git a/docs/api/jetstream.md b/docs/api/jetstream.md index c2b421f..9f516a7 100644 --- a/docs/api/jetstream.md +++ b/docs/api/jetstream.md @@ -14,7 +14,7 @@ subscribe to the jetstream websocket stream. | :--- | :--- | :--- | | `wantedCollections` | seq | list of collection NSIDs to receive (e.g. `app.bsky.feed.post`). supports namespace wildcards (e.g. `app.bsky.feed.*`). | | `wantedDids` | seq | list of DIDs to receive (e.g. `did:plc:abc123xyz`). | -| `wantedEventTypes` | seq | list of event types to receive: `live` or `historical`. if not specified, both are returned. | +| `wantedEventTypes` | seq | list of event types to receive: `live` or `historical`. if not specified, cursor replay returns both and the live tail only `live` (see [live stream scope](#additional-details)). | | `maxMessageSizeBytes` | integer | filters out events whose serialized JSON size exceeds this value. | | `cursor` | integer | unix microseconds timestamp (`time_us`) to replay historical events from. | | `compress` | boolean | if `true` (or if header `Socket-Encoding` contains `zstd`), compresses frames using zstd and sends them as binary websocket frames. | @@ -108,5 +108,7 @@ fired when a repository's status changes (active, deactivated, or deleted): ### additional details -- **live stream scope**: jetstream subscribers receive both live firehose events (`live: true`) and historical backfill / sync replay events (`live: false`) by default. be aware that historical backfills will interleave with live events unless filtered using the `wantedEventTypes` query parameter. +- **live stream scope**: by default, cursor replay returns both live firehose events (`live: true`) and historical backfill / sync replay events (`live: false`), interleaved in commit order, while the live tail carries only live events, even after the subscriber falls behind and catches up from the database. filter them with the `wantedEventTypes` query parameter. asking for `historical` explicitly also delivers backfills that finish while the socket is open, from the live tail. +- **ordering**: events are delivered in commit order, which is also `time_us` order, so a client that reconnects with its last `time_us` as the cursor misses nothing the socket would have sent. +- **falling behind**: a subscriber that falls behind the live broadcast catches up from the database at its own pace instead of being disconnected. - **slow consumers**: if a socket's send buffer remains full for longer than the configured timeout, the server sends a `{"type":"error","error":"ConsumerTooSlow"}` message and drops the connection. diff --git a/docs/api/stream.md b/docs/api/stream.md index 2b629d9..a61ef99 100644 --- a/docs/api/stream.md +++ b/docs/api/stream.md @@ -126,6 +126,8 @@ account events are live-only and are not available through cursor replay. - the server sends an empty ping every 30 seconds and responds to client ping frames with pong frames. - client text and binary frames are not accepted; hydrant closes the connection when it receives one. +- a subscriber that falls behind the live broadcast catches up from the database at its own pace. `record` events are never skipped, but `identity` and `account` events it fell behind on are not replayed. +- if hydrant can't read the next events from the database, it closes the connection rather than skip them. reconnect with the last `id` you received as the cursor. - if sending blocks for `HYDRANT_STREAM_SEND_TIMEOUT` (30 seconds by default), hydrant sends an error text frame and closes the connection: ```json diff --git a/src/config.rs b/src/config.rs index 1fa189a..723d112 100644 --- a/src/config.rs +++ b/src/config.rs @@ -310,7 +310,8 @@ pub struct Config { /// /// 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. + /// maximum number of live stream events buffered in memory per subscriber while it waits + /// to send, before it falls back to catching up from the database. /// set via `HYDRANT_STREAM_PENDING_EVENT_LIMIT`. pub stream_pending_event_limit: usize, /// maximum time a subscriber may block stream delivery before being disconnected. diff --git a/src/control/indexer.rs b/src/control/indexer.rs index 4a61769..ebd963a 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -5,7 +5,8 @@ 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). +/// the stream terminates when hydrant shuts down, or when it can't read the next +/// events from the database: resubscribe with the last event's id as the cursor. pub struct EventStream { receiver: mpsc::Receiver>, wake: tokio::sync::broadcast::Sender, @@ -95,8 +96,10 @@ 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. slow consumers receive - /// [`StreamError::ConsumerTooSlow`] before the stream terminates when possible. + /// the stream ends when the `EventStream` is dropped. a subscriber that falls behind + /// the live broadcast catches up from the database, skipping only the ephemeral + /// events it fell behind on. one that stops reading for `HYDRANT_STREAM_SEND_TIMEOUT` + /// receives [`StreamError::ConsumerTooSlow`] before the stream terminates when possible. pub fn subscribe(&self, cursor: Option) -> EventStream { let (tx, rx) = mpsc::channel(stream::STREAM_CHANNEL_CAPACITY); let state = self.state.clone(); diff --git a/src/control/jetstream.rs b/src/control/jetstream.rs index 89f9088..c4333a9 100644 --- a/src/control/jetstream.rs +++ b/src/control/jetstream.rs @@ -55,6 +55,17 @@ impl JetstreamFilter { pub(crate) fn wants(&self, event: &StoredJetstreamEvent<'_>) -> bool { self.inner.load().wants(event) } + + /// whether backfilled (`live: false`) commits that finish after + /// connecting are sent from the live tail. only when `historical` is + /// asked for explicitly; cursor replay sends them by default. + pub(crate) fn tails_historical(&self) -> bool { + self.inner + .load() + .wanted_event_types + .as_ref() + .is_some_and(|wanted| wanted.historical) + } } #[derive(Clone, Default)] @@ -101,11 +112,11 @@ impl JetstreamSubscriberOptions { pub(crate) fn wants(&self, event: &StoredJetstreamEvent<'_>) -> bool { if let Some(wanted) = &self.wanted_event_types { - let is_live = event.is_live(); - if is_live && !wanted.live { - return false; - } - if !is_live && !wanted.historical { + let wanted_type = match event.is_live() { + true => wanted.live, + false => wanted.historical, + }; + if !wanted_type { return false; } } @@ -303,6 +314,37 @@ mod tests { ); } + #[test] + #[cfg(feature = "indexer_stream")] + fn backfilled_commits_tail_only_when_asked_for() { + use jacquard_common::CowStr; + + let did = TrimmedDid::from(&Did::new_static("did:plc:abc123").unwrap()).into_static(); + let commit = |live| StoredJetstreamEvent::Commit { + did: did.clone(), + collection: CowStr::Borrowed("app.bsky.feed.post"), + event_id: 1, + live, + }; + let filter = |types: &[&str]| { + let types: Vec<_> = types.iter().map(|t| t.to_string()).collect(); + JetstreamFilter::new(JetstreamSubscriberOptions::parse(&[], &[], 0, &types).unwrap()) + }; + + // (types, passes live, passes historical, tails historical) + for (types, live, historical, tails) in [ + (&[][..], true, true, false), + (&["live"][..], true, false, false), + (&["historical"][..], false, true, true), + (&["live", "historical"][..], true, true, true), + ] { + let filter = filter(types); + assert_eq!(filter.wants(&commit(true)), live, "{types:?}"); + assert_eq!(filter.wants(&commit(false)), historical, "{types:?}"); + assert_eq!(filter.tails_historical(), tails, "{types:?}"); + } + } + #[test] fn wanted_did_limit_counts_unique_dids() { let dids = vec!["did:plc:abc123".to_string(); 10_001]; diff --git a/src/control/relay.rs b/src/control/relay.rs index 02de5a1..6c0eb68 100644 --- a/src/control/relay.rs +++ b/src/control/relay.rs @@ -37,8 +37,9 @@ impl Hydrant { /// - if `cursor` is `None`, streaming starts from the current head (live tail only). /// - if `cursor` is `Some(seq)`, all persisted events from that seq onward are replayed first. /// - /// slow consumers receive [`RelayStreamError::ConsumerTooSlow`] before the stream terminates - /// when possible. + /// a subscriber that falls behind the live broadcast catches up from the database. one + /// that stops reading for `HYDRANT_STREAM_SEND_TIMEOUT` receives + /// [`RelayStreamError::ConsumerTooSlow`] before the stream terminates when possible. pub fn subscribe_repos(&self, cursor: Option) -> RelayEventStream { let (tx, rx) = mpsc::channel(stream::STREAM_CHANNEL_CAPACITY); let state = self.state.clone(); diff --git a/src/control/stream.rs b/src/control/stream.rs index e6cdc66..1efb858 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -26,3 +26,6 @@ pub(crate) use jetstream::{ pub(crate) use engine::*; #[cfg(any(feature = "indexer_stream", feature = "relay", feature = "jetstream"))] pub(crate) use types::*; + +#[cfg(all(test, feature = "indexer_stream"))] +mod interleaving; diff --git a/src/control/stream/engine.rs b/src/control/stream/engine.rs index c8fcf34..d84e268 100644 --- a/src/control/stream/engine.rs +++ b/src/control/stream/engine.rs @@ -1,184 +1,231 @@ +//! the replay-then-tail loop every event stream runs. +//! +//! positions come from `db::sequencer`, so broadcasts arrive in position +//! order, and every committed row is either broadcast or covered by a later +//! persisted marker. a subscriber only has to take broadcasts in the order +//! they arrive, reading the db for a marker's rows, and one that falls behind +//! the broadcast channel can always resume from the db at its position. + +use std::collections::VecDeque; use std::fmt; use std::num::NonZeroUsize; use std::time::{Duration, Instant}; +use tokio::sync::broadcast::error::{RecvError, TryRecvError}; use tokio::sync::mpsc::error::TrySendError; use tokio::sync::{broadcast, mpsc}; -use tracing::warn; +use tracing::{debug, warn}; use super::types::{ - PendingLiveEvents, ReplayChunk, STREAM_SEND_RETRY_PAUSE, SendOutcome, StreamBroadcast, - StreamOptions, StreamTooSlow, + ReplayChunk, ReplayFailed, STREAM_SEND_RETRY_PAUSE, StreamOptions, StreamTooSlow, }; -pub(crate) fn run_ordered_stream( - tx: mpsc::Sender>, - mut event_rx: broadcast::Receiver, - mut current_seq: Option, - mut catch_up_target: Option, - opts: StreamOptions, - mut read_replay_chunk: impl FnMut(Option, u64, usize) -> ReplayChunk, - mut live_to_output: impl FnMut(B) -> Option, -) where - B: StreamBroadcast + Clone, - E: From + fmt::Display, -{ - let mut replay_gap_target = None; - let mut pending = PendingLiveEvents::new(opts.pending_event_limit); - let mut replay_blocked_since: Option = None; +/// an event broadcast after its commit. +pub(crate) trait StreamBroadcast: Clone { + /// where the event sits in commit order. + type Position: Ord + Copy; - loop { - if tx.is_closed() { - return; - } + fn position(&self) -> Self::Position; - if let Err(err) = drain_pending_broadcasts(&mut event_rx, &mut pending) { - send_stream_error(&tx, err.into()); - return; - } + /// a persisted marker carries no event: every row through its position is + /// committed, and the ones since the previous broadcast weren't broadcast. + fn is_marker(&self) -> bool; +} + +/// what one stream endpoint reads and sends. +pub(crate) trait StreamSource { + type Broadcast: StreamBroadcast; + type Output; + + /// up to `limit` outputs for the committed rows after `after`, through + /// `through` when it is set. + fn read( + &mut self, + after: Option>, + through: Option>, + limit: usize, + ) -> Result>, ReplayFailed>; + + /// the output for a live broadcast, or none when there's nothing to send. + fn render(&mut self, event: Self::Broadcast) -> Option; + + /// whether a marker's rows are worth reading. a subscriber that filters + /// them all out can skip the read. + fn wants_marker(&self, _marker: &Self::Broadcast) -> bool { + true + } +} - drop_delivered_pending(&mut pending, current_seq); +pub(crate) type Position = <::Broadcast as StreamBroadcast>::Position; - if let Some(target) = replay_gap_target { - advance_replay_gap(&mut current_seq, target, &pending); - if current_seq.is_some_and(|seq| seq >= target) { - replay_gap_target = None; - } +/// where a subscriber begins. +pub(crate) struct Start

{ + /// the last position it has seen, none for before the first event + after: Option

, + /// replay the db through this position before tailing + catch_up: Option

, +} + +impl Start

{ + /// tail live broadcasts after `head`. + pub(crate) fn live(head: Option

) -> Self { + Self { + after: head, + catch_up: None, } + } - if catch_up_target.is_none() { - catch_up_target = pending.take_persisted_after(current_seq); + /// replay the rows after `after` through `head`, then tail. + pub(crate) fn replay(after: Option

, head: Option

) -> Self { + Self { + after, + catch_up: head.filter(|head| after < Some(*head)), } + } +} - let pending_next_seq = pending.next_sequence(); - let pending_is_ready = - pending_next_seq.is_some_and(|seq| seq == next_expected_seq(current_seq)); +enum CatchUp

{ + Through(P), + ToEnd, +} - if let Some(target) = catch_up_target.filter(|_| !pending_is_ready) { - let effective_target = pending_next_seq - .and_then(|seq| seq.checked_sub(1)) - .map(|before_pending| before_pending.min(target)) - .unwrap_or(target); +/// the stream is over: the subscriber hung up, was sent a closing error, or +/// the db couldn't be read. +struct Ended; - let Some(chunk_size) = replay_chunk_size_for(&tx, opts.replay_chunk_size) else { - if let Err(err) = drain_pending_broadcasts(&mut event_rx, &mut pending) { - send_stream_error(&tx, err.into()); - return; - } - if let Err(err) = note_replay_blocked(&mut replay_blocked_since, opts.send_timeout) - { - send_stream_error(&tx, err.into()); - return; - } - std::thread::sleep(STREAM_SEND_RETRY_PAUSE); - continue; - }; - clear_replay_blocked(&mut replay_blocked_since); - - let chunk = read_replay_chunk(current_seq, effective_target, chunk_size); - current_seq = chunk.last_seen_seq.or(current_seq); - - 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.into()); - return; - } - } - } +/// subscribe `event_rx` before reading the head `start` was built from, so +/// nothing committed in between is missed. +pub(crate) fn run_stream( + tx: mpsc::Sender>, + event_rx: broadcast::Receiver, + start: Start>, + opts: StreamOptions, + mut source: S, +) where + S: StreamSource, + E: From + fmt::Display, +{ + let mut live = Live::new(event_rx, opts.pending_event_limit); + let mut after = start.after; + let mut catch_up = start.catch_up.map(CatchUp::Through); - if chunk.exhausted || current_seq.is_some_and(|seq| seq >= 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); - } + loop { + if tx.is_closed() { + return; + } + if let Some(target) = catch_up.take() { + if replay(&tx, &mut source, &mut live, &mut after, target, opts).is_err() { + return; + } continue; } - if let Some(event) = pending.pop_front() { - let seq = event.sequence(); - if !stream_seq_after(seq, current_seq) { + let event = match live.next() { + Next::Event(event) => event, + Next::Lagged => { + debug!("stream subscriber fell behind the broadcast, catching up from the db"); + catch_up = Some(CatchUp::ToEnd); continue; } + Next::Closed => return, + }; - let expected = next_expected_seq(current_seq); - if seq != expected { - catch_up_target = seq.checked_sub(1); - pending.push_front(event); - continue; + let position = event.position(); + if after.is_some_and(|after| position <= after) { + continue; + } + if event.is_marker() { + if source.wants_marker(&event) { + catch_up = Some(CatchUp::Through(position)); + } else { + after = Some(position); } + continue; + } + if let Some(output) = source.render(event) + && send(&tx, output, opts, &mut live).is_err() + { + return; + } + after = Some(position); + } +} - let Some(out_event) = live_to_output(event) else { - catch_up_target = Some(seq); - continue; - }; +/// send the db's rows after `after` until `target`, moving `after` along. +fn replay( + tx: &mpsc::Sender>, + source: &mut S, + live: &mut Live, + after: &mut Option>, + target: CatchUp>, + opts: StreamOptions, +) -> Result<(), Ended> +where + S: StreamSource, + E: From + fmt::Display, +{ + let through = match target { + CatchUp::Through(position) => Some(position), + CatchUp::ToEnd => None, + }; + let mut blocked_since: Option = None; - match send_stream_event(&tx, out_event, &mut event_rx, &mut pending, opts) { - Ok(SendOutcome::Sent) => { - current_seq = Some(seq); - } - Ok(SendOutcome::ReceiverDropped) => return, - Err(err) => { - send_stream_error(&tx, err.into()); - return; - } + loop { + // skip the db read while the output channel is saturated + let Some(limit) = replay_chunk_size_for(tx, opts.replay_chunk_size) else { + live.drain(); + if let Err(err) = note_replay_blocked(&mut blocked_since, opts.send_timeout) { + send_stream_error(tx, err.into()); + return Err(Ended); } + std::thread::sleep(STREAM_SEND_RETRY_PAUSE); continue; - } + }; + blocked_since = None; - match event_rx.blocking_recv() { - Ok(event) => { - if event.is_persisted_marker() { - let seq = event.sequence(); - if stream_seq_after(seq, current_seq) { - catch_up_target = Some(seq); - } - continue; - } + let chunk = source + .read(*after, through, limit) + .map_err(|ReplayFailed| Ended)?; + for output in chunk.events { + send(tx, output, opts, live)?; + } + *after = (*after).max(chunk.last_seen); - if let Err(err) = pending.push(event) { - send_stream_error(&tx, err.into()); - return; - } - } - Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { - let err = StreamTooSlow::lagged(skipped); - warn!(%err, "closing slow stream subscriber"); - send_stream_error(&tx, err.into()); - return; - } - Err(tokio::sync::broadcast::error::RecvError::Closed) => break, + if chunk.exhausted { + // every row through the target is committed, so none are left + *after = (*after).max(through); + return Ok(()); + } + if !opts.replay_chunk_pause.is_zero() { + // emergency compatibility knob; zero by default + std::thread::sleep(opts.replay_chunk_pause); } } } -fn send_stream_event( +/// send one output, waiting out a full channel for up to the send timeout. +fn send( tx: &mpsc::Sender>, - event: O, - event_rx: &mut broadcast::Receiver, - pending: &mut PendingLiveEvents, + output: O, opts: StreamOptions, -) -> Result + live: &mut Live, +) -> Result<(), Ended> where - B: StreamBroadcast + Clone, + B: Clone, + E: From + fmt::Display, { - let mut item = Ok(event); + let mut item = Ok(output); let started = Instant::now(); - loop { match tx.try_send(item) { - Ok(()) => return Ok(SendOutcome::Sent), - Err(TrySendError::Closed(_)) => return Ok(SendOutcome::ReceiverDropped), + Ok(()) => return Ok(()), + Err(TrySendError::Closed(_)) => return Err(Ended), Err(TrySendError::Full(returned)) => { item = returned; - drain_pending_broadcasts(event_rx, pending)?; + live.drain(); if started.elapsed() >= opts.send_timeout { - return Err(StreamTooSlow::send_timeout(opts.send_timeout)); + send_stream_error(tx, StreamTooSlow::send_timeout(opts.send_timeout).into()); + return Err(Ended); } std::thread::sleep(STREAM_SEND_RETRY_PAUSE); } @@ -186,68 +233,69 @@ where } } -pub(crate) fn send_stream_error(tx: &mpsc::Sender>, err: E) -where - E: fmt::Display, -{ - warn!(%err, "closing stream subscriber"); - let _ = tx.try_send(Err(err)); +/// the broadcasts a subscriber hasn't taken yet: the broadcast channel, plus +/// a bounded backlog drained from it while a send waits, so a briefly slow +/// subscriber doesn't lag the channel. +struct Live { + rx: broadcast::Receiver, + backlog: VecDeque, + limit: usize, + /// broadcasts after the backlog were lost + lagged: bool, } -fn drain_pending_broadcasts( - event_rx: &mut broadcast::Receiver, - pending: &mut PendingLiveEvents, -) -> Result<(), StreamTooSlow> -where - B: StreamBroadcast + Clone, -{ - 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(StreamTooSlow::lagged(skipped)); +enum Next { + Event(B), + Lagged, + Closed, +} + +impl Live { + fn new(rx: broadcast::Receiver, limit: usize) -> Self { + Self { + rx, + backlog: VecDeque::new(), + limit, + lagged: false, + } + } + + /// move waiting broadcasts into the backlog, up to its limit. past it the + /// channel holds them, until it lags. + fn drain(&mut self) { + while !self.lagged && self.backlog.len() < self.limit { + match self.rx.try_recv() { + Ok(event) => self.backlog.push_back(event), + Err(TryRecvError::Lagged(_)) => self.lagged = true, + Err(TryRecvError::Empty | TryRecvError::Closed) => return, } } } -} -fn drop_delivered_pending(pending: &mut PendingLiveEvents, current_id: Option) -where - B: StreamBroadcast, -{ - while pending - .next_sequence() - .is_some_and(|id| !stream_seq_after(id, current_id)) - { - pending.pop_front(); + fn next(&mut self) -> Next { + if let Some(event) = self.backlog.pop_front() { + return Next::Event(event); + } + if std::mem::take(&mut self.lagged) { + return Next::Lagged; + } + match self.rx.blocking_recv() { + Ok(event) => Next::Event(event), + Err(RecvError::Lagged(_)) => Next::Lagged, + Err(RecvError::Closed) => Next::Closed, + } } } -fn advance_replay_gap(current_id: &mut Option, target: u64, pending: &PendingLiveEvents) +fn send_stream_error(tx: &mpsc::Sender>, err: E) where - B: StreamBroadcast, + E: fmt::Display, { - let next_pending_id = pending.next_sequence().filter(|id| *id <= target); - let advance_to = next_pending_id - .and_then(|id| id.checked_sub(1)) - .unwrap_or(target); - - if stream_seq_after(advance_to, *current_id) { - *current_id = Some(advance_to); - } -} - -pub(crate) fn stream_seq_after(id: u64, current_id: Option) -> bool { - current_id.is_none_or(|current| id > current) -} - -pub(crate) fn next_expected_seq(current_id: Option) -> u64 { - current_id.map(|id| id.saturating_add(1)).unwrap_or(0) + warn!(%err, "closing stream subscriber"); + let _ = tx.try_send(Err(err)); } -pub(crate) fn replay_chunk_size_for( +fn replay_chunk_size_for( tx: &mpsc::Sender, configured: Option, ) -> Option { @@ -266,7 +314,7 @@ pub(crate) fn replay_chunk_size_for( Some(desired.min(available).max(1)) } -pub(crate) fn note_replay_blocked( +fn note_replay_blocked( blocked_since: &mut Option, timeout: Duration, ) -> Result<(), StreamTooSlow> { @@ -279,119 +327,438 @@ pub(crate) fn note_replay_blocked( Ok(()) } -pub(crate) fn clear_replay_blocked(blocked_since: &mut Option) { - *blocked_since = None; -} - #[cfg(all(test, any(feature = "indexer_stream", feature = "relay")))] mod tests { use super::*; + use std::collections::BTreeSet; + use std::sync::{Arc, Mutex}; + use std::thread::JoinHandle; #[derive(Clone, Debug)] - enum TestBroadcast { - Persisted(u64), - Live(u64), + enum Ev { + /// a committed row, broadcast live + Row(u64), + /// a broadcast that has no row + Ephemeral(u64), + Marker(u64), } - impl StreamBroadcast for TestBroadcast { - fn sequence(&self) -> u64 { + impl StreamBroadcast for Ev { + type Position = u64; + + fn position(&self) -> u64 { match self { - Self::Persisted(seq) | Self::Live(seq) => *seq, + Self::Row(id) | Self::Ephemeral(id) | Self::Marker(id) => *id, } } - fn is_persisted_marker(&self) -> bool { - matches!(self, Self::Persisted(_)) + fn is_marker(&self) -> bool { + matches!(self, Self::Marker(_)) } } - fn chunk_size(n: usize) -> Option { - NonZeroUsize::new(n) + /// the db and broadcast channel the way the sequencer drives them: + /// commit, then broadcast, in position order. + struct Log { + db: Mutex>, + tx: Mutex>>, + next: Mutex, + reads: Mutex, + reads_fail: Mutex, } - #[test] - fn ordered_stream_replays_chunks_then_live_tail() { - let opts = StreamOptions { - replay_chunk_size: chunk_size(2), + impl Log { + fn new(capacity: usize) -> (Arc, broadcast::Receiver) { + let (tx, rx) = broadcast::channel(capacity); + let log = Self { + db: Mutex::new(BTreeSet::new()), + tx: Mutex::new(Some(tx)), + next: Mutex::new(1), + reads: Mutex::new(0), + reads_fail: Mutex::new(false), + }; + (Arc::new(log), rx) + } + + fn take(&self) -> u64 { + let mut next = self.next.lock().unwrap(); + *next += 1; + *next - 1 + } + + fn broadcast(&self, ev: Ev) { + if let Some(tx) = &*self.tx.lock().unwrap() { + let _ = tx.send(ev); + } + } + + fn live(&self) -> u64 { + let id = self.take(); + self.db.lock().unwrap().insert(id); + self.broadcast(Ev::Row(id)); + id + } + + fn ephemeral(&self) -> u64 { + let id = self.take(); + self.broadcast(Ev::Ephemeral(id)); + id + } + + /// rows written without a broadcast, like a backfill + fn backfill(&self, n: usize) -> Vec { + let ids: Vec<_> = (0..n).map(|_| self.take()).collect(); + self.db.lock().unwrap().extend(&ids); + self.broadcast(Ev::Marker(*ids.last().unwrap())); + ids + } + + fn head(&self) -> Option { + self.next + .lock() + .unwrap() + .checked_sub(1) + .filter(|id| *id > 0) + } + + fn close(&self) { + self.tx.lock().unwrap().take(); + } + + fn fail_reads(&self) { + *self.reads_fail.lock().unwrap() = true; + } + + /// wait until the subscriber has taken every broadcast off the channel + fn wait_taken(&self) { + let deadline = Instant::now() + Duration::from_secs(5); + while self + .tx + .lock() + .unwrap() + .as_ref() + .is_some_and(|tx| !tx.is_empty()) + { + assert!( + Instant::now() < deadline, + "the subscriber never took a broadcast" + ); + std::thread::sleep(Duration::from_millis(1)); + } + } + } + + /// runs before each db read, to commit at exact points of a replay + type OnRead = Box; + + struct Source { + log: Arc, + wants_markers: bool, + on_read: OnRead, + } + + impl StreamSource for Source { + type Broadcast = Ev; + type Output = u64; + + fn read( + &mut self, + after: Option, + through: Option, + limit: usize, + ) -> Result, ReplayFailed> { + let call = { + let mut reads = self.log.reads.lock().unwrap(); + *reads += 1; + *reads - 1 + }; + (self.on_read)(&self.log, call); + if *self.log.reads_fail.lock().unwrap() { + return Err(ReplayFailed); + } + let start = after.map_or(0, |after| after + 1); + let end = through.unwrap_or(u64::MAX); + let db = self.log.db.lock().unwrap(); + let mut rows = db.range(start..=end).copied(); + let events: Vec<_> = rows.by_ref().take(limit).collect(); + Ok(ReplayChunk { + last_seen: events.last().copied(), + exhausted: rows.next().is_none(), + events, + }) + } + + fn render(&mut self, event: Ev) -> Option { + match event { + Ev::Row(id) | Ev::Ephemeral(id) => Some(id), + Ev::Marker(_) => unreachable!("markers are handled by the engine"), + } + } + + fn wants_marker(&self, _marker: &Ev) -> bool { + self.wants_markers + } + } + + type Out = mpsc::Receiver>; + + fn opts() -> StreamOptions { + StreamOptions { + replay_chunk_size: NonZeroUsize::new(2), replay_chunk_pause: Duration::ZERO, pending_event_limit: 4, - send_timeout: Duration::from_secs(1), + send_timeout: Duration::from_secs(5), + } + } + + fn spawn( + log: &Arc, + event_rx: broadcast::Receiver, + start: Start, + out_capacity: usize, + opts: StreamOptions, + on_read: impl FnMut(&Log, usize) + Send + 'static, + ) -> (JoinHandle<()>, Out) { + let (tx, out) = mpsc::channel(out_capacity); + let source = Source { + log: log.clone(), + wants_markers: true, + on_read: Box::new(on_read), }; - let (out_tx, mut out_rx) = mpsc::channel(16); - let (broadcast_tx, broadcast_rx) = broadcast::channel(16); - - let handle = std::thread::spawn(move || { - run_ordered_stream::( - out_tx, - broadcast_rx, - Some(0), - Some(5), - opts, - |current, target, chunk_size| { - let start = current.map(|seq| seq.saturating_add(1)).unwrap_or(0); - if start > target { - return ReplayChunk { - events: Vec::new(), - last_seen_seq: current, - exhausted: true, - }; - } + let handle = std::thread::spawn(move || run_stream(tx, event_rx, start, opts, source)); + (handle, out) + } - let end = target.min(start.saturating_add(chunk_size as u64).saturating_sub(1)); - let events = (start..=end).collect::>(); - ReplayChunk { - last_seen_seq: events.last().copied().or(current), - exhausted: end >= target, - events, + fn take(out: &mut Out, n: usize) -> Vec { + (0..n) + .map(|_| { + let deadline = Instant::now() + Duration::from_secs(5); + loop { + match out.try_recv() { + Ok(item) => break item.expect("stream error"), + Err(mpsc::error::TryRecvError::Empty) if Instant::now() < deadline => { + std::thread::sleep(Duration::from_millis(5)); + } + Err(err) => panic!("no output: {err}"), } - }, - |event| match event { - TestBroadcast::Persisted(_) => None, - TestBroadcast::Live(seq) => Some(seq), - }, - ); - }); + } + }) + .collect() + } - broadcast_tx.send(TestBroadcast::Live(6)).unwrap(); - broadcast_tx.send(TestBroadcast::Persisted(6)).unwrap(); - drop(broadcast_tx); + fn finish(log: &Log, handle: JoinHandle<()>, mut out: Out) { + log.close(); + handle.join().unwrap(); + assert!(out.try_recv().is_err(), "unexpected extra output"); + } - let mut out = Vec::new(); - while let Some(item) = out_rx.blocking_recv() { - out.push(item.unwrap()); + /// wait for the stream thread to end by itself + fn ends(handle: JoinHandle<()>) { + let deadline = Instant::now() + Duration::from_secs(5); + while !handle.is_finished() { + assert!(Instant::now() < deadline, "the stream never ended"); + std::thread::sleep(Duration::from_millis(5)); } handle.join().unwrap(); + } - assert_eq!(out, vec![1, 2, 3, 4, 5, 6]); + fn nothing_on_read(_: &Log, _: usize) {} + + #[test] + fn replays_in_chunks_then_tails() { + let (log, rx) = Log::new(16); + let rows = log.backfill(5); + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(None, log.head()), + 16, + opts(), + nothing_on_read, + ); + + assert_eq!(take(&mut out, 5), rows); + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + } + + #[test] + fn a_marker_reads_its_rows_from_the_db() { + let (log, rx) = Log::new(16); + let (handle, mut out) = spawn( + &log, + rx, + Start::live(log.head()), + 16, + opts(), + nothing_on_read, + ); + + let first = log.live(); + let rows = log.backfill(3); + let last = log.live(); + let expected: Vec<_> = [first].into_iter().chain(rows).chain([last]).collect(); + assert_eq!(take(&mut out, 5), expected); + finish(&log, handle, out); + } + + #[test] + fn broadcast_only_events_keep_their_place() { + let (log, rx) = Log::new(16); + let (handle, mut out) = spawn( + &log, + rx, + Start::live(log.head()), + 16, + opts(), + nothing_on_read, + ); + + let rows = log.backfill(2); + let ephemeral = log.ephemeral(); + let more = log.backfill(2); + let expected: Vec<_> = rows.into_iter().chain([ephemeral]).chain(more).collect(); + assert_eq!(take(&mut out, 5), expected); + finish(&log, handle, out); } #[test] - fn ordered_stream_closes_when_output_queue_is_full() { + fn commits_during_a_replay_are_sent_once() { + let (log, rx) = Log::new(16); + let rows = log.backfill(3); + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(None, log.head()), + 16, + opts(), + |log, call| { + if call == 0 { + // lands past the head the replay started with, broadcast + // before the replay reaches it + log.live(); + log.backfill(1); + } + }, + ); + + let sent = take(&mut out, 5); + assert_eq!(sent[..3], rows[..]); + assert_eq!(sent, (rows[0]..rows[0] + 5).collect::>()); + finish(&log, handle, out); + } + + #[test] + fn a_lagging_subscriber_catches_up_from_the_db() { + let (log, rx) = Log::new(2); + let opts = StreamOptions { + pending_event_limit: 1, + ..opts() + }; + let (handle, mut out) = spawn(&log, rx, Start::live(log.head()), 1, opts, nothing_on_read); + + let first = log.live(); + std::thread::sleep(Duration::from_millis(50)); + // the subscriber is stuck sending while these overflow the channel + // and its backlog + let rest: Vec<_> = (0..8).map(|_| log.live()).collect(); + let expected: Vec<_> = [first].into_iter().chain(rest).collect(); + assert_eq!(take(&mut out, 9), expected); + + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + } + + #[test] + fn the_backlog_holds_broadcasts_while_a_send_waits() { + let (log, rx) = Log::new(2); + let (handle, mut out) = spawn(&log, rx, Start::live(log.head()), 1, opts(), |_, _| { + panic!("a backlog within its limit needs no db read"); + }); + + // fill the output, then leave the subscriber waiting to send the next + let mut sent = vec![log.live()]; + log.wait_taken(); + sent.push(log.ephemeral()); + log.wait_taken(); + // more than the broadcast channel holds, taken into the backlog while + // that send waits. they have no rows, so a lag would lose them + for _ in 0..3 { + sent.push(log.ephemeral()); + log.wait_taken(); + } + assert_eq!(take(&mut out, 5), sent); + finish(&log, handle, out); + } + + #[test] + fn unwanted_markers_are_skipped_without_a_read() { + let (log, rx) = Log::new(16); + let (tx, mut out) = mpsc::channel(16); + let source = Source { + log: log.clone(), + wants_markers: false, + on_read: Box::new(|_, _| panic!("an unwanted marker needs no db read")), + }; + let handle = + std::thread::spawn(move || run_stream(tx, rx, Start::live(None), opts(), source)); + + log.backfill(3); + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + } + + #[test] + fn a_stuck_subscriber_is_closed() { + let (log, rx) = Log::new(16); + let rows = log.backfill(2); let opts = StreamOptions { - replay_chunk_size: chunk_size(2), - replay_chunk_pause: Duration::ZERO, - pending_event_limit: 4, send_timeout: Duration::ZERO, + ..opts() }; - let (out_tx, _out_rx) = mpsc::channel(1); - let (_broadcast_tx, broadcast_rx) = broadcast::channel(16); - - run_ordered_stream::( - out_tx, - broadcast_rx, - Some(0), - Some(2), + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(None, log.head()), + 1, opts, - |_, _, _| ReplayChunk { - events: vec![1, 2], - last_seen_seq: Some(2), - exhausted: true, - }, - |event| match event { - TestBroadcast::Persisted(_) => None, - TestBroadcast::Live(seq) => Some(seq), + nothing_on_read, + ); + handle.join().unwrap(); + + assert_eq!(out.try_recv().unwrap().unwrap(), rows[0]); + assert!(out.try_recv().is_err()); + } + + #[test] + fn a_failed_read_ends_the_stream() { + let (log, rx) = Log::new(16); + let rows = log.backfill(5); + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(None, log.head()), + 16, + opts(), + |log, call| { + if call == 1 { + log.fail_reads(); + } }, ); + + // rather than skip the rows it couldn't read and go on tailing + assert_eq!(take(&mut out, 2), rows[..2]); + ends(handle); + log.live(); + assert!(out.try_recv().is_err()); + } + + fn chunk_size(n: usize) -> Option { + NonZeroUsize::new(n) } #[test] diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index 3f73153..7e59cd4 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -14,7 +14,7 @@ use jacquard_common::{CowStr, IntoStatic, RawData}; use jacquard_repo::DAG_CBOR_CID_CODEC; use sha2::{Digest, Sha256}; -use super::{ReplayChunk, StreamBroadcast, StreamOptions, run_ordered_stream, stream_seq_after}; +use super::{ReplayChunk, ReplayFailed, Start, StreamOptions, StreamSource, run_stream}; use crate::control::{Event, StreamError}; pub(crate) fn event_stream_thread( @@ -23,112 +23,79 @@ pub(crate) fn event_stream_thread( cursor: Option, opts: StreamOptions, ) { - let db = &state.db; - let event_rx = db.stream.event_tx.subscribe(); - let head = db.stream.ids.head(); - let current_id = match cursor { - Some(c) => c.checked_sub(1), - None => head, + let event_rx = state.db.stream.event_tx.subscribe(); + let head = state.db.stream.ids.head(); + let start = match cursor { + Some(cursor) => Start::replay(cursor.checked_sub(1), head), + None => Start::live(head), }; - let catch_up_target = cursor - .and(head) - .filter(|target| stream_seq_after(*target, current_id)); - let replay_state = state.clone(); - - run_ordered_stream( - tx, - event_rx, - current_id, - catch_up_target, - opts, - move |current_id, target, chunk_size| { - read_event_replay_chunk(&replay_state, current_id, target, chunk_size) - }, - move |event| broadcast_to_event(&state, event), - ); + run_stream(tx, event_rx, start, opts, EventSource { state }); +} + +struct EventSource { + state: Arc, +} + +impl StreamSource for EventSource { + type Broadcast = BroadcastEvent; + type Output = Event; + + fn read( + &mut self, + after: Option, + through: Option, + limit: usize, + ) -> Result, ReplayFailed> { + read_event_replay_chunk(&self.state, after, through, limit) + } + + fn render(&mut self, event: BroadcastEvent) -> Option { + match event { + BroadcastEvent::Persisted(_) => None, + BroadcastEvent::LiveRecord(evt) => { + let stored = evt.stored.clone(); + stored_to_event(&self.state, evt.id, stored, evt.inline_block.clone()) + } + BroadcastEvent::Ephemeral(evt) => Some(*evt), + } + } } fn read_event_replay_chunk( state: &AppState, - 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 { + after: Option, + through: Option, + limit: usize, +) -> Result, ReplayFailed> { + let start = after.map_or(0, |id| id.saturating_add(1)); + let end = through.unwrap_or(u64::MAX); + if start > end { + return Ok(ReplayChunk { events: Vec::new(), - last_seen_seq: current_id, + last_seen: after, exhausted: true, - }; + }); } - let mut events = Vec::with_capacity(chunk_size); - let mut last_seen_seq = current_id; - let mut exhausted = false; - let max_scanned = chunk_size.saturating_mul(4).max(chunk_size); - let mut scanned = 0usize; - let mut iter = state + let rows = state .db .stream - .event_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_seq = 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_seq, - exhausted, - } -} - -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), - } + .event_range(keys::event_key(start)..=keys::event_key(end)); + ReplayChunk::read( + rows, + after, + limit, + |key| { + let id = key.try_into().map(u64::from_be_bytes); + id.inspect_err(|_| error!("failed to parse event id")).ok() + }, + |id, value| { + let stored = rmp_serde::from_slice::(value) + .inspect_err(|e| error!(err = %e, "failed to deserialize stored event")) + .ok()?; + stored_to_event(state, id, stored, None) + }, + ) } pub(crate) fn stored_to_event( @@ -457,8 +424,8 @@ mod tests { } let state = AppState::new(&config).unwrap(); - let replay = read_event_replay_chunk(&state, None, event_id, 8); - assert_eq!(replay.last_seen_seq, Some(event_id)); + let replay = read_event_replay_chunk(&state, None, Some(event_id), 8).unwrap(); + assert_eq!(replay.last_seen, Some(event_id)); assert_eq!(replay.events.len(), 1); let json = serde_json::to_value(&replay.events[0]).unwrap(); @@ -515,8 +482,8 @@ mod tests { } let state = AppState::new(&config).unwrap(); - let replay = read_event_replay_chunk(&state, None, event_id, 1); - assert_eq!(replay.last_seen_seq, Some(event_id)); + let replay = read_event_replay_chunk(&state, None, Some(event_id), 1).unwrap(); + assert_eq!(replay.last_seen, Some(event_id)); assert_eq!(replay.events.len(), 1); let json = serde_json::to_value(&replay.events[0]).unwrap(); assert_eq!(json["record"]["did"], DID); diff --git a/src/control/stream/interleaving.rs b/src/control/stream/interleaving.rs new file mode 100644 index 0000000..59e9b53 --- /dev/null +++ b/src/control/stream/interleaving.rs @@ -0,0 +1,282 @@ +//! stream subscribers against the real write path, with a live commit landing +//! while a backfill batch is still open: the interleaving that once made +//! both streams skip the backfill (hydrant-abq). + +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use jacquard_common::types::did::Did; +use jacquard_common::types::string::Tid; +use tokio::sync::mpsc; + +use crate::config::Config; +use crate::db::Txn; +use crate::db::types::{DbAction, DbRkey, DbTid}; +use crate::state::AppState; + +use super::StreamOptions; + +const BACKFILLED: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; +const LIVE: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz"; +const COLLECTION: &str = "app.bsky.feed.post"; +const REPO_SIZE: usize = 100; + +fn state() -> (tempfile::TempDir, Arc, StreamOptions) { + let tmp = tempfile::tempdir().unwrap(); + let config = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let state = Arc::new(AppState::new(&config).unwrap()); + (tmp, state, StreamOptions::from_config(&config)) +} + +fn put(records: &mut crate::db::RecordTxn<'_, '_, '_>, n: usize) { + let block = bytes::Bytes::from( + serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "n": n })).unwrap(), + ); + let cid = jacquard_repo::mst::util::compute_cid(block.as_ref()).unwrap(); + let rkey = DbRkey::Str(smol_str::format_smolstr!("r{n}")); + records + .put_record(COLLECTION, &rkey, &cid, &block, DbAction::Create) + .unwrap(); +} + +fn rev() -> DbTid { + DbTid::from(&Tid::now_0()) +} + +/// the backfill worker's write: a whole repo in one uncommitted batch. +fn stage_backfill(state: &AppState) -> Txn<'_> { + let did = Did::new_static(BACKFILLED).unwrap(); + let mut txn = Txn::new(&state.db); + let mut records = txn.backfill_records(state, &rev(), &did).unwrap(); + (0..REPO_SIZE).for_each(|n| put(&mut records, n)); + records.finish().unwrap(); + txn +} + +/// one live commit, the way the indexer shard writes it. +fn commit_live(state: &AppState) { + commit_live_records(state, REPO_SIZE..REPO_SIZE + 1); +} + +/// one live commit writing a record for each of `records`. +fn commit_live_records(state: &AppState, records: std::ops::Range) { + let did = Did::new_static(LIVE).unwrap(); + let mut txn = Txn::new(&state.db); + let mut staged = txn.records(state, &rev(), &did).unwrap(); + records.for_each(|n| put(&mut staged, n)); + staged.finish().unwrap(); + txn.commit().unwrap(); +} + +/// dids of what a subscriber sent, collected until `until` holds or a +/// deadline passes. +struct Seen { + rx: mpsc::Receiver, + did_of: fn(T) -> String, + dids: Vec, +} + +impl Seen { + fn until(&mut self, until: impl Fn(&[String]) -> bool) -> &mut Self { + let deadline = Instant::now() + Duration::from_secs(5); + while Instant::now() < deadline && !until(&self.dids) { + match self.rx.try_recv() { + Ok(item) => self.dids.push((self.did_of)(item)), + Err(_) => std::thread::sleep(Duration::from_millis(5)), + } + } + self + } + + fn live(&mut self) -> &mut Self { + self.until(|dids| dids.iter().any(|did| did == LIVE)) + } + + /// wait a moment longer, for anything sent that shouldn't have been + fn settle(&mut self) -> &mut Self { + std::thread::sleep(Duration::from_millis(100)); + while let Ok(item) = self.rx.try_recv() { + self.dids.push((self.did_of)(item)); + } + self + } + + fn counts(&self) -> (usize, usize) { + let count = |did| self.dids.iter().filter(|seen| *seen == did).count(); + (count(BACKFILLED), count(LIVE)) + } +} + +fn stream(state: &Arc, opts: StreamOptions, cursor: Option) -> Seen { + let (tx, rx) = mpsc::channel(4096); + let state = state.clone(); + std::thread::spawn(move || super::event_stream_thread(state, tx, cursor, opts)); + Seen { + rx, + did_of: |item| item.unwrap().record.unwrap().did.as_str().to_owned(), + dids: Vec::new(), + } +} + +type StreamItem = Result; + +/// wait until `n` live-only subscribers have subscribed, so they see what +/// follows. +fn wait_subscribed(n: usize, receivers: impl Fn() -> usize) { + let deadline = Instant::now() + Duration::from_secs(5); + while receivers() < n { + assert!(Instant::now() < deadline, "subscriber never subscribed"); + std::thread::sleep(Duration::from_millis(5)); + } +} + +#[test] +fn stream_tail_gets_a_backfill_committed_after_a_live_commit() { + let (_tmp, state, opts) = state(); + let mut seen = stream(&state, opts, None); + wait_subscribed(1, || state.db.stream.event_tx.receiver_count()); + + let backfill = stage_backfill(&state); + commit_live(&state); + seen.live(); + backfill.commit().unwrap(); + + seen.until(|dids| dids.len() > REPO_SIZE).settle(); + assert_eq!(seen.counts(), (REPO_SIZE, 1)); +} + +#[test] +fn stream_replay_gets_a_backfill_committed_after_a_live_commit() { + let (_tmp, state, opts) = state(); + let backfill = stage_backfill(&state); + commit_live(&state); + let mut seen = stream(&state, opts, Some(0)); + seen.live(); + backfill.commit().unwrap(); + + seen.until(|dids| dids.len() > REPO_SIZE).settle(); + assert_eq!(seen.counts(), (REPO_SIZE, 1)); +} + +#[cfg(feature = "jetstream")] +mod jetstream { + use super::*; + use crate::control::{JetstreamFilter, JetstreamSubscriberOptions}; + + type Item = Result; + + fn subscribe( + state: &Arc, + opts: StreamOptions, + cursor: Option, + wanted_event_types: &[&str], + ) -> Seen { + subscribe_with_capacity(state, opts, 4096, cursor, wanted_event_types) + } + + fn subscribe_with_capacity( + state: &Arc, + opts: StreamOptions, + capacity: usize, + cursor: Option, + wanted_event_types: &[&str], + ) -> Seen { + let (tx, rx) = mpsc::channel(capacity); + let state = state.clone(); + let wanted: Vec<_> = wanted_event_types.iter().map(|t| t.to_string()).collect(); + let filter = + JetstreamFilter::new(JetstreamSubscriberOptions::parse(&[], &[], 0, &wanted).unwrap()); + std::thread::spawn(move || { + super::super::jetstream_stream_thread(state, tx, cursor, filter, opts) + }); + Seen { + rx, + did_of: |item| { + let json: serde_json::Value = serde_json::from_slice(&item.unwrap()).unwrap(); + json["did"].as_str().unwrap().to_owned() + }, + dids: Vec::new(), + } + } + + fn a_second_ago() -> i64 { + chrono::Utc::now().timestamp_micros() - 1_000_000 + } + + #[test] + fn replay_gets_a_backfill_committed_after_a_live_commit() { + let (_tmp, state, opts) = state(); + let cursor = a_second_ago(); + let backfill = stage_backfill(&state); + commit_live(&state); + let mut seen = subscribe(&state, opts, Some(cursor), &["live", "historical"]); + seen.live(); + backfill.commit().unwrap(); + + seen.until(|dids| dids.len() > REPO_SIZE).settle(); + assert_eq!(seen.counts(), (REPO_SIZE, 1)); + } + + #[test] + fn tail_gets_backfills_only_when_asked() { + let (_tmp, state, opts) = state(); + let mut historical = subscribe(&state, opts, None, &["live", "historical"]); + let mut default = subscribe(&state, opts, None, &[]); + wait_subscribed(2, || state.db.jetstream.tx.receiver_count()); + + let backfill = stage_backfill(&state); + commit_live(&state); + backfill.commit().unwrap(); + + historical.until(|dids| dids.len() > REPO_SIZE).settle(); + assert_eq!(historical.counts(), (REPO_SIZE, 1)); + default.live().settle(); + assert_eq!(default.counts(), (0, 1)); + } + + #[test] + fn a_lagging_tail_gets_backfills_only_when_asked() { + // more live events than the broadcast channel holds, per commit + const LAG: usize = 600; + let (_tmp, state, opts) = state(); + let opts = StreamOptions { + pending_event_limit: 1, + ..opts + }; + let mut historical = + subscribe_with_capacity(&state, opts, 64, None, &["live", "historical"]); + let mut default = subscribe_with_capacity(&state, opts, 64, None, &[]); + wait_subscribed(2, || state.db.jetstream.tx.receiver_count()); + + // nobody reads, so both subscribers lag and catch up from the db + commit_live_records(&state, REPO_SIZE..REPO_SIZE + LAG); + stage_backfill(&state).commit().unwrap(); + commit_live_records(&state, REPO_SIZE + LAG..REPO_SIZE + 2 * LAG); + + historical + .until(|dids| dids.len() >= REPO_SIZE + 2 * LAG) + .settle(); + assert_eq!(historical.counts(), (REPO_SIZE, 2 * LAG)); + default.until(|dids| dids.len() >= 2 * LAG).settle(); + assert_eq!(default.counts(), (0, 2 * LAG)); + } + + #[test] + fn replay_sends_backfills_by_default() { + let (_tmp, state, opts) = state(); + let cursor = a_second_ago(); + stage_backfill(&state).commit().unwrap(); + commit_live(&state); + + let mut seen = subscribe(&state, opts, Some(cursor), &[]); + seen.until(|dids| dids.len() > REPO_SIZE).settle(); + assert_eq!(seen.counts(), (REPO_SIZE, 1)); + + let mut live_only = subscribe(&state, opts, Some(cursor), &["live"]); + live_only.live().settle(); + assert_eq!(live_only.counts(), (0, 1)); + } +} diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index b537b75..4efd4f1 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -1,293 +1,134 @@ use bytes::Bytes; -use std::collections::HashSet; -use std::fmt; use std::sync::Arc; -use std::time::Instant; -use tokio::sync::broadcast::error::{RecvError, TryRecvError}; -use tokio::sync::mpsc::error::TrySendError; -use tokio::sync::{broadcast, mpsc}; use tracing::error; use crate::db::keys; use crate::state::AppState; -use crate::types::{JetstreamBroadcast, StoredJetstreamEvent}; +use crate::types::{JetstreamBroadcast, JetstreamPosition, StoredJetstreamEvent}; #[cfg(feature = "relay")] use crate::ingest::stream::{SubscribeReposMessage, decode_frame}; #[cfg(feature = "indexer_stream")] use crate::types::StoredEvent; -use super::{ - ReplayChunk, STREAM_SEND_RETRY_PAUSE, StreamOptions, StreamTooSlow, clear_replay_blocked, - note_replay_blocked, replay_chunk_size_for, send_stream_error, -}; +use super::{ReplayChunk, ReplayFailed, Start, StreamOptions, StreamSource, run_stream}; use crate::control::{JetstreamFilter, JetstreamStreamError}; #[cfg(feature = "indexer_stream")] use super::indexer::stored_to_event; -type ReplayKey = [u8; 16]; - -/// an event's place in the time-ordered jetstream keyspace. -trait ReplayPosition { - fn id(&self) -> u64; - /// the first key a replay resuming after this event reads. - fn key_after(&self) -> ReplayKey; -} - -impl ReplayPosition for JetstreamBroadcast { - fn id(&self) -> u64 { - self.id - } - - fn key_after(&self) -> ReplayKey { - keys::jetstream_event_key(self.time_us as u64, self.id.saturating_add(1)) - } -} - -/// the stream is over: the subscriber hung up or was sent a closing error. -struct Ended; - -/// the tail lost broadcasts before sending any live event, so a replay can -/// recover them from the db. -struct Lagged; - pub(crate) fn jetstream_stream_thread( state: Arc, - tx: mpsc::Sender>, + tx: tokio::sync::mpsc::Sender>, cursor: Option, filter: JetstreamFilter, opts: StreamOptions, ) { let event_rx = state.db.jetstream.tx.subscribe(); - let replay_from = cursor - .filter(|cursor| *cursor <= chrono::Utc::now().timestamp_micros()) - .map(|cursor| keys::jetstream_event_key(cursor.max(0) as u64, 0)); + let head = match state.db.jetstream.head() { + Ok(head) => head, + Err(e) => { + error!(err = %e, "failed to read the Jetstream head"); + return; + } + }; + // a cursor in the future means live, like upstream jetstream + let start = match cursor.filter(|cursor| *cursor <= chrono::Utc::now().timestamp_micros()) { + Some(cursor) => Start::replay(JetstreamPosition::before(cursor.max(0) as u64), head), + None => Start::live(head), + }; + let source = JetstreamSource { + state, + filter, + head, + }; + run_stream(tx, event_rx, start, opts, source); +} - run_jetstream_stream( - &tx, - event_rx, - replay_from, - opts, - |start, chunk_size| read_jetstream_replay_chunk(&state, start, chunk_size), - |event| jetstream_event_to_bytes(&state, event, &filter), - ); +struct JetstreamSource { + state: Arc, + filter: JetstreamFilter, + /// the head when the subscriber connected. rows after it belong to the + /// live tail, even when a lag has them read from the db. + head: Option, } -/// replay the db from `replay_from`, then tail live broadcasts. -/// -/// broadcasts only follow commits, so a scan that starts after draining the -/// broadcast receiver sees every event the drain discarded. replay therefore -/// discards live events instead of buffering them, runs until a scan reaches -/// the end of the keyspace, and the tail skips only what that last scan read. -/// a tail that lags before sending a live event goes back to the db for the -/// same reason. after one live send, broadcast order no longer tells which -/// keys were delivered, so a lag closes the stream instead of skipping events. -fn run_jetstream_stream( - tx: &mpsc::Sender>, - mut event_rx: broadcast::Receiver, - mut replay_from: Option, - opts: StreamOptions, - mut read_chunk: impl FnMut(&ReplayKey, usize) -> ReplayChunk, - mut to_bytes: impl FnMut(B) -> Option, -) where - B: ReplayPosition + Clone, - E: From + fmt::Display, -{ - loop { - let last_scan = match replay_from.as_mut() { - Some(key) => { - match replay(tx, &mut event_rx, key, opts, &mut read_chunk, &mut to_bytes) { - Ok(ids) => ids, - Err(Ended) => return, - } - } - None => HashSet::new(), - }; - match tail( - tx, - &mut event_rx, - &last_scan, - &mut replay_from, - opts, - &mut to_bytes, - ) { - Ok(Lagged) => continue, - Err(Ended) => return, - } +impl JetstreamSource { + /// whether a stored row goes out: what the filter wants, except backfills + /// past the head for subscribers that don't tail them. + fn sends(&self, position: JetstreamPosition, event: &StoredJetstreamEvent<'_>) -> bool { + self.filter.wants(event) + && (event.is_live() || Some(position) <= self.head || self.filter.tails_historical()) } } -/// send the db from `key` until a scan comes back exhausted, returning the ids -/// that final scan read. -fn replay( - tx: &mpsc::Sender>, - event_rx: &mut broadcast::Receiver, - key: &mut ReplayKey, - opts: StreamOptions, - read_chunk: &mut impl FnMut(&ReplayKey, usize) -> ReplayChunk, - to_bytes: &mut impl FnMut(B) -> Option, -) -> Result, Ended> -where - B: ReplayPosition + Clone, - E: From + fmt::Display, -{ - let mut blocked_since: Option = None; - loop { - discard_broadcasts(event_rx); - - // skip the db read when the output channel is already saturated. - let Some(chunk_size) = replay_chunk_size_for(tx, opts.replay_chunk_size) else { - if let Err(err) = note_replay_blocked(&mut blocked_since, opts.send_timeout) { - send_stream_error(tx, err.into()); - return Err(Ended); - } - std::thread::sleep(STREAM_SEND_RETRY_PAUSE); - continue; - }; - clear_replay_blocked(&mut blocked_since); - - let chunk = read_chunk(key, chunk_size); - let mut scanned = HashSet::with_capacity(chunk.events.len()); - for event in chunk.events { - *key = event.key_after(); - scanned.insert(event.id()); - if let Some(bytes) = to_bytes(event) { - send_frame(tx, bytes, opts)?; - } - } - - if chunk.exhausted { - return Ok(scanned); - } - if !opts.replay_chunk_pause.is_zero() { - // emergency compatibility knob; zero by default. - std::thread::sleep(opts.replay_chunk_pause); - } +impl StreamSource for JetstreamSource { + type Broadcast = JetstreamBroadcast; + type Output = Bytes; + + fn read( + &mut self, + after: Option, + through: Option, + limit: usize, + ) -> Result, ReplayFailed> { + read_jetstream_replay_chunk(&self.state, after, through, limit, |position, event| { + self.sends(position, event) + }) } -} -/// send live broadcasts, skipping the ids the last replay scan already sent. -/// clears `replay_from` once a live event is sent. -fn tail( - tx: &mpsc::Sender>, - event_rx: &mut broadcast::Receiver, - last_scan: &HashSet, - replay_from: &mut Option, - opts: StreamOptions, - to_bytes: &mut impl FnMut(B) -> Option, -) -> Result -where - B: ReplayPosition + Clone, - E: From + fmt::Display, -{ - loop { - match event_rx.blocking_recv() { - Ok(event) => { - if last_scan.contains(&event.id()) { - continue; - } - let Some(bytes) = to_bytes(event) else { - continue; - }; - send_frame(tx, bytes, opts)?; - *replay_from = None; - } - Err(RecvError::Lagged(_)) if replay_from.is_some() => return Ok(Lagged), - Err(RecvError::Lagged(skipped)) => { - send_stream_error(tx, StreamTooSlow::lagged(skipped).into()); - return Err(Ended); - } - Err(RecvError::Closed) => return Err(Ended), + fn render(&mut self, event: JetstreamBroadcast) -> Option { + let JetstreamBroadcast::Live(live) = event else { + return None; + }; + if !self.filter.wants(&live.event) { + return None; } + // pre-serialized json avoids a db read + live.json.or_else(|| { + stored_event_to_bytes(&self.state, live.position.time_us as i64, &live.event) + }) } -} - -// drop everything already broadcast. only safe right before a db scan, which -// reads every committed event the drain dropped. -fn discard_broadcasts(event_rx: &mut broadcast::Receiver) { - while let Ok(_) | Err(TryRecvError::Lagged(_)) = event_rx.try_recv() {} -} -// send one frame, waiting out a full channel for up to the send timeout. the -// broadcast receiver is left alone so a tail that falls behind sees the lag. -fn send_frame( - tx: &mpsc::Sender>, - bytes: Bytes, - opts: StreamOptions, -) -> Result<(), Ended> -where - E: From + fmt::Display, -{ - let mut item = Ok(bytes); - let started = Instant::now(); - loop { - match tx.try_send(item) { - Ok(()) => return Ok(()), - Err(TrySendError::Closed(_)) => return Err(Ended), - Err(TrySendError::Full(returned)) => { - if started.elapsed() >= opts.send_timeout { - send_stream_error(tx, StreamTooSlow::send_timeout(opts.send_timeout).into()); - return Err(Ended); - } - item = returned; - std::thread::sleep(STREAM_SEND_RETRY_PAUSE); - } - } + /// markers only cover backfilled commits, which the tail sends only to + /// subscribers that asked for them, so the rest skip the read. + fn wants_marker(&self, _marker: &JetstreamBroadcast) -> bool { + self.filter.tails_historical() } } fn read_jetstream_replay_chunk( state: &AppState, - start_key: &ReplayKey, - chunk_size: usize, -) -> ReplayChunk { - let mut events = Vec::with_capacity(chunk_size); - let mut exhausted = false; - let mut iter = state.db.jetstream.events.range(start_key.as_slice()..); - - while events.len() < chunk_size { - let Some(item) = iter.next() else { - exhausted = true; - break; - }; - - let (k, v) = match item.into_inner() { - Ok(kv) => kv, - Err(e) => { - error!(err = %e, "failed to read Jetstream event from db"); - exhausted = true; - break; - } - }; - let (time_us, id) = match keys::parse_jetstream_event_key(&k) { - Ok(parsed) => parsed, - Err(e) => { - error!(err = %e, "failed to parse Jetstream event key"); - continue; - } - }; - - let event: StoredJetstreamEvent = match rmp_serde::from_slice(&v) { - Ok(event) => event, - Err(e) => { - error!(err = %e, "failed to deserialize Jetstream event"); - continue; + after: Option, + through: Option, + limit: usize, + sends: impl Fn(JetstreamPosition, &StoredJetstreamEvent<'_>) -> bool, +) -> Result, ReplayFailed> { + let start = after.map_or([0; 16], JetstreamPosition::key_after); + let rows = match through { + Some(through) => state.db.jetstream.events.range(start..=through.key()), + None => state.db.jetstream.events.range(start..), + }; + ReplayChunk::read( + rows, + after, + limit, + |key| { + let (time_us, id) = keys::parse_jetstream_event_key(key) + .inspect_err(|e| error!(err = %e, "failed to parse Jetstream event key")) + .ok()?; + Some(JetstreamPosition { time_us, id }) + }, + |position, value| { + let event = rmp_serde::from_slice::(value) + .inspect_err(|e| error!(err = %e, "failed to deserialize Jetstream event")) + .ok()?; + if !sends(position, &event) { + return None; } - }; - events.push(JetstreamBroadcast { - id, - time_us: time_us as i64, - event: event.into_static(), - ephemeral: None, - }); - } - - ReplayChunk { - events, - last_seen_seq: None, - exhausted, - } + stored_event_to_bytes(state, position.time_us as i64, &event) + }, + ) } #[derive(serde::Serialize)] @@ -346,21 +187,14 @@ pub(crate) struct JetstreamAccount<'a> { pub(crate) status: Option, } -fn jetstream_event_to_bytes( +/// render a stored event from the database, for a replay or a live event +/// broadcast without its json. +fn stored_event_to_bytes( state: &AppState, - event: JetstreamBroadcast, - filter: &JetstreamFilter, + time_us: i64, + event: &StoredJetstreamEvent<'_>, ) -> Option { - if !filter.wants(&event.event) { - return None; - } - - // live tailing: use pre-serialized ephemeral bytes to avoid db reads. - if let Some(bytes) = event.ephemeral { - return Some(bytes); - } - - match &event.event { + match event { #[cfg(feature = "indexer_stream")] StoredJetstreamEvent::Commit { event_id, live, .. } => { let bytes = state.db.stream.event(keys::event_key(*event_id)).ok()??; @@ -373,7 +207,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: did_str, - time_us: event.time_us, + time_us, payload: JetstreamPayload::Commit { commit: JetstreamCommit { rev: rec.rev.as_str(), @@ -426,7 +260,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: commit.repo.as_str(), - time_us: event.time_us, + time_us, payload: JetstreamPayload::Commit { commit: JetstreamCommit { rev: commit.rev.as_str(), @@ -460,7 +294,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: did_str, - time_us: event.time_us, + time_us, payload: JetstreamPayload::Account { account: JetstreamAccount { active: account.active, @@ -490,7 +324,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: did_str, - time_us: event.time_us, + time_us, payload: JetstreamPayload::Identity { identity: JetstreamIdentity { did: did_str.to_string(), @@ -514,7 +348,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: did_str.as_str(), - time_us: event.time_us, + time_us, payload: JetstreamPayload::Account { account: JetstreamAccount { active: *active, @@ -538,7 +372,7 @@ fn jetstream_event_to_bytes( let json_event = JetstreamEvent { did: did_str.as_str(), - time_us: event.time_us, + time_us, payload: JetstreamPayload::Identity { identity: JetstreamIdentity { did: did_str.as_str().to_string(), @@ -556,246 +390,31 @@ fn jetstream_event_to_bytes( #[cfg(test)] mod tests { use super::*; - use std::collections::BTreeMap; - use std::num::NonZeroUsize; - use std::sync::Mutex; - use std::thread::JoinHandle; - use std::time::Duration; - - #[derive(Clone)] - struct Ev { - time: u64, - id: u64, - } - - impl ReplayPosition for Ev { - fn id(&self) -> u64 { - self.id - } - - fn key_after(&self) -> ReplayKey { - keys::jetstream_event_key(self.time, self.id + 1) - } - } - - /// the db and broadcast channel as ingest drives them: commit, then broadcast. - struct Log { - db: Mutex>, - tx: Mutex>>, - } - - impl Log { - fn new(capacity: usize) -> (Arc, broadcast::Receiver) { - let (tx, rx) = broadcast::channel(capacity); - let log = Self { - db: Mutex::new(BTreeMap::new()), - tx: Mutex::new(Some(tx)), - }; - (Arc::new(log), rx) - } - - fn commit(&self, time: u64, id: u64) { - let ev = Ev { time, id }; - self.db - .lock() - .unwrap() - .insert(keys::jetstream_event_key(time, id), ev.clone()); - self.broadcast(ev); - } - - fn broadcast(&self, ev: Ev) { - if let Some(tx) = &*self.tx.lock().unwrap() { - let _ = tx.send(ev); - } - } - - fn read(&self, start: &ReplayKey, chunk_size: usize) -> ReplayChunk { - let events: Vec<_> = self - .db - .lock() - .unwrap() - .range(*start..) - .take(chunk_size) - .map(|(_, ev)| ev.clone()) - .collect(); - ReplayChunk { - exhausted: events.len() < chunk_size, - events, - last_seen_seq: None, - } - } - - fn close(&self) { - self.tx.lock().unwrap().take(); - } - } - - type Out = mpsc::Receiver>; - - fn opts(chunk_size: usize) -> StreamOptions { - StreamOptions { - replay_chunk_size: NonZeroUsize::new(chunk_size), - replay_chunk_pause: Duration::ZERO, - pending_event_limit: 1, - send_timeout: Duration::from_secs(5), - } - } - - // run the stream on its own thread. `read` sees the log and a call counter so a - // test can commit events at exact points of the replay. - fn spawn( - log: &Arc, - event_rx: broadcast::Receiver, - cursor: Option, - out_capacity: usize, - chunk_size: usize, - mut read: impl FnMut(&Log, usize, &ReplayKey, usize) -> ReplayChunk + Send + 'static, - ) -> (JoinHandle<()>, Out) { - let (tx, out) = mpsc::channel(out_capacity); - let log = log.clone(); - let handle = std::thread::spawn(move || { - let mut calls = 0; - run_jetstream_stream( - &tx, - event_rx, - cursor.map(|time| keys::jetstream_event_key(time, 0)), - opts(chunk_size), - |start, chunk_size| { - calls += 1; - read(&log, calls - 1, start, chunk_size) - }, - |ev: Ev| Some(Bytes::copy_from_slice(&ev.id.to_be_bytes())), - ); - }); - (handle, out) - } - - fn next_id(out: &mut Out) -> u64 { - let deadline = Instant::now() + Duration::from_secs(5); - let frame = loop { - match out.try_recv() { - Ok(frame) => break frame.expect("stream error"), - Err(mpsc::error::TryRecvError::Empty) if Instant::now() < deadline => { - std::thread::sleep(Duration::from_millis(5)); - } - Err(err) => panic!("no frame: {err}"), - } - }; - u64::from_be_bytes(frame.as_ref().try_into().unwrap()) - } - - fn take_ids(out: &mut Out, n: usize) -> Vec { - (0..n).map(|_| next_id(out)).collect() - } - - fn finish(log: &Log, handle: JoinHandle<()>, mut out: Out) { - log.close(); - handle.join().unwrap(); - assert!(out.try_recv().is_err(), "unexpected extra frame"); - } #[test] - fn commits_during_replay_are_sent() { - let (log, event_rx) = Log::new(16); - (0..3).for_each(|id| log.commit(id + 1, id)); - let (handle, mut out) = spawn(&log, event_rx, Some(0), 16, 2, |log, call, start, n| { - let chunk = log.read(start, n); - if call == 0 { - // lands past the head the replay started with - log.commit(10, 3); - } - chunk - }); - - assert_eq!(take_ids(&mut out, 4), [0, 1, 2, 3]); - log.commit(11, 4); - assert_eq!(next_id(&mut out), 4); - finish(&log, handle, out); - } - - #[test] - fn events_in_last_scan_and_tail_are_sent_once() { - let (log, event_rx) = Log::new(16); - log.commit(1, 0); - let (handle, mut out) = spawn(&log, event_rx, Some(0), 16, 4, |log, call, start, n| { - if call == 0 { - // broadcast after the drain, and visible to the scan - log.commit(5, 1); - } - log.read(start, n) - }); - - assert_eq!(take_ids(&mut out, 2), [0, 1]); - log.commit(6, 2); - assert_eq!(next_id(&mut out), 2); - finish(&log, handle, out); - } - - #[test] - fn tail_sends_broadcasts_that_arrive_out_of_id_order() { - let (log, event_rx) = Log::new(16); - let (handle, mut out) = spawn(&log, event_rx, None, 16, 4, |_, _, _, _| { - unreachable!("no cursor, no replay") - }); - - // another shard committed the later id first - log.commit(5, 11); - log.commit(4, 10); - assert_eq!(take_ids(&mut out, 2), [11, 10]); - finish(&log, handle, out); - } - - #[test] - fn tail_lagging_before_a_live_send_replays_from_the_db() { - let (log, event_rx) = Log::new(2); - log.commit(1, 0); - let (handle, mut out) = spawn(&log, event_rx, Some(0), 16, 4, |log, call, start, n| { - let chunk = log.read(start, n); - if call == 0 { - // overflows the broadcast channel before the tail starts - (1..=5).for_each(|id| log.commit(id + 1, id)); - } - chunk - }); - - assert_eq!(take_ids(&mut out, 6), [0, 1, 2, 3, 4, 5]); - log.commit(10, 6); - assert_eq!(next_id(&mut out), 6); - finish(&log, handle, out); + fn positions_order_like_their_keys() { + let positions = [ + JetstreamPosition { + time_us: 1, + id: u64::MAX, + }, + JetstreamPosition { time_us: 2, id: 0 }, + JetstreamPosition { time_us: 2, id: 7 }, + JetstreamPosition { time_us: 3, id: 1 }, + ]; + for pair in positions.windows(2) { + assert!(pair[0] < pair[1]); + assert!(pair[0].key() < pair[1].key()); + assert!(pair[0].key_after() <= pair[1].key()); + assert!(pair[0].key() < pair[0].key_after()); + } } #[test] - fn tail_lagging_after_a_live_send_closes_instead_of_skipping() { - let (log, event_rx) = Log::new(2); - let (handle, mut out) = spawn(&log, event_rx, None, 1, 4, |_, _, _, _| { - unreachable!("no cursor, no replay") - }); - - log.commit(1, 0); - log.commit(2, 1); - std::thread::sleep(Duration::from_millis(100)); - // the tail is stuck sending 1 while 2..=4 overflow the broadcast channel - (2..=4).for_each(|id| log.commit(id + 1, id)); - std::thread::sleep(Duration::from_millis(100)); - assert_eq!(next_id(&mut out), 0); - - let deadline = Instant::now() + Duration::from_secs(5); - while !handle.is_finished() { - assert!(Instant::now() < deadline, "lagged tail kept running"); - std::thread::sleep(Duration::from_millis(10)); - } - handle.join().unwrap(); - - let rest: Vec<_> = std::iter::from_fn(|| out.try_recv().ok()).collect(); - let sent: Vec = rest - .iter() - .filter_map(|frame| frame.as_ref().ok()) - .map(|frame| u64::from_be_bytes(frame.as_ref().try_into().unwrap())) - .collect(); - assert_eq!(sent, [1]); - assert!(rest.iter().all(|frame| match frame { - Ok(_) => true, - Err(err) => err.reason.contains("lagged"), - })); + fn a_cursor_starts_before_its_first_event() { + let before = JetstreamPosition::before(5).unwrap(); + assert!(before.key() < JetstreamPosition { time_us: 5, id: 0 }.key()); + assert_eq!(before.key_after(), keys::jetstream_event_key(5, 0)); + assert!(JetstreamPosition::before(0).is_none()); } } diff --git a/src/control/stream/relay.rs b/src/control/stream/relay.rs index f5bcfed..450ac3e 100644 --- a/src/control/stream/relay.rs +++ b/src/control/stream/relay.rs @@ -6,7 +6,7 @@ use crate::db::keys; use crate::state::AppState; use crate::types::RelayBroadcast; -use super::{ReplayChunk, StreamOptions, run_ordered_stream, stream_seq_after}; +use super::{ReplayChunk, ReplayFailed, Start, StreamOptions, StreamSource, run_stream}; use crate::control::RelayStreamError; pub(crate) fn relay_stream_thread( @@ -16,87 +16,65 @@ pub(crate) fn relay_stream_thread( opts: StreamOptions, ) { let relay_rx = state.db.relay.broadcast_tx.subscribe(); - let ks = state.db.relay.events.clone(); let head = state.db.relay.seqs.head(); - let current_seq = match cursor { - Some(c) => Some(c.saturating_sub(1)), - None => Some(head.unwrap_or(0)), + let start = match cursor { + Some(cursor) => Start::replay(cursor.checked_sub(1), head), + None => Start::live(head), }; - let catch_up_target = cursor - .and(head) - .filter(|target| stream_seq_after(*target, current_seq)); - - run_ordered_stream( - tx, - relay_rx, - current_seq, - catch_up_target, - opts, - move |current_seq, target, chunk_size| { - read_relay_replay_chunk(&ks, current_seq, target, chunk_size) - }, - relay_broadcast_to_frame, - ); + run_stream(tx, relay_rx, start, opts, RelaySource { state }); } -fn read_relay_replay_chunk( - ks: &fjall::Keyspace, - current_seq: Option, - target: u64, - chunk_size: usize, -) -> ReplayChunk { - let start = current_seq.map(|seq| seq.saturating_add(1)).unwrap_or(0); - if start > target { - return ReplayChunk { - events: Vec::new(), - last_seen_seq: current_seq, - exhausted: true, - }; - } - - let mut events = Vec::with_capacity(chunk_size); - let mut last_seen_seq = current_seq; - 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::relay_event_key(start)..=keys::relay_event_key(target)); +struct RelaySource { + state: Arc, +} - while events.len() < chunk_size && scanned < max_scanned { - let Some(item) = iter.next() else { - exhausted = true; - break; - }; - scanned += 1; +impl StreamSource for RelaySource { + type Broadcast = RelayBroadcast; + type Output = bytes::Bytes; - let (k, v) = match item.into_inner() { - Ok(kv) => kv, - Err(e) => { - error!(err = %e, "relay stream: failed to read relay_events"); - exhausted = true; - break; - } - }; - let seq = match k.as_ref().try_into().map(u64::from_be_bytes) { - Ok(seq) => seq, - Err(_) => { - error!("relay stream: failed to parse relay event seq"); - continue; - } - }; - last_seen_seq = Some(seq); - events.push(bytes::Bytes::copy_from_slice(&v)); + fn read( + &mut self, + after: Option, + through: Option, + limit: usize, + ) -> Result, ReplayFailed> { + read_relay_replay_chunk(&self.state.db.relay.events, after, through, limit) } - ReplayChunk { - events, - last_seen_seq, - exhausted, + fn render(&mut self, event: RelayBroadcast) -> Option { + match event { + RelayBroadcast::Persisted(_) => None, + RelayBroadcast::Ephemeral(_, frame) => Some(frame), + } } } -fn relay_broadcast_to_frame(event: RelayBroadcast) -> Option { - match event { - RelayBroadcast::Persisted(_) => None, - RelayBroadcast::Ephemeral(_, frame) => Some(frame), +fn read_relay_replay_chunk( + ks: &fjall::Keyspace, + after: Option, + through: Option, + limit: usize, +) -> Result, ReplayFailed> { + let start = after.map_or(0, |seq| seq.saturating_add(1)); + let end = through.unwrap_or(u64::MAX); + if start > end { + return Ok(ReplayChunk { + events: Vec::new(), + last_seen: after, + exhausted: true, + }); } + + let rows = ks.range(keys::relay_event_key(start)..=keys::relay_event_key(end)); + ReplayChunk::read( + rows, + after, + limit, + |key| { + let seq = key.try_into().map(u64::from_be_bytes); + seq.inspect_err(|_| error!("relay stream: failed to parse relay event seq")) + .ok() + }, + |_, frame| Some(bytes::Bytes::copy_from_slice(frame)), + ) } diff --git a/src/control/stream/types.rs b/src/control/stream/types.rs index d477736..0f2703b 100644 --- a/src/control/stream/types.rs +++ b/src/control/stream/types.rs @@ -1,11 +1,14 @@ -use std::collections::VecDeque; use std::fmt; use std::num::NonZeroUsize; use std::time::Duration; +use tracing::error; + use crate::config::Config; -pub(crate) const STREAM_SEND_RETRY_PAUSE: Duration = Duration::from_millis(10); +use super::engine::StreamBroadcast; + +pub(super) const STREAM_SEND_RETRY_PAUSE: Duration = Duration::from_millis(10); pub(crate) const STREAM_CHANNEL_CAPACITY: usize = 500; #[derive(Debug, Clone, Copy)] @@ -33,18 +36,6 @@ pub(crate) struct StreamTooSlow { } impl StreamTooSlow { - pub(crate) fn pending_limit(limit: usize) -> Self { - Self { - reason: format!("pending stream event buffer exceeded {limit} events"), - } - } - - pub(crate) fn lagged(skipped: u64) -> Self { - Self { - reason: format!("subscriber lagged past {skipped} broadcast events"), - } - } - pub(crate) fn send_timeout(timeout: Duration) -> Self { Self { reason: format!( @@ -82,112 +73,107 @@ impl From for super::super::JetstreamStreamError { } } -pub(crate) trait StreamBroadcast { - fn sequence(&self) -> u64; - fn is_persisted_marker(&self) -> bool; -} - #[cfg(feature = "indexer_stream")] impl StreamBroadcast for crate::types::BroadcastEvent { - fn sequence(&self) -> u64 { + type Position = u64; + + fn position(&self) -> u64 { match self { - crate::types::BroadcastEvent::Persisted(id) => *id, - crate::types::BroadcastEvent::LiveRecord(evt) => evt.id, - crate::types::BroadcastEvent::Ephemeral(evt) => evt.id, + Self::Persisted(id) => *id, + Self::LiveRecord(evt) => evt.id, + Self::Ephemeral(evt) => evt.id, } } - fn is_persisted_marker(&self) -> bool { - matches!(self, crate::types::BroadcastEvent::Persisted(_)) + fn is_marker(&self) -> bool { + matches!(self, Self::Persisted(_)) } } #[cfg(feature = "relay")] impl StreamBroadcast for crate::types::RelayBroadcast { - fn sequence(&self) -> u64 { + type Position = u64; + + fn position(&self) -> u64 { match self { Self::Persisted(seq) | Self::Ephemeral(seq, _) => *seq, } } - fn is_persisted_marker(&self) -> bool { + fn is_marker(&self) -> bool { matches!(self, Self::Persisted(_)) } } #[cfg(feature = "jetstream")] impl StreamBroadcast for crate::types::JetstreamBroadcast { - fn sequence(&self) -> u64 { - self.id - } + type Position = crate::types::JetstreamPosition; - fn is_persisted_marker(&self) -> bool { - false - } -} - -pub(crate) struct PendingLiveEvents { - pub(crate) queue: VecDeque, - pub(crate) persisted_head: Option, - pub(crate) limit: usize, -} - -impl PendingLiveEvents -where - T: StreamBroadcast, -{ - pub(crate) fn new(limit: usize) -> Self { - Self { - queue: VecDeque::new(), - persisted_head: None, - limit, - } - } - - pub(crate) fn push(&mut self, event: T) -> Result<(), StreamTooSlow> { - if event.is_persisted_marker() { - let seq = event.sequence(); - self.persisted_head = Some(self.persisted_head.unwrap_or(0).max(seq)); - return Ok(()); - } - - if self.queue.len() >= self.limit { - return Err(StreamTooSlow::pending_limit(self.limit)); + fn position(&self) -> Self::Position { + match self { + Self::Historical(position) => *position, + Self::Live(live) => live.position, } - self.queue.push_back(event); - Ok(()) - } - - pub(crate) fn next_sequence(&self) -> Option { - self.queue.front().map(StreamBroadcast::sequence) } - pub(crate) fn pop_front(&mut self) -> Option { - self.queue.pop_front() - } - - pub(crate) fn push_front(&mut self, event: T) { - self.queue.push_front(event); - } - - pub(crate) fn take_persisted_after(&mut self, current_id: Option) -> Option { - let head = self.persisted_head.take()?; - stream_seq_after(head, current_id).then_some(head) + fn is_marker(&self) -> bool { + matches!(self, Self::Historical(_)) } } -pub(crate) struct ReplayChunk { +/// a db read failed partway. the rows past it can't be skipped, so the +/// stream ends and its subscriber resumes from the last event it got. +#[derive(Debug)] +pub(crate) struct ReplayFailed; + +/// outputs for the rows one db read scanned. +pub(crate) struct ReplayChunk { pub(crate) events: Vec, - pub(crate) last_seen_seq: Option, + /// the last row scanned, sent or filtered out + pub(crate) last_seen: Option

, + /// no rows are left in the range read pub(crate) exhausted: bool, } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub(crate) enum SendOutcome { - Sent, - ReceiverDropped, -} - -fn stream_seq_after(id: u64, current_id: Option) -> bool { - current_id.is_none_or(|current| id > current) +impl ReplayChunk { + /// read `rows` after `after` into up to `limit` outputs. scans at most + /// four times as many rows, so a range the subscriber mostly filters out + /// still returns promptly. rows whose key doesn't parse are skipped. + pub(crate) fn read( + mut rows: fjall::Iter, + after: Option

, + limit: usize, + position: impl Fn(&[u8]) -> Option

, + mut output: impl FnMut(P, &[u8]) -> Option, + ) -> Result + where + P: Copy, + { + let mut chunk = Self { + events: Vec::with_capacity(limit), + last_seen: after, + exhausted: false, + }; + let max_scanned = limit.saturating_mul(4).max(limit); + let mut scanned = 0usize; + + while chunk.events.len() < limit && scanned < max_scanned { + let Some(row) = rows.next() else { + chunk.exhausted = true; + break; + }; + scanned += 1; + + let (key, value) = row.into_inner().map_err(|e| { + error!(err = %e, "failed to read stream rows from db, closing the stream"); + ReplayFailed + })?; + let Some(position) = position(&key) else { + continue; + }; + chunk.last_seen = Some(position); + chunk.events.extend(output(position, &value)); + } + Ok(chunk) + } } diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index 2e598da..d63496f 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -237,18 +237,21 @@ impl JetstreamDb { }) } + /// the last committed event's position. + pub(crate) fn head(&self) -> Result> { + let Some(guard) = self.events.iter().next_back() else { + return Ok(None); + }; + let key = guard.key().into_diagnostic()?; + let (time_us, id) = super::keys::parse_jetstream_event_key(&key)?; + Ok(Some(crate::types::JetstreamPosition { time_us, id })) + } + /// resume jetstream ids and time watermark after the last stored event. pub(super) fn init(&self) -> Result<()> { - let mut last_id = 0; - let mut last_time_us = 0; - if let Some(guard) = self.events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - let (time_us, id) = super::keys::parse_jetstream_event_key(&k)?; - last_id = id; - last_time_us = time_us; - } - self.ids.resume_at(last_id + 1); - self.clock.resume_after(last_time_us); + let head = self.head()?; + self.ids.resume_at(head.map_or(1, |head| head.id + 1)); + self.clock.resume_after(head.map_or(0, |head| head.time_us)); Ok(()) } } diff --git a/src/db/mod.rs b/src/db/mod.rs index 3620555..d65bcd3 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -52,6 +52,8 @@ pub use schema::Ks; #[cfg(feature = "indexer")] use lifecycle_counts::LifecycleCountBatch; +#[cfg(all(test, feature = "indexer_stream"))] +pub(crate) use txn::RecordTxn; pub(crate) use txn::{Txn, TxnCommitTimings}; use tracing::error; diff --git a/src/db/outbox/jetstream.rs b/src/db/outbox/jetstream.rs index 9fc4d66..579c73e 100644 --- a/src/db/outbox/jetstream.rs +++ b/src/db/outbox/jetstream.rs @@ -4,11 +4,10 @@ use fjall::OwnedWriteBatch; use miette::Result; use tracing::error; -use crate::db::keys; use crate::db::keyspaces::JetstreamDb; -use crate::db::sequencer::Assign; +use crate::db::sequencer::{Assign, publish_in_order}; use crate::jetstream::JetstreamEphemeral; -use crate::types::{JetstreamBroadcast, StoredJetstreamEvent}; +use crate::types::{JetstreamBroadcast, JetstreamLive, JetstreamPosition, StoredJetstreamEvent}; #[cfg(feature = "relay")] use super::relay::{RelayFrameRef as UpstreamRef, StagedRelay as StagedUpstream}; @@ -47,21 +46,22 @@ impl JetstreamOutbox { .into_iter() .filter_map(|(event, json)| { let event = event.resolve(|upstream_ref| upstream.id(upstream_ref)); - let time_us = db.clock.tick(assign); - let id = db.ids.take(assign); + let position = JetstreamPosition { + time_us: db.clock.tick(assign), + id: db.ids.take(assign), + }; // like a relay frame that fails to encode, only this event is lost let row = rmp_serde::to_vec(&event) .inspect_err(|e| { - error!(err = %e, time_us, id, "dropping a jetstream event that failed to serialize") + error!(err = %e, ?position, "dropping a jetstream event that failed to serialize") }) .ok()?; - batch.insert(&db.events, keys::jetstream_event_key(time_us, id), row); - let time_us = time_us as i64; - Some(JetstreamBroadcast { - id, - time_us, - ephemeral: json.and_then(|json| json.to_json(time_us)), + batch.insert(&db.events, position.key(), row); + let json = json.and_then(|json| json.to_json(position.time_us as i64)); + Some(JetstreamLive { + position, event, + json, }) }) .collect(); @@ -70,17 +70,18 @@ impl JetstreamOutbox { } pub(super) struct StagedJetstream { - events: Vec, + events: Vec, } impl StagedJetstream { - /// broadcast the live events. backfilled commits aren't broadcast, a - /// replay reads them. + /// broadcast the live events. backfilled commits are only announced by a + /// marker, for the subscribers that want them to read. pub(super) fn publish(self, db: &JetstreamDb) { - for event in self.events { - if event.event.is_live() { - let _ = db.tx.send(event); - } - } + let events = self.events.into_iter().map(|live| { + let position = live.position; + let broadcast = live.event.is_live().then(|| JetstreamBroadcast::Live(live)); + (position, broadcast) + }); + publish_in_order(&db.tx, events, JetstreamBroadcast::Historical); } } diff --git a/src/types/event/jetstream.rs b/src/types/event/jetstream.rs index bc4a572..4ff4109 100644 --- a/src/types/event/jetstream.rs +++ b/src/types/event/jetstream.rs @@ -1,4 +1,3 @@ -use jacquard_common::IntoStatic; use jacquard_common::{CowStr, types::string::Handle}; use serde::{Deserialize, Serialize}; @@ -69,14 +68,52 @@ pub(crate) enum StoredJetstreamEvent<'i, Id = u64> { }, } +/// a jetstream event's place in commit order, which is also its key in the +/// events keyspace. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] +pub(crate) struct JetstreamPosition { + pub(crate) time_us: u64, + pub(crate) id: u64, +} + +impl JetstreamPosition { + pub(crate) fn key(self) -> [u8; 16] { + crate::db::keys::jetstream_event_key(self.time_us, self.id) + } + + /// the first key after this position. + pub(crate) fn key_after(self) -> [u8; 16] { + match self.id.checked_add(1) { + Some(id) => crate::db::keys::jetstream_event_key(self.time_us, id), + None => crate::db::keys::jetstream_event_key(self.time_us.saturating_add(1), 0), + } + } + + /// the last position before any event at or after `time_us`. + pub(crate) fn before(time_us: u64) -> Option { + time_us.checked_sub(1).map(|time_us| Self { + time_us, + id: u64::MAX, + }) + } +} + +#[derive(Clone, Debug)] +pub(crate) enum JetstreamBroadcast { + /// every event through this position is committed, and the ones since + /// the previous broadcast are backfilled commits, which are never + /// broadcast themselves. + Historical(JetstreamPosition), + Live(JetstreamLive), +} + #[derive(Clone, Debug)] -pub(crate) struct JetstreamBroadcast { - pub id: u64, - pub time_us: i64, - pub event: StoredJetstreamEvent<'static>, - /// pre-serialized json bytes for live tailing. - /// replay events read from the database instead. - pub ephemeral: Option, +pub(crate) struct JetstreamLive { + pub(crate) position: JetstreamPosition, + pub(crate) event: StoredJetstreamEvent<'static>, + /// pre-serialized json, when a subscriber was listening as the event was + /// staged. otherwise it is rendered from the database. + pub(crate) json: Option, } impl<'i, Id> StoredJetstreamEvent<'i, Id> { @@ -175,69 +212,6 @@ impl<'i, Id> StoredJetstreamEvent<'i, Id> { }, } } - - pub(crate) fn into_static(self) -> StoredJetstreamEvent<'static, Id> { - match self { - #[cfg(feature = "indexer_stream")] - Self::Commit { - did, - collection, - event_id, - live, - } => StoredJetstreamEvent::Commit { - did: did.into_static(), - collection: collection.into_static(), - event_id, - live, - }, - #[cfg(feature = "relay")] - Self::RelayCommit { - did, - collection, - relay_seq, - op_index, - } => StoredJetstreamEvent::RelayCommit { - did: did.into_static(), - collection: collection.into_static(), - relay_seq, - op_index, - }, - #[cfg(feature = "relay")] - Self::RelayAccount { did, relay_seq } => StoredJetstreamEvent::RelayAccount { - did: did.into_static(), - relay_seq, - }, - #[cfg(feature = "relay")] - Self::RelayIdentity { did, relay_seq } => StoredJetstreamEvent::RelayIdentity { - did: did.into_static(), - relay_seq, - }, - Self::Account { - did, - active, - status, - seq, - time, - } => StoredJetstreamEvent::Account { - did: did.into_static(), - active, - status: status.map(IntoStatic::into_static), - seq, - time, - }, - Self::Identity { - did, - handle, - seq, - time, - } => StoredJetstreamEvent::Identity { - did: did.into_static(), - handle: handle.map(IntoStatic::into_static), - seq, - time, - }, - } - } } #[cfg(test)] -- 2.51.2