Something went wrong. Try again.
Identities for entities did.bot
agent llm did
Something went wrong. Try again.
32 kB · 791 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792//! 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<Arc<str>>,}
/// One repository's line: mutual exclusion plus the freeze transition.struct Line { state: Mutex<LineState>, 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<Line>, did: &str, capacity: usize) -> Result<Ticket, QueueError> { 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<str>) { 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<Line>` — 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<Line>, 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<R> { /// 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<String>,}
/// 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::Provisionerpub struct WriteQueue { lines: Mutex<HashMap<String, Arc<Line>>>, 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<Line> { 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<R>( &self, gate: &dyn PolicyGate, subject: &didbot_policy::Subject<'_>, admit: impl FnOnce() -> R, ) -> Result<Admitted<R>, 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<R>( &self, gate: &dyn PolicyGate, subjects: &[didbot_policy::Subject<'_>], admit: impl FnOnce() -> R, ) -> Result<Admitted<R>, 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<str> = 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<AtomicUsize>, release: Mutex<std::sync::mpsc::Receiver<()>>, } 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<Barrier>, } 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<AtomicUsize>, release_freezer: Mutex<std::sync::mpsc::Receiver<()>>, } 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"); }}