Author SHA1 Message Date
mollusk 64f98990c8 host/taint: uncertainty is not history — it never enters sticky state
Found by the phase-5 audit on the live graph, immediately after phase 3r
landed: a hardware sink carried a permanent `unresolved-ancestry` taint. The
cause was one link observed while its output node was still unbound — a
correct fail-closed answer — which was then written into sticky state, where
retirement requires every member object to be absent. A live sound card never
is, so the mark survived readiness, 21 recomputes and deliberate churn.

Phase 3r makes this systematic rather than rare: every node is now withheld
until its bind resolves, so any link seen across that gap raises
`UnresolvedAncestry` on its input side. It fires at startup, every startup.

User decision (2026-07-25): uncertainty-based taint retires once the
uncertainty is gone; evidence-based taint keeps the absence rule.

Retiring by reason *code* would not be enough, because uncertainty launders
itself — an unresolved node propagates `TaintedUpstream`, which is
indistinguishable from real contamination once recorded. So the split is by
**provenance**: `evaluate` runs the fixpoint twice. Pass 1 fails closed
exactly as before and is what every decision is made from; pass 2 raises no
uncertainty root at all, and is the only thing sticky state is built from.
Nothing derived from an uncertainty can reach the sticky path.

Decisions are unchanged by construction — all 57 existing taint tests pass
untouched, including the fail-closed and sticky-survival rows.

3 new tests, mutation-verified (pointing `build_sticky` back at the
fail-closed taint kills exactly the two new uncertainty tests and nothing
else): unresolved ancestry clears once resolved; taint laundered downstream
of an uncertainty clears with it; real taint still survives its topology
disappearing.

Live: the audit's post-readiness records now report taint 0 where they
reported a permanent sticky entry before. Recompute cost roughly doubles as
expected (two fixpoints) — 80 µs worst case observed, against a 47 Hz event
rate.

