Compare commits

...

5 Commits

22 changed files with 3807 additions and 12 deletions
+6 -5
View File
@@ -88,8 +88,6 @@ jobs:
# ADR-262 P1: `wifi-densepose-rufield` path-deps the `vendor/rufield`
# submodule. Without a recursive checkout the workspace build fails to
# resolve those path deps in CI even though it passes locally.
with:
submodules: recursive
# `wifi-densepose-desktop` is a Tauri v2 app — `glib-sys`, `gtk-sys`,
# `webkit2gtk-sys`, etc. need the Linux dev libraries via pkg-config or the
@@ -146,7 +144,10 @@ jobs:
env:
CARGO_PROFILE_DEV_DEBUG: "0"
CARGO_PROFILE_TEST_DEBUG: "0"
run: cargo test -p wifi-densepose-worldmodel --no-default-features
run: >-
cargo test
--manifest-path crates/worldgraph/wifi-densepose-worldmodel/Cargo.toml
--no-default-features
# ADR-134 CIR tests are behind the `cir` feature so the bench dependency
# (Criterion) only pulls when actually exercised. Run them as a separate
@@ -249,7 +250,7 @@ jobs:
continue-on-error: true
uses: codecov/codecov-action@v6
with:
file: ./coverage.xml
files: ./coverage.xml
flags: unittests
name: codecov-umbrella
@@ -500,4 +501,4 @@ jobs:
**Docker Image:**
`${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}:${{ github.sha }}`
draft: false
prerelease: false
prerelease: false
+7 -4
View File
@@ -115,7 +115,7 @@ jobs:
# RUN guard catches missing ones at build time, this re-checks the
# pushed artifact post-hoc as belt-and-braces).
# 2. /health is up.
# 3. /api/v1/info returns 200 with no auth (LAN-mode default).
# 3. /api/v1/info returns 200 with the explicit trusted-LAN opt-in.
# 4. With RUVIEW_API_TOKEN set, /api/v1/info returns 401 without a
# Bearer header, 200 with the correct one (the #443 auth middleware).
# ---------------------------------------------------------------------
@@ -124,11 +124,14 @@ jobs:
set -euo pipefail
IMAGE="ghcr.io/ruvnet/wifi-densepose:sha-${GITHUB_SHA::7}"
docker pull "$IMAGE"
docker run --rm "$IMAGE" sh -c \
docker run --rm --entrypoint sh "$IMAGE" -c \
'ls /app/ui/observatory.html /app/ui/pose-fusion.html /app/ui/index.html /app/ui/viz.html >/dev/null'
docker run --rm "$IMAGE" sh -c 'ls -d /app/ui/observatory /app/ui/pose-fusion >/dev/null'
docker run --rm --entrypoint sh "$IMAGE" -c 'ls -d /app/ui/observatory /app/ui/pose-fusion >/dev/null'
docker run -d --name sm -p 3000:3000 -e CSI_SOURCE=simulated "$IMAGE"
docker run -d --name sm -p 3000:3000 \
-e CSI_SOURCE=simulated \
-e RUVIEW_ALLOW_UNAUTHENTICATED=1 \
"$IMAGE"
# Wait up to 30 s for /health.
for _ in $(seq 1 30); do
if curl -fsS http://127.0.0.1:3000/health >/dev/null 2>&1; then break; fi
+2
View File
@@ -132,6 +132,8 @@ pip install "ruview[client]" # or: pip install "wifi-densepose[clie
> | **ESP32 Mesh** | 3-6× ESP32-S3 + WiFi router | ~$54 | Yes | Same capabilities as above without the persistent-memory features |
> | **ESP32-C6 research node** ([ADR-110](docs/adr/ADR-110-esp32-c6-firmware-extension.md), [witness](docs/WITNESS-LOG-110.md), [reviewer guide](docs/ADR-110-REVIEW-GUIDE.md), [firmware v0.7.0](https://github.com/ruvnet/RuView/releases/tag/v0.7.0-esp32)) | ESP32-C6-DevKit ($610) | ~$10 | Yes (Wi-Fi 6 capable) | Same CSI pipeline as S3 with the dual-target firmware. **Firmware-side ADR-110 substrate now closed** (v0.7.0): ESP-NOW cross-board mesh quantified at **99.56 % match / 104 µs smoothed offset stdev / 3.95× EMA suppression** over a 5-min two-board soak (witness §A0.10), 32-byte UDP sync packet with operator-tunable cadence (§A0.12), ADR-018 byte 19 bit 4 wire-fix sourced from the working ESP-NOW path (§A0.13). Wire format ready for HE-LTF PPDU tagging in ADR-018 bytes 18-19 (firmware encoder + Rust + Python decoders verified end-to-end across 23 unit tests). LP-core motion-gate RISC-V program and Wi-Fi 6 soft-AP with TWT Responder both ship as opt-in code paths (default off). **Hardware-gated for measurement**: HE-LTF live subcarrier capture needs an 11ax AP (IDF v5.4 doesn't expose AP-side HE config — §A0.6); ~5 µA LP-core hibernation needs an INA meter to capture; 802.15.4 raw RX is broken in IDF v5.4 (workaround: ESP-NOW transport, shipped + measured). See witness log for the empirical / claimed split. |
> | **Research NIC** | Intel 5300 / Atheros AR9580 | ~$50-100 | Yes | Full CSI with 3x3 MIMO |
> | **Qualcomm CSI beta** ([ADR-268](docs/adr/ADR-268-qualcomm-atheros-csi-platform.md)) | QCA9300 now; QCN9074/QCN9274 experimental | ~$30-200 | Simulator now; hardware adapter gated | Rust `QCS1` codec, deterministic replay, UDP/API integration; modern ath11k/ath12k profiles do not claim public CSI export |
> | **Vendor provider beta** ([ADR-270](docs/adr/ADR-270-vendor-rf-sensing-integration-program.md)) | Origin, Plume, Mist, NETGEAR, Electric Imp, RF Solutions, Luma, Nest, Linksys, Wifigarden | Varies | Capability-dependent | Bounded Rust adapters and deterministic fixtures; telemetry/network-only/unsupported states cannot masquerade as CSI |
> | **Any WiFi** | Windows, macOS, or Linux laptop | $0 | No | RSSI-only: coarse presence and motion (see [tutorial #36](https://github.com/ruvnet/RuView/issues/36)) |
>
> No hardware? Verify the signal processing pipeline with the deterministic reference signal: `python archive/v1/data/proof/verify.py`
@@ -0,0 +1,41 @@
# ADR-268: Qualcomm Atheros CSI Platform Strategy
- **Status**: accepted
- **Date**: 2026-07-18
- **Tags**: qualcomm, atheros, csi, ath9k, ath11k, ath12k, simulator
## Context
RuView needs a Qualcomm path that is useful before vendor hardware access while
remaining honest about firmware boundaries. QCA9300 has demonstrated CSI tooling
through ath9k/PicoScenes-class systems. QCN9074 and QCN9274 have upstream Linux
connectivity drivers, but upstream ath11k/ath12k support does not by itself prove
that raw per-packet complex CSI is exported by public firmware.
## Decision
1. Use QCA9300 as the first physical baseline: 802.11n, up to 3x3 MIMO and
20/40 MHz. Accept translated captures from established research tooling.
2. Model QCN9074 (Wi-Fi 6/6E, 4x4, up to 160 MHz) and QCN9274 (Wi-Fi 7, 4x4,
up to 160 MHz in protocol v1) as explicitly experimental simulator profiles.
3. Keep firmware/kernel formats behind a Rust adapter. RuView ingests only the
validated QCS1 application envelope defined by ADR-269.
4. Never label simulated frames as hardware. Physical support requires captured
fixtures, firmware provenance, antenna ordering, scaling and repeatability tests.
5. Prefer an upstream-reviewed Generic Netlink or relay-style export if modern
Qualcomm firmware exposes CFR/CSI; do not depend on undisclosed structs.
## Consequences
- Development, APIs and downstream sensing can be tested immediately.
- QCA9300 offers the shortest path to real Qualcomm data.
- Modern profiles may remain simulator-only until firmware cooperation exists.
- A translation copy is accepted in exchange for a stable, fuzzable boundary.
## Links
- [ADR-269: QCS1 wire protocol](ADR-269-qualcomm-csi-wire-protocol.md)
- [Linux ath11k supported devices](https://wireless.docs.kernel.org/en/latest/en/users/drivers/ath11k.html)
- [PicoScenes supported hardware](https://ps.zpj.io/manual/hardware.html)
- [ADR-270: vendor integration portfolio and acceptance gates](ADR-270-vendor-rf-sensing-integration-program.md)
@@ -0,0 +1,40 @@
# ADR-269: Qualcomm CSI Wire Protocol
- **Status**: accepted
- **Date**: 2026-07-18
- **Tags**: qualcomm, csi, protocol, rust, udp, replay
## Decision
Define `QCS1` version 1 as a vendor-boundary envelope, not a Qualcomm firmware ABI.
It uses a 72-byte little-endian header plus payload and CRC-32/IEEE. The header
records report kind, total length, sequence, monotonic timestamp, device ID,
chipset profile, center frequency, bandwidth, flags, Tx/Rx counts, numeric format,
PPDU type, subcarrier count, noise floor, scale, subcarrier spacing, calibration
ID and payload length.
CSI payloads contain one signed RSSI byte per receive chain followed by
`tx * rx * subcarriers` complex i16 or finite f32 values in Tx-major, Rx-major,
subcarrier-major order. Capability reports carry bounded opaque bytes. One QCS1
frame maps to one UDP datagram; replay files prefix each frame with a little-endian
u32 length.
Parsers fail closed on unknown enums, bad CRC, truncation, trailing datagram data,
non-finite values, inconsistent dimensions, chipset chain/bandwidth violations,
payload mismatches, arithmetic overflow and the IPv4 UDP payload ceiling. A
synthetic flag provides end-to-end simulator provenance.
Version 1 profiles are QCA9300, QCN9074 and QCN9274. QCA9300 is capped at three
chains and 40 MHz; modern profiles are capped at four chains and 160 MHz.
## Consequences
- Simulator, replay and future hardware adapters share one validated Rust API.
- No private firmware layout is represented or redistributed.
- 320 MHz/EHT matrices require segmentation or a later protocol revision.
## Links
- [ADR-268: Qualcomm platform strategy](ADR-268-qualcomm-atheros-csi-platform.md)
- [ADR-267: MediaTek MTC1 protocol](ADR-267-mediatek-mimo-csi-wire-protocol.md)
@@ -0,0 +1,122 @@
# ADR-270: Vendor RF Sensing Integration Program
- **Status**: accepted
- **Date**: 2026-07-18
- **Deciders**: RuView maintainers
- **Tags**: vendors, csi, telemetry, simulator, rust, hardware-validation
## Context
RuView is evaluating Qualcomm, RF Solutions, Origin AI, Plume, Linksys,
Electric Imp, Mist/Juniper, Luma, Google Nest, NETGEAR and Wifigarden. These
names do not represent equivalent integration surfaces: some expose raw CSI,
some expose derived sensing events or network telemetry, and some expose no
supported developer interface. A repeated implementation process must not turn
brand compatibility, Linux connectivity or synthetic fixtures into a false CSI
claim.
## Decision
Adopt a Rust-first provider portfolio with explicit capability negotiation:
- `ComplexCsi`: calibrated per-packet complex channel matrices.
- `DerivedSensing`: vendor-produced motion, occupancy or location events.
- `RfTelemetry`: RSSI, radio, client and topology observations.
- `NetworkOnly`: useful as excitation/AP infrastructure but not a sensor.
- `Unsupported`: no stable, lawful or supportable integration surface.
Every provider follows the same gated loop:
1. Verify an authoritative API/SDK, exact model/chipset and licensing boundary.
2. Write provider and wire/contract ADRs before coupling core code to a vendor.
3. Implement bounded Rust types, explicit capabilities and synthetic provenance.
4. Test deterministic replay, corruption, loss, reconnect, backpressure, schema
evolution and secrets handling.
5. Promote to hardware support only after lawful physical capture on an exact
model/firmware, calibration and repeatability tests, and fixture publication
rights. Simulator success never satisfies this gate.
6. Publish code/release and an upstream or vendor collaboration announcement
that states the measured-versus-simulated boundary.
### Portfolio decisions
| Provider | Classification | Decision |
|---|---|---|
| Qualcomm QCA9300 | `ComplexCsi` candidate | Implement first physical baseline via established ath9k research tooling; QCS1 adapter ships simulator-first. |
| Qualcomm QCN9074/QCN9274 | experimental `ComplexCsi` | Simulator and protocol now; require confirmed ath11k/ath12k firmware export before hardware claim. |
| Origin AI | commercial `DerivedSensing`, possible CSI | Pursue NDA sandbox/API and raw-data rights; isolate proprietary engine behind provider trait/service boundary. |
| Plume/OpenSync | `RfTelemetry`; Plume Sense is gated `DerivedSensing` | Build optional OVSDB/control-plane adapter; negotiate Sense separately and do not infer raw CSI. |
| Mist/Juniper | `RfTelemetry` + location | Conditional read-only REST/webhook adapter for occupancy, RSSI and coordinates; no CSI claim. |
| NETGEAR | partner-gated `RfTelemetry` | Insight adapter only after API access; exact legacy OpenWrt models remain community experiments. |
| Luma | discontinued OpenWrt salvage target | Generic OpenWrt telemetry/pcap fixture only when already owned; no procurement or Luma CSI source. |
| Google Nest Wifi | `NetworkOnly` | Use as traffic/AP infrastructure; Device Access does not expose router CSI or radio telemetry. |
| Linksys | `Unsupported` for sensing | Linksys Aware reached end of support in 2024; record capability probe only, if needed. |
| Electric Imp | scalar IoT/RSSI telemetry | Optional agent/impCentral bridge for existing fleets; reject as CSI acquisition hardware. |
| RF Solutions | non-Wi-Fi RF/IoT telemetry | Exclude from sensing backend; optional RIoT environmental fusion is a separate future concern. |
| Wifigarden | commercial OEM, capability unknown | Hold implementation pending chipset, schema, offline, calibration and data-rights disclosure. |
### Provider boundary
Core code consumes a vendor-neutral `RfSource`-style contract whose capability
set prevents RSSI, location or derived occupancy from being represented as CSI.
Cloud adapters use bounded async queues, regional endpoints, secret-provider
credentials and explicit data provenance. Proprietary device SDKs live behind a
feature-gated FFI or sidecar boundary and are never redistributed without rights.
## Consequences
### Positive
- The integration loop can be repeated without duplicating unsafe parsers.
- Product integrations remain useful even when only telemetry is available.
- Public releases make hardware confidence and simulator confidence distinct.
### Negative
- Several named vendors cannot produce a legitimate CSI implementation today.
- Commercial providers require contracts, subscriptions, test vectors or NDAs.
- Exact hardware revisions and firmware provenance increase validation effort.
### Neutral
- A no-go or telemetry-only ADR is a completed research outcome, not a failed port.
- Vendor status and APIs must be rechecked before each implementation begins.
## Implementation Status
The ADR-270 provider contract is implemented in Rust. Each portfolio entry has
a descriptor, bounded decoder or explicit fail-closed access state, deterministic
contract fixtures where lawful, registry coverage, and API exposure:
- Origin AI: contract-configured derived-sensing decoder and request plan.
- Plume/OpenSync: read-only OVSDB request plan and RF telemetry decoder.
- Mist/Juniper: regional request configuration, paginated RF/location decoder.
- NETGEAR Insight: regional partner request configuration and telemetry decoder.
- Electric Imp and RF Solutions: bounded scalar telemetry bridges.
- Luma: explicitly experimental generic OpenWrt telemetry bridge.
- Google Nest: network-only contract events; never represented as CSI.
- Linksys: `Unsupported` decoder because Linksys Aware is end-of-support.
- Wifigarden: `ContractRequired` decoder pending a disclosed SDK/schema.
`vendor-rf-sim` generates deterministic, provenance-labelled events for the
eight providers with a defined event contract and refuses to fabricate Linksys
or Wifigarden events. The sensing server exposes provider descriptors and latest
events under `/api/v1/rf/vendors` and accepts validated canonical simulator
events over its existing UDP port. Physical/vendor-cloud validation remains
separate from implementation completeness and is reflected by
`hardware_validated: false` until performed.
## Evidence and Links
- [ADR-268: Qualcomm strategy](ADR-268-qualcomm-atheros-csi-platform.md)
- [OpenSync developer sandbox](https://www.opensync.io/developer)
- [Origin AI Wi-Fi sensing architecture](https://www.originwirelessai.com/wifi-sensing/)
- [Juniper Mist webhook hierarchy](https://www.juniper.net/documentation/us/en/software/mist/automation-integration/topics/topic-map/webhook-hierarchy.html)
- [Linksys product end-of-life](https://www.linksys.com/pages/linksys-product-end-of-life)
- [Google Nest Device Access supported devices](https://developers.google.com/nest/device-access/supported-devices)
- [OpenWrt Luma WRTQ-329ACN](https://openwrt.org/toh/hwdata/luma/luma_wrtq-329acn)
- [NETGEAR Insight compatible devices](https://kb.netgear.com/000048452/What-devices-can-I-discover-monitor-and-manage-with-Insight)
- [Electric Imp imp005 hardware guide](https://developer.electricimp.com/hardware/imp/imp005_hardware_guide)
- [RF Solutions company portfolio](https://www.rfsolutions.co.uk/about-us-i1/)
- [Wifigarden service terms](https://policies.wifigarden.com/en-us/terms-of-service)
+22
View File
@@ -0,0 +1,22 @@
# RuView v0.9.2-qualcomm-beta.1
This simulator-first beta adds a Rust Qualcomm Atheros CSI boundary without
claiming modern Qualcomm firmware exports that have not been physically verified.
## Included
- ADR-268 selects QCA9300 as the first physical baseline and treats QCN9074 and
QCN9274 as experimental modern profiles.
- ADR-269 defines the bounded, versioned, CRC-protected `QCS1` protocol.
- `qualcomm-csi-sim` emits deterministic MIMO CSI over UDP or replay files.
- The sensing server validates QCS1 datagrams, broadcasts bounded summaries and
exposes `/api/v1/csi/qualcomm/latest`.
- `qualcomm:simulated` provenance is retained end to end.
## Validation boundary
Codec, corruption, truncation, finite-value, dimensions, chipset bandwidth,
determinism and prefix parsing are automated. Loopback UDP/API validation covers
all profiles. Physical QCA9300 comparison and modern firmware export validation
remain hardware gates and will be published with firmware and calibration details.
@@ -0,0 +1,21 @@
# RuView v0.9.3 Vendor Providers Beta 1
This beta implements ADR-270 as a capability-safe Rust provider program across
all ten researched vendors.
## Included
- Shared `VendorRfProvider` contract with bounded event validation.
- Origin AI, Plume/OpenSync, Mist/Juniper, NETGEAR Insight, Electric Imp,
RF Solutions, Luma/OpenWrt and Google Nest contract adapters.
- Explicit fail-closed Linksys (`Unsupported`) and Wifigarden
(`ContractRequired`) providers.
- Deterministic `vendor-rf-sim` JSONL/UDP fixtures for defined contracts.
- Provider registry, descriptors, latest-event REST endpoints and WebSocket
summaries through the sensing server.
## Boundary
This release implements and validates software contracts. It does not claim
vendor-cloud credentials, commercial SDK rights, physical hardware validation,
or complex CSI support for telemetry-only providers.
+60
View File
@@ -0,0 +1,60 @@
# ADR-270 Vendor RF Providers
RuView exposes a capability-safe Rust provider layer for vendor sensing and RF
telemetry. It never converts RSSI, occupancy, location or network inventory into
complex CSI.
## API
- `GET /api/v1/rf/vendors` — all provider descriptors and access states.
- `GET /api/v1/rf/vendors/latest` — latest validated event per vendor.
- `GET /api/v1/rf/vendors/:vendor/latest` — latest event for one stable vendor ID.
- `POST /api/v1/rf/vendors/:vendor/events` — ingest the vendor's documented
sidecar/webhook payload through its strict provider decoder. This `/api/v1/*`
route uses the server's bearer-token policy when configured.
Stable IDs are `origin_ai`, `plume`, `mist`, `netgear`, `electric_imp`,
`rf_solutions`, `linksys`, `luma`, `google_nest`, and `wifigarden`.
## Deterministic simulator
```bash
cd v2
cargo run -p wifi-densepose-hardware --bin vendor-rf-sim -- \
--vendor plume --frames 100 --output plume.jsonl
# Stream canonical synthetic events to the sensing server UDP port.
cargo run -p wifi-densepose-hardware --bin vendor-rf-sim -- \
--vendor mist --frames 100 --udp 127.0.0.1:5005 --realtime
```
Supported simulator names are `origin-ai`, `plume`, `mist`, `netgear`,
`electric-imp`, `rf-solutions`, `luma`, and `google-nest`. Linksys is refused
because its sensing service is discontinued. Wifigarden is refused until a
contracted event schema exists.
Every synthetic event includes `synthetic: true`, a deterministic sequence and
timestamp, and a source ending in `-sim-01`.
Canonical UDP JSON is accepted only when `synthetic: true`. Live vendor payloads
must use the HTTP ingestion route so provider-specific schemas, metric allowlists,
access states and bounds cannot be bypassed.
## Live/provider payloads
Provider decoders are strict, bounded and reject unknown schema fields. Origin
paths and credentials are supplied by the commercial contract. Plume uses a
read-only allow-listed OVSDB request plan. Mist and NETGEAR configurations use
regional HTTPS endpoints with redacted tokens. Electric Imp, RF Solutions and
Luma accept only allow-listed scalar metrics. Google Nest remains network-only.
Credentials are never embedded in fixtures or descriptors. Linksys returns
`Unsupported`; Wifigarden returns `ContractRequired`. These are usable,
test-covered provider outcomes—not simulated integrations.
## Hardware honesty
All descriptors remain `hardware_validated: false` until exact hardware/cloud
versions, lawful access, repeatable captures, calibration where applicable, and
fixture publication rights have been verified. Passing the simulator and API
tests validates RuView software only.
@@ -0,0 +1,147 @@
//! Deterministic Qualcomm Atheros MIMO CSI simulator (ADR-268/269).
use clap::{Parser, ValueEnum};
use std::{
fs::File,
io::{self, Write},
net::{SocketAddr, UdpSocket},
path::PathBuf,
thread,
time::Duration,
};
use wifi_densepose_hardware::qualcomm_csi::{
simulator::{QualcommCsiSimulator, SimulatorConfig},
ChipsetProfile, CsiFrame,
};
#[derive(Debug, Clone, Copy, ValueEnum)]
enum Profile {
Qca9300,
Qcn9074,
Qcn9274,
}
impl Profile {
fn chipset(self) -> ChipsetProfile {
match self {
Self::Qca9300 => ChipsetProfile::Qca9300,
Self::Qcn9074 => ChipsetProfile::Qcn9074,
Self::Qcn9274 => ChipsetProfile::Qcn9274,
}
}
fn default_chains(self) -> u8 {
match self {
Self::Qca9300 => 3,
Self::Qcn9074 | Self::Qcn9274 => 4,
}
}
fn default_bandwidth(self) -> u16 {
match self {
Self::Qca9300 => 40,
Self::Qcn9074 | Self::Qcn9274 => 80,
}
}
fn default_subcarriers(self) -> u16 {
match self {
Self::Qca9300 => 114,
Self::Qcn9074 | Self::Qcn9274 => 256,
}
}
}
#[derive(Debug, Parser)]
#[command(
name = "qualcomm-csi-sim",
about = "Emit synthetic ADR-269 Qualcomm Atheros MIMO CSI frames"
)]
struct Args {
#[arg(long, value_enum, default_value_t=Profile::Qca9300)]
profile: Profile,
#[arg(long, default_value_t = 100)]
frames: u32,
#[arg(long, default_value="0x5143414353490001", value_parser=parse_u64)]
seed: u64,
#[arg(long)]
bandwidth: Option<u16>,
#[arg(long, default_value_t = 2)]
tx: u8,
#[arg(long)]
rx: Option<u8>,
#[arg(long)]
subcarriers: Option<u16>,
#[arg(long, default_value_t = 20)]
interval_ms: u64,
#[arg(long)]
udp: Option<SocketAddr>,
/// Replay: little-endian u32 length followed by one ADR-269 envelope.
#[arg(long)]
output: Option<PathBuf>,
#[arg(long)]
realtime: bool,
}
fn parse_u64(v: &str) -> Result<u64, String> {
if let Some(h) = v.strip_prefix("0x").or_else(|| v.strip_prefix("0X")) {
u64::from_str_radix(h, 16).map_err(|e| e.to_string())
} else {
v.parse()
.map_err(|e: std::num::ParseIntError| e.to_string())
}
}
fn emit(
frame: CsiFrame,
socket: Option<&UdpSocket>,
destination: Option<SocketAddr>,
output: &mut Option<File>,
) -> Result<usize, Box<dyn std::error::Error>> {
let wire = frame.to_bytes()?;
if let (Some(s), Some(d)) = (socket, destination) {
if s.send_to(&wire, d)? != wire.len() {
return Err(io::Error::new(io::ErrorKind::WriteZero, "partial UDP datagram").into());
}
}
if let Some(f) = output {
f.write_all(&(wire.len() as u32).to_le_bytes())?;
f.write_all(&wire)?;
}
Ok(wire.len())
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let a = Args::parse();
if a.udp.is_none() && a.output.is_none() {
return Err("select at least one sink with --udp or --output".into());
}
let cfg = SimulatorConfig {
seed: a.seed,
chipset: a.profile.chipset(),
bandwidth_mhz: a.bandwidth.unwrap_or_else(|| a.profile.default_bandwidth()),
tx_count: a.tx,
rx_count: a.rx.unwrap_or_else(|| a.profile.default_chains()),
subcarriers: a
.subcarriers
.unwrap_or_else(|| a.profile.default_subcarriers()),
frame_period_us: a.interval_ms * 1000,
..Default::default()
};
let mut sim = QualcommCsiSimulator::new(cfg)?;
let socket = a.udp.map(|_| UdpSocket::bind("0.0.0.0:0")).transpose()?;
let mut output = a.output.as_ref().map(File::create).transpose()?;
let mut bytes = emit(
sim.capabilities_frame(),
socket.as_ref(),
a.udp,
&mut output,
)?;
for _ in 0..a.frames {
bytes += emit(sim.next_frame(), socket.as_ref(), a.udp, &mut output)?;
if a.realtime {
thread::sleep(Duration::from_millis(a.interval_ms));
}
}
eprintln!(
"emitted {} synthetic Qualcomm CSI frames ({} bytes, profile={}, seed={:#x})",
a.frames + 1,
bytes,
a.profile.chipset().name(),
a.seed
);
Ok(())
}
@@ -0,0 +1,232 @@
//! Deterministic ADR-270 vendor event simulator.
use clap::{Parser, ValueEnum};
use std::{
collections::BTreeMap,
fs::File,
io::{self, Write},
net::{SocketAddr, UdpSocket},
path::PathBuf,
thread,
time::Duration,
};
use wifi_densepose_hardware::vendor_rf::{
ProviderAvailability, ProviderDescriptor, RfCapability, VendorId, VendorRfEvent,
};
#[derive(Debug, Clone, Copy, ValueEnum)]
enum Vendor {
OriginAi,
Plume,
Mist,
Netgear,
ElectricImp,
RfSolutions,
Luma,
GoogleNest,
Linksys,
Wifigarden,
}
impl Vendor {
fn id(self) -> VendorId {
match self {
Self::OriginAi => VendorId::OriginAi,
Self::Plume => VendorId::Plume,
Self::Mist => VendorId::Mist,
Self::Netgear => VendorId::Netgear,
Self::ElectricImp => VendorId::ElectricImp,
Self::RfSolutions => VendorId::RfSolutions,
Self::Luma => VendorId::Luma,
Self::GoogleNest => VendorId::GoogleNest,
Self::Linksys => VendorId::Linksys,
Self::Wifigarden => VendorId::Wifigarden,
}
}
fn descriptor(self) -> ProviderDescriptor {
let (capabilities, availability, reason) = match self {
Self::OriginAi => (
vec![RfCapability::DerivedSensing],
ProviderAvailability::ContractRequired,
"Origin partner API/SDK contract required",
),
Self::Plume => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::CredentialsRequired,
"OpenSync telemetry; Plume Sense is separately gated",
),
Self::Mist => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::CredentialsRequired,
"Mist REST/webhook telemetry",
),
Self::Netgear => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::CredentialsRequired,
"Insight partner API telemetry",
),
Self::ElectricImp => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::CredentialsRequired,
"impCentral scalar telemetry only",
),
Self::RfSolutions => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::CredentialsRequired,
"RIoT environmental telemetry only",
),
Self::Luma => (
vec![RfCapability::RfTelemetry],
ProviderAvailability::Experimental,
"discontinued OpenWrt salvage fixture",
),
Self::GoogleNest => (
vec![RfCapability::NetworkOnly],
ProviderAvailability::Experimental,
"network infrastructure contract fixture only",
),
Self::Linksys => (
vec![RfCapability::Unsupported],
ProviderAvailability::Unsupported,
"Linksys Aware reached end of support",
),
Self::Wifigarden => (
vec![RfCapability::Unsupported],
ProviderAvailability::ContractRequired,
"technical SDK disclosure required",
),
};
ProviderDescriptor {
vendor: self.id(),
capabilities,
availability,
hardware_validated: false,
reason: reason.into(),
}
}
}
#[derive(Debug, Parser)]
#[command(
name = "vendor-rf-sim",
about = "Emit deterministic ADR-270 vendor RF events"
)]
struct Args {
#[arg(long, value_enum)]
vendor: Vendor,
#[arg(long, default_value_t = 100)]
frames: u64,
#[arg(long, default_value_t = 0x5255_5645_4e44_4f52)]
seed: u64,
#[arg(long, default_value_t = 100)]
interval_ms: u64,
#[arg(long)]
udp: Option<SocketAddr>,
#[arg(long)]
output: Option<PathBuf>,
#[arg(long)]
realtime: bool,
}
fn next_random(state: &mut u64) -> f64 {
*state ^= *state << 13;
*state ^= *state >> 7;
*state ^= *state << 17;
(*state >> 11) as f64 / ((1u64 << 53) as f64)
}
fn event(vendor: Vendor, sequence: u64, timestamp_us: u64, state: &mut u64) -> VendorRfEvent {
let capability = vendor.descriptor().capabilities[0];
let wave = (sequence as f64 * 0.17).sin();
let noise = next_random(state) - 0.5;
let metrics = match capability {
RfCapability::DerivedSensing => BTreeMap::from([
(
"motion_score".into(),
(0.5 + 0.35 * wave + 0.05 * noise).clamp(0.0, 1.0),
),
("occupancy_count".into(), if wave > 0.0 { 2.0 } else { 1.0 }),
("confidence".into(), 0.92),
]),
RfCapability::RfTelemetry => BTreeMap::from([
("rssi_dbm".into(), -52.0 + 5.0 * wave + noise),
("client_count".into(), if wave > 0.0 { 4.0 } else { 3.0 }),
(
"channel_utilization".into(),
(0.31 + 0.08 * wave).clamp(0.0, 1.0),
),
]),
RfCapability::NetworkOnly => {
BTreeMap::from([("reachable".into(), 1.0), ("device_count".into(), 5.0)])
}
_ => BTreeMap::new(),
};
VendorRfEvent {
vendor: vendor.id(),
capability,
sequence,
timestamp_us,
source_id: format!("{}-sim-01", vendor.id().as_str()),
synthetic: true,
metrics,
label: Some("deterministic_contract_fixture".into()),
}
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let args = Args::parse();
if args.udp.is_none() && args.output.is_none() {
return Err("select at least one sink with --udp or --output".into());
}
let descriptor = args.vendor.descriptor();
descriptor.validate()?;
if matches!(
descriptor.availability,
ProviderAvailability::Unsupported | ProviderAvailability::ContractRequired
) && descriptor.capabilities.contains(&RfCapability::Unsupported)
{
return Err(format!(
"{} has no simulatable event contract: {}",
descriptor.vendor.as_str(),
descriptor.reason
)
.into());
}
let socket = args.udp.map(|_| UdpSocket::bind("0.0.0.0:0")).transpose()?;
let mut output = args.output.as_ref().map(File::create).transpose()?;
let mut rng = args.seed;
let mut bytes = 0usize;
for sequence in 0..args.frames {
let value = event(
args.vendor,
sequence,
sequence * args.interval_ms * 1_000,
&mut rng,
);
value.validate(&descriptor)?;
let wire = serde_json::to_vec(&value)?;
if let (Some(socket), Some(destination)) = (&socket, args.udp) {
if socket.send_to(&wire, destination)? != wire.len() {
return Err(
io::Error::new(io::ErrorKind::WriteZero, "partial UDP datagram").into(),
);
}
}
if let Some(file) = &mut output {
file.write_all(&wire)?;
file.write_all(b"\n")?;
}
bytes += wire.len();
if args.realtime {
thread::sleep(Duration::from_millis(args.interval_ms));
}
}
eprintln!(
"emitted {} synthetic {} events ({} bytes, seed={:#x})",
args.frames,
args.vendor.id().as_str(),
bytes,
args.seed
);
Ok(())
}
@@ -47,6 +47,8 @@ mod esp32_parser;
// standardized report path until an OTA binding exists.
pub mod ieee80211bf;
pub mod sync_packet;
/// ADR-270 capability-safe vendor RF provider contract.
pub mod vendor_rf;
// ADR-081: Rust mirror of the firmware radio abstraction layer (L1) and
// mesh sensing plane (L3). Lets host tests, simulators, and future
@@ -55,6 +57,8 @@ pub mod sync_packet;
pub mod radio_ops;
/// ADR-267 vendor-neutral MediaTek Filogic MIMO CSI framing and simulator.
pub mod mediatek_csi;
/// ADR-269 vendor-neutral Qualcomm Atheros CSI framing and simulator.
pub mod qualcomm_csi;
/// ADR-264 host-side framing for Realtek RTL8720F CFR and FMCW radar reports.
/// This module has no dependency on the vendor SDK.
pub mod rtl8720f;
@@ -82,6 +86,13 @@ pub use mediatek_csi::{
PpduType as MediatekPpduType, ReportKind as MediatekReportKind,
MEDIATEK_CSI_HEADER_LEN, MEDIATEK_CSI_MAGIC, MEDIATEK_CSI_VERSION,
};
pub use qualcomm_csi::{
ChipsetProfile as QualcommChipsetProfile, CsiFlags as QualcommCsiFlags,
CsiFrame as QualcommCsiFrame, CsiParseError as QualcommCsiParseError,
CsiPayload as QualcommCsiPayload, ElementFormat as QualcommElementFormat,
PpduType as QualcommPpduType, ReportKind as QualcommReportKind,
QUALCOMM_CSI_HEADER_LEN, QUALCOMM_CSI_MAGIC, QUALCOMM_CSI_VERSION,
};
pub use rtl8720f::{
ElementFormat as Rtl8720fElementFormat, RadarFlags as Rtl8720fRadarFlags,
RadarFrame as Rtl8720fRadarFrame, RadarParseError as Rtl8720fRadarParseError,
@@ -91,3 +102,7 @@ pub use rtl8720f::{
pub use sync_packet::{
SyncPacket, SyncPacketFlags, SYNC_PACKET_MAGIC, SYNC_PACKET_PROTO_VER, SYNC_PACKET_SIZE,
};
pub use vendor_rf::{
ProviderAvailability, ProviderDescriptor, RfCapability, VendorEventError, VendorId,
VendorRfEvent, VendorRfProvider,
};
@@ -0,0 +1,700 @@
//! Vendor-neutral Qualcomm Atheros MIMO CSI transport and deterministic simulator.
//! This is not a Qualcomm firmware ABI; see ADR-268/269.
use serde::{Deserialize, Serialize};
use thiserror::Error;
pub const QUALCOMM_CSI_MAGIC: u32 = 0x3153_4351; // "QCS1" little endian
pub const QUALCOMM_CSI_VERSION: u8 = 1;
pub const QUALCOMM_CSI_HEADER_LEN: usize = 72;
pub const QUALCOMM_CSI_CRC_LEN: usize = 4;
pub const QUALCOMM_CSI_MAX_FRAME_LEN: usize = 65_507;
pub const QUALCOMM_CSI_MAX_ELEMENTS: usize = 16_384;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum ReportKind {
Csi = 1,
Capabilities = 2,
}
impl TryFrom<u8> for ReportKind {
type Error = CsiParseError;
fn try_from(value: u8) -> Result<Self, Self::Error> {
match value {
1 => Ok(Self::Csi),
2 => Ok(Self::Capabilities),
_ => Err(CsiParseError::UnknownReportKind(value)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u16)]
pub enum ChipsetProfile {
Qca9300 = 1,
Qcn9074 = 2,
Qcn9274 = 3,
}
impl TryFrom<u16> for ChipsetProfile {
type Error = CsiParseError;
fn try_from(value: u16) -> Result<Self, Self::Error> {
match value {
1 => Ok(Self::Qca9300),
2 => Ok(Self::Qcn9074),
3 => Ok(Self::Qcn9274),
_ => Err(CsiParseError::UnknownChipset(value)),
}
}
}
impl ChipsetProfile {
pub fn name(self) -> &'static str {
match self {
Self::Qca9300 => "qca9300",
Self::Qcn9074 => "qcn9074",
Self::Qcn9274 => "qcn9274",
}
}
pub fn max_chains(self) -> u8 {
match self {
Self::Qca9300 => 3,
Self::Qcn9074 | Self::Qcn9274 => 4,
}
}
pub fn max_bandwidth_mhz(self) -> u16 {
match self {
Self::Qca9300 => 40,
Self::Qcn9074 | Self::Qcn9274 => 160,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum ElementFormat {
ComplexI16 = 1,
ComplexF32 = 2,
Bytes = 3,
}
impl TryFrom<u8> for ElementFormat {
type Error = CsiParseError;
fn try_from(value: u8) -> Result<Self, Self::Error> {
match value {
1 => Ok(Self::ComplexI16),
2 => Ok(Self::ComplexF32),
3 => Ok(Self::Bytes),
_ => Err(CsiParseError::UnknownElementFormat(value)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum PpduType {
Ht = 1,
Vht = 2,
HeSu = 3,
HeMu = 4,
Eht = 5,
}
impl TryFrom<u8> for PpduType {
type Error = CsiParseError;
fn try_from(value: u8) -> Result<Self, Self::Error> {
match value {
1 => Ok(Self::Ht),
2 => Ok(Self::Vht),
3 => Ok(Self::HeSu),
4 => Ok(Self::HeMu),
5 => Ok(Self::Eht),
_ => Err(CsiParseError::UnknownPpduType(value)),
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct CsiFlags(pub u16);
impl CsiFlags {
pub const CALIBRATED: u16 = 1 << 0;
pub const SATURATED: u16 = 1 << 1;
pub const TIME_SYNCHRONIZED: u16 = 1 << 2;
pub const DROPPED_PREDECESSOR: u16 = 1 << 3;
pub const SYNTHETIC: u16 = 1 << 15;
pub fn contains(self, flag: u16) -> bool {
self.0 & flag != 0
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum CsiPayload {
ComplexI16 {
rssi_dbm: Vec<i8>,
values: Vec<[i16; 2]>,
},
ComplexF32 {
rssi_dbm: Vec<i8>,
values: Vec<[f32; 2]>,
},
Bytes(Vec<u8>),
}
impl CsiPayload {
pub fn len(&self) -> usize {
match self {
Self::ComplexI16 { values, .. } => values.len(),
Self::ComplexF32 { values, .. } => values.len(),
Self::Bytes(values) => values.len(),
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn rssi_dbm(&self) -> &[i8] {
match self {
Self::ComplexI16 { rssi_dbm, .. } | Self::ComplexF32 { rssi_dbm, .. } => rssi_dbm,
Self::Bytes(_) => &[],
}
}
fn format(&self) -> ElementFormat {
match self {
Self::ComplexI16 { .. } => ElementFormat::ComplexI16,
Self::ComplexF32 { .. } => ElementFormat::ComplexF32,
Self::Bytes(_) => ElementFormat::Bytes,
}
}
fn encoded_len(&self) -> usize {
match self {
Self::ComplexI16 { rssi_dbm, values } => rssi_dbm.len() + values.len() * 4,
Self::ComplexF32 { rssi_dbm, values } => rssi_dbm.len() + values.len() * 8,
Self::Bytes(values) => values.len(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CsiFrame {
pub report_kind: ReportKind,
pub sequence: u32,
pub timestamp_us: u64,
pub device_id: u64,
pub chipset: ChipsetProfile,
pub bandwidth_mhz: u16,
pub center_freq_khz: u32,
pub flags: CsiFlags,
pub tx_count: u8,
pub rx_count: u8,
pub ppdu_type: PpduType,
pub subcarrier_count: u16,
pub noise_floor_dbm: i8,
pub scale: f32,
pub subcarrier_spacing_hz: f32,
pub calibration_id: u32,
pub payload: CsiPayload,
}
impl CsiFrame {
pub fn to_bytes(&self) -> Result<Vec<u8>, CsiParseError> {
self.validate()?;
let payload_len = self.payload.encoded_len();
let frame_len = QUALCOMM_CSI_HEADER_LEN
.checked_add(payload_len)
.and_then(|n| n.checked_add(QUALCOMM_CSI_CRC_LEN))
.ok_or(CsiParseError::LengthOverflow)?;
if frame_len > QUALCOMM_CSI_MAX_FRAME_LEN {
return Err(CsiParseError::FrameTooLarge(frame_len));
}
let mut out = Vec::with_capacity(frame_len);
out.extend_from_slice(&QUALCOMM_CSI_MAGIC.to_le_bytes());
out.push(QUALCOMM_CSI_VERSION);
out.push(self.report_kind as u8);
out.extend_from_slice(&(QUALCOMM_CSI_HEADER_LEN as u16).to_le_bytes());
out.extend_from_slice(&(frame_len as u32).to_le_bytes());
out.extend_from_slice(&self.sequence.to_le_bytes());
out.extend_from_slice(&self.timestamp_us.to_le_bytes());
out.extend_from_slice(&self.device_id.to_le_bytes());
out.extend_from_slice(&(self.chipset as u16).to_le_bytes());
out.extend_from_slice(&self.bandwidth_mhz.to_le_bytes());
out.extend_from_slice(&self.center_freq_khz.to_le_bytes());
out.extend_from_slice(&self.flags.0.to_le_bytes());
out.push(self.tx_count);
out.push(self.rx_count);
out.push(self.payload.format() as u8);
out.push(self.ppdu_type as u8);
out.extend_from_slice(&self.subcarrier_count.to_le_bytes());
out.push(self.payload.rssi_dbm().len() as u8);
out.push(self.noise_floor_dbm as u8);
out.extend_from_slice(&0u16.to_le_bytes());
out.extend_from_slice(&self.scale.to_le_bytes());
out.extend_from_slice(&self.subcarrier_spacing_hz.to_le_bytes());
out.extend_from_slice(&self.calibration_id.to_le_bytes());
out.extend_from_slice(&(payload_len as u32).to_le_bytes());
out.extend_from_slice(&0u32.to_le_bytes());
debug_assert_eq!(out.len(), QUALCOMM_CSI_HEADER_LEN);
match &self.payload {
CsiPayload::ComplexI16 { rssi_dbm, values } => {
out.extend(rssi_dbm.iter().map(|v| *v as u8));
for [i, q] in values {
out.extend_from_slice(&i.to_le_bytes());
out.extend_from_slice(&q.to_le_bytes());
}
}
CsiPayload::ComplexF32 { rssi_dbm, values } => {
out.extend(rssi_dbm.iter().map(|v| *v as u8));
for [i, q] in values {
out.extend_from_slice(&i.to_le_bytes());
out.extend_from_slice(&q.to_le_bytes());
}
}
CsiPayload::Bytes(values) => out.extend_from_slice(values),
}
out.extend_from_slice(&crc32_ieee(&out).to_le_bytes());
Ok(out)
}
pub fn from_bytes(input: &[u8]) -> Result<(Self, usize), CsiParseError> {
if input.len() < QUALCOMM_CSI_HEADER_LEN {
return Err(CsiParseError::InsufficientData {
needed: QUALCOMM_CSI_HEADER_LEN,
got: input.len(),
});
}
let magic = u32_at(input, 0);
if magic != QUALCOMM_CSI_MAGIC {
return Err(CsiParseError::InvalidMagic(magic));
}
if input[4] != QUALCOMM_CSI_VERSION {
return Err(CsiParseError::UnsupportedVersion(input[4]));
}
let report_kind = ReportKind::try_from(input[5])?;
let header_len = u16_at(input, 6) as usize;
if header_len != QUALCOMM_CSI_HEADER_LEN {
return Err(CsiParseError::InvalidHeaderLength(header_len));
}
let frame_len = u32_at(input, 8) as usize;
if frame_len > QUALCOMM_CSI_MAX_FRAME_LEN {
return Err(CsiParseError::FrameTooLarge(frame_len));
}
if frame_len < header_len + QUALCOMM_CSI_CRC_LEN {
return Err(CsiParseError::InvalidFrameLength(frame_len));
}
if input.len() < frame_len {
return Err(CsiParseError::InsufficientData {
needed: frame_len,
got: input.len(),
});
}
let expected_crc = u32_at(input, frame_len - 4);
let actual_crc = crc32_ieee(&input[..frame_len - 4]);
if expected_crc != actual_crc {
return Err(CsiParseError::CrcMismatch {
expected: expected_crc,
actual: actual_crc,
});
}
let chipset = ChipsetProfile::try_from(u16_at(input, 32))?;
let format = ElementFormat::try_from(input[44])?;
let ppdu_type = PpduType::try_from(input[45])?;
let tx_count = input[42];
let rx_count = input[43];
let subcarrier_count = u16_at(input, 46);
let rssi_count = input[48] as usize;
let payload_len = u32_at(input, 64) as usize;
if header_len + payload_len + 4 != frame_len {
return Err(CsiParseError::PayloadLengthMismatch);
}
let payload_bytes = &input[header_len..header_len + payload_len];
let elements = (tx_count as usize)
.checked_mul(rx_count as usize)
.and_then(|n| n.checked_mul(subcarrier_count as usize))
.ok_or(CsiParseError::LengthOverflow)?;
let payload = match format {
ElementFormat::Bytes => CsiPayload::Bytes(payload_bytes.to_vec()),
ElementFormat::ComplexI16 => {
if rssi_count > payload_bytes.len()
|| payload_bytes.len() - rssi_count != elements * 4
{
return Err(CsiParseError::PayloadLengthMismatch);
}
let rssi_dbm = payload_bytes[..rssi_count]
.iter()
.map(|v| *v as i8)
.collect();
let values = payload_bytes[rssi_count..]
.chunks_exact(4)
.map(|b| {
[
i16::from_le_bytes([b[0], b[1]]),
i16::from_le_bytes([b[2], b[3]]),
]
})
.collect();
CsiPayload::ComplexI16 { rssi_dbm, values }
}
ElementFormat::ComplexF32 => {
if rssi_count > payload_bytes.len()
|| payload_bytes.len() - rssi_count != elements * 8
{
return Err(CsiParseError::PayloadLengthMismatch);
}
let rssi_dbm = payload_bytes[..rssi_count]
.iter()
.map(|v| *v as i8)
.collect();
let mut values = Vec::with_capacity(elements);
for b in payload_bytes[rssi_count..].chunks_exact(8) {
let i = f32::from_le_bytes(b[0..4].try_into().unwrap());
let q = f32::from_le_bytes(b[4..8].try_into().unwrap());
if !i.is_finite() || !q.is_finite() {
return Err(CsiParseError::NonFiniteValue);
}
values.push([i, q]);
}
CsiPayload::ComplexF32 { rssi_dbm, values }
}
};
let frame = Self {
report_kind,
sequence: u32_at(input, 12),
timestamp_us: u64_at(input, 16),
device_id: u64_at(input, 24),
chipset,
bandwidth_mhz: u16_at(input, 34),
center_freq_khz: u32_at(input, 36),
flags: CsiFlags(u16_at(input, 40)),
tx_count,
rx_count,
ppdu_type,
subcarrier_count,
noise_floor_dbm: input[49] as i8,
scale: f32_at(input, 52),
subcarrier_spacing_hz: f32_at(input, 56),
calibration_id: u32_at(input, 60),
payload,
};
frame.validate()?;
Ok((frame, frame_len))
}
fn validate(&self) -> Result<(), CsiParseError> {
if !matches!(self.bandwidth_mhz, 20 | 40 | 80 | 160)
|| self.bandwidth_mhz > self.chipset.max_bandwidth_mhz()
{
return Err(CsiParseError::InvalidBandwidth(self.bandwidth_mhz));
}
if self.tx_count == 0
|| self.rx_count == 0
|| self.tx_count > self.chipset.max_chains()
|| self.rx_count > self.chipset.max_chains()
{
return Err(CsiParseError::InvalidDimensions);
}
if !self.scale.is_finite()
|| self.scale <= 0.0
|| !self.subcarrier_spacing_hz.is_finite()
|| self.subcarrier_spacing_hz <= 0.0
{
return Err(CsiParseError::NonFiniteValue);
}
match (&self.report_kind, &self.payload) {
(ReportKind::Csi, CsiPayload::ComplexI16 { rssi_dbm, values }) => {
self.validate_csi(rssi_dbm, values.len())
}
(ReportKind::Csi, CsiPayload::ComplexF32 { rssi_dbm, values }) => {
if !values.iter().flatten().all(|v| v.is_finite()) {
return Err(CsiParseError::NonFiniteValue);
}
self.validate_csi(rssi_dbm, values.len())
}
(ReportKind::Capabilities, CsiPayload::Bytes(v)) if !v.is_empty() => Ok(()),
_ => Err(CsiParseError::PayloadTypeMismatch),
}
}
fn validate_csi(&self, rssi: &[i8], values: usize) -> Result<(), CsiParseError> {
let expected =
self.tx_count as usize * self.rx_count as usize * self.subcarrier_count as usize;
if expected == 0 || expected > QUALCOMM_CSI_MAX_ELEMENTS {
return Err(CsiParseError::InvalidDimensions);
}
if values != expected || rssi.len() != self.rx_count as usize {
return Err(CsiParseError::PayloadLengthMismatch);
}
Ok(())
}
}
#[derive(Debug, Error, PartialEq)]
pub enum CsiParseError {
#[error("insufficient data: needed {needed}, got {got}")]
InsufficientData { needed: usize, got: usize },
#[error("invalid magic {0:#010x}")]
InvalidMagic(u32),
#[error("unsupported version {0}")]
UnsupportedVersion(u8),
#[error("unknown report kind {0}")]
UnknownReportKind(u8),
#[error("unknown chipset profile {0}")]
UnknownChipset(u16),
#[error("unknown element format {0}")]
UnknownElementFormat(u8),
#[error("unknown PPDU type {0}")]
UnknownPpduType(u8),
#[error("invalid header length {0}")]
InvalidHeaderLength(usize),
#[error("invalid frame length {0}")]
InvalidFrameLength(usize),
#[error("frame too large: {0}")]
FrameTooLarge(usize),
#[error("length arithmetic overflow")]
LengthOverflow,
#[error("payload length mismatch")]
PayloadLengthMismatch,
#[error("payload type does not match report kind")]
PayloadTypeMismatch,
#[error("invalid MIMO dimensions")]
InvalidDimensions,
#[error("invalid bandwidth {0} MHz")]
InvalidBandwidth(u16),
#[error("non-finite or non-positive numeric metadata/value")]
NonFiniteValue,
#[error("CRC mismatch: expected {expected:#010x}, actual {actual:#010x}")]
CrcMismatch { expected: u32, actual: u32 },
}
pub mod simulator {
use super::*;
#[derive(Debug, Clone)]
pub struct SimulatorConfig {
pub seed: u64,
pub device_id: u64,
pub chipset: ChipsetProfile,
pub bandwidth_mhz: u16,
pub center_freq_khz: u32,
pub tx_count: u8,
pub rx_count: u8,
pub subcarriers: u16,
pub frame_period_us: u64,
}
impl Default for SimulatorConfig {
fn default() -> Self {
Self {
seed: 0x5143_4143_5349_0001,
device_id: 0x5255_5651_4341_3031,
chipset: ChipsetProfile::Qca9300,
bandwidth_mhz: 40,
center_freq_khz: 5_210_000,
tx_count: 2,
rx_count: 3,
subcarriers: 114,
frame_period_us: 20_000,
}
}
}
pub struct QualcommCsiSimulator {
config: SimulatorConfig,
rng: u64,
sequence: u32,
timestamp_us: u64,
motion_phase: f32,
}
impl QualcommCsiSimulator {
pub fn new(config: SimulatorConfig) -> Result<Self, CsiParseError> {
let s = Self {
rng: config.seed,
config,
sequence: 0,
timestamp_us: 0,
motion_phase: 0.0,
};
s.csi_frame()?.validate()?;
Ok(s)
}
pub fn capabilities_frame(&self) -> CsiFrame {
self.base(
ReportKind::Capabilities,
CsiPayload::Bytes(vec![
1,
1,
self.config.chipset.max_chains(),
2,
1,
0b0000_1111,
3,
2,
(self.config.subcarriers & 255) as u8,
(self.config.subcarriers >> 8) as u8,
]),
)
}
pub fn next_frame(&mut self) -> CsiFrame {
let frame = self.csi_frame().expect("validated simulator config");
self.sequence = self.sequence.wrapping_add(1);
self.timestamp_us = self.timestamp_us.wrapping_add(self.config.frame_period_us);
self.motion_phase += 0.037;
frame
}
fn csi_frame(&self) -> Result<CsiFrame, CsiParseError> {
let mut rng = self.rng ^ self.sequence as u64;
let count = self.config.tx_count as usize
* self.config.rx_count as usize
* self.config.subcarriers as usize;
let values = (0..count)
.map(|idx| {
rng ^= rng << 13;
rng ^= rng >> 7;
rng ^= rng << 17;
let noise = ((rng >> 48) as i16 % 24) as f32;
let sc = (idx % self.config.subcarriers as usize) as f32;
let chain = (idx / self.config.subcarriers as usize) as f32;
let phase = sc * 0.031 + chain * 0.23 + self.motion_phase;
[
((phase.cos() * 1800.0) + noise) as i16,
((phase.sin() * 1800.0) - noise) as i16,
]
})
.collect();
Ok(self.base(
ReportKind::Csi,
CsiPayload::ComplexI16 {
rssi_dbm: (0..self.config.rx_count)
.map(|i| -42 - i as i8 * 2)
.collect(),
values,
},
))
}
fn base(&self, kind: ReportKind, payload: CsiPayload) -> CsiFrame {
CsiFrame {
report_kind: kind,
sequence: self.sequence,
timestamp_us: self.timestamp_us,
device_id: self.config.device_id,
chipset: self.config.chipset,
bandwidth_mhz: self.config.bandwidth_mhz,
center_freq_khz: self.config.center_freq_khz,
flags: CsiFlags(CsiFlags::CALIBRATED | CsiFlags::SYNTHETIC),
tx_count: self.config.tx_count,
rx_count: self.config.rx_count,
ppdu_type: PpduType::HeSu,
subcarrier_count: self.config.subcarriers,
noise_floor_dbm: -95,
scale: 1.0 / 2048.0,
subcarrier_spacing_hz: 312_500.0,
calibration_id: 1,
payload,
}
}
}
}
fn u16_at(b: &[u8], o: usize) -> u16 {
u16::from_le_bytes([b[o], b[o + 1]])
}
fn u32_at(b: &[u8], o: usize) -> u32 {
u32::from_le_bytes(b[o..o + 4].try_into().unwrap())
}
fn u64_at(b: &[u8], o: usize) -> u64 {
u64::from_le_bytes(b[o..o + 8].try_into().unwrap())
}
fn f32_at(b: &[u8], o: usize) -> f32 {
f32::from_le_bytes(b[o..o + 4].try_into().unwrap())
}
fn crc32_ieee(data: &[u8]) -> u32 {
let mut crc = 0xffff_ffffu32;
for &byte in data {
crc ^= byte as u32;
for _ in 0..8 {
crc = (crc >> 1) ^ ((0u32.wrapping_sub(crc & 1)) & 0xedb8_8320);
}
}
!crc
}
#[cfg(test)]
mod tests {
use super::*;
use simulator::*;
#[test]
fn simulator_round_trip_is_deterministic() {
let cfg = SimulatorConfig::default();
let mut a = QualcommCsiSimulator::new(cfg.clone()).unwrap();
let mut b = QualcommCsiSimulator::new(cfg).unwrap();
let wa = a.next_frame().to_bytes().unwrap();
assert_eq!(wa, b.next_frame().to_bytes().unwrap());
let (decoded, n) = CsiFrame::from_bytes(&wa).unwrap();
assert_eq!(n, wa.len());
assert!(decoded.flags.contains(CsiFlags::SYNTHETIC));
assert_eq!(decoded.payload.len(), 2 * 3 * 114);
}
#[test]
fn capabilities_round_trip() {
let s = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let f = s.capabilities_frame();
let w = f.to_bytes().unwrap();
assert_eq!(CsiFrame::from_bytes(&w).unwrap().0, f);
}
#[test]
fn crc_corruption_is_rejected() {
let mut s = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let mut w = s.next_frame().to_bytes().unwrap();
w[80] ^= 1;
assert!(matches!(
CsiFrame::from_bytes(&w),
Err(CsiParseError::CrcMismatch { .. })
));
}
#[test]
fn truncation_is_rejected() {
let mut s = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let w = s.next_frame().to_bytes().unwrap();
assert!(matches!(
CsiFrame::from_bytes(&w[..w.len() - 1]),
Err(CsiParseError::InsufficientData { .. })
));
}
#[test]
fn invalid_dimensions_are_rejected() {
let cfg = SimulatorConfig {
rx_count: 4,
chipset: ChipsetProfile::Qca9300,
..Default::default()
};
assert!(matches!(
QualcommCsiSimulator::new(cfg),
Err(CsiParseError::InvalidDimensions)
));
}
#[test]
fn non_finite_float_is_rejected() {
let mut s = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let mut f = s.next_frame();
f.payload = CsiPayload::ComplexF32 {
rssi_dbm: vec![-40, -42, -44],
values: vec![[f32::NAN, 0.0]; 2 * 3 * 114],
};
assert_eq!(f.to_bytes().unwrap_err(), CsiParseError::NonFiniteValue);
}
#[test]
fn parser_never_panics_on_prefixes() {
let mut s = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let w = s.next_frame().to_bytes().unwrap();
for end in 0..w.len() {
let _ = CsiFrame::from_bytes(&w[..end]);
}
}
#[test]
fn qca9300_rejects_wifi6_bandwidths() {
let cfg = SimulatorConfig {
bandwidth_mhz: 80,
..Default::default()
};
assert!(matches!(
QualcommCsiSimulator::new(cfg),
Err(CsiParseError::InvalidBandwidth(80))
));
}
}
@@ -0,0 +1,207 @@
//! ADR-270 vendor RF provider contract.
//!
//! This model prevents RSSI, cloud occupancy, or network inventory from being
//! represented as complex CSI. Vendor adapters may emit only declared capabilities.
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use thiserror::Error;
pub const MAX_VENDOR_METRICS: usize = 64;
pub const MAX_VENDOR_KEY_LEN: usize = 64;
pub const MAX_VENDOR_TEXT_LEN: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum VendorId {
OriginAi,
Plume,
Mist,
Netgear,
ElectricImp,
RfSolutions,
Linksys,
Luma,
GoogleNest,
Wifigarden,
}
impl VendorId {
pub const ALL: [Self; 10] = [
Self::OriginAi,
Self::Plume,
Self::Mist,
Self::Netgear,
Self::ElectricImp,
Self::RfSolutions,
Self::Linksys,
Self::Luma,
Self::GoogleNest,
Self::Wifigarden,
];
pub fn as_str(self) -> &'static str {
match self {
Self::OriginAi => "origin_ai",
Self::Plume => "plume",
Self::Mist => "mist",
Self::Netgear => "netgear",
Self::ElectricImp => "electric_imp",
Self::RfSolutions => "rf_solutions",
Self::Linksys => "linksys",
Self::Luma => "luma",
Self::GoogleNest => "google_nest",
Self::Wifigarden => "wifigarden",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RfCapability {
ComplexCsi,
DerivedSensing,
RfTelemetry,
NetworkOnly,
Unsupported,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ProviderAvailability {
Available,
CredentialsRequired,
ContractRequired,
Experimental,
Unsupported,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ProviderDescriptor {
pub vendor: VendorId,
pub capabilities: Vec<RfCapability>,
pub availability: ProviderAvailability,
pub hardware_validated: bool,
pub reason: String,
}
impl ProviderDescriptor {
pub fn validate(&self) -> Result<(), VendorEventError> {
if self.capabilities.is_empty()
|| self.reason.is_empty()
|| self.reason.len() > MAX_VENDOR_TEXT_LEN
{
return Err(VendorEventError::InvalidDescriptor);
}
if self.hardware_validated && self.availability != ProviderAvailability::Available {
return Err(VendorEventError::InvalidDescriptor);
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct VendorRfEvent {
pub vendor: VendorId,
pub capability: RfCapability,
pub sequence: u64,
pub timestamp_us: u64,
pub source_id: String,
pub synthetic: bool,
pub metrics: BTreeMap<String, f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub label: Option<String>,
}
impl VendorRfEvent {
pub fn validate(&self, descriptor: &ProviderDescriptor) -> Result<(), VendorEventError> {
descriptor.validate()?;
if self.vendor != descriptor.vendor || !descriptor.capabilities.contains(&self.capability) {
return Err(VendorEventError::CapabilityMismatch);
}
if matches!(
self.capability,
RfCapability::ComplexCsi | RfCapability::Unsupported
) {
return Err(VendorEventError::InvalidEventCapability);
}
if self.source_id.is_empty()
|| self.source_id.len() > MAX_VENDOR_TEXT_LEN
|| self.metrics.is_empty()
|| self.metrics.len() > MAX_VENDOR_METRICS
|| self
.metrics
.iter()
.any(|(k, v)| k.is_empty() || k.len() > MAX_VENDOR_KEY_LEN || !v.is_finite())
|| self
.label
.as_ref()
.is_some_and(|v| v.len() > MAX_VENDOR_TEXT_LEN)
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
}
pub trait VendorRfProvider: Send + Sync {
fn descriptor(&self) -> ProviderDescriptor;
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError>;
}
#[derive(Debug, Error, PartialEq, Eq)]
pub enum VendorEventError {
#[error("invalid provider descriptor")]
InvalidDescriptor,
#[error("event capability does not match provider")]
CapabilityMismatch,
#[error("complex CSI and unsupported states cannot be represented as scalar vendor events")]
InvalidEventCapability,
#[error("invalid or unbounded vendor payload")]
InvalidPayload,
#[error("malformed provider payload: {0}")]
MalformedPayload(String),
#[error("provider credentials are required")]
CredentialsRequired,
#[error("commercial contract or SDK access is required")]
ContractRequired,
#[error("provider has no supported sensing interface")]
Unsupported,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn every_vendor_has_a_stable_identifier() {
let names: std::collections::BTreeSet<_> =
VendorId::ALL.iter().map(|v| v.as_str()).collect();
assert_eq!(names.len(), VendorId::ALL.len());
}
#[test]
fn scalar_contract_rejects_csi_masquerading() {
let descriptor = ProviderDescriptor {
vendor: VendorId::Plume,
capabilities: vec![RfCapability::RfTelemetry],
availability: ProviderAvailability::CredentialsRequired,
hardware_validated: false,
reason: "OpenSync telemetry".into(),
};
let event = VendorRfEvent {
vendor: VendorId::Plume,
capability: RfCapability::ComplexCsi,
sequence: 1,
timestamp_us: 1,
source_id: "pod-1".into(),
synthetic: true,
metrics: BTreeMap::from([("rssi_dbm".into(), -42.0)]),
label: None,
};
assert_eq!(
event.validate(&descriptor),
Err(VendorEventError::CapabilityMismatch)
);
}
}
@@ -34,3 +34,11 @@ pub mod sparse_inference;
#[allow(dead_code)]
pub mod trainer;
pub mod vital_signs;
/// ADR-270 Mist and NETGEAR telemetry providers.
pub mod vendor_mist_netgear;
/// ADR-270 Origin AI and Plume/OpenSync providers.
pub mod vendor_origin_plume;
/// ADR-270 scalar, network-only, and fail-closed vendor providers.
pub mod vendor_remaining;
/// ADR-270 provider registry and canonical event helpers.
pub mod vendor_rf;
@@ -18,6 +18,7 @@ mod field_localize;
mod model_format;
mod multistatic_bridge;
mod mediatek_csi;
mod qualcomm_csi;
mod realtek_radar;
pub mod pose;
mod rvf_container;
@@ -32,13 +33,14 @@ use wifi_densepose_sensing_server::{
};
use ruvector_mincut::{DynamicMinCut, MinCutBuilder};
use std::collections::{HashMap, VecDeque};
use std::collections::{BTreeMap, HashMap, VecDeque};
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use axum::{
body::Bytes,
extract::{
ws::{Message, WebSocket, WebSocketUpgrade},
Path, Query, State,
@@ -1038,6 +1040,12 @@ struct AppStateInner {
latest_mediatek_csi: Option<mediatek_csi::MediatekCsiSnapshot>,
/// Instant of the last validated MediaTek CSI UDP frame.
last_mediatek_frame: Option<std::time::Instant>,
/// Latest validated Qualcomm CSI summary; raw matrices are not retained here.
latest_qualcomm_csi: Option<qualcomm_csi::QualcommCsiSnapshot>,
/// Instant of the last validated Qualcomm CSI UDP frame.
last_qualcomm_frame: Option<std::time::Instant>,
/// Latest bounded ADR-270 event per vendor. Complex CSI uses dedicated transports.
latest_vendor_rf: BTreeMap<String, wifi_densepose_sensing_server::vendor_rf::VendorEventSnapshot>,
tx: broadcast::Sender<String>,
// ADR-099 D2/D3/D4: real-time CSI introspection tap. Per-frame state +
// a parallel broadcast topic (`/ws/introspection`) running alongside
@@ -1223,6 +1231,13 @@ impl AppStateInner {
}
}
}
if self.source.starts_with("qualcomm") {
if let Some(last) = self.last_qualcomm_frame {
if last.elapsed() > ESP32_OFFLINE_TIMEOUT {
return format!("{}:offline", self.source);
}
}
}
self.source.clone()
}
}
@@ -3391,6 +3406,97 @@ async fn latest_mediatek_csi(State(state): State<SharedState>) -> Json<serde_jso
}
}
async fn latest_qualcomm_csi(State(state): State<SharedState>) -> Json<serde_json::Value> {
let s = state.read().await;
match &s.latest_qualcomm_csi {
Some(snapshot) => Json(serde_json::to_value(snapshot).unwrap_or_default()),
None => Json(serde_json::json!({"status": "no Qualcomm CSI data yet"})),
}
}
async fn vendor_descriptors() -> Json<serde_json::Value> {
Json(
serde_json::to_value(wifi_densepose_sensing_server::vendor_rf::descriptors())
.unwrap_or_default(),
)
}
async fn latest_vendor_events(State(state): State<SharedState>) -> Json<serde_json::Value> {
let state = state.read().await;
Json(serde_json::to_value(&state.latest_vendor_rf).unwrap_or_default())
}
async fn latest_vendor_event(
State(state): State<SharedState>,
Path(vendor): Path<String>,
) -> impl IntoResponse {
let Some(vendor_id) = wifi_densepose_sensing_server::vendor_rf::vendor_from_str(&vendor) else {
return (StatusCode::NOT_FOUND, Json(serde_json::json!({"error": "unknown vendor"})));
};
let state = state.read().await;
let canonical_vendor = vendor_id.as_str();
match state.latest_vendor_rf.get(canonical_vendor) {
Some(snapshot) => (
StatusCode::OK,
Json(serde_json::to_value(snapshot).unwrap_or_default()),
),
None => (
StatusCode::OK,
Json(serde_json::json!({
"status": "no vendor RF data yet",
"vendor": canonical_vendor
})),
),
}
}
async fn ingest_vendor_events(
State(state): State<SharedState>,
Path(vendor): Path<String>,
payload: Bytes,
) -> impl IntoResponse {
let Some(vendor_id) = wifi_densepose_sensing_server::vendor_rf::vendor_from_str(&vendor) else {
return (StatusCode::NOT_FOUND, Json(serde_json::json!({"error": "unknown vendor"})));
};
let events = match wifi_densepose_sensing_server::vendor_rf::decode_provider(vendor_id, &payload) {
Ok(events) => events,
Err(error) => {
let status = match error {
wifi_densepose_hardware::vendor_rf::VendorEventError::Unsupported => StatusCode::NOT_IMPLEMENTED,
wifi_densepose_hardware::vendor_rf::VendorEventError::ContractRequired
| wifi_densepose_hardware::vendor_rf::VendorEventError::CredentialsRequired => StatusCode::FORBIDDEN,
_ => StatusCode::BAD_REQUEST,
};
return (status, Json(serde_json::json!({"error": error.to_string(), "vendor": vendor})));
}
};
let mut accepted = 0usize;
let mut state = state.write().await;
let canonical_vendor = vendor_id.as_str().to_string();
for event in events {
match wifi_densepose_sensing_server::vendor_rf::VendorEventSnapshot::from_event(event) {
Ok(snapshot) => {
let json = serde_json::to_string(&snapshot).ok();
state.source = snapshot.source.clone();
state
.latest_vendor_rf
.insert(canonical_vendor.clone(), snapshot);
if let Some(json) = json {
let _ = state.tx.send(json);
}
accepted += 1;
}
Err(error) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({"error": error.to_string(), "vendor": canonical_vendor})),
)
}
}
}
(StatusCode::ACCEPTED, Json(serde_json::json!({"vendor": vendor, "accepted": accepted})))
}
/// Generate WiFi-derived pose keypoints from sensing data.
///
/// Keypoint positions are modulated by real signal features rather than a pure
@@ -5485,7 +5591,7 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) {
let addr = format!("0.0.0.0:{udp_port}");
let socket = match UdpSocket::bind(&addr).await {
Ok(s) => {
info!("UDP listening on {addr} for ESP32 CSI, MediaTek CSI, and RTL8720F radar frames");
info!("UDP listening on {addr} for ESP32, MediaTek, Qualcomm CSI, and RTL8720F radar frames");
s
}
Err(e) => {
@@ -5498,6 +5604,44 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) {
loop {
match socket.recv_from(&mut buf).await {
Ok((len, src)) => {
if len > 0 && buf[0] == b'{' {
match serde_json::from_slice::<wifi_densepose_hardware::vendor_rf::VendorRfEvent>(&buf[..len])
.map_err(|error| error.to_string())
.and_then(|event| wifi_densepose_sensing_server::vendor_rf::VendorEventSnapshot::from_event(event).map_err(|error| error.to_string()))
{
Ok(snapshot) if snapshot.event.synthetic => {
debug!("Vendor RF event from {src}: vendor={} capability={:?} seq={}", snapshot.event.vendor.as_str(), snapshot.event.capability, snapshot.event.sequence);
let json = serde_json::to_string(&snapshot).ok();
let mut state = state.write().await;
state.source = snapshot.source.clone();
state.latest_vendor_rf.insert(snapshot.event.vendor.as_str().to_string(), snapshot);
if let Some(json) = json { let _ = state.tx.send(json); }
}
Ok(_) => warn!("Rejected non-synthetic canonical vendor event from {src}; live payloads must use the provider decoder HTTP route"),
Err(error) => warn!("Rejected ADR-270 vendor event from {src}: {error}"),
}
continue;
}
if len >= 4
&& u32::from_le_bytes(buf[..4].try_into().expect("four-byte slice"))
== wifi_densepose_hardware::qualcomm_csi::QUALCOMM_CSI_MAGIC
{
match wifi_densepose_hardware::qualcomm_csi::CsiFrame::from_bytes(&buf[..len]) {
Ok((frame, consumed)) if consumed == len => {
let snapshot = qualcomm_csi::QualcommCsiSnapshot::from_frame(&frame);
debug!("Qualcomm CSI from {src}: profile={} seq={} dimensions={}x{}x{}", snapshot.chipset, snapshot.sequence, snapshot.tx_count, snapshot.rx_count, snapshot.subcarrier_count);
let json = serde_json::to_string(&snapshot).ok();
let mut s = state.write().await;
s.source = snapshot.source.to_string();
s.last_qualcomm_frame = Some(std::time::Instant::now());
s.latest_qualcomm_csi = Some(snapshot);
if let Some(json) = json { let _ = s.tx.send(json); }
}
Ok((_, consumed)) => warn!("Qualcomm CSI datagram from {src} has trailing bytes: consumed={consumed} received={len}"),
Err(error) => warn!("Rejected Qualcomm CSI datagram from {src}: {error}"),
}
continue;
}
if len >= 4
&& u32::from_le_bytes(buf[..4].try_into().expect("four-byte slice"))
== wifi_densepose_hardware::mediatek_csi::MEDIATEK_CSI_MAGIC
@@ -7638,6 +7782,9 @@ async fn main() {
last_realtek_frame: None,
latest_mediatek_csi: None,
last_mediatek_frame: None,
latest_qualcomm_csi: None,
last_qualcomm_frame: None,
latest_vendor_rf: BTreeMap::new(),
tx,
intro: wifi_densepose_sensing_server::introspection::IntrospectionState::new(),
intro_tx,
@@ -7856,6 +8003,11 @@ async fn main() {
.route("/api/v1/sensing/latest", get(latest))
.route("/api/v1/radar/latest", get(latest_realtek_radar))
.route("/api/v1/csi/mediatek/latest", get(latest_mediatek_csi))
.route("/api/v1/csi/qualcomm/latest", get(latest_qualcomm_csi))
.route("/api/v1/rf/vendors", get(vendor_descriptors))
.route("/api/v1/rf/vendors/latest", get(latest_vendor_events))
.route("/api/v1/rf/vendors/:vendor/latest", get(latest_vendor_event))
.route("/api/v1/rf/vendors/:vendor/events", post(ingest_vendor_events))
// Per-node health endpoint
.route("/api/v1/nodes", get(nodes_endpoint))
// ADR-110 iter 29 — per-node mesh sync state for HTTP clients.
@@ -0,0 +1,127 @@
//! Bounded summaries for ADR-269 Qualcomm MIMO CSI frames.
use serde::Serialize;
use wifi_densepose_hardware::qualcomm_csi::{CsiFlags, CsiFrame, CsiPayload, ReportKind};
#[derive(Debug, Clone, PartialEq, Serialize)]
pub(crate) struct QualcommCsiSnapshot {
pub event_type: &'static str,
pub source: &'static str,
pub report_kind: &'static str,
pub sequence: u32,
pub timestamp_us: u64,
pub device_id: String,
pub chipset: &'static str,
pub center_freq_khz: u32,
pub bandwidth_mhz: u16,
pub tx_count: u8,
pub rx_count: u8,
pub subcarrier_count: u16,
pub element_count: usize,
pub ppdu_type: String,
pub rssi_dbm: Vec<i8>,
pub noise_floor_dbm: i8,
pub calibrated: bool,
pub synthetic: bool,
pub saturated: bool,
pub time_synchronized: bool,
pub dropped_predecessor: bool,
pub calibration_id: u32,
pub subcarrier_spacing_hz: f32,
pub mean_amplitude: Option<f32>,
pub peak_amplitude: Option<f32>,
}
impl QualcommCsiSnapshot {
pub(crate) fn from_frame(frame: &CsiFrame) -> Self {
let synthetic = frame.flags.contains(CsiFlags::SYNTHETIC);
let (mean_amplitude, peak_amplitude) = amplitude_summary(frame);
Self {
event_type: "qualcomm_csi",
source: if synthetic {
"qualcomm:simulated"
} else {
"qualcomm"
},
report_kind: match frame.report_kind {
ReportKind::Csi => "csi",
ReportKind::Capabilities => "capabilities",
},
sequence: frame.sequence,
timestamp_us: frame.timestamp_us,
device_id: format!("{:016x}", frame.device_id),
chipset: frame.chipset.name(),
center_freq_khz: frame.center_freq_khz,
bandwidth_mhz: frame.bandwidth_mhz,
tx_count: frame.tx_count,
rx_count: frame.rx_count,
subcarrier_count: frame.subcarrier_count,
element_count: frame.payload.len(),
ppdu_type: format!("{:?}", frame.ppdu_type).to_ascii_lowercase(),
rssi_dbm: frame.payload.rssi_dbm().to_vec(),
noise_floor_dbm: frame.noise_floor_dbm,
calibrated: frame.flags.contains(CsiFlags::CALIBRATED),
synthetic,
saturated: frame.flags.contains(CsiFlags::SATURATED),
time_synchronized: frame.flags.contains(CsiFlags::TIME_SYNCHRONIZED),
dropped_predecessor: frame.flags.contains(CsiFlags::DROPPED_PREDECESSOR),
calibration_id: frame.calibration_id,
subcarrier_spacing_hz: frame.subcarrier_spacing_hz,
mean_amplitude,
peak_amplitude,
}
}
}
fn amplitude_summary(frame: &CsiFrame) -> (Option<f32>, Option<f32>) {
let amplitudes: Vec<f32> = match &frame.payload {
CsiPayload::ComplexI16 { values, .. } => values
.iter()
.map(|[i, q]| (*i as f32).hypot(*q as f32) * frame.scale)
.collect(),
CsiPayload::ComplexF32 { values, .. } => values
.iter()
.map(|[i, q]| i.hypot(*q) * frame.scale)
.collect(),
CsiPayload::Bytes(_) => return (None, None),
};
if amplitudes.is_empty() {
return (None, None);
}
let mean = amplitudes.iter().sum::<f32>() / amplitudes.len() as f32;
let peak = amplitudes.into_iter().max_by(f32::total_cmp);
(Some(mean), peak)
}
#[cfg(test)]
mod tests {
use super::*;
use wifi_densepose_hardware::qualcomm_csi::simulator::{QualcommCsiSimulator, SimulatorConfig};
#[test]
fn simulator_summary_preserves_dimensions_and_provenance() {
let mut sim = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let snapshot = QualcommCsiSnapshot::from_frame(&sim.next_frame());
assert_eq!(snapshot.source, "qualcomm:simulated");
assert_eq!(
(
snapshot.tx_count,
snapshot.rx_count,
snapshot.subcarrier_count
),
(2, 3, 114)
);
assert_eq!(snapshot.element_count, 684);
assert!(snapshot.mean_amplitude.unwrap() > 0.0);
assert!(snapshot.peak_amplitude.unwrap() >= snapshot.mean_amplitude.unwrap());
}
#[test]
fn capability_summary_does_not_invent_signal_statistics() {
let sim = QualcommCsiSimulator::new(SimulatorConfig::default()).unwrap();
let snapshot = QualcommCsiSnapshot::from_frame(&sim.capabilities_frame());
assert_eq!(snapshot.report_kind, "capabilities");
assert_eq!(snapshot.mean_amplitude, None);
assert!(snapshot.rssi_dbm.is_empty());
}
}
@@ -0,0 +1,712 @@
//! ADR-270 adapters for Mist/Juniper and NETGEAR Insight.
//!
//! These cloud APIs expose client RF telemetry and location/network context.
//! They do not expose complex channel state information (CSI), and this module
//! deliberately cannot construct a `ComplexCsi` event.
use serde::Deserialize;
use serde_json::{Map, Value};
use std::collections::BTreeMap;
use std::fmt;
use wifi_densepose_hardware::vendor_rf::{
ProviderAvailability, ProviderDescriptor, RfCapability, VendorEventError, VendorId,
VendorRfEvent, VendorRfProvider,
};
pub const MAX_VENDOR_PAYLOAD_BYTES: usize = 1024 * 1024;
pub const MAX_EVENTS_PER_PAGE: usize = 1_000;
pub const MAX_CURSOR_BYTES: usize = 512;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MistRegion {
Global,
Europe,
AsiaPacific,
Australia,
}
impl MistRegion {
pub const fn base_url(self) -> &'static str {
match self {
Self::Global => "https://api.mist.com",
Self::Europe => "https://api.eu.mist.com",
Self::AsiaPacific => "https://api.ac2.mist.com",
Self::Australia => "https://api.gc1.mist.com",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NetgearRegion {
NorthAmerica,
Europe,
Australia,
}
impl NetgearRegion {
pub const fn base_url(self) -> &'static str {
match self {
Self::NorthAmerica => "https://insight.netgear.com",
Self::Europe => "https://eu.insight.netgear.com",
Self::Australia => "https://au.insight.netgear.com",
}
}
}
/// An authentication value whose `Debug` output is always redacted.
#[derive(Clone, PartialEq, Eq)]
pub struct SecretToken(String);
impl SecretToken {
pub fn new(value: impl Into<String>) -> Result<Self, VendorEventError> {
let value = value.into();
if value.is_empty() || value.len() > 4_096 || value.chars().any(char::is_control) {
return Err(VendorEventError::InvalidPayload);
}
Ok(Self(value))
}
/// Intended only for constructing an HTTP authorization header.
pub fn expose_for_header(&self) -> &str {
&self.0
}
}
impl fmt::Debug for SecretToken {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("SecretToken([REDACTED])")
}
}
#[derive(Clone, PartialEq, Eq)]
pub struct VendorRequestConfig {
pub base_url: &'static str,
pub path: String,
pub cursor: Option<String>,
pub token: SecretToken,
}
impl VendorRequestConfig {
pub fn validate(&self) -> Result<(), VendorEventError> {
if !self.base_url.starts_with("https://")
|| !self.path.starts_with('/')
|| self.path.contains("..")
|| self.path.chars().any(char::is_control)
|| self.cursor.as_deref().is_some_and(|cursor| {
cursor.is_empty()
|| cursor.len() > MAX_CURSOR_BYTES
|| cursor
.chars()
.any(|c| c.is_control() || c == '&' || c == '?' || c == '#')
})
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
}
impl fmt::Debug for VendorRequestConfig {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("VendorRequestConfig")
.field("base_url", &self.base_url)
.field("path", &self.path)
.field("cursor", &self.cursor.as_ref().map(|_| "[PRESENT]"))
.field("token", &"[REDACTED]")
.finish()
}
}
pub fn mist_request(
region: MistRegion,
site_id: &str,
cursor: Option<String>,
token: SecretToken,
) -> Result<VendorRequestConfig, VendorEventError> {
validate_identifier(site_id)?;
let request = VendorRequestConfig {
base_url: region.base_url(),
path: format!("/api/v1/sites/{site_id}/stats/clients"),
cursor,
token,
};
request.validate()?;
Ok(request)
}
pub fn netgear_request(
region: NetgearRegion,
location_id: &str,
cursor: Option<String>,
token: SecretToken,
) -> Result<VendorRequestConfig, VendorEventError> {
validate_identifier(location_id)?;
let request = VendorRequestConfig {
base_url: region.base_url(),
path: format!("/api/v1/locations/{location_id}/clients"),
cursor,
token,
};
request.validate()?;
Ok(request)
}
fn validate_identifier(value: &str) -> Result<(), VendorEventError> {
if value.is_empty()
|| value.len() > 128
|| !value
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_'))
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
#[derive(Debug, Clone, PartialEq)]
pub struct DecodedVendorPage {
pub events: Vec<VendorRfEvent>,
pub next_cursor: Option<String>,
}
#[derive(Debug, Default, Clone, Copy)]
pub struct MistProvider;
impl MistProvider {
pub fn decode_page(&self, payload: &[u8]) -> Result<DecodedVendorPage, VendorEventError> {
decode_page(
payload,
VendorId::Mist,
mist_descriptor(),
parse_mist_record,
)
}
}
impl VendorRfProvider for MistProvider {
fn descriptor(&self) -> ProviderDescriptor {
mist_descriptor()
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
Ok(self.decode_page(payload)?.events)
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NetgearInsightProvider;
impl NetgearInsightProvider {
pub fn decode_page(&self, payload: &[u8]) -> Result<DecodedVendorPage, VendorEventError> {
decode_page(
payload,
VendorId::Netgear,
netgear_descriptor(),
parse_netgear_record,
)
}
}
impl VendorRfProvider for NetgearInsightProvider {
fn descriptor(&self) -> ProviderDescriptor {
netgear_descriptor()
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
Ok(self.decode_page(payload)?.events)
}
}
pub fn mist_descriptor() -> ProviderDescriptor {
ProviderDescriptor {
vendor: VendorId::Mist,
capabilities: vec![RfCapability::RfTelemetry, RfCapability::NetworkOnly],
availability: ProviderAvailability::CredentialsRequired,
hardware_validated: false,
reason: "Mist cloud client RF telemetry and location context; never CSI".into(),
}
}
pub fn netgear_descriptor() -> ProviderDescriptor {
ProviderDescriptor {
vendor: VendorId::Netgear,
capabilities: vec![RfCapability::RfTelemetry, RfCapability::NetworkOnly],
availability: ProviderAvailability::CredentialsRequired,
hardware_validated: false,
reason: "NETGEAR Insight client RF and network telemetry; never CSI".into(),
}
}
type RecordParser = fn(&Map<String, Value>, usize) -> Result<VendorRfEvent, VendorEventError>;
fn decode_page(
payload: &[u8],
vendor: VendorId,
descriptor: ProviderDescriptor,
parser: RecordParser,
) -> Result<DecodedVendorPage, VendorEventError> {
if payload.is_empty() || payload.len() > MAX_VENDOR_PAYLOAD_BYTES {
return Err(VendorEventError::InvalidPayload);
}
let root: Value = serde_json::from_slice(payload)
.map_err(|e| VendorEventError::MalformedPayload(e.to_string()))?;
let (records, next_cursor) = extract_records_and_cursor(&root)?;
if records.is_empty() || records.len() > MAX_EVENTS_PER_PAGE {
return Err(VendorEventError::InvalidPayload);
}
let mut events = Vec::with_capacity(records.len());
for (index, value) in records.iter().enumerate() {
let object = value.as_object().ok_or(VendorEventError::InvalidPayload)?;
let event = parser(object, index)?;
if event.vendor != vendor || event.capability != RfCapability::RfTelemetry {
return Err(VendorEventError::CapabilityMismatch);
}
event.validate(&descriptor)?;
events.push(event);
}
Ok(DecodedVendorPage {
events,
next_cursor,
})
}
fn extract_records_and_cursor(
root: &Value,
) -> Result<(&[Value], Option<String>), VendorEventError> {
if let Some(records) = root.as_array() {
return Ok((records, None));
}
let object = root.as_object().ok_or(VendorEventError::InvalidPayload)?;
let records = ["results", "data", "clients", "events", "items"]
.iter()
.find_map(|key| object.get(*key).and_then(Value::as_array))
.ok_or(VendorEventError::InvalidPayload)?;
let cursor_value = object
.get("next_cursor")
.or_else(|| object.get("nextPageToken"))
.or_else(|| object.get("next_page_token"))
.or_else(|| {
object
.get("pagination")
.and_then(Value::as_object)
.and_then(|p| p.get("next").or_else(|| p.get("cursor")))
});
let next_cursor = match cursor_value {
None | Some(Value::Null) => None,
Some(Value::String(value)) if valid_cursor(value) => Some(value.clone()),
Some(_) => return Err(VendorEventError::InvalidPayload),
};
Ok((records, next_cursor))
}
fn valid_cursor(value: &str) -> bool {
!value.is_empty()
&& value.len() <= MAX_CURSOR_BYTES
&& !value
.chars()
.any(|c| c.is_control() || c == '&' || c == '?' || c == '#')
}
#[derive(Debug, Deserialize)]
struct MistRecord {
#[serde(alias = "client_id", alias = "mac")]
id: String,
#[serde(default, alias = "last_seen", alias = "lastSeen")]
timestamp: Option<Value>,
#[serde(default, alias = "rssi_dbm")]
rssi: Option<f64>,
#[serde(default)]
snr: Option<f64>,
#[serde(default)]
channel: Option<f64>,
#[serde(default)]
x: Option<f64>,
#[serde(default)]
y: Option<f64>,
#[serde(default)]
site_id: Option<String>,
#[serde(default)]
ap_id: Option<String>,
#[serde(default, alias = "event_type", alias = "type")]
event: Option<String>,
}
fn parse_mist_record(
object: &Map<String, Value>,
index: usize,
) -> Result<VendorRfEvent, VendorEventError> {
let record: MistRecord = serde_json::from_value(Value::Object(object.clone()))
.map_err(|e| VendorEventError::MalformedPayload(e.to_string()))?;
validate_source(&record.id)?;
validate_optional_text(record.site_id.as_deref())?;
validate_optional_text(record.ap_id.as_deref())?;
validate_optional_text(record.event.as_deref())?;
let timestamp_us = parse_timestamp(record.timestamp.as_ref())?;
let mut metrics = BTreeMap::new();
push_metric(&mut metrics, "rssi_dbm", record.rssi, -150.0, 20.0)?;
push_metric(&mut metrics, "snr_db", record.snr, -50.0, 100.0)?;
push_metric(&mut metrics, "channel", record.channel, 1.0, 7_000.0)?;
push_metric(&mut metrics, "x_m", record.x, -1_000_000.0, 1_000_000.0)?;
push_metric(&mut metrics, "y_m", record.y, -1_000_000.0, 1_000_000.0)?;
if metrics.is_empty() {
return Err(VendorEventError::InvalidPayload);
}
Ok(VendorRfEvent {
vendor: VendorId::Mist,
capability: RfCapability::RfTelemetry,
sequence: stable_sequence(&record.id, timestamp_us, index),
timestamp_us,
source_id: record.id,
synthetic: false,
metrics,
label: context_label(&[record.site_id, record.ap_id, record.event])?,
})
}
#[derive(Debug, Deserialize)]
struct NetgearRecord {
#[serde(
alias = "clientId",
alias = "mac",
alias = "macAddress",
alias = "deviceId"
)]
id: String,
#[serde(default, alias = "lastSeen", alias = "observedAt")]
timestamp: Option<Value>,
#[serde(default, alias = "signalStrength", alias = "rssi_dbm")]
rssi: Option<f64>,
#[serde(default, alias = "signalToNoiseRatio")]
snr: Option<f64>,
#[serde(default)]
channel: Option<f64>,
#[serde(default, alias = "txRateMbps", alias = "tx_rate")]
tx_rate_mbps: Option<f64>,
#[serde(default, alias = "rxRateMbps", alias = "rx_rate")]
rx_rate_mbps: Option<f64>,
#[serde(default, alias = "locationId")]
location_id: Option<String>,
#[serde(default, alias = "accessPointId", alias = "apId")]
ap_id: Option<String>,
#[serde(default, alias = "ssidName")]
ssid: Option<String>,
}
fn parse_netgear_record(
object: &Map<String, Value>,
index: usize,
) -> Result<VendorRfEvent, VendorEventError> {
let record: NetgearRecord = serde_json::from_value(Value::Object(object.clone()))
.map_err(|e| VendorEventError::MalformedPayload(e.to_string()))?;
validate_source(&record.id)?;
validate_optional_text(record.location_id.as_deref())?;
validate_optional_text(record.ap_id.as_deref())?;
validate_optional_text(record.ssid.as_deref())?;
let timestamp_us = parse_timestamp(record.timestamp.as_ref())?;
let mut metrics = BTreeMap::new();
push_metric(&mut metrics, "rssi_dbm", record.rssi, -150.0, 20.0)?;
push_metric(&mut metrics, "snr_db", record.snr, -50.0, 100.0)?;
push_metric(&mut metrics, "channel", record.channel, 1.0, 7_000.0)?;
push_metric(
&mut metrics,
"tx_rate_mbps",
record.tx_rate_mbps,
0.0,
100_000.0,
)?;
push_metric(
&mut metrics,
"rx_rate_mbps",
record.rx_rate_mbps,
0.0,
100_000.0,
)?;
if metrics.is_empty() {
return Err(VendorEventError::InvalidPayload);
}
Ok(VendorRfEvent {
vendor: VendorId::Netgear,
capability: RfCapability::RfTelemetry,
sequence: stable_sequence(&record.id, timestamp_us, index),
timestamp_us,
source_id: record.id,
synthetic: false,
metrics,
label: context_label(&[record.location_id, record.ap_id, record.ssid])?,
})
}
fn parse_timestamp(value: Option<&Value>) -> Result<u64, VendorEventError> {
let value = value.ok_or(VendorEventError::InvalidPayload)?;
if let Some(text) = value.as_str() {
if let Ok(raw) = text.parse::<u64>() {
return normalize_integer_timestamp(raw);
}
if let Ok(raw) = text.parse::<f64>() {
return normalize_fractional_timestamp(raw);
}
let parsed = chrono::DateTime::parse_from_rfc3339(text)
.map_err(|_| VendorEventError::InvalidPayload)?;
return u64::try_from(parsed.timestamp_micros())
.map_err(|_| VendorEventError::InvalidPayload);
}
if let Some(raw) = value.as_u64() {
return normalize_integer_timestamp(raw);
}
normalize_fractional_timestamp(value.as_f64().ok_or(VendorEventError::InvalidPayload)?)
}
fn normalize_integer_timestamp(raw: u64) -> Result<u64, VendorEventError> {
if raw == 0 {
return Err(VendorEventError::InvalidPayload);
}
// Normalize seconds, milliseconds, or microseconds to microseconds.
if raw < 10_000_000_000 {
raw.checked_mul(1_000_000)
.ok_or(VendorEventError::InvalidPayload)
} else if raw < 10_000_000_000_000 {
raw.checked_mul(1_000)
.ok_or(VendorEventError::InvalidPayload)
} else if raw < 10_000_000_000_000_000 {
Ok(raw)
} else {
Err(VendorEventError::InvalidPayload)
}
}
fn normalize_fractional_timestamp(raw: f64) -> Result<u64, VendorEventError> {
if !raw.is_finite() || raw <= 0.0 {
return Err(VendorEventError::InvalidPayload);
}
let micros = if raw < 10_000_000_000.0 {
raw * 1_000_000.0
} else if raw < 10_000_000_000_000.0 {
raw * 1_000.0
} else if raw < 10_000_000_000_000_000.0 {
raw
} else {
return Err(VendorEventError::InvalidPayload);
};
if !micros.is_finite() || micros > u64::MAX as f64 {
return Err(VendorEventError::InvalidPayload);
}
Ok(micros.round() as u64)
}
fn push_metric(
metrics: &mut BTreeMap<String, f64>,
name: &str,
value: Option<f64>,
minimum: f64,
maximum: f64,
) -> Result<(), VendorEventError> {
if let Some(value) = value {
if !value.is_finite() || !(minimum..=maximum).contains(&value) {
return Err(VendorEventError::InvalidPayload);
}
metrics.insert(name.into(), value);
}
Ok(())
}
fn validate_source(value: &str) -> Result<(), VendorEventError> {
if value.is_empty() || value.len() > 256 || value.chars().any(char::is_control) {
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
fn validate_optional_text(value: Option<&str>) -> Result<(), VendorEventError> {
if value.is_some_and(|v| v.is_empty() || v.len() > 256 || v.chars().any(char::is_control)) {
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
fn context_label(parts: &[Option<String>]) -> Result<Option<String>, VendorEventError> {
let label = parts
.iter()
.filter_map(Option::as_deref)
.collect::<Vec<_>>()
.join("/");
if label.len() > 256 {
return Err(VendorEventError::InvalidPayload);
}
Ok((!label.is_empty()).then_some(label))
}
fn stable_sequence(source: &str, timestamp_us: u64, index: usize) -> u64 {
// FNV-1a gives a deterministic correlation key without process-random state.
let mut hash = 0xcbf29ce484222325_u64;
for byte in source
.bytes()
.chain(timestamp_us.to_le_bytes())
.chain((index as u64).to_le_bytes())
{
hash ^= u64::from(byte);
hash = hash.wrapping_mul(0x100000001b3);
}
hash
}
#[cfg(test)]
mod tests {
use super::*;
const MIST_FIXTURE: &[u8] = br#"{
"results": [
{"id":"aa:bb:cc:dd:ee:ff","timestamp":1710000000,"rssi":-47,"snr":31,"channel":44,"x":12.5,"y":8.25,"site_id":"site-1","ap_id":"ap-7"},
{"client_id":"station-2","last_seen":1710000000123,"rssi_dbm":-62,"event_type":"client-info"}
],
"next_cursor":"page-2"
}"#;
const NETGEAR_FIXTURE: &[u8] = br#"{
"data": [
{"clientId":"client-1","observedAt":"1710000000000000","signalStrength":-53,"signalToNoiseRatio":24,"channel":149,"txRateMbps":866.7,"rxRateMbps":721.2,"locationId":"office","accessPointId":"ap-2","ssidName":"lab"}
],
"pagination":{"next":"cursor-2"}
}"#;
#[test]
fn descriptors_are_honest_and_valid() {
for descriptor in [mist_descriptor(), netgear_descriptor()] {
descriptor.validate().unwrap();
assert!(!descriptor.hardware_validated);
assert_eq!(
descriptor.availability,
ProviderAvailability::CredentialsRequired
);
assert!(!descriptor.capabilities.contains(&RfCapability::ComplexCsi));
}
}
#[test]
fn mist_rest_and_webhook_fixture_is_deterministic() {
let provider = MistProvider;
let first = provider.decode_page(MIST_FIXTURE).unwrap();
let second = provider.decode_page(MIST_FIXTURE).unwrap();
assert_eq!(first, second);
assert_eq!(first.next_cursor.as_deref(), Some("page-2"));
assert_eq!(first.events.len(), 2);
assert_eq!(first.events[0].timestamp_us, 1_710_000_000_000_000);
assert_eq!(first.events[1].timestamp_us, 1_710_000_000_123_000);
assert_eq!(first.events[0].metrics["x_m"], 12.5);
assert!(!first.events[0].synthetic);
}
#[test]
fn netgear_page_fixture_normalizes_aliases() {
let page = NetgearInsightProvider.decode_page(NETGEAR_FIXTURE).unwrap();
assert_eq!(page.next_cursor.as_deref(), Some("cursor-2"));
assert_eq!(page.events.len(), 1);
let event = &page.events[0];
assert_eq!(event.vendor, VendorId::Netgear);
assert_eq!(event.capability, RfCapability::RfTelemetry);
assert_eq!(event.metrics["tx_rate_mbps"], 866.7);
assert_eq!(event.label.as_deref(), Some("office/ap-2/lab"));
}
#[test]
fn top_level_array_is_supported_without_pagination() {
let payload = br#"[{"mac":"a","timestamp":1710000000,"rssi":-40}]"#;
let page = MistProvider.decode_page(payload).unwrap();
assert_eq!(page.events.len(), 1);
assert_eq!(page.next_cursor, None);
}
#[test]
fn payload_and_page_bounds_are_enforced() {
assert_eq!(
MistProvider.decode(&vec![b' '; MAX_VENDOR_PAYLOAD_BYTES + 1]),
Err(VendorEventError::InvalidPayload)
);
let values = (0..=MAX_EVENTS_PER_PAGE)
.map(|_| serde_json::json!({"id":"a","timestamp":1710000000,"rssi":-40}))
.collect::<Vec<_>>();
let bytes = serde_json::to_vec(&values).unwrap();
assert_eq!(
MistProvider.decode(&bytes),
Err(VendorEventError::InvalidPayload)
);
}
#[test]
fn missing_identity_timestamp_or_metrics_fails_closed() {
for payload in [
br#"[{"timestamp":1710000000,"rssi":-40}]"#.as_slice(),
br#"[{"id":"a","rssi":-40}]"#.as_slice(),
br#"[{"id":"a","timestamp":1710000000}]"#.as_slice(),
] {
assert!(MistProvider.decode(payload).is_err());
}
}
#[test]
fn invalid_metric_cursor_and_record_schema_fail_closed() {
assert!(MistProvider
.decode(br#"[{"id":"a","timestamp":1710000000,"rssi":999}]"#)
.is_err());
assert!(NetgearInsightProvider.decode(br#"{"data":[42]}"#).is_err());
assert!(MistProvider
.decode_page(
br#"{"results":[{"id":"a","timestamp":1710000000,"rssi":-40}],"next_cursor":"x&admin=true"}"#
)
.is_err());
}
#[test]
fn request_configuration_is_regional_validated_and_redacted() {
let token = SecretToken::new("super-secret").unwrap();
let mist = mist_request(
MistRegion::Europe,
"site_1",
Some("page-2".into()),
token.clone(),
)
.unwrap();
assert_eq!(mist.base_url, "https://api.eu.mist.com");
assert_eq!(mist.path, "/api/v1/sites/site_1/stats/clients");
let netgear = netgear_request(NetgearRegion::Australia, "location-7", None, token).unwrap();
assert_eq!(netgear.base_url, "https://au.insight.netgear.com");
assert!(!format!("{netgear:?}").contains("super-secret"));
assert!(!format!("{:?}", netgear.token).contains("super-secret"));
assert!(mist_request(
MistRegion::Global,
"../other-site",
None,
SecretToken::new("token").unwrap()
)
.is_err());
}
#[test]
fn token_rejects_header_injection() {
assert!(SecretToken::new("token\r\nX-Injected: yes").is_err());
assert!(SecretToken::new("").is_err());
}
#[test]
fn timestamp_units_are_normalized_and_extremes_rejected() {
assert_eq!(
parse_timestamp(Some(&serde_json::json!(1_710_000_000))).unwrap(),
1_710_000_000_000_000
);
assert_eq!(
parse_timestamp(Some(&serde_json::json!(1_710_000_000_123_u64))).unwrap(),
1_710_000_000_123_000
);
assert!(parse_timestamp(Some(&serde_json::json!(0))).is_err());
assert!(parse_timestamp(Some(&serde_json::json!(-1))).is_err());
assert!(parse_timestamp(Some(&serde_json::json!(u64::MAX))).is_err());
assert_eq!(
parse_timestamp(Some(&serde_json::json!("2024-03-09T16:00:00Z"))).unwrap(),
1_710_000_000_000_000
);
assert_eq!(
parse_timestamp(Some(&serde_json::json!(1710000000.25))).unwrap(),
1_710_000_000_250_000
);
}
}
@@ -0,0 +1,598 @@
//! Capability-safe ADR-270 adapters for Origin AI and Plume/OpenSync.
//!
//! Origin publishes derived sensing results through contract-gated APIs and
//! webhooks. Public documentation does not define stable paths, so paths are
//! supplied by the contracted deployment configuration. OpenSync exports RF
//! and network telemetry; neither adapter promotes scalar data to complex CSI.
use serde::Deserialize;
use serde_json::{json, Value};
use std::collections::BTreeMap;
use wifi_densepose_hardware::vendor_rf::{
ProviderAvailability, ProviderDescriptor, RfCapability, VendorEventError, VendorId,
VendorRfEvent, VendorRfProvider,
};
const MAX_PAYLOAD_BYTES: usize = 256 * 1024;
const MAX_EVENTS_PER_PAYLOAD: usize = 256;
const MAX_ENDPOINT_LEN: usize = 2048;
const MAX_ENV_NAME_LEN: usize = 128;
/// A request plan deliberately containing a credential *reference*, not a secret.
#[derive(Debug, Clone, PartialEq)]
pub struct VendorRequest {
pub method: &'static str,
pub endpoint: String,
pub headers: BTreeMap<String, String>,
pub credential_env: Option<String>,
pub body: Option<Value>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OriginAiConfig {
/// Contract-provided HTTPS sensing-server base URL.
pub base_url: String,
/// Contract-provided relative event API path. No public path is assumed.
pub event_path: String,
/// Name of the environment variable holding the partner bearer token.
pub token_env: String,
}
impl OriginAiConfig {
pub fn validate(&self) -> Result<(), VendorEventError> {
validate_https_base(&self.base_url)?;
validate_relative_path(&self.event_path)?;
validate_env_name(&self.token_env)
}
/// Constructs a GET plan. The caller resolves `credential_env` at execution time.
pub fn events_request(&self) -> Result<VendorRequest, VendorEventError> {
self.validate()?;
Ok(VendorRequest {
method: "GET",
endpoint: join_endpoint(&self.base_url, &self.event_path),
headers: BTreeMap::from([("accept".into(), "application/json".into())]),
credential_env: Some(self.token_env.clone()),
body: None,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PlumeOpenSyncConfig {
/// HTTPS northbound/sandbox URL supplied by the operator.
pub base_url: String,
/// Relative endpoint accepting OVSDB JSON-RPC for this deployment.
pub ovsdb_path: String,
/// Optional environment variable used by a protected northbound endpoint.
pub token_env: Option<String>,
}
impl PlumeOpenSyncConfig {
pub fn validate(&self) -> Result<(), VendorEventError> {
validate_https_base(&self.base_url)?;
validate_relative_path(&self.ovsdb_path)?;
if let Some(name) = &self.token_env {
validate_env_name(name)?;
}
Ok(())
}
/// Builds a read-only OVSDB `select` transaction for an allow-listed table.
pub fn select_request(&self, table: &str) -> Result<VendorRequest, VendorEventError> {
self.validate()?;
if !matches!(
table,
"Wifi_Radio_State" | "Wifi_VIF_State" | "Wifi_Associated_Clients"
) {
return Err(VendorEventError::InvalidPayload);
}
Ok(VendorRequest {
method: "POST",
endpoint: join_endpoint(&self.base_url, &self.ovsdb_path),
headers: BTreeMap::from([
("accept".into(), "application/json".into()),
("content-type".into(), "application/json".into()),
]),
credential_env: self.token_env.clone(),
body: Some(json!({
"jsonrpc": "2.0",
"id": "ruview-read-only",
"method": "transact",
"params": ["Open_vSwitch", {"op": "select", "table": table, "where": []}]
})),
})
}
}
#[derive(Debug, Clone, Default)]
pub struct OriginAiProvider;
impl OriginAiProvider {
pub fn synthetic_fixture(seed: u64, count: usize) -> Vec<VendorRfEvent> {
let count = count.min(MAX_EVENTS_PER_PAYLOAD);
(0..count)
.map(|index| {
let state = splitmix64(seed.wrapping_add(index as u64));
let motion = (state & 1) as f64;
let confidence = 0.70 + ((state >> 8) % 300) as f64 / 1000.0;
VendorRfEvent {
vendor: VendorId::OriginAi,
capability: RfCapability::DerivedSensing,
sequence: index as u64,
timestamp_us: 1_700_000_000_000_000 + index as u64 * 100_000,
source_id: format!("origin-sim-{:08x}", seed as u32),
synthetic: true,
metrics: BTreeMap::from([
("motion".into(), motion),
("confidence".into(), confidence),
]),
label: Some(if motion == 1.0 { "motion" } else { "clear" }.into()),
}
})
.collect()
}
}
impl VendorRfProvider for OriginAiProvider {
fn descriptor(&self) -> ProviderDescriptor {
ProviderDescriptor {
vendor: VendorId::OriginAi,
capabilities: vec![RfCapability::DerivedSensing],
availability: ProviderAvailability::ContractRequired,
hardware_validated: false,
reason: "Origin partner API/SDK access is contract-gated; adapter accepts derived sensing only"
.into(),
}
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
check_payload(payload)?;
let envelope: OriginEnvelope = decode_json(payload)?;
if envelope.events.is_empty() || envelope.events.len() > MAX_EVENTS_PER_PAYLOAD {
return Err(VendorEventError::InvalidPayload);
}
envelope
.events
.into_iter()
.map(|event| event.into_vendor_event(&self.descriptor()))
.collect()
}
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct OriginEnvelope {
events: Vec<OriginEvent>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct OriginEvent {
sequence: u64,
timestamp_us: u64,
source_id: String,
kind: OriginKind,
confidence: f64,
#[serde(default)]
value: Option<f64>,
#[serde(default)]
label: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "snake_case")]
enum OriginKind {
Motion,
Occupancy,
Presence,
Fall,
BreathingRate,
}
impl OriginEvent {
fn into_vendor_event(
self,
descriptor: &ProviderDescriptor,
) -> Result<VendorRfEvent, VendorEventError> {
if self.timestamp_us == 0 || !(0.0..=1.0).contains(&self.confidence) {
return Err(VendorEventError::InvalidPayload);
}
let (metric, requires_value) = match self.kind {
OriginKind::Motion => ("motion", false),
OriginKind::Occupancy => ("occupancy", false),
OriginKind::Presence => ("presence", false),
OriginKind::Fall => ("fall", false),
OriginKind::BreathingRate => ("breathing_rate_bpm", true),
};
let value = self.value.ok_or(VendorEventError::InvalidPayload)?;
if !value.is_finite()
|| (requires_value && !(1.0..=120.0).contains(&value))
|| (!requires_value && value != 0.0 && value != 1.0)
{
return Err(VendorEventError::InvalidPayload);
}
let event = VendorRfEvent {
vendor: VendorId::OriginAi,
capability: RfCapability::DerivedSensing,
sequence: self.sequence,
timestamp_us: self.timestamp_us,
source_id: self.source_id,
synthetic: false,
metrics: BTreeMap::from([
(metric.into(), value),
("confidence".into(), self.confidence),
]),
label: self.label,
};
event.validate(descriptor)?;
Ok(event)
}
}
#[derive(Debug, Clone, Default)]
pub struct PlumeOpenSyncProvider;
impl PlumeOpenSyncProvider {
pub fn synthetic_fixture(seed: u64, count: usize) -> Vec<VendorRfEvent> {
let count = count.min(MAX_EVENTS_PER_PAYLOAD);
(0..count)
.map(|index| {
let state = splitmix64(seed.wrapping_add(index as u64));
VendorRfEvent {
vendor: VendorId::Plume,
capability: RfCapability::RfTelemetry,
sequence: index as u64,
timestamp_us: 1_700_000_000_000_000 + index as u64 * 250_000,
source_id: format!("opensync-sim-{:08x}", seed as u32),
synthetic: true,
metrics: BTreeMap::from([
("rssi_dbm".into(), -30.0 - (state % 55) as f64),
("channel".into(), 1.0 + ((state >> 8) % 165) as f64),
(
"noise_floor_dbm".into(),
-100.0 + ((state >> 16) % 12) as f64,
),
]),
label: None,
}
})
.collect()
}
}
impl VendorRfProvider for PlumeOpenSyncProvider {
fn descriptor(&self) -> ProviderDescriptor {
ProviderDescriptor {
vendor: VendorId::Plume,
capabilities: vec![RfCapability::RfTelemetry],
availability: ProviderAvailability::CredentialsRequired,
hardware_validated: false,
reason: "OpenSync radio/client telemetry only; Plume Sense is a separate gated service"
.into(),
}
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
check_payload(payload)?;
let envelope: OpenSyncEnvelope = decode_json(payload)?;
if envelope.observations.is_empty() || envelope.observations.len() > MAX_EVENTS_PER_PAYLOAD
{
return Err(VendorEventError::InvalidPayload);
}
envelope
.observations
.into_iter()
.map(|event| event.into_vendor_event(&self.descriptor()))
.collect()
}
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct OpenSyncEnvelope {
observations: Vec<OpenSyncObservation>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct OpenSyncObservation {
sequence: u64,
timestamp_us: u64,
source_id: String,
rssi_dbm: f64,
channel: u16,
#[serde(default)]
noise_floor_dbm: Option<f64>,
#[serde(default)]
tx_rate_mbps: Option<f64>,
#[serde(default)]
rx_rate_mbps: Option<f64>,
#[serde(default)]
clients: Option<u16>,
}
impl OpenSyncObservation {
fn into_vendor_event(
self,
descriptor: &ProviderDescriptor,
) -> Result<VendorRfEvent, VendorEventError> {
if self.timestamp_us == 0
|| !(-127.0..=0.0).contains(&self.rssi_dbm)
|| self.channel == 0
|| self.channel > 233
{
return Err(VendorEventError::InvalidPayload);
}
let mut metrics = BTreeMap::from([
("rssi_dbm".into(), self.rssi_dbm),
("channel".into(), self.channel as f64),
]);
insert_optional(
&mut metrics,
"noise_floor_dbm",
self.noise_floor_dbm,
-127.0,
0.0,
)?;
insert_optional(
&mut metrics,
"tx_rate_mbps",
self.tx_rate_mbps,
0.0,
100_000.0,
)?;
insert_optional(
&mut metrics,
"rx_rate_mbps",
self.rx_rate_mbps,
0.0,
100_000.0,
)?;
if let Some(clients) = self.clients {
metrics.insert("clients".into(), clients as f64);
}
let event = VendorRfEvent {
vendor: VendorId::Plume,
capability: RfCapability::RfTelemetry,
sequence: self.sequence,
timestamp_us: self.timestamp_us,
source_id: self.source_id,
synthetic: false,
metrics,
label: None,
};
event.validate(descriptor)?;
Ok(event)
}
}
fn insert_optional(
metrics: &mut BTreeMap<String, f64>,
key: &str,
value: Option<f64>,
minimum: f64,
maximum: f64,
) -> Result<(), VendorEventError> {
if let Some(value) = value {
if !value.is_finite() || !(minimum..=maximum).contains(&value) {
return Err(VendorEventError::InvalidPayload);
}
metrics.insert(key.into(), value);
}
Ok(())
}
fn check_payload(payload: &[u8]) -> Result<(), VendorEventError> {
if payload.is_empty() || payload.len() > MAX_PAYLOAD_BYTES {
Err(VendorEventError::InvalidPayload)
} else {
Ok(())
}
}
fn decode_json<T: for<'de> Deserialize<'de>>(payload: &[u8]) -> Result<T, VendorEventError> {
serde_json::from_slice(payload).map_err(|error| {
let message = error.to_string();
VendorEventError::MalformedPayload(message.chars().take(160).collect())
})
}
fn validate_https_base(value: &str) -> Result<(), VendorEventError> {
let authority = value
.strip_prefix("https://")
.and_then(|rest| rest.split('/').next())
.unwrap_or_default();
if value.len() > MAX_ENDPOINT_LEN
|| authority.is_empty()
|| authority.contains('@')
|| authority.chars().any(char::is_whitespace)
|| value.contains(['\r', '\n', '#', '?'])
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
fn validate_relative_path(value: &str) -> Result<(), VendorEventError> {
if value.len() > MAX_ENDPOINT_LEN
|| !value.starts_with('/')
|| value.starts_with("//")
|| value.contains(['\r', '\n', '#', '?'])
|| value.split('/').any(|segment| segment == "..")
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
fn validate_env_name(value: &str) -> Result<(), VendorEventError> {
if value.is_empty()
|| value.len() > MAX_ENV_NAME_LEN
|| !value
.bytes()
.all(|byte| byte.is_ascii_uppercase() || byte.is_ascii_digit() || byte == b'_')
|| value.as_bytes()[0].is_ascii_digit()
{
return Err(VendorEventError::InvalidPayload);
}
Ok(())
}
fn join_endpoint(base: &str, path: &str) -> String {
format!("{}{}", base.trim_end_matches('/'), path)
}
fn splitmix64(mut value: u64) -> u64 {
value = value.wrapping_add(0x9e37_79b9_7f4a_7c15);
value = (value ^ (value >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
value = (value ^ (value >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
value ^ (value >> 31)
}
#[cfg(test)]
mod tests {
use super::*;
use wifi_densepose_hardware::vendor_rf::MAX_VENDOR_TEXT_LEN;
#[test]
fn descriptors_are_honest_about_capabilities_and_access() {
let origin = OriginAiProvider.descriptor();
assert_eq!(origin.capabilities, vec![RfCapability::DerivedSensing]);
assert_eq!(origin.availability, ProviderAvailability::ContractRequired);
let plume = PlumeOpenSyncProvider.descriptor();
assert_eq!(plume.capabilities, vec![RfCapability::RfTelemetry]);
assert_eq!(
plume.availability,
ProviderAvailability::CredentialsRequired
);
assert!(!origin.capabilities.contains(&RfCapability::ComplexCsi));
assert!(!plume.capabilities.contains(&RfCapability::ComplexCsi));
}
#[test]
fn origin_decodes_only_derived_sensing() {
let payload = br#"{"events":[{"sequence":7,"timestamp_us":42,"source_id":"zone-a","kind":"occupancy","confidence":0.91,"value":1,"label":"occupied"}]}"#;
let events = OriginAiProvider.decode(payload).unwrap();
assert_eq!(events[0].capability, RfCapability::DerivedSensing);
assert_eq!(events[0].metrics["occupancy"], 1.0);
assert!(!events[0].synthetic);
}
#[test]
fn origin_rejects_unknown_raw_csi_and_bad_values() {
let raw = br#"{"events":[{"sequence":1,"timestamp_us":2,"source_id":"z","kind":"motion","confidence":1,"value":1,"raw_csi":[[1,2]]}]}"#;
assert!(OriginAiProvider.decode(raw).is_err());
let bad_confidence = br#"{"events":[{"sequence":1,"timestamp_us":2,"source_id":"z","kind":"motion","confidence":1.1,"value":1}]}"#;
assert_eq!(
OriginAiProvider.decode(bad_confidence),
Err(VendorEventError::InvalidPayload)
);
}
#[test]
fn plume_decodes_telemetry_without_inventing_csi() {
let payload = br#"{"observations":[{"sequence":9,"timestamp_us":10,"source_id":"pod-a","rssi_dbm":-54,"channel":36,"noise_floor_dbm":-94,"tx_rate_mbps":1200,"clients":4}]}"#;
let events = PlumeOpenSyncProvider.decode(payload).unwrap();
assert_eq!(events[0].capability, RfCapability::RfTelemetry);
assert_eq!(events[0].metrics["rssi_dbm"], -54.0);
assert!(!events[0].metrics.contains_key("csi"));
}
#[test]
fn malformed_empty_oversized_and_unbounded_arrays_fail_closed() {
assert!(OriginAiProvider.decode(b"{").is_err());
assert_eq!(
OriginAiProvider.decode(b""),
Err(VendorEventError::InvalidPayload)
);
assert_eq!(
PlumeOpenSyncProvider.decode(&vec![b' '; MAX_PAYLOAD_BYTES + 1]),
Err(VendorEventError::InvalidPayload)
);
let event = r#"{"sequence":1,"timestamp_us":2,"source_id":"p","rssi_dbm":-50,"channel":1}"#;
let payload = format!("{{\"observations\":[{}]}}", vec![event; 257].join(","));
assert_eq!(
PlumeOpenSyncProvider.decode(payload.as_bytes()),
Err(VendorEventError::InvalidPayload)
);
}
#[test]
fn telemetry_ranges_are_enforced() {
let payload = br#"{"observations":[{"sequence":1,"timestamp_us":2,"source_id":"p","rssi_dbm":4,"channel":36}]}"#;
assert_eq!(
PlumeOpenSyncProvider.decode(payload),
Err(VendorEventError::InvalidPayload)
);
}
#[test]
fn fixtures_are_deterministic_bounded_and_marked_synthetic() {
let a = OriginAiProvider::synthetic_fixture(19, 4);
assert_eq!(a, OriginAiProvider::synthetic_fixture(19, 4));
assert!(a.iter().all(|event| event.synthetic));
let p = PlumeOpenSyncProvider::synthetic_fixture(23, usize::MAX);
assert_eq!(p.len(), MAX_EVENTS_PER_PAYLOAD);
assert!(p.iter().all(|event| event.synthetic));
assert_eq!(p, PlumeOpenSyncProvider::synthetic_fixture(23, usize::MAX));
}
#[test]
fn request_plans_reference_secrets_without_embedding_them() {
let origin = OriginAiConfig {
base_url: "https://partner.example".into(),
event_path: "/contract/v1/events".into(),
token_env: "ORIGIN_AI_TOKEN".into(),
}
.events_request()
.unwrap();
assert_eq!(
origin.endpoint,
"https://partner.example/contract/v1/events"
);
assert_eq!(origin.credential_env.as_deref(), Some("ORIGIN_AI_TOKEN"));
assert!(format!("{origin:?}").find("Bearer ").is_none());
let plume = PlumeOpenSyncConfig {
base_url: "https://sandbox.example".into(),
ovsdb_path: "/ovsdb".into(),
token_env: Some("OPENSYNC_TOKEN".into()),
}
.select_request("Wifi_Radio_State")
.unwrap();
assert_eq!(plume.body.as_ref().unwrap()["method"], "transact");
assert_eq!(plume.credential_env.as_deref(), Some("OPENSYNC_TOKEN"));
}
#[test]
fn request_validation_rejects_injection_and_write_tables() {
let config = PlumeOpenSyncConfig {
base_url: "https://sandbox.example".into(),
ovsdb_path: "/ovsdb".into(),
token_env: None,
};
assert_eq!(
config.select_request("AWLAN_Node"),
Err(VendorEventError::InvalidPayload)
);
let bad = OriginAiConfig {
base_url: "http://insecure.example".into(),
event_path: "/events\r\nx: y".into(),
token_env: "token".into(),
};
assert_eq!(bad.events_request(), Err(VendorEventError::InvalidPayload));
}
#[test]
fn labels_and_source_ids_obey_shared_contract_bounds() {
let label = "x".repeat(MAX_VENDOR_TEXT_LEN + 1);
let payload = format!(
"{{\"events\":[{{\"sequence\":1,\"timestamp_us\":2,\"source_id\":\"z\",\"kind\":\"motion\",\"confidence\":1,\"value\":1,\"label\":\"{label}\"}}]}}"
);
assert_eq!(
OriginAiProvider.decode(payload.as_bytes()),
Err(VendorEventError::InvalidPayload)
);
}
}
@@ -0,0 +1,481 @@
//! ADR-270 providers whose useful integration surface is scalar telemetry,
//! network-only metadata, or an explicit no-go decision.
//!
//! None of these providers emits complex CSI. The small JSON contract in this
//! module is intended for sidecars and deterministic replay fixtures; it is not
//! a claim that a vendor exposes this exact wire format.
use serde::Deserialize;
use std::collections::BTreeMap;
use wifi_densepose_hardware::vendor_rf::{
ProviderAvailability, ProviderDescriptor, RfCapability, VendorEventError, VendorId,
VendorRfEvent, VendorRfProvider,
};
/// Upper bound applied before JSON decoding, limiting parser allocation.
pub const MAX_REMAINING_VENDOR_PAYLOAD_BYTES: usize = 64 * 1024;
/// Upper bound on events accepted in one sidecar envelope.
pub const MAX_REMAINING_VENDOR_EVENTS: usize = 256;
const ELECTRIC_IMP_METRICS: &[MetricRule] = &[
MetricRule::new("battery_v", 0.0, 100.0, false),
MetricRule::new("humidity_percent", 0.0, 100.0, false),
MetricRule::new("rssi_dbm", -127.0, 0.0, false),
MetricRule::new("temperature_c", -100.0, 200.0, false),
MetricRule::new("voltage_v", 0.0, 1_000.0, false),
];
const RF_SOLUTIONS_METRICS: &[MetricRule] = &[
MetricRule::new("battery_v", 0.0, 100.0, false),
MetricRule::new("humidity_percent", 0.0, 100.0, false),
MetricRule::new("relay_state", 0.0, 1.0, true),
MetricRule::new("rssi_dbm", -127.0, 0.0, false),
MetricRule::new("temperature_c", -100.0, 200.0, false),
];
const LUMA_METRICS: &[MetricRule] = &[
MetricRule::new("client_count", 0.0, 1_000_000.0, true),
MetricRule::new("noise_dbm", -127.0, 0.0, false),
MetricRule::new("rssi_dbm", -127.0, 0.0, false),
MetricRule::new("rx_bytes", 0.0, 9_007_199_254_740_991.0, true),
MetricRule::new("tx_bytes", 0.0, 9_007_199_254_740_991.0, true),
];
const GOOGLE_NEST_METRICS: &[MetricRule] = &[
MetricRule::new("client_count", 0.0, 1_000_000.0, true),
MetricRule::new("probe_count", 0.0, 1_000_000_000.0, true),
MetricRule::new("rx_bytes", 0.0, 9_007_199_254_740_991.0, true),
MetricRule::new("tx_bytes", 0.0, 9_007_199_254_740_991.0, true),
];
#[derive(Debug, Clone, Copy)]
struct MetricRule {
name: &'static str,
minimum: f64,
maximum: f64,
integer: bool,
}
impl MetricRule {
const fn new(name: &'static str, minimum: f64, maximum: f64, integer: bool) -> Self {
Self {
name,
minimum,
maximum,
integer,
}
}
fn accepts(self, value: f64) -> bool {
value.is_finite()
&& (self.minimum..=self.maximum).contains(&value)
&& (!self.integer || value.fract() == 0.0)
}
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct ScalarEnvelope {
events: Vec<ScalarEvent>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct ScalarEvent {
sequence: u64,
timestamp_us: u64,
source_id: String,
synthetic: bool,
metrics: BTreeMap<String, f64>,
#[serde(default)]
label: Option<String>,
}
fn descriptor(
vendor: VendorId,
capability: RfCapability,
availability: ProviderAvailability,
reason: &str,
) -> ProviderDescriptor {
ProviderDescriptor {
vendor,
capabilities: vec![capability],
availability,
hardware_validated: false,
reason: reason.to_owned(),
}
}
fn decode_bounded_scalar_events(
payload: &[u8],
provider: &ProviderDescriptor,
capability: RfCapability,
allowed_metrics: &[MetricRule],
) -> Result<Vec<VendorRfEvent>, VendorEventError> {
provider.validate()?;
if !provider.capabilities.contains(&capability)
|| matches!(
capability,
RfCapability::ComplexCsi | RfCapability::Unsupported
)
{
return Err(VendorEventError::CapabilityMismatch);
}
if payload.is_empty() || payload.len() > MAX_REMAINING_VENDOR_PAYLOAD_BYTES {
return Err(VendorEventError::InvalidPayload);
}
let envelope: ScalarEnvelope = serde_json::from_slice(payload).map_err(|error| {
VendorEventError::MalformedPayload(error.to_string().chars().take(160).collect())
})?;
if envelope.events.is_empty() || envelope.events.len() > MAX_REMAINING_VENDOR_EVENTS {
return Err(VendorEventError::InvalidPayload);
}
envelope
.events
.into_iter()
.map(|event| {
// An allowlist keeps arbitrary scalar fields from silently changing
// the meaning of a provider contract (and rejects CSI-shaped data).
if event.timestamp_us == 0
|| event.source_id.chars().any(char::is_control)
|| event
.label
.as_deref()
.is_some_and(|label| label.chars().any(char::is_control))
{
return Err(VendorEventError::InvalidPayload);
}
for (key, value) in &event.metrics {
let rule = allowed_metrics
.iter()
.find(|rule| rule.name == key)
.ok_or(VendorEventError::CapabilityMismatch)?;
if !rule.accepts(*value) {
return Err(VendorEventError::InvalidPayload);
}
}
let event = VendorRfEvent {
vendor: provider.vendor,
capability,
sequence: event.sequence,
timestamp_us: event.timestamp_us,
source_id: event.source_id,
synthetic: event.synthetic,
metrics: event.metrics,
label: event.label,
};
event.validate(provider)?;
Ok(event)
})
.collect()
}
/// Electric Imp agent/impCentral scalar telemetry bridge for existing fleets.
#[derive(Debug, Clone, Copy, Default)]
pub struct ElectricImpProvider;
impl VendorRfProvider for ElectricImpProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::ElectricImp,
RfCapability::RfTelemetry,
ProviderAvailability::CredentialsRequired,
"Optional authenticated agent/impCentral scalar telemetry bridge; never CSI",
)
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
let descriptor = self.descriptor();
decode_bounded_scalar_events(
payload,
&descriptor,
RfCapability::RfTelemetry,
ELECTRIC_IMP_METRICS,
)
}
}
/// RF Solutions environmental/RIoT scalar telemetry boundary.
#[derive(Debug, Clone, Copy, Default)]
pub struct RfSolutionsProvider;
impl VendorRfProvider for RfSolutionsProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::RfSolutions,
RfCapability::RfTelemetry,
ProviderAvailability::Experimental,
"Optional non-Wi-Fi environmental telemetry fusion; excluded as a CSI source",
)
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
let descriptor = self.descriptor();
decode_bounded_scalar_events(
payload,
&descriptor,
RfCapability::RfTelemetry,
RF_SOLUTIONS_METRICS,
)
}
}
/// Generic OpenWrt telemetry fixture for already-owned discontinued Luma units.
#[derive(Debug, Clone, Copy, Default)]
pub struct LumaOpenWrtProvider;
impl VendorRfProvider for LumaOpenWrtProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::Luma,
RfCapability::RfTelemetry,
ProviderAvailability::Experimental,
"Generic OpenWrt scalar telemetry fixture for already-owned Luma hardware; no Luma CSI claim",
)
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
let descriptor = self.descriptor();
decode_bounded_scalar_events(
payload,
&descriptor,
RfCapability::RfTelemetry,
LUMA_METRICS,
)
}
}
/// Google Nest Wifi may participate as network infrastructure, not a sensor.
#[derive(Debug, Clone, Copy, Default)]
pub struct GoogleNestProvider;
impl VendorRfProvider for GoogleNestProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::GoogleNest,
RfCapability::NetworkOnly,
ProviderAvailability::Experimental,
"Network-only replay events; Device Access exposes no router CSI or RF telemetry",
)
}
fn decode(&self, payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
let descriptor = self.descriptor();
decode_bounded_scalar_events(
payload,
&descriptor,
RfCapability::NetworkOnly,
GOOGLE_NEST_METRICS,
)
}
}
/// Linksys Aware reached end of support; there is no supported sensing API.
#[derive(Debug, Clone, Copy, Default)]
pub struct LinksysProvider;
impl VendorRfProvider for LinksysProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::Linksys,
RfCapability::Unsupported,
ProviderAvailability::Unsupported,
"Linksys Aware reached end of support in 2024; no supported sensing interface",
)
}
fn decode(&self, _payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
Err(VendorEventError::Unsupported)
}
}
/// Wifigarden remains gated until its commercial SDK contract is disclosed.
#[derive(Debug, Clone, Copy, Default)]
pub struct WifigardenProvider;
impl VendorRfProvider for WifigardenProvider {
fn descriptor(&self) -> ProviderDescriptor {
descriptor(
VendorId::Wifigarden,
RfCapability::Unsupported,
ProviderAvailability::ContractRequired,
"Commercial SDK, chipset, schema, calibration and data-rights disclosure required",
)
}
fn decode(&self, _payload: &[u8]) -> Result<Vec<VendorRfEvent>, VendorEventError> {
Err(VendorEventError::ContractRequired)
}
}
/// Deterministic synthetic Electric Imp sidecar contract fixture.
pub const ELECTRIC_IMP_CONTRACT_FIXTURE: &[u8] = br#"{"events":[{"sequence":7,"timestamp_us":1700000000000000,"source_id":"imp005-fixture","synthetic":true,"metrics":{"rssi_dbm":-48.0,"temperature_c":21.5},"label":"lab"}]}"#;
/// Deterministic synthetic RF Solutions sidecar contract fixture.
pub const RF_SOLUTIONS_CONTRACT_FIXTURE: &[u8] = br#"{"events":[{"sequence":8,"timestamp_us":1700000000000100,"source_id":"riot-fixture","synthetic":true,"metrics":{"battery_v":3.1,"humidity_percent":44.0},"label":"lab"}]}"#;
/// Deterministic synthetic generic OpenWrt/Luma contract fixture.
pub const LUMA_OPENWRT_CONTRACT_FIXTURE: &[u8] = br#"{"events":[{"sequence":9,"timestamp_us":1700000000000200,"source_id":"luma-openwrt-fixture","synthetic":true,"metrics":{"client_count":3.0,"noise_dbm":-91.0,"rx_bytes":1024.0},"label":"openwrt"}]}"#;
/// Deterministic synthetic Google Nest network-only contract fixture.
pub const GOOGLE_NEST_CONTRACT_FIXTURE: &[u8] = br#"{"events":[{"sequence":10,"timestamp_us":1700000000000300,"source_id":"nest-fixture","synthetic":true,"metrics":{"client_count":4.0,"probe_count":2.0},"label":"network_activity"}]}"#;
#[cfg(test)]
mod tests {
use super::*;
fn assert_valid_fixture<P: VendorRfProvider>(
provider: P,
fixture: &[u8],
vendor: VendorId,
capability: RfCapability,
) {
let descriptor = provider.descriptor();
descriptor.validate().expect("descriptor must be valid");
let first = provider
.decode(fixture)
.expect("fixture must decode deterministically");
let second = provider.decode(fixture).expect("fixture must replay");
assert_eq!(first, second);
assert_eq!(first.len(), 1);
assert_eq!(first[0].vendor, vendor);
assert_eq!(first[0].capability, capability);
assert!(first[0].synthetic);
first[0]
.validate(&descriptor)
.expect("fixture event must satisfy provider contract");
}
#[test]
fn scalar_and_network_fixtures_are_deterministic_and_honest() {
assert_valid_fixture(
ElectricImpProvider,
ELECTRIC_IMP_CONTRACT_FIXTURE,
VendorId::ElectricImp,
RfCapability::RfTelemetry,
);
assert_valid_fixture(
RfSolutionsProvider,
RF_SOLUTIONS_CONTRACT_FIXTURE,
VendorId::RfSolutions,
RfCapability::RfTelemetry,
);
assert_valid_fixture(
LumaOpenWrtProvider,
LUMA_OPENWRT_CONTRACT_FIXTURE,
VendorId::Luma,
RfCapability::RfTelemetry,
);
assert_valid_fixture(
GoogleNestProvider,
GOOGLE_NEST_CONTRACT_FIXTURE,
VendorId::GoogleNest,
RfCapability::NetworkOnly,
);
}
#[test]
fn unavailable_providers_fail_before_interpreting_payloads() {
let linksys = LinksysProvider;
assert_eq!(
linksys.descriptor().availability,
ProviderAvailability::Unsupported
);
assert_eq!(
linksys.decode(b"not json"),
Err(VendorEventError::Unsupported)
);
let wifigarden = WifigardenProvider;
assert_eq!(
wifigarden.descriptor().availability,
ProviderAvailability::ContractRequired
);
assert_eq!(
wifigarden.decode(ELECTRIC_IMP_CONTRACT_FIXTURE),
Err(VendorEventError::ContractRequired)
);
}
#[test]
fn rejects_empty_oversized_and_excess_event_payloads() {
let provider = ElectricImpProvider;
assert_eq!(provider.decode(b""), Err(VendorEventError::InvalidPayload));
assert_eq!(
provider.decode(&vec![b' '; MAX_REMAINING_VENDOR_PAYLOAD_BYTES + 1]),
Err(VendorEventError::InvalidPayload)
);
let events = (0..=MAX_REMAINING_VENDOR_EVENTS)
.map(|sequence| {
format!(
r#"{{"sequence":{sequence},"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{{"rssi_dbm":-40.0}}}}"#
)
})
.collect::<Vec<_>>()
.join(",");
let payload = format!(r#"{{"events":[{events}]}}"#);
assert_eq!(
provider.decode(payload.as_bytes()),
Err(VendorEventError::InvalidPayload)
);
}
#[test]
fn rejects_csi_or_cross_provider_metric_masquerading() {
let fake_csi = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{"csi_real":1.0}}]}"#;
assert_eq!(
ElectricImpProvider.decode(fake_csi),
Err(VendorEventError::CapabilityMismatch)
);
let rf_metric_in_network_event = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{"rssi_dbm":-40.0}}]}"#;
assert_eq!(
GoogleNestProvider.decode(rf_metric_in_network_event),
Err(VendorEventError::CapabilityMismatch)
);
}
#[test]
fn rejects_invalid_values_bounds_and_schema_extensions() {
let non_finite = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{"rssi_dbm":1e999}}]}"#;
assert!(matches!(
ElectricImpProvider.decode(non_finite),
Err(VendorEventError::MalformedPayload(_)) | Err(VendorEventError::InvalidPayload)
));
let empty_metrics = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{}}]}"#;
assert_eq!(
ElectricImpProvider.decode(empty_metrics),
Err(VendorEventError::InvalidPayload)
);
let unknown_field = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","synthetic":true,"metrics":{"rssi_dbm":-40.0},"csi":[]}]}"#;
assert!(matches!(
ElectricImpProvider.decode(unknown_field),
Err(VendorEventError::MalformedPayload(_))
));
let missing_provenance = br#"{"events":[{"sequence":1,"timestamp_us":1,"source_id":"x","metrics":{"rssi_dbm":-40.0}}]}"#;
assert!(matches!(
ElectricImpProvider.decode(missing_provenance),
Err(VendorEventError::MalformedPayload(_))
));
}
#[test]
fn no_remaining_provider_claims_complex_csi_or_hardware_validation() {
let descriptors = [
ElectricImpProvider.descriptor(),
RfSolutionsProvider.descriptor(),
LumaOpenWrtProvider.descriptor(),
GoogleNestProvider.descriptor(),
LinksysProvider.descriptor(),
WifigardenProvider.descriptor(),
];
for descriptor in descriptors {
descriptor.validate().expect("descriptor must be valid");
assert!(!descriptor.hardware_validated);
assert!(!descriptor.capabilities.contains(&RfCapability::ComplexCsi));
}
}
}
@@ -0,0 +1,104 @@
//! ADR-270 provider registry and canonical event helpers.
use serde::Serialize;
use wifi_densepose_hardware::vendor_rf::{
ProviderDescriptor, VendorEventError, VendorId, VendorRfEvent, VendorRfProvider,
};
use crate::vendor_mist_netgear::{MistProvider, NetgearInsightProvider};
use crate::vendor_origin_plume::{OriginAiProvider, PlumeOpenSyncProvider};
use crate::vendor_remaining::{
ElectricImpProvider, GoogleNestProvider, LinksysProvider, LumaOpenWrtProvider,
RfSolutionsProvider, WifigardenProvider,
};
pub fn descriptor_for(vendor: VendorId) -> ProviderDescriptor {
match vendor {
VendorId::OriginAi => OriginAiProvider.descriptor(),
VendorId::Plume => PlumeOpenSyncProvider.descriptor(),
VendorId::Mist => MistProvider.descriptor(),
VendorId::Netgear => NetgearInsightProvider.descriptor(),
VendorId::ElectricImp => ElectricImpProvider.descriptor(),
VendorId::RfSolutions => RfSolutionsProvider.descriptor(),
VendorId::Linksys => LinksysProvider.descriptor(),
VendorId::Luma => LumaOpenWrtProvider.descriptor(),
VendorId::GoogleNest => GoogleNestProvider.descriptor(),
VendorId::Wifigarden => WifigardenProvider.descriptor(),
}
}
pub fn descriptors() -> Vec<ProviderDescriptor> {
VendorId::ALL.into_iter().map(descriptor_for).collect()
}
pub fn vendor_from_str(value: &str) -> Option<VendorId> {
VendorId::ALL
.into_iter()
.find(|vendor| vendor.as_str() == value)
}
pub fn decode_provider(
vendor: VendorId,
payload: &[u8],
) -> Result<Vec<VendorRfEvent>, VendorEventError> {
match vendor {
VendorId::OriginAi => OriginAiProvider.decode(payload),
VendorId::Plume => PlumeOpenSyncProvider.decode(payload),
VendorId::Mist => MistProvider.decode(payload),
VendorId::Netgear => NetgearInsightProvider.decode(payload),
VendorId::ElectricImp => ElectricImpProvider.decode(payload),
VendorId::RfSolutions => RfSolutionsProvider.decode(payload),
VendorId::Linksys => LinksysProvider.decode(payload),
VendorId::Luma => LumaOpenWrtProvider.decode(payload),
VendorId::GoogleNest => GoogleNestProvider.decode(payload),
VendorId::Wifigarden => WifigardenProvider.decode(payload),
}
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct VendorEventSnapshot {
pub event_type: &'static str,
pub source: String,
#[serde(flatten)]
pub event: VendorRfEvent,
}
impl VendorEventSnapshot {
pub fn from_event(event: VendorRfEvent) -> Result<Self, VendorEventError> {
let descriptor = descriptor_for(event.vendor);
event.validate(&descriptor)?;
let provenance = if event.synthetic { "simulated" } else { "live" };
Ok(Self {
event_type: "vendor_rf",
source: format!("vendor:{}:{provenance}", event.vendor.as_str()),
event,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn registry_has_exactly_one_valid_descriptor_per_vendor() {
let values = descriptors();
assert_eq!(values.len(), VendorId::ALL.len());
for (vendor, descriptor) in VendorId::ALL.into_iter().zip(values) {
assert_eq!(descriptor.vendor, vendor);
descriptor.validate().unwrap();
}
}
#[test]
fn unsupported_providers_fail_closed() {
assert_eq!(
decode_provider(VendorId::Linksys, b"{}"),
Err(VendorEventError::Unsupported)
);
assert_eq!(
decode_provider(VendorId::Wifigarden, b"{}"),
Err(VendorEventError::ContractRequired)
);
}
}