Skip to content

Repository files navigation

ObserverAI: Cognitive Observability & Forensics Platform

ObserverAI is a full-stack, distributed-tracing observability platform. It ingests OpenTelemetry traces/metrics/logs from any instrumented service, streams them through a real-time anomaly-detection pipeline, and surfaces results on a live dashboard with LLM-powered root-cause analysis.

It's built to demonstrate the same building blocks used by production observability platforms (Datadog, Honeycomb, Grafana-style stacks) end-to-end, from raw OTLP ingestion to a stream processor running online ML models, to an AI agent that explains why something broke:

  • Vendor-neutral ingestion — any service that speaks OTLP (any language, any framework) is picked up automatically, with zero code changes to that service.
  • Real-time stream processing — a Bytewax dataflow reconstructs distributed traces from a RabbitMQ stream and scores them in tumbling windows.
  • Hybrid anomaly detection — rule-based detectors (N+1 queries, bimodal latency, dependency-chain breaks) combined with five different online/statistical ML models (LOF, HS-Trees, Isolation Forest, Autoencoder, One-Class SVM).
  • PII redaction at the edge — the OTel Collector strips emails, card numbers, and author fields from log bodies before they ever leave the collector, and the platform tracks a "redaction density" security metric derived from that.
  • LLM-powered RCA — one click sends a reconstructed trace to Gemini, which returns a structured root-cause diagnosis and suggested fix.

