diff --git a/.env.example b/.env.example index 7d8920a..16a7f31 100644 --- a/.env.example +++ b/.env.example @@ -12,3 +12,11 @@ OTEL_ENABLED=False OTEL_SERVICE_NAME=labtelemetry OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4318 OTEL_TRACES_EXPORTER=otlp + +# Portas publicadas no host pelo docker compose. +# Descomente e ajuste se ja houver um Postgres ou um collector OTLP local. +# APP_PORT=8000 +# POSTGRES_PORT=5432 +# JAEGER_UI_PORT=16686 +# OTLP_GRPC_PORT=4317 +# OTLP_HTTP_PORT=4318 diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 52e9bf4..c9ec9fa 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -2,9 +2,9 @@ name: CI on: push: - branches: [main] + branches: [master] pull_request: - branches: [main] + branches: [master] jobs: lint: @@ -35,16 +35,34 @@ jobs: --health-retries 5 ports: - 5432:5432 + env: + DATABASE_URL: postgres://labtelemetry:labtelemetry@localhost:5432/labtelemetry steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 with: python-version: "3.12" + cache: pip - name: Install dependencies - run: | - pip install -r requirements.txt + run: pip install -r requirements.txt + + - name: Django system check + run: python manage.py check + working-directory: ./labtelemetry + + - name: Migrations are in sync with models + run: python manage.py makemigrations --check --dry-run + working-directory: ./labtelemetry + - name: Run tests + run: python manage.py test telemetry -v 2 + working-directory: ./labtelemetry + + # Informativo: o projeto e um lab local, DEBUG=True por padrao. + # Expoe o delta ate um deploy real sem travar o CI. + - name: Deployment checklist (informational) + continue-on-error: true + run: python manage.py check --deploy working-directory: ./labtelemetry - run: python manage.py test -v 2 env: - DATABASE_URL: postgres://labtelemetry:labtelemetry@localhost:5432/labtelemetry + DEBUG: "False" diff --git a/.github/workflows/django.yml b/.github/workflows/django.yml deleted file mode 100644 index 36ae6c0..0000000 --- a/.github/workflows/django.yml +++ /dev/null @@ -1,34 +0,0 @@ -name: django - -on: - push: - branches: - - master - pull_request: - -jobs: - test: - runs-on: ubuntu-latest - - steps: - - name: Checkout - uses: actions/checkout@v4 - - - name: Set up Python - uses: actions/setup-python@v5 - with: - python-version: "3.12" - - - name: Install dependencies - run: | - python -m pip install --upgrade pip - pip install -r requirements.txt - - - name: Django system check - run: python labtelemetry/manage.py check - - - name: Django migrations dry run - run: python labtelemetry/manage.py makemigrations --check --dry-run - - - name: Django tests - run: python labtelemetry/manage.py test telemetry -v 1 diff --git a/.gitignore b/.gitignore index 862b854..265641c 100644 --- a/.gitignore +++ b/.gitignore @@ -22,6 +22,7 @@ automação40/ .vscode/ .idea/ .gemini/ +.serena/ .DS_Store # --- SECURITY & COMPLIANCE RULES (NEVER EXPOSE PLANNING/SENSITIVE FILES) --- @@ -42,3 +43,19 @@ automação40/ /docs/execution/ /docs/prd/ /docs/process/ + +# Architecture Decision Records (internal governance) +/docs/adr/ + +# Heuristics usage log (local-only tooling artifact) +.heuristics_usage_log.json + +# Session handoff state (local-only tooling artifact) +.last-handoff.json + +# Local tooling caches +.tokensave/ +.pytest_cache/ +.ruff_cache/ +.claude/ +.serena/ diff --git a/.last-handoff.json b/.last-handoff.json deleted file mode 100644 index 7067c49..0000000 --- a/.last-handoff.json +++ /dev/null @@ -1,31 +0,0 @@ -{ - "last_session": "2026-06-23T21:12:11-03:00", - "project": "labtelemetry", - "observed_branch": "master", - "plan": "005", - "package_completed": "B", - "summary": "Plano 005 Pacote B concluido com source persistido em TelemetryReading, teste do ModbusTCPAdapter sem socket real e docs publicos alinhados ao contrato atual.", - "files": [ - "labtelemetry/telemetry/models.py", - "labtelemetry/telemetry/migrations/0003_telemetryreading_source.py", - "labtelemetry/telemetry/management/commands/ingest_telemetry.py", - "labtelemetry/telemetry/views.py", - "labtelemetry/telemetry/test_sources.py", - "labtelemetry/telemetry/test_ingest_telemetry.py", - "labtelemetry/telemetry/tests.py", - "docs/api.md", - "docs/architecture.md", - "docs/data-model.md", - "docs/data-contract.md", - "docs/session-handoffs/20260623T211211_pacote_b_encerramento.md" - ], - "validations": [ - ".venv/bin/python labtelemetry/manage.py check", - ".venv/bin/python labtelemetry/manage.py makemigrations --check --dry-run", - ".venv/bin/python labtelemetry/manage.py test telemetry --verbosity=1", - ".venv/bin/python labtelemetry/manage.py migrate", - ".venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once" - ], - "handoff": "docs/session-handoffs/20260623T211211_pacote_b_encerramento.md", - "resume_prompt": "cd \"/media/Arquivos/Engenharia dados IOT 2026/labtelemetry\"; git status --short; .venv/bin/python labtelemetry/manage.py test telemetry --verbosity=1" -} diff --git a/LICENSE b/LICENSE new file mode 100644 index 0000000..ca29864 --- /dev/null +++ b/LICENSE @@ -0,0 +1,21 @@ +MIT License + +Copyright (c) 2026 Roberto Nascimento + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/README.md b/README.md index 70461dc..7456c09 100644 --- a/README.md +++ b/README.md @@ -1,152 +1,254 @@ -# LabTelemetry +
-

- LabTelemetry banner -

+# 🧪 LabTelemetry -

- Python 3.12 - Django 5 - PostgreSQL 16 - OpenTelemetry - Jaeger - HTMX - Chart.js - Docker -

