Files
pixelpass/src/host/audit/sink.rs
T
molluskandClaude Opus 5 bbf6744444 host/audit: phase 5 — dry-run audit mode (read-only)
Runs phases 2-4 against the live PipeWire graph on every registry event and
reports the complete eligible/excluded candidate partition with stable reason
codes. Creates no links, loads no modules, changes no routing.

Impl plan §5. Two entry points behind the hidden PIXELPASS_AUDIO_AUDIT=1
trigger: inside a real `pixelpass host` run (the plan-literal reading, proves
the path phase 6 will mutate), and a hidden `--audit-audio` standalone mode
with no iroh endpoint or capture pipeline, which is what drives the §5.1
matrix.

The recompute runs inline on the observer thread via a new ProjectionSink
hook, once per applied event. Polling `latest()` was rejected: it coalesces,
and phase 4 detects a module unload by observing the empty gap before the next
module appears — with indices reused verbatim (v3.4 §5.2 correction 3), a
missed gap aliases a fresh module onto a dead identity. Running inline is what
makes phase 4's "one observe per graph event" contract true, and it puts the
cost where O5 can measure it.

Split as usual: the auditor and the metrics are pure and unit-tested; the
clock, the writer and the env parsing are the thin edge in `sink`/`run`.

- audit/mod.rs   Auditor: AEC validator + taint engine + record building.
                 The AEC gate and the engine's own reasons stay
                 distinguishable — a shut gate must not erase the reason codes
                 the §5.1 rows assert.
- audit/metrics.rs  O5: event rate, bucketed recompute distribution + exact
                 max, busy fraction, and a documented lower-bound queueing
                 proxy (libpipewire exposes no queue depth).
- audit/sink.rs  JSON Lines to stderr, or PIXELPASS_AUDIO_AUDIT_FILE. Never
                 stdout — peerspeak parses that stream.
- audit/run.rs   Env parsing; a malformed AEC value is fatal, matching phase
                 4's rule that it must not silently become "no AEC".

