From 9fb5af7cf268b685b06316c18595cc23e85007b5 Mon Sep 17 00:00:00 2001 From: ruv Date: Tue, 28 Jul 2026 22:42:18 -0400 Subject: [PATCH] feat(sensing-server): add secure opt-in OTLP log export Import and harden the OpenTelemetry logging work from #1382. Preserve default stderr behavior, register and validate RuView semantic conventions, attach a published schema, enable TLS roots, pin demo images, and fix first-CSI node lifecycle reporting. Supersedes #1382 Closes #1460 Co-authored-by: Jens Holdgaard Pedersen --- .github/workflows/semconv.yml | 69 +++++++ docker/Dockerfile.rust | 7 +- docker/otel-collector.yaml | 26 +++ docker/otel-compose.yml | 67 +++++++ docs/observability.md | 117 ++++++++++++ semconv/registry/attributes.yaml | 106 +++++++++++ semconv/registry/events.yaml | 106 +++++++++++ semconv/registry/manifest.yaml | 6 + semconv/schema/ruview-0.1.0.yaml | 11 ++ templates/registry/rust/semconv.rs.j2 | 82 +++++++++ templates/registry/rust/weaver.yaml | 10 ++ v2/Cargo.lock | 168 ++++++++++++++++++ .../wifi-densepose-sensing-server/Cargo.toml | 14 ++ .../wifi-densepose-sensing-server/src/lib.rs | 2 + .../wifi-densepose-sensing-server/src/main.rs | 154 +++++++++++++--- .../src/mqtt/publisher.rs | 30 +++- .../src/semconv.rs | 146 +++++++++++++++ .../src/telemetry.rs | 155 ++++++++++++++++ 18 files changed, 1250 insertions(+), 26 deletions(-) create mode 100644 .github/workflows/semconv.yml create mode 100644 docker/otel-collector.yaml create mode 100644 docker/otel-compose.yml create mode 100644 docs/observability.md create mode 100644 semconv/registry/attributes.yaml create mode 100644 semconv/registry/events.yaml create mode 100644 semconv/registry/manifest.yaml create mode 100644 semconv/schema/ruview-0.1.0.yaml create mode 100644 templates/registry/rust/semconv.rs.j2 create mode 100644 templates/registry/rust/weaver.yaml create mode 100644 v2/crates/wifi-densepose-sensing-server/src/semconv.rs create mode 100644 v2/crates/wifi-densepose-sensing-server/src/telemetry.rs diff --git a/.github/workflows/semconv.yml b/.github/workflows/semconv.yml new file mode 100644 index 00000000..880b8490 --- /dev/null +++ b/.github/workflows/semconv.yml @@ -0,0 +1,69 @@ +# Semantic-conventions gate: validates `semconv/registry/` with OpenTelemetry +# weaver and verifies the generated constants module +# (`v2/crates/wifi-densepose-sensing-server/src/semconv.rs`) is in sync with +# it (`weaver registry generate` + a no-diff check) — keeping RuView's +# telemetry names spec-adherent and drift-free. +name: semconv + +on: + push: + branches: [ main, develop ] + paths: + - 'semconv/**' + - 'templates/**' + - 'v2/crates/wifi-densepose-sensing-server/src/semconv.rs' + - '.github/workflows/semconv.yml' + pull_request: + paths: + - 'semconv/**' + - 'templates/**' + - 'v2/crates/wifi-densepose-sensing-server/src/semconv.rs' + - '.github/workflows/semconv.yml' + workflow_dispatch: + +jobs: + semconv: + name: semconv (weaver) + runs-on: ubuntu-latest + env: + WEAVER_VERSION: v0.23.0 + # sha256 of weaver-x86_64-unknown-linux-gnu.tar.xz for WEAVER_VERSION + # (open-telemetry/weaver release asset). Bump both together. + WEAVER_SHA256: a9822c712d6871bd89d6530f18c5df5cea3821f642e7b8e5e49e985917f7d12d + steps: + - name: Checkout code + uses: actions/checkout@v4 + with: + persist-credentials: false + - uses: dtolnay/rust-toolchain@stable + with: + components: rustfmt + - name: Install weaver + run: | + set -euo pipefail + tarball="weaver-x86_64-unknown-linux-gnu.tar.xz" + curl -fsSL -o "$RUNNER_TEMP/$tarball" \ + "https://github.com/open-telemetry/weaver/releases/download/${WEAVER_VERSION}/${tarball}" + echo "${WEAVER_SHA256} $RUNNER_TEMP/$tarball" | sha256sum -c - + tar xJf "$RUNNER_TEMP/$tarball" -C "$RUNNER_TEMP" + echo "$RUNNER_TEMP/weaver-x86_64-unknown-linux-gnu" >> "$GITHUB_PATH" + - run: weaver registry check -r semconv/registry --future + # Codegen no-diff: regenerate the semconv constants module from the + # registry and fail if the checked-in file drifts (the generated + # module is "do not hand-edit"; the registry is the source). + - name: Regenerate semconv constants + run: | + set -euo pipefail + weaver registry generate rust v2/crates/wifi-densepose-sensing-server/src \ + -t templates -r semconv/registry --future + rustfmt --edition 2021 v2/crates/wifi-densepose-sensing-server/src/semconv.rs + - name: Verify generated constants are in sync + run: | + set -euo pipefail + changes="$(git status --porcelain -- v2/crates/wifi-densepose-sensing-server/src/semconv.rs)" + if [ -n "$changes" ]; then + echo "::error::semconv.rs is out of sync with semconv/registry/. Regenerate (see the module header) and commit." + echo "$changes" + git diff -- v2/crates/wifi-densepose-sensing-server/src/semconv.rs + exit 1 + fi diff --git a/docker/Dockerfile.rust b/docker/Dockerfile.rust index 14c02efc..1332a0ae 100644 --- a/docker/Dockerfile.rust +++ b/docker/Dockerfile.rust @@ -29,7 +29,12 @@ COPY vendor/rufield/ /vendor/rufield/ # - homecore-server, the ADRs-126-134 HOMECORE native Rust port of # Home Assistant (HA-wire-compat REST + WebSocket on :8123, # SQLite + ruvector recorder, automation, assist, plugins, HAP) -RUN cargo build --release -p wifi-densepose-sensing-server --features mqtt 2>&1 \ +# +# SENSING_FEATURES lets a compose file extend the sensing-server feature +# set (docker/otel-compose.yml builds with `mqtt,otel` for OTLP log +# export) without forking this Dockerfile. +ARG SENSING_FEATURES=mqtt +RUN cargo build --release -p wifi-densepose-sensing-server --features "${SENSING_FEATURES}" 2>&1 \ && cargo build --release -p cog-ha-matter 2>&1 \ && cargo build --release -p homecore-server 2>&1 \ && strip target/release/sensing-server target/release/cog-ha-matter target/release/homecore-server diff --git a/docker/otel-collector.yaml b/docker/otel-collector.yaml new file mode 100644 index 00000000..61a64e9f --- /dev/null +++ b/docker/otel-collector.yaml @@ -0,0 +1,26 @@ +# OpenTelemetry Collector config for the RuView observability stack +# (docker/otel-compose.yml): receive OTLP from the sensing server, export +# OTLP to the Ourios log backend. See docs/observability.md. +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:4317 + http: + endpoint: 0.0.0.0:4318 + +processors: + batch: {} + +exporters: + otlp/ourios: + endpoint: ourios:4317 + tls: + insecure: true + +service: + pipelines: + logs: + receivers: [otlp] + processors: [batch] + exporters: [otlp/ourios] diff --git a/docker/otel-compose.yml b/docker/otel-compose.yml new file mode 100644 index 00000000..cd8fe388 --- /dev/null +++ b/docker/otel-compose.yml @@ -0,0 +1,67 @@ +# RuView → OpenTelemetry Collector → Ourios log backend. +# +# docker compose -f docker/otel-compose.yml up +# +# Brings up an OTLP pipeline for the sensing server's logs: the server +# (built with `--features otel` and pointed at the collector via +# OTEL_EXPORTER_OTLP_ENDPOINT) exports every tracing event as an OTel +# log record; the collector forwards them to Ourios, a Parquet + +# template-mining log backend that is OTLP-native on ingest. Query the +# logs at http://localhost:4319/v1/query — see docs/observability.md. +services: + sensing-server: + build: + context: .. + dockerfile: docker/Dockerfile.rust + args: + # The otel feature compiles the OTLP exporter in; export still + # only activates when OTEL_EXPORTER_OTLP_ENDPOINT is set. + SENSING_FEATURES: mqtt,otel + image: ruvnet/wifi-densepose:otel + ports: + - "3000:3000" # REST API + - "3001:3001" # WebSocket + - "5005:5005/udp" # ESP32 CSI (see docker-compose.yml for Windows notes) + environment: + - RUST_LOG=info + # Demo default: synthetic CSI so the pipeline produces events with + # no hardware attached. Set CSI_SOURCE=esp32 for live nodes. + - CSI_SOURCE=${CSI_SOURCE:-simulated} + - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4317 + depends_on: + - otel-collector + + otel-collector: + image: otel/opentelemetry-collector-contrib:0.116.0@sha256:70217a89d27c678ead44f196d80aa8c2717cb68d0301dbdc40331dbec0a3e605 + command: ["--config=/etc/otelcol-contrib/config.yaml"] + volumes: + - ./otel-collector.yaml:/etc/otelcol-contrib/config.yaml:ro + ports: + - "4317:4317" # OTLP gRPC (also reachable from the host) + - "4318:4318" # OTLP HTTP + depends_on: + - ourios + + # Ourios — OTLP-native log backend (Parquet + Drain-derived template + # mining + DataFusion). Local-disk storage; the tenant derives from the + # exported resource's service.name, so RuView's logs land in tenant + # "ruview". + ourios: + image: ghcr.io/jensholdgaard/ourios:0.4.0@sha256:9c88badb2089fe78dcdef317f28babba1cdd23984409439d4c4792f64a737ef0 + environment: + - OURIOS_BUCKET_ROOT=/data + - OURIOS_WAL_ROOT=/wal + - OURIOS_RECEIVER_ENABLED=1 + - OURIOS_RECEIVER_GRPC_ADDR=0.0.0.0:4317 + - OURIOS_RECEIVER_HTTP_ADDR=0.0.0.0:4318 + - OURIOS_QUERIER_ENABLED=1 + - OURIOS_QUERIER_HTTP_ADDR=0.0.0.0:4319 + ports: + - "4319:4319" # query endpoint (http://localhost:4319/v1/query) + volumes: + - ourios-data:/data + - ourios-wal:/wal + +volumes: + ourios-data: + ourios-wal: diff --git a/docs/observability.md b/docs/observability.md new file mode 100644 index 00000000..52b01bc3 --- /dev/null +++ b/docs/observability.md @@ -0,0 +1,117 @@ +# Observability: OTLP log export + +The sensing server can export every `tracing` log event as an +OpenTelemetry log record over OTLP, with a curated set of sensing events +(presence transitions, vitals estimates, node online/offline, fall +detections, CSI capture stats, MQTT errors, model loads) carrying +registry-backed event names and attributes under the `ruview.*` +namespace. + +## The event registry + +The names are not ad hoc: they are defined in a weaver-validated +semantic-conventions registry at `semconv/registry/` (attributes and log +event names, OpenTelemetry registry format). The Rust constants module +`v2/crates/wifi-densepose-sensing-server/src/semconv.rs` is **generated** +from that registry (`weaver registry generate`, template under +`templates/registry/rust/`) and CI (`.github/workflows/semconv.yml`) +fails if either the registry stops validating or the generated module +drifts. Executed Rust tests additionally reject any hard-coded +`ruview.*` instrumentation key that is absent from the generated registry. +Exported resources carry the registry's schema URL so downstream consumers +can identify the exact conventions version. + +Curated events: + +| Event | Emitted when | +| --- | --- | +| `ruview.node.online` | first frame from a sensing node (CSI or edge vitals) | +| `ruview.node.offline` | node evicted after 60 s without frames | +| `ruview.presence.changed` | smoothed presence classification flips (transition-only) | +| `ruview.vitals.estimate` | periodic breathing / heart-rate estimate (every 100 ticks) | +| `ruview.fall.detected` | edge-vitals fall flag rising edge, per node | +| `ruview.csi.stats` | periodic capture snapshot: frames processed, active nodes | +| `ruview.mqtt.error` | MQTT publish/connection error in the HA publisher | +| `ruview.model.loaded` | inference model loaded via the model API | + +## Enabling export + +Export is doubly gated so the default build and the default runtime are +both unaffected: + +1. **Build** with the `otel` cargo feature (compiles in the OTLP + exporter stack, same gating principle as `mqtt`): + + ```sh + cargo build --release -p wifi-densepose-sensing-server --features mqtt,otel + ``` + +2. **Run** with `OTEL_EXPORTER_OTLP_ENDPOINT` set (unset ⇒ the OTLP + pipeline is never constructed and logging behaves exactly as before): + + ```sh + OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \ + ./target/release/sensing-server --source simulated + ``` + +Use an `https://` collector endpoint outside a trusted local network. The +`otel` feature includes Rustls and native certificate roots; standard OTLP +environment variables can supply authentication headers. The Compose example +uses plaintext only for container-to-container traffic on its private network. + +Logs export with resource attribute `service.name = "ruview"` and schema URL +`https://raw.githubusercontent.com/ruvnet/RuView/main/semconv/schema/ruview-0.1.0.yaml`. +Curated sensing +events are emitted only after the configured exporter initializes +successfully; without it, the pre-existing stderr output is unchanged. + +## Full stack: `docker compose` + +`docker/otel-compose.yml` brings up the whole pipeline — +sensing server (synthetic CSI by default) → OpenTelemetry Collector → +[Ourios](https://github.com/jensholdgaard/ourios), an OTLP-native log +backend built on Parquet + online log-template mining + DataFusion: + +```sh +docker compose -f docker/otel-compose.yml up +``` + +The collector and backend image tags are pinned to immutable multi-platform +digests so the demo resolves to the reviewed images. + +Ourios derives the tenant from `service.name`, so all RuView logs land +in tenant `ruview`. + +## Example queries + +Ourios mines every log line into a stable `template_id` online at +ingest, which makes template-level questions cheap. Its query endpoint +speaks a small logs DSL: + +Which log templates dominate RuView's output? + +```sh +curl -s http://localhost:4319/v1/query \ + -H 'X-Ourios-Tenant: ruview' \ + -H 'Content-Type: text/plain' \ + -d 'severity >= trace | range(-1h, now) | count by template_id | sort count desc | limit 10' +``` + +Recent warnings and errors (fall detections, MQTT failures): + +```sh +curl -s http://localhost:4319/v1/query \ + -H 'X-Ourios-Tenant: ruview' \ + -H 'Content-Type: text/plain' \ + -d 'severity >= warn | limit 50' +``` + +Did a RuView deploy change what the service logs? Template drift between +two time windows (new / vanished / changed templates): + +```sh +curl -s http://localhost:4319/v1/query \ + -H 'X-Ourios-Tenant: ruview' \ + -H 'Content-Type: text/plain' \ + -d 'drift from -7d to now' +``` diff --git a/semconv/registry/attributes.yaml b/semconv/registry/attributes.yaml new file mode 100644 index 00000000..ffb12370 --- /dev/null +++ b/semconv/registry/attributes.yaml @@ -0,0 +1,106 @@ +groups: + - id: registry.ruview + type: attribute_group + display_name: RuView attributes + brief: Attributes shared across RuView sensing telemetry. + attributes: + - id: ruview.node.id + type: int + stability: development + brief: The ESP32 mesh node the event pertains to. + note: >- + The one-byte node id carried in the ESP32 CSI / edge-vitals + frame header. Simulated frames use node id 1. + examples: [1, 2] + - id: ruview.csi.source + type: string + stability: development + brief: The data source that produced the sensing cycle. + note: >- + One of the sensing server's source labels: `esp32` (live CSI or + edge-vitals frames over UDP), `wifi` (host WiFi RSSI scanning), + or `simulated` (the built-in synthetic frame generator). + examples: ["esp32", "wifi", "simulated"] + - id: ruview.csi.frames_total + type: int + stability: development + brief: Sensing frames processed since process start (the server's tick counter). + examples: [100, 42000] + - id: ruview.csi.nodes_active + type: int + stability: development + brief: Nodes that delivered a frame within the liveness window. + examples: [0, 3] + - id: ruview.presence.state + type: + members: + - id: present + value: "present" + stability: development + brief: The classifier reports at least one person present. + - id: absent + value: "absent" + stability: development + brief: The classifier reports the space as empty. + stability: development + brief: The presence classification after smoothing and any adaptive-model override. + - id: ruview.motion.level + type: + members: + - id: absent + value: "absent" + stability: development + brief: No presence detected. + - id: present_still + value: "present_still" + stability: development + brief: Presence with little or no motion. + - id: present_moving + value: "present_moving" + stability: development + brief: Presence with moderate motion. + - id: active + value: "active" + stability: development + brief: Presence with high motion energy. + stability: development + brief: >- + The motion-level class attached to a sensing update (the + adaptive classifier's class set). + - id: ruview.inference.confidence + type: double + stability: development + brief: Confidence of the presence/motion classification, in [0.0, 1.0]. + examples: [0.7, 0.95] + - id: ruview.persons.count + type: int + stability: development + brief: The estimated person count for the sensing cycle. + examples: [0, 2] + - id: ruview.vitals.breathing_rate_bpm + type: double + stability: development + brief: Estimated breathing rate in breaths per minute. + note: "`0.0` when no estimate was produced in this window — check the confidence attribute." + examples: [14.5] + - id: ruview.vitals.heart_rate_bpm + type: double + stability: development + brief: Estimated heart rate in beats per minute. + note: "`0.0` when no estimate was produced in this window — check the confidence attribute." + examples: [62.0] + - id: ruview.vitals.breathing_confidence + type: double + stability: development + brief: Confidence of the breathing-rate estimate, in [0.0, 1.0]. + examples: [0.7] + - id: ruview.vitals.heartbeat_confidence + type: double + stability: development + brief: Confidence of the heart-rate estimate, in [0.0, 1.0]. + examples: [0.7] + - id: ruview.model.id + type: string + stability: development + brief: The identifier of a loaded inference model. + examples: ["wifi-densepose-v1"] diff --git a/semconv/registry/events.yaml b/semconv/registry/events.yaml new file mode 100644 index 00000000..c93a2ba2 --- /dev/null +++ b/semconv/registry/events.yaml @@ -0,0 +1,106 @@ +groups: + # Log event names for the sensing server's curated telemetry — each + # instrumented `tracing` call site carries one of these as its explicit + # event name (never tracing's default `event :`), so the + # exported Logs signal stays registry-backed. Uncurated log lines keep + # their default names; only these events are part of the contract. + - id: event.ruview.node.online + type: event + name: ruview.node.online + stability: development + brief: > + A sensing node delivered its first frame (CSI or edge-vitals) — + either a new node joining the mesh or a previously evicted node + returning. + attributes: + - ref: ruview.node.id + requirement_level: required + - id: event.ruview.node.offline + type: event + name: ruview.node.offline + stability: development + brief: > + A sensing node was evicted after delivering no frames for the + staleness window (60 s). + attributes: + - ref: ruview.node.id + requirement_level: required + - id: event.ruview.csi.stats + type: event + name: ruview.csi.stats + stability: development + brief: > + Periodic CSI capture snapshot (every 100 sensing ticks): total + frames processed and currently active nodes. + attributes: + - ref: ruview.csi.frames_total + requirement_level: required + - ref: ruview.csi.nodes_active + requirement_level: required + - ref: ruview.csi.source + requirement_level: recommended + - id: event.ruview.presence.changed + type: event + name: ruview.presence.changed + stability: development + brief: > + The smoothed presence classification flipped between present and + absent. Emitted on transitions only, never per frame. + attributes: + - ref: ruview.presence.state + requirement_level: required + - ref: ruview.motion.level + requirement_level: recommended + - ref: ruview.inference.confidence + requirement_level: recommended + - ref: ruview.persons.count + requirement_level: recommended + - ref: ruview.csi.source + requirement_level: recommended + - id: event.ruview.vitals.estimate + type: event + name: ruview.vitals.estimate + stability: development + brief: > + Periodic vital-sign estimate (breathing / heart rate with + confidences), emitted on the ruview.csi.stats cadence when the + detector produced an estimate. + attributes: + - ref: ruview.vitals.breathing_rate_bpm + requirement_level: recommended + - ref: ruview.vitals.heart_rate_bpm + requirement_level: recommended + - ref: ruview.vitals.breathing_confidence + requirement_level: recommended + - ref: ruview.vitals.heartbeat_confidence + requirement_level: recommended + - ref: ruview.csi.source + requirement_level: recommended + - id: event.ruview.fall.detected + type: event + name: ruview.fall.detected + stability: development + brief: > + An ESP32 edge-vitals frame raised its fall flag. Edge-triggered on + the flag's rising edge per node, not re-emitted while it stays set. + attributes: + - ref: ruview.node.id + requirement_level: required + - id: event.ruview.mqtt.error + type: event + name: ruview.mqtt.error + stability: development + brief: > + An MQTT publish or connection error in the Home Assistant + discovery publisher; the publisher reconnects and retries. + attributes: + - ref: ruview.node.id + requirement_level: opt_in + - id: event.ruview.model.loaded + type: event + name: ruview.model.loaded + stability: development + brief: An inference model was loaded via the model-management API. + attributes: + - ref: ruview.model.id + requirement_level: required diff --git a/semconv/registry/manifest.yaml b/semconv/registry/manifest.yaml new file mode 100644 index 00000000..b00f8805 --- /dev/null +++ b/semconv/registry/manifest.yaml @@ -0,0 +1,6 @@ +name: ruview +description: RuView custom OpenTelemetry semantic conventions. +schema_url: https://raw.githubusercontent.com/ruvnet/RuView/main/semconv/schema/ruview-0.1.0.yaml +dependencies: + - name: otel + registry_path: https://github.com/open-telemetry/semantic-conventions/archive/refs/tags/v1.42.0.zip[model] diff --git a/semconv/schema/ruview-0.1.0.yaml b/semconv/schema/ruview-0.1.0.yaml new file mode 100644 index 00000000..23e863b6 --- /dev/null +++ b/semconv/schema/ruview-0.1.0.yaml @@ -0,0 +1,11 @@ +# OpenTelemetry schema file format version. This is independent of the RuView +# semantic convention version below. +file_format: 1.1.0 + +# The canonical URL where this schema file is published. +schema_url: https://raw.githubusercontent.com/ruvnet/RuView/main/semconv/schema/ruview-0.1.0.yaml + +versions: + # Initial RuView semantic convention release. There are no prior versions to + # transform from. + 0.1.0: diff --git a/templates/registry/rust/semconv.rs.j2 b/templates/registry/rust/semconv.rs.j2 new file mode 100644 index 00000000..cbee688a --- /dev/null +++ b/templates/registry/rust/semconv.rs.j2 @@ -0,0 +1,82 @@ +//! Generated OpenTelemetry semantic-convention name constants for +//! RuView's curated telemetry (event names and attribute keys). +//! +//! GENERATED from `semconv/registry/` by `weaver registry generate`. +//! Do not edit by hand: change the registry or the template at +//! `templates/registry/rust/`, then regenerate (the exact command CI +//! runs — note `--future`, matching `weaver registry check --future`) +//! from the repository root and commit the result: +//! +//! ```text +//! weaver registry generate rust v2/crates/wifi-densepose-sensing-server/src \ +//! -t templates -r semconv/registry --future +//! cargo fmt -p wifi-densepose-sensing-server +//! ``` +//! +//! The CI `semconv` workflow fails if this file drifts from the registry. + +/// The semantic-conventions schema URL these constants were generated from — +/// the registry manifest's `schema_url`, which carries the conventions +/// version. Attach it to a telemetry resource so consumers can resolve the +/// schema. +pub const SCHEMA_URL: &str = "{{ (ctx.groups | first).lineage.provenance.schema_url }}"; + +// Attribute keys. +{% for group in ctx.groups | selectattr("type", "equalto", "attribute_group") %} +{% for attr in group.attributes %} +/// `{{ attr.name }}` attribute key. +pub const {{ attr.name | screaming_snake_case }}: &str = "{{ attr.name }}"; +{% endfor %} +{% endfor %} + +/// Every attribute key registered for curated RuView events. +pub const ATTRIBUTE_KEYS: &[&str] = &[ +{% for group in ctx.groups | selectattr("type", "equalto", "attribute_group") %} +{% for attr in group.attributes %} + {{ attr.name | screaming_snake_case }}, +{% endfor %} +{% endfor %} +]; + +// Log event names (each instrumented `tracing` call site names its event +// with one of these so the exported Logs signal stays registry-backed). +{% for group in ctx.groups | selectattr("type", "equalto", "event") %} +/// `{{ group.name }}` log event name. +pub const EVENT_{{ group.name | screaming_snake_case }}: &str = "{{ group.name }}"; +{% endfor %} + +/// Every curated event name in the generated registry. +pub const EVENT_NAMES: &[&str] = &[ +{% for group in ctx.groups | selectattr("type", "equalto", "event") %} + EVENT_{{ group.name | screaming_snake_case }}, +{% endfor %} +]; + +#[cfg(test)] +mod tests { + use super::{ATTRIBUTE_KEYS, EVENT_NAMES}; + + #[test] + fn instrumentation_uses_only_registered_ruview_literals() { + let sources = [ + include_str!("main.rs"), + include_str!("mqtt/publisher.rs"), + ]; + + for source in sources { + let mut rest = source; + while let Some(start) = rest.find("\"ruview.") { + let value = &rest[start + 1..]; + let end = value + .find('"') + .expect("ruview string literal must have a closing quote"); + let literal = &value[..end]; + assert!( + ATTRIBUTE_KEYS.contains(&literal) || EVENT_NAMES.contains(&literal), + "instrumentation literal `{literal}` is absent from semconv/registry" + ); + rest = &value[end + 1..]; + } + } + } +} diff --git a/templates/registry/rust/weaver.yaml b/templates/registry/rust/weaver.yaml new file mode 100644 index 00000000..e688fa2e --- /dev/null +++ b/templates/registry/rust/weaver.yaml @@ -0,0 +1,10 @@ +# Weaver-forge config for the generated `semconv` module of the +# sensing server. +# `weaver registry generate rust v2/crates/wifi-densepose-sensing-server/src \ +# -t templates -r semconv/registry --future` +# renders the single template below over the whole resolved registry +# (`--future` matches the `weaver registry check --future` validation). +templates: + - pattern: semconv.rs.j2 + filter: . + application_mode: single diff --git a/v2/Cargo.lock b/v2/Cargo.lock index cf3e9851..3657d37e 100644 --- a/v2/Cargo.lock +++ b/v2/Cargo.lock @@ -3976,6 +3976,19 @@ dependencies = [ "webpki-roots 1.0.7", ] +[[package]] +name = "hyper-timeout" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" +dependencies = [ + "hyper 1.8.1", + "hyper-util", + "pin-project-lite", + "tokio", + "tower-service", +] + [[package]] name = "hyper-tls" version = "0.6.0" @@ -5752,6 +5765,73 @@ dependencies = [ "vcpkg", ] +[[package]] +name = "opentelemetry" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0142c63252a9e054e68a4c61a5778f7b14f576274d593f8ce883d191a099682" +dependencies = [ + "js-sys", +] + +[[package]] +name = "opentelemetry-appender-tracing" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" +dependencies = [ + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", +] + +[[package]] +name = "opentelemetry-otlp" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" +dependencies = [ + "http 1.4.0", + "opentelemetry", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost", + "thiserror 2.0.18", + "tokio", + "tonic", + "tonic-types", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost", + "tonic", + "tonic-prost", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b59f80e1ac4d5ff7a2db8fb6c80badb7f0f3f858211fba08dd9aaec750894f9" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "opentelemetry", + "portable-atomic", + "thiserror 2.0.18", + "tokio", + "tokio-stream", +] + [[package]] name = "option-ext" version = "0.2.0" @@ -6502,6 +6582,38 @@ dependencies = [ "unarray", ] +[[package]] +name = "prost" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1" +dependencies = [ + "bytes", + "prost-derive", +] + +[[package]] +name = "prost-derive" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" +dependencies = [ + "anyhow", + "itertools 0.14.0", + "proc-macro2", + "quote", + "syn 2.0.117", +] + +[[package]] +name = "prost-types" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f94967dc7688f3054c7fac87473ffae4cc4c3904800e2d9f5b857246d8963b0a" +dependencies = [ + "prost", +] + [[package]] name = "ptr_meta" version = "0.3.1" @@ -9881,6 +9993,56 @@ version = "1.1.1+spec-1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "756daf9b1013ebe47a8776667b466417e2d4c5679d441c26230efd9ef78692db" +[[package]] +name = "tonic" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef" +dependencies = [ + "async-trait", + "base64 0.22.1", + "bytes", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "hyper 1.8.1", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "rustls-native-certs 0.8.3", + "sync_wrapper 1.0.2", + "tokio", + "tokio-rustls 0.26.4", + "tokio-stream", + "tower 0.5.3", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "tonic-prost" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "50849f68853be452acf590cde0b146665b8d507b3b8af17261df47e02c209ea0" +dependencies = [ + "bytes", + "prost", + "tonic", +] + +[[package]] +name = "tonic-types" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73ab1b02061f83d519bba3caa167f88f261ef05720ab8ebc954ade70de3348e8" +dependencies = [ + "prost", + "prost-types", + "tonic", +] + [[package]] name = "torch-sys" version = "0.24.0" @@ -9922,9 +10084,12 @@ checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" dependencies = [ "futures-core", "futures-util", + "indexmap 2.13.0", "pin-project-lite", + "slab", "sync_wrapper 1.0.2", "tokio", + "tokio-util", "tower-layer", "tower-service", "tracing", @@ -11476,6 +11641,9 @@ dependencies = [ "jsonwebtoken", "midstreamer-attractor", "midstreamer-temporal-compare", + "opentelemetry-appender-tracing", + "opentelemetry-otlp", + "opentelemetry_sdk", "p256", "proptest", "rand 0.8.5", diff --git a/v2/crates/wifi-densepose-sensing-server/Cargo.toml b/v2/crates/wifi-densepose-sensing-server/Cargo.toml index bf752c46..b0f62cc0 100644 --- a/v2/crates/wifi-densepose-sensing-server/Cargo.toml +++ b/v2/crates/wifi-densepose-sensing-server/Cargo.toml @@ -110,8 +110,22 @@ rand = "0.8" # uses rustls). rumqttc = { version = "0.24", default-features = false, features = ["use-rustls"], optional = true } +# `otel` feature — OTLP log export (`telemetry` module). Same gating +# principle as `mqtt`: the heavy exporter stack (opentelemetry SDK + +# tonic) stays out of the default binary; with the feature built, +# export still only activates when OTEL_EXPORTER_OTLP_ENDPOINT is set. +# Curated event names / attribute keys live in `src/semconv.rs`, +# generated from the repo-root `semconv/registry/` by weaver. +opentelemetry_sdk = { version = "0.32", default-features = false, features = ["logs", "rt-tokio"], optional = true } +opentelemetry-otlp = { version = "0.32", default-features = false, features = ["logs", "grpc-tonic", "tls", "tls-roots"], optional = true } +opentelemetry-appender-tracing = { version = "0.32", default-features = false, optional = true } + [features] default = [] +# Enables OTLP log export from the `telemetry` module (dogfooding into an +# OTLP-native log backend). Without this feature the module falls back to +# the plain stderr fmt subscriber. +otel = ["dep:opentelemetry_sdk", "dep:opentelemetry-otlp", "dep:opentelemetry-appender-tracing"] # Enables the ADR-115 §2 MQTT auto-discovery publisher. Without this feature # all `--mqtt-*` CLI flags still parse (cli.rs declares them unconditionally), # but enabling `--mqtt` at runtime logs a `WARN` and the publisher is a no-op. diff --git a/v2/crates/wifi-densepose-sensing-server/src/lib.rs b/v2/crates/wifi-densepose-sensing-server/src/lib.rs index 7779f79a..2bcb849d 100644 --- a/v2/crates/wifi-densepose-sensing-server/src/lib.rs +++ b/v2/crates/wifi-densepose-sensing-server/src/lib.rs @@ -28,6 +28,8 @@ pub mod semantic; pub mod rufield_surface; pub mod rvf_container; pub mod rvf_pipeline; +pub mod semconv; +pub mod telemetry; #[allow(dead_code)] pub mod trainer; pub mod vital_signs; diff --git a/v2/crates/wifi-densepose-sensing-server/src/main.rs b/v2/crates/wifi-densepose-sensing-server/src/main.rs index 31b81acf..3f2332c2 100644 --- a/v2/crates/wifi-densepose-sensing-server/src/main.rs +++ b/v2/crates/wifi-densepose-sensing-server/src/main.rs @@ -35,7 +35,8 @@ mod vital_signs; // Training pipeline modules (exposed via lib.rs) use wifi_densepose_sensing_server::{ - dataset, embedding, error_response, graph_transformer, rufield_surface, trainer, + dataset, embedding, error_response, graph_transformer, rufield_surface, semconv, telemetry, + trainer, }; use ruvector_mincut::{DynamicMinCut, MinCutBuilder}; @@ -759,7 +760,11 @@ impl NodeState { }) } - pub(crate) fn observe_csi_frame_arrival(&mut self, now: std::time::Instant) { + /// Record a CSI data-frame arrival and return whether it was the node's + /// first sensing frame. Sync packets deliberately do not change this + /// result, so a sync-before-CSI sequence still produces `node.online`. + pub(crate) fn observe_csi_frame_arrival(&mut self, now: std::time::Instant) -> bool { + let first_sensing_frame = self.last_frame_time.is_none(); if let Some(prev) = self.last_frame_time { let dt = now.duration_since(prev).as_secs_f64(); // Burst arrivals (sub-floor dt, issue #1180): do NOT re-anchor on @@ -769,7 +774,7 @@ impl NodeState { // frames arrive in 36 µs bursts every 25 ms still reads ~40 fps, // not 27 kHz. if dt < MIN_PLAUSIBLE_CSI_DT_SEC { - return; + return false; } if let Some(new_ema) = update_csi_fps_ema(self.csi_fps_ema, dt) { self.csi_fps_ema = new_ema; @@ -777,6 +782,7 @@ impl NodeState { } } self.last_frame_time = Some(now); + first_sensing_frame } pub(crate) fn new() -> Self { @@ -2841,6 +2847,7 @@ async fn windows_wifi_task(state: SharedState, tick_ms: u64) { if let Ok(json) = serde_json::to_string(&update) { let _ = s.tx.send(json); } + observe_sensing_update(s.latest_update.as_ref(), &update); s.latest_update = Some(update); debug!( @@ -2996,6 +3003,7 @@ async fn windows_wifi_fallback_tick(state: &SharedState, seq: u32) { if let Ok(json) = serde_json::to_string(&update) { let _ = s.tx.send(json); } + observe_sensing_update(s.latest_update.as_ref(), &update); s.latest_update = Some(update); } @@ -4768,7 +4776,11 @@ async fn load_model( let mut s = state.write().await; s.active_model_id = Some(model_id.clone()); s.model_loaded = true; - info!("Model loaded: {model_id}"); + if telemetry::curated_events_enabled() { + info!(name: semconv::EVENT_RUVIEW_MODEL_LOADED, { "ruview.model.id" = %model_id }, "Model loaded: {model_id}"); + } else { + info!("Model loaded: {model_id}"); + } Json(serde_json::json!({ "success": true, "model_id": model_id })) } @@ -5824,7 +5836,18 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) { // ── Per-node state for edge vitals (issue #249) ────── let node_id = vitals.node_id; let ns = s.node_states.entry(node_id).or_insert_with(NodeState::new); + let first_sensing_frame = ns.last_frame_time.is_none(); ns.last_frame_time = Some(std::time::Instant::now()); + if first_sensing_frame && telemetry::curated_events_enabled() { + info!(name: semconv::EVENT_RUVIEW_NODE_ONLINE, { "ruview.node.id" = node_id }, "node {node_id} online (edge vitals)"); + } + // Edge-triggered on the fall flag's rising edge (against + // the node's previous edge-vitals frame), so a persisting + // flag does not re-emit every frame. + let prev_fall = ns.edge_vitals.as_ref().is_some_and(|v| v.fall_detected); + if vitals.fall_detected && !prev_fall && telemetry::curated_events_enabled() { + warn!(name: semconv::EVENT_RUVIEW_FALL_DETECTED, { "ruview.node.id" = node_id }, "fall detected by node {node_id}"); + } ns.edge_vitals = Some(vitals.clone()); ns.rssi_history.push_back(vitals.rssi as f64); if ns.rssi_history.len() > 60 { @@ -6037,6 +6060,7 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) { if let Ok(json) = serde_json::to_string(&update) { let _ = s.tx.send(json); } + observe_sensing_update(s.latest_update.as_ref(), &update); s.latest_update = Some(update); s.edge_vitals = Some(vitals); continue; @@ -6216,7 +6240,11 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) { // ADR-110 iter 19 — feed the per-node fps EMA from real // CSI arrivals. The helper sets `last_frame_time` as a // side effect, so the previous bare assignment is gone. - ns.observe_csi_frame_arrival(std::time::Instant::now()); + let first_sensing_frame = + ns.observe_csi_frame_arrival(std::time::Instant::now()); + if first_sensing_frame && telemetry::curated_events_enabled() { + info!(name: semconv::EVENT_RUVIEW_NODE_ONLINE, { "ruview.node.id" = node_id }, "node {node_id} online (CSI)"); + } // ADR-084 Pass 3: cluster-Pi novelty sensor. // Score this frame's feature vector against the per-node @@ -6481,21 +6509,31 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) { // held edge-local. `presence == false` ⇒ no phantom event. emit_rufield_event(&s, &update, node_id); + observe_sensing_update(s.latest_update.as_ref(), &update); s.latest_update = Some(update); // Evict stale nodes every 100 ticks to prevent memory leak. if tick % 100 == 0 { let stale = Duration::from_secs(60); - let before = s.node_states.len(); - s.node_states.retain(|_id, ns| { - ns.last_frame_time - .is_some_and(|t| now.duration_since(t) < stale) - }); - let evicted = before - s.node_states.len(); - if evicted > 0 { + let stale_ids: Vec = s + .node_states + .iter() + .filter(|(_, ns)| { + !ns.last_frame_time + .is_some_and(|t| now.duration_since(t) < stale) + }) + .map(|(&id, _)| id) + .collect(); + for id in &stale_ids { + s.node_states.remove(id); + if telemetry::curated_events_enabled() { + info!(name: semconv::EVENT_RUVIEW_NODE_OFFLINE, { "ruview.node.id" = *id }, "node {id} offline (no frames for 60s)"); + } + } + if !stale_ids.is_empty() { info!( "Evicted {} stale node(s), {} active", - evicted, + stale_ids.len(), s.node_states.len() ); } @@ -6510,6 +6548,67 @@ async fn udp_receiver_task(state: SharedState, udp_port: u16) { } } +/// Cadence, in sensing ticks, of the periodic `ruview.csi.stats` and +/// `ruview.vitals.estimate` telemetry snapshots. +const TELEMETRY_SNAPSHOT_TICKS: u64 = 100; + +/// Emit the curated telemetry events for a finished sensing cycle +/// (names and attribute keys from `semconv/registry/` at the repo root): +/// `ruview.presence.changed` on presence transitions against the +/// previously published update, plus the cadenced `ruview.csi.stats` and +/// `ruview.vitals.estimate` snapshots. Called from every path that +/// publishes a `SensingUpdate`, right before it lands in `latest_update` +/// — never per frame at full rate. +fn observe_sensing_update(prev: Option<&SensingUpdate>, update: &SensingUpdate) { + if !telemetry::curated_events_enabled() { + return; + } + + let presence = update.classification.presence; + if prev.map(|u| u.classification.presence) != Some(presence) { + let state = if presence { "present" } else { "absent" }; + info!( + name: semconv::EVENT_RUVIEW_PRESENCE_CHANGED, + { + "ruview.presence.state" = state, + "ruview.motion.level" = %update.classification.motion_level, + "ruview.inference.confidence" = update.classification.confidence, + "ruview.persons.count" = update.estimated_persons.unwrap_or(0) as u64, + "ruview.csi.source" = %update.source, + }, + "presence changed: {state}" + ); + } + if update.tick % TELEMETRY_SNAPSHOT_TICKS == 0 { + info!( + name: semconv::EVENT_RUVIEW_CSI_STATS, + { + "ruview.csi.frames_total" = update.tick, + "ruview.csi.nodes_active" = update.nodes.len() as u64, + "ruview.csi.source" = %update.source, + }, + "csi stats: {} frames processed, {} active node(s)", + update.tick, + update.nodes.len() + ); + if let Some(v) = &update.vital_signs { + if v.breathing_rate_bpm.is_some() || v.heart_rate_bpm.is_some() { + info!( + name: semconv::EVENT_RUVIEW_VITALS_ESTIMATE, + { + "ruview.vitals.breathing_rate_bpm" = v.breathing_rate_bpm.unwrap_or(0.0), + "ruview.vitals.heart_rate_bpm" = v.heart_rate_bpm.unwrap_or(0.0), + "ruview.vitals.breathing_confidence" = v.breathing_confidence, + "ruview.vitals.heartbeat_confidence" = v.heartbeat_confidence, + "ruview.csi.source" = %update.source, + }, + "vitals estimate" + ); + } + } + } +} + // ── Simulated data task ────────────────────────────────────────────────────── async fn simulated_data_task(state: SharedState, tick_ms: u64) { @@ -6659,6 +6758,7 @@ async fn simulated_data_task(state: SharedState, tick_ms: u64) { if let Ok(json) = serde_json::to_string(&update) { let _ = s.tx.send(json); } + observe_sensing_update(s.latest_update.as_ref(), &update); s.latest_update = Some(update); } } @@ -7060,13 +7160,11 @@ fn coalesce_ui_path(initial: std::path::PathBuf) -> std::path::PathBuf { #[tokio::main] async fn main() { - // Initialize tracing - tracing_subscriber::fmt() - .with_env_filter( - tracing_subscriber::EnvFilter::try_from_default_env() - .unwrap_or_else(|_| "info,tower_http=debug".into()), - ) - .init(); + // Initialize tracing; with the `otel` feature and + // OTEL_EXPORTER_OTLP_ENDPOINT set, logs also export over OTLP + // (service.name = "ruview") — see telemetry.rs. The guard flushes + // pending log records on exit. + let _telemetry = telemetry::init(); let mut args = Args::parse(); args.ui_path = coalesce_ui_path(args.ui_path); @@ -8577,6 +8675,22 @@ mod sync_snapshot_helper_tests { assert!(ns.sync_snapshot().is_some()); } + #[test] + fn sync_before_first_csi_still_marks_csi_as_first_sensing_frame() { + let mut ns = NodeState::new(); + let now = std::time::Instant::now(); + ns.apply_sync_packet(populated_sync(9), now); + + assert!( + ns.observe_csi_frame_arrival(now + std::time::Duration::from_millis(20)), + "a sync packet must not consume the first sensing-frame transition" + ); + assert!( + !ns.observe_csi_frame_arrival(now + std::time::Duration::from_millis(40)), + "subsequent CSI frames must not re-emit node.online" + ); + } + #[test] fn apply_sync_packet_overwrites_older_data() { // Subsequent packets must replace, not accumulate. Otherwise the diff --git a/v2/crates/wifi-densepose-sensing-server/src/mqtt/publisher.rs b/v2/crates/wifi-densepose-sensing-server/src/mqtt/publisher.rs index ad9177f3..6c3e3285 100644 --- a/v2/crates/wifi-densepose-sensing-server/src/mqtt/publisher.rs +++ b/v2/crates/wifi-densepose-sensing-server/src/mqtt/publisher.rs @@ -35,6 +35,26 @@ use tokio::sync::broadcast; use tokio::task::JoinHandle; use tracing::{error, info, warn}; +macro_rules! otel_error { + ($($arg:tt)*) => { + if crate::telemetry::curated_events_enabled() { + error!(name: crate::semconv::EVENT_RUVIEW_MQTT_ERROR, $($arg)*); + } else { + error!($($arg)*); + } + }; +} + +macro_rules! otel_warn { + ($($arg:tt)*) => { + if crate::telemetry::curated_events_enabled() { + warn!(name: crate::semconv::EVENT_RUVIEW_MQTT_ERROR, $($arg)*); + } else { + warn!($($arg)*); + } + }; +} + use super::config::{MqttConfig, TlsConfig}; use super::discovery::{DiscoveryBuilder, EntityKind}; use super::state::{RateLimiter, StateEncoder, StateMessage, VitalsSnapshot}; @@ -184,7 +204,7 @@ async fn run( match ev { Ok(_) => {} Err(e) => { - error!("[mqtt] event loop error, will reconnect: {e}"); + otel_error!("[mqtt] event loop error, will reconnect: {e}"); rate_limiter.reset(); // Brief backoff before next poll attempt. tokio::time::sleep(Duration::from_millis(500)).await; @@ -197,7 +217,7 @@ async fn run( if last_heartbeat.elapsed() >= AVAILABILITY_HEARTBEAT { for (_, na) in nodes.values() { if let Err(e) = publish_availability(&client, na, "online").await { - warn!("[mqtt] heartbeat publish failed: {e}"); + otel_warn!("[mqtt] heartbeat publish failed: {e}"); } } last_heartbeat = Instant::now(); @@ -207,7 +227,7 @@ async fn run( if let Err(e) = publish_all_discovery(&client, &nb.as_borrowed(), &entities).await { - warn!("[mqtt] discovery refresh failed: {e}"); + otel_warn!("[mqtt] discovery refresh failed: {e}"); } } last_refresh = Instant::now(); @@ -228,11 +248,11 @@ async fn run( if let Err(e) = publish_all_discovery(&client, &borrowed, &entities).await { - warn!("[mqtt] node {} discovery failed: {e}", snap.node_id); + otel_warn!("[mqtt] node {} discovery failed: {e}", snap.node_id); } let na = NodeAvailability::for_builder(&borrowed, &entities); if let Err(e) = publish_availability(&client, &na, "online").await { - warn!("[mqtt] node {} availability failed: {e}", snap.node_id); + otel_warn!("[mqtt] node {} availability failed: {e}", snap.node_id); } nodes.insert(snap.node_id.clone(), (nb, na)); } diff --git a/v2/crates/wifi-densepose-sensing-server/src/semconv.rs b/v2/crates/wifi-densepose-sensing-server/src/semconv.rs new file mode 100644 index 00000000..80ddda51 --- /dev/null +++ b/v2/crates/wifi-densepose-sensing-server/src/semconv.rs @@ -0,0 +1,146 @@ +//! Generated OpenTelemetry semantic-convention name constants for +//! RuView's curated telemetry (event names and attribute keys). +//! +//! GENERATED from `semconv/registry/` by `weaver registry generate`. +//! Do not edit by hand: change the registry or the template at +//! `templates/registry/rust/`, then regenerate (the exact command CI +//! runs — note `--future`, matching `weaver registry check --future`) +//! from the repository root and commit the result: +//! +//! ```text +//! weaver registry generate rust v2/crates/wifi-densepose-sensing-server/src \ +//! -t templates -r semconv/registry --future +//! cargo fmt -p wifi-densepose-sensing-server +//! ``` +//! +//! The CI `semconv` workflow fails if this file drifts from the registry. + +/// The semantic-conventions schema URL these constants were generated from — +/// the registry manifest's `schema_url`, which carries the conventions +/// version. Attach it to a telemetry resource so consumers can resolve the +/// schema. +pub const SCHEMA_URL: &str = + "https://raw.githubusercontent.com/ruvnet/RuView/main/semconv/schema/ruview-0.1.0.yaml"; + +// Attribute keys. + +/// `ruview.csi.frames_total` attribute key. +pub const RUVIEW_CSI_FRAMES_TOTAL: &str = "ruview.csi.frames_total"; + +/// `ruview.csi.nodes_active` attribute key. +pub const RUVIEW_CSI_NODES_ACTIVE: &str = "ruview.csi.nodes_active"; + +/// `ruview.csi.source` attribute key. +pub const RUVIEW_CSI_SOURCE: &str = "ruview.csi.source"; + +/// `ruview.inference.confidence` attribute key. +pub const RUVIEW_INFERENCE_CONFIDENCE: &str = "ruview.inference.confidence"; + +/// `ruview.model.id` attribute key. +pub const RUVIEW_MODEL_ID: &str = "ruview.model.id"; + +/// `ruview.motion.level` attribute key. +pub const RUVIEW_MOTION_LEVEL: &str = "ruview.motion.level"; + +/// `ruview.node.id` attribute key. +pub const RUVIEW_NODE_ID: &str = "ruview.node.id"; + +/// `ruview.persons.count` attribute key. +pub const RUVIEW_PERSONS_COUNT: &str = "ruview.persons.count"; + +/// `ruview.presence.state` attribute key. +pub const RUVIEW_PRESENCE_STATE: &str = "ruview.presence.state"; + +/// `ruview.vitals.breathing_confidence` attribute key. +pub const RUVIEW_VITALS_BREATHING_CONFIDENCE: &str = "ruview.vitals.breathing_confidence"; + +/// `ruview.vitals.breathing_rate_bpm` attribute key. +pub const RUVIEW_VITALS_BREATHING_RATE_BPM: &str = "ruview.vitals.breathing_rate_bpm"; + +/// `ruview.vitals.heart_rate_bpm` attribute key. +pub const RUVIEW_VITALS_HEART_RATE_BPM: &str = "ruview.vitals.heart_rate_bpm"; + +/// `ruview.vitals.heartbeat_confidence` attribute key. +pub const RUVIEW_VITALS_HEARTBEAT_CONFIDENCE: &str = "ruview.vitals.heartbeat_confidence"; + +/// Every attribute key registered for curated RuView events. +pub const ATTRIBUTE_KEYS: &[&str] = &[ + RUVIEW_CSI_FRAMES_TOTAL, + RUVIEW_CSI_NODES_ACTIVE, + RUVIEW_CSI_SOURCE, + RUVIEW_INFERENCE_CONFIDENCE, + RUVIEW_MODEL_ID, + RUVIEW_MOTION_LEVEL, + RUVIEW_NODE_ID, + RUVIEW_PERSONS_COUNT, + RUVIEW_PRESENCE_STATE, + RUVIEW_VITALS_BREATHING_CONFIDENCE, + RUVIEW_VITALS_BREATHING_RATE_BPM, + RUVIEW_VITALS_HEART_RATE_BPM, + RUVIEW_VITALS_HEARTBEAT_CONFIDENCE, +]; + +// Log event names (each instrumented `tracing` call site names its event +// with one of these so the exported Logs signal stays registry-backed). + +/// `ruview.csi.stats` log event name. +pub const EVENT_RUVIEW_CSI_STATS: &str = "ruview.csi.stats"; + +/// `ruview.fall.detected` log event name. +pub const EVENT_RUVIEW_FALL_DETECTED: &str = "ruview.fall.detected"; + +/// `ruview.model.loaded` log event name. +pub const EVENT_RUVIEW_MODEL_LOADED: &str = "ruview.model.loaded"; + +/// `ruview.mqtt.error` log event name. +pub const EVENT_RUVIEW_MQTT_ERROR: &str = "ruview.mqtt.error"; + +/// `ruview.node.offline` log event name. +pub const EVENT_RUVIEW_NODE_OFFLINE: &str = "ruview.node.offline"; + +/// `ruview.node.online` log event name. +pub const EVENT_RUVIEW_NODE_ONLINE: &str = "ruview.node.online"; + +/// `ruview.presence.changed` log event name. +pub const EVENT_RUVIEW_PRESENCE_CHANGED: &str = "ruview.presence.changed"; + +/// `ruview.vitals.estimate` log event name. +pub const EVENT_RUVIEW_VITALS_ESTIMATE: &str = "ruview.vitals.estimate"; + +/// Every curated event name in the generated registry. +pub const EVENT_NAMES: &[&str] = &[ + EVENT_RUVIEW_CSI_STATS, + EVENT_RUVIEW_FALL_DETECTED, + EVENT_RUVIEW_MODEL_LOADED, + EVENT_RUVIEW_MQTT_ERROR, + EVENT_RUVIEW_NODE_OFFLINE, + EVENT_RUVIEW_NODE_ONLINE, + EVENT_RUVIEW_PRESENCE_CHANGED, + EVENT_RUVIEW_VITALS_ESTIMATE, +]; + +#[cfg(test)] +mod tests { + use super::{ATTRIBUTE_KEYS, EVENT_NAMES}; + + #[test] + fn instrumentation_uses_only_registered_ruview_literals() { + let sources = [include_str!("main.rs"), include_str!("mqtt/publisher.rs")]; + + for source in sources { + let mut rest = source; + while let Some(start) = rest.find("\"ruview.") { + let value = &rest[start + 1..]; + let end = value + .find('"') + .expect("ruview string literal must have a closing quote"); + let literal = &value[..end]; + assert!( + ATTRIBUTE_KEYS.contains(&literal) || EVENT_NAMES.contains(&literal), + "instrumentation literal `{literal}` is absent from semconv/registry" + ); + rest = &value[end + 1..]; + } + } + } +} diff --git a/v2/crates/wifi-densepose-sensing-server/src/telemetry.rs b/v2/crates/wifi-densepose-sensing-server/src/telemetry.rs new file mode 100644 index 00000000..da104d93 --- /dev/null +++ b/v2/crates/wifi-densepose-sensing-server/src/telemetry.rs @@ -0,0 +1,155 @@ +//! Tracing bootstrap, with optional OTLP log export. +//! +//! Without the `otel` cargo feature (the default) this is exactly the +//! stderr `tracing_subscriber::fmt()` setup the server has always had. +//! With the feature, and only when `OTEL_EXPORTER_OTLP_ENDPOINT` is set +//! in the environment, [`init`] additionally installs an +//! `opentelemetry-appender-tracing` bridge so every `tracing` event is +//! exported as an OpenTelemetry log record (`service.name = "ruview"`) +//! to that endpoint over OTLP/gRPC. Endpoint unset means the OTLP +//! pipeline is never constructed and no curated events are emitted. +//! +//! The curated events named in [`crate::semconv`] carry registry-backed +//! event names and attribute keys (see `semconv/registry/` at the repo +//! root); everything else exports under tracing's default event names. + +#[cfg(feature = "otel")] +use std::sync::atomic::{AtomicBool, Ordering}; + +/// Service name reported in the OTLP resource. Log backends that derive +/// a tenant from `service.name` file all RuView logs under it. +#[cfg(feature = "otel")] +const SERVICE_NAME: &str = "ruview"; + +/// True only after an explicitly configured OTLP logger has been installed. +/// +/// Curated sensing events are new output, so call sites consult this flag +/// before emitting them. This preserves the server's pre-OTel stderr behavior +/// both in default builds and in `otel` builds without a working exporter. +#[cfg(feature = "otel")] +static CURATED_EVENTS_ENABLED: AtomicBool = AtomicBool::new(false); + +/// The server's long-standing default log filter. +const DEFAULT_FILTER: &str = "info,tower_http=debug"; + +/// Whether the explicitly configured OTLP pipeline is ready to receive curated +/// `ruview.*` events. +#[cfg(feature = "otel")] +pub fn curated_events_enabled() -> bool { + CURATED_EVENTS_ENABLED.load(Ordering::Acquire) +} + +/// Default builds cannot emit curated OTLP events. +#[cfg(not(feature = "otel"))] +pub const fn curated_events_enabled() -> bool { + false +} + +/// Owns the OTLP logger pipeline when one was installed. Hold it for the +/// process lifetime; dropping it flushes pending log records. +pub struct TelemetryGuard { + #[cfg(feature = "otel")] + logger: Option, +} + +impl Drop for TelemetryGuard { + fn drop(&mut self) { + #[cfg(feature = "otel")] + if let Some(logger) = &self.logger { + CURATED_EVENTS_ENABLED.store(false, Ordering::Release); + let _ = logger.shutdown(); + } + } +} + +/// Install the global `tracing` subscriber. Call once, at start-up, +/// inside the tokio runtime (the OTLP exporter runs on it). +#[must_use = "dropping the guard tears the OTLP log pipeline down"] +pub fn init() -> TelemetryGuard { + #[cfg(feature = "otel")] + if std::env::var_os("OTEL_EXPORTER_OTLP_ENDPOINT").is_some() { + match init_with_otlp() { + Ok(logger) => { + CURATED_EVENTS_ENABLED.store(true, Ordering::Release); + return TelemetryGuard { + logger: Some(logger), + }; + } + Err(e) => eprintln!("OTLP log export disabled (exporter build failed): {e}"), + } + } + + init_fmt_only(); + TelemetryGuard { + #[cfg(feature = "otel")] + logger: None, + } +} + +fn init_fmt_only() { + tracing_subscriber::fmt() + .with_env_filter( + tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| DEFAULT_FILTER.into()), + ) + .init(); +} + +#[cfg(feature = "otel")] +fn init_with_otlp( +) -> Result { + use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; + use opentelemetry_otlp::LogExporter; + use opentelemetry_sdk::logs::SdkLoggerProvider; + use opentelemetry_sdk::Resource; + use tracing_subscriber::layer::SubscriberExt as _; + use tracing_subscriber::util::SubscriberInitExt as _; + use tracing_subscriber::{EnvFilter, Layer as _}; + + // The exporter reads OTEL_EXPORTER_OTLP_ENDPOINT (and the other + // OTEL_EXPORTER_* variables) from the environment itself. + let exporter = LogExporter::builder().with_tonic().build()?; + let resource = Resource::builder() + .with_service_name(SERVICE_NAME) + .with_schema_url([], crate::semconv::SCHEMA_URL) + .build(); + let logger = SdkLoggerProvider::builder() + .with_batch_exporter(exporter) + .with_resource(resource) + .build(); + + // Telemetry-induced-telemetry loop guard: the OTLP exporter is itself + // a tonic/hyper client, so its internal tracing events must not + // re-enter the bridge (a failed export would emit records that + // trigger more exports). The `off` directives win over RUST_LOG. + let mut bridge_filter = + EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(DEFAULT_FILTER)); + for directive in [ + "hyper=off", + "tonic=off", + "h2=off", + "tower=off", + "opentelemetry=off", + "opentelemetry_sdk=off", + "opentelemetry_otlp=off", + ] { + if let Ok(directive) = directive.parse() { + bridge_filter = bridge_filter.add_directive(directive); + } + } + let bridge = OpenTelemetryTracingBridge::new(&logger).with_filter(bridge_filter); + let fmt = tracing_subscriber::fmt::layer().with_filter( + EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(DEFAULT_FILTER)), + ); + tracing_subscriber::registry().with(bridge).with(fmt).init(); + Ok(logger) +} + +#[cfg(test)] +mod tests { + #[cfg(not(feature = "otel"))] + #[test] + fn default_build_never_enables_curated_events() { + assert!(!super::curated_events_enabled()); + } +}