//! Per-repository serialization for policy-gated writes. //! //! `plan/policy.md`'s "Evaluation happens at the write" sets three //! constraints on where a write's policy check runs: //! //! * **The linearization unit is the repository, not the server.** One //! agent's commits chain, each naming its predecessor, so order within a //! repository is load-bearing. Order *between* two agents' repositories is //! observable by nobody, so serializing judgment across the whole server //! would make one agent's slow write hold up every other agent's, for no //! correctness reason. [`WriteQueue`] gives each repository its own line; //! different repositories never wait on each other here, and //! `crate::provision::Provisioner::writing_to` holds the same shape over the //! commit itself. //! * **A policy judgment must run outside the store's lock.** A judgment may //! be slow — cross a process boundary, hit a network resolution, run a //! model — and a store's lock (see `crate::durable` and //! `crate::provision::Provisioner::writing_to`) is not a lock a slow //! operation may hold: every commit to the repository goes through it. //! [`WriteQueue::submit`] calls //! [`PolicyGate::judge`] while holding //! only this repository's line — never a store lock — and only asks //! the caller to take one, inside `admit`, once the judgment is //! `Allow`. The sequence number a commit is stamped with is assigned //! inside `admit`, which is what "assigned at admission, not arrival" //! means in practice: nothing about a write's place in the chain is //! decided until policy has already let it through. //! * **A denied write must never need un-committing.** Because `admit` never //! runs for `Reject` or `Freeze`, there is nothing to undo — the //! write-ahead log never sees a denied write in the first place. //! //! # Freeze as a queue transition //! //! The first write that trips a freeze must drain every write already //! waiting for that repository's turn, rejecting each with the freeze's //! reason rather than letting it queue behind a repository that no longer //! accepts writes — see `plan/policy.md`'s "Outcomes". `Line` does this //! without an explicit walk over waiters: a ticket-lock's own wakeup order is //! already a queue, so setting `LineState::frozen` and releasing the //! ticket that tripped it wakes the next ticket, which finds the line frozen //! and immediately releases its own turn without ever calling //! [`PolicyGate::judge`] — which wakes the //! one after it, and so on. The cascade *is* the drain. //! //! # Bounding the queue //! //! An agent can submit writes faster than a slow evaluator drains them. //! [`WriteQueue`] bounds how many writes may be waiting for one repository's //! turn; past that bound, [`WriteQueue::submit`] returns //! [`QueueError::Full`] immediately, without waiting and without calling the //! gate. That is a refusal a caller can act on — the same shape as any other //! write refusal — rather than unbounded memory held for writes nothing is //! draining. use std::collections::HashMap; use std::sync::{Arc, Condvar, Mutex}; use crate::policy::{Outcome, PolicyGate, PolicyVersion}; /// How many writes may be waiting for one repository's turn at once, if a /// deployment does not choose its own. /// /// Loose on purpose: this bounds one repository's backlog, not the server's /// total concurrency, and every repository gets this budget independently. /// An agent that floods its own queue faster than policy drains it is a /// single misbehaving agent, and [`QueueError::Full`] is the answer it gets /// — not every other agent slowing down to match. /// /// This crate's value rather than a deployment's for the same reason it is /// loose: it bounds one repository's backlog, which follows from how fast /// policy drains a queue, not from how many accounts a deployment holds. A /// server sized for more accounts wants the same per-repository budget. pub const DEFAULT_CAPACITY: usize = 64; /// Why a submitted write never reached a verdict from /// [`PolicyGate::judge`]. #[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] pub enum QueueError { /// More writes were already waiting for this repository's turn than the /// queue allows. #[error( "{did} already has {pending} write(s) waiting for their turn, at the limit of {capacity}" )] Full { /// The repository whose queue is full. did: String, /// How many writes were already waiting. pending: usize, /// The configured limit. capacity: usize, }, /// A policy refused the write. See [`crate::policy::Outcome::Reject`]. #[error("write refused: {reason}")] Rejected { /// The reason a policy authored for this refusal. reason: String, }, /// The repository is frozen — by this write's own judgment, or by an /// earlier one that was still waiting its turn when a freeze landed. See /// [`crate::policy::Outcome::Freeze`]. #[error("account frozen: {reason}")] Frozen { /// The reason the freeze was tripped for. reason: String, }, } /// One repository's place in line, and whether it is frozen. /// /// Just two counters and an optional reason — no `VecDeque` of waiters is /// kept, because a ticket-lock does not need one: every waiter already knows /// its own ticket number, and "how many are waiting" is `next_ticket - /// serving`. #[derive(Default)] struct LineState { /// The next ticket not yet handed out. next_ticket: u64, /// The ticket currently allowed to judge-and-commit. serving: u64, /// Set the moment a write trips a freeze; cleared only by /// [`Line::unfreeze`]. See the module docs' "Freeze as a queue /// transition". frozen: Option>, } /// One repository's line: mutual exclusion plus the freeze transition. struct Line { state: Mutex, turn: Condvar, } impl Line { fn new() -> Self { Self { state: Mutex::new(LineState::default()), turn: Condvar::new(), } } /// Takes a ticket for `did`, refusing outright if the line is already at /// `capacity`, then waits for that ticket's turn. /// /// Returns [`QueueError::Frozen`] — without ever calling a gate — the /// moment this ticket's turn comes if the line is frozen by then, whether /// it already was when this ticket was taken or a write ahead of it /// tripped the freeze while this one waited. This is also, not /// incidentally, what keeps a metapolicy from ever watching its own /// freeze cause more of what it froze the agent for: the write that /// trips a freeze is the *last* one this line ever hands to `judge` for /// that freeze (see [`WriteQueue::submit`]'s `Outcome::Freeze` arm, /// which sets [`LineState::frozen`] before this ticket even leaves), and /// every ticket queued behind it is refused right here, before `judge` /// runs — so [`PolicyGate::observe`] /// never sees a denial this freeze itself caused. A "drain" implemented /// as a parallel pass that judged the backlog and rejected it after the /// fact would reopen exactly that loop; the correctness here depends on /// staying a refusal *before* judgment, not after it. fn enter(&self, handle: Arc, did: &str, capacity: usize) -> Result { let mut state = self .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); let pending = state.next_ticket.saturating_sub(state.serving); if pending as usize >= capacity { return Err(QueueError::Full { did: did.to_owned(), pending: pending as usize, capacity, }); } let ticket = state.next_ticket; state.next_ticket += 1; while state.serving != ticket { state = self .turn .wait(state) .unwrap_or_else(std::sync::PoisonError::into_inner); } let frozen = state.frozen.clone(); drop(state); if let Some(reason) = frozen { self.leave(ticket); return Err(QueueError::Frozen { reason: reason.to_string(), }); } Ok(Ticket { line: handle, ticket, }) } /// Advances the line past `ticket` and wakes whoever is waiting for the /// next one. /// /// `notify_all` rather than `notify_one`: only the ticket that is /// actually next can pass its own `while` check, so every other waiter /// that wakes just re-checks and sleeps again. That is a cheap, correct /// price for not tracking which condvar wait belongs to which ticket. fn leave(&self, ticket: u64) { let mut state = self .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); debug_assert_eq!(state.serving, ticket, "a ticket left out of turn"); state.serving = ticket + 1; drop(state); self.turn.notify_all(); } /// Trips the freeze. Called only by the ticket currently holding its /// turn — the one whose judgment produced `Outcome::Freeze` — so no /// other write can be mid-judgment when this happens: a line has exactly /// one ticket "at its turn" at any moment. fn freeze(&self, reason: Arc) { let mut state = self .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); state.frozen = Some(reason); } /// Lifts a freeze, so this repository's writes reach a gate again. /// /// Only ever called for an operator's explicit action — /// `plan/policy.md`'s "An evaluator can never unfreeze" — never by this /// module on its own. See [`WriteQueue::unfreeze`]. fn unfreeze(&self) { let mut state = self .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); state.frozen = None; } } /// One repository's exclusive turn, held until dropped. /// /// A guard rather than a bare ticket number so a panic during judgment or /// admission still releases the turn: an unwind drops this and its own /// `Drop` calls `Line::leave`, rather than leaving every future write to /// this repository waiting forever for a turn that panicked instead of /// finishing. Holds its own `Arc` — not just the ticket number — for /// exactly that reason: the guard has to be able to release the line /// without `submit` doing it explicitly on every path, which is what makes /// the panicking path and the ordinary path the same code. struct Ticket { line: Arc, ticket: u64, } /// The result of an admitted write, with the policy version that admitted /// it. /// /// `plan/policy.md`'s "Record which policy version admitted each write": /// this is where that version reaches a caller, so it can be attached to /// whatever this deployment's accountability record for the write turns out /// to be — wiring it into a persisted shape is left for the pass that /// connects this crate to `didbot-policy`'s real version type, since /// `plan/policy.md` leaves that wrapper shape unsettled too. #[derive(Debug, Clone)] pub struct Admitted { /// What `admit` returned. pub result: R, /// The version [`PolicyGate::judge`] /// ran under. pub version: PolicyVersion, /// The operator revision it ran under; see /// [`PolicyGate::judged_revision`]. pub revision: Option, } /// Per-repository write queues, keyed by agent DID. /// /// Holds one `Line` per repository that has ever submitted a write, kept /// for the life of the process. That is a small, bounded cost — one mutex /// and two integers per line — of the same order [`Provisioner`]'s own /// per-account bookkeeping already pays for every account this deployment /// has provisioned. /// /// [`Provisioner`]: crate::provision::Provisioner pub struct WriteQueue { lines: Mutex>>, capacity: usize, } impl std::fmt::Debug for WriteQueue { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("WriteQueue") .field("capacity", &self.capacity) .finish_non_exhaustive() } } impl WriteQueue { /// A queue bounding each repository's backlog to `capacity` writes. #[must_use] pub fn new(capacity: usize) -> Self { Self { lines: Mutex::new(HashMap::new()), capacity, } } /// A queue with [`DEFAULT_CAPACITY`]. #[must_use] pub fn with_default_capacity() -> Self { Self::new(DEFAULT_CAPACITY) } fn line(&self, did: &str) -> Arc { let mut lines = self .lines .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); lines .entry(did.to_owned()) .or_insert_with(|| Arc::new(Line::new())) .clone() } /// Judges `subject` against `gate` — outside this repository's store /// lock, whatever `admit` closes over to take one — then, only on /// [`Outcome::Allow`], runs `admit` to perform the write. /// /// `admit` is responsible for its own locking around the store; this /// only decides *when* it may run and *in what order* relative to other /// writes to `subject.account`'s repository. On any other outcome, `admit` /// never runs at all, which is what makes a denial cheap: nothing was /// ever appended for it, so there is nothing to undo. /// /// The [`PolicyVersion`] on a successful [`Admitted`] is read from `gate` /// immediately after `judge` returns `Allow`, before `admit` runs — so a /// policy set that reloads while `admit` is doing its (possibly /// non-trivial) work cannot leave the recorded version disagreeing with /// the one that actually judged this write. pub fn submit( &self, gate: &dyn PolicyGate, subject: &didbot_policy::Subject<'_>, admit: impl FnOnce() -> R, ) -> Result, QueueError> { self.submit_batch(gate, std::slice::from_ref(subject), admit) } /// Judges every subject in one batch under a *single* turn in this /// repository's line, and runs `admit` only if none of them denied. /// /// This is [`Self::submit`]'s shape for `com.atproto.repo.applyWrites`, /// where several writes become one commit. Three things follow from that /// and are the reason this is not a loop over `submit` at the call site: /// /// * **One turn, not one per operation.** A batch is one commit built on /// one head. Taking the line once is what keeps another write to the /// same repository from interleaving between the judgment of operation /// two and operation three, which would leave the batch judged against /// a repository state it is not committed against. /// * **All or nothing.** `admit` runs only when every operation is /// allowed. A batch is atomic in the record store already, and a /// partially applied batch would be a commit whose meaning no caller /// asked for; deny-only monotonicity says the answer to "some of this /// is refused" is to refuse. /// * **Severity, not position, decides the answer.** Outcomes fold /// through [`didbot_policy::worst`], so a freeze tripped by any /// operation reaches the caller as a freeze even when another /// operation merely rejected — and the account is frozen, which a /// generic batch refusal would have silently dropped. /// /// Judgment stops at the first `Freeze` and every later operation goes /// unjudged, for the same reason a freeze drains the writes queued /// behind it without judging them: a metapolicy must never observe /// denials its own freeze caused. A `Reject` does **not** stop judgment, /// so every operation the caller attempted is judged and observed — a /// batch must not be a cheaper way to hide attempts from a stateful /// evaluator than the same writes sent one at a time. /// /// Every subject must name the same repository; a batch spans exactly /// one, because it becomes exactly one commit. pub fn submit_batch( &self, gate: &dyn PolicyGate, subjects: &[didbot_policy::Subject<'_>], admit: impl FnOnce() -> R, ) -> Result, QueueError> { // This queue only ever serializes writes -- see this module's own // docs -- so a subject is always `Subject::Write`, never a `Grant` // request judged before any repository line exists to serialize it // against. let account = match subjects.first() { Some(didbot_policy::Subject::Write { account, .. }) => *account, Some(_) => { unreachable!("WriteQueue::submit_batch is only ever called with Subject::Write") } // A batch with nothing in it is nothing to judge. Callers refuse // to reach here — `Provisioner::apply_writes` answers an empty // batch before the queue — and there is no write for a gate to // have an opinion about if one ever did. None => unreachable!("WriteQueue::submit_batch is never called with no subjects"), }; debug_assert!( subjects .iter() .all(|subject| matches!(subject, didbot_policy::Subject::Write { account: other, .. } if *other == account)), "a batch spans exactly one repository, because it becomes exactly one commit" ); let line = self.line(account); // `ticket` releases this repository's turn when it drops — see // `Ticket`'s own docs — so nothing below calls `Line::leave` // explicitly. It drops wherever this function's control flow leaves // it: at the end of each match arm below on the ordinary paths, or // partway through one if `gate.judge` or `admit` panics. let _ticket = line.enter(line.clone(), account, self.capacity)?; // One set for every operation of this batch, and for the version // and revision read off the gate below. Dropped with this call, // however it leaves. See `PolicyGate::pin_set`. let _pinned = gate.pin_set(); let mut outcomes = Vec::with_capacity(subjects.len()); for subject in subjects { let outcome = gate.judge(subject); // Never short-circuits, unlike everything below it: a metapolicy // watching for repeated denials must see this attempt even though // `judge` already decided it, so `observe` runs for every outcome, // including the two that stop `admit` from ever running. gate.observe(subject, &outcome); let froze = matches!(outcome, Outcome::Freeze { .. }); outcomes.push(outcome); if froze { // Nothing after a freeze is judged or observed, exactly as // nothing queued behind one is. See this function's docs. break; } } match didbot_policy::worst(outcomes) { Outcome::Allow => { let version = gate.version(); let revision = gate.judged_revision(); let result = admit(); Ok(Admitted { result, version, revision, }) } Outcome::Reject { reason } => Err(QueueError::Rejected { reason }), Outcome::Freeze { reason } => { let reason: Arc = Arc::from(reason.as_str()); // Set before this ticket leaves: the next ticket to wake // must see it, or it would be judged as if nothing had // happened. `ticket` is still held here and does not drop // until this arm finishes, so this runs before that leave. line.freeze(reason.clone()); gate.froze(account, &reason); Err(QueueError::Frozen { reason: reason.to_string(), }) } } } /// Lifts a freeze on one repository's line, so its writes reach a gate /// again. /// /// This is only half of unfreezing an account: the account's own /// [`AccountState`](crate::account::AccountState) is a separate fact, /// changed elsewhere, that governs paths which do not go through this /// queue at all. A caller unfreezing an account is expected to change /// both. A DID this queue has never seen a write for is a no-op — there /// is no line to unfreeze. pub fn unfreeze(&self, did: &str) { let line = self .lines .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .get(did) .cloned(); if let Some(line) = line { line.unfreeze(); } } } impl Drop for Ticket { fn drop(&mut self) { // The only place `leave` is called. `submit` never calls it // directly — it lets this ticket fall out of scope instead, on // every path: `Allow` after `admit` returns, `Reject` and `Freeze` // at the end of their match arms, and a panic inside `gate.judge` or // `admit` unwinding through here same as any other drop. Without // this, a panicking write would leave `serving` stuck at its own // ticket forever, and every later write to the same repository // would wait for a turn that can never come — a permanent, silent // stall dressed up as a slow write. self.line.leave(self.ticket); } } #[cfg(test)] mod tests { use super::*; use crate::kind::AccountKind; use crate::policy::{write_subject, PolicyVersion, WriteAction}; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Barrier; use std::time::Duration; fn subject(account: &str) -> didbot_policy::Subject<'_> { write_subject( account, AccountKind::Agent, None, "did:web:pds.example", "pds.example", None, None, "app.bsky.feed.post", WriteAction::Create, &[], None, None, &[], ) } struct AllowGate; impl PolicyGate for AllowGate { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { Outcome::Allow } fn version(&self) -> PolicyVersion { PolicyVersion("test".to_owned()) } } #[test] fn an_admitted_write_carries_the_judging_version() { let queue = WriteQueue::new(4); let admitted = queue .submit(&AllowGate, &subject("did:plc:a"), || 42) .unwrap(); assert_eq!(admitted.result, 42); assert_eq!(admitted.version, PolicyVersion("test".to_owned())); } struct RejectGate; impl PolicyGate for RejectGate { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { Outcome::Reject { reason: "no".to_owned(), } } fn version(&self) -> PolicyVersion { PolicyVersion::none() } } #[test] fn a_rejected_write_never_runs_admit() { let queue = WriteQueue::new(4); let ran = AtomicUsize::new(0); let error = queue .submit(&RejectGate, &subject("did:plc:a"), || { ran.fetch_add(1, Ordering::SeqCst); }) .unwrap_err(); assert_eq!(ran.load(Ordering::SeqCst), 0); assert!(matches!(error, QueueError::Rejected { reason } if reason == "no")); } #[test] fn a_full_queue_refuses_without_waiting_or_judging() { let queue = WriteQueue::new(1); let judged = Arc::new(AtomicUsize::new(0)); struct CountingBlockGate { judged: Arc, release: Mutex>, } impl PolicyGate for CountingBlockGate { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { self.judged.fetch_add(1, Ordering::SeqCst); // Held until the test says the second submission has already // been refused, so the queue is provably still occupied when // that refusal happens. let release = self .release .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); let _ = release.recv(); Outcome::Allow } fn version(&self) -> PolicyVersion { PolicyVersion::none() } } let (tx, rx) = std::sync::mpsc::channel(); let gate = Arc::new(CountingBlockGate { judged: judged.clone(), release: Mutex::new(rx), }); let queue = Arc::new(queue); let holder = { let queue = queue.clone(); let gate = gate.clone(); std::thread::spawn(move || { queue .submit(gate.as_ref(), &subject("did:plc:a"), || ()) .unwrap(); }) }; // Give the holder time to reach `judge` and block there. while judged.load(Ordering::SeqCst) == 0 { std::thread::sleep(Duration::from_millis(5)); } let refused = queue.submit(gate.as_ref(), &subject("did:plc:a"), || ()); assert!(matches!( refused, Err(QueueError::Full { pending: 1, capacity: 1, .. }) )); // The refused write never reached the gate. assert_eq!(judged.load(Ordering::SeqCst), 1); tx.send(()).unwrap(); holder.join().unwrap(); } #[test] fn two_repositories_never_wait_on_each_other() { struct SlowGate { barrier: Arc, } impl PolicyGate for SlowGate { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { // Both threads must reach this before either proceeds. If // the queue serialized across repositories, the second // submission below would never reach its own `judge` call // concurrently with the first, and this would deadlock — // which is exactly what this test is checking does not // happen. self.barrier.wait(); Outcome::Allow } fn version(&self) -> PolicyVersion { PolicyVersion::none() } } let barrier = Arc::new(Barrier::new(2)); let gate = Arc::new(SlowGate { barrier: barrier.clone(), }); let queue = Arc::new(WriteQueue::new(4)); let a = { let queue = queue.clone(); let gate = gate.clone(); std::thread::spawn(move || { queue .submit(gate.as_ref(), &subject("did:plc:a"), || ()) .unwrap(); }) }; let b = { let queue = queue.clone(); let gate = gate.clone(); std::thread::spawn(move || { queue .submit(gate.as_ref(), &subject("did:plc:b"), || ()) .unwrap(); }) }; a.join().expect("repository a's write completed"); b.join().expect("repository b's write completed"); } #[test] fn a_freeze_drains_writes_already_waiting_without_judging_them() { struct OneShotFreezeGate { judged: Arc, release_freezer: Mutex>, } impl PolicyGate for OneShotFreezeGate { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { let seen = self.judged.fetch_add(1, Ordering::SeqCst); if seen == 0 { // The first ticket in: block until the test has queued // up followers behind it, then freeze. let release = self .release_freezer .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); let _ = release.recv(); Outcome::Freeze { reason: "tripped".to_owned(), } } else { // No follower should ever get here. Outcome::Allow } } fn version(&self) -> PolicyVersion { PolicyVersion::none() } } let judged = Arc::new(AtomicUsize::new(0)); let (tx, rx) = std::sync::mpsc::channel(); let gate = Arc::new(OneShotFreezeGate { judged: judged.clone(), release_freezer: Mutex::new(rx), }); let queue = Arc::new(WriteQueue::new(8)); // The write that will trip the freeze, occupying the line first. let leader = { let queue = queue.clone(); let gate = gate.clone(); std::thread::spawn(move || queue.submit(gate.as_ref(), &subject("did:plc:a"), || ())) }; while judged.load(Ordering::SeqCst) == 0 { std::thread::sleep(Duration::from_millis(5)); } // Three followers, queued up behind the leader's still-blocked // judgment. let followers: Vec<_> = (0..3) .map(|_| { let queue = queue.clone(); let gate = gate.clone(); std::thread::spawn(move || { queue.submit(gate.as_ref(), &subject("did:plc:a"), || ()) }) }) .collect(); // Give the followers time to enter the line and start waiting. std::thread::sleep(Duration::from_millis(50)); tx.send(()).unwrap(); let leader_result = leader.join().unwrap(); assert!(matches!(leader_result, Err(QueueError::Frozen { .. }))); for follower in followers { let result = follower.join().unwrap(); assert!( matches!(result, Err(QueueError::Frozen { .. })), "a follower queued behind a freeze should be refused without being judged" ); } // Exactly one judgment: the leader's. Every follower was drained. assert_eq!(judged.load(Ordering::SeqCst), 1); } #[test] fn unfreeze_lets_writes_reach_the_gate_again() { let queue = WriteQueue::new(4); struct FreezeThenAllow { first: AtomicUsize, } impl PolicyGate for FreezeThenAllow { fn judge(&self, _subject: &didbot_policy::Subject<'_>) -> Outcome { if self.first.fetch_add(1, Ordering::SeqCst) == 0 { Outcome::Freeze { reason: "tripped".to_owned(), } } else { Outcome::Allow } } fn version(&self) -> PolicyVersion { PolicyVersion::none() } } let gate = FreezeThenAllow { first: AtomicUsize::new(0), }; let frozen = queue.submit(&gate, &subject("did:plc:a"), || ()); assert!(matches!(frozen, Err(QueueError::Frozen { .. }))); let still_frozen = queue.submit(&gate, &subject("did:plc:a"), || ()); assert!( matches!(still_frozen, Err(QueueError::Frozen { .. })), "a later write must stay refused without reaching the gate again" ); queue.unfreeze("did:plc:a"); let admitted = queue.submit(&gate, &subject("did:plc:a"), || "ok").unwrap(); assert_eq!(admitted.result, "ok"); } }