Observer gains `EventKind` (derived from RegEvent, so a consumer's view of
"was this a real graph change?" cannot disagree with the model's) and
`Projection::readiness`, which distinguishes the three ways graph_ready can be
false. taint::fixture is now pub(crate) so audit tests share one graph
vocabulary with the taint tests.

33 new tests, 178 green, clippy -D warnings and fmt clean.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 15:43:10 -04:00

185 lines
7.0 KiB
Rust

//! The audit's I/O edge: timing, JSON Lines emission, O5 accounting.
//!
//! Everything impure about phase 5 lives here, and it is deliberately thin —
//! read the clock, call [`Auditor::observe`], write a line, fold a
//! [`metrics::Sample`]. The decisions are all upstream in the pure core, which
//! is why the matrix can be argued about in unit tests rather than only in front
//! of a live daemon.
//!
//! ## Why this runs on the observer thread
//!
//! [`AuditSink`] is a [`ProjectionSink`], invoked inline from the PipeWire
//! observer thread once per applied registry event. The obvious alternative —
//! a consumer task polling
//! [`RegistryObserverHandle::latest`](super::super::observer::adapter::RegistryObserverHandle::latest)
//! — was rejected: polling **coalesces**, and phase 4's revocation logic
//! detects a module unload by observing the *empty gap* before the next module
//! appears. Module indices are reused verbatim across an unload/reload (v3.4
//! §5.2 correction 3), so a poller that misses the gap silently aliases a fresh
//! module onto a dead module's validated identity. Running inline is what makes
//! "one `observe` per graph event, no coalescing" — the contract phase 4
//! documents as owed — actually true.
//!
//! The cost of that choice is that recompute and logging happen on the thread
//! servicing PipeWire, which is precisely the risk O5 asks about. That is not an
//! accident: this arrangement puts the cost exactly where the measurement can
//! see it. See [`metrics`].
//!
//! ## Output contract
//!
//! One JSON object per line, to **stderr** by default, each tagged with a `kind`
//! discriminator (`"audit"` or `"metrics"`). Never stdout: peerspeak parses
//! pixelpass's stdout event stream, and the impl plan §5 is explicit that
//! unstructured output must not go there. `PIXELPASS_AUDIO_AUDIT_FILE`
//! redirects the records to a file instead, which is how the §5.1 matrix is
//! driven — it separates the audit stream from interleaved `tracing` output
//! without needing either side to change format.
use std::io::Write;
use std::time::Instant;
use serde::Serialize;
use super::metrics::{self, Metrics, Summary};
use super::{AuditConfig, AuditRecord, Auditor};
use crate::host::observer::adapter::ProjectionSink;
use crate::host::observer::{EventKind, Millis, Projection};
/// Emit a rolling metrics line every this many ticks. Ticks are 250 ms, so this
/// is every 10 s — often enough that a run killed abruptly still leaves a
/// usable O5 record, rare enough that it does not crowd out the audit records.
const SUMMARY_INTERVAL_TICKS: u64 = 40;
/// The live audit: pure auditor + clock + writer.
pub struct AuditSink {
auditor: Auditor,
metrics: Metrics,
writer: Box<dyn Write + Send>,
/// Set once the first sample has completed, so the first event is not
/// counted as having queued behind a predecessor that does not exist.
last_completion_us: Option<u64>,
ticks_since_summary: u64,
/// Wall-clock origin for the microsecond timings. Only used for durations,
/// never for the AEC deadline — that runs on the observer's own clock,
/// handed in as `now_us`, so the validator and the readiness epoch cannot
/// disagree about what time it is.
epoch: Instant,
}
impl AuditSink {
pub fn new(config: AuditConfig, writer: Box<dyn Write + Send>) -> Self {
Self {
auditor: Auditor::new(config),
metrics: Metrics::default(),
writer,
last_completion_us: None,
ticks_since_summary: 0,
epoch: Instant::now(),
}
}
fn elapsed_us(&self) -> u64 {
u64::try_from(self.epoch.elapsed().as_micros()).unwrap_or(u64::MAX)
}
/// Write one line. Failures are logged once per occurrence and otherwise
/// ignored: a broken stderr must not take down the observer thread, and the
/// audit is diagnostic — losing a line is a worse audit, not a worse share.
fn write_line<T: Serialize>(&mut self, line: &T) {
match serde_json::to_string(line) {
Ok(json) => {
if let Err(e) = writeln!(self.writer, "{json}") {
tracing::warn!("audit: failed to write record: {e}");
}
}
Err(e) => tracing::warn!("audit: failed to serialise record: {e}"),
}
}
fn write_summary(&mut self, at_ms: Millis) {
let summary = self.metrics.summary();
self.write_line(&MetricsLine {
kind: "metrics",
at_ms,
summary: &summary,
});
let _ = self.writer.flush();
}
}
impl ProjectionSink for AuditSink {
fn on_projection(&mut self, projection: &Projection, kind: EventKind, now_us: u64) {
let at_us = self.elapsed_us();
let gap_us = self
.last_completion_us
.map(|previous| at_us.saturating_sub(previous))
.unwrap_or(0);
let recompute_start = self.elapsed_us();
let outcome = self.auditor.observe(projection, kind, now_us / 1_000);
let recompute_us = self.elapsed_us().saturating_sub(recompute_start);
let emit_us = if outcome.emit {
let emit_start = self.elapsed_us();
self.write_line(&AuditLine {
kind: "audit",
recompute_us,
record: &outcome.record,
});
// Flushed per record so a run ended with SIGKILL (or a matrix row
// that reads the file while the process is still up) still shows
// every decision made before that instant. The cost is measured, not
// assumed — it is inside `emit_us`.
let _ = self.writer.flush();
self.elapsed_us().saturating_sub(emit_start).max(1)
} else {
0
};
self.metrics.record(metrics::Sample {
at_us,
gap_us,
recompute_us,
emit_us,
kind,
});
self.last_completion_us = Some(self.elapsed_us());
if kind == EventKind::Tick {
self.ticks_since_summary += 1;
if self.ticks_since_summary >= SUMMARY_INTERVAL_TICKS {
self.ticks_since_summary = 0;
self.write_summary(now_us / 1_000);
}
}
}
}
impl Drop for AuditSink {
/// The final O5 record. The observer thread drops its sink when the main
/// loop quits, so an ordinary ctrl-c leaves a complete summary behind
/// without the runner having to ask for one.
fn drop(&mut self) {
let at_ms = self.elapsed_us() / 1_000;
self.write_summary(at_ms);
}
}
#[derive(Serialize)]
struct AuditLine<'a> {
kind: &'static str,
/// This record's own recompute cost, so a surprising row can be correlated
/// with a cost spike without cross-referencing the periodic summary.
recompute_us: u64,
#[serde(flatten)]
record: &'a AuditRecord,
}
#[derive(Serialize)]
struct MetricsLine<'a> {
kind: &'static str,
at_ms: Millis,
#[serde(flatten)]
summary: &'a Summary,
}