Author SHA1 Message Date
mollusk 471b8221ff host/observer: address the Codex phase-3r review (2 fixes, both verified)
Codex's adversarial review of the pure core found no *certain* P1. Two
findings taken, both mutation-verified (the fix reverted, the intended test
dies, nothing else moves):

**F2, certain, P2 — `device_props` tested the wrong kind of ambiguity.** It
required exactly one live *Device* on the claimed id rather than exactly one
live *global*. With `[Device, Port]` on one id — a missed removal, the same
precondition as every other recycled-id hazard — it kept answering from the
older Device, so a node claiming that id held a stale `session_device = true`.
That flag strips the node's owner keys and its fail-closed backstop, so a
forwarder wearing it can put its output leg back on the eligible side. Now:
one slot total, and it must be the Device.

**F3, worth checking, P3 — `device.api` was corroborating by presence.**
`device.api=v4l2` under `factory.name=api.alsa.pcm.sink` satisfied the
positive classifier. No truthful configuration produces that pair, which is
the argument for reading it as an observation gone wrong rather than as
corroboration. The API must now equal the one the factory allowlist is
written for, an empty value is not a value, and the two sides disagreeing
fails closed. Tied to the allowlist being ALSA-only via a named constant.

Two findings NOT fixed here, both pre-existing and neither introduced by
round 8 — raised to the design doc instead:

- **P1, worth checking: hardware playback-to-capture paths** (Stereo Mix,
  Digital Loopback) on a card whose driver is an ordinary `snd_hda_intel`.
  Both its sink and source classify `session_device`, taint cannot cross the
  hardware hop, and a capture app reading that source can re-emit the call.
  This is `snd_aloop` again in a form the driver name cannot detect;
  distinguishing it needs ALSA control inspection, which is a design change
  and a new I/O surface, not a local fix.
- **P3: the 2 s readiness budget** can in principle never see an
  obligation-free instant under sustained startup churn, and `TimedOut` is
  sticky by design, so the process would be silent for its lifetime.
  Measured here: readiness at ~3 ms with 19 binds, so the margin is three
  orders of magnitude — but it wants a calibration argument, not a guess.

197 unit + 3 live green, clippy -D warnings and fmt clean.
2026-07-25 18:48:53 -04:00
mollusk 306b601490 host/observer: phase 3r adapter — bind every Node and Device
The I/O half of round 8 (Codex, gpt-5.6-sol xhigh; reviewed, formatted and
extended here). The adapter now reads `object.serial` and nothing else off a
Node or Device global, binds the object, and takes every property the engine
reasons about from its `info` props.

- `BoundProxy` generalises `BoundLink` to Node/Device/Link, each holding its
  listener *before* its proxy so the listener is dropped first — the original
  Link variant had that order inverted.
- Bind attachment now finds its slot by never-recycled serial rather than
  taking the queue's back, so nested callback activity during a bind cannot
  attach one generation's proxy to another's slot on a recycled id. A proxy
  that finds no slot is returned to the caller and dropped after the borrow
  ends. Removal still pops oldest-first, matching the model's `live_ids`.
- An `info` is parsed and emitted on the first callback carrying props and
  thereafter only when `change_mask` contains PROPS. I considered emitting
  unconditionally and leaning on the model's suppression rule, and rejected
  it: if a state-only `info` ever delivered a partial props dict, that would
  overwrite a complete observation with an incomplete one — a worse failure
  than the one it guards against, and the same class as F1.
- Ports stay unbound (v3.5 §6.7 / impl plan §4 item 6).

Gates: exit-gate row 1 (live prop recovery) passes on this host — the tagged
null sink projects `peerspeak.owned`, `pulse.module.id`, `node.passthrough`,
the loopback legs share a `node.link-group`, and a real ALSA node classifies
`session_device`.

Added a second live test for the Device half. Row 1's `session_device`
assertion is satisfied by a *union*: WirePlumber 0.5.15 copies `device.api`
and `alsa.driver_name` onto ALSA nodes here, so it passes through the node
fallback and would keep passing if the Device bind delivered nothing —
leaving §6.7 decision 4 ungated on the development machine. The new test
binds every Device and requires an ALSA card to announce both keys.
Mutation-verified: breaking the Device-side driver read fails the new test
while row 1 still passes, which is the gap as claimed.

195 unit + 3 live green, clippy -D warnings and fmt clean.
2026-07-25 18:32:27 -04:00
mollusk b3d71724ae host/observer: phase 3r pure core — node/device props come from a bind
v3.5 §6.7. The registry `global` event announces only a fixed 13-key subset
of a Node's properties, and eight the engine depends on are never among them
(phase-5 gate failure F1/F2). The core now treats the global as an index and
takes every property from the object's bound `info`.

- `RegEvent::NodeAdded { serial, id }` is identity only; `RegEvent::NodeInfo`
  carries the properties and is both the first resolution and every later
  PROPS change for the node's lifetime (decision 2). Same split for Device
  (`DeviceAdded` / `DeviceInfo`).
- A node with no `info` is withheld from the snapshot and is a readiness
  obligation; an unresolvable bind ends in sticky `TimedOut`, fail closed
  (decision 3). Devices are keyed by serial too, so a recycled device id with
  two live claimants is ambiguous ⇒ withheld rather than guessed.
- One live-node map replaces the admitted/withheld pair; classification is
  recomputed at projection time from current inputs, since both sides of it
  now change over an object's lifetime.
