From cd26a6453ead9c26174ab18f1cf645be76d8d49f Mon Sep 17 00:00:00 2001 From: kellito Date: Sat, 25 Jul 2026 23:53:23 +0900 Subject: [PATCH] v5.0: AER/pciehp seq-based dedup + persistent restart state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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=: 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=:' 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. --- .../system/driver-manager/source/src/aer.rs | 114 ++++++++++-- .../driver-manager/source/src/pciehp.rs | 105 +++++++++-- .../source/src/unified_events.rs | 176 ++++++++++++++++-- local/sources/base | 2 +- 4 files changed, 352 insertions(+), 45 deletions(-) diff --git a/local/recipes/system/driver-manager/source/src/aer.rs b/local/recipes/system/driver-manager/source/src/aer.rs index 3a7d064cee..c32786a672 100644 --- a/local/recipes/system/driver-manager/source/src/aer.rs +++ b/local/recipes/system/driver-manager/source/src/aer.rs @@ -39,7 +39,7 @@ impl AerEvent { } } -pub fn read_aer_lines(path: &Path, last_seen: &mut u64) -> Option> { +pub fn read_aer_lines(path: &Path, last_seq: &mut u64) -> Option> { if !path.exists() { return None; } @@ -52,12 +52,16 @@ pub fn read_aer_lines(path: &Path, last_seen: &mut u64) -> Option> }; 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> } } +/// Extract the `seq=:` 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); + } } diff --git a/local/recipes/system/driver-manager/source/src/pciehp.rs b/local/recipes/system/driver-manager/source/src/pciehp.rs index b8abc5b619..98acb34d44 100644 --- a/local/recipes/system/driver-manager/source/src/pciehp.rs +++ b/local/recipes/system/driver-manager/source/src/pciehp.rs @@ -57,7 +57,7 @@ impl PciehpEventKind { } } -pub fn read_pciehp_lines(path: &Path, last_seen: &mut u64) -> Option> { +pub fn read_pciehp_lines(path: &Path, last_seq: &mut u64) -> Option> { if !path.exists() { return None; } @@ -70,12 +70,16 @@ pub fn read_pciehp_lines(path: &Path, last_seen: &mut u64) -> Option *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:` 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 { let mut kind = PciehpEventKind::Unknown; let mut device = String::new(); @@ -106,15 +119,6 @@ fn parse_pciehp_line(line: &str) -> Option { } } -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); + } } diff --git a/local/recipes/system/driver-manager/source/src/unified_events.rs b/local/recipes/system/driver-manager/source/src/unified_events.rs index 1aeef4fe6b..c0336527d1 100644 --- a/local/recipes/system/driver-manager/source/src/unified_events.rs +++ b/local/recipes/system/driver-manager/source/src/unified_events.rs @@ -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( ) where F: Fn(&str, ErrorSeverity) -> Option, { + 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); + } } diff --git a/local/sources/base b/local/sources/base index 29166263bb..d98330a7de 160000 --- a/local/sources/base +++ b/local/sources/base @@ -1 +1 @@ -Subproject commit 29166263bbf1d03026ede041e6662ca5a58f3530 +Subproject commit d98330a7de195a9b2cadf379830572f8cb07ad02