//! stream positions, assigned while committing. //! //! every event stream position (`/stream` ids, `/subscribe` keys, relay //! seqs) is handed out inside [`Sequencer::commit`], which holds one lock across assigning the //! positions, committing the batch that writes their rows, and broadcasting //! the events. so by construction: //! //! - position order is commit order. every position at or below a sequence's //! head is committed, so a missing row below the head is a permanent hole, //! never a write still in flight. //! - a broadcast is only sent once its row is readable. //! - broadcasts arrive in position order, across every writer. //! - positions never repeat, even across restarts: a commit that goes past //! a counter's reservation reserves the next block in the same batch, and //! reopening resumes after the reservation, so positions that leave no row //! or whose rows get pruned are never handed out again. a crash skips the //! rest of the block, which is just more holes. use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Mutex, PoisonError}; use fjall::OwnedWriteBatch; use miette::{IntoDiagnostic, Result, WrapErr}; use super::schema::{Cursors, Ks}; /// 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; /// the one place stream positions are handed out. pub(crate) struct Sequencer { lock: Mutex<()>, /// where each counter's next position is written down marks: Ks, } impl Sequencer { pub(crate) fn new(marks: Ks) -> Self { Self { lock: Mutex::new(()), marks, } } /// assign positions and stage their rows, commit, then publish, all under /// one lock. positions only advance if the batch commits. pub(crate) fn commit<'a, T, R>( &self, mut batch: OwnedWriteBatch, stage: impl FnOnce(&mut Assign<'a>, &mut OwnedWriteBatch) -> Result, publish: impl FnOnce(T) -> R, ) -> Result { let _guard = self.lock.lock().unwrap_or_else(PoisonError::into_inner); let mut assign = Assign { pending: Vec::new(), reserving: Vec::new(), }; let staged = stage(&mut assign, &mut batch)?; assign.mark(&self.marks, &mut batch); batch.commit().into_diagnostic()?; assign.advance(); Ok(publish(staged)) } } /// positions taken by one sequenced commit. only [`Sequencer::commit`] can /// make one, so positions can't be taken anywhere else. pub(crate) struct Assign<'a> { /// the value each touched counter moves to once the batch commits pending: Vec<(&'a Counter, u64)>, /// the reservation each counter moves to once the batch commits, for the /// ones this commit took past theirs reserving: Vec<(&'a Counter, u64)>, } 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 Counter, start: impl FnOnce(u64) -> u64) -> Result { let slot = match self .pending .iter() .position(|(pending, _)| std::ptr::eq(*pending, counter)) { Some(slot) => slot, None => { let next = start(counter.next.load(Ordering::Acquire)); self.pending.push((counter, next)); self.pending.len() - 1 } }; let taken = self.pending[slot].1; miette::ensure!(taken <= LAST_POSITION, "stream positions are exhausted"); self.pending[slot].1 = taken + 1; Ok(taken) } /// reserve the next block for each counter this commit took past its /// reservation, in the same batch. the rest write nothing, because a mark /// per commit fills the cursors memtable, and every flush of it holds up /// all writers on fsyncs fn mark(&mut self, marks: &Ks, batch: &mut OwnedWriteBatch) { for &(counter, next) in &self.pending { if next > counter.reserved.load(Ordering::Acquire) { let reserved = next.max(next.saturating_add(counter.block).min(LAST_POSITION)); batch.insert(marks, counter.mark, reserved.to_be_bytes()); self.reserving.push((counter, reserved)); } } } fn advance(self) { for (counter, next) in self.pending { counter.next.store(next, Ordering::Release); } for (counter, reserved) in self.reserving { counter.reserved.store(reserved, Ordering::Release); } } } /// a position counter. `mark` holds a reservation: no commit has taken a /// position at or past it, so reopening the database resumes there, past /// positions that left no row or whose rows were pruned. struct Counter { /// the first position not committed yet next: AtomicU64, /// the committed reservation, only moved inside a sequenced commit reserved: AtomicU64, /// how far past `next` a commit reserves when it runs out block: u64, mark: &'static [u8], } impl Counter { const fn new(mark: &'static [u8], block: u64) -> Self { Self { next: AtomicU64::new(0), reserved: AtomicU64::new(0), block, mark, } } /// continue at `next`, or at the reservation if a commit got further. fn resume_at(&self, marks: &Ks, next: u64) -> Result<()> { let marked = marks .get(self.mark) .into_diagnostic()? .map(|mark| { <[u8; 8]>::try_from(mark.as_ref()) .map(u64::from_be_bytes) .into_diagnostic() .wrap_err("a stream position mark must be 8 bytes") }) .transpose()?; let reserved = marked.unwrap_or(0); let next = next.max(reserved); // only a corrupted or hand-edited database gets here miette::ensure!( next <= LAST_POSITION, "the database's next stream position {next} leaves no room for more" ); self.next.store(next, Ordering::Release); self.reserved.store(reserved, Ordering::Release); Ok(()) } } /// a counter whose values are only taken inside a sequenced commit. pub(crate) struct Sequence(Counter); /// positions a sequence reserves at a time. big enough that its mark is /// rewritten once every thousand or so commits, so the cursors memtable /// practically never fills up and flushes, and still nothing next to the /// position space a crash can skip. pub(in crate::db) const SEQUENCE_BLOCK: u64 = 1024; impl Sequence { pub(crate) const fn new(mark: &'static [u8]) -> Self { Self(Counter::new(mark, SEQUENCE_BLOCK)) } /// the last committed position, if any. pub(crate) fn head(&self) -> Option { self.next().checked_sub(1) } /// the position the next commit will take first. pub(crate) fn next(&self) -> u64 { self.0.next.load(Ordering::Acquire) } /// 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, marks: &Ks, last: Option) -> Result<()> { self.0.resume_at(marks, next_after(last).max(1)) } #[cfg(test)] pub(in crate::db) fn resume_at(&self, next: u64) { self.0.next.store(next, Ordering::Release); } pub(crate) fn take<'a>(&'a self, assign: &mut Assign<'a>) -> Result { assign.take(&self.0, |next| next) } } /// microsecond timestamps that strictly increase in commit order, even when /// the wall clock stalls or steps back. a commit's timestamps are consecutive, /// starting at the wall time or just after the previous commit's last, /// whichever is later, so the wall clock is read once per commit. #[cfg(feature = "jetstream")] pub(crate) struct Clock(Counter); /// how far ahead a clock reserves, in microseconds. a commit only reserves /// when the wall clock passes the reservation, so this is about one mark a /// second, and after a crash timestamps run at most this far ahead of the /// wall clock until it catches up. #[cfg(feature = "jetstream")] const CLOCK_BLOCK_US: u64 = 1_000_000; #[cfg(feature = "jetstream")] impl Clock { pub(crate) const fn new(mark: &'static [u8]) -> Self { Self(Counter::new(mark, CLOCK_BLOCK_US)) } /// continue after `last`, the last stored timestamp, when opening the /// database. pub(in crate::db) fn resume_after(&self, marks: &Ks, last: Option) -> Result<()> { self.0.resume_at(marks, next_after(last)) } /// start the next commit `by` ahead of the wall clock, as if the wall /// clock had stepped back after it. #[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(); let ahead = now + u64::try_from(by.as_micros()).unwrap(); self.0.next.store(ahead, Ordering::Release); } pub(crate) fn tick<'a>(&'a self, assign: &mut Assign<'a>) -> Result { assign.take(&self.0, |next| { let now = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(0); next.max(now) }) } } /// the position after `last`, saturating so a stored `u64::MAX` fails /// resuming instead of wrapping. fn next_after(last: Option) -> u64 { last.map_or(0, |last| last.saturating_add(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. pub(crate) fn publish_in_order( tx: &tokio::sync::broadcast::Sender, events: impl IntoIterator)>, marker: impl Fn(P) -> B, ) { let mut unannounced = None; for (position, event) in events { let Some(event) = event else { unannounced = Some(position); continue; }; if let Some(position) = unannounced.take() { let _ = tx.send(marker(position)); } let _ = tx.send(event); } if let Some(position) = unannounced { let _ = tx.send(marker(position)); } } #[cfg(test)] mod tests { use super::*; use crate::config::Config; fn db() -> (tempfile::TempDir, crate::db::Db) { let tmp = tempfile::tempdir().unwrap(); let db = crate::db::Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); (tmp, db) } #[test] fn positions_advance_only_when_the_batch_commits() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_at(5); let failed = sequencer.commit( db.inner.batch(), |assign, _| { assert_eq!((seq.take(assign)?, seq.take(assign)?), (5, 6)); Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, ); assert!(failed.is_err()); assert_eq!(seq.next(), 5); let taken = sequencer .commit( db.inner.batch(), |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) }, |taken| taken, ) .unwrap(); assert_eq!(taken, [5, 6]); assert_eq!(seq.head(), Some(6)); } #[test] fn positions_stop_at_the_last_one() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); 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 (_tmp, db) = db(); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); assert_eq!(seq.next(), 1); seq.resume_after(&db.cursors, Some(LAST_POSITION - 1)) .unwrap(); assert_eq!(seq.next(), LAST_POSITION); assert!(seq.resume_after(&db.cursors, Some(LAST_POSITION)).is_err()); assert!(seq.resume_after(&db.cursors, Some(u64::MAX)).is_err()); } #[test] fn a_counter_resumes_past_what_its_commits_took() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); sequencer .commit( db.inner.batch(), |assign, _| (0..3).try_for_each(|_| seq.take(assign).map(drop)), |_| {}, ) .unwrap(); // none of those positions has a row to resume after, so it resumes at the reservation let reopened = Sequence::new(b"test_seq"); reopened.resume_after(&db.cursors, None).unwrap(); assert_eq!(reopened.next(), 4 + SEQUENCE_BLOCK); // rows further along still win reopened .resume_after(&db.cursors, Some(SEQUENCE_BLOCK + 9)) .unwrap(); assert_eq!(reopened.next(), SEQUENCE_BLOCK + 10); } #[test] fn positions_taken_past_a_reservation_are_not_repeated_after_a_crash() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); let mut last = 0; for _ in 0..SEQUENCE_BLOCK + 500 { last = sequencer .commit(db.inner.batch(), |assign, _| seq.take(assign), |t| t) .unwrap(); } // a crash keeps only what committed, so a fresh counter reading the marks is one let reopened = Sequence::new(b"test_seq"); reopened.resume_after(&db.cursors, None).unwrap(); let next = reopened.next(); assert!(next > last, "{next} was already taken"); assert!( next <= last + 1 + SEQUENCE_BLOCK, "{next} skipped more than a block" ); } #[test] fn an_exact_mark_from_before_reservations_resumes_where_it_stopped() { let (_tmp, db) = db(); // older databases wrote the next position itself on every commit db.cursors.insert(b"test_seq", 7_u64.to_be_bytes()).unwrap(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); assert_eq!(seq.next(), 7); let taken = sequencer .commit(db.inner.batch(), |assign, _| seq.take(assign), |t| t) .unwrap(); assert_eq!(taken, 7); let reopened = Sequence::new(b"test_seq"); reopened.resume_after(&db.cursors, None).unwrap(); assert_eq!(reopened.next(), 8 + SEQUENCE_BLOCK); } #[test] fn a_hundred_thousand_commits_dont_fill_the_cursors_memtable() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"stream_next|test_ids"); seq.resume_after(&db.cursors, None).unwrap(); for _ in 0..100_000 { sequencer .commit(db.inner.batch(), |assign, _| seq.take(assign), |_| {}) .unwrap(); } // a full memtable asks a background worker to rotate it, so give that a moment let deadline = std::time::Instant::now() + std::time::Duration::from_millis(500); while std::time::Instant::now() < deadline { assert_eq!( db.cursors.sealed_memtable_count() + db.cursors.table_count(), 0 ); std::thread::sleep(std::time::Duration::from_millis(10)); } } #[test] fn a_position_mark_is_read_big_endian() { let (_tmp, db) = db(); let seq = Sequence::new(b"test_seq"); db.cursors .insert(b"test_seq", [0, 0, 0, 0, 0, 0, 1, 2]) .unwrap(); seq.resume_after(&db.cursors, None).unwrap(); assert_eq!(seq.next(), 258); db.cursors.insert(b"test_seq", [0; 7]).unwrap(); assert!(seq.resume_after(&db.cursors, None).is_err()); } #[cfg(feature = "jetstream")] #[test] fn clock_increases_even_when_ahead_of_the_wall_clock() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let clock = Clock::new(b"test_clock"); let ahead = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap() + 60_000_000; clock.resume_after(&db.cursors, Some(ahead)).unwrap(); // a failed commit takes nothing let failed = sequencer.commit( db.inner.batch(), |assign, _| { clock.tick(assign)?; Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, ); assert!(failed.is_err()); let ticks = sequencer .commit( db.inner.batch(), |assign, _| Ok([clock.tick(assign)?, clock.tick(assign)?]), |ticks| ticks, ) .unwrap(); assert_eq!(ticks, [ahead + 1, ahead + 2]); } #[cfg(feature = "jetstream")] #[test] fn a_commit_reads_the_wall_clock_once() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let clock = Clock::new(b"test_clock"); let wall = || u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap(); let before = wall(); let ticks = sequencer .commit( db.inner.batch(), |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)?]) }, |ticks| ticks, ) .unwrap(); assert!(ticks[0] >= before); assert_eq!(ticks, [ticks[0], ticks[0] + 1, ticks[0] + 2]); // the next commit starts at the wall time, which has moved on let next = sequencer .commit(db.inner.batch(), |assign, _| clock.tick(assign), |t| t) .unwrap(); assert!(next >= ticks[0] + 5_000); } #[test] fn unbroadcast_runs_are_announced_before_the_next_broadcast() { let (tx, mut rx) = tokio::sync::broadcast::channel(16); publish_in_order( &tx, [ (1, None), (2, None), (3, Some("live 3")), (4, Some("live 4")), (5, None), ], |position| match position { 2 => "marker 2", 5 => "marker 5", _ => panic!("unexpected marker for {position}"), }, ); let received: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect(); assert_eq!(received, ["marker 2", "live 3", "live 4", "marker 5"]); } }