diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index c3083fb..62746d7 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -204,18 +204,7 @@ impl StreamDb { /// resume event ids after the last stored event. pub(super) fn init(&self) -> Result<()> { - let mut last_id = 0; - if let Some(guard) = self.events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - last_id = u64::from_be_bytes( - k.as_ref() - .try_into() - .into_diagnostic() - .wrap_err("expected to be id (8 bytes)")?, - ); - } - self.ids.resume_at(last_id + 1); - Ok(()) + self.ids.resume_after(self.event_head()?) } } @@ -256,9 +245,8 @@ impl JetstreamDb { /// resume jetstream ids and time watermark after the last stored event. pub(super) fn init(&self) -> Result<()> { 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(()) + self.ids.resume_after(head.map(|head| head.id))?; + self.clock.resume_after(head.map(|head| head.time_us)) } } @@ -286,17 +274,19 @@ impl RelayDb { /// resume relay sequence numbers after the last stored frame. pub(super) fn init(&self) -> Result<()> { - let mut last_relay_seq = 0u64; - if let Some(guard) = self.events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - last_relay_seq = u64::from_be_bytes( - k.as_ref() + let last_relay_seq = self + .events + .iter() + .next_back() + .map(|guard| { + let key = guard.key().into_diagnostic()?; + key.as_ref() .try_into() + .map(u64::from_be_bytes) .into_diagnostic() - .wrap_err("relay_events: invalid key length")?, - ); - } - self.seqs.resume_at(last_relay_seq + 1); - Ok(()) + .wrap_err("relay_events: invalid key length") + }) + .transpose()?; + self.seqs.resume_after(last_relay_seq) } } diff --git a/src/db/outbox.rs b/src/db/outbox.rs index c392252..68aeda1 100644 --- a/src/db/outbox.rs +++ b/src/db/outbox.rs @@ -110,7 +110,7 @@ impl Outbox { batch: &mut OwnedWriteBatch, ) -> Result { #[cfg(feature = "indexer_stream")] - let stream = self.stream.stage(&db.stream, assign, batch); + let stream = self.stream.stage(&db.stream, assign, batch)?; #[cfg(feature = "relay")] let relay = self.relay.stage(&db.relay, assign, batch)?; // jetstream rows refer to the upstream rows staged just above diff --git a/src/db/outbox/jetstream.rs b/src/db/outbox/jetstream.rs index adf6f8c..839ff9c 100644 --- a/src/db/outbox/jetstream.rs +++ b/src/db/outbox/jetstream.rs @@ -44,30 +44,31 @@ impl JetstreamOutbox { let events = self .events .into_iter() - .filter_map(|(event, json)| { + .map(|(event, json)| { let event = event.resolve(|upstream_ref| upstream.id(upstream_ref)); let position = JetstreamPosition { - time_us: db.clock.tick(assign), - id: db.ids.take(assign), + 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, ?position, "dropping a jetstream event that failed to serialize") - }) - .ok()?; + let Ok(row) = rmp_serde::to_vec(&event).inspect_err(|e| { + error!(err = %e, ?position, "dropping a jetstream event that failed to serialize") + }) else { + return Ok(None); + }; batch.insert(&db.events, position.key(), row); let (json, redaction_key) = json.map_or((None, None), |json| { (json.to_json(position.time_us as i64), json.redaction_key) }); - Some(JetstreamLive { + Ok(Some(JetstreamLive { position, event, json, redaction_key, - }) + })) }) - .collect(); + .filter_map(Result::transpose) + .collect::>()?; Ok(StagedJetstream { events }) } } diff --git a/src/db/outbox/relay.rs b/src/db/outbox/relay.rs index a629aab..0e4f4c3 100644 --- a/src/db/outbox/relay.rs +++ b/src/db/outbox/relay.rs @@ -79,22 +79,22 @@ impl RelayOutbox { .frames .into_iter() .map(|frame| { - let seq = db.seqs.take(assign); + let seq = db.seqs.take(assign)?; // encoding is deterministic, so failing the commit would lose // the rest of its writes for good. only this frame is lost. let encoded = frame.encode(seq).inspect_err( |e| error!(err = %e, seq, "dropping a relay frame that failed to encode"), ); let Ok((bytes, broadcast)) = encoded else { - return (seq, None); + return Ok((seq, None)); }; batch.insert(&db.events, keys::relay_event_key(seq), bytes.as_ref()); - ( + Ok(( seq, broadcast.then(|| RelayBroadcast::Ephemeral(seq, bytes)), - ) + )) }) - .collect(); + .collect::>()?; Ok(StagedRelay { frames }) } } diff --git a/src/db/outbox/stream.rs b/src/db/outbox/stream.rs index 0beed77..561b6c8 100644 --- a/src/db/outbox/stream.rs +++ b/src/db/outbox/stream.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use bytes::Bytes; use fjall::OwnedWriteBatch; +use miette::Result; use crate::db::keys; use crate::db::keyspaces::StreamDb; @@ -70,16 +71,16 @@ impl StreamOutbox { db: &'a StreamDb, assign: &mut Assign<'a>, batch: &mut OwnedWriteBatch, - ) -> StagedStream { + ) -> Result { let events = self .events .into_iter() .map(|event| { - let id = db.ids.take(assign); - (id, event.stage(id, db, batch)) + let id = db.ids.take(assign)?; + Ok((id, event.stage(id, db, batch))) }) - .collect(); - StagedStream { events } + .collect::>()?; + Ok(StagedStream { events }) } } diff --git a/src/db/sequencer.rs b/src/db/sequencer.rs index 2babf92..5e9f884 100644 --- a/src/db/sequencer.rs +++ b/src/db/sequencer.rs @@ -17,6 +17,10 @@ use std::sync::{Mutex, PoisonError}; use fjall::OwnedWriteBatch; use miette::{IntoDiagnostic, Result}; +/// the last position a counter hands out. every stream position fits in an +/// i64, the width relay seqs and jetstream `time_us` have on the wire. +const LAST_POSITION: u64 = i64::MAX as u64; + #[derive(Default)] pub(crate) struct Sequencer { lock: Mutex<()>, @@ -52,7 +56,7 @@ pub(crate) struct Assign<'a> { impl<'a> Assign<'a> { /// take `counter`'s next value. the first take in a commit begins at /// `start(committed next)` and later ones count up from there. - fn take(&mut self, counter: &'a AtomicU64, start: impl FnOnce(u64) -> u64) -> u64 { + fn take(&mut self, counter: &'a AtomicU64, start: impl FnOnce(u64) -> u64) -> Result { let slot = match self .pending .iter() @@ -66,8 +70,9 @@ impl<'a> Assign<'a> { } }; let taken = self.pending[slot].1; + miette::ensure!(taken <= LAST_POSITION, "stream positions are exhausted"); self.pending[slot].1 = taken + 1; - taken + Ok(taken) } fn advance(self) { @@ -95,12 +100,19 @@ impl Sequence { self.next.load(Ordering::Acquire) } - /// continue numbering at `next`, when opening the database. + /// continue numbering after `last`, the last stored position, or at 1 + /// for none, when opening the database. + pub(in crate::db) fn resume_after(&self, last: Option) -> Result<()> { + self.next.store(next_after(last)?.max(1), Ordering::Release); + Ok(()) + } + + #[cfg(test)] pub(in crate::db) fn resume_at(&self, next: u64) { self.next.store(next, Ordering::Release); } - pub(crate) fn take<'a>(&'a self, assign: &mut Assign<'a>) -> u64 { + pub(crate) fn take<'a>(&'a self, assign: &mut Assign<'a>) -> Result { assign.take(&self.next, |next| next) } } @@ -118,9 +130,11 @@ pub(crate) struct Clock { #[cfg(feature = "jetstream")] impl Clock { - /// continue after `last`, when opening the database. - pub(in crate::db) fn resume_after(&self, last: u64) { - self.next.store(last + 1, Ordering::Release); + /// continue after `last`, the last stored timestamp, when opening the + /// database. + pub(in crate::db) fn resume_after(&self, last: Option) -> Result<()> { + self.next.store(next_after(last)?, Ordering::Release); + Ok(()) } /// start the next commit `by` ahead of the wall clock, as if the wall @@ -128,10 +142,11 @@ impl Clock { #[cfg(all(test, feature = "indexer_stream"))] pub(crate) fn run_ahead(&self, by: std::time::Duration) { let now = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap(); - self.resume_after(now + u64::try_from(by.as_micros()).unwrap()); + self.resume_after(Some(now + u64::try_from(by.as_micros()).unwrap())) + .unwrap(); } - pub(crate) fn tick<'a>(&'a self, assign: &mut Assign<'a>) -> u64 { + pub(crate) fn tick<'a>(&'a self, assign: &mut Assign<'a>) -> Result { assign.take(&self.next, |next| { let now = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(0); next.max(now) @@ -139,6 +154,19 @@ impl Clock { } } +/// the position after `last`, refusing a stored position that leaves no room +/// for another, which only a corrupted or hand-edited database has. +fn next_after(last: Option) -> Result { + let Some(last) = last else { + return Ok(0); + }; + miette::ensure!( + last < LAST_POSITION, + "the database's last stream position {last} leaves no room for more" + ); + Ok(last + 1) +} + /// send `events` in position order. a run of rows that aren't broadcast is /// announced by one `marker` for its last position, sent before the next /// broadcast, so a subscriber reads them from the db before moving past them. @@ -188,7 +216,7 @@ mod tests { let failed = sequencer.commit( db.inner.batch(), |assign, _| { - assert_eq!((seq.take(assign), seq.take(assign)), (5, 6)); + assert_eq!((seq.take(assign)?, seq.take(assign)?), (5, 6)); Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, @@ -200,7 +228,7 @@ mod tests { .commit( db.inner.batch(), |assign, _| { - let taken = [seq.take(assign), seq.take(assign)]; + let taken = [seq.take(assign)?, seq.take(assign)?]; // readers never see positions that haven't committed assert_eq!(seq.head(), Some(4)); Ok(taken) @@ -212,6 +240,36 @@ mod tests { assert_eq!(seq.head(), Some(6)); } + #[test] + fn positions_stop_at_the_last_one() { + let (_tmp, db) = db(); + let sequencer = Sequencer::default(); + let seq = Sequence::default(); + seq.resume_at(LAST_POSITION); + + let failed = sequencer.commit( + db.inner.batch(), + |assign, _| { + assert_eq!(seq.take(assign)?, LAST_POSITION); + seq.take(assign) + }, + |_| {}, + ); + assert!(failed.is_err()); + assert_eq!(seq.next(), LAST_POSITION); + } + + #[test] + fn a_stored_position_needs_room_after_it() { + let seq = Sequence::default(); + seq.resume_after(None).unwrap(); + assert_eq!(seq.next(), 1); + seq.resume_after(Some(LAST_POSITION - 1)).unwrap(); + assert_eq!(seq.next(), LAST_POSITION); + assert!(seq.resume_after(Some(LAST_POSITION)).is_err()); + assert!(seq.resume_after(Some(u64::MAX)).is_err()); + } + #[cfg(feature = "jetstream")] #[test] fn clock_increases_even_when_ahead_of_the_wall_clock() { @@ -219,13 +277,13 @@ mod tests { let sequencer = Sequencer::default(); let clock = Clock::default(); let ahead = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap() + 60_000_000; - clock.resume_after(ahead); + clock.resume_after(Some(ahead)).unwrap(); // a failed commit takes nothing let failed = sequencer.commit( db.inner.batch(), |assign, _| { - clock.tick(assign); + clock.tick(assign)?; Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, @@ -235,7 +293,7 @@ mod tests { let ticks = sequencer .commit( db.inner.batch(), - |assign, _| Ok([clock.tick(assign), clock.tick(assign)]), + |assign, _| Ok([clock.tick(assign)?, clock.tick(assign)?]), |ticks| ticks, ) .unwrap(); @@ -255,10 +313,10 @@ mod tests { .commit( db.inner.batch(), |assign, _| { - let first = clock.tick(assign); + let first = clock.tick(assign)?; // slow staging doesn't spread the commit's timestamps out std::thread::sleep(std::time::Duration::from_millis(5)); - Ok([first, clock.tick(assign), clock.tick(assign)]) + Ok([first, clock.tick(assign)?, clock.tick(assign)?]) }, |ticks| ticks, ) @@ -268,7 +326,7 @@ mod tests { // the next commit starts at the wall time, which has moved on let next = sequencer - .commit(db.inner.batch(), |assign, _| Ok(clock.tick(assign)), |t| t) + .commit(db.inner.batch(), |assign, _| clock.tick(assign), |t| t) .unwrap(); assert!(next >= ticks[0] + 5_000); }