Something went wrong. Try again.
Identities for entities did.bot
agent llm did
Something went wrong. Try again.
12 kB · 330 lines
Rust
at commit 18ba4fe0
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331//! The stream sequence number `com.atproto.sync.subscribeRepos` cursors are.//!//! atproto's firehose cursor is a bare integer, and a consumer hands back the//! last one it applied. That makes the number a promise rather than a counter://! two different events must never wear the same number, across a restart//! included, or a consumer resuming from a cursor is handed events it has//! already applied and told they are new ones.//!//! `/firehose` — this server's own Server-Sent-Events record stream — answers//! that by prefixing its cursor with an instance identifier minted at boot, so//! a number from a previous run is visibly not one of this run's. The real//! lexicon has nowhere to put such a prefix, so the number itself has to//! survive the restart.//!//! # Reservations, not one append per event//!//! Writing the number down for every event would double the write-ahead log's//! traffic to keep a counter. Instead a run *reserves* a block of//! [`RESERVATION`] numbers, writes down the highest one in the block, and//! hands them out from memory; the next reservation is only written when the//! block runs out.//!//! A restart resumes at the last reserved number rather than at the last//! *assigned* one, so a crash skips the rest of the block. That is a gap in//! the numbering and not a repeat, which is the direction that matters: the//! lexicon requires the sequence to increase, not to be gapless, and a//! consumer whose cursor falls in the gap is behind a restart that emptied the//! replay buffer anyway — it is told `OutdatedCursor` for that reason and not//! for this one.
use std::sync::{Arc, Mutex};
use crate::wal::Entry;
/// How many numbers one reservation covers.////// A thousand and twenty-four: at the development swarm's few dozen writes a/// second that is one synchronous log append every half-minute or so, and the/// most a crash can skip is a block that no consumer had been handed.pub const RESERVATION: u64 = 1024;
/// Where a sequence reservation is written down.////// One method, because there is one fact: numbers up to and including/// `through` have been spoken for and must never be handed out by a later run/// of this server. Implemented over the write-ahead log by/// [`FileSequenceLog`](crate::durable::FileSequenceLog); the in-memory/// implementation writes nothing, which is the right behaviour for a/// deployment that keeps nothing else either.pub trait SequenceLog: Send + Sync + std::fmt::Debug { /// Records that no later run may hand out a number at or below `through`. /// /// Must not return until the fact is durable. It is called once per /// [`RESERVATION`] numbers, so it is allowed to be the expensive kind of /// write. fn reserve(&self, through: u64);}
/// A reservation log that forgets, for a deployment that keeps nothing.#[derive(Debug, Default)]pub struct MemorySequenceLog;
impl SequenceLog for MemorySequenceLog { fn reserve(&self, _through: u64) {}}
/// The highest reservation the journal carries.////// The store behind [`Entry::StreamReserved`], and the counterexample that/// keeps "append-only" out of the axes a durable store is sorted along: this/// one and the agent ledger both only ever grow, and where the ledger's/// replay appends and its checkpoint is the whole of it, this one's replay is/// a maximum and its checkpoint is a single entry.////// That single entry is the point. What a consumer can be handed after a/// restart is decided by `didbot_serve::subscribe`'s replay buffer, which/// lives in memory and is empty when a process starts: a cursor it cannot/// reach is answered `#info OutdatedCursor`. So what the journal owes the stream is a floor and not a/// history — the one number no later run may hand out again — and carrying/// the reservations that led to it would be carrying a history of a promise/// that has already been kept.#[derive(Debug, Default)]pub struct SequenceFloor { through: Mutex<u64>,}
impl SequenceFloor { /// A floor at zero: nothing has been reserved. #[must_use] pub fn new() -> Self { Self::default() }
/// The highest number any run has reserved. pub fn peek(&self) -> u64 { *self.lock() }
/// Raises the floor to `through`, if that is higher than where it stands. /// /// A maximum rather than an assignment, so that the order reservations /// arrive in cannot lower it: a floor that went down is a number handed /// out twice. pub fn raise(&self, through: u64) { let mut current = self.lock(); *current = (*current).max(through); }
fn lock(&self) -> std::sync::MutexGuard<'_, u64> { self.through .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) }}
impl crate::journal::Journaled for SequenceFloor { const STORE: &'static str = "sequence";
fn owns(entry: &Entry) -> bool { matches!(entry, Entry::StreamReserved { .. }) }
fn apply(&self, entry: Entry) { // Only what `owns` claims reaches here. if let Entry::StreamReserved { through } = entry { self.raise(through); } }
/// One entry, or none at a floor of zero: a stream that has numbered /// nothing needs nothing written to resume where it left off. fn checkpoint(&self) -> Vec<Entry> { let through = self.peek(); (through > 0) .then_some(Entry::StreamReserved { through }) .into_iter() .collect() }
fn durable_through(&self) -> crate::journal::Mark { crate::journal::Mark::ORIGIN }}
/// A monotonic stream sequence number that survives a restart.#[derive(Debug)]pub struct Sequence { state: Mutex<State>, log: Arc<dyn SequenceLog>,}
#[derive(Debug)]struct State { /// The highest number handed out. assigned: u64, /// The highest number the log says is spoken for. reserved: u64,}
impl Sequence { /// A sequence resuming above `floor`, writing its reservations to `log`. /// /// `floor` is the highest number any previous run reserved — what /// [`Durable::sequence_floor`](crate::Durable::sequence_floor) replays out /// of the log — and zero for a deployment that has never assigned one. pub fn new(log: Arc<dyn SequenceLog>, floor: u64) -> Self { Self { state: Mutex::new(State { assigned: floor, reserved: floor, }), log, } }
/// A sequence that starts at one and forgets when the process ends. #[must_use] pub fn in_memory() -> Self { Self::new(Arc::new(MemorySequenceLog), 0) }
/// The next number, reserving another block if this one is spent. /// /// Callers hold their own stream lock across this and the buffering of the /// event it numbers, so that two concurrent writes cannot reach the replay /// buffer in the opposite order to their numbers. pub fn next(&self) -> u64 { let mut state = self.lock(); let seq = state.assigned + 1; if seq > state.reserved { // Reserved before it is handed out, never after: a number that // reached a consumer and not the log is one a later run would // hand out again. let through = seq.saturating_add(RESERVATION - 1); self.log.reserve(through); state.reserved = through; } state.assigned = seq; seq }
/// The highest number handed out, or the floor if none has been. pub fn latest(&self) -> u64 { self.lock().assigned }
/// Takes the lock, recovering from a poisoned mutex. /// /// The same posture the stores take: everything under it is two integer /// updates, so a panic elsewhere cannot have left it half-written, and /// refusing to number events over an unrelated panic would stop the /// stream for good. fn lock(&self) -> std::sync::MutexGuard<'_, State> { self.state .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) }}
impl Default for Sequence { fn default() -> Self { Self::in_memory() }}
#[cfg(test)]mod tests { use super::*;
/// A log that remembers the last reservation, as a restart would read it /// back. #[derive(Debug, Default)] struct Recording(Mutex<Vec<u64>>);
impl Recording { fn floor(&self) -> u64 { self.0.lock().unwrap().last().copied().unwrap_or(0) } }
impl SequenceLog for Recording { fn reserve(&self, through: u64) { self.0.lock().unwrap().push(through); } }
#[test] fn numbers_start_at_one_and_do_not_skip_within_a_run() { let sequence = Sequence::in_memory(); let seqs: Vec<u64> = (0..5).map(|_| sequence.next()).collect(); assert_eq!(seqs, vec![1, 2, 3, 4, 5]); assert_eq!(sequence.latest(), 5); }
#[test] fn one_reservation_covers_a_whole_block() { let log = Arc::new(Recording::default()); let sequence = Sequence::new(log.clone(), 0); for _ in 0..RESERVATION { sequence.next(); } assert_eq!(*log.0.lock().unwrap(), vec![RESERVATION]); sequence.next(); assert_eq!( *log.0.lock().unwrap(), vec![RESERVATION, RESERVATION * 2], "the block ran out, so another was written down" ); }
/// The floor is a maximum, so a reservation arriving under one already /// applied leaves it where it is. Replay is in order, so this is what a /// log written by two runs cannot do rather than what one does; a floor /// that could go down is a number handed out twice. #[test] fn the_floor_only_ever_rises() { use crate::journal::Journaled;
let floor = SequenceFloor::new(); assert!( floor.checkpoint().is_empty(), "a stream that has numbered nothing writes nothing down" );
floor.apply(Entry::StreamReserved { through: 4096 }); floor.apply(Entry::StreamReserved { through: 1024 }); assert_eq!(floor.peek(), 4096); }
/// One entry, whatever the log held: the checkpoint is the floor and not /// the reservations that led to it. What a consumer may be handed after a /// restart is bounded by a replay buffer that a restart empties, so the /// history behind the floor is a promise already kept. #[test] fn a_checkpoint_of_the_floor_is_one_entry_and_replays_to_the_same_floor() { use crate::journal::Journaled;
let floor = SequenceFloor::new(); for through in [1024, 2048, 3072] { floor.apply(Entry::StreamReserved { through }); }
let fresh = SequenceFloor::new(); let (written, replayed) = crate::journal::round_trip(&floor, &fresh); assert_eq!(written, replayed); assert_eq!(fresh.peek(), 3072); assert_eq!( fresh.checkpoint().len(), 1, "three reservations came back as more than the floor they add up to" ); }
/// The property the whole module exists for: a number handed out before a /// restart is never handed out again after one. #[test] fn a_restart_resumes_above_every_number_it_could_have_handed_out() { let log = Arc::new(Recording::default()); let before = Sequence::new(log.clone(), 0); let handed: Vec<u64> = (0..7).map(|_| before.next()).collect(); let highest = *handed.last().expect("seven numbers were handed out"); assert_eq!(highest, before.latest());
let after = Sequence::new(log.clone(), log.floor()); assert!( after.next() > highest, "a restart handed out a number a previous run had already used" ); }}