`taint/tests.rs` keeps its one pre-existing hand-formatted line; everything
else in both files is rustfmt-clean.
2026-07-25 18:46:56 -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
6 changed files with 1829 additions and 377 deletions
+429 -70
View File
@@ -4,8 +4,10 @@
//! callbacks into [`RegEvent`]s, and publishes the latest [`Projection`] for //! callbacks into [`RegEvent`]s, and publishes the latest [`Projection`] for
//! consumers running outside the PipeWire thread. //! consumers running outside the PipeWire thread.
use super::classify::DeviceClaim; use super::classify::{DeviceClaim, DeviceProps};
use super::{EventKind, LinkEndpoints, NodeObservation, Projection, RegEvent, RegistryModel}; use super::{
EventKind, LinkEndpoints, NodeObservation, Outcome, Projection, RegEvent, RegistryModel,
};
use crate::host::audio::parse_object_serial; use crate::host::audio::parse_object_serial;
use crate::host::taint::snapshot::{ use crate::host::taint::snapshot::{
ClientSnapshot, GlobalId, MediaRole, NodeProps, PortDirection, PortSnapshot, Serial, ClientSnapshot, GlobalId, MediaRole, NodeProps, PortDirection, PortSnapshot, Serial,
@@ -104,14 +106,24 @@ impl Drop for RegistryObserverHandle {
} }
} }
struct BoundLink { enum BoundProxy {
_proxy: pw::link::Link, Node {
_listener: pw::link::LinkListener, _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 { struct LiveGlobal {
bound_link: Option<BoundLink>, serial: Serial,
bound_proxy: Option<BoundProxy>,
} }
struct ObserverState { 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 // Taken before the model consumes the event: the sink is told what kind
// of observation produced the projection, and deriving that from the // of observation produced the projection, and deriving that from the
// event itself is what stops the two from ever disagreeing. // event itself is what stops the two from ever disagreeing.
let kind = event.kind(); 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(); let candidate = self.model.pulse_pid_candidate();
if candidate != self.last_candidate { if candidate != self.last_candidate {
@@ -160,11 +173,19 @@ impl ObserverState {
// one registry event still yields exactly one sink call — the // one registry event still yields exactly one sink call — the
// no-coalescing contract cuts both ways, and a *duplicated* // no-coalescing contract cuts both ways, and a *duplicated*
// observation would make the O5 event rate a fiction. // 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) { 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 /// 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 /// index. Recording only on *applied* adds (never on unknown object types
/// or globals dropped for a missing serial) is what keeps the two id /// 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 /// queues the same length per id — otherwise a phantom slot could pop
/// bound Link would be popped on removal, leaking that Link's proxy. /// another generation's proxy after an id is recycled.
fn add(&mut self, id: GlobalId, event: RegEvent) { fn add(&mut self, serial: Serial, id: GlobalId, event: RegEvent) {
self.live_globals if self.apply(event) == Outcome::Applied {
.entry(id) self.live_globals
.or_default() .entry(id)
.push_back(LiveGlobal::default()); .or_default()
self.apply(event); .push_back(LiveGlobal {
serial,
bound_proxy: None,
});
}
} }
fn attach_bound_link(&mut self, id: GlobalId, bound_link: BoundLink) { /// Return a proxy that could not be attached so its listener is dropped
let Some(global) = self.live_globals.get_mut(&id).and_then(VecDeque::back_mut) else { /// 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!( tracing::warn!(
global_id = id.0, 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> { fn remove_global(&mut self, id: GlobalId) -> Option<BoundProxy> {
let (bound_link, empty) = { let (bound_proxy, empty) = {
let globals = self.live_globals.get_mut(&id)?; let globals = self.live_globals.get_mut(&id)?;
let bound_link = globals.pop_front().and_then(|global| global.bound_link); let bound_proxy = globals.pop_front().and_then(|global| global.bound_proxy);
(bound_link, globals.is_empty()) (bound_proxy, globals.is_empty())
}; };
if empty { if empty {
self.live_globals.remove(&id); self.live_globals.remove(&id);
} }
bound_link bound_proxy
} }
} }
@@ -273,6 +319,9 @@ fn run_observer(
match obj.type_ { match obj.type_ {
ObjectType::Node => { 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 { let Some(props) = obj.props.as_ref() else {
tracing::warn!( tracing::warn!(
node_id = obj.id, node_id = obj.id,
@@ -284,41 +333,59 @@ fn run_observer(
else { else {
return; return;
}; };
let node_props = NodeProps { state_for_global.borrow_mut().add(
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 {
serial, serial,
id, id,
name: props.get("node.name").map(str::to_owned), RegEvent::NodeAdded { serial, id },
role: MediaRole::parse(props.get("media.class")), );
props: node_props,
device_claim: DeviceClaim { let Some(registry) = registry_weak.upgrade() else {
device_id: props return;
.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),
},
}; };
state_for_global let node: pw::node::Node = match registry.bind(obj) {
.borrow_mut() Ok(node) => node,
.add(id, RegEvent::NodeAdded(observation)); 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 => { ObjectType::Port => {
let Some(props) = obj.props.as_ref() else { let Some(props) = obj.props.as_ref() else {
@@ -357,6 +424,7 @@ fn run_observer(
} }
}; };
state_for_global.borrow_mut().add( state_for_global.borrow_mut().add(
serial,
id, id,
RegEvent::PortAdded(PortSnapshot { RegEvent::PortAdded(PortSnapshot {
serial, serial,
@@ -381,6 +449,7 @@ fn run_observer(
return; return;
}; };
state_for_global.borrow_mut().add( state_for_global.borrow_mut().add(
serial,
id, id,
RegEvent::ClientAdded(ClientSnapshot { RegEvent::ClientAdded(ClientSnapshot {
serial, serial,
@@ -392,9 +461,72 @@ fn run_observer(
); );
} }
ObjectType::Device => { ObjectType::Device => {
state_for_global // Index only, exactly as for a Node: `device.api` and
.borrow_mut() // `alsa.driver_name` live on the bind's `info` (v3.5 §6.7
.add(id, RegEvent::DeviceAdded { id }); // 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 => { ObjectType::Link => {
let Some(props) = obj.props.as_ref() else { let Some(props) = obj.props.as_ref() else {
@@ -410,6 +542,7 @@ fn run_observer(
}; };
let endpoints = link_endpoints_from_props(props); let endpoints = link_endpoints_from_props(props);
state_for_global.borrow_mut().add( state_for_global.borrow_mut().add(
serial,
id, id,
RegEvent::LinkAdded { RegEvent::LinkAdded {
serial, serial,
@@ -456,24 +589,26 @@ fn run_observer(
} }
}) })
.register(); .register();
state_for_global.borrow_mut().attach_bound_link( let unattached = state_for_global.borrow_mut().attach_bound_proxy(
id, id,
BoundLink { serial,
_proxy: link, BoundProxy::Link {
_listener: listener, _listener: listener,
_proxy: link,
}, },
); );
drop(unattached);
} }
_ => {} _ => {}
} }
}) })
.global_remove(move |id| { .global_remove(move |id| {
let id = GlobalId(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 state_for_remove
.borrow_mut() .borrow_mut()
.apply(RegEvent::Removed { id }); .apply(RegEvent::Removed { id });
drop(bound_link); drop(bound_proxy);
}) })
.register(); .register();
@@ -517,6 +652,45 @@ fn truthy(value: Option<&str>) -> bool {
value.is_some_and(|value| value != "false" && value != "0") 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> { 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 output_node = props.get("link.output.node")?.parse::<u32>().ok()?;
let input_node = props.get("link.input.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)) .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] #[test]
#[ignore = "needs live pipewire"] #[ignore = "needs live pipewire"]
fn live_topology_diff_tracks_null_sink_and_loopback() { fn live_topology_diff_tracks_null_sink_and_loopback() {
+87 -37
View File
@@ -18,11 +18,21 @@
//! both. So the discriminator is `factory.name` on an **allowlist** of //! both. So the discriminator is `factory.name` on an **allowlist** of
//! real hardware-PCM factories, never a substring or a denylist: an unknown //! real hardware-PCM factories, never a substring or a denylist: an unknown
//! factory is not a device. //! factory is not a device.
//! - The backing Device must actually have been observed. A node that claims //! - The backing Device must actually have been **bound and resolved**. A node
//! a `device.id` we have not yet resolved is **withheld**, not admitted with //! that claims a `device.id` whose Device's properties we do not hold is
//! a provisional `false` — a provisional `false` during the not-ready //! **withheld**, not admitted with a provisional `false` — a provisional
//! window fuses sink and mic on the shared session client and that fusion //! `false` during the not-ready window fuses sink and mic on the shared
//! can persist as sticky over-exclusion (round-3 finding 3). //! 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; use crate::host::taint::snapshot::GlobalId;
@@ -67,8 +77,9 @@ const HARDWARE_PCM_FACTORIES: &[&str] = &[
/// does not couple playback to capture, so it is not a loopback hazard. /// does not couple playback to capture, so it is not a loopback hazard.
const NON_TERMINAL_ALSA_DRIVERS: &[&str] = &["snd_aloop"]; const NON_TERMINAL_ALSA_DRIVERS: &[&str] = &["snd_aloop"];
/// The three node properties the classifier reads, exactly as the adapter /// The node-side properties the classifier reads, exactly as the adapter
/// parsed them off the Node global. Kept separate from /// 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 /// [`super::super::taint::snapshot::NodeProps`] because these feed the
/// *decision* whose output is the `session_device` field — they are inputs, /// *decision* whose output is the `session_device` field — they are inputs,
/// not part of the graph the engine reasons over. /// not part of the graph the engine reasons over.
@@ -78,10 +89,12 @@ pub struct DeviceClaim {
/// `Stream/*` nodes, which is exactly why their absence means "not a /// `Stream/*` nodes, which is exactly why their absence means "not a
/// device", not "unknown". /// device", not "unknown".
pub device_id: Option<GlobalId>, pub device_id: Option<GlobalId>,
/// `device.api` — the access API of that Device (e.g. `alsa`, `bluez5`). /// `device.api` **as copied onto the node**, when it is — the access API
/// Its mere presence is **not** sufficient (a card-associated filter has /// of that Device (e.g. `alsa`, `bluez5`). Its mere presence is **not**
/// it too); required only as a corroborating signal alongside the factory /// sufficient (a card-associated filter has it too); required only as a
/// allowlist. /// corroborating signal alongside the factory allowlist. The
/// authoritative copy is [`DeviceProps::device_api`]; this is the
/// fallback.
pub device_api: Option<String>, pub device_api: Option<String>,
/// `factory.name` — the discriminator. Only an allowlisted hardware-PCM /// `factory.name` — the discriminator. Only an allowlisted hardware-PCM
/// factory earns `session_device`. /// factory earns `session_device`.
@@ -92,8 +105,27 @@ pub struct DeviceClaim {
/// shares the same factory. `session_device` requires this to be /// shares the same factory. `session_device` requires this to be
/// **present and not** on [`NON_TERMINAL_ALSA_DRIVERS`]; a driver on the /// **present and not** on [`NON_TERMINAL_ALSA_DRIVERS`]; a driver on the
/// denylist, or an absent value, both fail closed (see [`classify`]). /// denylist, or an absent value, both fail closed (see [`classify`]).
/// May be absent on non-ALSA backends or on version pairings that do not /// Frequently absent here — PipeWire ≥ 1.2.6 with WirePlumber < 0.5.13
/// copy `alsa.*` onto the node. /// 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>, pub alsa_driver_name: Option<String>,
} }
@@ -102,9 +134,10 @@ pub struct DeviceClaim {
pub enum Classification { pub enum Classification {
/// No `device.id` — a `Stream/*` node. Admit with `session_device=false`. /// No `device.id` — a `Stream/*` node. Admit with `session_device=false`.
NotADevice, NotADevice,
/// A `device.id` is claimed but the backing Device has not been resolved /// A `device.id` is claimed but the backing Device's properties are not
/// yet. **Withhold the node and keep the readiness epoch not-ready**; /// held: never observed, its bind still outstanding, or its global id
/// re-classify when the Device is observed. /// 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 }, Withhold { device_id: GlobalId },
/// Positively a passive hardware terminal. Admit with /// Positively a passive hardware terminal. Admit with
/// `session_device=true`. /// `session_device=true`.
@@ -115,42 +148,59 @@ pub enum Classification {
NotSessionDevice, 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 /// `device` is the bound Device's properties, and `None` means the claim is
/// as a Device global; it is only consulted when a `device_id` is present. /// **unresolved** — never observed, bind outstanding, or an ambiguous
/// Pure: the model supplies `device_resolved` from its resolved-Device set, /// recycled id. It is only consulted when a `device_id` is present. Pure: the
/// and the I/O of *binding* the Device lives in the adapter. /// model looks the Device up, and the I/O of *binding* it lives in the
pub fn classify(claim: &DeviceClaim, device_resolved: bool) -> Classification { /// 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 { let Some(device_id) = claim.device_id else {
// No backing Device: a stream. Not withheld, not a device. // No backing Device: a stream. Not withheld, not a device.
return Classification::NotADevice; return Classification::NotADevice;
}; };
if !device_resolved { let Some(device) = device else {
// Backed by a Device we have not seen — the one case that blocks // Backed by a Device we have not resolved — the one case that blocks
// readiness. A provisional answer here is the leak the contract // readiness. A provisional answer here is the leak the contract
// forbids. // forbids.
return Classification::Withhold { device_id }; return Classification::Withhold { device_id };
} };
let on_factory_allowlist = claim let on_factory_allowlist = claim
.factory_name .factory_name
.as_deref() .as_deref()
.is_some_and(|f| HARDWARE_PCM_FACTORIES.contains(&f)); .is_some_and(|f| HARDWARE_PCM_FACTORIES.contains(&f));
// A **present, non-denied** ALSA driver is required — absence fails closed // A **present, non-denied** ALSA driver is required — absence fails closed
// (Codex phase-3 re-review). `alsa.driver_name` is not copied onto the // (Codex phase-3 re-review). The factory allowlist cannot tell a real card
// node on every PipeWire/WirePlumber version pairing (PipeWire ≥1.2.6 // from `snd_aloop`, which presents the same `api.alsa.pcm.*` factory, so a
// stopped overwriting node props with card props; WirePlumber only began // *missing* value must not be read as "not a loopback". Round 8 makes the
// copying `alsa.*` onto nodes in 0.5.13), so a *missing* value must not be // bound Device the primary source, so a real card is no longer
// read as "not a loopback" — that is exactly the hole an `snd_aloop` node // over-excluded merely because the session manager did not copy `alsa.*`
// without the property would slip through. A real card whose node lacks // onto its node.
// the driver is instead over-excluded (keeps its owner keys — safe); let driver = device
// 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
.alsa_driver_name .alsa_driver_name
.as_deref() .as_deref()
.is_some_and(|d| !NON_TERMINAL_ALSA_DRIVERS.contains(&d)); .or(claim.alsa_driver_name.as_deref());
let is_hardware_pcm = claim.device_api.is_some() && on_factory_allowlist && driver_ok; 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;
let api_present = device.device_api.is_some() || claim.device_api.is_some();
let is_hardware_pcm = api_present && on_factory_allowlist && driver_ok;
if is_hardware_pcm { if is_hardware_pcm {
Classification::SessionDevice Classification::SessionDevice
} else { } else {
+320 -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 //! 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 //! typed [`RegEvent`]s into a live model of the PipeWire graph and projects
//! the [`GraphSnapshot`] + context the taint engine (phase 2) consumes. **No //! the [`GraphSnapshot`] + context the taint engine (phase 2) consumes. **No
//! PipeWire types appear here** — the I/O adapter (Codex's half) translates //! PipeWire types appear here** — the I/O adapter (Codex's half) translates
//! live registry callbacks, Link/Device binds, `/proc` reads, and the //! live registry callbacks, binds, `/proc` reads, and the `core.sync`/`done`
//! `core.sync`/`done` round-trip into these events and feeds them in. Every //! round-trip into these events and feeds them in. Every test in this module
//! test in this module builds the event stream by hand. //! 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: //! 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 //! id, and those recycle. The model keeps an insertion-ordered index per id
//! so a removal accounts for the *oldest* generation first, and the //! so a removal accounts for the *oldest* generation first, and the
//! snapshot projection treats any id still claimed by two live objects as //! 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 //! - **The readiness epoch.** `graph_ready` is false until the initial graph
//! is fully observed: the server has synced **and** no binds/withheld nodes //! is fully observed: the server has synced **and** no binds/withheld nodes
//! remain outstanding. A bounded timeout makes it fail closed. It gates //! remain outstanding. A bounded timeout makes it fail closed. It gates
//! sticky *retirement* only; withholding after completion is per-object. //! sticky *retirement* only; withholding after completion is per-object.
//! - **Withholding on unresolved devices.** A node claiming a `device.id` //! - **Withholding on unresolved input.** A node with no `info` yet, or one
//! whose Device we have not observed is held out of the snapshot entirely //! claiming a `device.id` whose Device we have not resolved, is held out of
//! rather than admitted with a provisional `session_device` (see //! the snapshot entirely rather than admitted with provisional ownership
//! [`classify`]). //! (see [`classify`]).
//! //!
//! **Two accepted limitations (Codex phase-3 review, findings 3 and 4), both //! **Three accepted limitations, all low-reachability, owed to a later
//! low-reachability, owed to a later hardening round:** //! hardening round:**
//! //!
//! - *A Link dropped for a missing `object.serial`/props is unrepresented.* //! - *A Link dropped for a missing `object.serial`/props is unrepresented.*
//! The adapter drops such a global before it reaches [`RegistryModel`], so //! 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. //! silently drop `global_remove`, so this needs callback loss to trigger.
//! The snapshot treats the two-claimant window as [`IdLookup::Ambiguous`] //! The snapshot treats the two-claimant window as [`IdLookup::Ambiguous`]
//! (fail closed) meanwhile. //! (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. #![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, ClientSnapshot, GlobalId, GraphSnapshot, LinkSnapshot, MediaRole, NodeProps, NodeSnapshot,
PortSnapshot, Serial, PortSnapshot, Serial,
}; };
use classify::{Classification, DeviceClaim}; use classify::{Classification, DeviceClaim, DeviceProps};
use std::collections::{BTreeMap, VecDeque}; use std::collections::{BTreeMap, VecDeque};
/// A monotonic millisecond clock value, supplied by the adapter via /// 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. /// [`std::time::Instant`] so the readiness timeout is deterministic in tests.
pub type Millis = u64; pub type Millis = u64;
/// A Node as observed off the registry, before `session_device` has been /// A Node's **bound `info` properties** — the sole source of node properties
/// decided. The adapter fills [`NodeProps`] with everything it can parse and /// (v3.5 §6.7), delivered by [`RegEvent::NodeInfo`].
/// leaves `session_device` at its `false` default; the model overwrites it ///
/// from the [`classify`] result once the backing Device (if any) is resolved. /// 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)] #[derive(Clone, Debug, PartialEq, Eq)]
pub struct NodeObservation { pub struct NodeObservation {
pub serial: Serial,
pub id: GlobalId,
pub name: Option<String>, pub name: Option<String>,
pub role: MediaRole, pub role: MediaRole,
pub props: NodeProps, pub props: NodeProps,
@@ -97,19 +131,39 @@ pub struct LinkEndpoints {
/// model consumes them in [`RegistryModel::apply`]. /// model consumes them in [`RegistryModel::apply`].
#[derive(Clone, Debug, PartialEq, Eq)] #[derive(Clone, Debug, PartialEq, Eq)]
pub enum RegEvent { pub enum RegEvent {
/// A Node global appeared. Admitted immediately unless it claims an /// A Node global appeared. **Index only** — the global's properties are a
/// unresolved Device (then withheld — see [`classify`]). /// filtered subset and are not read (v3.5 §6.7). The node is withheld
NodeAdded(NodeObservation), /// 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. /// A Port global appeared.
PortAdded(PortSnapshot), PortAdded(PortSnapshot),
/// A Client global appeared. Feeds pulse-PID derivation via `sec_pid`. /// A Client global appeared. Feeds pulse-PID derivation via `sec_pid`.
ClientAdded(ClientSnapshot), ClientAdded(ClientSnapshot),
/// A Device global appeared. Resolves any nodes withheld on its id. /// A Device global appeared. Index only, exactly as for a Node: it does
DeviceAdded { id: GlobalId }, /// 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 /// 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 optimisation) and `None` when the adapter must bind to learn
/// them (the correctness path) — the latter is an outstanding obligation /// them (the correctness path) — the latter is an outstanding obligation
/// until a matching [`RegEvent::LinkEndpointsResolved`] arrives. /// 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 { LinkAdded {
serial: Serial, serial: Serial,
id: GlobalId, id: GlobalId,
@@ -143,7 +197,7 @@ pub enum RegEvent {
/// one is worth suppressing on a tick but never on a graph event. /// one is worth suppressing on a tick but never on a graph event.
#[derive(Clone, Copy, Debug, PartialEq, Eq)] #[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum EventKind { 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. /// probe, or the server sync.
Graph, Graph,
/// The periodic clock sample. Carries no graph information; it exists so the /// 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 /// 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 /// only the id, so the index remembers what each id currently holds. Every
/// slot's serial may live in either the admitted or the withheld map. /// slot names its object by never-recycled serial.
#[derive(Clone, Copy, Debug, PartialEq, Eq)] #[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Slot { enum Slot {
Node(Serial), Node(Serial),
Port(Serial), Port(Serial),
Link(Serial), Link(Serial),
Client(Serial), Client(Serial),
Device, Device(Serial),
} }
/// The readiness epoch. A one-time transition out of [`Readiness::Waiting`]; /// The readiness epoch. A one-time transition out of [`Readiness::Waiting`];
@@ -220,25 +297,45 @@ pub struct Projection {
pub readiness: Readiness, 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`]. /// The live model. Folds [`RegEvent`]s; project with [`RegistryModel::project`].
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
pub struct RegistryModel { pub struct RegistryModel {
// Admitted objects, keyed by their never-recycled serial. /// **Every** live Node, keyed by serial — admitted or withheld. Admission
nodes: BTreeMap<Serial, NodeSnapshot>, /// 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>, ports: BTreeMap<Serial, PortSnapshot>,
links: BTreeMap<Serial, LinkSnapshot>, links: BTreeMap<Serial, LinkSnapshot>,
clients: BTreeMap<Serial, ClientSnapshot>, 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 /// Links whose endpoints the adapter is still binding; the id is kept so
/// removal and resolution can find them. /// removal and resolution can find them.
pending_links: BTreeMap<Serial, GlobalId>, 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` /// Insertion-ordered holders of each live global id. `global_remove`
/// accounts for the oldest generation first (v3.4 §6.1.3). /// accounts for the oldest generation first (v3.4 §6.1.3).
live_ids: BTreeMap<GlobalId, VecDeque<Slot>>, live_ids: BTreeMap<GlobalId, VecDeque<Slot>>,
@@ -259,12 +356,11 @@ impl RegistryModel {
pub fn new(now: Millis, timeout: Millis) -> Self { pub fn new(now: Millis, timeout: Millis) -> Self {
Self { Self {
nodes: BTreeMap::new(), nodes: BTreeMap::new(),
devices: BTreeMap::new(),
ports: BTreeMap::new(), ports: BTreeMap::new(),
links: BTreeMap::new(), links: BTreeMap::new(),
clients: BTreeMap::new(), clients: BTreeMap::new(),
withheld: BTreeMap::new(),
pending_links: BTreeMap::new(), pending_links: BTreeMap::new(),
resolved_devices: BTreeMap::new(),
live_ids: BTreeMap::new(), live_ids: BTreeMap::new(),
probed_comm: BTreeMap::new(), probed_comm: BTreeMap::new(),
server_synced: false, server_synced: false,
@@ -283,21 +379,22 @@ impl RegistryModel {
/// ///
/// This is **dynamic**, not the sticky [`Readiness::Complete`] flag: it is /// This is **dynamic**, not the sticky [`Readiness::Complete`] flag: it is
/// true only when the initial enumeration has completed **and** there are /// true only when the initial enumeration has completed **and** there are
/// no current obligations outstanding (a node withheld on an unresolved /// no current obligations outstanding (a node whose bind is outstanding, a
/// Device, or a Link still being bound). The distinction is the fix for /// node withheld on an unresolved Device, or a Link still being bound).
/// Codex phase-3 review finding 1: a Link whose endpoints are still /// The distinction is the fix for Codex phase-3 review finding 1: a Link
/// resolving is an **invisible edge** — it is absent from the snapshot, /// whose endpoints are still resolving is an **invisible edge** — it is
/// not merely dangling — so a decision made while one exists can miss real /// absent from the snapshot, not merely dangling — so a decision made
/// tainted ancestry and wrongly report a candidate eligible. Unresolved /// while one exists can miss real tainted ancestry and wrongly report a
/// ancestry ⇒ fail closed is the governing invariant (v3.4 §6.1), and an /// candidate eligible. Unresolved ancestry ⇒ fail closed is the governing
/// unresolved Link is unresolved ancestry, so `graph_ready` must drop back /// invariant (v3.4 §6.1), and round 8 adds the far more common case: an
/// to false whenever one is pending — even after the initial epoch. /// 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 /// [`Readiness::Complete`] stays sticky (it records that the initial
/// enumeration happened, for logging and to distinguish "not started" from /// enumeration happened, for logging and to distinguish "not started" from
/// "momentarily churning"); `graph_ready` layers the dynamic obligation /// "momentarily churning"); `graph_ready` layers the dynamic obligation
/// check on top. Downstream (phase 6) may debounce the brief blips a /// 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 { pub fn graph_ready(&self) -> bool {
matches!(self.readiness, Readiness::Complete) && !self.obligations_outstanding() matches!(self.readiness, Readiness::Complete) && !self.obligations_outstanding()
} }
@@ -310,102 +407,120 @@ impl RegistryModel {
pulse_pid::candidate(&clients) pulse_pid::candidate(&clients)
} }
/// Fold one observation into the model. /// Fold one observation into the model. The returned [`Outcome`] tells the
pub fn apply(&mut self, event: RegEvent) { /// 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 { 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) => { RegEvent::PortAdded(port) => {
self.push_id(port.id, Slot::Port(port.serial)); self.push_id(port.id, Slot::Port(port.serial));
self.ports.insert(port.serial, port); self.ports.insert(port.serial, port);
Outcome::Applied
} }
RegEvent::ClientAdded(client) => { RegEvent::ClientAdded(client) => {
self.push_id(client.id, Slot::Client(client.serial)); self.push_id(client.id, Slot::Client(client.serial));
self.clients.insert(client.serial, client); self.clients.insert(client.serial, client);
// A new client can change the pulse candidate; the adapter // A new client can change the pulse candidate; the adapter
// learns that via `pulse_pid_candidate`. No readiness effect. // 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 { RegEvent::LinkAdded {
serial, serial,
id, id,
endpoints, endpoints,
} => self.on_link_added(serial, id, endpoints), } => {
self.on_link_added(serial, id, endpoints);
Outcome::Applied
}
RegEvent::LinkEndpointsResolved { serial, endpoints } => { RegEvent::LinkEndpointsResolved { serial, endpoints } => {
self.on_link_resolved(serial, endpoints) self.on_link_resolved(serial, endpoints)
} }
RegEvent::ProcCommProbed { pid, comm } => { 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::Removed { id } => self.on_removed(id),
RegEvent::ServerSynced => { RegEvent::ServerSynced => {
let already = self.server_synced;
self.server_synced = true; self.server_synced = true;
self.maybe_complete(); self.maybe_complete();
if already {
Outcome::Suppressed
} else {
Outcome::Applied
}
} }
RegEvent::Tick { now } => { RegEvent::Tick { now } => {
self.last_now = now; self.last_now = now;
self.maybe_timeout(now); self.maybe_timeout(now);
Outcome::Applied
} }
} }
} }
fn on_node_added(&mut self, obs: NodeObservation) { /// First resolution *and* every later property change (v3.5 §6.7
self.push_id(obs.id, Slot::Node(obs.serial)); /// decision 2). The model distinguishes them by what it already holds, so
let resolved = obs /// the adapter can forward every `info` callback unconditionally.
.device_claim fn on_node_info(&mut self, serial: Serial, observation: NodeObservation) -> Outcome {
.device_id let Some(entry) = self.nodes.get_mut(&serial) else {
.is_some_and(|id| self.device_resolved(id)); // A late `info` for a node already removed. Re-inserting it here
match classify::classify(&obs.device_claim, resolved) { // would resurrect a dead node with no id index behind it.
Classification::Withhold { .. } => { tracing::debug!(serial = serial.0, "observer: node info for an unknown node");
self.withheld.insert(obs.serial, obs); return Outcome::Suppressed;
} };
Classification::SessionDevice => self.admit_node(obs, true), if entry.obs.as_ref() == Some(&observation) {
Classification::NotADevice | Classification::NotSessionDevice => { // The state-only `info` callbacks PipeWire emits constantly: same
self.admit_node(obs, false) // properties, so the projection is provably identical.
} return Outcome::Suppressed;
} }
// Withholding a node adds an obligation; admitting one can never entry.obs = Some(observation);
// complete readiness on its own, but re-check is cheap and keeps the // The first `info` retires this node's obligation, which can be the
// invariant local. // last one outstanding.
self.maybe_complete(); self.maybe_complete();
Outcome::Applied
} }
fn admit_node(&mut self, obs: NodeObservation, session_device: bool) { fn on_device_info(&mut self, serial: Serial, props: DeviceProps) -> Outcome {
let mut props = obs.props; let Some(entry) = self.devices.get_mut(&serial) else {
props.session_device = session_device; tracing::debug!(
self.nodes.insert( serial = serial.0,
obs.serial, "observer: device info for an unknown device"
NodeSnapshot { );
serial: obs.serial, return Outcome::Suppressed;
id: obs.id, };
name: obs.name, if entry.props.as_ref() == Some(&props) {
role: obs.role, return Outcome::Suppressed;
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);
}
} }
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(); self.maybe_complete();
Outcome::Applied
} }
fn on_link_added(&mut self, serial: Serial, id: GlobalId, endpoints: Option<LinkEndpoints>) { fn on_link_added(&mut self, serial: Serial, id: GlobalId, endpoints: Option<LinkEndpoints>) {
@@ -423,20 +538,23 @@ impl RegistryModel {
self.maybe_complete(); 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 // `remove` also guards against a stale resolution for a Link already
// gone: unknown serial ⇒ ignore. // gone: unknown serial ⇒ ignore.
if let Some(id) = self.pending_links.remove(&serial) { if let Some(id) = self.pending_links.remove(&serial) {
self.links self.links
.insert(serial, link_snapshot(serial, id, endpoints)); .insert(serial, link_snapshot(serial, id, endpoints));
self.maybe_complete(); 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 { let Some(queue) = self.live_ids.get_mut(&id) else {
tracing::warn!(global_id = id.0, "observer: remove for an id we never saw"); 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 // Oldest generation first — the id may be shared during a
// missed-removal window. // missed-removal window.
@@ -446,10 +564,7 @@ impl RegistryModel {
} }
match slot { match slot {
Some(Slot::Node(serial)) => { Some(Slot::Node(serial)) => {
if self.nodes.remove(&serial).is_none() { self.nodes.remove(&serial);
// Was still withheld — drop the obligation.
self.withheld.remove(&serial);
}
} }
Some(Slot::Port(serial)) => { Some(Slot::Port(serial)) => {
self.ports.remove(&serial); self.ports.remove(&serial);
@@ -461,35 +576,67 @@ impl RegistryModel {
Some(Slot::Client(serial)) => { Some(Slot::Client(serial)) => {
self.clients.remove(&serial); self.clients.remove(&serial);
} }
Some(Slot::Device) => { Some(Slot::Device(serial)) => {
if let Some(count) = self.resolved_devices.get_mut(&id) { self.devices.remove(&serial);
*count -= 1;
if *count == 0 {
self.resolved_devices.remove(&id);
}
}
} }
None => { None => {
tracing::warn!(global_id = id.0, "observer: empty id slot on remove"); 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 // A removal can drain the last obligation (an unbound node, a node
// link vanished before it resolved). // withheld on a Device, or a pending link vanished before it
// resolved).
self.maybe_complete(); self.maybe_complete();
Outcome::Applied
} }
fn push_id(&mut self, id: GlobalId, slot: Slot) { fn push_id(&mut self, id: GlobalId, slot: Slot) {
self.live_ids.entry(id).or_default().push_back(slot); self.live_ids.entry(id).or_default().push_back(slot);
} }
fn device_resolved(&self, id: GlobalId) -> bool { /// The bound properties of the Device a node claims by global id, or
self.resolved_devices.get(&id).is_some_and(|&n| n > 0) /// `None` when that claim is unresolved — which covers all three
/// fail-closed cases at once: no such Device observed, its bind still
/// outstanding, or **two live Devices sharing the recycled id**, where
/// there is no way to tell whose properties these are (v3.4 §6.1.3).
fn device_props(&self, id: GlobalId) -> Option<&DeviceProps> {
let mut found: Option<Serial> = None;
for slot in self.live_ids.get(&id)? {
if let Slot::Device(serial) = slot {
if found.is_some() {
return None; // Ambiguous ⇒ unresolved ⇒ withheld.
}
found = Some(*serial);
}
}
self.devices.get(&found?)?.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: /// 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 { 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 /// Completion needs no clock — only the sync flag and an empty obligation
@@ -512,13 +659,34 @@ impl RegistryModel {
if now >= self.deadline { if now >= self.deadline {
self.readiness = Readiness::TimedOut; self.readiness = Readiness::TimedOut;
tracing::warn!( tracing::warn!(
withheld = self.withheld.len(), unbound_nodes = self.unbound_node_count(),
withheld = self.withheld_node_count(),
pending_links = self.pending_links.len(), pending_links = self.pending_links.len(),
"observer: readiness epoch timed out with obligations outstanding — fail closed" "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 /// pipewire-pulse's PID from the current clients, validated against the
/// probed `comm`. `None` whenever anything is ambiguous or unconfirmed — /// probed `comm`. `None` whenever anything is ambiguous or unconfirmed —
/// the safe answer (key 4 unusable). /// the safe answer (key 4 unusable).
@@ -529,9 +697,34 @@ impl RegistryModel {
} }
/// Project the current state into the taint engine's inputs. /// 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 { 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( let snapshot = GraphSnapshot::new(
self.nodes.values().cloned().collect(), nodes,
self.ports.values().cloned().collect(), self.ports.values().cloned().collect(),
self.links.values().cloned().collect(), self.links.values().cloned().collect(),
self.clients.values().cloned().collect(), self.clients.values().cloned().collect(),
+733 -109
View File
File diff suppressed because it is too large Load Diff
+129 -34
View File
@@ -381,36 +381,28 @@ pub fn evaluate(
let components = OwnerComponents::build(snapshot, ctx.pipewire_pulse_pid); let components = OwnerComponents::build(snapshot, ctx.pipewire_pulse_pid);
let keys = owner::OwnerKeyIndex::build(snapshot, ctx.pipewire_pulse_pid); let keys = owner::OwnerKeyIndex::build(snapshot, ctx.pipewire_pulse_pid);
let mut taint: BTreeMap<Serial, Reason> = BTreeMap::new(); // Pass 1 — the fail-closed view. Every decision is made from this one, so
let mut sticky_serials: BTreeSet<Serial> = BTreeSet::new(); // "we could not see" counts as taint.
let (taint, sticky_serials) = compute_taint(
seed_local_roots(snapshot, ctx, &mut taint);
seed_sticky(
snapshot, snapshot,
ctx,
&keys, &keys,
prior,
&components, &components,
&mut taint, prior,
&mut sticky_serials, Uncertainty::FailsClosed,
); );
// Monotone fixpoint: every step only adds taint, or lowers a node's
// reason priority, both of which are bounded. Link propagation and the
// owner bridge feed each other — a bridged output leg has downstream
// links, and a downstream monitor reader bridges to its own siblings —
// so neither can be run once.
let edges = downstream_edges(snapshot, &mut taint);
loop {
let mut changed = false;
changed |= propagate_links(&edges.edges, &mut taint);
changed |= propagate_owner_bridge(&keys, &components, &edges, &mut taint);
changed |= propagate_unresolved_owner(snapshot, &keys, &edges, &mut taint);
if !changed {
break;
}
}
let decisions = build_decisions(snapshot, ctx, &taint, &sticky_serials); let decisions = build_decisions(snapshot, ctx, &taint, &sticky_serials);
// Pass 2 — the evidence-only view, and the only thing sticky state is
// ever built from (see [`Uncertainty`]).
let (evidence, _) = compute_taint(
snapshot,
ctx,
&keys,
&components,
prior,
Uncertainty::Ignored,
);
// ⚠️ Readiness gates **retirement only**, never addition (Codex rounds // ⚠️ Readiness gates **retirement only**, never addition (Codex rounds
// 1 and 2, which caught the two halves of this in turn). An object // 1 and 2, which caught the two halves of this in turn). An object
// missing from an untrustworthy snapshot has not been observed to // missing from an untrustworthy snapshot has not been observed to
@@ -419,10 +411,105 @@ pub fn evaluate(
// *observed* during a not-ready epoch is real — a reader can consume // *observed* during a not-ready epoch is real — a reader can consume
// and buffer the call and then vanish before readiness — so discarding // and buffer the call and then vanish before readiness — so discarding
// additions was the same defect pointing the other way. // additions was the same defect pointing the other way.
let next_sticky = build_sticky(snapshot, &keys, &components, &taint, prior, ctx.graph_ready); let next_sticky = build_sticky(
snapshot,
&keys,
&components,
&evidence,
prior,
ctx.graph_ready,
);
(decisions, next_sticky) (decisions, next_sticky)
} }
/// Whether a pass treats "we could not see" as taint.
///
/// **Both passes exist because stickiness is a claim about history, and
/// uncertainty is not history.** A node tainted only because the graph was
/// mid-enumeration has had nothing observed about it; remembering that as
/// taint forever is over-exclusion with no evidence behind it, and phase 3r's
/// bind-everything observer makes the window it happens in systematically
/// wide (every node is withheld until its bind resolves, so any link observed
/// across that gap raises [`Reason::UnresolvedAncestry`] on its input side).
/// Measured on a live desktop: a hardware sink acquired a permanent sticky
/// taint at every startup, from one link seen while its output node was still
/// unbound.
///
/// Retiring by *reason code* is not enough, because uncertainty launders
/// itself: an unresolved node propagates [`Reason::TaintedUpstream`] to its
/// downstream, and that reason is indistinguishable from real contamination
/// once recorded. So the split is by **provenance** — the sticky pass never
/// raises an uncertainty root at all, and nothing derived from one can reach
/// it. Decisions are unaffected: they are made from the fail-closed pass,
/// which is unchanged.
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum Uncertainty {
/// Unresolved ancestry and an unbounded tainted reader are taint
/// (v3.4 §6.1, §6.1.1, §6.1.4).
FailsClosed,
/// Only positively observed contamination counts.
Ignored,
}
/// One taint fixpoint over the snapshot. The `uncertainty` mode decides
/// whether absence of evidence is treated as evidence of contamination.
fn compute_taint(
snapshot: &GraphSnapshot,
ctx: &ExclusionCtx,
keys: &owner::OwnerKeyIndex,
components: &OwnerComponents,
prior: &StickyState,
uncertainty: Uncertainty,
) -> (BTreeMap<Serial, Reason>, BTreeSet<Serial>) {
let fails_closed = uncertainty == Uncertainty::FailsClosed;
let mut taint: BTreeMap<Serial, Reason> = BTreeMap::new();
let mut sticky_serials: BTreeSet<Serial> = BTreeSet::new();
seed_local_roots(snapshot, ctx, &mut taint);
if fails_closed {
for serial in ambiguous_id_nodes(snapshot) {
raise(&mut taint, serial, Reason::UnresolvedAncestry);
}
}
seed_sticky(
snapshot,
keys,
prior,
components,
&mut taint,
&mut sticky_serials,
);
// Edges are built identically in both passes — receiver status is a
// topological fact and must not depend on the mode, or the owner bridge
// would see two different graphs.
let mut unresolved_input: BTreeSet<Serial> = BTreeSet::new();
let edges = downstream_edges(snapshot, &mut unresolved_input);
if fails_closed {
for serial in unresolved_input {
raise(&mut taint, serial, Reason::UnresolvedAncestry);
}
}
// Monotone fixpoint: every step only adds taint, or lowers a node's
// reason priority, both of which are bounded. Link propagation and the
// owner bridge feed each other — a bridged output leg has downstream
// links, and a downstream monitor reader bridges to its own siblings —
// so neither can be run once.
loop {
let mut changed = false;
changed |= propagate_links(&edges.edges, &mut taint);
changed |= propagate_owner_bridge(keys, components, &edges, &mut taint);
if fails_closed {
changed |= propagate_unresolved_owner(snapshot, keys, &edges, &mut taint);
}
if !changed {
break;
}
}
(taint, sticky_serials)
}
/// Roots that are visible on the node itself. /// Roots that are visible on the node itself.
fn seed_local_roots( fn seed_local_roots(
snapshot: &GraphSnapshot, snapshot: &GraphSnapshot,
@@ -433,14 +520,20 @@ fn seed_local_roots(
if let Some(reason) = local_root_reason(node, ctx) { if let Some(reason) = local_root_reason(node, ctx) {
raise(taint, node.serial, reason); raise(taint, node.serial, reason);
} }
// A node whose own global id is ambiguous cannot be the reliable
// endpoint of any link, so its ancestry is unresolvable.
if snapshot.node_by_id(node.id) == Some(IdLookup::Ambiguous) {
raise(taint, node.serial, Reason::UnresolvedAncestry);
}
} }
} }
/// Nodes whose own global id is ambiguous: they cannot be the reliable
/// endpoint of any link, so their ancestry is unresolvable. Uncertainty, not
/// evidence — see [`Uncertainty`].
fn ambiguous_id_nodes(snapshot: &GraphSnapshot) -> BTreeSet<Serial> {
snapshot
.nodes()
.filter(|node| snapshot.node_by_id(node.id) == Some(IdLookup::Ambiguous))
.map(|node| node.serial)
.collect()
}
fn local_root_reason(node: &NodeSnapshot, ctx: &ExclusionCtx) -> Option<Reason> { fn local_root_reason(node: &NodeSnapshot, ctx: &ExclusionCtx) -> Option<Reason> {
if node.props.peerspeak_owned { if node.props.peerspeak_owned {
return Some(Reason::PeerspeakOwned); return Some(Reason::PeerspeakOwned);
@@ -563,7 +656,7 @@ fn nodes_of_client(
/// `output node → input nodes`, resolving snapshot-local ids. An endpoint /// `output node → input nodes`, resolving snapshot-local ids. An endpoint
/// that does not resolve taints the *other* end as unresolved ancestry when /// that does not resolve taints the *other* end as unresolved ancestry when
/// that other end is the input side — we cannot know what is feeding it. /// that other end is the input side — we cannot know what is feeding it.
fn downstream_edges(snapshot: &GraphSnapshot, taint: &mut BTreeMap<Serial, Reason>) -> Edges { fn downstream_edges(snapshot: &GraphSnapshot, unresolved_input: &mut BTreeSet<Serial>) -> Edges {
let mut edges: BTreeMap<Serial, Vec<Serial>> = BTreeMap::new(); let mut edges: BTreeMap<Serial, Vec<Serial>> = BTreeMap::new();
let mut receivers: BTreeSet<Serial> = BTreeSet::new(); let mut receivers: BTreeSet<Serial> = BTreeSet::new();
for link in snapshot.links() { for link in snapshot.links() {
@@ -575,8 +668,10 @@ fn downstream_edges(snapshot: &GraphSnapshot, taint: &mut BTreeMap<Serial, Reaso
receivers.insert(to); receivers.insert(to);
} }
(_, Some(IdLookup::Unique(to))) => { (_, Some(IdLookup::Unique(to))) => {
// Something feeds this node and we cannot say what. // Something feeds this node and we cannot say what. Reported
raise(taint, to, Reason::UnresolvedAncestry); // rather than raised here, because whether "cannot say" is
// taint depends on which pass is running ([`Uncertainty`]).
unresolved_input.insert(to);
receivers.insert(to); receivers.insert(to);
} }
(_, Some(IdLookup::Ambiguous)) => { (_, Some(IdLookup::Ambiguous)) => {
+131
View File
@@ -820,6 +820,137 @@ fn an_ambiguous_recycled_global_id_fails_closed() {
); );
} }
// ──────────────────────────────────────────────────────────────────────
// Uncertainty is not history — it never enters sticky state
// (round 9, from a live phase-5 audit run; see `Uncertainty` in mod.rs)
// ──────────────────────────────────────────────────────────────────────
#[test]
fn unresolved_ancestry_does_not_survive_being_resolved() {
// Measured live on a desktop: a link is observed while its output node is
// still unbound, the input side fails closed — correctly — and then that
// fail-closed mark became *sticky*, so a hardware sink stayed excluded for
// the process lifetime even after the node resolved and turned out to be
// an ordinary game. Phase 3r's bind-everything observer widens that window
// to every node, so this must clear.
let mut graph = Graph::new();
let ghost = graph.dangling_id();
let client = graph.client_of_app(6000);
let victim = graph.node("victim-in", MediaRole::StreamInput, app(client, 6000));
let sibling = graph.node("victim-out", MediaRole::StreamOutput, app(client, 6000));
graph.link_ids(ghost, victim.id);
let firefox = graph.app_node("firefox", MediaRole::StreamOutput, 11114);
let c = ctx();
// While the ancestry is genuinely unresolved, the decision is unchanged:
// fail closed, both the victim and its sibling excluded.
let (first, sticky) = evaluate(&graph.build(), &c, &StickyState::default());
assert_partition(
&first,
&[("firefox", firefox)],
&[("victim-out", sibling, "tainted-owner-bridge")],
);
assert_tainted(&first, victim, "unresolved-ancestry");
// The node behind that id turns up — nothing tainted, it was simply not
// observed yet. The uncertainty is gone, so nothing may remain of it.
let late_client = graph.client_of_app(7100);
let resolved = graph.node_with_id(
"was-unbound",
MediaRole::StreamOutput,
ghost,
app(late_client, 7100),
);
let (second, _) = evaluate(&graph.build(), &c, &sticky);
assert_partition(
&second,
&[
("firefox", firefox),
("victim-out", sibling),
("was-unbound", resolved),
],
&[],
);
}
#[test]
fn uncertainty_laundered_into_downstream_taint_is_not_sticky_either() {
// Retiring by reason *code* would not be enough: an unresolved node
// propagates `tainted-upstream`, which is indistinguishable from real
// contamination once recorded. The split has to be by provenance, so a
// node two hops from the uncertainty must clear too.
let mut graph = Graph::new();
let ghost = graph.dangling_id();
let forwarder_client = graph.client_of_app(6100);
let forwarder_in = graph.node(
"fwd-in",
MediaRole::StreamInput,
app(forwarder_client, 6100),
);
let forwarder_out = graph.node(
"fwd-out",
MediaRole::StreamOutput,
app(forwarder_client, 6100),
);
let downstream_client = graph.client_of_app(6200);
let downstream = graph.node("downstream", MediaRole::Sink, app(downstream_client, 6200));
let downstream_leg = graph.node(
"downstream-out",
MediaRole::StreamOutput,
app(downstream_client, 6200),
);
graph.link_ids(ghost, forwarder_in.id);
graph.link(forwarder_out, downstream);
let c = ctx();
let (first, sticky) = evaluate(&graph.build(), &c, &StickyState::default());
assert_tainted(&first, forwarder_in, "unresolved-ancestry");
assert_tainted(&first, downstream, "tainted-upstream");
assert!(
first.candidates[&downstream_leg.serial].reason().is_some(),
"while the ancestry is unresolved the downstream owner is excluded too"
);
let late_client = graph.client_of_app(7200);
graph.node_with_id(
"was-unbound",
MediaRole::StreamOutput,
ghost,
app(late_client, 7200),
);
let (second, _) = evaluate(&graph.build(), &c, &sticky);
assert_eq!(
second.candidates[&downstream_leg.serial].reason(),
None,
"nothing derived from the uncertainty may outlive it"
);
assert_eq!(
second.candidates[&forwarder_out.serial].reason(),
None,
"including the unresolved node's own owner siblings"
);
}
#[test]
fn real_taint_is_still_sticky_when_its_topology_goes_away() {
// The other half of the same rule, stated positively: *evidence* is
// history and must survive. This is the guard on the change above — if
// provenance splitting ever leaks into the evidence path, peerspeak's own
// audio starts escaping.
let (graph, call, rec_in, rec_out, firefox) = sticky_scene();
let c = ctx();
let (_, sticky) = evaluate(&graph.build(), &c, &StickyState::default());
let (second, _) = evaluate(&graph.build_without(&[rec_in]), &c, &sticky);
assert_partition(
&second,
&[("firefox", firefox)],
&[
("call", call, "peerspeak-owned"),
("rec-out", rec_out, "tainted-owner-bridge"),
],
);
}
// ────────────────────────────────────────────────────────────────────── // ──────────────────────────────────────────────────────────────────────
// Stickiness and lifetime-awareness (v3.4 §6.1.3) // Stickiness and lifetime-awareness (v3.4 §6.1.3)
// ────────────────────────────────────────────────────────────────────── // ──────────────────────────────────────────────────────────────────────