//! What this host has put into the Pulse module table — and, where it cannot be //! sure, what it must go and find out before doing anything else. //! //! # The defect this exists for //! //! Module ids used to live in three `Option`s, two of them shared with the //! event task behind an `Arc>`. That representation cannot express *a //! load is in flight*, and the gap is reachable today: teardown calls //! `event_task.abort()` and then reads the ids, but the task's work is a //! synchronous `pactl` call with no await point inside it, so the abort cannot //! land until the load has already returned. Teardown therefore sees `None`, //! unloads the sink, and the still-running task stores the new module's id into a //! mutex nobody will ever read again — an orphan loopback pointing at a sink that //! no longer exists. //! //! A slot is consequently not an `Option` but a small state machine, and the //! transitions that matter are the ones that admit *we do not know*: //! //! ```text //! begin_load commit //! Vacant ───────────────► Loading ───────────────► Loaded //! ▲ │ │ │ //! │ abandon │ │ permit dropped │ begin_unload //! └────────────────────┘ │ unsettled ▼ //! ▲ ▼ Unloading //! │ Resolution::Nothing Ambiguous ◄──────────────┘ //! └───────────────────────► │ ▲ uncertain outcome //! │ │ //! Resolution::Conflict▼ └── Resolution::Adopt ──► Loaded //! Poisoned //! ``` //! //! # Why the permit is affine //! //! [`LoadPermit`] is not `Clone`, is consumed by value to settle, and its [`Drop`] //! marks the slot [`Reconcile::Load`] when it was never settled. The permit moves //! into a registered blocking worker before any await can detach that work; either //! the worker settles it, or dropping the worker marks the question ambiguous. //! Teardown waits for all such permits before reconciling. Two permitted loads for //! one slot cannot both commit, because [`ModuleLedger::begin_load`] issues a permit //! only for a `Vacant` slot and every other state — including the ambiguous one — //! refuses. //! //! # Why an ambiguous slot blocks the next load //! //! An unresolved load may or may not have created a module carrying our sink name. //! Starting another one on top of it risks two sinks sharing a `node.name`, which //! is measurably not an error the server reports: it accepts both, and //! `pulsesrc device=.monitor` attaches to the *older* one. A second capture //! session would then be silently stolen by the debris of the first. Refusing to //! proceed until the question is answered is the whole point. use std::collections::BTreeMap; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Condvar, Mutex}; use anyhow::{Context as _, Result}; use crate::repair::introspect::PulseSession; use crate::repair::plan::{ Fingerprint, ModuleObservation, OwnerToken, Shape, classify, recorded_argument, }; /// PulseAudio's "no such index" — `PA_INVALID_INDEX`, `(uint32_t) -1`. /// /// Every Pulse load callback reports failure by handing back this value rather /// than an index. We load through `pactl`, which reports failure by exiting /// non-zero instead, so seeing it on stdout would mean the tool printed a /// sentinel we must not mistake for a module: unloading it would be a request /// about something that cannot exist. Treated as a load whose outcome is unknown, /// never as a successful one. pub const PA_INVALID_INDEX: u32 = u32::MAX; /// What an unresolved slot needs looked up before it can be trusted again. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Reconcile { /// A load whose outcome is unknown: `pactl` was killed, timed out, or printed /// something we could not read as an index. The server may or may not have /// created the module, so it is looked for **by owner token** — the nonce is /// minted per load, so it names this attempt and no other. Load { shape: Shape, pid: u32, token: OwnerToken, }, /// An unload whose outcome is unknown. The module is looked for by its full /// fingerprint: still present means the unload did not happen, absent means it /// did. An id alone would not do, because Pulse reuses module indices verbatim. Unload { fp: Fingerprint }, } /// One shape's slot in the module table. #[derive(Debug, Clone, PartialEq, Eq)] pub enum SlotState { /// Nothing loaded, and nothing outstanding. Vacant, /// A load is in flight under `permit`. The id is meaningless until it commits. Loading { token: OwnerToken, permit: u64 }, /// A module we loaded and can name exactly. Loaded { fp: Fingerprint }, /// An unload is in flight. The fingerprint is retained deliberately: an unload /// that times out must not leave the module unrecorded. Unloading { fp: Fingerprint, permit: u64 }, /// The slot's real state is unknown and must be resolved against the server. Ambiguous(Reconcile), /// Resolution found something we refuse to act on. Terminal. Poisoned { reason: String }, } impl SlotState { /// A short, stable label for logs and refusal messages. fn label(&self) -> &'static str { match self { SlotState::Vacant => "vacant", SlotState::Loading { .. } => "loading", SlotState::Loaded { .. } => "loaded", SlotState::Unloading { .. } => "unloading", SlotState::Ambiguous(_) => "ambiguous", SlotState::Poisoned { .. } => "poisoned", } } } /// Why the ledger refused an operation. #[derive(Debug, Clone, PartialEq, Eq)] pub enum LedgerError { /// The slot was not free. Carries the state that refused, for the log line. Busy { shape: Shape, state: &'static str }, /// The slot is poisoned and will not be used again this session. Poisoned { shape: Shape, reason: String }, /// `pactl` reported `PA_INVALID_INDEX` where an index was expected. InvalidIndex { shape: Shape }, /// Teardown has begun, so event-driven work may no longer start. Closed { shape: Shape }, /// The two inverse loopbacks must never coexist. Incompatible { shape: Shape, occupied: Shape, state: &'static str, }, /// The capture sink cannot be removed while a loopback may still reference it. Referenced { shape: Shape, dependency: Shape, state: &'static str, }, } impl std::fmt::Display for LedgerError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { LedgerError::Busy { shape, state } => { write!(f, "the {} slot is {state}, not free", shape.label()) } LedgerError::Poisoned { shape, reason } => { write!(f, "the {} slot is poisoned: {reason}", shape.label()) } LedgerError::InvalidIndex { shape } => write!( f, "pactl reported PA_INVALID_INDEX for the {} module", shape.label() ), LedgerError::Closed { shape } => write!( f, "the module ledger is closing; refusing new work on the {} slot", shape.label() ), LedgerError::Incompatible { shape, occupied, state, } => write!( f, "cannot load the {} while the incompatible {} slot is {state}", shape.label(), occupied.label() ), LedgerError::Referenced { shape, dependency, state, } => write!( f, "cannot unload the {} while the dependent {} slot is {state}", shape.label(), dependency.label() ), } } } impl std::error::Error for LedgerError {} /// The outcome of an attempted unload, as the caller observed it. #[derive(Debug, Clone, PartialEq, Eq)] pub enum UnloadOutcome { /// The server acknowledged the unload. Confirmed, /// The unload may or may not have happened: the command failed, was killed, /// or its result could not be read. Uncertain(String), } /// What a reconciliation found. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Resolution { /// Exactly one module answers the question. Adopt it. Adopt(Fingerprint), /// No module answers it: the load never happened, or the unload did. Nothing, /// More than one module answers it. Refuse to guess — see [`Resolution`]'s use /// in [`ModuleLedger::apply`], which poisons the slot rather than picking. Conflict(usize), } /// The ledger proper: one slot per [`Shape`], plus the permit counter. /// /// Held behind an `Arc` because the permit needs a way back to it from [`Drop`], /// which is what makes a cancelled load impossible to lose. pub struct ModuleLedger { slots: Mutex>, next_permit: AtomicU64, closed: AtomicBool, operations: OperationCoordinator, } /// Registers module mutations before they can be detached from their caller. /// /// `spawn_blocking` work continues when the awaiting future is cancelled. The /// count is therefore registered synchronously, while the slot lock is held, and /// is released only by the affine permit's `Drop`. Teardown closes the ledger and /// waits on this condition variable before it asks the server what exists. struct OperationCoordinator { active: Mutex, idle: Condvar, server: Mutex<()>, } impl OperationCoordinator { fn new() -> Self { Self { active: Mutex::new(0), idle: Condvar::new(), server: Mutex::new(()), } } fn register(&self) { let mut active = self.active.lock().unwrap_or_else(|e| e.into_inner()); *active = active .checked_add(1) .expect("module operation count overflow"); } fn finish(&self) { let mut active = self.active.lock().unwrap_or_else(|e| e.into_inner()); *active = active .checked_sub(1) .expect("module operation count underflow"); if *active == 0 { self.idle.notify_all(); } } fn wait_idle(&self) { let mut active = self.active.lock().unwrap_or_else(|e| e.into_inner()); while *active != 0 { active = self.idle.wait(active).unwrap_or_else(|e| e.into_inner()); } } fn run(&self, operation: impl FnOnce() -> T) -> T { let _serial = self.server.lock().unwrap_or_else(|e| e.into_inner()); operation() } } impl ModuleLedger { pub fn new() -> Arc { Arc::new(Self { slots: Mutex::new(BTreeMap::new()), next_permit: AtomicU64::new(1), closed: AtomicBool::new(false), operations: OperationCoordinator::new(), }) } /// Stop event-driven mutations and wait until every already-registered load /// or unload has settled its affine permit. /// /// The slots-lock hand-off closes the only race that matters: an operation /// that saw `closed == false` has registered itself before this method can /// pass the hand-off and begin waiting. pub(super) fn close(&self) { self.closed.store(true, Ordering::SeqCst); drop(self.slots.lock().unwrap_or_else(|e| e.into_inner())); } pub(super) fn wait_for_operations(&self) { self.operations.wait_idle(); } pub(super) fn close_and_wait(&self) { self.close(); self.wait_for_operations(); } /// Serialize one server mutation or observation with every other such action. pub(super) fn with_server_operation(&self, operation: impl FnOnce() -> T) -> T { self.operations.run(operation) } /// The current state of one slot. Absent keys read as [`SlotState::Vacant`]. /// /// Test-only: production code never needs to look a slot up, because every /// decision that depends on one is made *by* the ledger — `begin_load` /// refuses a busy slot and says which state refused, `begin_unload` returns /// nothing for a slot holding nothing. An accessor callers could branch on /// would invite exactly the check-then-act races the permit removes. #[cfg(test)] pub fn state(&self, shape: Shape) -> SlotState { self.slots .lock() .unwrap() .get(&shape) .cloned() .unwrap_or(SlotState::Vacant) } /// Take permission to load `shape`, moving the slot to /// [`SlotState::Loading`]. /// /// Refuses anything but a `Vacant` slot. In particular an *ambiguous* slot /// refuses, so an unresolved load blocks the next one until it is reconciled. pub fn begin_load( self: &Arc, shape: Shape, pid: u32, token: OwnerToken, ) -> Result { let mut slots = self.slots.lock().unwrap(); if self.closed.load(Ordering::SeqCst) { return Err(LedgerError::Closed { shape }); } let current = slots.get(&shape).cloned().unwrap_or(SlotState::Vacant); match current { SlotState::Vacant => {} SlotState::Poisoned { reason } => { return Err(LedgerError::Poisoned { shape, reason }); } other => { return Err(LedgerError::Busy { shape, state: other.label(), }); } } if let Some(occupied) = inverse_loopback(shape) { let state = slots.get(&occupied).cloned().unwrap_or(SlotState::Vacant); if !matches!(state, SlotState::Vacant) { return Err(LedgerError::Incompatible { shape, occupied, state: state.label(), }); } } let permit = self.next_permit.fetch_add(1, Ordering::Relaxed); self.operations.register(); slots.insert( shape, SlotState::Loading { token: token.clone(), permit, }, ); drop(slots); Ok(LoadPermit { ledger: Arc::clone(self), shape, pid, token, permit, settled: false, }) } /// Move a `Loaded` slot to `Unloading` and hand back what to unload. /// /// `Ok(None)` means confirmed vacancy. Busy, poisoned, closed, and still- /// referenced states are errors so callers cannot mistake uncertainty for /// absence and load an inverse loopback or destroy the capture sink. pub fn begin_unload( self: &Arc, shape: Shape, ) -> Result, LedgerError> { self.begin_unload_inner(shape, false) } /// Teardown-only form of [`ModuleLedger::begin_unload`]. New event work is /// closed by then, but cleanup must still be able to remove tracked modules. pub(super) fn begin_cleanup_unload( self: &Arc, shape: Shape, ) -> Result, LedgerError> { self.begin_unload_inner(shape, true) } fn begin_unload_inner( self: &Arc, shape: Shape, cleanup: bool, ) -> Result, LedgerError> { let mut slots = self.slots.lock().unwrap(); if !cleanup && self.closed.load(Ordering::SeqCst) { return Err(LedgerError::Closed { shape }); } let current = slots.get(&shape).cloned().unwrap_or(SlotState::Vacant); let fp = match current { SlotState::Vacant => return Ok(None), SlotState::Loaded { fp } => fp, SlotState::Poisoned { reason } => { return Err(LedgerError::Poisoned { shape, reason }); } other => { return Err(LedgerError::Busy { shape, state: other.label(), }); } }; if shape == Shape::LegacyCaptureSink { for dependency in [Shape::LoopbackIntoCapture, Shape::LoopbackOutOfCapture] { let state = slots.get(&dependency).cloned().unwrap_or(SlotState::Vacant); if !matches!(state, SlotState::Vacant) { return Err(LedgerError::Referenced { shape, dependency, state: state.label(), }); } } } let permit = self.next_permit.fetch_add(1, Ordering::Relaxed); self.operations.register(); slots.insert( shape, SlotState::Unloading { fp: fp.clone(), permit, }, ); Ok(Some(UnloadPermit { ledger: Arc::clone(self), shape, fp, permit, settled: false, })) } /// Every slot currently holding a module, in [`Shape`] declaration order — /// which is unload order: the loopbacks that reference the capture sink come /// before the sink itself. pub fn loaded(&self) -> Vec { self.slots .lock() .unwrap() .values() .filter_map(|state| match state { SlotState::Loaded { fp } => Some(fp.clone()), _ => None, }) .collect() } /// Every question outstanding against the server, in shape order. pub fn pending(&self) -> Vec<(Shape, Reconcile)> { self.slots .lock() .unwrap() .iter() .filter_map(|(shape, state)| match state { SlotState::Ambiguous(r) => Some((*shape, r.clone())), _ => None, }) .collect() } /// Is every slot in a state we can explain? False while an operation is in /// flight, its outcome is ambiguous, or a conflict poisoned the slot. pub fn is_settled(&self) -> bool { self.slots .lock() .unwrap() .values() .all(|state| matches!(state, SlotState::Vacant | SlotState::Loaded { .. })) } /// Has teardown removed every module rather than merely accounted for it? pub fn is_clean(&self) -> bool { self.slots .lock() .unwrap() .values() .all(|state| matches!(state, SlotState::Vacant)) } /// A stable snapshot used to stop cleanup when another pass made no progress. pub(super) fn snapshot(&self) -> BTreeMap { self.slots.lock().unwrap().clone() } /// Apply a reconciliation result to an ambiguous slot. /// /// A slot that is no longer ambiguous is left alone: the answer is stale, and /// overwriting a live state with it would be worse than ignoring it. pub fn apply(&self, shape: Shape, resolution: Resolution) { let mut slots = self.slots.lock().unwrap(); if !matches!(slots.get(&shape), Some(SlotState::Ambiguous(_))) { return; } let next = match resolution { Resolution::Adopt(fp) => { tracing::info!( shape = shape.label(), module = fp.id, "audio ledger: reconciled — adopting the module the server actually has" ); SlotState::Loaded { fp } } Resolution::Nothing => { tracing::info!( shape = shape.label(), "audio ledger: reconciled — the server has no such module" ); SlotState::Vacant } Resolution::Conflict(n) => { let reason = format!("{n} modules answer to this slot's token; refusing to choose one"); tracing::error!(shape = shape.label(), "audio ledger: {reason}"); SlotState::Poisoned { reason } } }; slots.insert(shape, next); } } fn inverse_loopback(shape: Shape) -> Option { match shape { Shape::LoopbackIntoCapture => Some(Shape::LoopbackOutOfCapture), Shape::LoopbackOutOfCapture => Some(Shape::LoopbackIntoCapture), Shape::LegacyCaptureSink => None, } } /// Permission to perform exactly one load, which must be settled by value. /// /// Dropping it unsettled is not an error — it is the *reporting* path for a /// cancelled or panicking load, and it marks the slot ambiguous so the module the /// server may have created is looked for rather than forgotten. #[must_use = "an unsettled permit marks the slot ambiguous when dropped"] pub struct LoadPermit { ledger: Arc, shape: Shape, pid: u32, token: OwnerToken, permit: u64, settled: bool, } impl LoadPermit { /// Run the server-facing part of this load in the same serialized domain as /// unload verification and reconciliation. pub fn with_server_operation(&self, operation: impl FnOnce() -> T) -> T { self.ledger.with_server_operation(operation) } /// Record that the server created the module at `index`. /// /// Fails on [`PA_INVALID_INDEX`], leaving the permit unsettled so that /// dropping it marks the slot ambiguous — a sentinel where an index belongs /// means the load's outcome is precisely what we do not know. pub fn commit(mut self, index: u32) -> Result { if index == PA_INVALID_INDEX { return Err(LedgerError::InvalidIndex { shape: self.shape }); } let fp = expected_fingerprint(self.shape, self.pid, &self.token, index); // A permit outlives its slot's `Loading` state only if something else has // already moved the slot on — in which case this answer is stale and the // live state wins. let mut slots = self.ledger.slots.lock().unwrap(); match slots.get(&self.shape) { Some(SlotState::Loading { permit, .. }) if *permit == self.permit => { slots.insert(self.shape, SlotState::Loaded { fp: fp.clone() }); } other => { let state = other.map(SlotState::label).unwrap_or("vacant"); tracing::warn!( shape = self.shape.label(), module = index, "audio ledger: a load committed against a slot that is now {state}; \ leaving the live state alone" ); } } drop(slots); self.settled = true; Ok(fp) } /// Record that the server definitively created nothing. /// /// Only for a load that failed *cleanly* — `pactl` exiting non-zero of its own /// accord, having reported the server's refusal. A killed or timed-out load is /// not this: drop the permit instead and let the slot go ambiguous. pub fn abandon(mut self) { let mut slots = self.ledger.slots.lock().unwrap(); if let Some(SlotState::Loading { permit, .. }) = slots.get(&self.shape) && *permit == self.permit { slots.insert(self.shape, SlotState::Vacant); } drop(slots); self.settled = true; } } impl Drop for LoadPermit { fn drop(&mut self) { if !self.settled { let mut slots = self.ledger.slots.lock().unwrap(); // Only claim the slot if it is still *our* load. Anything else already // moved past this permit. if let Some(SlotState::Loading { permit, .. }) = slots.get(&self.shape) && *permit == self.permit { tracing::warn!( shape = self.shape.label(), "audio ledger: a load was cancelled before its outcome was known; \ the slot needs reconciling" ); slots.insert( self.shape, SlotState::Ambiguous(Reconcile::Load { shape: self.shape, pid: self.pid, token: self.token.clone(), }), ); } } self.ledger.operations.finish(); } } /// Permission to perform exactly one unload, which must be settled by value. /// /// Like [`LoadPermit`], this is affine. Cancellation drops it unsettled, retaining /// the full fingerprint as a reconciliation question instead of leaving the slot /// stuck in an invisible `Unloading` state. #[must_use = "an unsettled permit marks the unload ambiguous when dropped"] pub struct UnloadPermit { ledger: Arc, shape: Shape, fp: Fingerprint, permit: u64, settled: bool, } impl UnloadPermit { pub fn fingerprint(&self) -> &Fingerprint { &self.fp } pub fn with_server_operation(&self, operation: impl FnOnce() -> T) -> T { self.ledger.with_server_operation(operation) } /// Record how the server operation ended. An uncertain answer retains the /// fingerprint and makes the next action reconcile it before doing anything. pub fn finish(mut self, outcome: UnloadOutcome) { let mut slots = self.ledger.slots.lock().unwrap(); let Some(SlotState::Unloading { fp, permit }) = slots.get(&self.shape).cloned() else { self.settled = true; return; }; if permit != self.permit { self.settled = true; return; } let next = match outcome { UnloadOutcome::Confirmed => SlotState::Vacant, UnloadOutcome::Uncertain(why) => { tracing::warn!( shape = fp.shape.label(), module = fp.id, "audio ledger: unload outcome uncertain ({why}); slot needs reconciling" ); SlotState::Ambiguous(Reconcile::Unload { fp }) } }; slots.insert(self.shape, next); drop(slots); self.settled = true; } } impl Drop for UnloadPermit { fn drop(&mut self) { if !self.settled { let mut slots = self.ledger.slots.lock().unwrap(); if let Some(SlotState::Unloading { permit, .. }) = slots.get(&self.shape) && *permit == self.permit { tracing::warn!( shape = self.shape.label(), module = self.fp.id, "audio ledger: an unload was cancelled before its outcome was known; \ the slot needs reconciling" ); slots.insert( self.shape, SlotState::Ambiguous(Reconcile::Unload { fp: self.fp.clone(), }), ); } } self.ledger.operations.finish(); } } /// The fingerprint a successful load of `shape` for `pid` under `token` must have. /// /// Built from the shape's own renderer, exactly as [`classify`] rebuilds it from /// an observation, so an adopted module and a committed one are the same value. fn expected_fingerprint(shape: Shape, pid: u32, token: &OwnerToken, id: u32) -> Fingerprint { Fingerprint { id, module_name: shape.module_name().to_string(), args: recorded_argument(&shape.render_args(pid, Some(token))), pid, shape, owner: Some(token.clone()), } } /// Answer one outstanding question against a snapshot of the module table. /// /// Pure: the snapshot is the only input, so every branch is testable without a /// Pulse server. pub fn resolve(reconcile: &Reconcile, observations: &[ModuleObservation]) -> Resolution { let matches: Vec = match reconcile { // The nonce is minted per load, so token equality names this attempt and // no other — including a previous load of the same shape by the same pid. Reconcile::Load { shape, token, .. } => observations .iter() .filter_map(classify) .filter(|fp| fp.shape == *shape && fp.owner.as_ref() == Some(token)) .collect(), // A full fingerprint match, not an id: Pulse reuses module indices // verbatim, so "something is at that index" is not "our module is". Reconcile::Unload { fp } => observations .iter() .filter(|obs| fp.still_matches(obs)) .filter_map(classify) .collect(), }; match matches.len() { 0 => Resolution::Nothing, 1 => Resolution::Adopt(matches.into_iter().next().expect("length checked")), n => Resolution::Conflict(n), } } /// Resolve every outstanding question against one snapshot of the module table. /// /// One listing answers all of them, so the slots are reconciled against a single /// server response rather than several that could disagree. /// /// ⚠️ The session is deliberately short-lived — connect, list, disconnect — which /// is `--repair`'s pattern and not the long-lived host session round 19 sketched. /// `repair::introspect` documents why: on a request timeout the binding leaks the /// boxed callback until the context disconnects, which is bounded and harmless for /// a session that ends immediately, and is not acceptable for one held open for /// the life of a share. Reconciliation is rare and off the hot path, so paying a /// connection for it costs nothing that matters. pub async fn reconcile_pending(ledger: &Arc) -> Result<()> { let ledger = Arc::clone(ledger); tokio::task::spawn_blocking(move || reconcile_pending_blocking(&ledger)) .await .context("the Pulse reconciliation task failed to run")? } /// Synchronous reconciliation for [`Routing::drop`](crate::host::audio::Routing). /// The caller closes and quiesces the ledger first; serialization here also makes /// ordinary async reconciliation wait behind any server operation already running. pub(super) fn reconcile_pending_blocking(ledger: &Arc) -> Result<()> { ledger.with_server_operation(|| { let pending = ledger.pending(); if pending.is_empty() { return Ok(()); } tracing::info!( n = pending.len(), "audio ledger: reconciling unresolved module slots against the server" ); let mut session = PulseSession::connect()?; let observations = session .list_modules() .context("could not list Pulse modules to reconcile the audio ledger")?; for (shape, reconcile) in pending { ledger.apply(shape, resolve(&reconcile, &observations)); } Ok(()) }) } #[cfg(test)] mod tests { use super::*; fn token(nonce: u64) -> OwnerToken { OwnerToken { machine: "abc123".to_string(), boot: "def456".to_string(), pid_ns: 4_026_531_836, nonce, } } /// A module observation as the server would report it for one of our loads. fn observed(id: u32, shape: Shape, pid: u32, token: &OwnerToken) -> ModuleObservation { ModuleObservation::new( id, shape.module_name(), &recorded_argument(&shape.render_args(pid, Some(token))), ) } #[test] fn a_committed_load_becomes_loaded() { let ledger = ModuleLedger::new(); let permit = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"); assert!(matches!( ledger.state(Shape::LegacyCaptureSink), SlotState::Loading { .. } )); let fp = permit.commit(9).expect("9 is a real index"); assert_eq!(fp.id, 9); assert_eq!( ledger.state(Shape::LegacyCaptureSink), SlotState::Loaded { fp } ); } #[test] fn dropping_a_permit_unsettled_marks_the_slot_ambiguous() { // The orphan race in one test: a load that is cancelled between spawning // pactl and reading its id must leave the slot asking a question, never // looking empty. let ledger = ModuleLedger::new(); let permit = ledger .begin_load(Shape::LoopbackOutOfCapture, 42, token(7)) .expect("a vacant slot issues a permit"); drop(permit); assert_eq!( ledger.state(Shape::LoopbackOutOfCapture), SlotState::Ambiguous(Reconcile::Load { shape: Shape::LoopbackOutOfCapture, pid: 42, token: token(7), }), "a cancelled load must be remembered as a question, not as a vacancy" ); assert!(!ledger.is_settled()); } #[test] fn abandoning_a_cleanly_failed_load_returns_the_slot_to_vacant() { let ledger = ModuleLedger::new(); ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(7)) .expect("a vacant slot issues a permit") .abandon(); assert_eq!( ledger.state(Shape::LoopbackIntoCapture), SlotState::Vacant, "a load the server refused created nothing to reconcile" ); assert!(ledger.is_settled()); } #[test] fn a_second_permit_is_refused_while_a_load_is_in_flight() { let ledger = ModuleLedger::new(); let _first = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"); let second = ledger.begin_load(Shape::LegacyCaptureSink, 42, token(8)); assert_eq!( second.err(), Some(LedgerError::Busy { shape: Shape::LegacyCaptureSink, state: "loading" }), "two loads for one slot must not both be permitted" ); } #[test] fn an_ambiguous_slot_refuses_the_next_load_until_it_is_reconciled() { // Why this matters: two sinks may share a node.name, and pulsesrc attaches // to the older one — so loading over unresolved debris silently steals the // next session's capture. let ledger = ModuleLedger::new(); drop( ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"), ); assert_eq!( ledger .begin_load(Shape::LegacyCaptureSink, 42, token(8)) .err(), Some(LedgerError::Busy { shape: Shape::LegacyCaptureSink, state: "ambiguous" }) ); // Reconciling to "nothing was created" frees it again. ledger.apply(Shape::LegacyCaptureSink, Resolution::Nothing); assert!( ledger .begin_load(Shape::LegacyCaptureSink, 42, token(9)) .is_ok() ); } #[test] fn an_invalid_index_is_not_adopted_and_leaves_a_question_behind() { let ledger = ModuleLedger::new(); let permit = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"); assert_eq!( permit.commit(PA_INVALID_INDEX).err(), Some(LedgerError::InvalidIndex { shape: Shape::LegacyCaptureSink }) ); assert!( matches!( ledger.state(Shape::LegacyCaptureSink), SlotState::Ambiguous(_) ), "a sentinel where an index belongs is exactly the unknown outcome" ); } #[test] fn an_uncertain_unload_is_not_forgotten() { let ledger = ModuleLedger::new(); let fp = ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(7)) .expect("a vacant slot issues a permit") .commit(3) .expect("3 is a real index"); let permit = ledger .begin_unload(Shape::LoopbackIntoCapture) .expect("the ledger accepts the unload") .expect("the loaded slot issues a permit"); assert_eq!(permit.fingerprint(), &fp); permit.finish(UnloadOutcome::Uncertain("pactl was killed".to_string())); assert_eq!( ledger.state(Shape::LoopbackIntoCapture), SlotState::Ambiguous(Reconcile::Unload { fp }), "an unload whose outcome is unknown must keep naming the module" ); } #[test] fn a_confirmed_unload_empties_the_slot() { let ledger = ModuleLedger::new(); ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(7)) .expect("a vacant slot issues a permit") .commit(3) .expect("3 is a real index"); ledger .begin_unload(Shape::LoopbackIntoCapture) .expect("the ledger accepts the unload") .expect("the loaded slot issues a permit") .finish(UnloadOutcome::Confirmed); assert_eq!(ledger.state(Shape::LoopbackIntoCapture), SlotState::Vacant); assert!(ledger.is_settled()); assert!(ledger.is_clean()); } #[test] fn dropping_an_unload_permit_marks_the_slot_ambiguous() { let ledger = ModuleLedger::new(); let fp = ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(7)) .expect("a vacant slot issues a permit") .commit(3) .expect("3 is a real index"); let permit = ledger .begin_unload(Shape::LoopbackIntoCapture) .expect("the ledger accepts the unload") .expect("the loaded slot issues a permit"); assert!(matches!( ledger.state(Shape::LoopbackIntoCapture), SlotState::Unloading { .. } )); assert!(!ledger.is_settled(), "in-flight unloads are not settled"); drop(permit); assert_eq!( ledger.state(Shape::LoopbackIntoCapture), SlotState::Ambiguous(Reconcile::Unload { fp }) ); } #[test] fn inverse_loopbacks_cannot_coexist_or_cross_an_unknown_boundary() { let ledger = ModuleLedger::new(); ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(7)) .expect("the first loopback may load") .commit(3) .expect("3 is a real index"); assert_eq!( ledger .begin_load(Shape::LoopbackOutOfCapture, 42, token(8)) .err(), Some(LedgerError::Incompatible { shape: Shape::LoopbackOutOfCapture, occupied: Shape::LoopbackIntoCapture, state: "loaded" }) ); let unload = ledger .begin_unload(Shape::LoopbackIntoCapture) .expect("the ledger accepts the unload") .expect("the loaded slot issues a permit"); drop(unload); assert_eq!( ledger .begin_load(Shape::LoopbackOutOfCapture, 42, token(9)) .err(), Some(LedgerError::Incompatible { shape: Shape::LoopbackOutOfCapture, occupied: Shape::LoopbackIntoCapture, state: "ambiguous" }), "an unconfirmed absence must block the inverse loopback" ); } #[test] fn capture_sink_unload_waits_for_every_dependent_slot_to_be_vacant() { let ledger = ModuleLedger::new(); ledger .begin_load(Shape::LegacyCaptureSink, 42, token(1)) .expect("the sink may load") .commit(1) .expect("1 is a real index"); ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(2)) .expect("the mirror may load") .commit(2) .expect("2 is a real index"); assert_eq!( ledger.begin_unload(Shape::LegacyCaptureSink).err(), Some(LedgerError::Referenced { shape: Shape::LegacyCaptureSink, dependency: Shape::LoopbackIntoCapture, state: "loaded" }) ); ledger .begin_unload(Shape::LoopbackIntoCapture) .expect("the ledger accepts the unload") .expect("the mirror issues an unload permit") .finish(UnloadOutcome::Confirmed); assert!( ledger .begin_unload(Shape::LegacyCaptureSink) .expect("the sink is no longer referenced") .is_some() ); } #[test] fn closing_waits_for_a_previously_registered_operation() { use std::sync::mpsc; use std::time::Duration; let ledger = ModuleLedger::new(); let permit = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("the operation registers before it can be detached"); let (entered_tx, entered_rx) = mpsc::channel(); let (done_tx, done_rx) = mpsc::channel(); let waiter_ledger = Arc::clone(&ledger); let waiter = std::thread::spawn(move || { entered_tx.send(()).unwrap(); waiter_ledger.close_and_wait(); done_tx.send(()).unwrap(); }); entered_rx.recv().unwrap(); assert!( matches!( done_rx.recv_timeout(Duration::from_millis(30)), Err(mpsc::RecvTimeoutError::Timeout) ), "the close barrier must not pass a live affine permit" ); permit.abandon(); done_rx .recv_timeout(Duration::from_secs(1)) .expect("settling the operation releases teardown"); waiter.join().unwrap(); assert!(ledger.is_clean()); } #[test] fn a_closed_ledger_refuses_event_work_but_allows_cleanup() { let ledger = ModuleLedger::new(); ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("the sink may load before close") .commit(3) .expect("3 is a real index"); ledger.close_and_wait(); assert_eq!( ledger .begin_load(Shape::LoopbackIntoCapture, 42, token(8)) .err(), Some(LedgerError::Closed { shape: Shape::LoopbackIntoCapture }) ); assert_eq!( ledger.begin_unload(Shape::LegacyCaptureSink).err(), Some(LedgerError::Closed { shape: Shape::LegacyCaptureSink }) ); assert!( ledger .begin_cleanup_unload(Shape::LegacyCaptureSink) .expect("teardown retains its cleanup authority") .is_some() ); } #[test] fn begin_unload_distinguishes_vacant_from_busy() { let ledger = ModuleLedger::new(); assert!( ledger .begin_unload(Shape::LegacyCaptureSink) .expect("vacancy is not an error") .is_none() ); let _permit = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"); assert_eq!( ledger.begin_unload(Shape::LegacyCaptureSink).err(), Some(LedgerError::Busy { shape: Shape::LegacyCaptureSink, state: "loading" }), "an in-flight load must not look like a confirmed vacancy" ); } #[test] fn resolve_adopts_the_one_module_carrying_our_token() { let mine = token(7); let reconcile = Reconcile::Load { shape: Shape::LegacyCaptureSink, pid: 42, token: mine.clone(), }; let observations = vec![ // A previous load by the same pid and shape: same everything but the // per-load nonce, which is the whole reason the nonce exists. observed(4, Shape::LegacyCaptureSink, 42, &token(6)), observed(5, Shape::LegacyCaptureSink, 42, &mine), // Somebody else's module entirely. ModuleObservation::new(6, "module-null-sink", "sink_name=other"), ]; assert_eq!( resolve(&reconcile, &observations), Resolution::Adopt(expected_fingerprint(Shape::LegacyCaptureSink, 42, &mine, 5)) ); } #[test] fn resolve_reports_nothing_when_the_server_never_created_it() { let reconcile = Reconcile::Load { shape: Shape::LegacyCaptureSink, pid: 42, token: token(7), }; let observations = vec![observed(4, Shape::LegacyCaptureSink, 42, &token(6))]; assert_eq!(resolve(&reconcile, &observations), Resolution::Nothing); } #[test] fn resolve_fails_closed_when_more_than_one_module_answers() { let mine = token(7); let reconcile = Reconcile::Load { shape: Shape::LegacyCaptureSink, pid: 42, token: mine.clone(), }; // Two modules carrying the same token should be impossible. If it ever // happens, guessing which to keep is how a live sink gets unloaded. let observations = vec![ observed(4, Shape::LegacyCaptureSink, 42, &mine), observed(5, Shape::LegacyCaptureSink, 42, &mine), ]; assert_eq!(resolve(&reconcile, &observations), Resolution::Conflict(2)); } #[test] fn a_conflict_poisons_the_slot_and_it_stays_poisoned() { let ledger = ModuleLedger::new(); drop( ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit"), ); ledger.apply(Shape::LegacyCaptureSink, Resolution::Conflict(2)); let SlotState::Poisoned { reason } = ledger.state(Shape::LegacyCaptureSink) else { panic!("a conflict must poison the slot"); }; assert!(!ledger.is_settled()); assert_eq!( ledger .begin_load(Shape::LegacyCaptureSink, 42, token(8)) .err(), Some(LedgerError::Poisoned { shape: Shape::LegacyCaptureSink, reason }), "poisoning is terminal for the session" ); } #[test] fn resolve_finds_a_module_an_uncertain_unload_left_behind() { let mine = token(7); let fp = expected_fingerprint(Shape::LoopbackIntoCapture, 42, &mine, 5); let observations = vec![observed(5, Shape::LoopbackIntoCapture, 42, &mine)]; assert_eq!( resolve(&Reconcile::Unload { fp: fp.clone() }, &observations), Resolution::Adopt(fp.clone()), "still present means the unload did not happen" ); assert_eq!( resolve(&Reconcile::Unload { fp }, &[]), Resolution::Nothing, "absent means it did" ); } #[test] fn an_unload_reconcile_ignores_a_stranger_at_the_same_index() { // Pulse reuses module indices verbatim, so "something is at index 5" must // not be read as "our module is still at index 5". // // The stranger has to be a module that *classifies* — another pixelpass // host's canonical loopback — or this gate is vacuous: a comparator using // the id alone would still return `Nothing` for junk, because junk is // discarded by `classify` regardless of how the match was made. let fp = expected_fingerprint(Shape::LoopbackIntoCapture, 42, &token(7), 5); let another_hosts = observed(5, Shape::LoopbackIntoCapture, 99, &token(3)); assert_eq!( resolve( &Reconcile::Unload { fp: fp.clone() }, std::slice::from_ref(&another_hosts) ), Resolution::Nothing, "another host's module at our old index is not our module" ); // And plain junk at that index is ignored too. let junk = ModuleObservation::new(5, "module-loopback", "source=some_mic sink=theirs"); assert_eq!( resolve(&Reconcile::Unload { fp }, std::slice::from_ref(&junk)), Resolution::Nothing ); } #[test] fn a_stale_answer_does_not_overwrite_a_live_slot() { let ledger = ModuleLedger::new(); // Nothing is ambiguous, so an answer arriving late is meaningless. let fp = ledger .begin_load(Shape::LegacyCaptureSink, 42, token(7)) .expect("a vacant slot issues a permit") .commit(3) .expect("3 is a real index"); ledger.apply(Shape::LegacyCaptureSink, Resolution::Nothing); assert_eq!( ledger.state(Shape::LegacyCaptureSink), SlotState::Loaded { fp }, "a live state must win over a stale reconciliation" ); } #[test] fn loaded_lists_a_loopback_before_the_capture_sink() { // Declaration order is unload order: the loopbacks that reference the // capture sink must come before the sink itself. let ledger = ModuleLedger::new(); for (shape, id) in [ (Shape::LegacyCaptureSink, 1), (Shape::LoopbackOutOfCapture, 3), ] { ledger .begin_load(shape, 42, token(u64::from(id))) .expect("a vacant slot issues a permit") .commit(id) .expect("a real index"); } let order: Vec = ledger.loaded().into_iter().map(|fp| fp.shape).collect(); assert_eq!( order, vec![Shape::LoopbackOutOfCapture, Shape::LegacyCaptureSink] ); assert_eq!( order.last(), Some(&Shape::LegacyCaptureSink), "the sink must be unloaded after everything that references it" ); } }