- `classify` takes the bound Device's props: presence is a union with the
  Device winning (this recovers a real card whose node was never given
  `alsa.driver_name` — the phase-3 review's owed fix), while the
  non-terminal-driver denylist is a union in the safe direction.
- `apply` returns `Outcome`, the only sound place to enforce the suppression
  rule: a property update is dropped only when model state provably did not
  change, i.e. the resulting projection is identical.

Adapter: stops reading properties off Node/Device globals and emits the new
index events. Binding every Node and Device — the I/O half — is the next
commit (Codex's), so until then every node is withheld and readiness times
out by design.

Tests: 55 observer (was 38) — the prop-update matrix, readiness with node
binds, and recycled-Node-id churn (phase 3r gate rows 2–4). 195 green,
clippy -D warnings and fmt clean.
2026-07-25 18:17:48 -04:00
molluskandClaude Opus 5 a1ac7ea8d5 Merge phase 5: dry-run audit mode (read-only)
The audit machinery is complete and verified live. The §5.1 gate itself
FAILED — see peerspeak docs/screenshare-audio-exclusion-phase5-results.md —
but both findings are defects in phase 3's observation boundary, not in this
code, and round 8 needs the audit tool on main to re-run the matrix.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 15:43:39 -04:00
4 changed files with 1684 additions and 343 deletions
+429 -70
View File
@@ -4,8 +4,10 @@
//! callbacks into [`RegEvent`]s, and publishes the latest [`Projection`] for
//! consumers running outside the PipeWire thread.
use super::classify::DeviceClaim;
use super::{EventKind, LinkEndpoints, NodeObservation, Projection, RegEvent, RegistryModel};
use super::classify::{DeviceClaim, DeviceProps};
use super::{
EventKind, LinkEndpoints, NodeObservation, Outcome, Projection, RegEvent, RegistryModel,
};
use crate::host::audio::parse_object_serial;
use crate::host::taint::snapshot::{
ClientSnapshot, GlobalId, MediaRole, NodeProps, PortDirection, PortSnapshot, Serial,
@@ -104,14 +106,24 @@ impl Drop for RegistryObserverHandle {
}
}
struct BoundLink {
_proxy: pw::link::Link,
_listener: pw::link::LinkListener,
enum BoundProxy {
Node {
_listener: pw::node::NodeListener,
_proxy: pw::node::Node,
},
Device {
_listener: pw::device::DeviceListener,
_proxy: pw::device::Device,
},
Link {
_listener: pw::link::LinkListener,
_proxy: pw::link::Link,
},
}
#[derive(Default)]
struct LiveGlobal {
bound_link: Option<BoundLink>,
serial: Serial,
bound_proxy: Option<BoundProxy>,
}
struct ObserverState {
@@ -142,12 +154,13 @@ impl ObserverState {
}
}
fn apply(&mut self, event: RegEvent) {
fn apply(&mut self, event: RegEvent) -> Outcome {
// Taken before the model consumes the event: the sink is told what kind
// of observation produced the projection, and deriving that from the
// event itself is what stops the two from ever disagreeing.
let kind = event.kind();
self.model.apply(event);
let event_outcome = self.model.apply(event);
let mut outcome = event_outcome;
let candidate = self.model.pulse_pid_candidate();
if candidate != self.last_candidate {
@@ -160,11 +173,19 @@ impl ObserverState {
// one registry event still yields exactly one sink call — the
// no-coalescing contract cuts both ways, and a *duplicated*
// observation would make the O5 event rate a fiction.
self.model.apply(RegEvent::ProcCommProbed { pid, comm });
if self.model.apply(RegEvent::ProcCommProbed { pid, comm }) == Outcome::Applied {
outcome = Outcome::Applied;
}
}
}
self.publish(kind);
// v3.5 §6.7 decision 2: a projection the model proved identical is not
// published. Only the model can make that claim soundly, which is why
// it is [`Outcome`] and not a diff of two snapshots here.
if outcome == Outcome::Applied {
self.publish(kind);
}
event_outcome
}
fn publish(&mut self, kind: EventKind) {
@@ -183,40 +204,65 @@ impl ObserverState {
}
/// Record the global's id and apply its add event as one step, so the
/// bound-link FIFO stays provably lockstep with the model's own `live_ids`
/// bound-proxy FIFO stays provably lockstep with the model's own `live_ids`
/// index. Recording only on *applied* adds (never on unknown object types
/// or globals dropped for a missing serial) is what keeps the two id
/// queues the same length per id — otherwise a phantom slot ahead of a
/// bound Link would be popped on removal, leaking that Link's proxy.
fn add(&mut self, id: GlobalId, event: RegEvent) {
self.live_globals
.entry(id)
.or_default()
.push_back(LiveGlobal::default());
self.apply(event);
/// queues the same length per id — otherwise a phantom slot could pop
/// another generation's proxy after an id is recycled.
fn add(&mut self, serial: Serial, id: GlobalId, event: RegEvent) {
if self.apply(event) == Outcome::Applied {
self.live_globals
.entry(id)
.or_default()
.push_back(LiveGlobal {
serial,
bound_proxy: None,
});
}
}
fn attach_bound_link(&mut self, id: GlobalId, bound_link: BoundLink) {
let Some(global) = self.live_globals.get_mut(&id).and_then(VecDeque::back_mut) else {
/// Return a proxy that could not be attached so its listener is dropped
/// after the caller releases the `RefCell` borrow.
fn attach_bound_proxy(
&mut self,
id: GlobalId,
serial: Serial,
bound_proxy: BoundProxy,
) -> Option<BoundProxy> {
let Some(global) = self
.live_globals
.get_mut(&id)
.and_then(|globals| globals.iter_mut().find(|global| global.serial == serial))
else {
tracing::warn!(
global_id = id.0,
"registry observer: link bind completed without a live global slot"
serial = serial.0,
"registry observer: bind completed without a live global slot"
);
return;
return Some(bound_proxy);
};
global.bound_link = Some(bound_link);
if global.bound_proxy.is_some() {
tracing::warn!(
global_id = id.0,
serial = serial.0,
"registry observer: live global slot already has a bound proxy"
);
return Some(bound_proxy);
}
global.bound_proxy = Some(bound_proxy);
None
}
fn remove_global(&mut self, id: GlobalId) -> Option<BoundLink> {
let (bound_link, empty) = {
fn remove_global(&mut self, id: GlobalId) -> Option<BoundProxy> {
let (bound_proxy, empty) = {
let globals = self.live_globals.get_mut(&id)?;
let bound_link = globals.pop_front().and_then(|global| global.bound_link);
(bound_link, globals.is_empty())
let bound_proxy = globals.pop_front().and_then(|global| global.bound_proxy);
(bound_proxy, globals.is_empty())
};
if empty {
self.live_globals.remove(&id);
}
bound_link
bound_proxy
}
}
@@ -273,6 +319,9 @@ fn run_observer(
match obj.type_ {
ObjectType::Node => {
// ⚠️ v3.5 §6.7: the global is an INDEX. Only `object.serial`
// is read here; every property the engine reasons about
// comes from the bind's `info` (phase 3r).
let Some(props) = obj.props.as_ref() else {
tracing::warn!(
node_id = obj.id,
@@ -284,41 +333,59 @@ fn run_observer(
else {
return;
};
let node_props = NodeProps {
peerspeak_owned: truthy(props.get("peerspeak.owned")),
pulse_module_id: props
.get("pulse.module.id")
.and_then(|value| value.parse::<u64>().ok()),
link_group: props.get("node.link-group").map(str::to_owned),
client_id: props
.get("client.id")
.and_then(|value| value.parse::<u32>().ok())
.map(GlobalId),
process_id: props
.get("application.process.id")
.and_then(|value| value.parse::<u32>().ok()),
passthrough: truthy(props.get("node.passthrough")),
session_device: false,
};
let observation = NodeObservation {
state_for_global.borrow_mut().add(
serial,
id,
name: props.get("node.name").map(str::to_owned),
role: MediaRole::parse(props.get("media.class")),
props: node_props,
device_claim: DeviceClaim {
device_id: props
.get("device.id")
.and_then(|value| value.parse::<u32>().ok())
.map(GlobalId),
device_api: props.get("device.api").map(str::to_owned),
factory_name: props.get("factory.name").map(str::to_owned),
alsa_driver_name: props.get("alsa.driver_name").map(str::to_owned),
},
RegEvent::NodeAdded { serial, id },
);
let Some(registry) = registry_weak.upgrade() else {
return;
};
state_for_global
.borrow_mut()
.add(id, RegEvent::NodeAdded(observation));
let node: pw::node::Node = match registry.bind(obj) {
Ok(node) => node,
Err(e) => {
tracing::warn!(
node_id = obj.id,
"registry observer: failed to bind Node for properties: {e}"
);
return;
}
};
// This bit only recognizes the initial callback for the
// change-mask fast path. Admission vs update remains
// entirely the model's decision.
let first_info = Cell::new(true);
let state_for_info = Rc::downgrade(&state_for_global);
let listener = node
.add_listener_local()
.info(move |info| {
let Some(props) = info.props() else {
return;
};
let first = first_info.replace(false);
if !first
&& !info.change_mask().contains(pw::node::NodeChangeMask::PROPS)
{
return;
}
if let Some(state) = state_for_info.upgrade() {
state.borrow_mut().apply(RegEvent::NodeInfo {
serial,
observation: node_observation_from_props(props),
});
}
})
.register();
let unattached = state_for_global.borrow_mut().attach_bound_proxy(
id,
serial,
BoundProxy::Node {
_listener: listener,
_proxy: node,
},
);
drop(unattached);
}
ObjectType::Port => {
let Some(props) = obj.props.as_ref() else {
@@ -357,6 +424,7 @@ fn run_observer(
}
};
state_for_global.borrow_mut().add(
serial,
id,
RegEvent::PortAdded(PortSnapshot {
serial,
@@ -381,6 +449,7 @@ fn run_observer(
return;
};
state_for_global.borrow_mut().add(
serial,
id,
RegEvent::ClientAdded(ClientSnapshot {
serial,
@@ -392,9 +461,72 @@ fn run_observer(
);
}
ObjectType::Device => {
state_for_global
.borrow_mut()
.add(id, RegEvent::DeviceAdded { id });
// Index only, exactly as for a Node: `device.api` and
// `alsa.driver_name` live on the bind's `info` (v3.5 §6.7
// decision 4), not here.
let Some(props) = obj.props.as_ref() else {
tracing::warn!(
device_id = obj.id,
"registry observer: Device has no properties; dropping"
);
return;
};
let Some(serial) = parse_serial(obj.id, "Device", props.get("object.serial"))
else {
return;
};
state_for_global.borrow_mut().add(
serial,
id,
RegEvent::DeviceAdded { serial, id },
);
let Some(registry) = registry_weak.upgrade() else {
return;
};
let device: pw::device::Device = match registry.bind(obj) {
Ok(device) => device,
Err(e) => {
tracing::warn!(
device_id = obj.id,
"registry observer: failed to bind Device for properties: {e}"
);
return;
}
};
let first_info = Cell::new(true);
let state_for_info = Rc::downgrade(&state_for_global);
let listener = device
.add_listener_local()
.info(move |info| {
let Some(props) = info.props() else {
return;
};
let first = first_info.replace(false);
if !first
&& !info
.change_mask()
.contains(pw::device::DeviceChangeMask::PROPS)
{
return;
}
if let Some(state) = state_for_info.upgrade() {
state.borrow_mut().apply(RegEvent::DeviceInfo {
serial,
props: device_props_from_props(props),
});
}
})
.register();
let unattached = state_for_global.borrow_mut().attach_bound_proxy(
id,
serial,
BoundProxy::Device {
_listener: listener,
_proxy: device,
},
);
drop(unattached);
}
ObjectType::Link => {
let Some(props) = obj.props.as_ref() else {
@@ -410,6 +542,7 @@ fn run_observer(
};
let endpoints = link_endpoints_from_props(props);
state_for_global.borrow_mut().add(
serial,
id,
RegEvent::LinkAdded {
serial,
@@ -456,24 +589,26 @@ fn run_observer(
}
})
.register();
state_for_global.borrow_mut().attach_bound_link(
let unattached = state_for_global.borrow_mut().attach_bound_proxy(
id,
BoundLink {
_proxy: link,
serial,
BoundProxy::Link {
_listener: listener,
_proxy: link,
},
);
drop(unattached);
}
_ => {}
}
})
.global_remove(move |id| {
let id = GlobalId(id);
let bound_link = state_for_remove.borrow_mut().remove_global(id);
let bound_proxy = state_for_remove.borrow_mut().remove_global(id);
state_for_remove
.borrow_mut()
.apply(RegEvent::Removed { id });
drop(bound_link);
drop(bound_proxy);
})
.register();
@@ -517,6 +652,45 @@ fn truthy(value: Option<&str>) -> bool {
value.is_some_and(|value| value != "false" && value != "0")
}
fn node_observation_from_props(props: &pw::spa::utils::dict::DictRef) -> NodeObservation {
NodeObservation {
name: props.get("node.name").map(str::to_string),
role: MediaRole::parse(props.get("media.class")),
props: NodeProps {
peerspeak_owned: truthy(props.get("peerspeak.owned")),
pulse_module_id: props
.get("pulse.module.id")
.and_then(|value| value.parse::<u64>().ok()),
link_group: props.get("node.link-group").map(str::to_string),
client_id: props
.get("client.id")
.and_then(|value| value.parse::<u32>().ok())
.map(GlobalId),
process_id: props
.get("application.process.id")
.and_then(|value| value.parse::<u32>().ok()),
passthrough: truthy(props.get("node.passthrough")),
session_device: false,
},
device_claim: DeviceClaim {
device_id: props
.get("device.id")
.and_then(|value| value.parse::<u32>().ok())
.map(GlobalId),
device_api: props.get("device.api").map(str::to_string),
factory_name: props.get("factory.name").map(str::to_string),
alsa_driver_name: props.get("alsa.driver_name").map(str::to_string),
},
}
}
fn device_props_from_props(props: &pw::spa::utils::dict::DictRef) -> DeviceProps {
DeviceProps {
device_api: props.get("device.api").map(str::to_string),
alsa_driver_name: props.get("alsa.driver_name").map(str::to_string),
}
}
fn link_endpoints_from_props(props: &pw::spa::utils::dict::DictRef) -> Option<LinkEndpoints> {
let output_node = props.get("link.output.node")?.parse::<u32>().ok()?;
let input_node = props.get("link.input.node")?.parse::<u32>().ok()?;
@@ -617,6 +791,191 @@ mod tests {
.any(|node| node.name.as_deref() == Some(name))
}
/// Phase 3r exit-gate row 1, the Device half — and the reason it needs its
/// own test.
///
/// `live_bound_properties_recover_node_and_device_inputs` asserts
/// `session_device`, which the classifier grants on a **union**:
/// `device.api` and `alsa.driver_name` may come from the bound Device *or*
/// from the node's own copies. On this host (WirePlumber 0.5.15 ≥ 0.5.13)
/// the session manager *does* copy both onto ALSA nodes, so that assertion
/// passes through the node fallback and would keep passing if the Device
/// bind delivered nothing at all — leaving v3.5 §6.7 decision 4, the whole
/// authoritative path, ungated on the machine we develop on.
///
/// So assert the Device side directly: bind every Device global and require
/// that at least one ALSA card announces **both** keys on its `info` props.
/// A failure here means the fix for the phase-3 review's owed finding (a
/// real card over-excluded on installs that do not copy `alsa.*` onto the
/// node) rests on nothing.
#[test]
#[ignore = "needs live pipewire"]
fn live_device_bind_carries_api_and_driver_name() {
pw::init();
let main_loop = pw::main_loop::MainLoopRc::new(None).expect("pw main loop");
let context = pw::context::ContextRc::new(&main_loop, None).expect("pw context");
let core = context.connect_rc(None).expect("pw core connect");
let registry = core.get_registry_rc().expect("pw registry");
// Devices bound off the registry, each holding its proxy + listener so
// the callback lives long enough to fire, exactly as the adapter does.
let bound: Rc<RefCell<Vec<(pw::device::Device, pw::device::DeviceListener)>>> =
Rc::new(RefCell::new(Vec::new()));
let observed: Rc<RefCell<Vec<DeviceProps>>> = Rc::new(RefCell::new(Vec::new()));
let bound_for_global = Rc::clone(&bound);
let observed_for_global = Rc::clone(&observed);
let registry_weak = registry.downgrade();
let _listener = registry
.add_listener_local()
.global(move |obj| {
if obj.type_ != ObjectType::Device {
return;
}
let Some(registry) = registry_weak.upgrade() else {
return;
};
let Ok(device) = registry.bind::<pw::device::Device, _>(obj) else {
return;
};
let observed_for_info = Rc::clone(&observed_for_global);
let listener = device
.add_listener_local()
.info(move |info| {
if let Some(props) = info.props() {
observed_for_info
.borrow_mut()
.push(device_props_from_props(props));
}
})
.register();
bound_for_global.borrow_mut().push((device, listener));
})
.register();
// Two seconds is the same budget the observer gives its own binds.
let main_loop_for_timer = main_loop.clone();
let timer = main_loop
.loop_()
.add_timer(move |_| main_loop_for_timer.quit());
timer
.update_timer(Some(Duration::from_secs(2)), None)
.into_result()
.expect("arm the test deadline");
main_loop.run();
let observed = observed.borrow();
assert!(
!observed.is_empty(),
"no Device delivered info props at all — the Device bind path is dead"
);
assert!(
observed.iter().any(|props| {
props.device_api.as_deref() == Some("alsa") && props.alsa_driver_name.is_some()
}),
"no bound Device carried both device.api=alsa and alsa.driver_name; \
observed: {observed:?}"
);
}
// Phase 3r exit-gate row 1: failure means the observation boundary regressed.
#[test]
#[ignore = "needs live pipewire"]
fn live_bound_properties_recover_node_and_device_inputs() {
pw::init();
let observer = RegistryObserverHandle::spawn().expect("observer thread must spawn");
wait_for(&observer, |projection| projection.graph_ready);
let unique = format!("pixelpass_observer_props_test_{}", std::process::id());
let capture_name = format!("{unique}_capture");
let playback_name = format!("{unique}_playback");
let null_sink = PactlModule::load(
"module-null-sink",
&[
format!("sink_name={unique}"),
"sink_properties=peerspeak.owned=true node.passthrough=true".to_string(),
],
);
let null_sink_id = null_sink.id.expect("null-sink module must have an id");
let loopback = PactlModule::load(
"module-loopback",
&[
format!("source={unique}.monitor"),
format!("sink={unique}"),
format!("source_output_properties=node.name={capture_name}"),
format!("sink_input_properties=node.name={playback_name}"),
],
);
let projection = wait_for(&observer, |projection| {
projection.graph_ready
&& has_node(projection, &unique)
&& has_node(projection, &capture_name)
&& has_node(projection, &playback_name)
});
let tagged_sink = projection
.snapshot
.nodes()
.find(|node| node.name.as_deref() == Some(&unique))
.expect("tagged null sink must be projected");
assert!(tagged_sink.props.peerspeak_owned);
assert!(tagged_sink.props.passthrough);
assert_eq!(
tagged_sink.props.pulse_module_id,
Some(u64::from(null_sink_id))
);
let capture = projection
.snapshot
.nodes()
.find(|node| node.name.as_deref() == Some(&capture_name))
.expect("loopback capture leg must be projected");
let playback = projection
.snapshot
.nodes()
.find(|node| node.name.as_deref() == Some(&playback_name))
.expect("loopback playback leg must be projected");
let capture_group = capture
.props
.link_group
.as_ref()
.expect("loopback capture leg must carry node.link-group");
let playback_group = playback
.props
.link_group
.as_ref()
.expect("loopback playback leg must carry node.link-group");
assert_eq!(capture_group, playback_group);
assert!(
projection
.snapshot
.nodes()
.any(|node| node.props.process_id.is_some()),
"at least one projected node must carry application.process.id"
);
let session_device = projection.snapshot.nodes().find(|node| {
node.props.session_device
&& (node
.name
.as_deref()
.is_some_and(|name| name.contains("alsa"))
|| matches!(node.role, MediaRole::Sink | MediaRole::Source))
});
assert!(
session_device.is_some(),
"a named ALSA or Audio/Sink/Audio/Source node must classify as a session device"
);
assert!(projection.graph_ready);
loopback.unload();
null_sink.unload();
wait_for(&observer, |projection| {
!has_node(projection, &unique)
&& !has_node(projection, &capture_name)
&& !has_node(projection, &playback_name)
});
}
#[test]
#[ignore = "needs live pipewire"]
fn live_topology_diff_tracks_null_sink_and_loopback() {
+104 -37
View File
@@ -18,11 +18,21 @@
//! both. So the discriminator is `factory.name` on an **allowlist** of
//! real hardware-PCM factories, never a substring or a denylist: an unknown
//! factory is not a device.
//! - The backing Device must actually have been observed. A node that claims
//! a `device.id` we have not yet resolved is **withheld**, not admitted with
//! a provisional `false` — a provisional `false` during the not-ready
//! window fuses sink and mic on the shared session client and that fusion
//! can persist as sticky over-exclusion (round-3 finding 3).
//! - The backing Device must actually have been **bound and resolved**. A node
//! that claims a `device.id` whose Device's properties we do not hold is
//! **withheld**, not admitted with a provisional `false` — a provisional
//! `false` during the not-ready window fuses sink and mic on the shared
//! session client and that fusion can persist as sticky over-exclusion
//! (round-3 finding 3).
//!
//! **Round 8 (v3.5 §6.7 decision 4): the Device is the authority on
//! `device.api` and `alsa.driver_name`.** Both are absent from the Node
//! *global* and both are present on the **bound Device**'s `info` props
//! (measured 2026-07-25). Reading them from the Device closes the phase-3
//! review's owed fix: on PipeWire ≥ 1.2.6 with WirePlumber < 0.5.13 the driver
//! name is not copied onto the node, and the fail-closed "absent driver ⇒ not
//! a session device" rule would over-exclude real sound cards. `factory.name`
//! exists only on the node, which is why the node bind is required regardless.
use crate::host::taint::snapshot::GlobalId;
@@ -54,6 +64,11 @@ const HARDWARE_PCM_FACTORIES: &[&str] = &[
"api.alsa.pcm.source",
];
/// The `device.api` every entry in [`HARDWARE_PCM_FACTORIES`] belongs to.
/// A single value rather than a list, because the allowlist is ALSA-only;
/// this constant is the thing to change when that stops being true.
const HARDWARE_PCM_API: &str = "alsa";
/// ALSA drivers that expose a hardware-PCM `factory.name` but are **not**
/// passive terminals — audio written in reappears on their capture side
/// through a path the PipeWire Link graph cannot see, so classifying them
@@ -67,8 +82,9 @@ const HARDWARE_PCM_FACTORIES: &[&str] = &[
/// does not couple playback to capture, so it is not a loopback hazard.
const NON_TERMINAL_ALSA_DRIVERS: &[&str] = &["snd_aloop"];
/// The three node properties the classifier reads, exactly as the adapter
/// parsed them off the Node global. Kept separate from
/// The node-side properties the classifier reads, exactly as the adapter
/// parsed them off the **bound Node's `info`** (never off the registry
/// global — v3.5 §6.7). Kept separate from
/// [`super::super::taint::snapshot::NodeProps`] because these feed the
/// *decision* whose output is the `session_device` field — they are inputs,
/// not part of the graph the engine reasons over.
@@ -78,10 +94,12 @@ pub struct DeviceClaim {
/// `Stream/*` nodes, which is exactly why their absence means "not a
/// device", not "unknown".
pub device_id: Option<GlobalId>,
/// `device.api` — the access API of that Device (e.g. `alsa`, `bluez5`).
/// Its mere presence is **not** sufficient (a card-associated filter has
/// it too); required only as a corroborating signal alongside the factory
/// allowlist.
/// `device.api` **as copied onto the node**, when it is — the access API
/// of that Device (e.g. `alsa`, `bluez5`). Its mere presence is **not**
/// sufficient (a card-associated filter has it too); required only as a
/// corroborating signal alongside the factory allowlist. The
/// authoritative copy is [`DeviceProps::device_api`]; this is the
/// fallback.
pub device_api: Option<String>,
/// `factory.name` — the discriminator. Only an allowlisted hardware-PCM
/// factory earns `session_device`.
@@ -92,8 +110,27 @@ pub struct DeviceClaim {
/// shares the same factory. `session_device` requires this to be
/// **present and not** on [`NON_TERMINAL_ALSA_DRIVERS`]; a driver on the
/// denylist, or an absent value, both fail closed (see [`classify`]).
/// May be absent on non-ALSA backends or on version pairings that do not
/// copy `alsa.*` onto the node.
/// Frequently absent here — PipeWire ≥ 1.2.6 with WirePlumber < 0.5.13
/// does not copy `alsa.*` onto the node — which is why the authoritative
/// copy is [`DeviceProps::alsa_driver_name`] and this is only the
/// fallback.
pub alsa_driver_name: Option<String>,
}
/// The **bound Device's** `info` properties — the authoritative half of the
/// `session_device` decision (v3.5 §6.7 decision 4).
///
/// Absent from the Device *registry global* exactly as the node's properties
/// are absent from the Node global; both are recovered by binding. A node
/// claiming a `device.id` is withheld until this struct exists for that
/// Device (see [`Classification::Withhold`]).
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct DeviceProps {
/// `device.api` on the Device — `alsa`, `bluez5`, `v4l2`, …
pub device_api: Option<String>,
/// `alsa.driver_name` on the Device — the kernel driver behind the card,
/// authoritative regardless of whether the session manager copied it onto
/// the node.
pub alsa_driver_name: Option<String>,
}
@@ -102,9 +139,10 @@ pub struct DeviceClaim {
pub enum Classification {
/// No `device.id` — a `Stream/*` node. Admit with `session_device=false`.
NotADevice,
/// A `device.id` is claimed but the backing Device has not been resolved
/// yet. **Withhold the node and keep the readiness epoch not-ready**;
/// re-classify when the Device is observed.
/// A `device.id` is claimed but the backing Device's properties are not
/// held: never observed, its bind still outstanding, or its global id
/// ambiguously shared by two live Devices. **Withhold the node and keep
/// the readiness epoch not-ready**; re-classify when the Device resolves.
Withhold { device_id: GlobalId },
/// Positively a passive hardware terminal. Admit with
/// `session_device=true`.
@@ -115,42 +153,71 @@ pub enum Classification {
NotSessionDevice,
}
/// Classify a node's device claim.
/// Classify a node's device claim against its backing Device.
///
/// `device_resolved` is whether [`DeviceClaim::device_id`] has been observed
/// as a Device global; it is only consulted when a `device_id` is present.
/// Pure: the model supplies `device_resolved` from its resolved-Device set,
/// and the I/O of *binding* the Device lives in the adapter.
pub fn classify(claim: &DeviceClaim, device_resolved: bool) -> Classification {
/// `device` is the bound Device's properties, and `None` means the claim is
/// **unresolved** — never observed, bind outstanding, or an ambiguous
/// recycled id. It is only consulted when a `device_id` is present. Pure: the
/// model looks the Device up, and the I/O of *binding* it lives in the
/// adapter.
///
/// Where the two sides disagree the rule is deliberately asymmetric, and
/// safety picks the direction (v3.5 §6.7 decision 4):
///
/// - **Presence: the Device wins, the node is the fallback.** That is what
/// recovers a real card whose node was never given `alsa.driver_name`.
/// - **The denylist is a union.** If *either* side names a non-terminal
/// driver the node is not a session device. A disagreement here is not
/// expected on any measured configuration, and treating it as "the Device
/// says it is fine" would be the one reading that can leak.
pub fn classify(claim: &DeviceClaim, device: Option<&DeviceProps>) -> Classification {
let Some(device_id) = claim.device_id else {
// No backing Device: a stream. Not withheld, not a device.
return Classification::NotADevice;
};
if !device_resolved {
// Backed by a Device we have not seen — the one case that blocks
let Some(device) = device else {
// Backed by a Device we have not resolved — the one case that blocks
// readiness. A provisional answer here is the leak the contract
// forbids.
return Classification::Withhold { device_id };
}
};
let on_factory_allowlist = claim
.factory_name
.as_deref()
.is_some_and(|f| HARDWARE_PCM_FACTORIES.contains(&f));
// A **present, non-denied** ALSA driver is required — absence fails closed
// (Codex phase-3 re-review). `alsa.driver_name` is not copied onto the
// node on every PipeWire/WirePlumber version pairing (PipeWire ≥1.2.6
// stopped overwriting node props with card props; WirePlumber only began
// copying `alsa.*` onto nodes in 0.5.13), so a *missing* value must not be
// read as "not a loopback" — that is exactly the hole an `snd_aloop` node
// without the property would slip through. A real card whose node lacks
// the driver is instead over-excluded (keeps its owner keys — safe);
// recovering `session_device` for it needs reading the driver from the
// backing Device global, which is owed to a later round.
let driver_ok = claim
// (Codex phase-3 re-review). The factory allowlist cannot tell a real card
// from `snd_aloop`, which presents the same `api.alsa.pcm.*` factory, so a
// *missing* value must not be read as "not a loopback". Round 8 makes the
// bound Device the primary source, so a real card is no longer
// over-excluded merely because the session manager did not copy `alsa.*`
// onto its node.
let driver = device
.alsa_driver_name
.as_deref()
.is_some_and(|d| !NON_TERMINAL_ALSA_DRIVERS.contains(&d));
let is_hardware_pcm = claim.device_api.is_some() && on_factory_allowlist && driver_ok;
.or(claim.alsa_driver_name.as_deref());
let driver_denied = [
device.alsa_driver_name.as_deref(),
claim.alsa_driver_name.as_deref(),
]
.into_iter()
.flatten()
.any(|d| NON_TERMINAL_ALSA_DRIVERS.contains(&d));
let driver_ok = driver.is_some() && !driver_denied;
// The API must positively be the one the factory allowlist is written
// for, not merely present (Codex phase-3r review, finding 3). "Present"
// admitted `device.api=v4l2` alongside `factory.name=api.alsa.pcm.sink`
// — a contradiction no truthful configuration produces, which is exactly
// why it should be read as an observation gone wrong rather than as
// corroboration. Disagreement between the two sides fails closed for the
// same reason. ⚠️ Tied to [`HARDWARE_PCM_FACTORIES`] being ALSA-only:
// adding a BlueZ factory means allowing `bluez5` here too.
let api_ok = match (device.device_api.as_deref(), claim.device_api.as_deref()) {
(Some(from_device), Some(from_node)) if from_device != from_node => false,
(Some(api), _) | (None, Some(api)) => api == HARDWARE_PCM_API,
(None, None) => false,
};
let is_hardware_pcm = api_ok && on_factory_allowlist && driver_ok;
if is_hardware_pcm {
Classification::SessionDevice
} else {
+330 -127
View File
@@ -1,12 +1,34 @@
//! The registry observer's **pure core** (impl plan §4, phase 3).
//! The registry observer's **pure core** (impl plan §4, phases 3 and 3r).
//!
//! This is my half of the phase-3 split: a reducer that folds a stream of
//! typed [`RegEvent`]s into a live model of the PipeWire graph and projects
//! the [`GraphSnapshot`] + context the taint engine (phase 2) consumes. **No
//! PipeWire types appear here** — the I/O adapter (Codex's half) translates
//! live registry callbacks, Link/Device binds, `/proc` reads, and the
//! `core.sync`/`done` round-trip into these events and feeds them in. Every
//! test in this module builds the event stream by hand.
//! live registry callbacks, binds, `/proc` reads, and the `core.sync`/`done`
//! round-trip into these events and feeds them in. Every test in this module
//! builds the event stream by hand.
//!
//! ## 🔴 Round 8 (v3.5 §6.7): the global is an INDEX, not a source of truth
//!
//! Phase 3 shipped reading node properties off the registry `global` event.
//! The registry announces only a fixed 13-key subset for a Node, and **eight
//! properties this feature depends on are never among them** — they read as
//! absent rather than failing, so the engine was silently, permanently
//! starved of both its primary taint root and every strong owner key (the
//! phase-5 gate failure, F1/F2). The rule that replaces it:
//!
//! > A node's properties come from a **bind**, never from the global. The
//! > global tells us an object exists, its id and its serial. Everything
//! > else — including `node.name` and `media.class`, so there is exactly one
//! > source — arrives on [`RegEvent::NodeInfo`]. Same for `Device`
//! > ([`RegEvent::DeviceInfo`]).
//!
//! Consequences visible in this file: a Node is admitted to the snapshot
//! **only** once its `info` has arrived (until then it is withheld and is a
//! readiness obligation); a Device resolves a node's claim only once *its*
//! `info` has arrived; and `info` may fire again for the lifetime of the
//! object, so [`RegEvent::NodeInfo`] is both the first resolution and every
//! later property change (v3.5 §6.7 decisions 14).
//!
//! Three things this core is shaped to get right, each an exit-gate row:
//!
@@ -14,18 +36,20 @@
//! id, and those recycle. The model keeps an insertion-ordered index per id
//! so a removal accounts for the *oldest* generation first, and the
//! snapshot projection treats any id still claimed by two live objects as
//! [`IdLookup::Ambiguous`] — fail closed (v3.4 §6.1.3).
//! [`IdLookup::Ambiguous`] — fail closed (v3.4 §6.1.3). Everything the
//! model *owns* is keyed by never-recycled `object.serial`; ids are only
//! ever a lookup.
//! - **The readiness epoch.** `graph_ready` is false until the initial graph
//! is fully observed: the server has synced **and** no binds/withheld nodes
//! remain outstanding. A bounded timeout makes it fail closed. It gates
//! sticky *retirement* only; withholding after completion is per-object.
//! - **Withholding on unresolved devices.** A node claiming a `device.id`
//! whose Device we have not observed is held out of the snapshot entirely
//! rather than admitted with a provisional `session_device` (see
//! [`classify`]).
//! - **Withholding on unresolved input.** A node with no `info` yet, or one
//! claiming a `device.id` whose Device we have not resolved, is held out of
//! the snapshot entirely rather than admitted with provisional ownership
//! (see [`classify`]).
//!
//! **Two accepted limitations (Codex phase-3 review, findings 3 and 4), both
//! low-reachability, owed to a later hardening round:**
//! **Three accepted limitations, all low-reachability, owed to a later
//! hardening round:**
//!
//! - *A Link dropped for a missing `object.serial`/props is unrepresented.*
//! The adapter drops such a global before it reaches [`RegistryModel`], so
@@ -45,6 +69,13 @@
//! silently drop `global_remove`, so this needs callback loss to trigger.
//! The snapshot treats the two-claimant window as [`IdLookup::Ambiguous`]
//! (fail closed) meanwhile.
//! - *An unresolvable bind takes the whole graph down, not just its node*
//! (v3.5 §6.7 decision 3). A node whose `info` never arrives keeps
//! readiness false until the deadline, then sticky-[`Readiness::TimedOut`]
//! — no fan-out at all, identical to a never-resolving Link bind. Per-node
//! quarantine (that node ineligible **and** taint-bearing, the rest of the
//! graph still working) is strictly better and is deferred because it is a
//! new concept in the *pure engine*, not a fix to the observer.
#![allow(dead_code)] // Wired by the phase-3 adapter (Codex's half) and consumed by later phases.
@@ -59,7 +90,7 @@ use crate::host::taint::snapshot::{
ClientSnapshot, GlobalId, GraphSnapshot, LinkSnapshot, MediaRole, NodeProps, NodeSnapshot,
PortSnapshot, Serial,
};
use classify::{Classification, DeviceClaim};
use classify::{Classification, DeviceClaim, DeviceProps};
use std::collections::{BTreeMap, VecDeque};
/// A monotonic millisecond clock value, supplied by the adapter via
@@ -67,14 +98,17 @@ use std::collections::{BTreeMap, VecDeque};
/// [`std::time::Instant`] so the readiness timeout is deterministic in tests.
pub type Millis = u64;
/// A Node as observed off the registry, before `session_device` has been
/// decided. The adapter fills [`NodeProps`] with everything it can parse and
/// leaves `session_device` at its `false` default; the model overwrites it
/// from the [`classify`] result once the backing Device (if any) is resolved.
/// A Node's **bound `info` properties** — the sole source of node properties
/// (v3.5 §6.7), delivered by [`RegEvent::NodeInfo`].
///
/// This carries no identity: the serial names the node on the event and the
/// global id was recorded by [`RegEvent::NodeAdded`], so the adapter cannot
/// contradict the index it already published. `session_device` inside
/// [`NodeObservation::props`] is left at its `false` default; the model
/// overwrites it from the [`classify`] result at projection time, once the
/// backing Device (if any) is resolved.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct NodeObservation {
pub serial: Serial,
pub id: GlobalId,
pub name: Option<String>,
pub role: MediaRole,
pub props: NodeProps,
@@ -97,19 +131,39 @@ pub struct LinkEndpoints {
/// model consumes them in [`RegistryModel::apply`].
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RegEvent {
/// A Node global appeared. Admitted immediately unless it claims an
/// unresolved Device (then withheld — see [`classify`]).
NodeAdded(NodeObservation),
/// A Node global appeared. **Index only** — the global's properties are a
/// filtered subset and are not read (v3.5 §6.7). The node is withheld
/// from the snapshot and is a readiness obligation until its
/// [`RegEvent::NodeInfo`] arrives.
NodeAdded { serial: Serial, id: GlobalId },
/// A bound Node's `info` properties. **Both** the first resolution and
/// every later `PROPS` change for the node's lifetime — the model tells
/// them apart, so the adapter holds no per-node "have I seen info yet?"
/// state to get wrong. An `info` for a serial we do not hold (a node
/// already removed) is ignored.
NodeInfo {
serial: Serial,
observation: NodeObservation,
},
/// A Port global appeared.
PortAdded(PortSnapshot),
/// A Client global appeared. Feeds pulse-PID derivation via `sec_pid`.
ClientAdded(ClientSnapshot),
/// A Device global appeared. Resolves any nodes withheld on its id.
DeviceAdded { id: GlobalId },
/// A Device global appeared. Index only, exactly as for a Node: it does
/// not resolve anything until [`RegEvent::DeviceInfo`] arrives.
DeviceAdded { serial: Serial, id: GlobalId },
/// A bound Device's `info` properties — the **authoritative** source of
/// `device.api` and `alsa.driver_name` (v3.5 §6.7 decision 4). Resolves
/// every node withheld on this Device's id.
DeviceInfo { serial: Serial, props: DeviceProps },
/// A Link global appeared. `endpoints` is `Some` when the global carried
/// them (the optimisation) and `None` when the adapter must bind to learn
/// them (the correctness path) — the latter is an outstanding obligation
/// until a matching [`RegEvent::LinkEndpointsResolved`] arrives.
///
/// Unlike Nodes and Devices, Link endpoint props **are** announced on the
/// global (measured, phase-5 results F1), so this asymmetry is real and
/// deliberate.
LinkAdded {
serial: Serial,
id: GlobalId,
@@ -143,7 +197,7 @@ pub enum RegEvent {
/// one is worth suppressing on a tick but never on a graph event.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum EventKind {
/// A registry observation: an add, a removal, a link resolution, a `/proc`
/// A registry observation: an add, a removal, a bind resolution, a `/proc`
/// probe, or the server sync.
Graph,
/// The periodic clock sample. Carries no graph information; it exists so the
@@ -169,16 +223,39 @@ impl RegEvent {
}
}
/// Whether an applied event could have changed the projection.
///
/// The suppression rule of v3.5 §6.7 decision 2, in the one place that can
/// enforce it: **a property update may be dropped only when the resulting
/// [`Projection`] is identical to the current one.** The projection is a pure
/// function of model state, so "state provably unchanged" *is* "projection
/// identical" — which is what [`Outcome::Suppressed`] means and why the check
/// is a cheap field comparison rather than building and diffing two snapshots.
///
/// Anything looser (dropping updates that do change state) breaks phase 4's
/// no-coalescing contract, which needs to see the empty gap between an AEC
/// module unload and a reload that reuses the index. Anything stricter
/// (publishing on every `info`, including the state-only changes PipeWire
/// emits constantly) inflates the O5 event rate with non-events.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Outcome {
/// Model state may have changed; the caller must publish the projection.
Applied,
/// Model state provably did not change; publishing is optional and the
/// adapter skips it.
Suppressed,
}
/// Which slot in the id index a live object occupies. `global_remove` gives
/// only the id, so the index remembers what each id currently holds. A Node
/// slot's serial may live in either the admitted or the withheld map.
/// only the id, so the index remembers what each id currently holds. Every
/// slot names its object by never-recycled serial.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Slot {
Node(Serial),
Port(Serial),
Link(Serial),
Client(Serial),
Device,
Device(Serial),
}
/// The readiness epoch. A one-time transition out of [`Readiness::Waiting`];
@@ -220,25 +297,45 @@ pub struct Projection {
pub readiness: Readiness,
}
/// A live Node: its global id (for link endpoint lookup) plus its bound
/// properties once they arrive.
#[derive(Clone, Debug, PartialEq, Eq)]
struct NodeEntry {
id: GlobalId,
/// `None` while the bind is outstanding — withheld from the snapshot and
/// an outstanding readiness obligation (v3.5 §6.7 decision 3).
obs: Option<NodeObservation>,
}
/// A live Device: its global id plus its bound properties once they arrive.
#[derive(Clone, Debug, PartialEq, Eq)]
struct DeviceEntry {
id: GlobalId,
/// `None` while the bind is outstanding. A node claiming this Device
/// stays withheld until it is `Some` — the Device's `device.api` and
/// `alsa.driver_name` are the authoritative inputs to `session_device`
/// (v3.5 §6.7 decision 4), so classifying without them would be the same
/// provisional answer the contract forbids.
props: Option<DeviceProps>,
}
/// The live model. Folds [`RegEvent`]s; project with [`RegistryModel::project`].
#[derive(Clone, Debug)]
pub struct RegistryModel {
// Admitted objects, keyed by their never-recycled serial.
nodes: BTreeMap<Serial, NodeSnapshot>,
/// **Every** live Node, keyed by serial — admitted or withheld. Admission
/// is decided at projection time from the entry's own state, so there is
/// no admitted/withheld pair of maps to drift apart.
nodes: BTreeMap<Serial, NodeEntry>,
/// Every live Device, keyed by serial.
devices: BTreeMap<Serial, DeviceEntry>,
ports: BTreeMap<Serial, PortSnapshot>,
links: BTreeMap<Serial, LinkSnapshot>,
clients: BTreeMap<Serial, ClientSnapshot>,
/// Nodes held out of the snapshot pending their Device's resolution.
withheld: BTreeMap<Serial, NodeObservation>,
/// Links whose endpoints the adapter is still binding; the id is kept so
/// removal and resolution can find them.
pending_links: BTreeMap<Serial, GlobalId>,
/// Live Device global ids, ref-counted so a recycled id is only
/// considered resolved while a Device actually holds it.
resolved_devices: BTreeMap<GlobalId, usize>,
/// Insertion-ordered holders of each live global id. `global_remove`
/// accounts for the oldest generation first (v3.4 §6.1.3).
live_ids: BTreeMap<GlobalId, VecDeque<Slot>>,
@@ -259,12 +356,11 @@ impl RegistryModel {
pub fn new(now: Millis, timeout: Millis) -> Self {
Self {
nodes: BTreeMap::new(),
devices: BTreeMap::new(),
ports: BTreeMap::new(),
links: BTreeMap::new(),
clients: BTreeMap::new(),
withheld: BTreeMap::new(),
pending_links: BTreeMap::new(),
resolved_devices: BTreeMap::new(),
live_ids: BTreeMap::new(),
probed_comm: BTreeMap::new(),
server_synced: false,
@@ -283,21 +379,22 @@ impl RegistryModel {
///
/// This is **dynamic**, not the sticky [`Readiness::Complete`] flag: it is
/// true only when the initial enumeration has completed **and** there are
/// no current obligations outstanding (a node withheld on an unresolved
/// Device, or a Link still being bound). The distinction is the fix for
/// Codex phase-3 review finding 1: a Link whose endpoints are still
/// resolving is an **invisible edge** — it is absent from the snapshot,
/// not merely dangling — so a decision made while one exists can miss real
/// tainted ancestry and wrongly report a candidate eligible. Unresolved
/// ancestry ⇒ fail closed is the governing invariant (v3.4 §6.1), and an
/// unresolved Link is unresolved ancestry, so `graph_ready` must drop back
/// to false whenever one is pending — even after the initial epoch.
/// no current obligations outstanding (a node whose bind is outstanding, a
/// node withheld on an unresolved Device, or a Link still being bound).
/// The distinction is the fix for Codex phase-3 review finding 1: a Link
/// whose endpoints are still resolving is an **invisible edge** — it is
/// absent from the snapshot, not merely dangling — so a decision made
/// while one exists can miss real tainted ancestry and wrongly report a
/// candidate eligible. Unresolved ancestry ⇒ fail closed is the governing
/// invariant (v3.4 §6.1), and round 8 adds the far more common case: an
/// unbound node is an invisible *vertex*, which hides everything the edge
/// case hides and its ownership besides.
///
/// [`Readiness::Complete`] stays sticky (it records that the initial
/// enumeration happened, for logging and to distinguish "not started" from
/// "momentarily churning"); `graph_ready` layers the dynamic obligation
/// check on top. Downstream (phase 6) may debounce the brief blips a
/// normal Link bind causes; the observer's job is to report the truth.
/// normal bind causes; the observer's job is to report the truth.
pub fn graph_ready(&self) -> bool {
matches!(self.readiness, Readiness::Complete) && !self.obligations_outstanding()
}
@@ -310,102 +407,120 @@ impl RegistryModel {
pulse_pid::candidate(&clients)
}
/// Fold one observation into the model.
pub fn apply(&mut self, event: RegEvent) {
/// Fold one observation into the model. The returned [`Outcome`] tells the
/// caller whether the projection can have changed; see [`Outcome`] for why
/// that is the only sound place to enforce the suppression rule.
pub fn apply(&mut self, event: RegEvent) -> Outcome {
match event {
RegEvent::NodeAdded(obs) => self.on_node_added(obs),
RegEvent::NodeAdded { serial, id } => {
self.push_id(id, Slot::Node(serial));
self.nodes.insert(serial, NodeEntry { id, obs: None });
// A node awaiting its bind is a fresh obligation, so this can
// only ever *hold* readiness, never complete it — but the
// re-check is cheap and keeps the invariant local.
self.maybe_complete();
Outcome::Applied
}
RegEvent::NodeInfo {
serial,
observation,
} => self.on_node_info(serial, observation),
RegEvent::PortAdded(port) => {
self.push_id(port.id, Slot::Port(port.serial));
self.ports.insert(port.serial, port);
Outcome::Applied
}
RegEvent::ClientAdded(client) => {
self.push_id(client.id, Slot::Client(client.serial));
self.clients.insert(client.serial, client);
// A new client can change the pulse candidate; the adapter
// learns that via `pulse_pid_candidate`. No readiness effect.
Outcome::Applied
}
RegEvent::DeviceAdded { id } => self.on_device_added(id),
RegEvent::DeviceAdded { serial, id } => {
self.push_id(id, Slot::Device(serial));
self.devices.insert(serial, DeviceEntry { id, props: None });
self.maybe_complete();
Outcome::Applied
}
RegEvent::DeviceInfo { serial, props } => self.on_device_info(serial, props),
RegEvent::LinkAdded {
serial,
id,
endpoints,
} => self.on_link_added(serial, id, endpoints),
} => {
self.on_link_added(serial, id, endpoints);
Outcome::Applied
}
RegEvent::LinkEndpointsResolved { serial, endpoints } => {
self.on_link_resolved(serial, endpoints)
}
RegEvent::ProcCommProbed { pid, comm } => {
self.probed_comm.insert(pid, comm);
let previous = self.probed_comm.insert(pid, comm.clone());
if previous.as_ref() == Some(&comm) {
Outcome::Suppressed
} else {
Outcome::Applied
}
}
RegEvent::Removed { id } => self.on_removed(id),
RegEvent::ServerSynced => {
let already = self.server_synced;
self.server_synced = true;
self.maybe_complete();
if already {
Outcome::Suppressed
} else {
Outcome::Applied
}
}
RegEvent::Tick { now } => {
self.last_now = now;
self.maybe_timeout(now);
Outcome::Applied
}
}
}
fn on_node_added(&mut self, obs: NodeObservation) {
self.push_id(obs.id, Slot::Node(obs.serial));
let resolved = obs
.device_claim
.device_id
.is_some_and(|id| self.device_resolved(id));
match classify::classify(&obs.device_claim, resolved) {
Classification::Withhold { .. } => {
self.withheld.insert(obs.serial, obs);
}
Classification::SessionDevice => self.admit_node(obs, true),
Classification::NotADevice | Classification::NotSessionDevice => {
self.admit_node(obs, false)
}
/// First resolution *and* every later property change (v3.5 §6.7
/// decision 2). The model distinguishes them by what it already holds, so
/// the adapter can forward every `info` callback unconditionally.
fn on_node_info(&mut self, serial: Serial, observation: NodeObservation) -> Outcome {
let Some(entry) = self.nodes.get_mut(&serial) else {
// A late `info` for a node already removed. Re-inserting it here
// would resurrect a dead node with no id index behind it.
tracing::debug!(serial = serial.0, "observer: node info for an unknown node");
return Outcome::Suppressed;
};
if entry.obs.as_ref() == Some(&observation) {
// The state-only `info` callbacks PipeWire emits constantly: same
// properties, so the projection is provably identical.
return Outcome::Suppressed;
}
// Withholding a node adds an obligation; admitting one can never
// complete readiness on its own, but re-check is cheap and keeps the
// invariant local.
entry.obs = Some(observation);
// The first `info` retires this node's obligation, which can be the
// last one outstanding.
self.maybe_complete();
Outcome::Applied
}
fn admit_node(&mut self, obs: NodeObservation, session_device: bool) {
let mut props = obs.props;
props.session_device = session_device;
self.nodes.insert(
obs.serial,
NodeSnapshot {
serial: obs.serial,
id: obs.id,
name: obs.name,
role: obs.role,
props,
},
);
}
fn on_device_added(&mut self, id: GlobalId) {
self.push_id(id, Slot::Device);
*self.resolved_devices.entry(id).or_insert(0) += 1;
// Admit every node that was withheld waiting on exactly this Device.
let ready: Vec<Serial> = self
.withheld
.iter()
.filter(|(_, obs)| obs.device_claim.device_id == Some(id))
.map(|(&serial, _)| serial)
.collect();
for serial in ready {
if let Some(obs) = self.withheld.remove(&serial) {
// Resolved now, so classify yields a terminal answer, never
// Withhold again.
let session_device = matches!(
classify::classify(&obs.device_claim, true),
Classification::SessionDevice
);
self.admit_node(obs, session_device);
}
fn on_device_info(&mut self, serial: Serial, props: DeviceProps) -> Outcome {
let Some(entry) = self.devices.get_mut(&serial) else {
tracing::debug!(
serial = serial.0,
"observer: device info for an unknown device"
);
return Outcome::Suppressed;
};
if entry.props.as_ref() == Some(&props) {
return Outcome::Suppressed;
}
entry.props = Some(props);
// Resolving a Device admits every node that was withheld on it —
// which happens at projection time; here it can only retire
// obligations.
self.maybe_complete();
Outcome::Applied
}
fn on_link_added(&mut self, serial: Serial, id: GlobalId, endpoints: Option<LinkEndpoints>) {
@@ -423,20 +538,23 @@ impl RegistryModel {
self.maybe_complete();
}
fn on_link_resolved(&mut self, serial: Serial, endpoints: LinkEndpoints) {
fn on_link_resolved(&mut self, serial: Serial, endpoints: LinkEndpoints) -> Outcome {
// `remove` also guards against a stale resolution for a Link already
// gone: unknown serial ⇒ ignore.
if let Some(id) = self.pending_links.remove(&serial) {
self.links
.insert(serial, link_snapshot(serial, id, endpoints));
self.maybe_complete();
Outcome::Applied
} else {
Outcome::Suppressed
}
}
fn on_removed(&mut self, id: GlobalId) {
fn on_removed(&mut self, id: GlobalId) -> Outcome {
let Some(queue) = self.live_ids.get_mut(&id) else {
tracing::warn!(global_id = id.0, "observer: remove for an id we never saw");
return;
return Outcome::Suppressed;
};
// Oldest generation first — the id may be shared during a
// missed-removal window.
@@ -446,10 +564,7 @@ impl RegistryModel {
}
match slot {
Some(Slot::Node(serial)) => {
if self.nodes.remove(&serial).is_none() {
// Was still withheld — drop the obligation.
self.withheld.remove(&serial);
}
self.nodes.remove(&serial);
}
Some(Slot::Port(serial)) => {
self.ports.remove(&serial);
@@ -461,35 +576,77 @@ impl RegistryModel {
Some(Slot::Client(serial)) => {
self.clients.remove(&serial);
}
Some(Slot::Device) => {
if let Some(count) = self.resolved_devices.get_mut(&id) {
*count -= 1;
if *count == 0 {
self.resolved_devices.remove(&id);
}
}
Some(Slot::Device(serial)) => {
self.devices.remove(&serial);
}
None => {
tracing::warn!(global_id = id.0, "observer: empty id slot on remove");
return Outcome::Suppressed;
}
}
// A removal can drain the last obligation (a withheld node or pending
// link vanished before it resolved).
// A removal can drain the last obligation (an unbound node, a node
// withheld on a Device, or a pending link vanished before it
// resolved).
self.maybe_complete();
Outcome::Applied
}
fn push_id(&mut self, id: GlobalId, slot: Slot) {
self.live_ids.entry(id).or_default().push_back(slot);
}
fn device_resolved(&self, id: GlobalId) -> bool {
self.resolved_devices.get(&id).is_some_and(|&n| n > 0)
/// The bound properties of the Device a node claims by global id, or
/// `None` when that claim is unresolved — which covers every fail-closed
/// case at once: no such Device observed, its bind still outstanding, or
/// **the id claimed by more than one live global**, where there is no way
/// to tell whose properties these are (v3.4 §6.1.3).
///
/// ⚠️ The ambiguity test is "**exactly one** live global holds this id",
/// not "exactly one live *Device*" (Codex phase-3r review, finding 2).
/// The weaker test looks equivalent and is not: with `[Device, Port]` on
/// one id — a missed removal, the same precondition as every other
/// recycled-id hazard — it keeps answering with the older Device's
/// properties, so a node claiming that id holds a stale
/// `session_device = true`. That flag *removes* the node's owner keys and
/// its fail-closed backstop, so a forwarder wearing it can put its output
/// leg back on the eligible side: echo, from a lookup that was merely
/// looking at the wrong object type.
fn device_props(&self, id: GlobalId) -> Option<&DeviceProps> {
let slots = self.live_ids.get(&id)?;
if slots.len() != 1 {
return None; // Ambiguous ⇒ unresolved ⇒ withheld.
}
let Slot::Device(serial) = slots.front()? else {
// The id is live, but it is not a Device any more.
return None;
};
self.devices.get(serial)?.props.as_ref()
}
/// Classify one node's device claim against the currently resolved
/// Devices. Recomputed per projection rather than cached at admission:
/// the inputs (this node's props, its Device's props) both change over an
/// object's lifetime now, and a cached classification is exactly the kind
/// of stale provisional answer §6.1.3 forbids.
fn classification(&self, obs: &NodeObservation) -> Classification {
let device = obs
.device_claim
.device_id
.and_then(|id| self.device_props(id));
classify::classify(&obs.device_claim, device)
}
/// Every obligation that must clear before the initial graph is trusted:
/// no node withheld on an unresolved Device, no Link awaiting its bind.
/// no node awaiting its bind, no node withheld on an unresolved Device,
/// no Link awaiting its bind.
fn obligations_outstanding(&self) -> bool {
!self.withheld.is_empty() || !self.pending_links.is_empty()
if !self.pending_links.is_empty() {
return true;
}
self.nodes.values().any(|entry| match &entry.obs {
None => true,
Some(obs) => matches!(self.classification(obs), Classification::Withhold { .. }),
})
}
/// Completion needs no clock — only the sync flag and an empty obligation
@@ -512,13 +669,34 @@ impl RegistryModel {
if now >= self.deadline {
self.readiness = Readiness::TimedOut;
tracing::warn!(
withheld = self.withheld.len(),
unbound_nodes = self.unbound_node_count(),
withheld = self.withheld_node_count(),
pending_links = self.pending_links.len(),
"observer: readiness epoch timed out with obligations outstanding — fail closed"
);
}
}
/// Nodes whose bind has not delivered `info` yet — diagnostics only.
fn unbound_node_count(&self) -> usize {
self.nodes
.values()
.filter(|entry| entry.obs.is_none())
.count()
}
/// Nodes held out on an unresolved Device — diagnostics only.
fn withheld_node_count(&self) -> usize {
self.nodes
.values()
.filter(|entry| {
entry.obs.as_ref().is_some_and(|obs| {
matches!(self.classification(obs), Classification::Withhold { .. })
})
})
.count()
}
/// pipewire-pulse's PID from the current clients, validated against the
/// probed `comm`. `None` whenever anything is ambiguous or unconfirmed —
/// the safe answer (key 4 unusable).
@@ -529,9 +707,34 @@ impl RegistryModel {
}
/// Project the current state into the taint engine's inputs.
///
/// A node enters the snapshot only if its bind has delivered `info`
/// **and** its device claim classifies terminally; anything else is
/// withheld (and is already holding `graph_ready` false).
pub fn project(&self) -> Projection {
let nodes: Vec<NodeSnapshot> = self
.nodes
.iter()
.filter_map(|(&serial, entry)| {
let obs = entry.obs.as_ref()?;
let session_device = match self.classification(obs) {
Classification::Withhold { .. } => return None,
Classification::SessionDevice => true,
Classification::NotADevice | Classification::NotSessionDevice => false,
};
let mut props = obs.props.clone();
props.session_device = session_device;
Some(NodeSnapshot {
serial,
id: entry.id,
name: obs.name.clone(),
role: obs.role,
props,
})
})
.collect();
let snapshot = GraphSnapshot::new(
self.nodes.values().cloned().collect(),
nodes,
self.ports.values().cloned().collect(),
self.links.values().cloned().collect(),
self.clients.values().cloned().collect(),
+821 -109
View File
File diff suppressed because it is too large Load Diff