fix(audio): make module teardown cancellation-safe
This commit is contained in:
+412
-130
@@ -33,33 +33,27 @@
|
|||||||
use anyhow::{Context, Result, bail};
|
use anyhow::{Context, Result, bail};
|
||||||
use std::cell::RefCell;
|
use std::cell::RefCell;
|
||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
use std::process::Command;
|
use std::io::{self, Read};
|
||||||
|
use std::process::{Child, Command, ExitStatus, Stdio};
|
||||||
use std::rc::Rc;
|
use std::rc::Rc;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::thread::JoinHandle;
|
use std::thread::JoinHandle;
|
||||||
use std::time::Duration;
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
use crate::cli::HostOpts;
|
use crate::cli::HostOpts;
|
||||||
use crate::host::ledger::{self, LedgerError, ModuleLedger, UnloadOutcome};
|
use crate::host::ledger::{self, LedgerError, ModuleLedger, UnloadOutcome};
|
||||||
use crate::repair::plan::{self as repair_plan, Fingerprint, Shape};
|
use crate::repair::plan::{self as repair_plan, Fingerprint, Shape};
|
||||||
|
|
||||||
/// How long a single `pactl load-module` / `unload-module` may take.
|
/// How long a `pactl load-module` worker may run before it is killed and reaped.
|
||||||
///
|
///
|
||||||
/// A bound, not a calibration: a local Pulse socket answers in milliseconds, and
|
/// A bound, not a calibration: a local Pulse socket answers in milliseconds, and
|
||||||
/// this exists only so a wedged server cannot hang teardown forever. It is
|
/// this exists only so a wedged server cannot hang teardown forever. It is
|
||||||
/// deliberately generous because exceeding it is no longer destructive — the
|
/// deliberately generous because exceeding it is no longer destructive: the
|
||||||
/// ledger records the attempt, and reconciliation finds whatever the server
|
/// ledger records the attempt, and reconciliation finds whatever the server
|
||||||
/// actually did.
|
/// actually did. Unloads use PulseSession's independently bounded native
|
||||||
|
/// connect/list/unload requests instead of a second pactl connection.
|
||||||
const PACTL_BUDGET: Duration = Duration::from_secs(5);
|
const PACTL_BUDGET: Duration = Duration::from_secs(5);
|
||||||
|
const PACTL_REAP_BUDGET: Duration = Duration::from_secs(1);
|
||||||
/// How many reconcile-then-unload rounds teardown runs.
|
|
||||||
///
|
|
||||||
/// Two, because one round can create exactly one new question: an ambiguous load
|
|
||||||
/// resolves to a module that then needs unloading, and an uncertain unload
|
|
||||||
/// resolves to a module that is either gone or still there. A second round
|
|
||||||
/// settles either. Anything still unresolved after that is left to `--repair`
|
|
||||||
/// rather than looped over.
|
|
||||||
const TEARDOWN_ROUNDS: usize = 2;
|
|
||||||
|
|
||||||
/// Owns the pactl-loaded modules plus, when filtering is active, the
|
/// Owns the pactl-loaded modules plus, when filtering is active, the
|
||||||
/// libpipewire stream-router thread. Drop unloads modules as a backstop;
|
/// libpipewire stream-router thread. Drop unloads modules as a backstop;
|
||||||
@@ -86,6 +80,17 @@ impl Routing {
|
|||||||
let pid = std::process::id();
|
let pid = std::process::id();
|
||||||
let sink_name = repair_plan::sink_name_for(pid);
|
let sink_name = repair_plan::sink_name_for(pid);
|
||||||
let ledger = ModuleLedger::new();
|
let ledger = ModuleLedger::new();
|
||||||
|
// Construct the owner before the first mutation. Any error or cancellation
|
||||||
|
// below now drops a real `Routing`, whose backstop closes, quiesces,
|
||||||
|
// reconciles, and unloads this ledger. Previously the owner did not exist
|
||||||
|
// until both initial modules had loaded, so constructor failure leaked
|
||||||
|
// everything loaded up to that point.
|
||||||
|
let mut routing = Self {
|
||||||
|
ledger: Arc::clone(&ledger),
|
||||||
|
sink_name: sink_name.clone(),
|
||||||
|
stream_router: None,
|
||||||
|
event_task: None,
|
||||||
|
};
|
||||||
|
|
||||||
// Every module this host loads carries an ownership token, minted per
|
// Every module this host loads carries an ownership token, minted per
|
||||||
// load, so `--repair` can tell whose pid the name refers to instead of
|
// load, so `--repair` can tell whose pid the name refers to instead of
|
||||||
@@ -119,13 +124,6 @@ impl Routing {
|
|||||||
"audio routing: null-sink ready (loopback skipped in strict app mode)"
|
"audio routing: null-sink ready (loopback skipped in strict app mode)"
|
||||||
);
|
);
|
||||||
|
|
||||||
let mut routing = Self {
|
|
||||||
ledger: Arc::clone(&ledger),
|
|
||||||
sink_name: sink_name.clone(),
|
|
||||||
stream_router: None,
|
|
||||||
event_task: None,
|
|
||||||
};
|
|
||||||
|
|
||||||
if let Some(app) = &opts.app {
|
if let Some(app) = &opts.app {
|
||||||
let (router, mut event_rx) = StreamRouter::spawn(app.clone(), sink_name.clone())?;
|
let (router, mut event_rx) = StreamRouter::spawn(app.clone(), sink_name.clone())?;
|
||||||
let ledger_for_task = Arc::clone(&ledger);
|
let ledger_for_task = Arc::clone(&ledger);
|
||||||
@@ -138,6 +136,7 @@ impl Routing {
|
|||||||
tracing::info!(
|
tracing::info!(
|
||||||
"audio routing: first stream routed → unloading default-sink loopback"
|
"audio routing: first stream routed → unloading default-sink loopback"
|
||||||
);
|
);
|
||||||
|
let mirror_absent =
|
||||||
unload_module(&ledger_for_task, Shape::LoopbackIntoCapture).await;
|
unload_module(&ledger_for_task, Shape::LoopbackIntoCapture).await;
|
||||||
// Mirror the routed app back to the sharer's own
|
// Mirror the routed app back to the sharer's own
|
||||||
// speakers so they hear the content they're sharing.
|
// speakers so they hear the content they're sharing.
|
||||||
@@ -146,7 +145,15 @@ impl Routing {
|
|||||||
// sourced from the null-sink monitor — the chosen app
|
// sourced from the null-sink monitor — the chosen app
|
||||||
// only, never the desktop/call — so it can't echo into
|
// only, never the desktop/call — so it can't echo into
|
||||||
// the capture.
|
// the capture.
|
||||||
ensure_loaded(&ledger_for_task, Shape::LoopbackOutOfCapture, pid).await;
|
if mirror_absent {
|
||||||
|
ensure_loaded(&ledger_for_task, Shape::LoopbackOutOfCapture, pid)
|
||||||
|
.await;
|
||||||
|
} else {
|
||||||
|
tracing::warn!(
|
||||||
|
"audio routing: default-sink mirror absence was not confirmed; \
|
||||||
|
refusing to load the inverse local monitor"
|
||||||
|
);
|
||||||
|
}
|
||||||
// Tell the front-end the chosen app's audio is live.
|
// Tell the front-end the chosen app's audio is live.
|
||||||
output::emit(output::Event::AppAudio {
|
output::emit(output::Event::AppAudio {
|
||||||
state: AppAudioState::Routed,
|
state: AppAudioState::Routed,
|
||||||
@@ -161,6 +168,7 @@ impl Routing {
|
|||||||
// The shared app is gone, so its null-sink is silent:
|
// The shared app is gone, so its null-sink is silent:
|
||||||
// stop mirroring it to the sharer's speakers. Re-loads
|
// stop mirroring it to the sharer's speakers. Re-loads
|
||||||
// on the next FirstRoutedStream if the app resumes.
|
// on the next FirstRoutedStream if the app resumes.
|
||||||
|
let local_monitor_absent =
|
||||||
unload_module(&ledger_for_task, Shape::LoopbackOutOfCapture).await;
|
unload_module(&ledger_for_task, Shape::LoopbackOutOfCapture).await;
|
||||||
if strict {
|
if strict {
|
||||||
// Strict mode: do NOT restore the whole-desktop
|
// Strict mode: do NOT restore the whole-desktop
|
||||||
@@ -172,6 +180,13 @@ impl Routing {
|
|||||||
);
|
);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
if !local_monitor_absent {
|
||||||
|
tracing::warn!(
|
||||||
|
"audio routing: local-monitor absence was not confirmed; \
|
||||||
|
refusing to restore the inverse default-sink mirror"
|
||||||
|
);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
// Best-effort mode: restore the default-sink loopback
|
// Best-effort mode: restore the default-sink loopback
|
||||||
// so the viewer hears system audio again instead of
|
// so the viewer hears system audio again instead of
|
||||||
// silence. Already loaded is not an error — the ledger
|
// silence. Already loaded is not an error — the ledger
|
||||||
@@ -216,6 +231,10 @@ impl Routing {
|
|||||||
/// module — but only a path that then reconciles can actually clean it up.
|
/// module — but only a path that then reconciles can actually clean it up.
|
||||||
/// `Drop` cannot await, which is why it is the narrower backstop.
|
/// `Drop` cannot await, which is why it is the narrower backstop.
|
||||||
pub async fn shutdown(mut self) {
|
pub async fn shutdown(mut self) {
|
||||||
|
// Closing is synchronous and happens first: after this point the event
|
||||||
|
// task cannot register another mutation even if it receives one last
|
||||||
|
// router event while shutdown is in progress.
|
||||||
|
self.ledger.close();
|
||||||
if let Some(router) = self.stream_router.take() {
|
if let Some(router) = self.stream_router.take() {
|
||||||
// ⚠️ Still an unbounded join: a wedged PipeWire thread parks this
|
// ⚠️ Still an unbounded join: a wedged PipeWire thread parks this
|
||||||
// task indefinitely. That is the pre-existing defect S3b exists for.
|
// task indefinitely. That is the pre-existing defect S3b exists for.
|
||||||
@@ -240,22 +259,23 @@ impl Routing {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
for _ in 0..TEARDOWN_ROUNDS {
|
// A cancelled `spawn_blocking` await detaches its worker. Every worker is
|
||||||
if let Err(e) = ledger::reconcile_pending(&self.ledger).await {
|
// registered before it can be spawned, so this is a real ordering
|
||||||
tracing::warn!("audio routing: could not reconcile the module ledger: {e:#}");
|
// boundary: reconciliation cannot overtake late module creation/removal.
|
||||||
}
|
let ledger_for_wait = Arc::clone(&self.ledger);
|
||||||
let loaded = self.ledger.loaded();
|
if let Err(e) = tokio::task::spawn_blocking(move || {
|
||||||
if loaded.is_empty() {
|
ledger_for_wait.wait_for_operations();
|
||||||
break;
|
})
|
||||||
}
|
.await
|
||||||
for fp in loaded {
|
{
|
||||||
unload_module(&self.ledger, fp.shape).await;
|
tracing::warn!("audio routing: module-operation wait task failed: {e}");
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
cleanup_modules(&self.ledger).await;
|
||||||
|
|
||||||
if !self.ledger.is_settled() {
|
if !self.ledger.is_clean() {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
"audio routing: some audio modules could not be accounted for; \
|
settled = self.ledger.is_settled(),
|
||||||
|
"audio routing: some audio modules could not be removed safely; \
|
||||||
`pixelpass --repair` will clean up anything left behind"
|
`pixelpass --repair` will clean up anything left behind"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
@@ -270,22 +290,19 @@ impl Drop for Routing {
|
|||||||
///
|
///
|
||||||
/// After a completed `shutdown` the ledger holds nothing and this does nothing.
|
/// After a completed `shutdown` the ledger holds nothing and this does nothing.
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
|
self.ledger.close();
|
||||||
if let Some(router) = self.stream_router.take() {
|
if let Some(router) = self.stream_router.take() {
|
||||||
router.shutdown();
|
router.shutdown();
|
||||||
}
|
}
|
||||||
if let Some(task) = self.event_task.take() {
|
if let Some(task) = self.event_task.take() {
|
||||||
task.abort();
|
task.abort();
|
||||||
}
|
}
|
||||||
for fp in self.ledger.loaded() {
|
self.ledger.close_and_wait();
|
||||||
if self.ledger.begin_unload(fp.shape).is_none() {
|
cleanup_modules_blocking(&self.ledger);
|
||||||
continue;
|
if !self.ledger.is_clean() {
|
||||||
}
|
|
||||||
let outcome = blocking_unload(fp.id);
|
|
||||||
self.ledger.finish_unload(fp.shape, outcome);
|
|
||||||
}
|
|
||||||
if !self.ledger.is_settled() {
|
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
"audio routing: torn down without settling the module ledger; \
|
settled = self.ledger.is_settled(),
|
||||||
|
"audio routing: torn down with modules that could not be removed safely; \
|
||||||
run `pixelpass --repair` to clean up anything left behind"
|
run `pixelpass --repair` to clean up anything left behind"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
@@ -408,55 +425,48 @@ async fn load_module(ledger: &Arc<ModuleLedger>, shape: Shape, pid: u32) -> Resu
|
|||||||
.begin_load(shape, pid, owner.clone())
|
.begin_load(shape, pid, owner.clone())
|
||||||
.map_err(anyhow::Error::new)
|
.map_err(anyhow::Error::new)
|
||||||
.with_context(|| format!("cannot load the {} module", shape.label()))?;
|
.with_context(|| format!("cannot load the {} module", shape.label()))?;
|
||||||
|
let args = shape.render_args(pid, Some(&owner));
|
||||||
|
|
||||||
let mut cmd = tokio::process::Command::new("pactl");
|
// The affine permit moves into the blocking worker. Dropping this await does
|
||||||
cmd.arg("load-module")
|
// not cancel `spawn_blocking`; the worker remains registered, owns and reaps
|
||||||
|
// its child, and settles the slot before teardown's quiescence barrier opens.
|
||||||
|
tokio::task::spawn_blocking(move || -> Result<Fingerprint> {
|
||||||
|
let output = permit.with_server_operation(|| {
|
||||||
|
let mut command = Command::new("pactl");
|
||||||
|
command
|
||||||
|
.arg("load-module")
|
||||||
.arg(shape.module_name())
|
.arg(shape.module_name())
|
||||||
.args(shape.render_args(pid, Some(&owner)))
|
.args(args);
|
||||||
.kill_on_drop(true);
|
bounded_output(&mut command, PACTL_BUDGET)
|
||||||
|
});
|
||||||
// Deliberately **not** `select!`ed against a cancellation signal: a completed
|
let output = match output {
|
||||||
// load whose index was then dropped on the floor is precisely the defect the
|
Ok(output) => output,
|
||||||
// ledger exists to prevent. Cancellation here happens by dropping this whole
|
Err(e) => {
|
||||||
// future, and the permit's `Drop` turns that into a question reconciliation
|
// Whether the server was reached is unknown. Leaving the permit
|
||||||
// can answer, rather than into silence.
|
// unsettled makes its Drop create a reconciliation question.
|
||||||
let output = match tokio::time::timeout(PACTL_BUDGET, cmd.output()).await {
|
|
||||||
Ok(Ok(output)) => output,
|
|
||||||
Ok(Err(e)) => {
|
|
||||||
// Spawning or reading failed. Whether the server was ever reached is
|
|
||||||
// not knowable from here, so leave the permit unsettled: an
|
|
||||||
// unnecessary reconcile costs one listing, a missed one costs an
|
|
||||||
// orphan.
|
|
||||||
drop(permit);
|
drop(permit);
|
||||||
return Err(e).context("failed to run pactl load-module");
|
return Err(e).context("failed to run pactl load-module");
|
||||||
}
|
}
|
||||||
Err(_) => {
|
};
|
||||||
|
|
||||||
|
if output.timed_out {
|
||||||
drop(permit);
|
drop(permit);
|
||||||
bail!("pactl load-module did not finish within {PACTL_BUDGET:?}");
|
bail!("pactl load-module did not finish within {PACTL_BUDGET:?}");
|
||||||
}
|
}
|
||||||
};
|
|
||||||
|
|
||||||
if !output.status.success() {
|
if !output.status.success() {
|
||||||
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
|
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
|
||||||
if output.status.code().is_some() {
|
if output.status.code().is_some() {
|
||||||
// pactl exited of its own accord, having reported the server's
|
// pactl exited normally and reported the server's refusal.
|
||||||
// refusal: nothing was created, so there is nothing to reconcile.
|
|
||||||
permit.abandon();
|
permit.abandon();
|
||||||
} else {
|
} else {
|
||||||
// Killed by a signal, which may have arrived *after* the server
|
|
||||||
// created the module.
|
|
||||||
drop(permit);
|
drop(permit);
|
||||||
}
|
}
|
||||||
bail!("pactl load-module failed: {stderr}");
|
bail!("pactl load-module failed: {stderr}");
|
||||||
}
|
}
|
||||||
|
|
||||||
let id_str = String::from_utf8_lossy(&output.stdout).trim().to_string();
|
let id_str = String::from_utf8_lossy(&output.stdout).trim().to_string();
|
||||||
// Genuinely 32-bit, unlike `object.serial`: this is a PulseAudio module
|
// Genuinely 32-bit: this is a Pulse module index, not object.serial.
|
||||||
// index (`pa_module.index`, `uint32_t`), which `pactl unload-module` takes
|
|
||||||
// back verbatim. Do not widen it.
|
|
||||||
let Ok(index) = id_str.parse::<u32>() else {
|
let Ok(index) = id_str.parse::<u32>() else {
|
||||||
// The load may well have succeeded — we simply cannot say which module it
|
|
||||||
// produced, which is exactly what reconciliation is for.
|
|
||||||
drop(permit);
|
drop(permit);
|
||||||
bail!("pactl returned unexpected module ID: {id_str:?}");
|
bail!("pactl returned unexpected module ID: {id_str:?}");
|
||||||
};
|
};
|
||||||
@@ -467,6 +477,9 @@ async fn load_module(ledger: &Arc<ModuleLedger>, shape: Shape, pid: u32) -> Resu
|
|||||||
"audio routing: loaded pactl module"
|
"audio routing: loaded pactl module"
|
||||||
);
|
);
|
||||||
Ok(fp)
|
Ok(fp)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.context("the pactl load worker failed")?
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Load `shape` unless the slot already holds it.
|
/// Load `shape` unless the slot already holds it.
|
||||||
@@ -496,49 +509,273 @@ async fn ensure_loaded(ledger: &Arc<ModuleLedger>, shape: Shape, pid: u32) {
|
|||||||
/// A no-op for a slot holding nothing. An outcome that cannot be confirmed is
|
/// A no-op for a slot holding nothing. An outcome that cannot be confirmed is
|
||||||
/// recorded as uncertain rather than assumed done, so the module keeps being
|
/// recorded as uncertain rather than assumed done, so the module keeps being
|
||||||
/// named until the server is asked about it.
|
/// named until the server is asked about it.
|
||||||
async fn unload_module(ledger: &Arc<ModuleLedger>, shape: Shape) {
|
async fn unload_module(ledger: &Arc<ModuleLedger>, shape: Shape) -> bool {
|
||||||
let Some(fp) = ledger.begin_unload(shape) else {
|
unload_module_inner(ledger, shape, false).await
|
||||||
return;
|
|
||||||
};
|
|
||||||
let mut cmd = tokio::process::Command::new("pactl");
|
|
||||||
cmd.arg("unload-module")
|
|
||||||
.arg(fp.id.to_string())
|
|
||||||
.kill_on_drop(true);
|
|
||||||
let outcome = match tokio::time::timeout(PACTL_BUDGET, cmd.output()).await {
|
|
||||||
Ok(Ok(output)) if output.status.success() => {
|
|
||||||
tracing::info!(module = fp.id, "audio routing: unloaded pactl module");
|
|
||||||
UnloadOutcome::Confirmed
|
|
||||||
}
|
|
||||||
Ok(Ok(output)) => UnloadOutcome::Uncertain(format!(
|
|
||||||
"pactl unload-module exited {}: {}",
|
|
||||||
output.status,
|
|
||||||
String::from_utf8_lossy(&output.stderr).trim()
|
|
||||||
)),
|
|
||||||
Ok(Err(e)) => UnloadOutcome::Uncertain(format!("could not run pactl unload-module: {e}")),
|
|
||||||
Err(_) => {
|
|
||||||
UnloadOutcome::Uncertain(format!("pactl unload-module exceeded {PACTL_BUDGET:?}"))
|
|
||||||
}
|
|
||||||
};
|
|
||||||
ledger.finish_unload(shape, outcome);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The blocking unload [`Drop`] uses, since it has no runtime to await on.
|
async fn unload_module_inner(ledger: &Arc<ModuleLedger>, shape: Shape, cleanup: bool) -> bool {
|
||||||
fn blocking_unload(id: u32) -> UnloadOutcome {
|
let begun = if cleanup {
|
||||||
match Command::new("pactl")
|
ledger.begin_cleanup_unload(shape)
|
||||||
.arg("unload-module")
|
} else {
|
||||||
.arg(id.to_string())
|
ledger.begin_unload(shape)
|
||||||
.output()
|
};
|
||||||
{
|
let permit = match begun {
|
||||||
Ok(output) if output.status.success() => {
|
Ok(Some(permit)) => permit,
|
||||||
tracing::info!(module = id, "audio routing: unloaded pactl module");
|
Ok(None) => return true,
|
||||||
UnloadOutcome::Confirmed
|
Err(e) => {
|
||||||
|
tracing::warn!(
|
||||||
|
shape = shape.label(),
|
||||||
|
"audio routing: refusing module unload: {e}"
|
||||||
|
);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
match tokio::task::spawn_blocking(move || finish_verified_unload(permit)).await {
|
||||||
|
Ok(confirmed_absent) => confirmed_absent,
|
||||||
|
Err(e) => {
|
||||||
|
// A panicking worker drops its affine permit and therefore leaves an
|
||||||
|
// unload reconciliation question behind.
|
||||||
|
tracing::warn!(
|
||||||
|
shape = shape.label(),
|
||||||
|
"audio routing: unload worker failed: {e}"
|
||||||
|
);
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn finish_verified_unload(permit: ledger::UnloadPermit) -> bool {
|
||||||
|
let fp = permit.fingerprint().clone();
|
||||||
|
let result = permit.with_server_operation(|| verified_unload(&fp));
|
||||||
|
match result {
|
||||||
|
Ok(()) => {
|
||||||
|
permit.finish(UnloadOutcome::Confirmed);
|
||||||
|
true
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
permit.finish(UnloadOutcome::Uncertain(format!("{e:#}")));
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
enum UnloadPresence {
|
||||||
|
Exact,
|
||||||
|
Absent,
|
||||||
|
Replaced,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn unload_presence(fp: &Fingerprint, current: &[repair_plan::ModuleObservation]) -> UnloadPresence {
|
||||||
|
match current.iter().find(|module| module.id == fp.id) {
|
||||||
|
None => UnloadPresence::Absent,
|
||||||
|
Some(observed) if fp.still_matches(observed) => UnloadPresence::Exact,
|
||||||
|
Some(_) => UnloadPresence::Replaced,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Re-list and unload through one verified-local Pulse connection. The exact
|
||||||
|
/// fingerprint is checked immediately before destruction; if the index vanished
|
||||||
|
/// or was reused, our module is already absent and the replacement is left alone.
|
||||||
|
fn verified_unload(fp: &Fingerprint) -> Result<()> {
|
||||||
|
let mut session = crate::repair::introspect::PulseSession::connect()
|
||||||
|
.context("could not connect to verify a module unload")?;
|
||||||
|
let current = session
|
||||||
|
.list_modules()
|
||||||
|
.context("could not list modules immediately before unload")?;
|
||||||
|
match unload_presence(fp, ¤t) {
|
||||||
|
UnloadPresence::Exact => {}
|
||||||
|
UnloadPresence::Absent => {
|
||||||
|
tracing::info!(
|
||||||
|
module = fp.id,
|
||||||
|
shape = fp.shape.label(),
|
||||||
|
"audio routing: tracked module was already absent"
|
||||||
|
);
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
UnloadPresence::Replaced => {
|
||||||
|
tracing::warn!(
|
||||||
|
module = fp.id,
|
||||||
|
shape = fp.shape.label(),
|
||||||
|
"audio routing: module index was reused; leaving the replacement alone"
|
||||||
|
);
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
session
|
||||||
|
.unload_module(fp.id)
|
||||||
|
.with_context(|| format!("the server did not confirm unloading module #{}", fp.id))?;
|
||||||
|
tracing::info!(
|
||||||
|
module = fp.id,
|
||||||
|
shape = fp.shape.label(),
|
||||||
|
"audio routing: unloaded verified Pulse module"
|
||||||
|
);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// One bounded child result. On deadline the child is killed and synchronously
|
||||||
|
/// reaped before this returns, so server reconciliation cannot overtake a late
|
||||||
|
/// `pactl` request merely because its async waiter was cancelled.
|
||||||
|
struct BoundedOutput {
|
||||||
|
status: ExitStatus,
|
||||||
|
stdout: Vec<u8>,
|
||||||
|
stderr: Vec<u8>,
|
||||||
|
timed_out: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Ensures every early-return and panic after spawn kills and reaps the child.
|
||||||
|
/// The normal path marks it reaped after `try_wait`/`wait` obtained the status.
|
||||||
|
struct ReapedChild {
|
||||||
|
child: Child,
|
||||||
|
reaped: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for ReapedChild {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if self.reaped {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let _ = self.child.kill();
|
||||||
|
match self.reap_within(PACTL_REAP_BUDGET) {
|
||||||
|
Ok(Some(_)) => {}
|
||||||
|
Ok(None) => tracing::warn!(
|
||||||
|
"audio routing: killed pactl child was not reaped within {PACTL_REAP_BUDGET:?}"
|
||||||
|
),
|
||||||
|
Err(e) => tracing::warn!("audio routing: could not reap killed pactl child: {e}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ReapedChild {
|
||||||
|
fn reap_within(&mut self, budget: Duration) -> io::Result<Option<ExitStatus>> {
|
||||||
|
let deadline = Instant::now() + budget;
|
||||||
|
loop {
|
||||||
|
if let Some(status) = self.child.try_wait()? {
|
||||||
|
self.reaped = true;
|
||||||
|
return Ok(Some(status));
|
||||||
|
}
|
||||||
|
if Instant::now() >= deadline {
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
std::thread::sleep(Duration::from_millis(5));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn bounded_output(command: &mut Command, budget: Duration) -> io::Result<BoundedOutput> {
|
||||||
|
command.stdout(Stdio::piped()).stderr(Stdio::piped());
|
||||||
|
let mut child = ReapedChild {
|
||||||
|
child: command.spawn()?,
|
||||||
|
reaped: false,
|
||||||
|
};
|
||||||
|
let stdout = child
|
||||||
|
.child
|
||||||
|
.stdout
|
||||||
|
.take()
|
||||||
|
.ok_or_else(|| io::Error::other("pactl stdout was not piped"))?;
|
||||||
|
let stderr = child
|
||||||
|
.child
|
||||||
|
.stderr
|
||||||
|
.take()
|
||||||
|
.ok_or_else(|| io::Error::other("pactl stderr was not piped"))?;
|
||||||
|
let stdout_reader = std::thread::Builder::new()
|
||||||
|
.name("pixelpass-pactl-stdout".to_string())
|
||||||
|
.spawn(move || {
|
||||||
|
let mut bytes = Vec::new();
|
||||||
|
let mut stdout = stdout;
|
||||||
|
stdout.read_to_end(&mut bytes).map(|_| bytes)
|
||||||
|
})?;
|
||||||
|
let stderr_reader = std::thread::Builder::new()
|
||||||
|
.name("pixelpass-pactl-stderr".to_string())
|
||||||
|
.spawn(move || {
|
||||||
|
let mut bytes = Vec::new();
|
||||||
|
let mut stderr = stderr;
|
||||||
|
stderr.read_to_end(&mut bytes).map(|_| bytes)
|
||||||
|
})?;
|
||||||
|
|
||||||
|
let deadline = Instant::now() + budget;
|
||||||
|
let (status, timed_out) = loop {
|
||||||
|
if let Some(status) = child.child.try_wait()? {
|
||||||
|
child.reaped = true;
|
||||||
|
break (status, false);
|
||||||
|
}
|
||||||
|
if Instant::now() >= deadline {
|
||||||
|
let _ = child.child.kill();
|
||||||
|
let Some(status) = child.reap_within(PACTL_REAP_BUDGET)? else {
|
||||||
|
return Err(io::Error::new(
|
||||||
|
io::ErrorKind::TimedOut,
|
||||||
|
format!("pactl did not exit within {PACTL_REAP_BUDGET:?} after SIGKILL"),
|
||||||
|
));
|
||||||
|
};
|
||||||
|
break (status, true);
|
||||||
|
}
|
||||||
|
std::thread::sleep(Duration::from_millis(5));
|
||||||
|
};
|
||||||
|
let stdout = stdout_reader
|
||||||
|
.join()
|
||||||
|
.map_err(|_| io::Error::other("pactl stdout reader panicked"))??;
|
||||||
|
let stderr = stderr_reader
|
||||||
|
.join()
|
||||||
|
.map_err(|_| io::Error::other("pactl stderr reader panicked"))??;
|
||||||
|
Ok(BoundedOutput {
|
||||||
|
status,
|
||||||
|
stdout,
|
||||||
|
stderr,
|
||||||
|
timed_out,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Async teardown: keep reconciling and unloading while a pass changes state.
|
||||||
|
/// This replaces the arbitrary two-round count. With loads closed, the state
|
||||||
|
/// graph is monotonic except for `Ambiguous(Unload) -> Loaded -> Ambiguous` when
|
||||||
|
/// the same unload remains uncertain; that produces an identical snapshot and
|
||||||
|
/// stops here for `--repair` rather than spinning.
|
||||||
|
async fn cleanup_modules(ledger: &Arc<ModuleLedger>) {
|
||||||
|
loop {
|
||||||
|
let before = ledger.snapshot();
|
||||||
|
if ledger.is_clean() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
if let Err(e) = ledger::reconcile_pending(ledger).await {
|
||||||
|
tracing::warn!("audio routing: could not reconcile the module ledger: {e:#}");
|
||||||
|
}
|
||||||
|
for fp in ledger.loaded() {
|
||||||
|
unload_module_inner(ledger, fp.shape, true).await;
|
||||||
|
}
|
||||||
|
if ledger.is_clean() || ledger.snapshot() == before {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Synchronous teardown backstop. PulseSession bounds every connect/list/unload
|
||||||
|
/// request, and `close_and_wait` has already drained any registered pactl worker.
|
||||||
|
fn cleanup_modules_blocking(ledger: &Arc<ModuleLedger>) {
|
||||||
|
loop {
|
||||||
|
let before = ledger.snapshot();
|
||||||
|
if ledger.is_clean() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
if let Err(e) = ledger::reconcile_pending_blocking(ledger) {
|
||||||
|
tracing::warn!("audio routing: could not reconcile the module ledger: {e:#}");
|
||||||
|
}
|
||||||
|
for fp in ledger.loaded() {
|
||||||
|
let permit = match ledger.begin_cleanup_unload(fp.shape) {
|
||||||
|
Ok(Some(permit)) => permit,
|
||||||
|
Ok(None) => continue,
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(
|
||||||
|
shape = fp.shape.label(),
|
||||||
|
"audio routing: refusing blocking module unload: {e}"
|
||||||
|
);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
finish_verified_unload(permit);
|
||||||
|
}
|
||||||
|
if ledger.is_clean() || ledger.snapshot() == before {
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
Ok(output) => UnloadOutcome::Uncertain(format!(
|
|
||||||
"pactl unload-module exited {}: {}",
|
|
||||||
output.status,
|
|
||||||
String::from_utf8_lossy(&output.stderr).trim()
|
|
||||||
)),
|
|
||||||
Err(e) => UnloadOutcome::Uncertain(format!("could not run pactl unload-module: {e}")),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -835,6 +1072,7 @@ fn try_flush(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
use crate::host::ledger::SlotState;
|
||||||
use crate::repair::plan::{ModuleObservation, classify};
|
use crate::repair::plan::{ModuleObservation, classify};
|
||||||
|
|
||||||
/// Whole-desktop routing: no app filter, so no PipeWire thread and no event
|
/// Whole-desktop routing: no app filter, so no PipeWire thread and no event
|
||||||
@@ -868,6 +1106,41 @@ mod tests {
|
|||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn normal_unload_requires_the_full_fingerprint_not_only_the_index() {
|
||||||
|
fn observation(id: u32, pid: u32, nonce: u64) -> ModuleObservation {
|
||||||
|
let token = repair_plan::OwnerToken {
|
||||||
|
machine: "abc123".to_string(),
|
||||||
|
boot: "def456".to_string(),
|
||||||
|
pid_ns: 4_026_531_836,
|
||||||
|
nonce,
|
||||||
|
};
|
||||||
|
ModuleObservation::new(
|
||||||
|
id,
|
||||||
|
Shape::LoopbackIntoCapture.module_name(),
|
||||||
|
&repair_plan::recorded_argument(
|
||||||
|
&Shape::LoopbackIntoCapture.render_args(pid, Some(&token)),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
let ours = observation(5, 42, 7);
|
||||||
|
let fp = classify(&ours).expect("the fixture is canonical");
|
||||||
|
assert_eq!(
|
||||||
|
unload_presence(&fp, std::slice::from_ref(&ours)),
|
||||||
|
UnloadPresence::Exact
|
||||||
|
);
|
||||||
|
assert_eq!(unload_presence(&fp, &[]), UnloadPresence::Absent);
|
||||||
|
|
||||||
|
// Same live index, but another perfectly canonical host module. This is
|
||||||
|
// the non-vacuous reuse case: an id-only normal unload would destroy it.
|
||||||
|
let replacement = observation(5, 99, 8);
|
||||||
|
assert_eq!(
|
||||||
|
unload_presence(&fp, std::slice::from_ref(&replacement)),
|
||||||
|
UnloadPresence::Replaced
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
/// A/B against the live graph: routing must leave the module table exactly
|
/// A/B against the live graph: routing must leave the module table exactly
|
||||||
/// as it found it. The same shape as `--repair`'s field gate, because the
|
/// as it found it. The same shape as `--repair`'s field gate, because the
|
||||||
/// property is the same one — nothing of ours outlives the session.
|
/// property is the same one — nothing of ours outlives the session.
|
||||||
@@ -925,33 +1198,42 @@ mod tests {
|
|||||||
let task = tokio::spawn(async move {
|
let task = tokio::spawn(async move {
|
||||||
let _ = load_module(&ledger_for_task, Shape::LegacyCaptureSink, pid).await;
|
let _ = load_module(&ledger_for_task, Shape::LegacyCaptureSink, pid).await;
|
||||||
});
|
});
|
||||||
// Let the task run up to its first await — the spawned `pactl` — so the
|
// Wait until the affine permit is registered. A single yield is not a
|
||||||
// abort lands mid-flight rather than before the load ever started, which
|
// scheduling guarantee and made the old version of this gate capable of
|
||||||
// would make this gate vacuous.
|
// aborting before the task had started.
|
||||||
|
tokio::time::timeout(PACTL_BUDGET, async {
|
||||||
|
while matches!(ledger.state(Shape::LegacyCaptureSink), SlotState::Vacant) {
|
||||||
tokio::task::yield_now().await;
|
tokio::task::yield_now().await;
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("the load registers its operation");
|
||||||
task.abort();
|
task.abort();
|
||||||
let _ = task.await;
|
let _ = task.await;
|
||||||
|
|
||||||
// Non-vacuity: the permit is taken *before* `pactl` is spawned, so a
|
// Cancellation detaches `spawn_blocking`; closing plus this registered-
|
||||||
// cancelled load must leave a question behind whichever side of the spawn
|
// operation wait is the ordering boundary that prevents reconciliation
|
||||||
// the abort landed on. Without this the gate could pass while the abort
|
// from overtaking the worker's late server mutation.
|
||||||
// fired before the load ever began, proving nothing.
|
ledger.close();
|
||||||
|
let ledger_for_wait = Arc::clone(&ledger);
|
||||||
|
tokio::task::spawn_blocking(move || ledger_for_wait.wait_for_operations())
|
||||||
|
.await
|
||||||
|
.expect("the operation wait runs");
|
||||||
|
|
||||||
|
// The worker either committed a known module or left a reconciliation
|
||||||
|
// question. Both are correct; vacancy here would mean the live operation
|
||||||
|
// disappeared from the ledger.
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ledger.pending().len(),
|
ledger.pending().len() + ledger.loaded().len(),
|
||||||
1,
|
1,
|
||||||
"the cancelled load must have left exactly one question behind"
|
"the cancelled load must retain exactly one tracked outcome"
|
||||||
);
|
);
|
||||||
|
|
||||||
ledger::reconcile_pending(&ledger)
|
cleanup_modules(&ledger).await;
|
||||||
.await
|
|
||||||
.expect("the ledger reconciles against the server");
|
|
||||||
for fp in ledger.loaded() {
|
|
||||||
unload_module(&ledger, fp.shape).await;
|
|
||||||
}
|
|
||||||
|
|
||||||
assert!(
|
assert!(
|
||||||
ledger.is_settled(),
|
ledger.is_clean(),
|
||||||
"every slot must end in a state we can explain"
|
"every tracked module must be removed, not merely explained"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
module_snapshot(),
|
module_snapshot(),
|
||||||
|
|||||||
+522
-63
@@ -33,12 +33,13 @@
|
|||||||
//! # Why the permit is affine
|
//! # Why the permit is affine
|
||||||
//!
|
//!
|
||||||
//! [`LoadPermit`] is not `Clone`, is consumed by value to settle, and its [`Drop`]
|
//! [`LoadPermit`] is not `Clone`, is consumed by value to settle, and its [`Drop`]
|
||||||
//! marks the slot [`Reconcile::Load`] when it was never settled. That is what makes
|
//! marks the slot [`Reconcile::Load`] when it was never settled. The permit moves
|
||||||
//! the guarantee structural rather than a discipline: a cancelled task drops its
|
//! into a registered blocking worker before any await can detach that work; either
|
||||||
//! locals, so an aborted load *cannot* silently forget a module the server may
|
//! the worker settles it, or dropping the worker marks the question ambiguous.
|
||||||
//! already have created. Two permitted loads for one slot cannot both commit,
|
//! Teardown waits for all such permits before reconciling. Two permitted loads for
|
||||||
//! because [`ModuleLedger::begin_load`] issues a permit only for a `Vacant` slot
|
//! one slot cannot both commit, because [`ModuleLedger::begin_load`] issues a permit
|
||||||
//! and every other state — including the ambiguous one — refuses.
|
//! only for a `Vacant` slot and every other state — including the ambiguous one —
|
||||||
|
//! refuses.
|
||||||
//!
|
//!
|
||||||
//! # Why an ambiguous slot blocks the next load
|
//! # Why an ambiguous slot blocks the next load
|
||||||
//!
|
//!
|
||||||
@@ -50,8 +51,8 @@
|
|||||||
//! proceed until the question is answered is the whole point.
|
//! proceed until the question is answered is the whole point.
|
||||||
|
|
||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
use std::sync::atomic::{AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Condvar, Mutex};
|
||||||
|
|
||||||
use anyhow::{Context as _, Result};
|
use anyhow::{Context as _, Result};
|
||||||
|
|
||||||
@@ -99,7 +100,7 @@ pub enum SlotState {
|
|||||||
Loaded { fp: Fingerprint },
|
Loaded { fp: Fingerprint },
|
||||||
/// An unload is in flight. The fingerprint is retained deliberately: an unload
|
/// An unload is in flight. The fingerprint is retained deliberately: an unload
|
||||||
/// that times out must not leave the module unrecorded.
|
/// that times out must not leave the module unrecorded.
|
||||||
Unloading { fp: Fingerprint },
|
Unloading { fp: Fingerprint, permit: u64 },
|
||||||
/// The slot's real state is unknown and must be resolved against the server.
|
/// The slot's real state is unknown and must be resolved against the server.
|
||||||
Ambiguous(Reconcile),
|
Ambiguous(Reconcile),
|
||||||
/// Resolution found something we refuse to act on. Terminal.
|
/// Resolution found something we refuse to act on. Terminal.
|
||||||
@@ -129,6 +130,20 @@ pub enum LedgerError {
|
|||||||
Poisoned { shape: Shape, reason: String },
|
Poisoned { shape: Shape, reason: String },
|
||||||
/// `pactl` reported `PA_INVALID_INDEX` where an index was expected.
|
/// `pactl` reported `PA_INVALID_INDEX` where an index was expected.
|
||||||
InvalidIndex { shape: Shape },
|
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 {
|
impl std::fmt::Display for LedgerError {
|
||||||
@@ -145,6 +160,31 @@ impl std::fmt::Display for LedgerError {
|
|||||||
"pactl reported PA_INVALID_INDEX for the {} module",
|
"pactl reported PA_INVALID_INDEX for the {} module",
|
||||||
shape.label()
|
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()
|
||||||
|
),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -180,6 +220,59 @@ pub enum Resolution {
|
|||||||
pub struct ModuleLedger {
|
pub struct ModuleLedger {
|
||||||
slots: Mutex<BTreeMap<Shape, SlotState>>,
|
slots: Mutex<BTreeMap<Shape, SlotState>>,
|
||||||
next_permit: AtomicU64,
|
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<usize>,
|
||||||
|
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<T>(&self, operation: impl FnOnce() -> T) -> T {
|
||||||
|
let _serial = self.server.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
|
operation()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ModuleLedger {
|
impl ModuleLedger {
|
||||||
@@ -187,9 +280,36 @@ impl ModuleLedger {
|
|||||||
Arc::new(Self {
|
Arc::new(Self {
|
||||||
slots: Mutex::new(BTreeMap::new()),
|
slots: Mutex::new(BTreeMap::new()),
|
||||||
next_permit: AtomicU64::new(1),
|
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<T>(&self, operation: impl FnOnce() -> T) -> T {
|
||||||
|
self.operations.run(operation)
|
||||||
|
}
|
||||||
|
|
||||||
/// The current state of one slot. Absent keys read as [`SlotState::Vacant`].
|
/// 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
|
/// Test-only: production code never needs to look a slot up, because every
|
||||||
@@ -219,6 +339,9 @@ impl ModuleLedger {
|
|||||||
token: OwnerToken,
|
token: OwnerToken,
|
||||||
) -> Result<LoadPermit, LedgerError> {
|
) -> Result<LoadPermit, LedgerError> {
|
||||||
let mut slots = self.slots.lock().unwrap();
|
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);
|
let current = slots.get(&shape).cloned().unwrap_or(SlotState::Vacant);
|
||||||
match current {
|
match current {
|
||||||
SlotState::Vacant => {}
|
SlotState::Vacant => {}
|
||||||
@@ -232,7 +355,18 @@ impl ModuleLedger {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
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);
|
let permit = self.next_permit.fetch_add(1, Ordering::Relaxed);
|
||||||
|
self.operations.register();
|
||||||
slots.insert(
|
slots.insert(
|
||||||
shape,
|
shape,
|
||||||
SlotState::Loading {
|
SlotState::Loading {
|
||||||
@@ -253,36 +387,76 @@ impl ModuleLedger {
|
|||||||
|
|
||||||
/// Move a `Loaded` slot to `Unloading` and hand back what to unload.
|
/// Move a `Loaded` slot to `Unloading` and hand back what to unload.
|
||||||
///
|
///
|
||||||
/// `None` for any other state: there is nothing to unload, or the slot is not
|
/// `Ok(None)` means confirmed vacancy. Busy, poisoned, closed, and still-
|
||||||
/// in a condition to be acted on.
|
/// referenced states are errors so callers cannot mistake uncertainty for
|
||||||
pub fn begin_unload(&self, shape: Shape) -> Option<Fingerprint> {
|
/// absence and load an inverse loopback or destroy the capture sink.
|
||||||
let mut slots = self.slots.lock().unwrap();
|
pub fn begin_unload(
|
||||||
let SlotState::Loaded { fp } = slots.get(&shape).cloned()? else {
|
self: &Arc<Self>,
|
||||||
return None;
|
shape: Shape,
|
||||||
};
|
) -> Result<Option<UnloadPermit>, LedgerError> {
|
||||||
slots.insert(shape, SlotState::Unloading { fp: fp.clone() });
|
self.begin_unload_inner(shape, false)
|
||||||
Some(fp)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Record how an unload turned out. An uncertain one becomes ambiguous rather
|
/// Teardown-only form of [`ModuleLedger::begin_unload`]. New event work is
|
||||||
/// than being assumed done — the module is not forgotten either way.
|
/// closed by then, but cleanup must still be able to remove tracked modules.
|
||||||
pub fn finish_unload(&self, shape: Shape, outcome: UnloadOutcome) {
|
pub(super) fn begin_cleanup_unload(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
shape: Shape,
|
||||||
|
) -> Result<Option<UnloadPermit>, LedgerError> {
|
||||||
|
self.begin_unload_inner(shape, true)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn begin_unload_inner(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
shape: Shape,
|
||||||
|
cleanup: bool,
|
||||||
|
) -> Result<Option<UnloadPermit>, LedgerError> {
|
||||||
let mut slots = self.slots.lock().unwrap();
|
let mut slots = self.slots.lock().unwrap();
|
||||||
let Some(SlotState::Unloading { fp }) = slots.get(&shape).cloned() else {
|
if !cleanup && self.closed.load(Ordering::SeqCst) {
|
||||||
return;
|
return Err(LedgerError::Closed { shape });
|
||||||
};
|
}
|
||||||
let next = match outcome {
|
let current = slots.get(&shape).cloned().unwrap_or(SlotState::Vacant);
|
||||||
UnloadOutcome::Confirmed => SlotState::Vacant,
|
let fp = match current {
|
||||||
UnloadOutcome::Uncertain(why) => {
|
SlotState::Vacant => return Ok(None),
|
||||||
tracing::warn!(
|
SlotState::Loaded { fp } => fp,
|
||||||
shape = fp.shape.label(),
|
SlotState::Poisoned { reason } => {
|
||||||
module = fp.id,
|
return Err(LedgerError::Poisoned { shape, reason });
|
||||||
"audio ledger: unload outcome uncertain ({why}); slot needs reconciling"
|
}
|
||||||
);
|
other => {
|
||||||
SlotState::Ambiguous(Reconcile::Unload { fp })
|
return Err(LedgerError::Busy {
|
||||||
|
shape,
|
||||||
|
state: other.label(),
|
||||||
|
});
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
slots.insert(shape, next);
|
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 —
|
/// Every slot currently holding a module, in [`Shape`] declaration order —
|
||||||
@@ -313,14 +487,28 @@ impl ModuleLedger {
|
|||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Is every slot in a state we can explain? False while anything is ambiguous
|
/// Is every slot in a state we can explain? False while an operation is in
|
||||||
/// or poisoned.
|
/// flight, its outcome is ambiguous, or a conflict poisoned the slot.
|
||||||
pub fn is_settled(&self) -> bool {
|
pub fn is_settled(&self) -> bool {
|
||||||
self.slots
|
self.slots
|
||||||
.lock()
|
.lock()
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.values()
|
.values()
|
||||||
.all(|state| !matches!(state, SlotState::Ambiguous(_) | SlotState::Poisoned { .. }))
|
.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<Shape, SlotState> {
|
||||||
|
self.slots.lock().unwrap().clone()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Apply a reconciliation result to an ambiguous slot.
|
/// Apply a reconciliation result to an ambiguous slot.
|
||||||
@@ -359,6 +547,14 @@ impl ModuleLedger {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn inverse_loopback(shape: Shape) -> Option<Shape> {
|
||||||
|
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.
|
/// 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
|
/// Dropping it unsettled is not an error — it is the *reporting* path for a
|
||||||
@@ -375,6 +571,12 @@ pub struct LoadPermit {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl LoadPermit {
|
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<T>(&self, operation: impl FnOnce() -> T) -> T {
|
||||||
|
self.ledger.with_server_operation(operation)
|
||||||
|
}
|
||||||
|
|
||||||
/// Record that the server created the module at `index`.
|
/// Record that the server created the module at `index`.
|
||||||
///
|
///
|
||||||
/// Fails on [`PA_INVALID_INDEX`], leaving the permit unsettled so that
|
/// Fails on [`PA_INVALID_INDEX`], leaving the permit unsettled so that
|
||||||
@@ -427,9 +629,7 @@ impl LoadPermit {
|
|||||||
|
|
||||||
impl Drop for LoadPermit {
|
impl Drop for LoadPermit {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
if self.settled {
|
if !self.settled {
|
||||||
return;
|
|
||||||
}
|
|
||||||
let mut slots = self.ledger.slots.lock().unwrap();
|
let mut slots = self.ledger.slots.lock().unwrap();
|
||||||
// Only claim the slot if it is still *our* load. Anything else already
|
// Only claim the slot if it is still *our* load. Anything else already
|
||||||
// moved past this permit.
|
// moved past this permit.
|
||||||
@@ -451,6 +651,85 @@ impl Drop for LoadPermit {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
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<ModuleLedger>,
|
||||||
|
shape: Shape,
|
||||||
|
fp: Fingerprint,
|
||||||
|
permit: u64,
|
||||||
|
settled: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl UnloadPermit {
|
||||||
|
pub fn fingerprint(&self) -> &Fingerprint {
|
||||||
|
&self.fp
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn with_server_operation<T>(&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.
|
/// The fingerprint a successful load of `shape` for `pid` under `token` must have.
|
||||||
@@ -509,6 +788,17 @@ pub fn resolve(reconcile: &Reconcile, observations: &[ModuleObservation]) -> Res
|
|||||||
/// the life of a share. Reconciliation is rare and off the hot path, so paying a
|
/// the life of a share. Reconciliation is rare and off the hot path, so paying a
|
||||||
/// connection for it costs nothing that matters.
|
/// connection for it costs nothing that matters.
|
||||||
pub async fn reconcile_pending(ledger: &Arc<ModuleLedger>) -> Result<()> {
|
pub async fn reconcile_pending(ledger: &Arc<ModuleLedger>) -> 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<ModuleLedger>) -> Result<()> {
|
||||||
|
ledger.with_server_operation(|| {
|
||||||
let pending = ledger.pending();
|
let pending = ledger.pending();
|
||||||
if pending.is_empty() {
|
if pending.is_empty() {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -517,18 +807,16 @@ pub async fn reconcile_pending(ledger: &Arc<ModuleLedger>) -> Result<()> {
|
|||||||
n = pending.len(),
|
n = pending.len(),
|
||||||
"audio ledger: reconciling unresolved module slots against the server"
|
"audio ledger: reconciling unresolved module slots against the server"
|
||||||
);
|
);
|
||||||
let observations = tokio::task::spawn_blocking(|| -> Result<Vec<ModuleObservation>> {
|
|
||||||
let mut session = PulseSession::connect()?;
|
let mut session = PulseSession::connect()?;
|
||||||
session.list_modules()
|
let observations = session
|
||||||
})
|
.list_modules()
|
||||||
.await
|
|
||||||
.context("the Pulse listing task failed to run")?
|
|
||||||
.context("could not list Pulse modules to reconcile the audio ledger")?;
|
.context("could not list Pulse modules to reconcile the audio ledger")?;
|
||||||
|
|
||||||
for (shape, reconcile) in pending {
|
for (shape, reconcile) in pending {
|
||||||
ledger.apply(shape, resolve(&reconcile, &observations));
|
ledger.apply(shape, resolve(&reconcile, &observations));
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -684,14 +972,12 @@ mod tests {
|
|||||||
.expect("a vacant slot issues a permit")
|
.expect("a vacant slot issues a permit")
|
||||||
.commit(3)
|
.commit(3)
|
||||||
.expect("3 is a real index");
|
.expect("3 is a real index");
|
||||||
assert_eq!(
|
let permit = ledger
|
||||||
ledger.begin_unload(Shape::LoopbackIntoCapture),
|
.begin_unload(Shape::LoopbackIntoCapture)
|
||||||
Some(fp.clone())
|
.expect("the ledger accepts the unload")
|
||||||
);
|
.expect("the loaded slot issues a permit");
|
||||||
ledger.finish_unload(
|
assert_eq!(permit.fingerprint(), &fp);
|
||||||
Shape::LoopbackIntoCapture,
|
permit.finish(UnloadOutcome::Uncertain("pactl was killed".to_string()));
|
||||||
UnloadOutcome::Uncertain("pactl was killed".to_string()),
|
|
||||||
);
|
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ledger.state(Shape::LoopbackIntoCapture),
|
ledger.state(Shape::LoopbackIntoCapture),
|
||||||
SlotState::Ambiguous(Reconcile::Unload { fp }),
|
SlotState::Ambiguous(Reconcile::Unload { fp }),
|
||||||
@@ -707,23 +993,194 @@ mod tests {
|
|||||||
.expect("a vacant slot issues a permit")
|
.expect("a vacant slot issues a permit")
|
||||||
.commit(3)
|
.commit(3)
|
||||||
.expect("3 is a real index");
|
.expect("3 is a real index");
|
||||||
ledger.begin_unload(Shape::LoopbackIntoCapture);
|
ledger
|
||||||
ledger.finish_unload(Shape::LoopbackIntoCapture, UnloadOutcome::Confirmed);
|
.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_eq!(ledger.state(Shape::LoopbackIntoCapture), SlotState::Vacant);
|
||||||
assert!(ledger.is_settled());
|
assert!(ledger.is_settled());
|
||||||
|
assert!(ledger.is_clean());
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn begin_unload_only_acts_on_a_loaded_slot() {
|
fn dropping_an_unload_permit_marks_the_slot_ambiguous() {
|
||||||
let ledger = ModuleLedger::new();
|
let ledger = ModuleLedger::new();
|
||||||
assert_eq!(ledger.begin_unload(Shape::LegacyCaptureSink), None);
|
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
|
let _permit = ledger
|
||||||
.begin_load(Shape::LegacyCaptureSink, 42, token(7))
|
.begin_load(Shape::LegacyCaptureSink, 42, token(7))
|
||||||
.expect("a vacant slot issues a permit");
|
.expect("a vacant slot issues a permit");
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ledger.begin_unload(Shape::LegacyCaptureSink),
|
ledger.begin_unload(Shape::LegacyCaptureSink).err(),
|
||||||
None,
|
Some(LedgerError::Busy {
|
||||||
"a load in flight has no id to unload yet"
|
shape: Shape::LegacyCaptureSink,
|
||||||
|
state: "loading"
|
||||||
|
}),
|
||||||
|
"an in-flight load must not look like a confirmed vacancy"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -864,13 +1321,12 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn loaded_lists_modules_in_shape_declaration_order() {
|
fn loaded_lists_a_loopback_before_the_capture_sink() {
|
||||||
// Declaration order is unload order: the loopbacks that reference the
|
// Declaration order is unload order: the loopbacks that reference the
|
||||||
// capture sink must come before the sink itself.
|
// capture sink must come before the sink itself.
|
||||||
let ledger = ModuleLedger::new();
|
let ledger = ModuleLedger::new();
|
||||||
for (shape, id) in [
|
for (shape, id) in [
|
||||||
(Shape::LegacyCaptureSink, 1),
|
(Shape::LegacyCaptureSink, 1),
|
||||||
(Shape::LoopbackIntoCapture, 2),
|
|
||||||
(Shape::LoopbackOutOfCapture, 3),
|
(Shape::LoopbackOutOfCapture, 3),
|
||||||
] {
|
] {
|
||||||
ledger
|
ledger
|
||||||
@@ -880,7 +1336,10 @@ mod tests {
|
|||||||
.expect("a real index");
|
.expect("a real index");
|
||||||
}
|
}
|
||||||
let order: Vec<Shape> = ledger.loaded().into_iter().map(|fp| fp.shape).collect();
|
let order: Vec<Shape> = ledger.loaded().into_iter().map(|fp| fp.shape).collect();
|
||||||
assert_eq!(order, crate::repair::plan::ALL_SHAPES.to_vec());
|
assert_eq!(
|
||||||
|
order,
|
||||||
|
vec![Shape::LoopbackOutOfCapture, Shape::LegacyCaptureSink]
|
||||||
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
order.last(),
|
order.last(),
|
||||||
Some(&Shape::LegacyCaptureSink),
|
Some(&Shape::LegacyCaptureSink),
|
||||||
|
|||||||
Reference in New Issue
Block a user