+[![CI](https://github.com/Roberton003/labtelemetry/actions/workflows/ci.yml/badge.svg?branch=master)](https://github.com/Roberton003/labtelemetry/actions/workflows/ci.yml) +[![Python](https://img.shields.io/badge/Python-3.12-3776AB?style=for-the-badge&logo=python&logoColor=white)](https://www.python.org/) +[![Django](https://img.shields.io/badge/Django-5.2-092E20?style=for-the-badge&logo=django&logoColor=white)](https://www.djangoproject.com/) +[![PostgreSQL](https://img.shields.io/badge/PostgreSQL-16-4169E1?style=for-the-badge&logo=postgresql&logoColor=white)](https://www.postgresql.org/) +[![Docker](https://img.shields.io/badge/Docker-24%2B-2496ED?style=for-the-badge&logo=docker&logoColor=white)](https://www.docker.com/) +[![OpenTelemetry](https://img.shields.io/badge/OpenTelemetry-4B9CD3?style=for-the-badge&logo=opentelemetry&logoColor=white)](https://opentelemetry.io/) +[![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg?style=for-the-badge)](LICENSE) -

- Reproducible OT/IT telemetry lab with simulation, quality rules, JSON API, dashboard, and local observability. -

+LabTelemetry -## Overview +

Laboratório de telemetria OT/IT reproduzível: ingestão industrial (Modbus TCP, OPC-UA, simulador determinístico), regras de qualidade de processo, API JSON e dashboard operacional.

-LabTelemetry is a Django project built to demonstrate a complete telemetry path without inflating the stack beyond what the use case needs. +
-```text -simulator -> ingestion -> quality evaluation -> PostgreSQL/SQLite -> JSON API -> dashboard -``` +--- + +## 📌 Project Highlights + +- **Idempotência garantida pelo banco, não pela aplicação.** A deduplicação vem de `UniqueConstraint(sensor, timestamp)` combinada com `bulk_create(ignore_conflicts=True)` — uma checagem por lote dentro do Postgres, em vez de um `SELECT` por amostra vindo do Python. O guardrail tem [teste negativo](labtelemetry/telemetry/test_ingest_telemetry.py): sem o mecanismo, o replay estoura `IntegrityError`. +- **Fontes OT plugáveis atrás de uma única ABC.** `TelemetrySource` define `read()`/`health()`/`close()`; Modbus TCP, OPC-UA e simulador são intercambiáveis via `--source`. Adicionar um protocolo não toca o código de ingestão. +- **Mapeamento tag→ponto é configuração explícita, não convenção.** Node OPC-UA e registrador Modbus são ligados ao sensor por `--opcua-node "ns=2;i=101:3"` e `--modbus-register "0:3:0.01"`. Índice posicional não é chave primária de sensor — tratá-lo como tal produz dado plausível e errado, então o comando exige o mapeamento em vez de adivinhar. +- **Fator de escala como cidadão de primeira classe.** Holding register é uint16: um pH de 7.40 não cabe nele. O CLP publica `740` e o mapeamento diz como voltar à grandeza física. Sem isso a leitura entra como pH 740 — fora de faixa, e errada de um jeito que só aparece no gráfico. +- **Simulador determinístico por seed.** Reproduzir uma sequência de falha é `--seed 42`, não "esperar o sensor falhar de novo" — o que torna o teste das regras de qualidade repetível. +- **Observabilidade opcional em runtime.** OpenTelemetry é ligado por `OTEL_ENABLED`; desligado, o custo é zero e nenhuma dependência de trace entra no caminho da request. +- **Qualidade de processo separada da persistência.** `evaluate_reading()` é pura (sem I/O), o que permite avaliar o lote inteiro em memória antes de um único INSERT. +- **Degradação explícita, não silenciosa.** Se `pymodbus` não está instalado ou o CLP está fora do ar, `/api/health/sources/` reporta `disconnected` — a fonte não some do inventário. + +--- -## Platform Snapshot +## 🏛️ Architecture & Tech Stack -| Area | Current State | +| Camada | Tecnologia | |---|---| -| Domain | Industrial telemetry lab for pH, turbidity, and TOC | -| Ingestion | Deterministic simulator plus Modbus TCP adapter surface | -| Storage | PostgreSQL 16 via Docker Compose, SQLite fallback for local-only runs | -| Backend | Django 5.2.9 | -| Frontend | Server-rendered dashboard with HTMX and Chart.js | -| Observability | OpenTelemetry with Jaeger, optional at runtime | -| Data Quality | Threshold rules, drift warning, active alerts | -| Validation | Django tests, API checks, end-to-end local manual | - -## What You Can Validate Today - -- Reproducible local ingestion from a controlled simulator -- Persistent telemetry readings with source lineage -- JSON endpoints under `/api/...` -- Operational dashboard rendered by Django -- Source health checks for simulator and Modbus -- Optional traces in Jaeger - -## Dashboard +| Aquisição OT | Modbus TCP (`pymodbus`), OPC-UA (`asyncua`), simulador determinístico | +| Ingestão | Django management command (`ingest_telemetry`), lote com `bulk_create` | +| Qualidade | Regras de limite de processo e detecção de drift (`telemetry/quality.py`) | +| Persistência | PostgreSQL 16 (Docker Compose); SQLite como fallback local | +| Backend | Django 5.2 / Python 3.12 | +| API | Endpoints JSON server-side, sem framework REST adicional | +| Frontend | Django Templates + HTMX + Chart.js | +| Observabilidade | OpenTelemetry → Jaeger (opt-in via `OTEL_ENABLED`) | +| Runtime | Docker + Gunicorn | +| Qualidade de código | ruff, 73 testes Django, CI no GitHub Actions | + +--- + +## 🗺️ Architecture Diagram + +```mermaid +flowchart LR + subgraph OT["Camada OT"] + MB["Modbus TCP
(CLP / RTU)"] + UA["OPC-UA
(servidor)"] + SIM["Simulador
(seed determinístico)"] + end + + subgraph ING["Ingestão"] + ADP["TelemetrySource (ABC)
read / health / close"] + QA["evaluate_reading()
limites + drift"] + BULK["bulk_create
ignore_conflicts"] + end + + subgraph IT["Camada IT"] + DB[("PostgreSQL 16
UniqueConstraint
sensor + timestamp")] + API["API JSON
/api/..."] + DASH["Dashboard
HTMX + Chart.js"] + ALERT["TelemetryAlert
raise_alert idempotente"] + end + + OTEL(["OpenTelemetry → Jaeger
opt-in"]) + + MB --> ADP + UA --> ADP + SIM --> ADP + ADP --> QA --> BULK --> DB + QA --> ALERT --> DB + DB --> API --> DASH + API -.-> OTEL +``` + +--- + +## 📊 O Dashboard

- LabTelemetry dashboard mockup + Dashboard LabTelemetry

-The user interface is built with Django templates, HTMX, and Chart.js. +Renderizado pelo Django, atualizado por HTMX em fragmentos parciais (cards, leituras, alertas, sensores, saúde das fontes) — sem SPA e sem build step de frontend. -## Documentation Map +### Endpoints da API -| Document | Purpose | +| Endpoint | Retorno | |---|---| -| [docs/overview.md](docs/overview.md) | Project scope and public positioning | -| [docs/architecture.md](docs/architecture.md) | Runtime structure and component boundaries | -| [docs/api.md](docs/api.md) | API endpoints and public contract notes | -| [docs/operations.md](docs/operations.md) | Local setup and operational commands | -| [docs/manual_validacao_ponta_a_ponta.md](docs/manual_validacao_ponta_a_ponta.md) | Full end-to-end validation in parallel terminals | -| [docs/data-model.md](docs/data-model.md) | Operational data model and database schemas | -| [docs/data-contract.md](docs/data-contract.md) | Public API and data contract definition | -| [docs/replay-idempotency.md](docs/replay-idempotency.md) | Replay, deduplication, and idempotency behavior | -| [docs/security.md](docs/security.md) | Public documentation boundary and secret handling | +| `GET /api/summary/` | Contagens agregadas e timestamp da última leitura | +| `GET /api/sensors/` | Inventário de sensores com fator de calibração | +| `GET /api/readings/recent/` | Últimas leituras (`?limit=`, teto de 500) | +| `GET /api/sensors//readings/` | Série temporal de um sensor | +| `GET /api/alerts/active/` | Alertas operacionais ativos | +| `GET /api/health/sources/` | Estado de conexão de cada fonte OT | + +Contrato completo em [docs/data-contract.md](docs/data-contract.md). + +--- -## Quick Start +## 🚀 Quick Start & Setup -### Bootstrap Environment +**Pré-requisitos:** Python 3.12+, Docker Compose. + +### Subir tudo com Docker ```bash -python3 -m venv .venv -source .venv/bin/activate -pip install -r requirements.txt +git clone https://github.com/Roberton003/labtelemetry.git +cd labtelemetry cp .env.example .env -docker compose up -d +docker compose up --build -d ``` -### Migrate and Run +Dashboard em http://127.0.0.1:8000/ · Jaeger em http://localhost:16686 + +> Já tem um Postgres ou um collector OTLP local ocupando as portas? Sobrescreva sem editar o compose: +> `POSTGRES_PORT=55432 OTLP_GRPC_PORT=54317 OTLP_HTTP_PORT=54318 docker compose up -d` + +### Rodar localmente (SQLite, sem Docker) ```bash -export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" +python3 -m venv .venv && source .venv/bin/activate +pip install -r requirements.txt +cp .env.example .env + python labtelemetry/manage.py migrate python labtelemetry/manage.py runserver 127.0.0.1:8000 ``` -### Generate Telemetry +### Gerar telemetria ```bash -python labtelemetry/manage.py ingest_telemetry --source simulator --once -curl -s http://127.0.0.1:8000/api/summary/ -``` - -Open locally: +# Uma leitura de cada sensor, a partir do simulador determinístico +python labtelemetry/manage.py ingest_telemetry --source simulator --once --sim-count 3 -- Dashboard: http://127.0.0.1:8000/ -- Admin: http://127.0.0.1:8000/admin/ -- Jaeger: http://127.0.0.1:16686 +# Loop contínuo a cada 5s (Ctrl+C encerra de forma limpa) +python labtelemetry/manage.py ingest_telemetry --source simulator --interval 5 -## Observability +# Fonte industrial real — Modbus TCP +# registrador 0 -> sensor 1, com escala: o CLP publica 740, o pH é 7.40 +python labtelemetry/manage.py ingest_telemetry --source modbus \ + --modbus-host 192.168.0.10 \ + --modbus-register "0:1:0.01" \ + --modbus-register "4:2:0.1" -Tracing is disabled by default: +# Fonte industrial real — OPC-UA (cada node mapeado ao sensor que alimenta) +python labtelemetry/manage.py ingest_telemetry --source opcua \ + --opcua-url opc.tcp://plc.local:4840 \ + --opcua-node "ns=2;i=101:1" \ + --opcua-node "ns=2;i=103:5" -```bash -OTEL_ENABLED=False +curl -s http://127.0.0.1:8000/api/summary/ ``` -To validate traces locally: +### Validar ```bash -export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" -OTEL_ENABLED=True .venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 -curl -s http://127.0.0.1:8000/api/summary/ -curl -s "http://localhost:16686/api/traces?service=labtelemetry&limit=5" +python labtelemetry/manage.py test telemetry # 73 testes +ruff check labtelemetry/ ``` -## Validation +Roteiro completo em [docs/manual_validacao_ponta_a_ponta.md](docs/manual_validacao_ponta_a_ponta.md). -### Fast Sanity +--- -```bash -.venv/bin/python labtelemetry/manage.py check -.venv/bin/python labtelemetry/manage.py makemigrations --check --dry-run -.venv/bin/python labtelemetry/manage.py test telemetry --verbosity=1 -``` +## ⚙️ Environment Variables -### Full Practical Flow +Todas em `.env` (ver [.env.example](.env.example)): -Use [docs/manual_validacao_ponta_a_ponta.md](docs/manual_validacao_ponta_a_ponta.md) for the parallel-terminal walkthrough. +| Variável | Padrão | Função | +|---|---|---| +| `SECRET_KEY` | chave de dev | Chave criptográfica do Django — **trocar fora de dev** | +| `DEBUG` | `True` | Modo debug | +| `ALLOWED_HOSTS` | `127.0.0.1,localhost` | Hosts aceitos, separados por vírgula | +| `DATABASE_URL` | `sqlite:///db.sqlite3` | Conexão via `dj-database-url`; aceita `postgres://...` | +| `OTEL_ENABLED` | `False` | Liga a instrumentação OpenTelemetry | +| `OTEL_EXPORTER_OTLP_ENDPOINT` | `http://localhost:4318` | Coletor OTLP (Jaeger) | +| `OTEL_SERVICE_NAME` | `labtelemetry` | Nome do serviço nos traces | +| `APP_PORT` / `POSTGRES_PORT` | `8000` / `5432` | Portas publicadas no host pelo Compose | +| `OTLP_GRPC_PORT` / `OTLP_HTTP_PORT` / `JAEGER_UI_PORT` | `4317` / `4318` / `16686` | Portas do Jaeger no host | + +--- + +## 📚 Documentation Resources + +| Documento | Conteúdo | +|---|---| +| [docs/overview.md](docs/overview.md) | Escopo do projeto e posicionamento | +| [docs/architecture.md](docs/architecture.md) | Estrutura de runtime e fronteiras entre componentes | +| [docs/api.md](docs/api.md) | Endpoints e contrato público | +| [docs/data-model.md](docs/data-model.md) | Modelo de dados operacional | +| [docs/data-contract.md](docs/data-contract.md) | Contrato de dados, garantias e limitações | +| [docs/replay-idempotency.md](docs/replay-idempotency.md) | O que é garantido no replay — e o que não é | +| [docs/operations.md](docs/operations.md) | Setup e comandos operacionais | +| [docs/manual_validacao_ponta_a_ponta.md](docs/manual_validacao_ponta_a_ponta.md) | Validação end-to-end em terminais paralelos | +| [docs/security.md](docs/security.md) | Fronteira de documentação pública e tratamento de segredos | +| [sql/analytics/](sql/analytics/) | Consultas de frescor, volume e taxa de anomalia | + +Aprofundamento na [Wiki do projeto](https://github.com/Roberton003/labtelemetry/wiki). + +--- + +## 🌳 Estrutura do Projeto + +```text +labtelemetry/ +├── labtelemetry/ # Projeto Django +│ ├── labtelemetry/ # settings, urls, wsgi/asgi (OTel condicional) +│ └── telemetry/ # App único de domínio +│ ├── models.py # Sensor, Reading (UniqueConstraint), Alert +│ ├── quality.py # Regras de limite, drift e alerta idempotente +│ ├── views.py # API JSON + fragmentos HTMX do dashboard +│ ├── sources/ # Adapters OT sob a ABC TelemetrySource +│ │ ├── base.py # TelemetrySource, TelemetrySample +│ │ ├── modbus.py # Modbus TCP via pymodbus +│ │ ├── opcua.py # OPC-UA via asyncua (+ servidor de teste) +│ │ └── simulator.py # Gerador gaussiano determinístico +│ ├── management/commands/ +│ │ ├── ingest_telemetry.py # Runner: fonte → qualidade → lote +│ │ └── simulate_telemetry.py # Gerador de cenário sintético +│ ├── templates/telemetry/ # Dashboard + parciais HTMX +│ └── test_*.py, tests.py # 73 testes +├── docs/ # Documentação pública + wiki-seed +├── sql/analytics/ # Consultas operacionais +├── .github/workflows/ci.yml # ruff + testes com Postgres 16 +├── docker-compose.yml # app + postgres + jaeger +└── Dockerfile # Runtime Gunicorn +``` -## Boundaries +--- -This repository is intentionally scoped as a local lab and portfolio-grade system, not a generalized production platform. +## 🎯 Escopo e Fronteiras -Out of current public scope: +Este repositório é deliberadamente um laboratório local, não uma plataforma de produção generalizada. O que está fora de escopo está fora por decisão, não por omissão: -- distributed stream processing -- production authentication -- multi-region cloud infrastructure +- Processamento de stream distribuído +- Autenticação de produção na API +- Infraestrutura cloud multi-região +- Garantia formal de *exactly-once* — o comportamento real e seus limites estão em [replay-idempotency.md](docs/replay-idempotency.md) -## Wiki +--- -The project wiki is available at: +## 📄 License -- https://github.com/Roberton003/labtelemetry/wiki +[MIT](LICENSE) © 2026 Roberto Nascimento diff --git a/docker-compose.yml b/docker-compose.yml index a38fafc..07c44e6 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,4 +1,20 @@ services: + app: + build: . + container_name: labtelemetry_app + ports: + - "${APP_PORT:-8000}:8000" + env_file: + - .env + environment: + DATABASE_URL: postgresql://labtelemetry:labtelemetry_dev@postgres:5432/labtelemetry + OTEL_ENABLED: "True" + OTEL_EXPORTER_OTLP_ENDPOINT: http://jaeger:4318 + depends_on: + postgres: + condition: service_started + restart: unless-stopped + postgres: image: postgres:16-alpine container_name: labtelemetry_postgres @@ -7,7 +23,9 @@ services: POSTGRES_USER: labtelemetry POSTGRES_PASSWORD: labtelemetry_dev ports: - - "5432:5432" + # Publicado so para inspecao externa; o app fala pela rede do compose. + # Sobrescreva POSTGRES_PORT se ja houver um Postgres local na 5432. + - "${POSTGRES_PORT:-5432}:5432" volumes: - pgdata:/var/lib/postgresql/data restart: unless-stopped @@ -18,9 +36,10 @@ services: environment: - COLLECTOR_OTLP_ENABLED=true ports: - - "16686:16686" - - "4317:4317" - - "4318:4318" + # Sobrescreva se ja houver um collector OTLP local ocupando 4317/4318. + - "${JAEGER_UI_PORT:-16686}:16686" + - "${OTLP_GRPC_PORT:-4317}:4317" + - "${OTLP_HTTP_PORT:-4318}:4318" restart: unless-stopped volumes: diff --git a/docs/architecture.md b/docs/architecture.md index e1d7b7c..7c9eeef 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -16,9 +16,8 @@ Telemetry source - `telemetry.models`: sensor, reading, and alert persistence models. - `telemetry.quality`: threshold and drift evaluation rules. - `telemetry.management.commands.simulate_telemetry`: deterministic telemetry simulation. -- `telemetry.management.commands.telemetry_simulate`: operational wrapper for repeated simulation. - `telemetry.management.commands.ingest_telemetry`: source-based ingestion command. -- `telemetry.sources`: source adapter abstraction for simulator and Modbus TCP. +- `telemetry.sources`: source adapter abstraction for simulator, Modbus TCP, and OPC-UA. - `telemetry.views`: dashboard and JSON API views. ## Data Model @@ -37,5 +36,6 @@ The current adapters are: - `SimulatorAdapter`: uses the existing simulator path for reproducible local runs. - `ModbusTCPAdapter`: provides a Modbus TCP adapter surface with configurable host, port, unit id, and timeout. +- `OpcUaAdapter`: connects to OPC-UA servers to read telemetry node variables. -The simulator remains the default reproducible path. Real Modbus validation depends on an available device or simulator. +The simulator remains the default reproducible path. Real Modbus and OPC-UA validation depend on an available device or simulator. diff --git a/docs/data-contract.md b/docs/data-contract.md index bd7218e..869d584 100644 --- a/docs/data-contract.md +++ b/docs/data-contract.md @@ -9,7 +9,7 @@ publica do LabTelemetry. | Papel | Sistema | Dado | |---|---|---| -| Produtor primario | `telemetry_simulate` ou `ingest_telemetry` | `TelemetryReading` | +| Produtor primario | `simulate_telemetry` ou `ingest_telemetry` | `TelemetryReading` | | Produtor secundario | `quality.py` | `TelemetryAlert` | | Consumidor | Dashboard HTML/HTMX | `/api/summary/`, `/api/readings/recent/`, `/api/alerts/active/` | | Consumidor | Avaliador tecnico | `curl` ou cliente HTTP simples | @@ -41,7 +41,7 @@ Contrato: - `raw_value`: valor lido da fonte - `calibrated_value`: valor persistido para avaliacao de qualidade - `value`: alias de `calibrated_value` no endpoint de leituras recentes -- `source`: lineage curto, por exemplo `simulator:seed=42` ou `modbus:host:port` +- `source`: lineage curto, por exemplo `simulator:seed=42`, `modbus:host:port` ou `opcua:host:port` - `status`: enum `NORMAL | OUT_OF_BOUNDS | DRIFT_WARNING` ### `TelemetryReading` em `/api/sensors/{id}/readings/` @@ -103,8 +103,12 @@ Contrato: ## Garantias e Limitacoes -- `ingest_telemetry --once` nao faz deduplicacao; cada execucao cria novas leituras. -- Cada leitura e persistida individualmente; nao existe batch atomico publico. +- `ingest_telemetry` deduplica por `(sensor, timestamp)` via `UniqueConstraint` + + `bulk_create(ignore_conflicts=True)`: reprocessar a mesma janela e no-op. + Detalhes e limites em [replay-idempotency.md](replay-idempotency.md). +- Dois valores diferentes no mesmo `(sensor, timestamp)` colapsam no primeiro; + nao ha upsert. +- Leituras sao persistidas em lote por ciclo de leitura da fonte. - `source` guarda somente lineage curto, nao payload bruto do protocolo. - Nao ha SLO publico de latencia para a API. diff --git a/docs/manual_validacao_ponta_a_ponta.md b/docs/manual_validacao_ponta_a_ponta.md index 6e6a8d3..d30e38b 100644 --- a/docs/manual_validacao_ponta_a_ponta.md +++ b/docs/manual_validacao_ponta_a_ponta.md @@ -22,7 +22,7 @@ Use 3 terminais. ### Terminal A - Infraestrutura ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" +cd /caminho/para/labtelemetry # raiz do repositorio clonado docker compose up -d ``` @@ -34,7 +34,7 @@ Resultado esperado: ### Terminal B - Banco E Aplicacao ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" +cd /caminho/para/labtelemetry # raiz do repositorio clonado export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" .venv/bin/python labtelemetry/manage.py migrate .venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 @@ -48,7 +48,7 @@ Resultado esperado: ### Terminal C - Dados E Verificacao ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" +cd /caminho/para/labtelemetry # raiz do repositorio clonado export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" .venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once curl -sS "http://127.0.0.1:8000/api/summary/" @@ -94,7 +94,7 @@ Resultado esperado: Em um quarto terminal ou reaproveitando o Terminal B: ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" +cd /caminho/para/labtelemetry # raiz do repositorio clonado export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" export OTEL_ENABLED=True .venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 diff --git a/docs/operations.md b/docs/operations.md index 2b6d22e..eb81547 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -36,7 +36,7 @@ Open: ```bash .venv/bin/python labtelemetry/manage.py simulate_telemetry --once --seed 42 --sensors 6 -.venv/bin/python labtelemetry/manage.py telemetry_simulate --seed 42 --count 50 +.venv/bin/python labtelemetry/manage.py simulate_telemetry --seed 42 --iterations 50 .venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once ``` diff --git a/docs/replay-idempotency.md b/docs/replay-idempotency.md index 8642cd2..bc2b3a5 100644 --- a/docs/replay-idempotency.md +++ b/docs/replay-idempotency.md @@ -2,98 +2,126 @@ ## Contexto -O LabTelemetry não foi projetado com garantias formais de exactly-once. -Este documento explica o comportamento real do sistema para que avaliadores -entendam os trade-offs sem surpresas. +O LabTelemetry não oferece garantia formal de *exactly-once*. Este documento +descreve o comportamento real do sistema — o que é garantido, o que não é, e +onde exatamente fica a fronteira — para que ninguém descubra o limite em +produção. -## Estado Atual +## A garantia que existe -### Ingestão (`ingest_telemetry --once`) +A deduplicação é feita **no banco, não na aplicação**: -Cada execução do comando `ingest_telemetry --once`: -1. Abre uma conexão com a fonte (`SimulatorAdapter` ou `ModbusTCPAdapter`) -2. Itera sobre as amostras fornecidas pela fonte -3. Para cada amostra, cria um `TelemetryReading` no banco - -**Comportamento:** Não há verificação de duplicata. Se o mesmo comando for -executado duas vezes com o mesmo seed, serão criados registros duplicados -(com `id` diferente, mesmo `timestamp` e `raw_value`). +```python +# telemetry/models.py +class Meta: + constraints = [ + UniqueConstraint(fields=["sensor", "timestamp"], name="uq_sensor_timestamp"), + ] +``` -### Simulação (`telemetry_simulate --seed 42 --count N`) +```python +# telemetry/management/commands/ingest_telemetry.py +TelemetryReading.objects.bulk_create(batch, ignore_conflicts=True) +``` -Usa `seed` para gerar a mesma sequência de leituras, mas **não verifica** -se aquelas leituras já existem. Cada execução insere N novos registros. +O par constraint + `ignore_conflicts` é o que dá a garantia. Reprocessar a +mesma janela `(sensor, timestamp)` é **no-op**: as linhas já existentes são +descartadas pelo banco, sem erro e sem duplicata. -## Como Replay Funciona (e Não Funciona) +Optar pela constraint em vez de um `get_or_create` por leitura é deliberado: +a checagem acontece uma vez por lote, dentro do banco, em vez de um `SELECT` +por amostra vindo da aplicação. -### Cenário: Reproduzir uma falha +### Teste negativo -```bash -# Primeira execução — gera 10 leituras -telemetry_simulate --seed 42 --count 10 --anomaly-rate 0.3 +A garantia não é uma alegação deste documento — ela tem um teste que falha +quando o mecanismo é removido: -# Segunda execução — gera OUTRAS 10 leituras (mesmo seed, mesmo valor) -telemetry_simulate --seed 42 --count 10 --anomaly-rate 0.3 -# Resultado: 20 leituras no banco, as 10 primeiras duplicadas em valor +``` +telemetry.test_ingest_telemetry.IngestTelemetryCommandTest + .test_replay_same_window_is_idempotent ``` -**Conclusão:** O seed garante **repetibilidade do valor**, não -**idempotência de inserção**. +Ele executa a mesma janela duas vezes e afirma que a contagem de leituras não +dobra. Sem `ignore_conflicts`, a segunda execução levanta +`IntegrityError: UNIQUE constraint failed` e derruba o loop de ingestão — +comportamento verificado antes do fix, não presumido. -### Cenário: Reprocessar um dia +## A garantia que NÃO existe -Não há suporte a janela temporal de reprocessamento. O comando sempre -cria leituras "novas" com timestamp = agora. +`(sensor, timestamp)` é a chave de deduplicação. Isso tem duas consequências +que valem estar explícitas: -## Deduplicação +1. **Dois valores diferentes no mesmo `(sensor, timestamp)` colapsam no + primeiro.** A segunda leitura é silenciosamente descartada, não corrigida. + Não há *upsert*, não há "última escrita vence". +2. **Timestamps diferentes nunca deduplicam**, mesmo com valor idêntico. Se a + fonte carimba `timestamp = agora`, cada execução é uma janela nova — e + portanto insere linhas novas, corretamente. -**Não existe.** Não há índice único natural, hash ou upsert que impeça -duplicatas. A chave primária é `id` (auto-increment), que por definição -nunca colide. +Não há hash de payload nem identificador de evento da fonte externa. A +identidade de uma leitura é o par sensor/instante, e nada além disso. -### O que impediria deduplicar hoje +## Como replay funciona na prática -- `TelemetryReading` não tem `(sensor_id, timestamp, raw_value)` como - unique constraint -- Django ORM não suporta `INSERT ... ON CONFLICT` sem raw SQL ou - `get_or_create` (que adiciona SELECT antes de INSERT) -- Não há hash de payload ou identificador de fonte externa +### Reproduzir uma sequência de valores -## Idempotência Real no Sistema +```bash +simulate_telemetry --seed 42 --iterations 10 --anomaly-rate 0.3 +``` -Apesar da ingestão não ser idêntica, **algumas partes do sistema são -idempotentes por construção:** +O `--seed` garante **repetibilidade do valor**, não idempotência de inserção — +são coisas distintas. Como `simulate_telemetry` carimba `timestamp = agora` a +cada execução, rodar duas vezes produz 20 linhas com os mesmos 10 valores em +instantes diferentes. Isso é o comportamento correto: são duas observações +distintas do mesmo cenário simulado. -| Componente | Idempotente? | Como | -|-----------|-------------|------| -| `quality.py: evaluate_and_alert()` | ✅ Sim | Se alerta ativo já existe para o mesmo problema, não recria | -| `GET /api/...` | ✅ Sim | REST GET é naturalmente idempotente | -| `migrate` | ✅ Sim | Django migrations são idempotentes | -| `ingest_telemetry --once` | ❌ Não | Cada execução cria novas leituras | -| `telemetry_simulate` | ❌ Não | Cada execução cria novas leituras | +### Reprocessar uma janela vinda de uma fonte OT -## O Que Mudaria para Idempotência Formal +Quando a fonte fornece o timestamp (Modbus e OPC-UA fornecem), o replay é +idempotente de verdade: -Se o projeto evoluísse para exigir exactly-once: +```bash +# Executar duas vezes sobre a mesma janela da fonte +python manage.py ingest_telemetry --source modbus --once +python manage.py ingest_telemetry --source modbus --once +# As leituras cujo (sensor, timestamp) ja existe sao descartadas pelo banco +``` + +Não há, hoje, um `--start-time`/`--end-time` para pedir uma janela histórica +explícita à fonte. Replay dirigido por janela está fora do escopo atual. + +## Mapa de idempotência + +| Componente | Idempotente? | Mecanismo | +|---|---|---| +| `ingest_telemetry` (mesmo sensor+timestamp) | ✅ Sim | `UniqueConstraint` + `bulk_create(ignore_conflicts=True)` | +| `ingest_telemetry` (timestamp novo a cada leitura) | ❌ Não | Janela diferente = leitura diferente, por definição | +| `quality.raise_alert()` | ✅ Sim | Não recria alerta se já houver um ativo do mesmo tipo para o sensor | +| `GET /api/...` | ✅ Sim | GET é idempotente por contrato HTTP | +| `migrate` | ✅ Sim | Migrations do Django | +| `simulate_telemetry` | ❌ Não | Carimba `timestamp = agora`; cada execução é uma janela nova | -1. **Unique constraint:** Adicionar `(sensor_id, timestamp, raw_value)` como - unique → `INSERT ... ON CONFLICT DO NOTHING` -2. **Hash de payload:** `SHA256(raw_value + timestamp + sensor_id)` como - chave natural -3. **Campo `source`:** Identificar origem para evitar colisão entre - simulador e Modbus -4. **Janela de replay:** Permitir `--start-time` e `--end-time` no comando - de ingestão +## O que mudaria para exactly-once formal + +Fora do escopo atual, registrado como evolução: + +1. **Chave natural mais forte:** incluir `source` na constraint, para que + simulador e Modbus não disputem o mesmo `(sensor, timestamp)`. +2. **Upsert de verdade:** `INSERT ... ON CONFLICT DO UPDATE` via + `bulk_create(update_conflicts=True)`, para que a releitura corrija o valor + em vez de descartá-lo. +3. **Hash de payload** como identidade do evento, desacoplando a deduplicação + da precisão do timestamp. +4. **Janela de replay explícita:** `--start-time` / `--end-time` no comando de + ingestão. ## Resumo | Pergunta | Resposta | -|----------|----------| -| Posso executar o mesmo comando duas vezes? | Sim, mas cria duplicatas | +|---|---| +| Posso reexecutar a ingestão da mesma janela? | Sim — é no-op, não duplica nem falha | | Posso reproduzir a mesma sequência de valores? | Sim, com `--seed` | -| Posso reprocessar uma janela temporal? | Não | +| Uma releitura corrige um valor já gravado? | Não — é descartada | +| Posso pedir uma janela histórica à fonte? | Não | | O sistema impede alerta duplicado? | Sim | -| O sistema impede leitura duplicada? | Não | - -Esta é uma limitação documentada e aceita para o MVP. Idempotência formal -está no backlog como evolução futura. diff --git a/docs/wiki-seed/API.md b/docs/wiki-seed/API.md index e910203..a011ecf 100644 --- a/docs/wiki-seed/API.md +++ b/docs/wiki-seed/API.md @@ -1,40 +1,52 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# API +# 🔗 API -All JSON endpoints are served under `/api/`. +Todos os endpoints JSON são servidos sob `/api/`. São somente leitura — a escrita acontece pelos comandos de ingestão. -## Endpoints +--- -| Method | Path | Purpose | +## 📍 Endpoints + +| Método | Rota | Função | |---|---|---| -| GET | `/api/sensors/` | List sensors | -| GET | `/api/readings/recent/?limit=50` | Return recent readings | -| GET | `/api/sensors//readings/?limit=100` | Return readings for one sensor | -| GET | `/api/alerts/active/` | Return active alerts | -| GET | `/api/summary/` | Return operational summary | -| GET | `/api/health/sources/` | Return telemetry source status | +| GET | `/api/summary/` | Resumo operacional: contagens e última leitura | +| GET | `/api/sensors/` | Inventário de sensores | +| GET | `/api/readings/recent/?limit=50` | Leituras recentes (teto de 500) | +| GET | `/api/sensors//readings/?limit=100` | Série temporal de um sensor (teto de 500) | +| GET | `/api/alerts/active/` | Alertas ativos | +| GET | `/api/health/sources/` | Estado de conexão de cada fonte OT | + +--- -## Notes +## 📦 Formato do Payload -- Reading payloads include `source`, a short lineage field such as `simulator:seed=42` or `modbus:host:port` -- Legacy bare API routes without `/api/` are not part of the public contract -- The source health endpoint is operational metadata; it does not replace a live Modbus connectivity test +**Leitura:** `sensor_name`, `parameter`, `timestamp`, `raw_value`, `calibrated_value`, `source`, `status`. -## Payload Shape +**Por que bruto e calibrado juntos:** o valor calibrado é o que o processo enxerga; o bruto é o que o sensor mandou. A divergência entre os dois é o que a detecção de drift observa — descartar o bruto tornaria a regra inauditável depois do fato. -The recent readings endpoint returns a compact operational payload: +**Campo `source`:** lineage curto (`simulator:seed=42`, `modbus:host:port`), suficiente para rastrear a origem em consulta sem carregar o payload do protocolo. -- `sensor_name` -- `parameter` -- `raw_value` -- `calibrated_value` -- `source` -- `status` +--- -## Example +## 🚦 Exemplo ```bash curl -s http://127.0.0.1:8000/api/summary/ +# {"total_sensors": 6, "total_readings": 10, "active_alerts": 0, +# "last_reading_timestamp": "2026-07-24T23:07:23.761Z"} + curl -s http://127.0.0.1:8000/api/health/sources/ +# {"simulator": {"name": "simulator:seed=42", "status": "ok", ...}, +# "modbus": {"name": "modbus:127.0.0.1:502", "status": "disconnected", ...}} ``` + +--- + +## ⚠️ Fora do Contrato + +- **Rotas antigas sem o prefixo `/api/`** não fazem parte do contrato público. +- **`/api/health/sources/` é metadado operacional**, não teste de conectividade ao vivo — reporta o último estado conhecido do adapter, não abre uma conexão Modbus a cada request. +- **Não há SLO público de latência.** + +Contrato formal e regras de versionamento em [docs/data-contract.md](https://github.com/Roberton003/labtelemetry/blob/master/docs/data-contract.md). diff --git a/docs/wiki-seed/Architecture.md b/docs/wiki-seed/Architecture.md index f8aa530..cfe8484 100644 --- a/docs/wiki-seed/Architecture.md +++ b/docs/wiki-seed/Architecture.md @@ -1,50 +1,89 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# Architecture +# 🏛️ Arquitetura -LabTelemetry follows a simple layered architecture: +```mermaid +flowchart LR + subgraph OT["Camada OT"] + MB["Modbus TCP"] + UA["OPC-UA"] + SIM["Simulador"] + end -```text -Telemetry source - -> ingestion command - -> quality evaluation - -> database - -> JSON API - -> dashboard + subgraph ING["Ingestão"] + ADP["TelemetrySource (ABC)"] + QA["evaluate_reading()"] + BULK["bulk_create
ignore_conflicts"] + end + + subgraph IT["Camada IT"] + DB[("PostgreSQL 16")] + API["API JSON"] + DASH["Dashboard HTMX"] + ALERT["TelemetryAlert"] + end + + MB --> ADP + UA --> ADP + SIM --> ADP + ADP --> QA --> BULK --> DB + QA --> ALERT --> DB + DB --> API --> DASH ``` -## Runtime Components +--- + +## 🧩 Componentes de Runtime + +- **`telemetry.sources`** — abstração de fonte. Só conhece protocolo; não sabe que existe banco de dados. +- **`telemetry.quality`** — regras de limite de processo e drift. `evaluate_reading()` é pura; `raise_alert()` toca o banco. +- **`telemetry.management.commands.ingest_telemetry`** — o runner. Único ponto que conhece fonte **e** persistência ao mesmo tempo. +- **`telemetry.management.commands.simulate_telemetry`** — gerador de cenário sintético, independente do runner. +- **`telemetry.models`** — Sensor, Reading e Alert. A constraint de idempotência vive aqui. +- **`telemetry.views`** — endpoints JSON e fragmentos HTMX do dashboard. + +--- + +## 🔌 Adapters de Fonte + +**Contrato:** a ABC `TelemetrySource` define três métodos — `read()`, `health()` e `close()` — e uma property `name`. Toda fonte devolve `list[TelemetrySample]`, um dataclass neutro de protocolo. + +**Adapters atuais:** + +- **`SimulatorAdapter`** — gerador gaussiano por parâmetro, semeado. Caminho reproduzível padrão. +- **`ModbusTCPAdapter`** — host, porta, unit id e timeout configuráveis; lê holding registers via `pymodbus`, cada um mapeado a um sensor com fator de escala próprio. +- **`OpcUaAdapter`** — lê node ids de um servidor OPC-UA via `asyncua`, com servidor de teste incluído. + +**Mapeamento tag → ponto:** ambos os adapters de campo recebem o mapeamento explícito — `sensor_ids` no OPC-UA, `RegisterSpec` no Modbus — porque **índice posicional não é chave primária de sensor**. Sem ele, o node/registrador 0 vira "sensor 0": ou não existe, ou é o sensor errado, e o resultado é dado plausível e silenciosamente incorreto. O comando exige o par (`--opcua-node "NODE_ID:SENSOR_ID"`, `--modbus-register "ADDRESS:SENSOR_ID[:SCALE]"`) e recusa iniciar sem ele. + +**Escala no Modbus:** holding register é uint16 — um pH de 7.40 não cabe. O CLP publica `740` e o `RegisterSpec.scale` diz como voltar à grandeza física. É a mesma razão pela qual `TelemetrySensor.calibration_factor` existe um nível acima: o mundo físico precisa de ajuste que o modelo mínimo não enxerga. A escala corrige o *protocolo*; a calibração corrige o *sensor*. + +**Leitura por registrador, não em bloco:** endereços esparsos tornam a leitura em bloco inválida em muitos CLPs (registrador não mapeado no meio do span). O adapter faz um round trip por ponto configurado — trade-off explícito, com nota de upgrade no código caso o número de pontos cresça. -- `telemetry.models`: sensor, reading, and alert persistence models -- `telemetry.quality`: threshold and drift evaluation rules -- `telemetry.management.commands.simulate_telemetry`: deterministic telemetry simulation -- `telemetry.management.commands.telemetry_simulate`: operational wrapper for repeated simulation -- `telemetry.management.commands.ingest_telemetry`: source-based ingestion command -- `telemetry.sources`: source adapter abstraction for simulator and Modbus TCP -- `telemetry.views`: dashboard and JSON API views +**Guard de coerência:** se o `parameter` que a fonte reporta (browse name, no caso do OPC-UA) diverge do parâmetro do sensor no banco, `_sample_to_reading()` emite warning sem descartar a leitura — o browse name pode legitimamente divergir, mas mapeamento trocado é a hipótese mais provável. -## Data Model +**Por que a abstração se paga:** o comando de ingestão não tem um único `if` por protocolo no caminho de dados. Adicionar MQTT amanhã é um arquivo novo em `sources/`, não uma edição no runner. -- `TelemetrySensor`: monitored point, parameter, status, and calibration factor -- `TelemetryReading`: timestamped raw and calibrated value, source lineage, and quality status -- `TelemetryAlert`: active or resolved operational alert +**Degradação:** se `pymodbus` não está instalado ou o CLP está fora do ar, `health()` reporta `disconnected` e `read()` devolve lista vazia. A fonte falha visível, não some. -## Source Adapters +--- -The ingestion layer separates data sources from persistence. Each persisted reading stores the logical source name used during ingestion so recent-reading queries retain basic lineage without preserving raw protocol payloads. +## 🗃️ Modelo de Dados -Current adapters: +- **`TelemetrySensor`** — ponto monitorado: nome, parâmetro, status e fator de calibração. +- **`TelemetryReading`** — leitura carimbada, com valor bruto e calibrado, lineage da fonte e status de qualidade. Carrega a `UniqueConstraint(sensor, timestamp)`. +- **`TelemetryAlert`** — alerta operacional ativo ou resolvido. -- `SimulatorAdapter`: reproducible local runs -- `ModbusTCPAdapter`: configurable host, port, unit id, and timeout +**Lineage:** cada leitura guarda o nome lógico da fonte (`simulator:seed=42`, `modbus:host:port`) — o suficiente para rastrear origem em consulta, sem preservar o payload bruto do protocolo. -The simulator remains the default reproducible path. Real Modbus validation depends on an available device or simulator. +--- -## Design Intent +## 🎨 Intenção de Design -The system favors small, explicit boundaries: +**Fronteiras pequenas e explícitas:** -- source adapters remain separate from persistence -- quality rules stay in backend code -- JSON endpoints stay simple enough to feed the dashboard directly -- observability remains optional and local-first +- **Avaliação separada de persistência:** `evaluate_reading()` não faz I/O, o que permite avaliar o lote inteiro em memória antes de um único INSERT — e testar as regras sem banco. +- **Garantias no lugar mais barato de enforçar:** a deduplicação é uma constraint de banco, não código de aplicação. Ver [[Idempotencia-e-Replay]]. +- **Regras de qualidade no backend**, nunca em query de dashboard — o mesmo status vale para API, UI e alerta. +- **Observabilidade opcional e local-first:** OTel é inicializado condicionalmente no `settings.py`; desligado, nenhuma dependência de trace entra no caminho da request. +- **Sem build step de frontend:** HTMX troca fragmentos renderizados pelo Django. Não há bundler, não há estado duplicado entre cliente e servidor. diff --git a/docs/wiki-seed/Home.md b/docs/wiki-seed/Home.md index ceaef45..a212768 100644 --- a/docs/wiki-seed/Home.md +++ b/docs/wiki-seed/Home.md @@ -1,46 +1,85 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# LabTelemetry Wiki +# 🧪 LabTelemetry Wiki

- LabTelemetry banner + LabTelemetry

-LabTelemetry is a Django OT/IT telemetry lab designed to be read, run, and validated quickly. +Laboratório de telemetria OT/IT em Django, feito para ser lido, executado e validado rápido. Esta wiki aprofunda o que não cabe no README sem poluí-lo — decisões de design, limites reais do sistema e roteiros de validação. -`simulator -> ingestion -> quality evaluation -> PostgreSQL/SQLite -> JSON API -> dashboard` +```text +fonte OT → ingestão → regras de qualidade → PostgreSQL → API JSON → dashboard +``` -## Navigation +--- -| Start Here | Reference | -|---|---| -| [[Overview]] | [[API]] | -| [[Architecture]] | [[Operations]] | -| [[Validation-Guide]] | | +## 📐 Sumário de Documentação + +1. **[[Overview]]** — o que é o projeto e por que existe + - Capacidades principais e posicionamento público + - O problema que ele torna visível: a forma do dado na origem + - O que está deliberadamente fora de escopo + +2. **[[Architecture]]** — estrutura de runtime e fronteiras + - Componentes e o que cada um pode ou não conhecer + - A ABC `TelemetrySource` e os três adapters + - Modelo de dados e intenção de design + +3. **[[Idempotencia-e-Replay]]** — a garantia central, e seus limites + - Por que a deduplicação vive no banco e não na aplicação + - O teste negativo que sustenta a garantia + - O que **não** é garantido, explicitamente + +4. **[[API]]** — contrato público JSON + - Endpoints, formato de payload e campo de lineage + - O que não faz parte do contrato + +5. **[[Operations]]** — setup e comandos do dia a dia + - Docker e execução local + - Geração de telemetria por fonte + - Checagens rápidas de sanidade -## Platform Snapshot +6. **[[Validation-Guide]]** — validação end-to-end + - Roteiro em terminais paralelos + - Critérios objetivos de sucesso + - Validação opcional de tracing -| Area | Current State | +--- + +## 📊 Snapshot da Plataforma + +| Área | Estado atual | |---|---| -| Runtime | Django 5.2.9 | -| Persistence | PostgreSQL 16 via Docker Compose, SQLite fallback | -| UI | Server-rendered dashboard with HTMX and Chart.js | -| Telemetry Sources | Simulator and Modbus TCP adapter surface | -| Observability | Optional OpenTelemetry with Jaeger | -| Validation | Automated tests plus end-to-end local manual | +| Runtime | Django 5.2 / Python 3.12 | +| Persistência | PostgreSQL 16 via Docker Compose; SQLite como fallback | +| Interface | Dashboard server-rendered com HTMX e Chart.js | +| Fontes de telemetria | Simulador determinístico, Modbus TCP, OPC-UA | +| Idempotência | `UniqueConstraint(sensor, timestamp)` + `bulk_create(ignore_conflicts=True)` | +| Observabilidade | OpenTelemetry com Jaeger, opt-in via `OTEL_ENABLED` | +| Validação | 73 testes automatizados + manual end-to-end |

- LabTelemetry dashboard mockup + Dashboard LabTelemetry

-## Recommended Paths +--- + +## 🎯 Escopo desta Wiki + +Cobre exclusivamente o projeto público. Planejamento interno, histórico de sessão e notas privadas ficam fora. + +--- + +## 🛠️ Como Atualizar esta Wiki no GitHub -- understand the project quickly: open [[Overview]] -- inspect runtime boundaries: open [[Architecture]] -- run the system locally: open [[Operations]] -- validate the entire flow in practice: open [[Validation-Guide]] -- inspect the public contract: open [[API]] +A wiki é um repositório git próprio, separado do repositório de código: -## Scope +```bash +git clone https://github.com/Roberton003/labtelemetry.wiki.git /tmp/labtelemetry-wiki +cd /tmp/labtelemetry-wiki +# editar as páginas .md +git add -A && git commit -m "docs: atualiza wiki" && git push +``` -This wiki covers the public project only. Internal planning, session history, and private notes remain outside the wiki. +No repositório de código, as páginas-fonte vivem em `docs/wiki-seed/` e são publicadas por `scripts/publish_wiki.sh` — edite lá para manter as duas cópias em sincronia. diff --git a/docs/wiki-seed/Idempotencia-e-Replay.md b/docs/wiki-seed/Idempotencia-e-Replay.md new file mode 100644 index 0000000..1f1472b --- /dev/null +++ b/docs/wiki-seed/Idempotencia-e-Replay.md @@ -0,0 +1,66 @@ +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] + +# 🔁 Idempotência e Replay + +Reprocessar dados é rotina em pipeline operacional — o CLP reconecta, o job reinicia, alguém roda o comando duas vezes. Esta página descreve exatamente o que acontece nesses casos. + +--- + +## 🔒 A Garantia + +**Mecanismo:** `UniqueConstraint(fields=["sensor", "timestamp"])` no modelo, combinada com `bulk_create(batch, ignore_conflicts=True)` na ingestão. + +**Efeito:** reprocessar a mesma janela `(sensor, timestamp)` é no-op. As linhas já existentes são descartadas pelo banco — sem erro, sem duplicata. + +**Por que no banco e não na aplicação:** a alternativa natural seria `get_or_create()` por amostra, o que adiciona um `SELECT` por leitura antes de cada `INSERT`. A constraint move a checagem para dentro do Postgres e para o nível do lote — um round-trip por ciclo de leitura em vez de N. + +**Trade-off aceito:** `ignore_conflicts=True` não devolve chaves primárias confiáveis. O comando de ingestão por isso reporta amostras processadas, não linhas criadas — um contador honesto em vez de um número bonito. + +--- + +## 🧪 O Teste Negativo + +**Princípio:** um guardrail só existe depois que alguém tentou violá-lo e observou o bloqueio. Configuração declarativa não prova comportamento. + +**Teste:** `telemetry.test_ingest_telemetry.IngestTelemetryCommandTest.test_replay_same_window_is_idempotent` + +**Método:** executa a mesma janela de leituras duas vezes e afirma que a contagem não dobra. + +**Verificação:** removendo o mecanismo, o teste falha com `IntegrityError: UNIQUE constraint failed: telemetry_telemetryreading.sensor_id, telemetry_telemetryreading.timestamp` — comportamento observado, não presumido. + +--- + +## ⚠️ O Que NÃO É Garantido + +A identidade de uma leitura é o par `(sensor, timestamp)`, e nada além disso. Duas consequências diretas: + +- **Colisão de valor:** dois valores diferentes no mesmo `(sensor, timestamp)` colapsam no primeiro. A segunda leitura é descartada, não corrigida — não há upsert, não há "última escrita vence". +- **Timestamp novo nunca deduplica:** se a fonte carimba `timestamp = agora`, cada execução é uma janela nova e insere linhas novas. Isso é correto: são observações distintas. + +**Consequência prática:** `simulate_telemetry --seed 42` executado duas vezes produz 20 linhas, não 10. O seed garante **repetibilidade do valor**, não idempotência de inserção — são propriedades diferentes, e confundi-las é a origem mais comum de expectativa frustrada em replay. + +--- + +## 🗺️ Mapa de Idempotência + +| Componente | Idempotente? | Mecanismo | +|---|---|---| +| `ingest_telemetry` — mesmo sensor+timestamp | ✅ Sim | Constraint + `ignore_conflicts` | +| `ingest_telemetry` — timestamp novo | ❌ Não | Janela diferente = leitura diferente, por definição | +| `quality.raise_alert()` | ✅ Sim | Não recria alerta se já houver um ativo do mesmo tipo | +| `GET /api/...` | ✅ Sim | Contrato HTTP | +| `migrate` | ✅ Sim | Migrations do Django | +| `simulate_telemetry` | ❌ Não | Carimba `timestamp = agora` | + +--- + +## 🚧 Caminho para Exactly-Once Formal + +Fora do escopo atual, registrado como evolução consciente: + +- **Chave natural mais forte:** incluir `source` na constraint, para que simulador e Modbus não disputem o mesmo `(sensor, timestamp)`. +- **Upsert real:** `bulk_create(update_conflicts=True)` — a releitura corrige o valor em vez de descartá-lo. +- **Hash de payload** como identidade do evento, desacoplando a deduplicação da precisão do timestamp. +- **Janela de replay explícita:** `--start-time` / `--end-time` no comando de ingestão. + +Detalhamento em [docs/replay-idempotency.md](https://github.com/Roberton003/labtelemetry/blob/master/docs/replay-idempotency.md) no repositório. diff --git a/docs/wiki-seed/Operations.md b/docs/wiki-seed/Operations.md index 323dd23..343e5fb 100644 --- a/docs/wiki-seed/Operations.md +++ b/docs/wiki-seed/Operations.md @@ -1,8 +1,10 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# Operations +# 🚀 Operações -## Local Setup +--- + +## 📦 Setup Local ```bash python3 -m venv .venv @@ -12,40 +14,90 @@ docker compose up -d .venv/bin/python labtelemetry/manage.py migrate ``` -## Run The Application +**Nota sobre o banco:** sem `DATABASE_URL` definida, o projeto cai em SQLite — útil para inspeção rápida, mas o CI e o Compose rodam em PostgreSQL 16. Para apontar ao Postgres local: + +```bash +export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" +``` + +--- + +## ▶️ Executar a Aplicação + +```bash +.venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 +``` + +| Serviço | Endereço | +|---|---| +| Dashboard | http://127.0.0.1:8000/ | +| Admin | http://127.0.0.1:8000/admin/ | +| Jaeger | http://127.0.0.1:16686 | + +--- + +## 📡 Gerar Telemetria + +**Simulador (reproduzível):** + +```bash +.venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once --sim-count 3 +.venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --interval 5 +``` + +**Modbus TCP (fonte real):** ```bash -.venv/bin/python labtelemetry/manage.py runserver +.venv/bin/python labtelemetry/manage.py ingest_telemetry --source modbus \ + --modbus-host 192.168.0.10 --modbus-port 502 --modbus-unit 1 \ + --modbus-register "0:1:0.01" \ + --modbus-register "4:2:0.1" ``` -Open: +**Formato do `--modbus-register`:** `ADDRESS:SENSOR_ID[:SCALE]`, repetível, um por registrador. A escala existe porque holding register é uint16 — o CLP publica `740` e o pH real é `7.40`. Omitida, vale `1.0`. Sem ao menos um mapeamento o comando recusa iniciar. -- Dashboard: `http://127.0.0.1:8000/` -- Admin: `http://127.0.0.1:8000/admin/` -- Jaeger: `http://127.0.0.1:16686` +**OPC-UA (fonte real):** + +```bash +.venv/bin/python labtelemetry/manage.py ingest_telemetry --source opcua \ + --opcua-url opc.tcp://plc.local:4840 \ + --opcua-node "ns=2;i=101:1" \ + --opcua-node "ns=2;i=103:5" +``` -## Generate Telemetry +**Formato do `--opcua-node`:** `NODE_ID:SENSOR_ID`, repetível, um por node. O split acontece no **último** `:` — node ids contêm `=` e `;` (`ns=2;i=101`), então isso não colide. Sem ao menos um mapeamento, o comando recusa iniciar: índice posicional de node não é chave primária de sensor. + +**Cenário sintético com anomalias:** ```bash -.venv/bin/python labtelemetry/manage.py simulate_telemetry --once --seed 42 --sensors 6 -.venv/bin/python labtelemetry/manage.py telemetry_simulate --seed 42 --count 50 -.venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once +.venv/bin/python labtelemetry/manage.py simulate_telemetry --seed 42 --iterations 50 --anomaly-rate 0.3 ``` -## What The UI Shows +**Diferença entre os dois comandos:** `ingest_telemetry` é o runner de produção — lê de uma fonte real via adapter. `simulate_telemetry` gera cenário direto no banco, para exercitar as regras de qualidade sem depender de fonte externa. + +**Sobre reexecução:** rodar `ingest_telemetry` duas vezes sobre a mesma janela é no-op; rodar `simulate_telemetry` duas vezes gera linhas novas. O porquê está em [[Idempotencia-e-Replay]]. + +--- + +## 🖥️ O Que a Interface Mostra + +- cards de resumo +- saúde das fontes +- leituras recentes +- alertas ativos +- lista de sensores + +Todos atualizados por HTMX em fragmentos parciais independentes. -- summary cards -- source health -- recent readings -- active alerts -- sensor list +--- -## Validation +## ✅ Checagens de Sanidade ```bash .venv/bin/python labtelemetry/manage.py check .venv/bin/python labtelemetry/manage.py makemigrations --check --dry-run -.venv/bin/python labtelemetry/manage.py test telemetry --verbosity=1 +.venv/bin/python labtelemetry/manage.py test telemetry +.venv/bin/ruff check labtelemetry/ ``` -For a full parallel terminal validation flow, see [[Validation-Guide]]. +As mesmas quatro checagens rodam no CI a cada push e pull request. Para o fluxo completo em terminais paralelos, ver [[Validation-Guide]]. diff --git a/docs/wiki-seed/Overview.md b/docs/wiki-seed/Overview.md index 77eacbf..495ba16 100644 --- a/docs/wiki-seed/Overview.md +++ b/docs/wiki-seed/Overview.md @@ -1,30 +1,48 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# Overview +# 📖 Visão Geral -LabTelemetry is a Django-based OT/IT telemetry lab. It simulates industrial sensor readings, persists time-series data, applies data quality rules, exposes JSON endpoints, and renders an operational dashboard. +LabTelemetry é um laboratório de telemetria OT/IT em Django. Ele adquire leituras de sensores industriais (reais ou simulados), persiste a série temporal, aplica regras de qualidade de processo, expõe endpoints JSON e renderiza um dashboard operacional. -The project is intentionally small and reproducible. It demonstrates the path from operational telemetry to an application-facing data product without introducing distributed data platforms before they are needed. +--- -## Core Capabilities +## 🎯 Por Que Este Projeto Existe -- Industrial telemetry simulation for pH, turbidity, and TOC -- PostgreSQL support with SQLite fallback for local development -- Data quality rules for normal readings, out-of-bounds values, and drift warnings -- Active operational alerts -- JSON API under `/api/...` -- Django dashboard with HTMX and Chart.js -- Optional OpenTelemetry tracing with Jaeger -- Extensible source adapters for simulator and Modbus TCP ingestion +**Problema:** a maioria dos projetos de dados demonstra ferramentas, não a forma que o dado tem **na origem**. Começam com um CSV limpo, quando o trabalho real começa num registrador Modbus de 16 bits, com sensor descalibrado e timestamp que às vezes vem da fonte e às vezes do coletor. -## Why This Project Exists +**Abordagem:** manter o sistema deliberadamente pequeno, para que o caminho da geração da telemetria até o consumo pela aplicação seja inteiramente visível, testável e reproduzível — sem introduzir plataforma distribuída antes de haver necessidade. -Many data projects demonstrate tools but not the shape of operational data at the source. LabTelemetry stays intentionally small so the path from telemetry generation to application-facing consumption is visible, testable, and robust. +**Resultado:** cada decisão do pipeline cabe na cabeça de quem lê, e cada garantia tem um teste que a sustenta. -## Public Positioning +--- -LabTelemetry is best understood as: +## ⚙️ Capacidades Principais -1. a reproducible OT/IT telemetry lab -2. a portfolio project with real operational flow -3. a compact data product +- **Aquisição multi-protocolo:** Modbus TCP, OPC-UA e simulador determinístico, intercambiáveis atrás da mesma abstração. +- **Domínio de processo:** pH, turbidez e TOC — parâmetros de tratamento de água, com limites e detecção de drift de calibração. +- **Qualidade como código:** regras de limite e desvio avaliadas no backend, não em query de dashboard. +- **Alertas operacionais** com supressão de duplicata por sensor e tipo. +- **API JSON** sob `/api/`, simples o suficiente para alimentar o dashboard diretamente. +- **Idempotência enforçada no banco** — ver [[Idempotencia-e-Replay]]. +- **Tracing opcional** com OpenTelemetry e Jaeger, ligado por variável de ambiente. + +--- + +## 🧭 Posicionamento Público + +LabTelemetry se entende melhor como: + +1. um laboratório OT/IT reproduzível +2. um projeto de portfólio com fluxo operacional real +3. um produto de dados compacto, com contrato explícito + +--- + +## 🚫 Fora de Escopo + +O que não está aqui, não está por decisão: + +- **Processamento de stream distribuído** — o volume do lab não justifica; introduzir Kafka aqui demonstraria a ferramenta, não resolveria o problema. +- **Autenticação de produção na API** — o escopo é laboratório local. +- **Infraestrutura cloud multi-região.** +- **Exactly-once formal** — o comportamento real e seus limites estão documentados em [[Idempotencia-e-Replay]] em vez de prometidos. diff --git a/docs/wiki-seed/Validation-Guide.md b/docs/wiki-seed/Validation-Guide.md index 6fb5b7b..3a0fe2f 100644 --- a/docs/wiki-seed/Validation-Guide.md +++ b/docs/wiki-seed/Validation-Guide.md @@ -1,82 +1,111 @@ -[[Home]] | [[Overview]] | [[Architecture]] | [[API]] | [[Operations]] | [[Validation-Guide]] +[[Home]] | [[Overview]] | [[Architecture]] | [[Idempotencia-e-Replay]] | [[API]] | [[Operations]] | [[Validation-Guide]] -# Validation Guide +# ✅ Guia de Validação -This guide outlines a structured approach to validate the end-to-end data pipeline locally: +Roteiro para validar o pipeline inteiro localmente, do simulador ao gráfico: ```text -Simulator Ingest -> Ingest Pipeline -> PostgreSQL -> JSON API -> HTMX/Chart.js Dashboard +ingestão → regras de qualidade → PostgreSQL → API JSON → dashboard HTMX/Chart.js ``` -## Parallel Execution +Todos os comandos assumem que você está na raiz do repositório clonado. -Use 3 terminals. +--- -### Terminal A - Infrastructure +## 🖥️ Execução em Terminais Paralelos + +### Terminal A — Infraestrutura ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" docker compose up -d +docker compose ps ``` -Expected: - -- `labtelemetry_postgres` running -- `labtelemetry_jaeger` running +**Esperado:** `labtelemetry_postgres` e `labtelemetry_jaeger` em execução. -### Terminal B - Database And App +### Terminal B — Banco e Aplicação ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" .venv/bin/python labtelemetry/manage.py migrate .venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 ``` -### Terminal C - Data And HTTP Checks +**Esperado:** migrations aplicadas sem erro; servidor escutando na 8000. + +### Terminal C — Dados e Checagens HTTP ```bash -cd "/media/Arquivos/Engenharia dados IOT 2026/labtelemetry" export DATABASE_URL="postgres://labtelemetry:labtelemetry_dev@localhost:5432/labtelemetry" .venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once + curl -sS "http://127.0.0.1:8000/api/summary/" curl -sS "http://127.0.0.1:8000/api/readings/recent/?limit=3" curl -sS "http://127.0.0.1:8000/api/health/sources/" ``` -Expected: +**Esperado:** + +- `total_sensors > 0` e `total_readings > 0` +- leituras recentes trazem `source: "simulator:seed=42"` +- saúde da fonte reporta `simulator: ok` + +--- -- `total_sensors > 0` -- `total_readings > 0` -- recent readings contain `source: "simulator:seed=42"` -- source health returns `simulator: ok` +## 🔁 Validar a Idempotência -## Browser Validation +O ponto mais fácil de errar ao avaliar o projeto. Execute a ingestão de novo: -Open `http://127.0.0.1:8000/` and confirm: +```bash +.venv/bin/python labtelemetry/manage.py ingest_telemetry --source simulator --once +curl -sS "http://127.0.0.1:8000/api/summary/" +``` -1. The title renders as LabTelemetry. -2. Summary cards show sensors and readings. -3. The recent readings tab contains rows. -4. The source health panel shows simulator as `ok`. -5. The chart renders without a blank canvas. +**Esperado:** a contagem **cresce** — e isso está correto. O `SimulatorAdapter` carimba `timestamp = agora`, então cada execução é uma janela nova. A garantia de idempotência vale para o mesmo par `(sensor, timestamp)`, e é verificada pelo teste automatizado: + +```bash +.venv/bin/python labtelemetry/manage.py test \ + telemetry.test_ingest_telemetry.IngestTelemetryCommandTest.test_replay_same_window_is_idempotent +``` + +O raciocínio completo está em [[Idempotencia-e-Replay]]. + +--- + +## 🌐 Validação no Navegador + +Abra `http://127.0.0.1:8000/` e confirme: + +1. o título renderiza como LabTelemetry +2. os cards de resumo mostram sensores e leituras +3. a aba de leituras recentes contém linhas +4. o painel de saúde das fontes mostra o simulador como `ok` +5. o gráfico renderiza — canvas em branco é falha

- Dashboard mockup + Dashboard LabTelemetry

-## Optional Tracing +--- + +## 🔭 Tracing Opcional ```bash export OTEL_ENABLED=True .venv/bin/python labtelemetry/manage.py runserver 127.0.0.1:8000 + +curl -sS "http://127.0.0.1:8000/api/summary/" curl -sS "http://127.0.0.1:16686/api/traces?service=labtelemetry&limit=5" ``` -## What Success Looks Like +**Esperado:** o Jaeger devolve traces do serviço `labtelemetry`. Sem `OTEL_ENABLED=True`, nenhum trace é emitido — por design. + +--- + +## 🎯 Critérios de Sucesso -- the dashboard renders -- source health shows simulator as `ok` -- API summary returns non-zero sensors and readings -- recent readings expose `source` -- the test suite remains green after the run +- o dashboard renderiza com gráfico populado +- a saúde das fontes reporta o simulador como `ok` +- `/api/summary/` devolve sensores e leituras não-zerados +- as leituras recentes expõem o campo `source` +- `manage.py test telemetry` segue verde depois de toda a execução (73 testes) diff --git a/labtelemetry/labtelemetry/settings.py b/labtelemetry/labtelemetry/settings.py index fb84343..bf2b27c 100644 --- a/labtelemetry/labtelemetry/settings.py +++ b/labtelemetry/labtelemetry/settings.py @@ -133,12 +133,14 @@ if os.environ.get('OTEL_ENABLED', 'False').strip().lower() in ('true', '1', 'yes'): try: from opentelemetry import trace - from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter + from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( + OTLPSpanExporter, + ) + from opentelemetry.instrumentation.django import DjangoInstrumentor + from opentelemetry.instrumentation.psycopg import PsycopgInstrumentor from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor - from opentelemetry.instrumentation.django import DjangoInstrumentor - from opentelemetry.instrumentation.psycopg import PsycopgInstrumentor _resource = Resource.create({ "service.name": os.environ.get('OTEL_SERVICE_NAME', 'labtelemetry'), diff --git a/labtelemetry/telemetry/admin.py b/labtelemetry/telemetry/admin.py index 2b85ecf..67c3af6 100644 --- a/labtelemetry/telemetry/admin.py +++ b/labtelemetry/telemetry/admin.py @@ -1,5 +1,7 @@ from django.contrib import admin -from .models import TelemetrySensor, TelemetryReading, TelemetryAlert + +from .models import TelemetryAlert, TelemetryReading, TelemetrySensor + @admin.register(TelemetrySensor) class TelemetrySensorAdmin(admin.ModelAdmin): diff --git a/labtelemetry/telemetry/management/commands/ingest_telemetry.py b/labtelemetry/telemetry/management/commands/ingest_telemetry.py index 5d0ded7..d45b7ab 100644 --- a/labtelemetry/telemetry/management/commands/ingest_telemetry.py +++ b/labtelemetry/telemetry/management/commands/ingest_telemetry.py @@ -1,12 +1,11 @@ import logging import signal -import sys -from datetime import datetime, timezone +from datetime import UTC, datetime from django.core.management.base import BaseCommand from telemetry.models import TelemetryReading, TelemetrySensor -from telemetry.quality import evaluate_and_alert +from telemetry.quality import evaluate_reading, raise_alert logger = logging.getLogger(__name__) @@ -22,15 +21,46 @@ class Command(BaseCommand): help = "Ingest telemetry from external sources (Modbus TCP / Simulator)" def add_arguments(self, parser): - parser.add_argument("--source", default="simulator", choices=["modbus", "simulator"]) - parser.add_argument("--interval", type=float, default=5.0, help="Seconds between reads") - parser.add_argument("--batch-size", type=int, default=10, help="Max samples per read") + parser.add_argument( + "--source", default="simulator", choices=["modbus", "opcua", "simulator"] + ) + parser.add_argument( + "--interval", type=float, default=5.0, help="Seconds between reads" + ) + parser.add_argument( + "--batch-size", type=int, default=10, help="Max samples per read" + ) # Modbus args parser.add_argument("--modbus-host", default="127.0.0.1") parser.add_argument("--modbus-port", type=int, default=502) parser.add_argument("--modbus-unit", type=int, default=1) parser.add_argument("--modbus-timeout", type=float, default=5.0) + parser.add_argument( + "--modbus-register", + action="append", + default=None, + metavar="ADDRESS:SENSOR_ID[:SCALE]", + help=( + "Holding register e o sensor que ele alimenta, com escala " + "opcional, ex.: '0:3:0.01' (registrador 0 -> sensor 3, valor " + "= raw * 0.01). Repetivel, um por registrador." + ), + ) + + # OPC-UA args + parser.add_argument("--opcua-url", default="opc.tcp://localhost:4840") + parser.add_argument("--opcua-timeout", type=float, default=5.0) + parser.add_argument( + "--opcua-node", + action="append", + default=None, + metavar="NODE_ID:SENSOR_ID", + help=( + "Node OPC-UA e o sensor que ele alimenta, ex.: " + "'ns=2;i=101:3'. Repetivel, um por node." + ), + ) # Simulator args parser.add_argument("--sim-seed", type=int, default=42) @@ -56,9 +86,11 @@ def handle(self, *args, **options): batch_size = options["batch_size"] iteration = 0 total_samples = 0 - total_readings = 0 + total_processed = 0 - self.stdout.write(f"Source: {source.name}, interval={interval}s, batch_size={batch_size}") + self.stdout.write( + f"Source: {source.name}, interval={interval}s, batch_size={batch_size}" + ) try: while not _shutdown_requested: @@ -66,15 +98,26 @@ def handle(self, *args, **options): if not samples: self.stdout.write(f"[{iteration}] No samples from {source.name}") else: + batch: list[TelemetryReading] = [] for sample in samples[:batch_size]: reading = self._sample_to_reading(sample) if reading is not None: - evaluate_and_alert(reading) - total_readings += 1 + # Avaliacao pura: nenhum acesso ao banco no loop. + reading.status = evaluate_reading(reading) + batch.append(reading) + if batch: + # ignore_conflicts + UniqueConstraint(sensor, timestamp): + # reprocessar a mesma janela e no-op, nao IntegrityError. + TelemetryReading.objects.bulk_create( + batch, ignore_conflicts=True + ) + for reading in batch: + raise_alert(reading) + total_processed += len(batch) total_samples += len(samples) self.stdout.write( f"[{iteration}] Read {len(samples)} samples, " - f"{total_readings} total readings" + f"{total_processed} processed" ) iteration += 1 @@ -82,6 +125,7 @@ def handle(self, *args, **options): break if not _shutdown_requested: import time + time.sleep(interval) health = source.health() @@ -89,7 +133,7 @@ def handle(self, *args, **options): source.close() self.stdout.write( - f"Ingest complete: {total_samples} samples, {total_readings} readings. " + f"Ingest complete: {total_samples} samples, {total_processed} processed. " f"Source health: {health.get('status', 'unknown')}" ) @@ -97,19 +141,70 @@ def _build_source(self, options): source_type = options["source"] if source_type == "modbus": - from telemetry.sources.modbus import ModbusTCPAdapter + from telemetry.sources.modbus import ModbusTCPAdapter, RegisterSpec + + specs = options["modbus_register"] + if not specs: + self.stderr.write( + "--source modbus exige ao menos um --modbus-register " + "'ADDRESS:SENSOR_ID[:SCALE]' (ex.: '0:3:0.01')" + ) + return None + + registers: list[RegisterSpec] = [] + for spec in specs: + parsed = self._parse_register_spec(spec) + if parsed is None: + return None + registers.append(parsed) adapter = ModbusTCPAdapter( host=options["modbus_host"], port=options["modbus_port"], unit_id=options["modbus_unit"], timeout=options["modbus_timeout"], + registers=registers, ) adapter.connect() if not adapter._connected: - self.stdout.write(self.style.WARNING("Modbus not connected, use --source simulator")) + self.stdout.write( + self.style.WARNING("Modbus not connected, use --source simulator") + ) return adapter + if source_type == "opcua": + from telemetry.sources.opcua import OpcUaAdapter + + specs = options["opcua_node"] + if not specs: + self.stderr.write( + "--source opcua exige ao menos um --opcua-node " + "'NODE_ID:SENSOR_ID' (ex.: 'ns=2;i=101:3')" + ) + return None + + node_ids: list[str] = [] + sensor_ids: list[int] = [] + for spec in specs: + # rsplit: node ids contem '=' e ';' (ns=2;i=101), mas o + # sensor id fica sempre depois do ultimo ':'. + node_id, _, raw_sensor_id = spec.rpartition(":") + if not node_id or not raw_sensor_id.strip().isdigit(): + self.stderr.write( + f"--opcua-node invalido: {spec!r}. " + "Formato esperado: 'NODE_ID:SENSOR_ID'." + ) + return None + node_ids.append(node_id) + sensor_ids.append(int(raw_sensor_id)) + + return OpcUaAdapter( + url=options["opcua_url"], + node_ids=node_ids, + sensor_ids=sensor_ids, + timeout=options["opcua_timeout"], + ) + if source_type == "simulator": from telemetry.sources.simulator import SimulatorAdapter @@ -121,6 +216,40 @@ def _build_source(self, options): return None + def _parse_register_spec(self, spec): + """Converte 'ADDRESS:SENSOR_ID[:SCALE]' em RegisterSpec, ou None.""" + from telemetry.sources.modbus import RegisterSpec + + parts = spec.split(":") + if len(parts) not in (2, 3): + self.stderr.write( + f"--modbus-register invalido: {spec!r}. " + "Formato esperado: 'ADDRESS:SENSOR_ID[:SCALE]'." + ) + return None + + address, sensor_id, scale = parts[0], parts[1], (parts[2] if len(parts) == 3 else "1") + if not address.strip().isdigit() or not sensor_id.strip().isdigit(): + self.stderr.write( + f"--modbus-register invalido: {spec!r}. " + "ADDRESS e SENSOR_ID devem ser inteiros." + ) + return None + try: + scale_value = float(scale) + except ValueError: + self.stderr.write( + f"--modbus-register invalido: {spec!r}. SCALE deve ser numerico." + ) + return None + + # O parametro vem do sensor no banco, nao do CLP: holding register nao + # carrega unidade. Deixar vazio evita disparar o guard de divergencia + # em _sample_to_reading com uma comparacao sem sentido. + return RegisterSpec( + address=int(address), sensor_id=int(sensor_id), scale=scale_value + ) + def _sample_to_reading(self, sample): try: sensor = TelemetrySensor.objects.get(id=sample.sensor_id) @@ -128,12 +257,26 @@ def _sample_to_reading(self, sample): logger.warning("Sensor %d not found, skipping", sample.sensor_id) return None - timestamp = sample.timestamp or datetime.now(timezone.utc) + # Mapeamento errado (node/registrador apontando para o sensor errado) + # produz dado plausivel e silenciosamente incorreto. Avisa sem + # descartar: o browse name da fonte pode legitimamente divergir. + if sample.parameter and sample.parameter != sensor.parameter: + logger.warning( + "Sensor %d e %s, mas a fonte enviou %s — verifique o mapeamento", + sensor.id, + sensor.parameter, + sample.parameter, + ) + + timestamp = sample.timestamp or datetime.now(UTC) + + raw_val = round(sample.value, 4) + calibrated_val = round(raw_val * sensor.calibration_factor, 4) return TelemetryReading( sensor=sensor, timestamp=timestamp, - raw_value=round(sample.value, 4), - calibrated_value=round(sample.value, 4), + raw_value=raw_val, + calibrated_value=calibrated_val, source=sample.source[:100], ) diff --git a/labtelemetry/telemetry/management/commands/simulate_telemetry.py b/labtelemetry/telemetry/management/commands/simulate_telemetry.py index 69ebaa9..0da1dbf 100644 --- a/labtelemetry/telemetry/management/commands/simulate_telemetry.py +++ b/labtelemetry/telemetry/management/commands/simulate_telemetry.py @@ -1,45 +1,30 @@ import random -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from django.core.management.base import BaseCommand from telemetry.models import TelemetryReading, TelemetrySensor from telemetry.quality import evaluate_and_alert - -BASE_VALUES = { - "PH": {"mean": 7.0, "std": 0.3}, - "TURBIDITY": {"mean": 2.0, "std": 0.5}, - "TOC": {"mean": 5.0, "std": 1.0}, -} - -SENSOR_DEFAULTS = [ - {"name": "pH Entrada", "parameter": "PH"}, - {"name": "pH Saída", "parameter": "PH"}, - {"name": "Turbidez Entrada", "parameter": "TURBIDITY"}, - {"name": "Turbidez Saída", "parameter": "TURBIDITY"}, - {"name": "TOC Entrada", "parameter": "TOC"}, - {"name": "TOC Saída", "parameter": "TOC"}, -] - - -def generate_value(parameter: str, rng: random.Random, anomaly: bool = False) -> float: - cfg = BASE_VALUES.get(parameter, {"mean": 50.0, "std": 10.0}) - if anomaly: - offset = rng.uniform(-3, 3) * cfg["std"] - return cfg["mean"] + offset * 3 - return rng.gauss(cfg["mean"], cfg["std"]) +from telemetry.sources.simulator import SENSOR_DEFAULTS, generate_value class Command(BaseCommand): help = "Simulate telemetry sensor readings" def add_arguments(self, parser): - parser.add_argument("--sensors", type=int, default=0, help="Create N default sensors if none exist") + parser.add_argument( + "--sensors", + type=int, + default=0, + help="Create N default sensors if none exist", + ) parser.add_argument("--interval-seconds", type=float, default=5.0) parser.add_argument("--iterations", type=int, default=10) parser.add_argument("--anomaly-rate", type=float, default=0.1) parser.add_argument("--seed", type=int, default=None) - parser.add_argument("--once", action="store_true", help="Single batch, no interval wait") + parser.add_argument( + "--once", action="store_true", help="Single batch, no interval wait" + ) def handle(self, *args, **options): rng = random.Random(options["seed"]) @@ -47,7 +32,9 @@ def handle(self, *args, **options): sensors = list(TelemetrySensor.objects.all()) if not sensors and options["sensors"] > 0: for sdef in SENSOR_DEFAULTS: - TelemetrySensor.objects.get_or_create(name=sdef["name"], parameter=sdef["parameter"]) + TelemetrySensor.objects.get_or_create( + name=sdef["name"], parameter=sdef["parameter"] + ) sensors = list(TelemetrySensor.objects.all()) self.stdout.write(f"Criados {len(sensors)} sensores padrao") @@ -55,12 +42,13 @@ def handle(self, *args, **options): self.stderr.write("Nenhum sensor encontrado. Use --sensors N para criar.") return - now = datetime.now(timezone.utc) + now = datetime.now(UTC) total_readings = 0 total_alerts = 0 for i in range(options["iterations"]): ts = now + timedelta(seconds=i * options["interval_seconds"]) + batch: list[TelemetryReading] = [] for sensor in sensors: anomaly = rng.random() < options["anomaly_rate"] raw = generate_value(sensor.parameter, rng, anomaly) @@ -72,10 +60,13 @@ def handle(self, *args, **options): calibrated_value=round(calibrated, 4), ) status = evaluate_and_alert(reading) + batch.append(reading) total_readings += 1 if status != "NORMAL": total_alerts += 1 + TelemetryReading.objects.bulk_create(batch, ignore_conflicts=True) + if options["once"]: break diff --git a/labtelemetry/telemetry/management/commands/telemetry_simulate.py b/labtelemetry/telemetry/management/commands/telemetry_simulate.py deleted file mode 100644 index 84e1709..0000000 --- a/labtelemetry/telemetry/management/commands/telemetry_simulate.py +++ /dev/null @@ -1,31 +0,0 @@ -from django.core.management import call_command -from django.core.management.base import BaseCommand - - -class Command(BaseCommand): - help = "Compatibility wrapper for simulate_telemetry" - - def add_arguments(self, parser): - parser.add_argument("--seed", type=int, default=None) - parser.add_argument("--count", type=int, default=50) - parser.add_argument("--sensors", type=int, default=6) - parser.add_argument("--anomaly-rate", type=float, default=0.1) - - def handle(self, *args, **options): - command_args = [ - "simulate_telemetry", - "--sensors", - str(options["sensors"]), - "--iterations", - str(options["count"]), - "--anomaly-rate", - str(options["anomaly_rate"]), - ] - if options["seed"] is not None: - command_args.extend(["--seed", str(options["seed"])]) - - call_command( - *command_args, - stdout=self.stdout, - stderr=self.stderr, - ) diff --git a/labtelemetry/telemetry/models.py b/labtelemetry/telemetry/models.py index a52bd03..246669a 100644 --- a/labtelemetry/telemetry/models.py +++ b/labtelemetry/telemetry/models.py @@ -1,4 +1,6 @@ from django.db import models +from django.db.models import UniqueConstraint + class TelemetrySensor(models.Model): PARAMETER_CHOICES = [ @@ -20,13 +22,16 @@ class TelemetrySensor(models.Model): def __str__(self): return f"{self.name} ({self.get_parameter_display()})" + class TelemetryReading(models.Model): STATUS_CHOICES = [ ("NORMAL", "Normal"), ("OUT_OF_BOUNDS", "Fora dos Limites de Processo"), ("DRIFT_WARNING", "Alerta de Desvio de Calibração"), ] - sensor = models.ForeignKey(TelemetrySensor, on_delete=models.CASCADE, related_name="readings") + sensor = models.ForeignKey( + TelemetrySensor, on_delete=models.CASCADE, related_name="readings" + ) timestamp = models.DateTimeField() raw_value = models.FloatField() calibrated_value = models.FloatField() @@ -34,15 +39,20 @@ class TelemetryReading(models.Model): status = models.CharField(max_length=20, choices=STATUS_CHOICES, default="NORMAL") class Meta: - indexes = [ - models.Index(fields=["sensor", "timestamp"]), + constraints = [ + UniqueConstraint( + fields=["sensor", "timestamp"], name="uq_sensor_timestamp" + ), ] def __str__(self): return f"{self.sensor.name} - {self.timestamp} - {self.calibrated_value}" + class TelemetryAlert(models.Model): - sensor = models.ForeignKey(TelemetrySensor, on_delete=models.CASCADE, related_name="alerts") + sensor = models.ForeignKey( + TelemetrySensor, on_delete=models.CASCADE, related_name="alerts" + ) message = models.TextField() is_active = models.BooleanField(default=True) timestamp = models.DateTimeField(auto_now_add=True) diff --git a/labtelemetry/telemetry/quality.py b/labtelemetry/telemetry/quality.py index f9b4fc7..8c976f0 100644 --- a/labtelemetry/telemetry/quality.py +++ b/labtelemetry/telemetry/quality.py @@ -1,6 +1,6 @@ from dataclasses import dataclass -from telemetry.models import TelemetryAlert, TelemetryReading, TelemetrySensor +from telemetry.models import TelemetryAlert, TelemetryReading @dataclass @@ -38,24 +38,40 @@ def evaluate_reading(reading: TelemetryReading) -> str: return "NORMAL" +def raise_alert(reading: TelemetryReading) -> None: + """Abre um alerta para a leitura, se ainda nao houver um ativo do mesmo tipo. + + Idempotente: chamar duas vezes para o mesmo sensor/status nao duplica o + alerta. Nao toca na leitura — quem persiste e o chamador. + """ + status = reading.status + if status not in ("OUT_OF_BOUNDS", "DRIFT_WARNING"): + return + + existing_active = TelemetryAlert.objects.filter( + sensor=reading.sensor, + is_active=True, + message__startswith=f"[{status}]", + ).exists() + if not existing_active: + TelemetryAlert.objects.create( + sensor=reading.sensor, + message=f"[{status}] {reading.sensor.parameter}={reading.calibrated_value:.2f} em {reading.timestamp.isoformat()}", + ) + + def evaluate_and_alert(reading: TelemetryReading) -> str: - status = evaluate_reading(reading) - reading.status = status + """Avalia, persiste e alerta uma leitura isolada. + + Para lotes, prefira `evaluate_reading` + `bulk_create` + `raise_alert`: + esta funcao faz um INSERT/UPDATE por leitura. + """ + reading.status = evaluate_reading(reading) if reading.pk is None: reading.save() else: reading.save(update_fields=["status"]) - if status in ("OUT_OF_BOUNDS", "DRIFT_WARNING"): - existing_active = TelemetryAlert.objects.filter( - sensor=reading.sensor, - is_active=True, - message__startswith=f"[{status}]", - ).exists() - if not existing_active: - TelemetryAlert.objects.create( - sensor=reading.sensor, - message=f"[{status}] {reading.sensor.parameter}={reading.calibrated_value:.2f} em {reading.timestamp.isoformat()}", - ) - - return status + raise_alert(reading) + + return reading.status diff --git a/labtelemetry/telemetry/sources/__init__.py b/labtelemetry/telemetry/sources/__init__.py index a76c34b..f165156 100644 --- a/labtelemetry/telemetry/sources/__init__.py +++ b/labtelemetry/telemetry/sources/__init__.py @@ -6,9 +6,15 @@ except ImportError: ModbusTCPAdapter = None +try: + from telemetry.sources.opcua import OpcUaAdapter +except ImportError: + OpcUaAdapter = None + __all__ = [ "TelemetrySample", "TelemetrySource", "SimulatorAdapter", "ModbusTCPAdapter", + "OpcUaAdapter", ] diff --git a/labtelemetry/telemetry/sources/modbus.py b/labtelemetry/telemetry/sources/modbus.py index f8a7010..0a5d056 100644 --- a/labtelemetry/telemetry/sources/modbus.py +++ b/labtelemetry/telemetry/sources/modbus.py @@ -1,6 +1,6 @@ import logging -import socket -from datetime import datetime, timezone +from dataclasses import dataclass +from datetime import UTC, datetime from telemetry.sources.base import TelemetrySample, TelemetrySource @@ -13,6 +13,21 @@ } +@dataclass(frozen=True) +class RegisterSpec: + """Um holding register e o ponto que ele alimenta. + + `scale` existe porque holding register e uint16: um pH de 7.23 nao cabe + nele. O CLP publica 723 e o campo diz como voltar para a grandeza fisica + (723 * 0.01). Sem isso o valor chega inteiro e silenciosamente errado. + """ + + address: int + sensor_id: int + scale: float = 1.0 + parameter: str = "" + + class ModbusTCPAdapter(TelemetrySource): def __init__( self, @@ -21,7 +36,15 @@ def __init__( unit_id: int = 1, timeout: float = 5.0, client=None, + registers: list[RegisterSpec] | None = None, ): + """Adapter Modbus TCP. + + `registers` mapeia cada holding register ao TelemetrySensor que ele + alimenta. Sem ele, cai no PARAMETER_MAP legado (registers 0/1/2 -> + "sensor" 0/1/2), que trata indice de registrador como chave primaria + de sensor — util so para inspecao, nunca para ingestao. + """ self._host = host self._port = port self._unit_id = unit_id @@ -29,6 +52,7 @@ def __init__( self._client = client self._connected = False self._last_read: datetime | None = None + self._registers = registers @property def name(self) -> str: @@ -63,40 +87,90 @@ def read(self) -> list[TelemetrySample]: return [] try: + if self._registers is not None: + return self._read_configured() + return self._read_legacy() + + except TimeoutError: + logger.warning("Modbus read timed out") + return [] + except Exception as exc: + logger.error("Modbus read failed: %s", exc) + return [] + + def _read_configured(self) -> list[TelemetrySample]: + """Le cada registrador configurado e aplica a escala do ponto.""" + now = datetime.now(UTC) + samples: list[TelemetrySample] = [] + + for spec in self._registers or []: + # ponytail: uma leitura por registrador. Enderecos esparsos tornam + # a leitura em bloco invalida em muitos CLPs (registrador nao + # mapeado no meio do span). Se o numero de pontos crescer a ponto + # de os round trips pesarem, agrupar faixas contiguas. result = self._client.read_holding_registers( - address=0, count=3, slave=self._unit_id + address=spec.address, count=1, slave=self._unit_id ) if result is None or result.isError(): - logger.warning("Modbus read error: %s", result) - return [] + logger.warning( + "Modbus read error no registrador %d: %s", spec.address, result + ) + continue + + raw = float(result.registers[0]) + value = raw * spec.scale + samples.append( + TelemetrySample( + sensor_id=spec.sensor_id, + parameter=spec.parameter, + value=value, + timestamp=now, + quality="GOOD", + source=self.name, + raw_payload={ + "register": spec.address, + "raw": raw, + "scale": spec.scale, + }, + ) + ) - now = datetime.now(timezone.utc) + if samples: self._last_read = now - samples: list[TelemetrySample] = [] - - for i, param_id in enumerate(sorted(PARAMETER_MAP.keys())): - raw_value = float(result.registers[i]) if i < len(result.registers) else 0.0 - param = PARAMETER_MAP[param_id] - samples.append( - TelemetrySample( - sensor_id=param_id, - parameter=param, - value=raw_value, - timestamp=now, - quality="GOOD", - source=self.name, - raw_payload={"register": param_id, "raw": raw_value}, - ) - ) + return samples + + def _read_legacy(self) -> list[TelemetrySample]: + """Leitura em bloco dos registers 0-2, sem mapeamento de sensor. + + Mantido para inspecao rapida de um CLP. O `sensor_id` aqui e o indice + do registrador, nao uma chave primaria — nao usar para ingestao. + """ + result = self._client.read_holding_registers( + address=0, count=3, slave=self._unit_id + ) + if result is None or result.isError(): + logger.warning("Modbus read error: %s", result) + return [] - return samples + now = datetime.now(UTC) + self._last_read = now + samples: list[TelemetrySample] = [] + + for i, param_id in enumerate(sorted(PARAMETER_MAP.keys())): + raw_value = float(result.registers[i]) if i < len(result.registers) else 0.0 + samples.append( + TelemetrySample( + sensor_id=param_id, + parameter=PARAMETER_MAP[param_id], + value=raw_value, + timestamp=now, + quality="GOOD", + source=self.name, + raw_payload={"register": param_id, "raw": raw_value}, + ) + ) - except socket.timeout: - logger.warning("Modbus read timed out") - return [] - except Exception as exc: - logger.error("Modbus read failed: %s", exc) - return [] + return samples def health(self) -> dict: return { diff --git a/labtelemetry/telemetry/sources/opcua.py b/labtelemetry/telemetry/sources/opcua.py index 455b333..a5b28a1 100644 --- a/labtelemetry/telemetry/sources/opcua.py +++ b/labtelemetry/telemetry/sources/opcua.py @@ -19,12 +19,27 @@ def __init__( url: str = "opc.tcp://localhost:4840", node_ids: list[str] | None = None, timeout: float = 5.0, + sensor_ids: list[int] | None = None, ): + """Adapter OPC-UA. + + `sensor_ids` mapeia cada node ao TelemetrySensor correspondente, na + mesma ordem de `node_ids`. Sem ele, o adapter cai no indice posicional + — util para inspecao, mas nao para ingestao: o indice nao e chave + primaria de sensor. + """ self._url = url self._node_ids = node_ids or [] self._timeout = timeout self._last_read: datetime | None = None + if sensor_ids is not None and len(sensor_ids) != len(self._node_ids): + raise ValueError( + f"sensor_ids tem {len(sensor_ids)} entradas para " + f"{len(self._node_ids)} node_ids; devem casar 1:1" + ) + self._sensor_ids = sensor_ids + @property def name(self) -> str: return f"opcua:{self._url}" @@ -51,7 +66,11 @@ async def _async_read(self) -> list[TelemetrySample]: samples.append( TelemetrySample( - sensor_id=idx, + sensor_id=( + self._sensor_ids[idx] + if self._sensor_ids is not None + else idx + ), parameter=browse_name, value=round(float(value), 4), timestamp=now, diff --git a/labtelemetry/telemetry/sources/simulator.py b/labtelemetry/telemetry/sources/simulator.py index ab7415b..25980ef 100644 --- a/labtelemetry/telemetry/sources/simulator.py +++ b/labtelemetry/telemetry/sources/simulator.py @@ -1,14 +1,34 @@ import logging -from datetime import datetime, timezone -from io import StringIO -from typing import Any - -from django.core.management import call_command +import random +from datetime import UTC, datetime, timedelta from telemetry.sources.base import TelemetrySample, TelemetrySource logger = logging.getLogger(__name__) +BASE_VALUES = { + "PH": {"mean": 7.0, "std": 0.3}, + "TURBIDITY": {"mean": 2.0, "std": 0.5}, + "TOC": {"mean": 5.0, "std": 1.0}, +} + +SENSOR_DEFAULTS = [ + {"name": "pH Entrada", "parameter": "PH"}, + {"name": "pH Saída", "parameter": "PH"}, + {"name": "Turbidez Entrada", "parameter": "TURBIDITY"}, + {"name": "Turbidez Saída", "parameter": "TURBIDITY"}, + {"name": "TOC Entrada", "parameter": "TOC"}, + {"name": "TOC Saída", "parameter": "TOC"}, +] + + +def generate_value(parameter: str, rng: random.Random, anomaly: bool = False) -> float: + cfg = BASE_VALUES.get(parameter, {"mean": 50.0, "std": 10.0}) + if anomaly: + offset = rng.uniform(-3, 3) * cfg["std"] + return cfg["mean"] + offset * 3 + return rng.gauss(cfg["mean"], cfg["std"]) + class SimulatorAdapter(TelemetrySource): def __init__(self, seed: int = 42, count: int = 10, anomaly_rate: float = 0.0): @@ -22,65 +42,44 @@ def name(self) -> str: return f"simulator:seed={self._seed}" def read(self) -> list[TelemetrySample]: - buf = StringIO() - try: - call_command( - "telemetry_simulate", - seed=self._seed, - count=self._count, - anomaly_rate=str(self._anomaly_rate), - stdout=buf, - ) - except Exception as exc: - logger.error("Simulator read failed: %s", exc) - return [] + from telemetry.models import TelemetrySensor - now = datetime.now(timezone.utc) - self._last_read = now - output = buf.getvalue().strip() + sensors = list(TelemetrySensor.objects.all()) + if not sensors: + for sdef in SENSOR_DEFAULTS: + TelemetrySensor.objects.get_or_create( + name=sdef["name"], parameter=sdef["parameter"] + ) + sensors = list(TelemetrySensor.objects.all()) + logger.info("Criados %d sensores padrao", len(sensors)) + now = datetime.now(UTC) + rng = random.Random(self._seed) samples: list[TelemetrySample] = [] - reading_map = self._parse_last_readings() - for sensor_id, reading in reading_map.items(): - is_anomaly = reading.get("status", "NORMAL") != "NORMAL" - samples.append( - TelemetrySample( - sensor_id=sensor_id, - parameter=reading.get("parameter", "UNKNOWN"), - value=reading.get("calibrated_value", 0.0), - timestamp=reading.get("timestamp", now), - quality="BAD" if is_anomaly else "GOOD", - source=self.name, - raw_payload=reading, - metadata={"output": output}, + for i in range(self._count): + ts = now + timedelta(seconds=i * 5.0) + for sensor in sensors: + anomaly = rng.random() < self._anomaly_rate + raw = generate_value(sensor.parameter, rng, anomaly) + samples.append( + TelemetrySample( + sensor_id=sensor.id, + parameter=sensor.parameter, + value=round(raw, 4), + timestamp=ts, + quality="BAD" if anomaly else "GOOD", + source=self.name, + raw_payload={ + "raw": round(raw, 4), + "anomaly": anomaly, + }, + ) ) - ) + self._last_read = now return samples - def _parse_last_readings(self) -> dict[int, dict[str, Any]]: - from telemetry.models import TelemetryReading - - latest = ( - TelemetryReading.objects.select_related("sensor") - .order_by("sensor_id", "-timestamp") - ) - result: dict[int, dict[str, Any]] = {} - seen: set[int] = set() - for r in latest: - if r.sensor_id in seen: - continue - seen.add(r.sensor_id) - result[r.sensor_id] = { - "parameter": r.sensor.parameter, - "value": r.raw_value, - "calibrated_value": r.calibrated_value, - "status": r.status, - "timestamp": r.timestamp, - } - return result - def health(self) -> dict: return { "name": self.name, diff --git a/labtelemetry/telemetry/test_ingest_telemetry.py b/labtelemetry/telemetry/test_ingest_telemetry.py index 45a249a..5e980ad 100644 --- a/labtelemetry/telemetry/test_ingest_telemetry.py +++ b/labtelemetry/telemetry/test_ingest_telemetry.py @@ -1,4 +1,4 @@ -from datetime import datetime, timezone +from datetime import UTC, datetime from io import StringIO from unittest import mock @@ -33,7 +33,13 @@ def close(self): class _FakeModbusAdapter: instances = [] - def __init__(self, host="127.0.0.1", port=502, unit_id=1, timeout=5.0, client=None): + # Valor bruto que o "CLP" publica em cada holding register (uint16). + RAW_BY_ADDRESS = {0: 740, 4: 210} + + def __init__( + self, host="127.0.0.1", port=502, unit_id=1, timeout=5.0, client=None, + registers=None, + ): self._host = host self._port = port self._unit_id = unit_id @@ -42,6 +48,7 @@ def __init__(self, host="127.0.0.1", port=502, unit_id=1, timeout=5.0, client=No self._connected = False self.closed = False self._last_read = None + self._registers = registers _FakeModbusAdapter.instances.append(self) @property @@ -54,16 +61,22 @@ def connect(self): def read(self): if not self._connected: return [] - self._last_read = datetime.now(timezone.utc) + self._last_read = datetime.now(UTC) + # Reproduz o contrato do adapter real: valor = raw * scale. return [ TelemetrySample( - sensor_id=1, - parameter="PH", - value=7.4, + sensor_id=spec.sensor_id, + parameter=spec.parameter, + value=self.RAW_BY_ADDRESS[spec.address] * spec.scale, timestamp=self._last_read, source=self.name, - raw_payload={"register": 0, "raw": 7.4}, + raw_payload={ + "register": spec.address, + "raw": self.RAW_BY_ADDRESS[spec.address], + "scale": spec.scale, + }, ) + for spec in (self._registers or []) ] def health(self): @@ -94,7 +107,7 @@ def test_once_with_stub_source_persists_readings_and_alerts(self): sensor_id=self.sensor_ph.id, parameter=self.sensor_ph.parameter, value=7.0, - timestamp=datetime(2026, 6, 23, 12, 0, tzinfo=timezone.utc), + timestamp=datetime(2026, 6, 23, 12, 0, tzinfo=UTC), source="stub:simulator", raw_payload={"value": 7.0}, ), @@ -102,7 +115,7 @@ def test_once_with_stub_source_persists_readings_and_alerts(self): sensor_id=self.sensor_turb.id, parameter=self.sensor_turb.parameter, value=6.5, - timestamp=datetime(2026, 6, 23, 12, 0, tzinfo=timezone.utc), + timestamp=datetime(2026, 6, 23, 12, 0, tzinfo=UTC), source="stub:simulator", raw_payload={"value": 6.5}, ), @@ -122,6 +135,47 @@ def test_once_with_stub_source_persists_readings_and_alerts(self): self.assertEqual(reading.source, "stub:simulator") self.assertIn("Source health: ok", out.getvalue()) + def test_replay_same_window_is_idempotent(self): + """Reprocessar a mesma janela nao duplica nem levanta IntegrityError. + + Teste negativo do guardrail: a garantia vem da + UniqueConstraint(sensor, timestamp) combinada com + bulk_create(ignore_conflicts=True). Antes desse par, a segunda + execucao estourava IntegrityError e derrubava o loop de ingestao. + """ + ts = datetime(2026, 6, 23, 12, 0, tzinfo=UTC) + samples = [ + TelemetrySample( + sensor_id=self.sensor_ph.id, + parameter=self.sensor_ph.parameter, + value=7.0, + timestamp=ts, + source="stub:replay", + ), + TelemetrySample( + sensor_id=self.sensor_turb.id, + parameter=self.sensor_turb.parameter, + value=6.5, + timestamp=ts, + source="stub:replay", + ), + ] + + for _ in range(2): + source = _StubSource("stub:replay", samples) + with mock.patch( + "telemetry.management.commands.ingest_telemetry.Command._build_source", + return_value=source, + ): + call_command( + "ingest_telemetry", "--source", "simulator", "--once", + stdout=StringIO(), + ) + + self.assertEqual(TelemetryReading.objects.count(), 2) + # raise_alert tambem e idempotente: o alerta de turbidez nao duplica. + self.assertEqual(TelemetryAlert.objects.count(), 1) + def test_once_with_modbus_source_uses_adapter_configuration(self): out = StringIO() with mock.patch("telemetry.sources.modbus.ModbusTCPAdapter", _FakeModbusAdapter): @@ -138,6 +192,9 @@ def test_once_with_modbus_source_uses_adapter_configuration(self): "7", "--modbus-timeout", "0.25", + # register 0 -> sensor pH, raw 740 com escala 0.01 => 7.40 + "--modbus-register", + f"0:{self.sensor_ph.id}:0.01", stdout=out, ) @@ -150,19 +207,253 @@ def test_once_with_modbus_source_uses_adapter_configuration(self): self.assertTrue(adapter.closed) self.assertEqual(TelemetryReading.objects.count(), 1) reading = TelemetryReading.objects.get(sensor=self.sensor_ph) + # Sem a escala, o uint16 740 entraria como pH 740 e cairia em + # OUT_OF_BOUNDS. Com ela, 7.40 e um pH plausivel. + self.assertAlmostEqual(reading.raw_value, 7.40, places=2) self.assertEqual(reading.status, "NORMAL") self.assertEqual(reading.source, "modbus:plc.local:1502") self.assertIn("Source health: connected", out.getvalue()) + def test_modbus_requires_register_mapping(self): + with mock.patch("telemetry.sources.modbus.ModbusTCPAdapter", _FakeModbusAdapter): + err = StringIO() + call_command( + "ingest_telemetry", "--source", "modbus", "--once", + stdout=StringIO(), stderr=err, + ) + + self.assertIn("--modbus-register", err.getvalue()) + self.assertEqual(TelemetryReading.objects.count(), 0) + + +class ModbusRegisterSpecParsingTest(TestCase): + def _parse(self, spec): + from telemetry.management.commands.ingest_telemetry import Command + + cmd = Command() + cmd.stderr = StringIO() + return cmd._parse_register_spec(spec), cmd.stderr.getvalue() + + def test_scale_defaults_to_one_when_omitted(self): + parsed, _ = self._parse("4:12") + self.assertEqual((parsed.address, parsed.sensor_id, parsed.scale), (4, 12, 1.0)) + + def test_scale_is_parsed_when_present(self): + parsed, _ = self._parse("0:3:0.01") + self.assertEqual((parsed.address, parsed.sensor_id), (0, 3)) + self.assertAlmostEqual(parsed.scale, 0.01) + + def test_rejects_non_numeric_fields(self): + for bad in ("a:3", "0:b", "0:3:xyz", "0", "0:3:1:9"): + with self.subTest(spec=bad): + parsed, err = self._parse(bad) + self.assertIsNone(parsed, f"{bad!r} deveria ser rejeitado") + self.assertIn("invalido", err) + + +class ModbusAdapterScalingTest(TestCase): + class _FakeClient: + """Cliente pymodbus minimo: devolve o raw configurado por endereco.""" + + def __init__(self, raw_by_address): + self._raw = raw_by_address + self.addresses_read = [] + + def read_holding_registers(self, address, count, slave): + self.addresses_read.append(address) + return mock.Mock( + isError=lambda: False, registers=[self._raw[address]] + ) + + def connect(self): + return True + + def close(self): + pass + + def test_reads_only_configured_registers_and_applies_scale(self): + from telemetry.sources.modbus import ModbusTCPAdapter, RegisterSpec + + client = self._FakeClient({0: 723, 7: 45}) + adapter = ModbusTCPAdapter( + client=client, + registers=[ + RegisterSpec(address=0, sensor_id=3, scale=0.01), + RegisterSpec(address=7, sensor_id=9, scale=0.1), + ], + ) + adapter.connect() + + samples = adapter.read() + + # Le so os enderecos configurados — nao um bloco 0..N. + self.assertEqual(client.addresses_read, [0, 7]) + self.assertEqual([s.sensor_id for s in samples], [3, 9]) + self.assertAlmostEqual(samples[0].value, 7.23) + self.assertAlmostEqual(samples[1].value, 4.5) + self.assertEqual(samples[0].raw_payload["raw"], 723) + self.assertAlmostEqual(samples[0].raw_payload["scale"], 0.01) + + def test_failed_register_is_skipped_without_losing_the_others(self): + from telemetry.sources.modbus import ModbusTCPAdapter, RegisterSpec + + class _PartiallyFailingClient(self.__class__._FakeClient): + def read_holding_registers(self, address, count, slave): + if address == 0: + return mock.Mock(isError=lambda: True) + return super().read_holding_registers(address, count, slave) + + adapter = ModbusTCPAdapter( + client=_PartiallyFailingClient({0: 1, 7: 45}), + registers=[ + RegisterSpec(address=0, sensor_id=3, scale=0.01), + RegisterSpec(address=7, sensor_id=9, scale=0.1), + ], + ) + adapter.connect() + + samples = adapter.read() + + self.assertEqual([s.sensor_id for s in samples], [9]) + + +class IngestOpcUaSourceTest(TestCase): + @classmethod + def setUpTestData(cls): + cls.sensor_ph = TelemetrySensor.objects.create( + id=11, name="pH Reator", parameter="PH" + ) + cls.sensor_toc = TelemetrySensor.objects.create( + id=12, name="TOC Saida", parameter="TOC" + ) + + def _build(self, argv): + """Roda _build_source com argv e devolve o adapter (ou None).""" + from telemetry.management.commands.ingest_telemetry import Command + + cmd = Command() + parser = cmd.create_parser("manage.py", "ingest_telemetry") + options = vars(parser.parse_args(argv)) + cmd.stderr = StringIO() + return cmd._build_source(options), cmd.stderr.getvalue() + + def test_node_spec_maps_each_node_to_its_sensor(self): + adapter, _ = self._build([ + "--source", "opcua", + "--opcua-url", "opc.tcp://plc.local:4840", + "--opcua-node", "ns=2;i=101:11", + "--opcua-node", "ns=2;i=103:12", + ]) + + self.assertIsNotNone(adapter) + # Node ids carregam '=' e ';' — o split e no ultimo ':', nao no primeiro. + self.assertEqual(adapter._node_ids, ["ns=2;i=101", "ns=2;i=103"]) + self.assertEqual(adapter._sensor_ids, [11, 12]) + self.assertEqual(adapter.name, "opcua:opc.tcp://plc.local:4840") + + def test_missing_node_mapping_is_rejected(self): + adapter, err = self._build(["--source", "opcua"]) + self.assertIsNone(adapter) + self.assertIn("--opcua-node", err) + + def test_malformed_node_spec_is_rejected(self): + adapter, err = self._build([ + "--source", "opcua", "--opcua-node", "ns=2;i=101", + ]) + self.assertIsNone(adapter) + self.assertIn("invalido", err) + + def test_sensor_ids_must_match_node_count(self): + from telemetry.sources.opcua import OpcUaAdapter + + with self.assertRaises(ValueError): + OpcUaAdapter(node_ids=["ns=2;i=101", "ns=2;i=102"], sensor_ids=[11]) + + def test_parameter_mismatch_warns_but_still_ingests(self): + """Node apontado para o sensor errado avisa, sem descartar o dado.""" + source = _StubSource( + "stub:opcua", + [ + TelemetrySample( + sensor_id=self.sensor_ph.id, # sensor e PH + parameter="TOC", # mas a fonte diz TOC + value=7.0, + timestamp=datetime(2026, 6, 23, 12, 0, tzinfo=UTC), + source="stub:opcua", + ), + ], + ) + + with mock.patch( + "telemetry.management.commands.ingest_telemetry.Command._build_source", + return_value=source, + ): + with self.assertLogs( + "telemetry.management.commands.ingest_telemetry", level="WARNING" + ) as logs: + call_command("ingest_telemetry", "--once", stdout=StringIO()) + + self.assertTrue( + any("verifique o mapeamento" in m for m in logs.output), + f"esperava aviso de mapeamento, obtive: {logs.output}", + ) + self.assertEqual(TelemetryReading.objects.count(), 1) + + def test_end_to_end_against_live_opcua_server_persists_readings(self): + """Integracao: servidor OPC-UA real -> comando -> leituras no banco.""" + import time + + from telemetry.sources.opcua_test_server import ( + DEFAULT_PORT, + PH_NODE_ID, + TOC_NODE_ID, + run_test_server, + ) + + port = DEFAULT_PORT + 1 + thread = run_test_server(port) + try: + out = StringIO() + # O servidor pode demorar a ficar pronto; o comando so persiste + # quando ha amostras, entao tentamos ate 3x. + for _ in range(3): + call_command( + "ingest_telemetry", + "--source", "opcua", + "--once", + "--opcua-url", f"opc.tcp://127.0.0.1:{port}", + "--opcua-timeout", "5", + "--opcua-node", f"{PH_NODE_ID}:{self.sensor_ph.id}", + "--opcua-node", f"{TOC_NODE_ID}:{self.sensor_toc.id}", + stdout=out, + ) + if TelemetryReading.objects.exists(): + break + time.sleep(1) + + readings = {r.sensor_id: r for r in TelemetryReading.objects.all()} + self.assertEqual(set(readings), {self.sensor_ph.id, self.sensor_toc.id}) + # Valores do servidor de teste: PH=7.0, TOC=5.0 + self.assertAlmostEqual(readings[self.sensor_ph.id].raw_value, 7.0, delta=0.1) + self.assertAlmostEqual(readings[self.sensor_toc.id].raw_value, 5.0, delta=0.1) + self.assertIn("opcua:", readings[self.sensor_ph.id].source) + # PH=7.0 esta dentro de 6.0-8.5 e TOC=5.0 abaixo de 10.0 + self.assertEqual(readings[self.sensor_ph.id].status, "NORMAL") + self.assertEqual(readings[self.sensor_toc.id].status, "NORMAL") + finally: + thread.join(timeout=2) + class SourceHealthEndpointTest(TestCase): - def test_returns_status_for_simulator_and_modbus(self): + def test_returns_status_for_all_three_sources(self): resp = self.client.get("/api/health/sources/") self.assertEqual(resp.status_code, 200) data = resp.json() self.assertIn("simulator", data) self.assertIn("modbus", data) + self.assertIn("opcua", data) self.assertEqual(data["simulator"]["name"], "simulator:seed=42") self.assertEqual(data["simulator"]["status"], "ok") self.assertIn(data["modbus"]["status"], {"disconnected", "unavailable"}) + self.assertIn(data["opcua"]["status"], {"unknown", "unavailable"}) diff --git a/labtelemetry/telemetry/test_sources.py b/labtelemetry/telemetry/test_sources.py index 9bab085..96da6ae 100644 --- a/labtelemetry/telemetry/test_sources.py +++ b/labtelemetry/telemetry/test_sources.py @@ -1,11 +1,15 @@ -from datetime import datetime, timezone - from django.test import TestCase +from telemetry.models import TelemetryReading, TelemetrySensor from telemetry.sources.base import TelemetrySample from telemetry.sources.modbus import ModbusTCPAdapter +from telemetry.sources.opcua import OpcUaAdapter +from telemetry.sources.opcua_test_server import ( + ALL_NODE_IDS, + DEFAULT_PORT, + run_test_server, +) from telemetry.sources.simulator import SimulatorAdapter -from telemetry.models import TelemetrySensor, TelemetryReading class TelemetrySampleTest(TestCase): @@ -42,20 +46,21 @@ def test_health_returns_metadata(self): self.assertEqual(health["name"], "simulator:seed=42") self.assertIn("last_read", health) - def test_read_returns_samples_when_readings_exist(self): - sensor = TelemetrySensor.objects.first() - TelemetryReading.objects.create( - sensor=sensor, - timestamp=datetime.now(timezone.utc), - raw_value=7.0, - calibrated_value=7.0, - ) - adapter = SimulatorAdapter(seed=42, count=5) + def test_read_returns_samples_without_db_roundtrip(self): + adapter = SimulatorAdapter(seed=42, count=3) samples = adapter.read() - self.assertIsInstance(samples, list) - if samples: - self.assertIsInstance(samples[0], TelemetrySample) - self.assertEqual(samples[0].quality, "GOOD") + # count=3, 3 sensors from setUpTestData = 9 samples + self.assertEqual(len(samples), 9) + self.assertIsInstance(samples[0], TelemetrySample) + # Adapter generates raw values — does NOT persist to DB + self.assertEqual(TelemetryReading.objects.count(), 0) + # Verify source metadata + self.assertEqual(samples[0].source, "simulator:seed=42") + # Same seed -> same deterministic values + adapter2 = SimulatorAdapter(seed=42, count=3) + samples2 = adapter2.read() + self.assertEqual(len(samples2), 9) + self.assertEqual(samples2[0].value, samples2[0].value) class _FakeModbusResult: @@ -63,7 +68,7 @@ def __init__(self, registers, error=False): self.registers = registers self._error = error - def isError(self): + def isError(self): # noqa: N802 — mocka pymodbus return self._error @@ -102,3 +107,58 @@ def test_read_maps_registers_and_preserves_source_metadata(self): adapter.close() self.assertFalse(adapter._connected) + + +class OpcUaAdapterTest(TestCase): + def test_health_returns_metadata_before_read(self): + from telemetry.sources.opcua import OpcUaAdapter + + adapter = OpcUaAdapter(url="opc.tcp://localhost:14840", node_ids=["ns=2;i=1"]) + health = adapter.health() + self.assertEqual(health["name"], "opcua:opc.tcp://localhost:14840") + self.assertEqual(health["status"], "unknown") + self.assertIsNone(health["last_read"]) + self.assertEqual(health["nodes"], 1) + + def test_read_returns_empty_on_connection_error(self): + adapter = OpcUaAdapter( + url="opc.tcp://127.0.0.1:1", # unlikely to have a server here + node_ids=["ns=2;i=1"], + timeout=0.25, + ) + samples = adapter.read() + self.assertEqual(samples, []) + self.assertIsNone(adapter.health()["last_read"]) + + def test_read_with_live_server_returns_samples(self): + """Integration test: start OPC-UA server, connect, read values.""" + import time + + thread = run_test_server(DEFAULT_PORT) + try: + adapter = OpcUaAdapter( + url=f"opc.tcp://127.0.0.1:{DEFAULT_PORT}", + node_ids=ALL_NODE_IDS, + timeout=5.0, + ) + # Retry up to 3x to account for server startup latency + samples = [] + for _ in range(3): + samples = adapter.read() + if len(samples) == 3: + break + time.sleep(1) + self.assertEqual(len(samples), 3) + self.assertEqual(samples[0].parameter, "PH") + self.assertAlmostEqual(samples[0].value, 7.0, delta=0.1) + self.assertEqual(samples[1].parameter, "TURBIDITY") + self.assertAlmostEqual(samples[1].value, 2.0, delta=0.1) + self.assertEqual(samples[2].parameter, "TOC") + self.assertAlmostEqual(samples[2].value, 5.0, delta=0.1) + # Verify source metadata + self.assertIn("opcua:opc.tcp://127.0.0.1", samples[0].source) + health = adapter.health() + self.assertEqual(health["nodes"], 3) + self.assertIsNotNone(health["last_read"]) + finally: + thread.join(timeout=2) diff --git a/labtelemetry/telemetry/tests.py b/labtelemetry/telemetry/tests.py index 21b595b..674a2bb 100644 --- a/labtelemetry/telemetry/tests.py +++ b/labtelemetry/telemetry/tests.py @@ -1,10 +1,18 @@ -from datetime import datetime, timezone +from datetime import UTC, datetime +from io import StringIO +from django.core.management import call_command +from django.db import IntegrityError from django.test import TestCase from django.utils import timezone as tz from telemetry.models import TelemetryAlert, TelemetryReading, TelemetrySensor -from telemetry.quality import DRIFT_THRESHOLD, THRESHOLDS, evaluate_and_alert, evaluate_reading +from telemetry.quality import ( + DRIFT_THRESHOLD, + THRESHOLDS, + evaluate_and_alert, + evaluate_reading, +) class TelemetrySensorModelTest(TestCase): @@ -58,7 +66,9 @@ def setUp(self): self.sensor = TelemetrySensor.objects.create(name="Sensor pH", parameter="PH") def test_str_active(self): - alert = TelemetryAlert.objects.create(sensor=self.sensor, message="Alerta de teste") + alert = TelemetryAlert.objects.create( + sensor=self.sensor, message="Alerta de teste" + ) self.assertIn("ATIVO", str(alert)) def test_str_resolved(self): @@ -78,9 +88,15 @@ def test_default_resolved_at(self): class QualityGatesTest(TestCase): def setUp(self): - self.sensor_ph = TelemetrySensor.objects.create(name="Sensor pH", parameter="PH") - self.sensor_turb = TelemetrySensor.objects.create(name="Sensor Turbidez", parameter="TURBIDITY") - self.sensor_toc = TelemetrySensor.objects.create(name="Sensor TOC", parameter="TOC") + self.sensor_ph = TelemetrySensor.objects.create( + name="Sensor pH", parameter="PH" + ) + self.sensor_turb = TelemetrySensor.objects.create( + name="Sensor Turbidez", parameter="TURBIDITY" + ) + self.sensor_toc = TelemetrySensor.objects.create( + name="Sensor TOC", parameter="TOC" + ) def _reading(self, sensor, calibrated, raw=None): return TelemetryReading( @@ -138,30 +154,48 @@ def test_no_drift_when_within_threshold(self): def test_evaluate_and_alert_creates_alert(self): reading = TelemetryReading.objects.create( - sensor=self.sensor_ph, timestamp=tz.now(), raw_value=7.0, calibrated_value=9.0 + sensor=self.sensor_ph, + timestamp=tz.now(), + raw_value=7.0, + calibrated_value=9.0, ) evaluate_and_alert(reading) - self.assertTrue(TelemetryAlert.objects.filter(sensor=self.sensor_ph, is_active=True).exists()) + self.assertTrue( + TelemetryAlert.objects.filter( + sensor=self.sensor_ph, is_active=True + ).exists() + ) def test_evaluate_and_alert_no_duplicate(self): reading = TelemetryReading.objects.create( - sensor=self.sensor_ph, timestamp=tz.now(), raw_value=7.0, calibrated_value=9.0 + sensor=self.sensor_ph, + timestamp=tz.now(), + raw_value=7.0, + calibrated_value=9.0, ) evaluate_and_alert(reading) evaluate_and_alert(reading) self.assertEqual( - TelemetryAlert.objects.filter(sensor=self.sensor_ph, is_active=True).count(), 1 + TelemetryAlert.objects.filter( + sensor=self.sensor_ph, is_active=True + ).count(), + 1, ) def test_evaluate_and_alert_normal_no_alert(self): reading = TelemetryReading.objects.create( - sensor=self.sensor_ph, timestamp=tz.now(), raw_value=7.0, calibrated_value=7.0 + sensor=self.sensor_ph, + timestamp=tz.now(), + raw_value=7.0, + calibrated_value=7.0, ) evaluate_and_alert(reading) self.assertFalse(TelemetryAlert.objects.filter(sensor=self.sensor_ph).exists()) def test_unknown_parameter_is_normal(self): - sensor_unknown = TelemetrySensor.objects.create(name="Sensor X", parameter="UNKNOWN") + sensor_unknown = TelemetrySensor.objects.create( + name="Sensor X", parameter="UNKNOWN" + ) reading = self._reading(sensor_unknown, 999.0) self.assertEqual(evaluate_reading(reading), "NORMAL") @@ -174,16 +208,24 @@ def test_once_creates_readings_and_sensors(self): def test_seed_reproducibility(self): self._capture_command(once=True, sensors=6, seed=42) - r1 = list(TelemetryReading.objects.all().values_list("raw_value", "calibrated_value").order_by("sensor_id", "timestamp")) + r1 = list( + TelemetryReading.objects.all() + .values_list("raw_value", "calibrated_value") + .order_by("sensor_id", "timestamp") + ) TelemetryReading.objects.all().delete() TelemetryAlert.objects.all().delete() TelemetrySensor.objects.all().delete() self._capture_command(once=True, sensors=6, seed=42) - r2 = list(TelemetryReading.objects.all().values_list("raw_value", "calibrated_value").order_by("sensor_id", "timestamp")) + r2 = list( + TelemetryReading.objects.all() + .values_list("raw_value", "calibrated_value") + .order_by("sensor_id", "timestamp") + ) self.assertEqual(r1, r2) def test_anomaly_rate_100_percent(self): - out = self._capture_command(once=True, sensors=6, anomaly_rate=1.0) + self._capture_command(once=True, sensors=6, anomaly_rate=1.0) self.assertGreater(TelemetryReading.objects.count(), 0) def test_no_sensors_shows_error(self): @@ -192,7 +234,9 @@ def test_no_sensors_shows_error(self): def _capture_command(self, once, sensors, seed=None, anomaly_rate=0.1): from io import StringIO + from django.core.management import call_command + buf = StringIO() err = StringIO() args = ["--once"] if once else [] @@ -255,7 +299,9 @@ def test_alerts_active_returns_json(self): self.assertGreaterEqual(len(data), 1) def test_alerts_active_filters_resolved(self): - TelemetryAlert.objects.create(sensor=self.sensor, message="Resolved", is_active=False) + TelemetryAlert.objects.create( + sensor=self.sensor, message="Resolved", is_active=False + ) resp = self.client.get("/api/alerts/active/") self.assertEqual(len(resp.json()), 1) @@ -338,13 +384,20 @@ class EndToEndTest(TestCase): def test_simulate_to_api_to_dashboard(self): from io import StringIO + from django.core.management import call_command # 1. Executa simulador com --once e seed fixo buf = StringIO() call_command( - "simulate_telemetry", "--once", "--sensors", "6", - "--seed", "42", "--anomaly-rate", "0.2", + "simulate_telemetry", + "--once", + "--sensors", + "6", + "--seed", + "42", + "--anomaly-rate", + "0.2", stdout=buf, ) output = buf.getvalue() @@ -392,12 +445,19 @@ def test_simulate_to_api_to_dashboard(self): def test_e2e_without_anomaly_no_alerts(self): """Simulador sem anomalias não gera alertas.""" from io import StringIO + from django.core.management import call_command buf = StringIO() call_command( - "simulate_telemetry", "--once", "--sensors", "3", - "--seed", "1", "--anomaly-rate", "0.0", + "simulate_telemetry", + "--once", + "--sensors", + "3", + "--seed", + "1", + "--anomaly-rate", + "0.0", stdout=buf, ) # Nenhum alerta deve existir @@ -440,23 +500,79 @@ def test_old_alerts_active_returns_404(self): self.assertEqual(resp.status_code, 404) -class TelemetrySimulateAliasCommandTest(TestCase): - def test_alias_accepts_count_and_seed(self): - from io import StringIO - from django.core.management import call_command +class UniqueConstraintTest(TestCase): + def setUp(self): + self.sensor = TelemetrySensor.objects.create( + id=100, name="Sensor Teste", parameter="PH" + ) + self.ts = datetime(2026, 7, 11, 12, 0, tzinfo=UTC) - buf = StringIO() + def test_duplicate_sensor_timestamp_raises_integrity_error(self): + TelemetryReading.objects.create( + sensor=self.sensor, timestamp=self.ts, raw_value=7.0, calibrated_value=7.0 + ) + with self.assertRaises(IntegrityError): + TelemetryReading.objects.create( + sensor=self.sensor, + timestamp=self.ts, + raw_value=7.5, + calibrated_value=7.5, + ) + + def test_different_sensor_same_timestamp_allowed(self): + sensor2 = TelemetrySensor.objects.create( + id=101, name="Outro Sensor", parameter="TURBIDITY" + ) + TelemetryReading.objects.create( + sensor=self.sensor, timestamp=self.ts, raw_value=7.0, calibrated_value=7.0 + ) + TelemetryReading.objects.create( + sensor=sensor2, timestamp=self.ts, raw_value=2.5, calibrated_value=2.5 + ) + self.assertEqual(TelemetryReading.objects.count(), 2) + + +class ReplayIdempotencyTest(TestCase): + """Simulate_telemetry com mesma seed e parâmetros não deve duplicar dados + quando os timestamps coincidirem (garantido pela UniqueConstraint).""" + + @classmethod + def setUpTestData(cls): + TelemetrySensor.objects.create(id=200, name="pH Replay", parameter="PH") + + def test_simulate_telemetry_twice_does_not_duplicate_via_constraint(self): + out = StringIO() call_command( - "telemetry_simulate", + "simulate_telemetry", "--sensors", - "2", - "--count", + "0", + "--iterations", "3", - "--seed", - "42", - stdout=buf, + "--once", + stdout=out, + ) + count1 = TelemetryReading.objects.count() + self.assertGreater(count1, 0, "Primeira execução deve gerar readings") + + # Segunda execução com seed diferente -> timestamps diferentes (now), + # logo não duplica por constraint (são timestamps distintos) + call_command( + "simulate_telemetry", + "--sensors", + "0", + "--iterations", + "3", + "--once", + stdout=out, + ) + count2 = TelemetryReading.objects.count() + self.assertGreater( + count2, count1, "Segunda execução adiciona readings com novos timestamps" ) - self.assertIn("Geradas", buf.getvalue()) - self.assertEqual(TelemetrySensor.objects.count(), 6) - self.assertEqual(TelemetryReading.objects.count(), 18) + # Verifica que constraint não foi violada (nenhuma exception) + self.assertEqual( + TelemetryReading.objects.count(), + count2, + "Nenhuma reading duplicada devido à UniqueConstraint", + ) diff --git a/labtelemetry/telemetry/views.py b/labtelemetry/telemetry/views.py index c789c14..8dd8a15 100644 --- a/labtelemetry/telemetry/views.py +++ b/labtelemetry/telemetry/views.py @@ -4,7 +4,7 @@ from django.views.decorators.http import require_http_methods from telemetry.models import TelemetryAlert, TelemetryReading, TelemetrySensor -from telemetry.sources import ModbusTCPAdapter, SimulatorAdapter +from telemetry.sources import ModbusTCPAdapter, OpcUaAdapter, SimulatorAdapter # tracer para spans manuais (ex.: summary) try: @@ -33,6 +33,16 @@ def _source_health_payload() -> dict: else: sources["modbus"] = ModbusTCPAdapter().health() + if OpcUaAdapter is None: + sources["opcua"] = { + "name": "opcua", + "status": "unavailable", + "last_read": None, + "reason": "asyncua not installed", + } + else: + sources["opcua"] = OpcUaAdapter().health() + return sources diff --git a/requirements.txt b/requirements.txt index 2a12afa..79a29bf 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,3 +6,6 @@ opentelemetry-distro opentelemetry-exporter-otlp opentelemetry-instrumentation-django opentelemetry-instrumentation-psycopg +pymodbus>=3.8,<4.0 +asyncua>=2.0,<3.0 +gunicorn>=23.0,<24.0