Live Demo (what you'll see)

Add a screenshot of your running dashboard here before publishing (./start.sh, run ./traffic.sh 60 mixed 5, then screenshot http://localhost:5173 into docs/screenshot-dashboard.png and reference it as ![Dashboard](docs/screenshot-dashboard.png)).

Throughput, P99 latency, and anomaly rate update in real time as traffic flows through the system; the incident stream shows every anomaly the pipeline caught, tagged with which detector fired.

Architecture

graph TD
    subgraph "Any Instrumented Service"
        AGW["API Gateway (Node.js)"]
        QS["Quote Service (Python)"]
        WS["Web Storefront (Node.js) — onboarded with zero code changes"]
    end

    subgraph "Instrumentation Wrappers"
        NW["Node Wrapper"]
        PW["Python Wrapper"]
    end

    subgraph "Data Pipeline"
        OC["OTel Collector (Docker) — PII redaction"]
        RMQ["RabbitMQ Stream"]
        BW["Bytewax Stream Processor — windowing + anomaly scoring"]
    end

    subgraph "Storage & Intelligence"
        DB["SQLite (telemetry.db)"]
        FAST["FastAPI Backend"]
        GEM["Gemini (AI RCA)"]
    end

    subgraph "Frontend"
        REACT["React Dashboard (Vite + Recharts)"]
    end

    AGW & QS & WS --> NW & PW
    NW & PW --> OC
    OC -->|Redacted OTLP| RMQ
    RMQ --> BW
    BW -->|Windowed Metrics & Traces| FAST
    FAST --> DB
    REACT -->|Live Metrics / Alerts| FAST
    REACT -->|Request RCA| FAST
    FAST -->|Normalize & Analyze| GEM
    GEM -->|Structured RCA| FAST
Loading

Forensic trace reconstruction flow

sequenceDiagram
    participant U as User Traffic
    participant S as Services
    participant C as OTel Collector
    participant B as Bytewax (Stream)
    participant BE as Backend
    participant AI as Gemini AI

    U->>S: Request (e.g. /api/proxy-slow-quote)
    S->>S: Multiple spans generated
    S->>C: Push spans (OTLP)
    C->>C: Redact PII (email / card / author)
    C->>B: Stream to RabbitMQ
    B->>B: 10s tumbling window
    B->>B: Reconstruct trace by trace_id + score anomaly
    B->>BE: POST /api/traces (full inventory)
    B->>BE: POST /api/alerts (if anomalous)
    BE->>BE: Save to telemetry.db, broadcast over WebSocket
    Note over BE: User clicks "Investigate Root Cause"
    BE->>AI: Send normalized full trace
    AI->>BE: Return structured root cause + fix
    BE->>U: Display diagnostic panel
Loading

Anomaly Detection Engine

Every completed trace is scored by a composite detector (stream-processor/detectors.py) combining:

Type Detector What it catches
Rule-based N+1 Query Regression A span fan-out (e.g. child DB calls) that regresses past a rolling mean + k·σ threshold
Rule-based Bimodal Latency A route whose latency distribution splits into a fast/slow mixture
Rule-based Dependency Chain Break A downstream span failing/timing out and breaking the causal chain
Rule-based PII Redaction Density Spike in redaction-token counts in log bodies (security signal, §IV-C)
ML (online) Statistical Outlier (LOF) Local Outlier Factor over streaming feature vectors
ML (online) Statistical Outlier (HS-Trees) Half-Space Trees (River) streaming isolation model
ML (online) Statistical Outlier (Isolation Forest) Scikit-learn isolation forest over a rolling feature window
ML (online) Reconstruction Anomaly (Autoencoder) Reconstruction error from a small online autoencoder
ML (online) Boundary Anomaly (SVM) One-Class SVM decision boundary over latency/feature space

Detectors run inside the Bytewax dataflow (stream-processor/dataflow.py) on a 10-second tumbling window, keyed by trace_id, so scoring happens on fully reconstructed cross-service traces rather than isolated spans.

Tech Stack

Layer Technology
Instrumentation OpenTelemetry auto-instrumentation (Node + Python), zero-code wrapper scripts
Collector otel/opentelemetry-collector-contrib (Docker), transform processor for PII redaction
Message bus RabbitMQ (native stream queue type, AMQP 0.9.1)
Stream processing Bytewax (Rust-backed Python dataflow engine)
ML River (online ML), scikit-learn, custom composite scorer
Backend FastAPI, aiosqlite, WebSockets, Gemini API
Frontend React 19, Vite, Recharts, Tailwind CSS
Demo services Node/Express (API Gateway, Web Storefront), FastAPI (Quote Service)

Project Structure

ObserveX/
├── dashboard/
│   ├── backend/          FastAPI server, SQLite storage, Gemini RCA
│   └── frontend/         React + Vite dashboard
├── infra/otel-collector/ Collector config + docker-compose
├── instrumentation/      Zero-code launcher scripts (node-wrapper, python-wrapper)
├── microservices/
│   ├── api-gateway/      Node.js demo service
│   ├── quote-service/    Python/FastAPI demo service
│   └── web-storefront/   Node.js demo service — proves onboarding a brand-new service works
├── stream-processor/      Bytewax dataflow, detectors, ML scorer
├── start.sh / stop.sh / status.sh / traffic.sh   Stack lifecycle scripts
└── trigger_traffic.py     Synthetic traffic generator (paper-calibrated traffic profiles)

Prerequisites

  • Docker Desktop (runs the OTel Collector)
  • RabbitMQ installed locally, with the management plugin enabled — see below for both Windows and Linux/macOS install paths
  • Python 3.10–3.12 (Bytewax does not yet ship Windows/Linux wheels for 3.13+) — this project uses its own virtualenv, so it's fine if your system Python is newer
  • Node.js 18+
  • A Gemini API key (optional — the AI root-cause feature degrades gracefully without one; everything else works fine)

Installing RabbitMQ

  • Windows: install the RabbitMQ Windows installer (it bundles Erlang and registers a RabbitMQ Windows service that starts automatically). Then enable the management plugin:
    rabbitmq-plugins enable rabbitmq_management
  • macOS: brew install rabbitmq && brew services start rabbitmq
  • Linux: sudo apt install rabbitmq-server && sudo rabbitmq-plugins enable rabbitmq_management

Then create the telemetry user the rest of the stack expects (management API, default port 15672, default admin user is guest/guest on a fresh install):

curl -u guest:guest -X PUT "http://localhost:15672/api/users/telemetry" \
  -H "content-type: application/json" -d "{\"password\":\"telemetry_password\",\"tags\":\"management\"}"
curl -u guest:guest -X PUT "http://localhost:15672/api/permissions/%2F/telemetry" \
  -H "content-type: application/json" -d "{\"configure\":\".*\",\"write\":\".*\",\"read\":\".*\"}"

Setup

  1. Clone the repo and create the Python virtualenv:

    git clone https://github.com/Swikarjadhav14/ObserveX.git && cd ObserveX
    python -m venv venv

    On Windows this creates venv\Scripts\; on Linux/macOS, venv/bin/. All scripts in this repo detect which one exists automatically.

  2. Install Python dependencies (Windows shown; use venv/bin/python on Linux/macOS):

    venv/Scripts/python.exe -m pip install --upgrade pip
    venv/Scripts/python.exe -m pip install -r stream-processor/requirements.txt -r microservices/quote-service/requirements.txt -r dashboard/backend/requirements.txt
    venv/Scripts/opentelemetry-bootstrap.exe -a install

    If your machine blocks pip's generated .exe shims (see Troubleshooting), run instead: venv/Scripts/python.exe -c "import sys; sys.argv=['opentelemetry-bootstrap','-a','install']; from opentelemetry.instrumentation.bootstrap import run; run()"

  3. Install Node dependencies for each Node service:

    (cd microservices/api-gateway && npm install)
    (cd microservices/web-storefront && npm install)
    (cd dashboard/frontend && npm install)
  4. Configure the AI RCA feature (optional). Create dashboard/backend/.env:

    GEMINI_API_KEY=your-gemini-api-key-here

    Get a free key at aistudio.google.com/apikey. Leave the placeholder in place to run everything except AI RCA.

Running the Stack

./start.sh      # brings up RabbitMQ check, OTel Collector, backend, Bytewax, all microservices, and the frontend
./status.sh     # health check across every component
./stop.sh       # stops everything this script started (RabbitMQ itself is left running, since it's a system service)

Then open http://localhost:5173 for the dashboard and http://localhost:8000/docs for the backend's OpenAPI docs.

Generating traffic

./traffic.sh 60 mixed 5      # 60s of paper-calibrated traffic (80% normal, 10% N+1, 5% bimodal, 5% PII) at ~5 req/s
./traffic.sh 30 pii 3        # PII-heavy traffic, to populate the redaction-density panel
./traffic.sh 60 burst        # ramping 20→100 rps, good for a live "watch it detect anomalies" demo

Watch the dashboard's incident stream populate in real time as the traffic runs.

Onboarding a new service (the interview-ready part)

microservices/web-storefront exists specifically to prove this platform is generic, not hardcoded to its two original demo services. Its index.js contains zero OpenTelemetry code — it's onboarded purely by how it's launched:

OTEL_SERVICE_NAME=web-storefront \
  ./instrumentation/node-wrapper/run_instrumented.sh node microservices/web-storefront/index.js

The wrapper script injects OTel auto-instrumentation via NODE_OPTIONS, points it at the same collector endpoint, and that's it — the service shows up in the dashboard's service list (dynamically discovered from telemetry via GET /api/services, not hardcoded in the frontend) with its own traces, metrics, and anomalies, including a real cross-service distributed trace when it calls the existing Quote Service.

The pitch: the OTel Collector is protocol-based, not app-specific. Anything that speaks OTLP to :4317/:4318 gets ingested, redacted, streamed, scored, and displayed — regardless of language or framework — because service identity flows end-to-end from the OTel service.name resource attribute, never from a hardcoded list.

Troubleshooting

  • ModuleNotFoundError for bytewax / wheel not found: your Python is too new. Bytewax publishes wheels for 3.8–3.12; install one of those versions and recreate venv/.
  • ACCESS_REFUSED connecting to RabbitMQ: the telemetry user doesn't exist yet or lacks permissions — re-run the curl commands under Installing RabbitMQ.
  • UnicodeEncodeError on Windows running a script that prints emoji/λ symbols: set PYTHONIOENCODING=utf-8 (already set inside start.sh/traffic.sh).
  • A generated venv/Scripts/*.exe is "blocked by your organization's Device Guard policy": some locked-down Windows machines refuse to run pip's generated console-script launcher .exe files. Invoke the underlying module directly instead, e.g. python -m uvicorn ... instead of uvicorn.exe ... (this is exactly why instrumentation/python-wrapper/otel_launch.py exists — it calls opentelemetry-instrument's entry point via python.exe rather than its blocked .exe shim).
  • stop.sh/status.sh say a service is down but it's actually still responding: on Windows Git Bash, a PID captured via $! for a backgrounded process can be an MSYS-internal id that doesn't map to the real Windows PID once the process has gone through an exec chain. The scripts route around this by resolving the real PID from the bound TCP port (netstat) and using taskkill //T //F instead of plain kill.

Notes on the paper reference

Traffic profiles and detector thresholds in trigger_traffic.py / detectors.py are calibrated against a companion write-up, "ObserveX: A Cognitive Observability Pipeline" (§IV-C, §V, §VI.A) — the corpus distribution (80% normal / 10% N+1 / 5% bimodal / 5% PII) and the PII redaction-density metric both come from that spec.

About

Cognitive observability platform: OTel ingestion, PII redaction, Bytewax stream processing, hybrid rule+ML anomaly detection, and Gemini-powered root cause analysis.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages