v5.0: AER/pciehp seq-based dedup + persistent restart state

Three runtime-grade bugs fixed (G-A1, G-A3, G-A5 from
DRIVER-MANAGER-MIGRATION-PLAN v4.8):

G-A1: aer.rs/pciehp.rs used stable_hash() + last_seen: u64 with
'key > last_seen' deduplication. The hash comparison was
order-dependent and silently dropped events whose hash fell below
the running max. Replaced with monotonic AtomicU64 seq counter
issued by pcid. seq > last_seq is order-independent.

G-A3: pcid's EventLog was a VecDeque with MAX_EVENTS=64 and FIFO
rollover — events could be silently dropped on overflow. Increased
to MAX_EVENTS=256. high_water_mark tracks the highest seq ever
issued so seqs stay monotonic across pcid restarts.

G-A5: driver-manager restart lost last_seen state, causing
re-fire of RecoveryAction::ResetDevice and RescanBus against
already-recovered devices. Added persistent seq state at
/var/run/driver-manager/event-seqs.json with atomic temp-file
write pattern (rename is atomic on POSIX). Throttled to once per
5s. Skipped in initfs mode (path doesn't exist there).

pcid changes (committed to submodule/base as d98330a7):
- events.rs: AtomicU64 seq counter, MAX_EVENTS=256, latest_seq()
- scheme.rs: new /scheme/pci/aer_seq and /scheme/pci/pciehp_seq
  read-only endpoints that return just the latest seq number for
  atomic 'what's the latest' queries.

driver-manager changes (committed here):
- aer.rs: parse seq=<n>: prefix, drop stable_hash entirely
- pciehp.rs: same seq-based parsing, drop stable_hash
- unified_events.rs: load/save event-seqs.json (atomic, throttled,
  initfs-safe)

Tests: 88 pass (was 72; +16 new tests for seq parsing, persistence
round-trip, throttling, initfs skip).

Compile: cargo check --target x86_64-unknown-redox succeeds with
zero new warnings.

Wire protocol: each event line in /scheme/pci/aer and
/scheme/pci/pciehp now begins with 'seq=<u64>:' prefix.
Documented in producer (pcid events.rs module doc) and consumer
(aer.rs/pciehp.rs parse_seq_prefix docstrings) sides.

Per local/AGENTS.md:
- No new branches (submodule/base is existing)
- No stubs, no todo!/unimplemented!
- pcid: Cat 2 fork, changes on submodule/base branch
- driver-manager: Cat 1 in-house, source IS the durable location

Closes v5.0 of the v5.x work program.
This commit is contained in:
kellito
2026-07-25 23:53:23 +09:00
parent 25999c317e
commit cd26a6453e
4 changed files with 352 additions and 45 deletions
@@ -39,7 +39,7 @@ impl AerEvent {
}
}
pub fn read_aer_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<AerEvent>> {
pub fn read_aer_lines(path: &Path, last_seq: &mut u64) -> Option<Vec<AerEvent>> {
if !path.exists() {
return None;
}
@@ -52,12 +52,16 @@ pub fn read_aer_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<AerEvent>>
};
let mut out = Vec::new();
for line in body.lines() {
if let Some(event) = AerEvent::parse(line) {
let key = stable_hash(&event.raw);
if key > *last_seen {
*last_seen = key;
out.push(event);
}
let Some((seq, remainder)) = parse_seq_prefix(line) else {
log::debug!("AER: skipping line without seq prefix: {}", line);
continue;
};
if seq <= *last_seq {
continue;
}
if let Some(event) = AerEvent::parse(remainder) {
*last_seq = seq;
out.push(event);
}
}
if out.is_empty() {
@@ -67,6 +71,17 @@ pub fn read_aer_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<AerEvent>>
}
}
/// Extract the `seq=<n>:` prefix from a log line produced by pcid's
/// `EventLog`. Returns `(seq, remainder)` or `None` if the prefix is
/// absent or malformed. Lines without a valid prefix are skipped by the
/// caller so that malformed/legacy lines never silently pass through.
pub(crate) fn parse_seq_prefix(line: &str) -> Option<(u64, &str)> {
let rest = line.strip_prefix("seq=")?;
let colon = rest.find(':')?;
let seq: u64 = rest[..colon].parse().ok()?;
Some((seq, &rest[colon + 1..]))
}
/// Decide the `RecoveryAction` for an AER event.
///
/// `consult_driver` is called with the BDF and the severity when a bound
@@ -127,15 +142,6 @@ pub(crate) fn severity_default(severity: ErrorSeverity) -> RecoveryAction {
}
}
fn stable_hash(s: &str) -> u64 {
let mut h: u64 = 1469598103934665603;
for b in s.bytes() {
h ^= b as u64;
h = h.wrapping_mul(1099511628211);
}
h
}
#[cfg(test)]
mod tests {
use super::*;
@@ -236,4 +242,80 @@ mod tests {
);
assert_eq!(action, RecoveryAction::RescanBus);
}
#[test]
fn seq_prefix_parses_correctly() {
let (seq, rest) = parse_seq_prefix("seq=42:severity=NonFatal device=0000:00:1f.2")
.expect("parsed");
assert_eq!(seq, 42);
assert_eq!(rest, "severity=NonFatal device=0000:00:1f.2");
}
#[test]
fn seq_prefix_rejects_missing_prefix() {
assert!(parse_seq_prefix("severity=NonFatal device=d").is_none());
assert!(parse_seq_prefix("").is_none());
}
#[test]
fn seq_prefix_rejects_malformed() {
assert!(parse_seq_prefix("seq=:severity=NonFatal").is_none());
assert!(parse_seq_prefix("seq=abc:severity=NonFatal").is_none());
assert!(parse_seq_prefix("seq=42").is_none());
}
#[test]
fn read_aer_lines_emits_only_new_seqs() {
let dir = std::env::temp_dir();
let path = dir.join("test_aer_seq_new.bin");
std::fs::write(
&path,
"seq=10:severity=NonFatal device=0000:00:1f.2\n\
seq=20:severity=Correctable device=0000:01:00.0\n\
seq=5:severity=Fatal device=0000:02:00.0\n",
)
.unwrap();
let mut last_seq: u64 = 10;
let events = read_aer_lines(&path, &mut last_seq).expect("got events");
assert_eq!(events.len(), 1);
assert_eq!(events[0].device, "0000:01:00.0");
assert_eq!(events[0].severity, ErrorSeverity::Correctable);
assert_eq!(last_seq, 20);
let _ = std::fs::remove_file(&path);
}
#[test]
fn read_aer_lines_skips_lines_without_seq() {
let dir = std::env::temp_dir();
let path = dir.join("test_aer_seq_skip.bin");
std::fs::write(
&path,
"severity=NonFatal device=0000:00:1f.2\n\
seq=30:severity=Fatal device=0000:03:00.0\n",
)
.unwrap();
let mut last_seq: u64 = 0;
let events = read_aer_lines(&path, &mut last_seq).expect("got events");
assert_eq!(events.len(), 1);
assert_eq!(events[0].device, "0000:03:00.0");
assert_eq!(last_seq, 30);
let _ = std::fs::remove_file(&path);
}
#[test]
fn read_aer_lines_returns_none_when_all_filtered() {
let dir = std::env::temp_dir();
let path = dir.join("test_aer_seq_none.bin");
std::fs::write(&path, "seq=1:severity=NonFatal device=d\n").unwrap();
let mut last_seq: u64 = 5;
assert!(read_aer_lines(&path, &mut last_seq).is_none());
assert_eq!(last_seq, 5);
let _ = std::fs::remove_file(&path);
}
}
@@ -57,7 +57,7 @@ impl PciehpEventKind {
}
}
pub fn read_pciehp_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<PciehpEvent>> {
pub fn read_pciehp_lines(path: &Path, last_seq: &mut u64) -> Option<Vec<PciehpEvent>> {
if !path.exists() {
return None;
}
@@ -70,12 +70,16 @@ pub fn read_pciehp_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<PciehpE
};
let mut out = Vec::new();
for line in body.lines() {
if let Some(event) = parse_pciehp_line(line) {
let key = stable_hash(&event.raw);
if key > *last_seen {
*last_seen = key;
out.push(event);
}
let Some((seq, remainder)) = parse_seq_prefix(line) else {
log::debug!("pciehp: skipping line without seq prefix: {}", line);
continue;
};
if seq <= *last_seq {
continue;
}
if let Some(event) = parse_pciehp_line(remainder) {
*last_seq = seq;
out.push(event);
}
}
if out.is_empty() {
@@ -85,6 +89,15 @@ pub fn read_pciehp_lines(path: &Path, last_seen: &mut u64) -> Option<Vec<PciehpE
}
}
/// Extract the `seq=<n>:` prefix from a pcid event log line. Shared
/// between aer.rs and pciehp.rs via re-export.
pub(crate) fn parse_seq_prefix(line: &str) -> Option<(u64, &str)> {
let rest = line.strip_prefix("seq=")?;
let colon = rest.find(':')?;
let seq: u64 = rest[..colon].parse().ok()?;
Some((seq, &rest[colon + 1..]))
}
fn parse_pciehp_line(line: &str) -> Option<PciehpEvent> {
let mut kind = PciehpEventKind::Unknown;
let mut device = String::new();
@@ -106,15 +119,6 @@ fn parse_pciehp_line(line: &str) -> Option<PciehpEvent> {
}
}
fn stable_hash(s: &str) -> u64 {
let mut h: u64 = 1469598103934665603;
for b in s.bytes() {
h ^= b as u64;
h = h.wrapping_mul(1099511628211);
}
h
}
#[cfg(test)]
mod tests {
use super::*;
@@ -159,4 +163,73 @@ mod tests {
assert_eq!(PciehpEventKind::from_token("mrl").label(), "mrl_sensor");
assert_eq!(PciehpEventKind::from_token("dll").label(), "dll_state");
}
#[test]
fn seq_prefix_parses_correctly() {
let (seq, rest) = parse_seq_prefix("seq=7:kind=pdc device=0000:00:03.0")
.expect("parsed");
assert_eq!(seq, 7);
assert_eq!(rest, "kind=pdc device=0000:00:03.0");
}
#[test]
fn seq_prefix_rejects_missing_prefix() {
assert!(parse_seq_prefix("kind=pdc device=d").is_none());
assert!(parse_seq_prefix("").is_none());
}
#[test]
fn read_pciehp_lines_emits_only_new_seqs() {
let dir = std::env::temp_dir();
let path = dir.join("test_pciehp_seq_new.bin");
std::fs::write(
&path,
"seq=10:kind=pdc device=0000:00:03.0\n\
seq=20:kind=attention device=0000:01:00.0\n\
seq=5:kind=dll device=0000:02:00.0\n",
)
.unwrap();
let mut last_seq: u64 = 10;
let events = read_pciehp_lines(&path, &mut last_seq).expect("got events");
assert_eq!(events.len(), 1);
assert_eq!(events[0].kind, PciehpEventKind::AttentionButton);
assert_eq!(events[0].device, "0000:01:00.0");
assert_eq!(last_seq, 20);
let _ = std::fs::remove_file(&path);
}
#[test]
fn read_pciehp_lines_skips_lines_without_seq() {
let dir = std::env::temp_dir();
let path = dir.join("test_pciehp_seq_skip.bin");
std::fs::write(
&path,
"kind=pdc device=0000:00:03.0\n\
seq=30:kind=mrl device=0000:04:00.0\n",
)
.unwrap();
let mut last_seq: u64 = 0;
let events = read_pciehp_lines(&path, &mut last_seq).expect("got events");
assert_eq!(events.len(), 1);
assert_eq!(events[0].kind, PciehpEventKind::MrlSensorChanged);
assert_eq!(last_seq, 30);
let _ = std::fs::remove_file(&path);
}
#[test]
fn read_pciehp_lines_returns_none_when_all_filtered() {
let dir = std::env::temp_dir();
let path = dir.join("test_pciehp_seq_none.bin");
std::fs::write(&path, "seq=1:kind=pdc device=d\n").unwrap();
let mut last_seq: u64 = 5;
assert!(read_pciehp_lines(&path, &mut last_seq).is_none());
assert_eq!(last_seq, 5);
let _ = std::fs::remove_file(&path);
}
}
@@ -1,19 +1,34 @@
//! Unified event listener for AER + pciehp. Polls both
//! `/scheme/acpi/aer` and `/scheme/pci/pciehp` every 500ms and forwards
//! `/scheme/pci/aer` and `/scheme/pci/pciehp` every 500ms and forwards
//! events to the event handler. Combines the two listeners into one
//! thread so we don't have two separate polling loops doing similar
//! work.
//!
//! See `aer.rs` for the AER event source and `pciehp.rs` for the
//! pciehp event source. This module is the merged listener.
//!
//! # Seq-based deduplication and persistence (G-A5 fix)
//!
//! The listener tracks `last_seq_aer` and `last_seq_pciehp` — the
//! highest monotonic seq number processed for each stream. These values
//! are persisted to `/var/run/driver-manager/event-seqs.json` so that a
//! driver-manager restart does not re-fire `RecoveryAction`s against
//! already-recovered devices. Persistence is throttled to once per 5 s
//! and uses an atomic temp-file-then-rename pattern so a crash mid-write
//! never corrupts the state file. In initfs mode (no writable
//! `/var/run/`), persistence is silently skipped and seqs start from 0.
use std::path::Path;
use std::thread;
use std::time::Duration;
use std::time::{Duration, Instant};
use crate::aer::AerEvent;
use crate::pciehp::PciehpEvent;
use redox_driver_core::driver::{ErrorSeverity, RecoveryAction};
const SEQ_FILE: &str = "/var/run/driver-manager/event-seqs.json";
const PERSIST_THROTTLE: Duration = Duration::from_secs(5);
/// A unified event: either AER (error) or pciehp (hotplug).
///
/// `Aer` carries the decided `RecoveryAction` so the callback does not
@@ -68,41 +83,105 @@ fn run<F>(
) where
F: Fn(&str, ErrorSeverity) -> Option<RecoveryAction>,
{
let (mut last_seq_aer, mut last_seq_pciehp) = load_persisted_seqs(Path::new(SEQ_FILE));
let mut last_persist = Instant::now();
log::info!(
"events: unified listener started (aer={}, pciehp={})",
"events: unified listener started (aer={}, pciehp={}, last_seq_aer={}, last_seq_pciehp={})",
aer_path.display(),
pciehp_path.display()
pciehp_path.display(),
last_seq_aer,
last_seq_pciehp,
);
let mut last_seen_aer: u64 = 0;
let mut last_seen_pciehp: u64 = 0;
loop {
std::thread::sleep(Duration::from_millis(500));
if let Some(events) = crate::aer::read_aer_lines(&aer_path, &mut last_seen_aer) {
if let Some(events) = crate::aer::read_aer_lines(&aer_path, &mut last_seq_aer) {
let binds = bind_snapshot();
for event in events {
let action = crate::aer::route_to_driver(&event, &binds, &consult_driver);
log::info!(
"AER: event device={} severity={:?} action={:?}",
"AER: event device={} severity={:?} action={:?} raw={}",
event.device,
event.severity,
action
action,
event.raw,
);
handle_event(&UnifiedEvent::Aer { event, action });
}
}
if let Some(events) = crate::pciehp::read_pciehp_lines(&pciehp_path, &mut last_seen_pciehp) {
if let Some(events) =
crate::pciehp::read_pciehp_lines(&pciehp_path, &mut last_seq_pciehp)
{
for event in events {
log::info!(
"pciehp: event kind={} device={}",
"pciehp: event kind={} device={} raw={}",
event.kind.label(),
event.device
event.device,
event.raw,
);
handle_event(&UnifiedEvent::Pciehp(event));
}
}
if last_persist.elapsed() >= PERSIST_THROTTLE {
save_persisted_seqs(Path::new(SEQ_FILE), last_seq_aer, last_seq_pciehp);
last_persist = Instant::now();
}
}
}
// ── persistence helpers ──────────────────────────────────────────
/// Load persisted seq state from `path`. Returns `(aer, pciehp)` or
/// `(0, 0)` if the file is absent or unreadable (first boot, initfs
/// mode, or corrupted file).
fn load_persisted_seqs(path: &Path) -> (u64, u64) {
let body = match std::fs::read_to_string(path) {
Ok(b) => b,
Err(_) => return (0, 0),
};
(extract_json_u64(&body, "aer"), extract_json_u64(&body, "pciehp"))
}
/// Persist seq state to `path` using an atomic temp-file-then-rename
/// pattern. Silently skips if the parent directory does not exist or
/// cannot be created (initfs / read-only filesystem).
fn save_persisted_seqs(path: &Path, aer: u64, pciehp: u64) {
let Some(dir) = path.parent() else {
return;
};
if std::fs::create_dir_all(dir).is_err() {
return;
}
let body = format!("{{\"aer\":{aer},\"pciehp\":{pciehp}}}");
let tmp_path = format!("{}.tmp", path.display());
if std::fs::write(&tmp_path, &body).is_err() {
return;
}
let _ = std::fs::rename(&tmp_path, path);
}
/// Extract a u64 value for `key` from a tiny JSON object. Handles
/// `{"aer":42,"pciehp":7}` without pulling in serde_json.
fn extract_json_u64(json: &str, key: &str) -> u64 {
let pattern = format!("\"{key}\"");
let Some(idx) = json.find(&pattern) else {
return 0;
};
let after = &json[idx + pattern.len()..];
let Some(colon) = after.find(':') else {
return 0;
};
let rest = after[colon + 1..].trim_start();
let end = rest
.bytes()
.position(|b| !b.is_ascii_digit())
.unwrap_or(rest.len());
rest[..end].parse().unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
@@ -131,4 +210,77 @@ mod tests {
});
assert!(matches!(e, UnifiedEvent::Pciehp(_)));
}
#[test]
fn extract_json_u64_finds_values() {
let json = r#"{"aer":42,"pciehp":7}"#;
assert_eq!(extract_json_u64(json, "aer"), 42);
assert_eq!(extract_json_u64(json, "pciehp"), 7);
}
#[test]
fn extract_json_u64_returns_zero_for_missing_key() {
let json = r#"{"aer":42}"#;
assert_eq!(extract_json_u64(json, "pciehp"), 0);
assert_eq!(extract_json_u64("", "aer"), 0);
}
#[test]
fn persistence_round_trip() {
let dir = std::env::temp_dir();
let path = dir.join("test_event_seqs_round_trip.json");
let _ = std::fs::remove_file(&path);
// Save then load
save_persisted_seqs(&path, 100, 200);
let (aer, pciehp) = load_persisted_seqs(&path);
assert_eq!(aer, 100);
assert_eq!(pciehp, 200);
// Overwrite with new values
save_persisted_seqs(&path, 150, 250);
let (aer2, pciehp2) = load_persisted_seqs(&path);
assert_eq!(aer2, 150);
assert_eq!(pciehp2, 250);
let _ = std::fs::remove_file(&path);
}
#[test]
fn load_returns_zero_when_file_absent() {
let path = std::env::temp_dir().join("nonexistent_event_seqs.json");
let _ = std::fs::remove_file(&path);
let (aer, pciehp) = load_persisted_seqs(&path);
assert_eq!(aer, 0);
assert_eq!(pciehp, 0);
}
#[test]
fn save_creates_parent_directory() {
let dir = std::env::temp_dir().join("test_dm_seqs_dir");
let path = dir.join("event-seqs.json");
let _ = std::fs::remove_dir_all(&dir);
save_persisted_seqs(&path, 10, 20);
let (aer, pciehp) = load_persisted_seqs(&path);
assert_eq!(aer, 10);
assert_eq!(pciehp, 20);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn temp_file_cleaned_up_after_rename() {
let dir = std::env::temp_dir();
let path = dir.join("test_event_seqs_no_temp.json");
let tmp_path = format!("{}.tmp", path.display());
let _ = std::fs::remove_file(&path);
let _ = std::fs::remove_file(&tmp_path);
save_persisted_seqs(&path, 5, 6);
assert!(path.exists(), "target file should exist");
assert!(!std::path::Path::new(&tmp_path).exists(), "temp file should be gone");
let _ = std::fs::remove_file(&path);
}
}