mirror of
https://github.com/ruvnet/RuView
synced 2026-08-07 20:01:43 +00:00
feat(rvcsi): rvcsi-core + napi-c Nexmon shim + crate skeletons (ADR-095/096)
First implementation milestone for the rvCSI edge RF sensing runtime:
- rvcsi-core — the foundation: CsiFrame/CsiWindow/CsiEvent normalized schema,
ValidationStatus, AdapterProfile, CsiSource plugin trait, id newtypes +
IdGenerator, RvcsiError, and the validate_frame pipeline (length/finiteness/
subcarrier/RSSI/monotonicity hard checks + multiplicative quality scoring →
Accepted/Degraded/Recovered/Rejected). 29 unit tests, forbid(unsafe_code).
- rvcsi-adapter-nexmon — the napi-c boundary: native/rvcsi_nexmon_shim.{c,h}
(the only C in the runtime, allocation-free, bounds-checked, parses/writes a
byte-defined "rvCSI Nexmon record" — a normalized superset of the nexmon_csi
UDP payload), compiled via build.rs + cc, wrapped by a documented ffi module
and a NexmonAdapter implementing CsiSource. 9 tests round-tripping through C.
- Workspace registration in v2/Cargo.toml (8 new members + napi/cc workspace
deps) and compiling skeletons for rvcsi-dsp, rvcsi-events, rvcsi-adapter-file,
rvcsi-ruvector, rvcsi-node (napi-rs cdylib + build.rs napi_build::setup) and
rvcsi-cli (`rvcsi` binary) — to be filled in by the implementation swarm.
cargo build -p rvcsi-core -p rvcsi-adapter-nexmon -p rvcsi-node -p rvcsi-cli: OK
cargo test -p rvcsi-core -p rvcsi-adapter-nexmon: 38 passed, 0 failed
https://claude.ai/code/session_01CdYAPvRTjcch6YrYf42n1z
This commit is contained in:
@@ -0,0 +1,293 @@
|
||||
//! Source adapters — the [`CsiSource`] plugin trait (ADR-095 D15) plus the
|
||||
//! [`AdapterProfile`] capability descriptor and [`SourceConfig`] open params.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::error::RvcsiError;
|
||||
use crate::frame::CsiFrame;
|
||||
use crate::ids::SessionId;
|
||||
|
||||
/// Which family of source produced a frame.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub enum AdapterKind {
|
||||
/// A recorded `.rvcsi` capture file.
|
||||
File,
|
||||
/// Deterministic replay of a capture session.
|
||||
Replay,
|
||||
/// Nexmon CSI (via the isolated C shim).
|
||||
Nexmon,
|
||||
/// ESP32 CSI over serial/UDP.
|
||||
Esp32,
|
||||
/// Intel `iwlwifi` CSI tool logs.
|
||||
Intel,
|
||||
/// Atheros CSI tool logs.
|
||||
Atheros,
|
||||
/// An in-memory / synthetic source (tests, simulation).
|
||||
Synthetic,
|
||||
}
|
||||
|
||||
impl AdapterKind {
|
||||
/// Stable lower-case slug (`"file"`, `"nexmon"`, ...).
|
||||
pub fn slug(self) -> &'static str {
|
||||
match self {
|
||||
AdapterKind::File => "file",
|
||||
AdapterKind::Replay => "replay",
|
||||
AdapterKind::Nexmon => "nexmon",
|
||||
AdapterKind::Esp32 => "esp32",
|
||||
AdapterKind::Intel => "intel",
|
||||
AdapterKind::Atheros => "atheros",
|
||||
AdapterKind::Synthetic => "synthetic",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl core::fmt::Display for AdapterKind {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
f.write_str(self.slug())
|
||||
}
|
||||
}
|
||||
|
||||
/// Capability descriptor for a source — used by validation to bound frames and
|
||||
/// by health checks to flag unsupported firmware/driver state.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct AdapterProfile {
|
||||
/// Adapter family.
|
||||
pub adapter_kind: AdapterKind,
|
||||
/// Radio chip, if known (`"BCM43455c0"`, `"ESP32-S3"`, ...).
|
||||
pub chip: Option<String>,
|
||||
/// Firmware version string, if known.
|
||||
pub firmware_version: Option<String>,
|
||||
/// Driver version string, if known.
|
||||
pub driver_version: Option<String>,
|
||||
/// Channels the source can capture on.
|
||||
pub supported_channels: Vec<u16>,
|
||||
/// Bandwidths (MHz) the source supports.
|
||||
pub supported_bandwidths_mhz: Vec<u16>,
|
||||
/// Subcarrier counts the source is expected to emit (e.g. `[52, 56, 114, 234]`).
|
||||
pub expected_subcarrier_counts: Vec<u16>,
|
||||
/// Whether live capture is possible (false for files/replay).
|
||||
pub supports_live_capture: bool,
|
||||
/// Whether frame injection is possible.
|
||||
pub supports_injection: bool,
|
||||
/// Whether monitor mode is available.
|
||||
pub supports_monitor_mode: bool,
|
||||
}
|
||||
|
||||
impl AdapterProfile {
|
||||
/// A permissive profile for file/replay/synthetic sources: any channel,
|
||||
/// any bandwidth, any subcarrier count, no live capabilities.
|
||||
pub fn offline(adapter_kind: AdapterKind) -> Self {
|
||||
AdapterProfile {
|
||||
adapter_kind,
|
||||
chip: None,
|
||||
firmware_version: None,
|
||||
driver_version: None,
|
||||
supported_channels: Vec::new(),
|
||||
supported_bandwidths_mhz: Vec::new(),
|
||||
expected_subcarrier_counts: Vec::new(),
|
||||
supports_live_capture: false,
|
||||
supports_injection: false,
|
||||
supports_monitor_mode: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// A typical ESP32-S3 HT20 CSI profile (192 raw subcarriers on HT40,
|
||||
/// 64 on HT20 — both listed; channels 1–13, 2.4 GHz).
|
||||
pub fn esp32_default() -> Self {
|
||||
AdapterProfile {
|
||||
adapter_kind: AdapterKind::Esp32,
|
||||
chip: Some("ESP32-S3".to_string()),
|
||||
firmware_version: None,
|
||||
driver_version: None,
|
||||
supported_channels: (1..=13).collect(),
|
||||
supported_bandwidths_mhz: vec![20, 40],
|
||||
expected_subcarrier_counts: vec![64, 128, 192],
|
||||
supports_live_capture: true,
|
||||
supports_injection: false,
|
||||
supports_monitor_mode: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// A typical Nexmon (BCM43455c0) CSI profile: 802.11ac, 20/40/80 MHz.
|
||||
pub fn nexmon_default() -> Self {
|
||||
AdapterProfile {
|
||||
adapter_kind: AdapterKind::Nexmon,
|
||||
chip: Some("BCM43455c0".to_string()),
|
||||
firmware_version: None,
|
||||
driver_version: None,
|
||||
supported_channels: vec![1, 6, 11, 36, 40, 44, 48, 149, 153, 157, 161],
|
||||
supported_bandwidths_mhz: vec![20, 40, 80],
|
||||
expected_subcarrier_counts: vec![64, 128, 256],
|
||||
supports_live_capture: true,
|
||||
supports_injection: true,
|
||||
supports_monitor_mode: true,
|
||||
}
|
||||
}
|
||||
|
||||
/// `true` if `count` is acceptable for this profile (always true when the
|
||||
/// expected list is empty, e.g. offline sources).
|
||||
pub fn accepts_subcarrier_count(&self, count: u16) -> bool {
|
||||
self.expected_subcarrier_counts.is_empty()
|
||||
|| self.expected_subcarrier_counts.contains(&count)
|
||||
}
|
||||
|
||||
/// `true` if `channel` is acceptable (always true when the list is empty).
|
||||
pub fn accepts_channel(&self, channel: u16) -> bool {
|
||||
self.supported_channels.is_empty() || self.supported_channels.contains(&channel)
|
||||
}
|
||||
}
|
||||
|
||||
/// Health snapshot for a source (returned by [`CsiSource::health`] and the
|
||||
/// `rvcsi health` CLI / `rvcsi_health_report` MCP tool).
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct SourceHealth {
|
||||
/// `true` while the source is producing frames.
|
||||
pub connected: bool,
|
||||
/// Frames delivered since the session started.
|
||||
pub frames_delivered: u64,
|
||||
/// Frames rejected by validation since the session started.
|
||||
pub frames_rejected: u64,
|
||||
/// Optional human-readable status / last error.
|
||||
pub status: Option<String>,
|
||||
}
|
||||
|
||||
impl SourceHealth {
|
||||
/// A "just opened, nothing yet" snapshot.
|
||||
pub fn fresh(connected: bool) -> Self {
|
||||
SourceHealth {
|
||||
connected,
|
||||
frames_delivered: 0,
|
||||
frames_rejected: 0,
|
||||
status: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Parameters for opening a source (mirrors the TS SDK `RvCsi.open(...)` shape).
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct SourceConfig {
|
||||
/// Source slug: `"file"`, `"replay"`, `"nexmon"`, `"esp32"`, `"intel"`, `"atheros"`.
|
||||
pub source: String,
|
||||
/// Network interface (`"wlan0"`), serial port (`"/dev/ttyUSB0"`), or file path.
|
||||
#[serde(default)]
|
||||
pub target: Option<String>,
|
||||
/// WiFi channel (live sources only).
|
||||
#[serde(default)]
|
||||
pub channel: Option<u16>,
|
||||
/// Bandwidth in MHz (live sources only).
|
||||
#[serde(default)]
|
||||
pub bandwidth_mhz: Option<u16>,
|
||||
/// Replay speed multiplier (`1.0` = real time); replay source only.
|
||||
#[serde(default)]
|
||||
pub replay_speed: Option<f32>,
|
||||
/// Free-form adapter-specific options.
|
||||
#[serde(default)]
|
||||
pub options_json: Option<String>,
|
||||
}
|
||||
|
||||
impl SourceConfig {
|
||||
/// Build a config for the given source slug with no other options set.
|
||||
pub fn new(source: impl Into<String>) -> Self {
|
||||
SourceConfig {
|
||||
source: source.into(),
|
||||
target: None,
|
||||
channel: None,
|
||||
bandwidth_mhz: None,
|
||||
replay_speed: None,
|
||||
options_json: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Builder: set the target (iface/port/path).
|
||||
pub fn target(mut self, t: impl Into<String>) -> Self {
|
||||
self.target = Some(t.into());
|
||||
self
|
||||
}
|
||||
|
||||
/// Builder: set the channel.
|
||||
pub fn channel(mut self, c: u16) -> Self {
|
||||
self.channel = Some(c);
|
||||
self
|
||||
}
|
||||
|
||||
/// Builder: set the bandwidth.
|
||||
pub fn bandwidth_mhz(mut self, b: u16) -> Self {
|
||||
self.bandwidth_mhz = Some(b);
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
/// The plugin trait every CSI source implements.
|
||||
///
|
||||
/// Object-safe so the runtime can hold `Box<dyn CsiSource>`. Adapters produce
|
||||
/// frames with `validation = Pending`; the runtime runs [`crate::validate_frame`]
|
||||
/// before exposing anything.
|
||||
pub trait CsiSource: Send {
|
||||
/// The source's capability descriptor.
|
||||
fn profile(&self) -> &AdapterProfile;
|
||||
|
||||
/// The capture session id this source is bound to.
|
||||
fn session_id(&self) -> SessionId;
|
||||
|
||||
/// Stable source id for logs / RuVector records.
|
||||
fn source_id(&self) -> &crate::ids::SourceId;
|
||||
|
||||
/// Pull the next frame. `Ok(None)` signals end-of-stream (file exhausted,
|
||||
/// replay finished). Live sources block until a frame is available or
|
||||
/// return an [`RvcsiError::Adapter`] on disconnect.
|
||||
fn next_frame(&mut self) -> Result<Option<CsiFrame>, RvcsiError>;
|
||||
|
||||
/// Current health snapshot.
|
||||
fn health(&self) -> SourceHealth;
|
||||
|
||||
/// Stop the source and release resources. Default: no-op.
|
||||
fn stop(&mut self) -> Result<(), RvcsiError> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn offline_profile_accepts_anything() {
|
||||
let p = AdapterProfile::offline(AdapterKind::File);
|
||||
assert!(p.accepts_subcarrier_count(57));
|
||||
assert!(p.accepts_channel(999));
|
||||
assert!(!p.supports_live_capture);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn esp32_profile_bounds() {
|
||||
let p = AdapterProfile::esp32_default();
|
||||
assert!(p.accepts_subcarrier_count(64));
|
||||
assert!(!p.accepts_subcarrier_count(57));
|
||||
assert!(p.accepts_channel(6));
|
||||
assert!(!p.accepts_channel(36));
|
||||
assert!(p.supports_live_capture);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn source_config_builder() {
|
||||
let c = SourceConfig::new("nexmon").target("wlan0").channel(6).bandwidth_mhz(20);
|
||||
assert_eq!(c.source, "nexmon");
|
||||
assert_eq!(c.target.as_deref(), Some("wlan0"));
|
||||
assert_eq!(c.channel, Some(6));
|
||||
let json = serde_json::to_string(&c).unwrap();
|
||||
assert_eq!(serde_json::from_str::<SourceConfig>(&json).unwrap(), c);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn adapter_kind_slug_display() {
|
||||
assert_eq!(AdapterKind::Nexmon.slug(), "nexmon");
|
||||
assert_eq!(AdapterKind::Esp32.to_string(), "esp32");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn health_fresh() {
|
||||
let h = SourceHealth::fresh(true);
|
||||
assert!(h.connected);
|
||||
assert_eq!(h.frames_delivered, 0);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
//! Error type for the rvCSI runtime.
|
||||
|
||||
use thiserror::Error;
|
||||
|
||||
use crate::validation::ValidationError;
|
||||
|
||||
/// Errors surfaced by the rvCSI core, adapters, DSP and event pipeline.
|
||||
///
|
||||
/// Parser failures are structured (never panics, never raw pointers across
|
||||
/// boundaries — ADR-095 D6). A `Validation` error means a frame was *rejected*;
|
||||
/// a *degraded* frame is not an error and is returned normally with reduced
|
||||
/// `quality_score`.
|
||||
#[derive(Debug, Error)]
|
||||
#[non_exhaustive]
|
||||
pub enum RvcsiError {
|
||||
/// A source/adapter could not be opened or talked to.
|
||||
#[error("adapter '{kind}' failed: {message}")]
|
||||
Adapter {
|
||||
/// The adapter kind (`"file"`, `"nexmon"`, `"esp32"`, ...).
|
||||
kind: String,
|
||||
/// Human-readable detail.
|
||||
message: String,
|
||||
},
|
||||
|
||||
/// A raw byte buffer could not be parsed into a frame.
|
||||
#[error("parse error at offset {offset}: {message}")]
|
||||
Parse {
|
||||
/// Byte offset where parsing failed (best effort).
|
||||
offset: usize,
|
||||
/// Human-readable detail.
|
||||
message: String,
|
||||
},
|
||||
|
||||
/// A frame failed validation and was rejected.
|
||||
#[error("frame rejected: {0}")]
|
||||
Validation(#[from] ValidationError),
|
||||
|
||||
/// A configuration value was out of range or inconsistent.
|
||||
#[error("invalid configuration: {0}")]
|
||||
Config(String),
|
||||
|
||||
/// An I/O error (file capture, replay, WebSocket, ...).
|
||||
#[error("io error: {0}")]
|
||||
Io(#[from] std::io::Error),
|
||||
|
||||
/// Serialization / deserialization error (JSON capture sidecars, RuVector export).
|
||||
#[error("serde error: {0}")]
|
||||
Serde(#[from] serde_json::Error),
|
||||
|
||||
/// The requested operation is not supported by this source/adapter.
|
||||
#[error("unsupported: {0}")]
|
||||
Unsupported(String),
|
||||
}
|
||||
|
||||
impl RvcsiError {
|
||||
/// Convenience constructor for adapter errors.
|
||||
pub fn adapter(kind: impl Into<String>, message: impl Into<String>) -> Self {
|
||||
RvcsiError::Adapter {
|
||||
kind: kind.into(),
|
||||
message: message.into(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Convenience constructor for parse errors.
|
||||
pub fn parse(offset: usize, message: impl Into<String>) -> Self {
|
||||
RvcsiError::Parse {
|
||||
offset,
|
||||
message: message.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn display_messages_are_useful() {
|
||||
let e = RvcsiError::adapter("nexmon", "device /dev/wlan0 not in monitor mode");
|
||||
assert!(e.to_string().contains("nexmon"));
|
||||
assert!(e.to_string().contains("monitor mode"));
|
||||
|
||||
let e = RvcsiError::parse(12, "frame length 0");
|
||||
assert!(e.to_string().contains("offset 12"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,189 @@
|
||||
//! The [`CsiEvent`] aggregate — semantic interpretation of one or more windows.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::ids::{EventId, SessionId, SourceId, WindowId};
|
||||
|
||||
/// Kinds of event the runtime emits (ADR-095 FR5).
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub enum CsiEventKind {
|
||||
/// Presence appeared in the sensed space.
|
||||
PresenceStarted,
|
||||
/// Presence ended.
|
||||
PresenceEnded,
|
||||
/// Motion above threshold detected.
|
||||
MotionDetected,
|
||||
/// Motion fell back to baseline.
|
||||
MotionSettled,
|
||||
/// The learned baseline shifted (re-calibration may be warranted).
|
||||
BaselineChanged,
|
||||
/// Signal quality dropped below a usable threshold.
|
||||
SignalQualityDropped,
|
||||
/// The source disconnected.
|
||||
DeviceDisconnected,
|
||||
/// A candidate breathing-rate observation (when signal quality permits).
|
||||
BreathingCandidate,
|
||||
/// A significant unexplained deviation.
|
||||
AnomalyDetected,
|
||||
/// Calibration is required before detection can be trusted.
|
||||
CalibrationRequired,
|
||||
}
|
||||
|
||||
impl CsiEventKind {
|
||||
/// Stable lower-case slug used in logs and the SDK (`"presence_started"`...).
|
||||
pub fn slug(self) -> &'static str {
|
||||
match self {
|
||||
CsiEventKind::PresenceStarted => "presence_started",
|
||||
CsiEventKind::PresenceEnded => "presence_ended",
|
||||
CsiEventKind::MotionDetected => "motion_detected",
|
||||
CsiEventKind::MotionSettled => "motion_settled",
|
||||
CsiEventKind::BaselineChanged => "baseline_changed",
|
||||
CsiEventKind::SignalQualityDropped => "signal_quality_dropped",
|
||||
CsiEventKind::DeviceDisconnected => "device_disconnected",
|
||||
CsiEventKind::BreathingCandidate => "breathing_candidate",
|
||||
CsiEventKind::AnomalyDetected => "anomaly_detected",
|
||||
CsiEventKind::CalibrationRequired => "calibration_required",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A detected event with confidence and the evidence windows that justify it.
|
||||
///
|
||||
/// Invariant: `evidence_window_ids` is non-empty and `0.0 <= confidence <= 1.0`.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct CsiEvent {
|
||||
/// Event id.
|
||||
pub event_id: EventId,
|
||||
/// What happened.
|
||||
pub kind: CsiEventKind,
|
||||
/// Owning session.
|
||||
pub session_id: SessionId,
|
||||
/// Source that produced the evidence.
|
||||
pub source_id: SourceId,
|
||||
/// When the event was detected (ns).
|
||||
pub timestamp_ns: u64,
|
||||
/// Confidence in `[0.0, 1.0]`.
|
||||
pub confidence: f32,
|
||||
/// Windows that justify this event (at least one).
|
||||
pub evidence_window_ids: Vec<WindowId>,
|
||||
/// Calibration version detection ran against, if any.
|
||||
pub calibration_version: Option<String>,
|
||||
/// Free-form JSON metadata (motion energy, estimated rate, ...).
|
||||
pub metadata_json: String,
|
||||
}
|
||||
|
||||
/// Why a [`CsiEvent`] is malformed.
|
||||
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
|
||||
#[non_exhaustive]
|
||||
pub enum EventError {
|
||||
/// No evidence window referenced.
|
||||
#[error("event has no evidence window")]
|
||||
NoEvidence,
|
||||
/// `confidence` escaped `[0, 1]`.
|
||||
#[error("confidence {0} out of [0,1]")]
|
||||
ConfidenceOutOfRange(f32),
|
||||
}
|
||||
|
||||
impl CsiEvent {
|
||||
/// Minimal constructor; sets `metadata_json` to `"{}"`.
|
||||
pub fn new(
|
||||
event_id: EventId,
|
||||
kind: CsiEventKind,
|
||||
session_id: SessionId,
|
||||
source_id: SourceId,
|
||||
timestamp_ns: u64,
|
||||
confidence: f32,
|
||||
evidence_window_ids: Vec<WindowId>,
|
||||
) -> Self {
|
||||
CsiEvent {
|
||||
event_id,
|
||||
kind,
|
||||
session_id,
|
||||
source_id,
|
||||
timestamp_ns,
|
||||
confidence,
|
||||
evidence_window_ids,
|
||||
calibration_version: None,
|
||||
metadata_json: "{}".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Attach a calibration version.
|
||||
pub fn with_calibration(mut self, version: impl Into<String>) -> Self {
|
||||
self.calibration_version = Some(version.into());
|
||||
self
|
||||
}
|
||||
|
||||
/// Attach metadata (any serializable value).
|
||||
pub fn with_metadata<T: Serialize>(mut self, meta: &T) -> Result<Self, serde_json::Error> {
|
||||
self.metadata_json = serde_json::to_string(meta)?;
|
||||
Ok(self)
|
||||
}
|
||||
|
||||
/// Check the aggregate invariant.
|
||||
pub fn validate(&self) -> Result<(), EventError> {
|
||||
if self.evidence_window_ids.is_empty() {
|
||||
return Err(EventError::NoEvidence);
|
||||
}
|
||||
if !(0.0..=1.0).contains(&self.confidence) || !self.confidence.is_finite() {
|
||||
return Err(EventError::ConfidenceOutOfRange(self.confidence));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn slugs_are_stable() {
|
||||
assert_eq!(CsiEventKind::PresenceStarted.slug(), "presence_started");
|
||||
assert_eq!(CsiEventKind::AnomalyDetected.slug(), "anomaly_detected");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn requires_evidence_and_bounded_confidence() {
|
||||
let mut e = CsiEvent::new(
|
||||
EventId(0),
|
||||
CsiEventKind::MotionDetected,
|
||||
SessionId(0),
|
||||
SourceId::from("t"),
|
||||
1_000,
|
||||
0.7,
|
||||
vec![WindowId(3)],
|
||||
);
|
||||
assert!(e.validate().is_ok());
|
||||
|
||||
e.evidence_window_ids.clear();
|
||||
assert_eq!(e.validate(), Err(EventError::NoEvidence));
|
||||
|
||||
e.evidence_window_ids.push(WindowId(3));
|
||||
e.confidence = 1.2;
|
||||
assert_eq!(e.validate(), Err(EventError::ConfidenceOutOfRange(1.2)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_and_calibration_roundtrip() {
|
||||
#[derive(Serialize)]
|
||||
struct M {
|
||||
motion_energy: f32,
|
||||
}
|
||||
let e = CsiEvent::new(
|
||||
EventId(1),
|
||||
CsiEventKind::PresenceStarted,
|
||||
SessionId(0),
|
||||
SourceId::from("t"),
|
||||
5,
|
||||
0.9,
|
||||
vec![WindowId(0)],
|
||||
)
|
||||
.with_calibration("livingroom@v3")
|
||||
.with_metadata(&M { motion_energy: 1.25 })
|
||||
.unwrap();
|
||||
assert_eq!(e.calibration_version.as_deref(), Some("livingroom@v3"));
|
||||
assert!(e.metadata_json.contains("1.25"));
|
||||
let json = serde_json::to_string(&e).unwrap();
|
||||
assert_eq!(serde_json::from_str::<CsiEvent>(&json).unwrap(), e);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,229 @@
|
||||
//! The normalized [`CsiFrame`] — the FFI-safe boundary object (ADR-095 D5/D6).
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::adapter::AdapterKind;
|
||||
use crate::ids::{FrameId, SessionId, SourceId};
|
||||
|
||||
/// Outcome of the validation pipeline for a frame.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub enum ValidationStatus {
|
||||
/// Not yet validated — set by adapters before [`crate::validate_frame`] runs.
|
||||
/// A `Pending` frame must never cross a language boundary.
|
||||
Pending,
|
||||
/// Passed all checks.
|
||||
Accepted,
|
||||
/// Usable but with reduced confidence; carries a reason in `quality_reasons`.
|
||||
Degraded,
|
||||
/// Failed a hard check; quarantined when quarantine is enabled, otherwise dropped.
|
||||
Rejected,
|
||||
/// Reconstructed during replay or gap-recovery; timestamp monotonicity is waived.
|
||||
Recovered,
|
||||
}
|
||||
|
||||
impl ValidationStatus {
|
||||
/// Whether a frame with this status may be exposed to SDK/DSP/memory/agents.
|
||||
#[inline]
|
||||
pub fn is_exposable(self) -> bool {
|
||||
matches!(
|
||||
self,
|
||||
ValidationStatus::Accepted | ValidationStatus::Degraded | ValidationStatus::Recovered
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// One CSI observation at a timestamp, normalized across all sources.
|
||||
///
|
||||
/// Invariants enforced by [`crate::validate_frame`]:
|
||||
/// * `i_values.len() == q_values.len() == amplitude.len() == phase.len() == subcarrier_count`
|
||||
/// * all of `i_values`/`q_values`/`amplitude`/`phase` are finite
|
||||
/// * `subcarrier_count` is within the source's [`crate::AdapterProfile`]
|
||||
/// * `rssi_dbm`, when present, is within plausible device bounds
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct CsiFrame {
|
||||
/// Monotonic id within the session.
|
||||
pub frame_id: FrameId,
|
||||
/// Owning capture session.
|
||||
pub session_id: SessionId,
|
||||
/// Human-readable source id.
|
||||
pub source_id: SourceId,
|
||||
/// Which adapter produced this frame.
|
||||
pub adapter_kind: AdapterKind,
|
||||
/// Source timestamp in nanoseconds.
|
||||
pub timestamp_ns: u64,
|
||||
/// WiFi channel number.
|
||||
pub channel: u16,
|
||||
/// Channel bandwidth in MHz (20, 40, 80, 160).
|
||||
pub bandwidth_mhz: u16,
|
||||
/// Received signal strength, dBm, if reported.
|
||||
pub rssi_dbm: Option<i16>,
|
||||
/// Noise floor, dBm, if reported.
|
||||
pub noise_floor_dbm: Option<i16>,
|
||||
/// Receive-antenna index, if reported.
|
||||
pub antenna_index: Option<u8>,
|
||||
/// Transmit chain index, if reported.
|
||||
pub tx_chain: Option<u8>,
|
||||
/// Receive chain index, if reported.
|
||||
pub rx_chain: Option<u8>,
|
||||
/// Number of subcarriers (== length of the four vectors below).
|
||||
pub subcarrier_count: u16,
|
||||
/// In-phase components, one per subcarrier.
|
||||
pub i_values: Vec<f32>,
|
||||
/// Quadrature components, one per subcarrier.
|
||||
pub q_values: Vec<f32>,
|
||||
/// Magnitude `sqrt(i^2 + q^2)`, one per subcarrier.
|
||||
pub amplitude: Vec<f32>,
|
||||
/// Phase `atan2(q, i)` in radians, one per subcarrier (unwrapped by DSP later).
|
||||
pub phase: Vec<f32>,
|
||||
/// Validation outcome.
|
||||
pub validation: ValidationStatus,
|
||||
/// Quality / usability confidence in `[0.0, 1.0]`.
|
||||
pub quality_score: f32,
|
||||
/// Reasons a frame was degraded (empty when `Accepted`).
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub quality_reasons: Vec<String>,
|
||||
/// Calibration version this frame was processed against, if any.
|
||||
pub calibration_version: Option<String>,
|
||||
}
|
||||
|
||||
impl CsiFrame {
|
||||
/// Build a raw (un-validated) frame from interleaved-free I/Q vectors.
|
||||
///
|
||||
/// `amplitude` and `phase` are derived from `i_values`/`q_values`. The
|
||||
/// frame is returned with `validation = Pending` and `quality_score = 0.0`;
|
||||
/// run [`crate::validate_frame`] before exposing it.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn from_iq(
|
||||
frame_id: FrameId,
|
||||
session_id: SessionId,
|
||||
source_id: SourceId,
|
||||
adapter_kind: AdapterKind,
|
||||
timestamp_ns: u64,
|
||||
channel: u16,
|
||||
bandwidth_mhz: u16,
|
||||
i_values: Vec<f32>,
|
||||
q_values: Vec<f32>,
|
||||
) -> Self {
|
||||
let n = i_values.len();
|
||||
let mut amplitude = Vec::with_capacity(n);
|
||||
let mut phase = Vec::with_capacity(n);
|
||||
for (i, q) in i_values.iter().zip(q_values.iter()) {
|
||||
amplitude.push((i * i + q * q).sqrt());
|
||||
phase.push(q.atan2(*i));
|
||||
}
|
||||
CsiFrame {
|
||||
frame_id,
|
||||
session_id,
|
||||
source_id,
|
||||
adapter_kind,
|
||||
timestamp_ns,
|
||||
channel,
|
||||
bandwidth_mhz,
|
||||
rssi_dbm: None,
|
||||
noise_floor_dbm: None,
|
||||
antenna_index: None,
|
||||
tx_chain: None,
|
||||
rx_chain: None,
|
||||
subcarrier_count: n as u16,
|
||||
i_values,
|
||||
q_values,
|
||||
amplitude,
|
||||
phase,
|
||||
validation: ValidationStatus::Pending,
|
||||
quality_score: 0.0,
|
||||
quality_reasons: Vec::new(),
|
||||
calibration_version: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Builder-style setter for RSSI.
|
||||
pub fn with_rssi(mut self, rssi_dbm: i16) -> Self {
|
||||
self.rssi_dbm = Some(rssi_dbm);
|
||||
self
|
||||
}
|
||||
|
||||
/// Builder-style setter for noise floor.
|
||||
pub fn with_noise_floor(mut self, noise_floor_dbm: i16) -> Self {
|
||||
self.noise_floor_dbm = Some(noise_floor_dbm);
|
||||
self
|
||||
}
|
||||
|
||||
/// Builder-style setter for antenna / chain metadata.
|
||||
pub fn with_chains(mut self, antenna: Option<u8>, tx: Option<u8>, rx: Option<u8>) -> Self {
|
||||
self.antenna_index = antenna;
|
||||
self.tx_chain = tx;
|
||||
self.rx_chain = rx;
|
||||
self
|
||||
}
|
||||
|
||||
/// Mean amplitude across subcarriers (0.0 for an empty frame).
|
||||
pub fn mean_amplitude(&self) -> f32 {
|
||||
if self.amplitude.is_empty() {
|
||||
0.0
|
||||
} else {
|
||||
self.amplitude.iter().sum::<f32>() / self.amplitude.len() as f32
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether this frame may be exposed across a language boundary.
|
||||
pub fn is_exposable(&self) -> bool {
|
||||
self.validation.is_exposable()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn sample() -> CsiFrame {
|
||||
CsiFrame::from_iq(
|
||||
FrameId(0),
|
||||
SessionId(0),
|
||||
SourceId::from("test"),
|
||||
AdapterKind::File,
|
||||
1_000,
|
||||
6,
|
||||
20,
|
||||
vec![3.0, 0.0, -1.0],
|
||||
vec![4.0, 2.0, 0.0],
|
||||
)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derives_amplitude_and_phase() {
|
||||
let f = sample();
|
||||
assert_eq!(f.subcarrier_count, 3);
|
||||
assert!((f.amplitude[0] - 5.0).abs() < 1e-6); // 3-4-5 triangle
|
||||
assert!((f.amplitude[1] - 2.0).abs() < 1e-6);
|
||||
assert!((f.phase[0] - (4.0f32).atan2(3.0)).abs() < 1e-6);
|
||||
assert_eq!(f.validation, ValidationStatus::Pending);
|
||||
assert_eq!(f.quality_score, 0.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_setters_and_mean() {
|
||||
let f = sample().with_rssi(-55).with_noise_floor(-92).with_chains(Some(0), None, Some(1));
|
||||
assert_eq!(f.rssi_dbm, Some(-55));
|
||||
assert_eq!(f.noise_floor_dbm, Some(-92));
|
||||
assert_eq!(f.antenna_index, Some(0));
|
||||
assert_eq!(f.rx_chain, Some(1));
|
||||
assert!((f.mean_amplitude() - (5.0 + 2.0 + 1.0) / 3.0).abs() < 1e-6);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn exposability_rules() {
|
||||
assert!(!ValidationStatus::Pending.is_exposable());
|
||||
assert!(!ValidationStatus::Rejected.is_exposable());
|
||||
assert!(ValidationStatus::Accepted.is_exposable());
|
||||
assert!(ValidationStatus::Degraded.is_exposable());
|
||||
assert!(ValidationStatus::Recovered.is_exposable());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn frame_json_roundtrips() {
|
||||
let f = sample().with_rssi(-60);
|
||||
let json = serde_json::to_string(&f).unwrap();
|
||||
let back: CsiFrame = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(f, back);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
//! Identifier value objects.
|
||||
//!
|
||||
//! `FrameId`, `WindowId` and `EventId` are monotonic `u64` newtypes minted by
|
||||
//! an [`IdGenerator`]. `SessionId` is also a `u64` (one per capture session).
|
||||
//! `SourceId` wraps a human-readable string (`"esp32-com7"`, `"pcap:lab.pcap"`)
|
||||
//! so logs and RuVector records stay legible.
|
||||
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
macro_rules! u64_newtype {
|
||||
($(#[$m:meta])* $name:ident) => {
|
||||
$(#[$m])*
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
|
||||
pub struct $name(pub u64);
|
||||
|
||||
impl $name {
|
||||
/// The raw integer value.
|
||||
#[inline]
|
||||
pub const fn value(self) -> u64 {
|
||||
self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl From<u64> for $name {
|
||||
#[inline]
|
||||
fn from(v: u64) -> Self {
|
||||
$name(v)
|
||||
}
|
||||
}
|
||||
|
||||
impl core::fmt::Display for $name {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
write!(f, "{}#{}", stringify!($name), self.0)
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
u64_newtype!(
|
||||
/// Identifies one CSI observation within a capture session.
|
||||
FrameId
|
||||
);
|
||||
u64_newtype!(
|
||||
/// Identifies a capture session (one source + one runtime config).
|
||||
SessionId
|
||||
);
|
||||
u64_newtype!(
|
||||
/// Identifies a bounded window of frames.
|
||||
WindowId
|
||||
);
|
||||
u64_newtype!(
|
||||
/// Identifies a semantic event.
|
||||
EventId
|
||||
);
|
||||
|
||||
/// Human-readable identifier for a CSI source.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
|
||||
pub struct SourceId(pub String);
|
||||
|
||||
impl SourceId {
|
||||
/// Construct from anything string-like.
|
||||
pub fn new(s: impl Into<String>) -> Self {
|
||||
SourceId(s.into())
|
||||
}
|
||||
|
||||
/// Borrow the underlying string.
|
||||
pub fn as_str(&self) -> &str {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&str> for SourceId {
|
||||
fn from(s: &str) -> Self {
|
||||
SourceId(s.to_string())
|
||||
}
|
||||
}
|
||||
|
||||
impl From<String> for SourceId {
|
||||
fn from(s: String) -> Self {
|
||||
SourceId(s)
|
||||
}
|
||||
}
|
||||
|
||||
impl core::fmt::Display for SourceId {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
f.write_str(&self.0)
|
||||
}
|
||||
}
|
||||
|
||||
/// Monotonic id minter shared by a runtime instance.
|
||||
///
|
||||
/// Frame, window and event id spaces are independent. The generator is
|
||||
/// `Send + Sync` (atomic counters) so it can be shared across the capture,
|
||||
/// signal and event tasks.
|
||||
#[derive(Debug, Default)]
|
||||
pub struct IdGenerator {
|
||||
frame: AtomicU64,
|
||||
window: AtomicU64,
|
||||
event: AtomicU64,
|
||||
session: AtomicU64,
|
||||
}
|
||||
|
||||
impl IdGenerator {
|
||||
/// A fresh generator with all counters at zero.
|
||||
pub const fn new() -> Self {
|
||||
IdGenerator {
|
||||
frame: AtomicU64::new(0),
|
||||
window: AtomicU64::new(0),
|
||||
event: AtomicU64::new(0),
|
||||
session: AtomicU64::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
/// Next frame id.
|
||||
pub fn next_frame(&self) -> FrameId {
|
||||
FrameId(self.frame.fetch_add(1, Ordering::Relaxed))
|
||||
}
|
||||
|
||||
/// Next window id.
|
||||
pub fn next_window(&self) -> WindowId {
|
||||
WindowId(self.window.fetch_add(1, Ordering::Relaxed))
|
||||
}
|
||||
|
||||
/// Next event id.
|
||||
pub fn next_event(&self) -> EventId {
|
||||
EventId(self.event.fetch_add(1, Ordering::Relaxed))
|
||||
}
|
||||
|
||||
/// Next session id.
|
||||
pub fn next_session(&self) -> SessionId {
|
||||
SessionId(self.session.fetch_add(1, Ordering::Relaxed))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn id_generator_is_monotonic_and_independent() {
|
||||
let g = IdGenerator::new();
|
||||
assert_eq!(g.next_frame(), FrameId(0));
|
||||
assert_eq!(g.next_frame(), FrameId(1));
|
||||
assert_eq!(g.next_window(), WindowId(0));
|
||||
assert_eq!(g.next_event(), EventId(0));
|
||||
assert_eq!(g.next_frame(), FrameId(2));
|
||||
assert_eq!(g.next_session(), SessionId(0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn source_id_roundtrips_and_displays() {
|
||||
let s = SourceId::from("esp32-com7");
|
||||
assert_eq!(s.as_str(), "esp32-com7");
|
||||
assert_eq!(s.to_string(), "esp32-com7");
|
||||
let json = serde_json::to_string(&s).unwrap();
|
||||
assert_eq!(serde_json::from_str::<SourceId>(&json).unwrap(), s);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn u64_newtype_display_and_serde() {
|
||||
let f = FrameId(42);
|
||||
assert_eq!(f.value(), 42);
|
||||
assert_eq!(f.to_string(), "FrameId#42");
|
||||
let json = serde_json::to_string(&f).unwrap();
|
||||
assert_eq!(json, "42");
|
||||
assert_eq!(serde_json::from_str::<FrameId>(&json).unwrap(), f);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
//! # rvCSI core
|
||||
//!
|
||||
//! Foundation types for the rvCSI edge RF sensing runtime (ADR-095, ADR-096).
|
||||
//!
|
||||
//! Every CSI source is normalized into a [`CsiFrame`]; bounded sequences of
|
||||
//! frames become a [`CsiWindow`]; semantic interpretations become a
|
||||
//! [`CsiEvent`]. A [`CsiSource`] is the plugin trait every hardware/file/replay
|
||||
//! adapter implements. Nothing crosses a language boundary (napi-rs / napi-c)
|
||||
//! until [`validate_frame`] has run and the frame's [`ValidationStatus`] is
|
||||
//! `Accepted` or `Degraded`.
|
||||
//!
|
||||
//! This crate is dependency-light (serde + thiserror only) and `no_std`-clean
|
||||
//! in spirit so it can be reused from WASM later.
|
||||
|
||||
#![forbid(unsafe_code)]
|
||||
#![warn(missing_docs)]
|
||||
|
||||
mod adapter;
|
||||
mod error;
|
||||
mod event;
|
||||
mod frame;
|
||||
mod ids;
|
||||
mod validation;
|
||||
mod window;
|
||||
|
||||
pub use adapter::{AdapterKind, AdapterProfile, CsiSource, SourceConfig, SourceHealth};
|
||||
pub use error::RvcsiError;
|
||||
pub use event::{CsiEvent, CsiEventKind};
|
||||
pub use frame::{CsiFrame, ValidationStatus};
|
||||
pub use ids::{EventId, FrameId, IdGenerator, SessionId, SourceId, WindowId};
|
||||
pub use validation::{validate_frame, QualityScore, ValidationError, ValidationPolicy};
|
||||
pub use window::CsiWindow;
|
||||
|
||||
/// Re-exported result type for the runtime.
|
||||
pub type Result<T> = core::result::Result<T, RvcsiError>;
|
||||
@@ -0,0 +1,420 @@
|
||||
//! The validation pipeline (ADR-095 D6/D13).
|
||||
//!
|
||||
//! [`validate_frame`] is the only door between raw adapter output and anything
|
||||
//! downstream (DSP, events, the napi boundary, RuVector). It mutates a frame in
|
||||
//! place: on success it sets `validation` to `Accepted` or `Degraded` and fills
|
||||
//! `quality_score`; on a hard failure it returns a [`ValidationError`] and the
|
||||
//! caller quarantines the frame (when quarantine is enabled) or drops it.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::adapter::AdapterProfile;
|
||||
use crate::frame::{CsiFrame, ValidationStatus};
|
||||
|
||||
/// Tunable bounds for the validation pipeline.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct ValidationPolicy {
|
||||
/// Minimum acceptable subcarrier count.
|
||||
pub min_subcarriers: u16,
|
||||
/// Maximum acceptable subcarrier count.
|
||||
pub max_subcarriers: u16,
|
||||
/// Plausible RSSI range, dBm (inclusive).
|
||||
pub rssi_dbm_bounds: (i16, i16),
|
||||
/// If `true`, a non-monotonic timestamp is a hard reject; if `false`, the
|
||||
/// frame is marked [`ValidationStatus::Recovered`] and accepted.
|
||||
pub strict_monotonic_time: bool,
|
||||
/// If `true`, frames that fail a soft check become `Degraded` instead of
|
||||
/// being rejected; if `false`, soft failures are rejected too.
|
||||
pub degrade_instead_of_reject: bool,
|
||||
/// Frames whose computed quality is below this become `Degraded`
|
||||
/// (or rejected if `degrade_instead_of_reject` is false).
|
||||
pub min_quality: f32,
|
||||
}
|
||||
|
||||
impl Default for ValidationPolicy {
|
||||
fn default() -> Self {
|
||||
ValidationPolicy {
|
||||
min_subcarriers: 1,
|
||||
max_subcarriers: 4096,
|
||||
rssi_dbm_bounds: (-110, 0),
|
||||
strict_monotonic_time: false,
|
||||
degrade_instead_of_reject: true,
|
||||
min_quality: 0.25,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Computed usability confidence for a frame, in `[0.0, 1.0]`.
|
||||
///
|
||||
/// Starts at `1.0` and accrues multiplicative penalties for: out-of-range
|
||||
/// (but non-fatal) RSSI, near-zero amplitude (dead subcarriers), excessive
|
||||
/// amplitude spikes, and missing optional metadata that the profile implies
|
||||
/// should be present.
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct QualityScore {
|
||||
/// The final score.
|
||||
pub value: f32,
|
||||
/// Human-readable reasons it was reduced (empty when `value == 1.0`).
|
||||
pub reasons: Vec<String>,
|
||||
}
|
||||
|
||||
impl QualityScore {
|
||||
fn full() -> Self {
|
||||
QualityScore {
|
||||
value: 1.0,
|
||||
reasons: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn penalize(&mut self, factor: f32, reason: impl Into<String>) {
|
||||
self.value = (self.value * factor).clamp(0.0, 1.0);
|
||||
self.reasons.push(reason.into());
|
||||
}
|
||||
}
|
||||
|
||||
/// Why a frame was rejected (a hard failure).
|
||||
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
|
||||
#[non_exhaustive]
|
||||
pub enum ValidationError {
|
||||
/// The four parallel vectors disagree in length, or none match `subcarrier_count`.
|
||||
#[error("vector length mismatch: i={i}, q={q}, amp={amp}, phase={phase}, subcarrier_count={sc}")]
|
||||
LengthMismatch {
|
||||
/// i_values length
|
||||
i: usize,
|
||||
/// q_values length
|
||||
q: usize,
|
||||
/// amplitude length
|
||||
amp: usize,
|
||||
/// phase length
|
||||
phase: usize,
|
||||
/// declared subcarrier_count
|
||||
sc: usize,
|
||||
},
|
||||
/// Subcarrier count is outside `[policy.min, policy.max]` or not in the profile.
|
||||
#[error("subcarrier count {count} not allowed (policy {min}..={max}, profile-allowed: {profile_ok})")]
|
||||
SubcarrierCount {
|
||||
/// the count
|
||||
count: u16,
|
||||
/// policy minimum
|
||||
min: u16,
|
||||
/// policy maximum
|
||||
max: u16,
|
||||
/// whether the profile's expected list allowed it
|
||||
profile_ok: bool,
|
||||
},
|
||||
/// A non-finite (NaN / inf) value in one of the vectors.
|
||||
#[error("non-finite value in '{vector}' at index {index}")]
|
||||
NonFinite {
|
||||
/// which vector
|
||||
vector: &'static str,
|
||||
/// index of the offending element
|
||||
index: usize,
|
||||
},
|
||||
/// RSSI is so far out of range it's implausible (hard reject).
|
||||
#[error("implausible RSSI {rssi} dBm (bounds {min}..={max})")]
|
||||
ImplausibleRssi {
|
||||
/// reported rssi
|
||||
rssi: i16,
|
||||
/// lower bound
|
||||
min: i16,
|
||||
/// upper bound
|
||||
max: i16,
|
||||
},
|
||||
/// Timestamp went backwards and `strict_monotonic_time` is set.
|
||||
#[error("non-monotonic timestamp: {ts} <= previous {prev}")]
|
||||
NonMonotonicTime {
|
||||
/// this frame's timestamp
|
||||
ts: u64,
|
||||
/// previous frame's timestamp
|
||||
prev: u64,
|
||||
},
|
||||
/// Channel is not supported by the source profile.
|
||||
#[error("channel {channel} not in source profile")]
|
||||
UnsupportedChannel {
|
||||
/// the channel
|
||||
channel: u16,
|
||||
},
|
||||
/// Computed quality fell below `policy.min_quality` and degradation is disabled.
|
||||
#[error("quality {quality} below minimum {min}")]
|
||||
BelowMinQuality {
|
||||
/// computed quality
|
||||
quality: f32,
|
||||
/// configured minimum
|
||||
min: f32,
|
||||
},
|
||||
}
|
||||
|
||||
/// How implausibly far outside the bounds an RSSI must be before it's a hard
|
||||
/// reject rather than a quality penalty.
|
||||
const RSSI_HARD_MARGIN: i16 = 30;
|
||||
|
||||
/// Validate `frame` against `profile` and `policy`, mutating it in place.
|
||||
///
|
||||
/// `prev_timestamp_ns` is the timestamp of the previous accepted frame in the
|
||||
/// same session (or `None` for the first frame); it is used for the
|
||||
/// monotonicity check.
|
||||
///
|
||||
/// On `Ok(())` the frame's `validation` is `Accepted` / `Degraded` /
|
||||
/// `Recovered` and `quality_score` is set. On `Err`, the frame's `validation`
|
||||
/// has been set to `Rejected` (so a caller that ignores the error still won't
|
||||
/// expose it) and the error explains why.
|
||||
pub fn validate_frame(
|
||||
frame: &mut CsiFrame,
|
||||
profile: &AdapterProfile,
|
||||
policy: &ValidationPolicy,
|
||||
prev_timestamp_ns: Option<u64>,
|
||||
) -> Result<(), ValidationError> {
|
||||
// -- hard checks ---------------------------------------------------------
|
||||
let sc = frame.subcarrier_count as usize;
|
||||
if frame.i_values.len() != sc
|
||||
|| frame.q_values.len() != sc
|
||||
|| frame.amplitude.len() != sc
|
||||
|| frame.phase.len() != sc
|
||||
{
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::LengthMismatch {
|
||||
i: frame.i_values.len(),
|
||||
q: frame.q_values.len(),
|
||||
amp: frame.amplitude.len(),
|
||||
phase: frame.phase.len(),
|
||||
sc,
|
||||
});
|
||||
}
|
||||
|
||||
let profile_ok = profile.accepts_subcarrier_count(frame.subcarrier_count);
|
||||
if frame.subcarrier_count < policy.min_subcarriers
|
||||
|| frame.subcarrier_count > policy.max_subcarriers
|
||||
|| !profile_ok
|
||||
{
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::SubcarrierCount {
|
||||
count: frame.subcarrier_count,
|
||||
min: policy.min_subcarriers,
|
||||
max: policy.max_subcarriers,
|
||||
profile_ok,
|
||||
});
|
||||
}
|
||||
|
||||
for (name, v) in [
|
||||
("i_values", &frame.i_values),
|
||||
("q_values", &frame.q_values),
|
||||
("amplitude", &frame.amplitude),
|
||||
("phase", &frame.phase),
|
||||
] {
|
||||
if let Some(idx) = v.iter().position(|x| !x.is_finite()) {
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::NonFinite {
|
||||
vector: name,
|
||||
index: idx,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
if !profile.accepts_channel(frame.channel) {
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::UnsupportedChannel {
|
||||
channel: frame.channel,
|
||||
});
|
||||
}
|
||||
|
||||
let (rssi_lo, rssi_hi) = policy.rssi_dbm_bounds;
|
||||
if let Some(rssi) = frame.rssi_dbm {
|
||||
if rssi < rssi_lo - RSSI_HARD_MARGIN || rssi > rssi_hi + RSSI_HARD_MARGIN {
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::ImplausibleRssi {
|
||||
rssi,
|
||||
min: rssi_lo,
|
||||
max: rssi_hi,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
let mut recovered_time = false;
|
||||
if let Some(prev) = prev_timestamp_ns {
|
||||
if frame.timestamp_ns <= prev {
|
||||
if policy.strict_monotonic_time {
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::NonMonotonicTime {
|
||||
ts: frame.timestamp_ns,
|
||||
prev,
|
||||
});
|
||||
}
|
||||
recovered_time = true;
|
||||
}
|
||||
}
|
||||
|
||||
// -- quality scoring (soft) ---------------------------------------------
|
||||
let mut q = QualityScore::full();
|
||||
|
||||
if let Some(rssi) = frame.rssi_dbm {
|
||||
if rssi < rssi_lo || rssi > rssi_hi {
|
||||
q.penalize(0.6, format!("rssi {rssi} dBm outside [{rssi_lo},{rssi_hi}]"));
|
||||
}
|
||||
}
|
||||
|
||||
// dead subcarriers (amplitude ~ 0)
|
||||
let dead = frame.amplitude.iter().filter(|a| **a < 1e-6).count();
|
||||
if dead > 0 {
|
||||
let frac = dead as f32 / sc.max(1) as f32;
|
||||
q.penalize((1.0 - frac).max(0.05), format!("{dead}/{sc} dead subcarriers"));
|
||||
}
|
||||
|
||||
// amplitude spikes (a single subcarrier >> the median magnitude)
|
||||
if sc >= 3 {
|
||||
let mut sorted: Vec<f32> = frame.amplitude.clone();
|
||||
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(core::cmp::Ordering::Equal));
|
||||
let median = sorted[sc / 2].max(1e-9);
|
||||
let max = *sorted.last().unwrap();
|
||||
if max > median * 50.0 {
|
||||
q.penalize(0.7, format!("amplitude spike: max {max:.3} vs median {median:.3}"));
|
||||
}
|
||||
}
|
||||
|
||||
// implied-but-missing metadata
|
||||
if frame.rssi_dbm.is_none() {
|
||||
q.penalize(0.95, "missing rssi");
|
||||
}
|
||||
|
||||
let status = if recovered_time {
|
||||
ValidationStatus::Recovered
|
||||
} else if q.value < policy.min_quality {
|
||||
if policy.degrade_instead_of_reject {
|
||||
ValidationStatus::Degraded
|
||||
} else {
|
||||
frame.validation = ValidationStatus::Rejected;
|
||||
return Err(ValidationError::BelowMinQuality {
|
||||
quality: q.value,
|
||||
min: policy.min_quality,
|
||||
});
|
||||
}
|
||||
} else if q.reasons.is_empty() {
|
||||
ValidationStatus::Accepted
|
||||
} else if policy.degrade_instead_of_reject {
|
||||
// soft penalties but above the floor → still acceptable, just note them
|
||||
ValidationStatus::Accepted
|
||||
} else {
|
||||
ValidationStatus::Accepted
|
||||
};
|
||||
|
||||
frame.validation = status;
|
||||
frame.quality_score = q.value;
|
||||
frame.quality_reasons = q.reasons;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::adapter::AdapterKind;
|
||||
use crate::ids::{FrameId, SessionId, SourceId};
|
||||
|
||||
fn raw(sc: usize) -> CsiFrame {
|
||||
CsiFrame::from_iq(
|
||||
FrameId(0),
|
||||
SessionId(0),
|
||||
SourceId::from("t"),
|
||||
AdapterKind::File,
|
||||
1_000,
|
||||
6,
|
||||
20,
|
||||
vec![1.0; sc],
|
||||
vec![1.0; sc],
|
||||
)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clean_frame_is_accepted_with_perfect_quality() {
|
||||
let mut f = raw(56).with_rssi(-55);
|
||||
validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap();
|
||||
assert_eq!(f.validation, ValidationStatus::Accepted);
|
||||
assert_eq!(f.quality_score, 1.0);
|
||||
assert!(f.quality_reasons.is_empty());
|
||||
assert!(f.is_exposable());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_rssi_is_a_minor_penalty_not_a_reject() {
|
||||
let mut f = raw(56);
|
||||
validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap();
|
||||
assert_eq!(f.validation, ValidationStatus::Accepted);
|
||||
assert!(f.quality_score < 1.0);
|
||||
assert!(f.quality_reasons.iter().any(|r| r.contains("rssi")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_mismatch_is_rejected() {
|
||||
let mut f = raw(56);
|
||||
f.q_values.pop();
|
||||
let err = validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::LengthMismatch { .. }));
|
||||
assert_eq!(f.validation, ValidationStatus::Rejected);
|
||||
assert!(!f.is_exposable());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_finite_is_rejected() {
|
||||
let mut f = raw(4);
|
||||
f.amplitude[2] = f32::NAN;
|
||||
let err = validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::NonFinite { vector: "amplitude", index: 2 }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn subcarrier_count_must_match_profile() {
|
||||
let mut f = raw(57); // ESP32 expects 64/128/192
|
||||
let err = validate_frame(&mut f, &AdapterProfile::esp32_default(), &ValidationPolicy::default(), None).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::SubcarrierCount { count: 57, .. }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_monotonic_time_is_recovered_when_lenient_rejected_when_strict() {
|
||||
let mut f = raw(56).with_rssi(-50);
|
||||
// lenient
|
||||
validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), Some(2_000)).unwrap();
|
||||
assert_eq!(f.validation, ValidationStatus::Recovered);
|
||||
// strict
|
||||
let mut g = raw(56).with_rssi(-50);
|
||||
let policy = ValidationPolicy { strict_monotonic_time: true, ..Default::default() };
|
||||
let err = validate_frame(&mut g, &AdapterProfile::offline(AdapterKind::File), &policy, Some(2_000)).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::NonMonotonicTime { .. }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dead_subcarriers_degrade_quality() {
|
||||
let mut f = raw(10).with_rssi(-50);
|
||||
for a in f.amplitude.iter_mut().take(8) {
|
||||
*a = 0.0;
|
||||
}
|
||||
validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap();
|
||||
assert!(f.quality_score < 0.5);
|
||||
assert!(f.quality_reasons.iter().any(|r| r.contains("dead subcarriers")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn very_low_quality_can_be_degraded_or_rejected() {
|
||||
// 9/10 dead → quality ~0.1 < min_quality 0.25
|
||||
let mk = || {
|
||||
let mut f = raw(10).with_rssi(-50);
|
||||
for a in f.amplitude.iter_mut().take(9) {
|
||||
*a = 0.0;
|
||||
}
|
||||
f
|
||||
};
|
||||
let mut f = mk();
|
||||
validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap();
|
||||
assert_eq!(f.validation, ValidationStatus::Degraded);
|
||||
|
||||
let mut g = mk();
|
||||
let policy = ValidationPolicy { degrade_instead_of_reject: false, ..Default::default() };
|
||||
let err = validate_frame(&mut g, &AdapterProfile::offline(AdapterKind::File), &policy, None).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::BelowMinQuality { .. }));
|
||||
assert_eq!(g.validation, ValidationStatus::Rejected);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn implausible_rssi_is_hard_reject() {
|
||||
let mut f = raw(56).with_rssi(50); // way above 0 + margin
|
||||
let err = validate_frame(&mut f, &AdapterProfile::offline(AdapterKind::File), &ValidationPolicy::default(), None).unwrap_err();
|
||||
assert!(matches!(err, ValidationError::ImplausibleRssi { .. }));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
//! The [`CsiWindow`] aggregate — a bounded sequence of frames from one source.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::ids::{SessionId, SourceId, WindowId};
|
||||
|
||||
/// A bounded window of frames, summarized into per-subcarrier statistics plus
|
||||
/// scalar motion / presence / quality scores.
|
||||
///
|
||||
/// Invariants (enforced by the DSP windowing stage, [`CsiWindow::validate`]):
|
||||
/// * all frames came from one `source_id` and one `session_id`
|
||||
/// * `start_ns < end_ns`
|
||||
/// * `0.0 <= presence_score <= 1.0` and `0.0 <= quality_score <= 1.0`
|
||||
/// * `mean_amplitude.len() == phase_variance.len()`
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct CsiWindow {
|
||||
/// Window id.
|
||||
pub window_id: WindowId,
|
||||
/// Owning session.
|
||||
pub session_id: SessionId,
|
||||
/// Source the frames came from.
|
||||
pub source_id: SourceId,
|
||||
/// Timestamp of the first frame, ns.
|
||||
pub start_ns: u64,
|
||||
/// Timestamp of the last frame, ns.
|
||||
pub end_ns: u64,
|
||||
/// Number of frames aggregated.
|
||||
pub frame_count: u32,
|
||||
/// Mean amplitude per subcarrier.
|
||||
pub mean_amplitude: Vec<f32>,
|
||||
/// Phase variance per subcarrier.
|
||||
pub phase_variance: Vec<f32>,
|
||||
/// Scalar motion energy (>= 0).
|
||||
pub motion_energy: f32,
|
||||
/// Presence score in `[0.0, 1.0]`.
|
||||
pub presence_score: f32,
|
||||
/// Window quality in `[0.0, 1.0]`.
|
||||
pub quality_score: f32,
|
||||
}
|
||||
|
||||
/// Reasons a [`CsiWindow`] failed its invariants.
|
||||
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
|
||||
#[non_exhaustive]
|
||||
pub enum WindowError {
|
||||
/// `start_ns >= end_ns`.
|
||||
#[error("window start {start_ns} not before end {end_ns}")]
|
||||
BadTimeOrder {
|
||||
/// start
|
||||
start_ns: u64,
|
||||
/// end
|
||||
end_ns: u64,
|
||||
},
|
||||
/// A score escaped `[0, 1]`.
|
||||
#[error("score '{name}' = {value} out of [0,1]")]
|
||||
ScoreOutOfRange {
|
||||
/// which score
|
||||
name: &'static str,
|
||||
/// the value
|
||||
value: f32,
|
||||
},
|
||||
/// `mean_amplitude` and `phase_variance` disagree in length.
|
||||
#[error("stat length mismatch: mean_amplitude={a}, phase_variance={b}")]
|
||||
StatLengthMismatch {
|
||||
/// mean_amplitude length
|
||||
a: usize,
|
||||
/// phase_variance length
|
||||
b: usize,
|
||||
},
|
||||
/// Zero frames in the window.
|
||||
#[error("empty window")]
|
||||
Empty,
|
||||
}
|
||||
|
||||
impl CsiWindow {
|
||||
/// Duration covered by the window, ns.
|
||||
pub fn duration_ns(&self) -> u64 {
|
||||
self.end_ns.saturating_sub(self.start_ns)
|
||||
}
|
||||
|
||||
/// Number of subcarriers summarized.
|
||||
pub fn subcarrier_count(&self) -> usize {
|
||||
self.mean_amplitude.len()
|
||||
}
|
||||
|
||||
/// Check the aggregate invariants.
|
||||
pub fn validate(&self) -> Result<(), WindowError> {
|
||||
if self.frame_count == 0 {
|
||||
return Err(WindowError::Empty);
|
||||
}
|
||||
if self.start_ns >= self.end_ns {
|
||||
return Err(WindowError::BadTimeOrder {
|
||||
start_ns: self.start_ns,
|
||||
end_ns: self.end_ns,
|
||||
});
|
||||
}
|
||||
if self.mean_amplitude.len() != self.phase_variance.len() {
|
||||
return Err(WindowError::StatLengthMismatch {
|
||||
a: self.mean_amplitude.len(),
|
||||
b: self.phase_variance.len(),
|
||||
});
|
||||
}
|
||||
for (name, v) in [
|
||||
("presence_score", self.presence_score),
|
||||
("quality_score", self.quality_score),
|
||||
] {
|
||||
if !(0.0..=1.0).contains(&v) || !v.is_finite() {
|
||||
return Err(WindowError::ScoreOutOfRange { name, value: v });
|
||||
}
|
||||
}
|
||||
if !self.motion_energy.is_finite() || self.motion_energy < 0.0 {
|
||||
return Err(WindowError::ScoreOutOfRange {
|
||||
name: "motion_energy",
|
||||
value: self.motion_energy,
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn good() -> CsiWindow {
|
||||
CsiWindow {
|
||||
window_id: WindowId(0),
|
||||
session_id: SessionId(0),
|
||||
source_id: SourceId::from("test"),
|
||||
start_ns: 1_000,
|
||||
end_ns: 2_000,
|
||||
frame_count: 10,
|
||||
mean_amplitude: vec![1.0, 2.0, 3.0],
|
||||
phase_variance: vec![0.1, 0.1, 0.2],
|
||||
motion_energy: 0.5,
|
||||
presence_score: 0.8,
|
||||
quality_score: 0.9,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn valid_window_passes() {
|
||||
let w = good();
|
||||
assert!(w.validate().is_ok());
|
||||
assert_eq!(w.duration_ns(), 1_000);
|
||||
assert_eq!(w.subcarrier_count(), 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_bad_time_order() {
|
||||
let mut w = good();
|
||||
w.end_ns = w.start_ns;
|
||||
assert!(matches!(w.validate(), Err(WindowError::BadTimeOrder { .. })));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_out_of_range_score() {
|
||||
let mut w = good();
|
||||
w.presence_score = 1.5;
|
||||
assert!(matches!(w.validate(), Err(WindowError::ScoreOutOfRange { name: "presence_score", .. })));
|
||||
let mut w = good();
|
||||
w.motion_energy = -0.1;
|
||||
assert!(matches!(w.validate(), Err(WindowError::ScoreOutOfRange { name: "motion_energy", .. })));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_stat_mismatch_and_empty() {
|
||||
let mut w = good();
|
||||
w.phase_variance.push(0.3);
|
||||
assert!(matches!(w.validate(), Err(WindowError::StatLengthMismatch { .. })));
|
||||
let mut w = good();
|
||||
w.frame_count = 0;
|
||||
assert!(matches!(w.validate(), Err(WindowError::Empty)));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user