diff --git a/lab-11-rust-ingestor/Cargo.toml b/lab-11-rust-ingestor/Cargo.toml new file mode 100644 index 0000000..03113cb --- /dev/null +++ b/lab-11-rust-ingestor/Cargo.toml @@ -0,0 +1,30 @@ +[package] +name = "rust-ingestor" +version = "0.1.0" +edition = "2021" +description = "Microservicio de ingestión de telemetría GPS: MQTT → validación → Redis" + +[[bin]] +name = "ingestor" +path = "src/main.rs" + +[[bin]] +name = "publisher" +path = "src/bin/publisher.rs" + +[lib] +path = "src/lib.rs" + +[dependencies] +anyhow = "1" +chrono = { version = "0.4", default-features = false, features = ["std", "clock"] } +redis = { version = "0.27", default-features = false, features = ["tokio-comp"] } +rumqttc = "0.24" +serde = { version = "1", features = ["derive"] } +serde_json = "1" +tokio = { version = "1", features = ["full"] } +tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["env-filter"] } + +[dev-dependencies] +tokio = { version = "1", features = ["full"] } diff --git a/lab-11-rust-ingestor/Dockerfile b/lab-11-rust-ingestor/Dockerfile new file mode 100644 index 0000000..30fd1ca --- /dev/null +++ b/lab-11-rust-ingestor/Dockerfile @@ -0,0 +1,60 @@ +# ───────────────────────────────────────────────────────────────────────────── +# Etapa 1 — Compilación +# Compila todos los binarios del workspace en modo release. +# ───────────────────────────────────────────────────────────────────────────── +FROM rust:1-slim AS builder + +WORKDIR /app + +# Copiar manifiestos primero para aprovechar la caché de capas de Docker: +# si solo cambia el código fuente (no Cargo.toml/Cargo.lock), el paso de +# descarga de dependencias se reutiliza. +COPY Cargo.toml Cargo.lock* ./ + +# Crear stubs mínimos para que `cargo fetch` pueda resolver el grafo de deps. +RUN mkdir -p src/bin && \ + echo 'fn main() {}' > src/main.rs && \ + echo 'fn main() {}' > src/bin/publisher.rs && \ + printf '' > src/lib.rs && \ + printf '' > src/models.rs && \ + printf '' > src/mqtt.rs && \ + printf '' > src/redis_sink.rs + +RUN cargo build --release --bin ingestor --bin publisher 2>/dev/null || true + +# Copiar el código fuente real y recompilar solo los artefactos modificados. +COPY src/ src/ +RUN touch src/lib.rs src/main.rs src/models.rs src/mqtt.rs src/redis_sink.rs \ + src/bin/publisher.rs && \ + cargo build --release --bin ingestor --bin publisher + +# ───────────────────────────────────────────────────────────────────────────── +# Etapa 2a — Imagen del ingestor +# ───────────────────────────────────────────────────────────────────────────── +FROM debian:bookworm-slim AS ingestor + +RUN apt-get update && apt-get install -y --no-install-recommends \ + ca-certificates \ + && rm -rf /var/lib/apt/lists/* + +COPY --from=builder /app/target/release/ingestor /usr/local/bin/ingestor + +ENV MQTT_HOST=localhost \ + MQTT_PORT=1883 \ + REDIS_URL=redis://localhost:6379 \ + RUST_LOG=info + +CMD ["ingestor"] + +# ───────────────────────────────────────────────────────────────────────────── +# Etapa 2b — Imagen del publisher (simulador GPS) +# ───────────────────────────────────────────────────────────────────────────── +FROM debian:bookworm-slim AS publisher + +COPY --from=builder /app/target/release/publisher /usr/local/bin/publisher + +ENV MQTT_HOST=localhost \ + MQTT_PORT=1883 \ + PUBLISH_INTERVAL_MS=1000 + +CMD ["publisher"] diff --git a/lab-11-rust-ingestor/README.md b/lab-11-rust-ingestor/README.md new file mode 100644 index 0000000..f5f397c --- /dev/null +++ b/lab-11-rust-ingestor/README.md @@ -0,0 +1,252 @@ +# Lab 11 — Rust: Microservicio de Ingestión de Telemetría GPS + +Microservicio de alta frecuencia escrito en Rust que consume posiciones GPS de vehículos desde MQTT (NanoMQ), las valida y las publica en Redis. Demuestra el modelo de ownership, tipos seguros y concurrencia sin data races en un contexto de producción real. + +--- + +## Stack + +| Componente | Tecnología | +|---|---| +| **Lenguaje** | Rust (edition 2021) | +| **Runtime async** | Tokio | +| **MQTT client** | rumqttc (puro Rust) | +| **Redis client** | redis-rs (tokio-comp) | +| **Broker MQTT** | NanoMQ | +| **Almacén** | Redis 7 | +| **Serialización** | serde + serde_json | +| **Logging** | tracing + tracing-subscriber | +| **Infraestructura** | Docker · Docker Compose | + +--- + +## Conceptos Demostrados + +### Ownership y borrowing +`RedisSink` es dueño del `redis::Client`. Las posiciones se pasan por referencia (`&VehiclePosition`) a `publish()` — sin clonación innecesaria. + +### Tipos seguros como documentación ejecutable +`VehiclePosition` encapsula las reglas del dominio. `ValidationError` es un enum exhaustivo que el compilador obliga a manejar. Una posición inválida no puede llegar a Redis — es imposible por construcción. + +### Concurrencia sin data races +Tokio gestiona el event loop de `rumqttc` y las conexiones Redis en un runtime async. No hay `Mutex` visible ni `unsafe`. El compilador verifica en tiempo de compilación que no hay acceso concurrente inseguro. + +### Manejo de errores explícito (sin excepciones) +El operador `?` propaga errores. Las posiciones con payload inválido o fuera de rango se descartan con un log estructurado — nunca causan panic ni detienen el servicio. + +--- + +## Estructura + +``` +lab-11-rust-ingestor/ +├── Cargo.toml # Package: dos binarios + lib +├── Dockerfile # Multi-stage: builder → ingestor | publisher +├── docker-compose.yml # NanoMQ + Redis + ingestor + publisher +├── demo.sh # Ciclo completo: up → observar → mostrar Redis → down +├── src/ +│ ├── lib.rs # Exporta los módulos públicos +│ ├── models.rs # VehiclePosition + ValidationError + validate() +│ ├── mqtt.rs # create_client() + constantes POSITION_TOPIC / TOPIC_QOS +│ ├── redis_sink.rs # RedisSink: SET latest (TTL 5min) + XADD stream +│ ├── main.rs # Binario ingestor: EventLoop → validate → RedisSink +│ └── bin/ +│ └── publisher.rs # Binario simulador GPS: 3 buses, LCG, 1 pos/s +├── tests/ +│ └── integration_test.rs # Tests de validación + JSON round-trip +└── docs/ + └── architecture.md +``` + +--- + +## Requisitos del Host + +| Requisito | Detalle | +|---|---| +| Docker Desktop | Motor de contenedores | +| ~1 GB RAM libre | NanoMQ + Redis + compilación Rust | +| ~500 MB disco | Imagen `rust:1-slim` para compilación | + +> **Primera compilación**: la imagen Docker descarga el toolchain de Rust y compila el proyecto (~3-5 min). Las compilaciones siguientes reutilizan la caché de capas de Docker. + +--- + +## Inicio Rápido + +### Demo automatizado + +```bash +cd lab-11-rust-ingestor +bash demo.sh +``` + +Levanta el stack, espera 15 segundos de telemetría, muestra las posiciones almacenadas en Redis y baja el stack. + +```bash +bash demo.sh --keep # no baja el stack al terminar +``` + +### Manual + +#### 1. Levantar el stack + +```bash +cd lab-11-rust-ingestor +docker compose up -d +``` + +> La primera vez construye las imágenes Rust (compilación release). Puede tardar varios minutos. + +#### 2. Ver logs en tiempo real + +```bash +# Ingestor recibiendo y almacenando posiciones: +docker compose logs -f ingestor + +# Publisher enviando posiciones simuladas: +docker compose logs -f publisher +``` + +Salida esperada del ingestor: +``` +INFO rust_ingestor: ingestor iniciando mqtt_host=nanomq mqtt_port=1883 ... +INFO rust_ingestor: conectado al broker MQTT +INFO rust_ingestor::redis_sink: posición almacenada vehicle_id=bus-001 route_id=R01 lat=9.9281 lon=-84.0907 speed_kmh=45.0 +INFO rust_ingestor::redis_sink: posición almacenada vehicle_id=bus-002 route_id=R02 lat=9.9571 lon=-84.1309 speed_kmh=38.2 +... +``` + +#### 3. Consultar Redis + +```bash +# Última posición de cada bus: +docker compose exec redis redis-cli GET transit:vehicle:bus-001:latest + +# Todas las claves de posición: +docker compose exec redis redis-cli --scan --pattern "transit:vehicle:*" + +# Longitud del stream histórico: +docker compose exec redis redis-cli XLEN transit:positions + +# Últimas 5 entradas del stream: +docker compose exec redis redis-cli XREVRANGE transit:positions + - COUNT 5 +``` + +#### 4. Bajar el stack + +```bash +docker compose down -v # -v elimina volúmenes +``` + +--- + +## Tests + +### Tests unitarios (sin servicios externos) + +```bash +# Dentro del contenedor de build (no requiere Rust instalado en el host): +docker compose run --rm --entrypoint cargo ingestor test --lib + +# O si Rust está instalado localmente: +cargo test +``` + +Cubre 18 casos: validación (campos vacíos, rangos lat/lon/speed/bearing), JSON round-trip, formato de clave Redis. + +``` +running 18 tests +test models::tests::valid_position_passes ... ok +test models::tests::bearing_none_is_valid ... ok +test models::tests::redis_key_format ... ok +test integration_test::valid_position_passes_validation ... ok +test integration_test::empty_vehicle_id_fails ... ok +test integration_test::latitude_above_90_fails ... ok +... +test result: ok. 18 passed; 0 failed +``` + +### Tests de integración con Redis (requieren stack levantado) + +```bash +# Con el stack corriendo: +docker compose up -d redis +cargo test -- --ignored +``` + +--- + +## Flujo de Datos + +``` +publisher (Rust) + 3 buses × 1 pos/s + bus-001 (R01) lat≈9.93 lon≈-84.09 + bus-002 (R02) lat≈9.96 lon≈-84.13 + bus-003 (R03) lat≈9.90 lon≈-84.15 + │ + │ MQTT QoS1 transit/vehicles/{id}/position + ▼ + NanoMQ :1883 + │ + │ MQTT QoS1 transit/vehicles/+/position (wildcard) + ▼ + ingestor (Rust) + parse JSON → VehiclePosition + validate() → Ok / Err (descarte con log) + RedisSink.publish() + │ + │ SET ... EX 300 + XADD MAXLEN ~ 1000 + ▼ + Redis :6379 + transit:vehicle:bus-001:latest → JSON (TTL 5 min) + transit:vehicle:bus-002:latest → JSON (TTL 5 min) + transit:vehicle:bus-003:latest → JSON (TTL 5 min) + transit:positions → stream ~1000 entradas +``` + +--- + +## Puertos + +| Servicio | Puerto (host) | Descripción | +|---|---|---| +| `nanomq` | 1883 | MQTT broker | +| `nanomq` | 8083 | MQTT sobre WebSocket | +| `redis` | 6379 | Redis | + +--- + +## Variables de Entorno + +### ingestor + +| Variable | Default | Descripción | +|---|---|---| +| `MQTT_HOST` | `localhost` | Hostname del broker MQTT | +| `MQTT_PORT` | `1883` | Puerto MQTT | +| `REDIS_URL` | `redis://localhost:6379` | URL de conexión Redis | +| `RUST_LOG` | `info` | Nivel de logging (`trace`, `debug`, `info`, `warn`, `error`) | + +### publisher + +| Variable | Default | Descripción | +|---|---|---| +| `MQTT_HOST` | `localhost` | Hostname del broker MQTT | +| `MQTT_PORT` | `1883` | Puerto MQTT | +| `PUBLISH_INTERVAL_MS` | `1000` | Intervalo entre publicaciones (ms) | + +--- + +## Qué Demuestra Este Laboratorio + +- **Rust** como lenguaje para microservicios de alta frecuencia: sin GC, latencia predecible +- **Ownership y borrowing** en código de producción real: no hay `clone()` innecesario +- **Enums exhaustivos** (`ValidationError`) que el compilador obliga a manejar +- **Async/await con Tokio**: event loop MQTT + conexiones Redis concurrentes sin data races +- **rumqttc**: cliente MQTT puro Rust con reconexión automática y backpressure via QoS1 +- **Redis Streams** (`XADD MAXLEN`) como buffer circular para telemetría +- **Multi-binary Cargo package**: `ingestor` y `publisher` en el mismo crate, compartiendo `VehiclePosition` +- **Dockerfile multi-stage con múltiples targets**: un solo builder, dos imágenes mínimas de runtime +- **Tests sin mocks**: la lógica de validación se prueba directamente; los tests de integración usan Redis real con `#[ignore]` diff --git a/lab-11-rust-ingestor/demo.sh b/lab-11-rust-ingestor/demo.sh new file mode 100644 index 0000000..2bb51da --- /dev/null +++ b/lab-11-rust-ingestor/demo.sh @@ -0,0 +1,129 @@ +#!/usr/bin/env bash +# demo.sh — Ciclo completo: up → espera → muestra Redis → down +# +# Uso: +# ./demo.sh → ciclo completo (baja el stack al terminar) +# ./demo.sh --keep → no baja el stack al terminar +# +# Requiere: Docker, bash, python3 + +set -euo pipefail + +KEEP=false +[[ "${1:-}" == "--keep" ]] && KEEP=true + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +cd "${SCRIPT_DIR}" + +GREEN='\033[0;32m'; YELLOW='\033[1;33m'; CYAN='\033[0;36m'; NC='\033[0m'; BOLD='\033[1m' + +log() { echo -e "${GREEN}[demo]${NC} $*"; } +warn() { echo -e "${YELLOW}[demo]${NC} $*"; } +section() { echo -e "\n${CYAN}${BOLD}══ $* ══${NC}"; } + +cleanup() { + if $KEEP; then + warn "Stack conservado (--keep). Bajar con: docker compose down -v" + return + fi + log "Bajando el stack..." + docker compose down -v 2>&1 || true +} +trap cleanup EXIT + +# ── Paso 1: Build + Up ─────────────────────────────────────────────────────── + +section "Paso 1 — Build + Up" + +log "Construyendo imágenes (primera vez puede tardar varios minutos por la compilación Rust)..." +docker compose build --quiet + +log "Iniciando servicios..." +docker compose up -d + +# ── Paso 2: Esperar servicios ──────────────────────────────────────────────── + +section "Paso 2 — Esperando disponibilidad de servicios" + +wait_redis() { + local elapsed=0 + printf " Esperando Redis" + until docker compose exec -T redis redis-cli ping 2>/dev/null | grep -q PONG; do + [[ $elapsed -ge 30 ]] && { echo ""; warn "Redis no respondió en 30s"; return 1; } + printf "."; sleep 2; elapsed=$((elapsed + 2)) + done + echo ""; log "Redis listo (${elapsed}s)" +} + +wait_ingestor() { + local elapsed=0 + printf " Esperando ingestor" + until docker compose ps --format json 2>/dev/null \ + | python3 -c " +import sys, json +for line in sys.stdin: + line = line.strip() + if not line: continue + try: + d = json.loads(line) + if 'ingestor' in d.get('Service','') and d.get('State','') == 'running': + print('ok'); break + except: pass +" 2>/dev/null | grep -q ok; do + [[ $elapsed -ge 60 ]] && { echo ""; warn "ingestor no arrancó en 60s"; return 1; } + printf "."; sleep 3; elapsed=$((elapsed + 3)) + done + echo ""; log "ingestor corriendo (${elapsed}s)" +} + +wait_redis +wait_ingestor + +# ── Paso 3: Observar 15 segundos ──────────────────────────────────────────── + +section "Paso 3 — Observando 15 segundos de telemetría" + +log "Publicando posiciones GPS (publisher → NanoMQ → ingestor → Redis)..." +log "Logs del ingestor:" +echo "" +docker compose logs --tail=5 ingestor 2>/dev/null || true +sleep 15 + +# ── Paso 4: Consultar Redis ────────────────────────────────────────────────── + +section "Paso 4 — Estado de Redis" + +log "Claves transit:vehicle:*:latest:" +docker compose exec -T redis redis-cli --scan --pattern "transit:vehicle:*:latest" \ + | while read -r key; do + val=$(docker compose exec -T redis redis-cli GET "$key" 2>/dev/null || echo "") + if [[ -n "$val" ]]; then + echo "$val" | python3 -c " +import sys, json +d = json.load(sys.stdin) +print(f' {d[\"vehicle_id\"]:10s} route={d[\"route_id\"]} lat={d[\"lat\"]} lon={d[\"lon\"]} speed={d[\"speed_kmh\"]} km/h') +" 2>/dev/null || echo " $key → $val" + fi + done + +echo "" +STREAM_LEN=$(docker compose exec -T redis redis-cli XLEN transit:positions 2>/dev/null | tr -d '[:space:]' || echo "0") +log "Stream transit:positions: ${STREAM_LEN} entradas" + +# ── Resultado ──────────────────────────────────────────────────────────────── + +echo "" +echo -e "${BOLD}────────────────────────────────────────${NC}" +if [[ "${STREAM_LEN}" -gt 0 ]]; then + echo -e " ${GREEN}OK${NC} Pipeline operativo: MQTT → Redis" +else + echo -e " ${YELLOW}WARN${NC} Stream vacío — ¿están corriendo publisher e ingestor?" +fi +echo "" +echo " NanoMQ MQTT → localhost:1883" +echo " Redis → localhost:6379" +echo "" +echo " Logs en tiempo real:" +echo " docker compose logs -f ingestor" +echo " docker compose logs -f publisher" +echo -e "${BOLD}────────────────────────────────────────${NC}" diff --git a/lab-11-rust-ingestor/docker-compose.yml b/lab-11-rust-ingestor/docker-compose.yml new file mode 100644 index 0000000..b1125e7 --- /dev/null +++ b/lab-11-rust-ingestor/docker-compose.yml @@ -0,0 +1,86 @@ +services: + + # ───────────────────────────────────────────── + # MQTT broker — NanoMQ + # Puerto 1883: MQTT estándar + # Puerto 8083: MQTT sobre WebSocket (opcional) + # ───────────────────────────────────────────── + + nanomq: + image: emqx/nanomq:latest + ports: + - "1883:1883" # MQTT + - "8083:8083" # MQTT/WS + networks: + - transit-net + restart: unless-stopped + + # ───────────────────────────────────────────── + # Redis 7 — almacén de posiciones + # transit:vehicle:{id}:latest → última posición (TTL 5 min) + # transit:positions → stream histórico (MAXLEN ~ 1000) + # ───────────────────────────────────────────── + + redis: + image: redis:7-alpine + ports: + - "6379:6379" + networks: + - transit-net + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 3s + retries: 5 + restart: unless-stopped + + # ───────────────────────────────────────────── + # Ingestor — consume MQTT, valida, escribe Redis + # ───────────────────────────────────────────── + + ingestor: + build: + context: . + dockerfile: Dockerfile + target: ingestor + environment: + MQTT_HOST: nanomq + MQTT_PORT: "1883" + REDIS_URL: redis://redis:6379 + RUST_LOG: info + depends_on: + nanomq: + condition: service_started + redis: + condition: service_healthy + networks: + - transit-net + restart: unless-stopped + + # ───────────────────────────────────────────── + # Publisher — simulador GPS (3 buses, 1 pos/s) + # ───────────────────────────────────────────── + + publisher: + build: + context: . + dockerfile: Dockerfile + target: publisher + environment: + MQTT_HOST: nanomq + MQTT_PORT: "1883" + PUBLISH_INTERVAL_MS: "1000" + depends_on: + nanomq: + condition: service_started + networks: + - transit-net + restart: unless-stopped + +# ───────────────────────────────────────────── +# REDES +# ───────────────────────────────────────────── + +networks: + transit-net: + driver: bridge diff --git a/lab-11-rust-ingestor/docs/architecture.md b/lab-11-rust-ingestor/docs/architecture.md new file mode 100644 index 0000000..4404939 --- /dev/null +++ b/lab-11-rust-ingestor/docs/architecture.md @@ -0,0 +1,224 @@ +# Arquitectura — Lab 11: Rust Ingestor de Telemetría GPS + +## Visión General + +Lab 11 implementa un microservicio de ingestión de telemetría GPS escrito en Rust. Su responsabilidad es consumir mensajes de posición de vehículos publicados vía MQTT, validarlos semánticamente y almacenarlos en Redis como estado actual y stream histórico. + +El servicio demuestra los conceptos clave de Rust en un contexto de producción real: +- **Ownership y borrowing**: `RedisSink` posee el `redis::Client`; las posiciones se pasan por referencia. +- **Tipos seguros**: `VehiclePosition` encapsula las reglas de dominio; `ValidationError` es un enum exhaustivo. +- **Concurrencia sin data races**: Tokio gestiona el event loop de rumqttc y las conexiones Redis de forma segura. +- **Manejo de errores explícito**: El operador `?` propaga errores; las posiciones inválidas se descartan con log, no con panic. + +--- + +## Diagrama de Componentes + +``` +┌─────────────────────────────────────────────────────────────────┐ +│ STACK LAB 11 │ +│ │ +│ ┌────────────────────────────────────────────────────────┐ │ +│ │ publisher (Rust binary) │ │ +│ │ │ │ +│ │ Bus { vehicle_id, route_id, lat, lon, speed, bearing }│ │ +│ │ LCG → variación GPS simulada │ │ +│ │ serde_json → JSON payload │ │ +│ │ rumqttc → publica cada 1s │ │ +│ └───────────────────────────┬────────────────────────────┘ │ +│ │ MQTT QoS1 │ +│ topic: transit/vehicles/{id}/position │ +│ │ │ +│ ┌───────────────────────────▼────────────────────────────┐ │ +│ │ NanoMQ (broker MQTT) │ │ +│ │ Puerto 1883 │ │ +│ └───────────────────────────┬────────────────────────────┘ │ +│ │ MQTT QoS1 │ +│ topic: transit/vehicles/+/position │ +│ │ │ +│ ┌───────────────────────────▼────────────────────────────┐ │ +│ │ ingestor (Rust binary) │ │ +│ │ │ │ +│ │ rumqttc EventLoop.poll() │ │ +│ │ └── Packet::Publish │ │ +│ │ └── serde_json::from_slice:: │ │ +│ │ └── VehiclePosition.validate() │ │ +│ │ ├── Ok → RedisSink.publish() │ │ +│ │ └── Err → warn! + descarte │ │ +│ └───────────────────────────┬────────────────────────────┘ │ +│ │ TCP (multiplexed) │ +│ │ │ +│ ┌───────────────────────────▼────────────────────────────┐ │ +│ │ Redis 7 │ │ +│ │ Puerto 6379 │ │ +│ │ │ │ +│ │ SET transit:vehicle:{id}:latest EX 300 │ │ +│ │ → Última posición por vehículo, TTL 5 min │ │ +│ │ │ │ +│ │ XADD transit:positions MAXLEN ~ 1000 * ... │ │ +│ │ → Stream global de posiciones, cap ~1000 entradas │ │ +│ └────────────────────────────────────────────────────────┘ │ +└─────────────────────────────────────────────────────────────────┘ +``` + +--- + +## Estructura del Crate + +``` +rust-ingestor (package) +│ +├── [lib] src/lib.rs — exporta los módulos públicos +│ +├── [bin] src/main.rs — ingestor: EventLoop → validate → RedisSink +├── [bin] src/bin/publisher.rs— simulador GPS: 3 buses, LCG, 1 pos/s +│ +├── src/models.rs — VehiclePosition + ValidationError +├── src/mqtt.rs — create_client() + constantes de topic +└── src/redis_sink.rs — RedisSink: SET latest + XADD stream +``` + +### Dependencias principales + +| Crate | Propósito | +|---|---| +| `tokio` | Runtime async: gestiona el event loop MQTT y las conexiones Redis | +| `rumqttc` | Cliente MQTT puro Rust; gestiona reconexión automáticamente | +| `redis` (tokio-comp) | Cliente Redis async sobre Tokio; conexión multiplexada | +| `serde` + `serde_json` | Serialización/deserialización del payload JSON | +| `anyhow` | Propagación de errores en el binario (vs `thiserror` en libs) | +| `tracing` | Logging estructurado con campos clave-valor | +| `chrono` | Timestamp Unix en milisegundos para el publisher | + +--- + +## Flujo de Datos + +``` +publisher + └── LCG.step() → coordenadas GPS ±0.0005° (~55 m) + └── VehiclePosition { vehicle_id, route_id, lat, lon, speed_kmh, bearing, ts } + └── serde_json::to_string() + └── MQTT PUBLISH transit/vehicles/{id}/position QoS1 + +NanoMQ + └── retiene mensajes hasta que el ingestor los consuma (QoS1) + +ingestor + └── rumqttc::EventLoop::poll() + └── Event::Incoming(Packet::Publish) + └── serde_json::from_slice::() + ├── Err → warn!(topic, error) — payload ignorado + └── Ok + └── VehiclePosition::validate() + ├── Err → warn!(vehicle_id, error) — descartado + └── Ok + └── RedisSink::publish() + ├── SET transit:vehicle:{id}:latest EX 300 + └── XADD transit:positions MAXLEN ~ 1000 * ... + +Redis + ├── transit:vehicle:bus-001:latest → {"vehicle_id":"bus-001","lat":9.9281,...} + ├── transit:vehicle:bus-002:latest → {...} + ├── transit:vehicle:bus-003:latest → {...} + └── transit:positions → stream: 1703...-0 {"vehicle_id":"bus-001",...} + 1703...-1 {"vehicle_id":"bus-002",...} + ... +``` + +--- + +## Reglas de Validación + +| Campo | Regla | Error | +|---|---|---| +| `vehicle_id` | No vacío (trim) | `EmptyVehicleId` | +| `route_id` | No vacío (trim) | `EmptyRouteId` | +| `lat` | `[-90.0, 90.0]` | `InvalidLatitude` | +| `lon` | `[-180.0, 180.0]` | `InvalidLongitude` | +| `speed_kmh` | `[0.0, 200.0]` | `InvalidSpeed` | +| `bearing` | `[0.0, 360.0)` si presente | `InvalidBearing` | + +Las posiciones inválidas se descartan con un `warn!` estructurado — nunca causan panic ni detienen el servicio. + +--- + +## Almacenamiento Redis + +### Última posición (`SET ... EX 300`) + +``` +KEY transit:vehicle:bus-001:latest +VALUE {"vehicle_id":"bus-001","route_id":"R01","lat":9.9281,"lon":-84.0907, + "speed_kmh":45.0,"bearing":270.0,"ts":1700000000000} +TTL 300 s (se renueva en cada actualización) +``` + +Uso: consulta rápida del estado actual de un vehículo concreto. Si el vehículo deja de reportar, la clave expira automáticamente. + +### Stream histórico (`XADD MAXLEN ~ 1000`) + +``` +KEY transit:positions +TYPE Redis Stream +CAP ~ 1000 entradas (memoria acotada) + +ENTRY 1703000000000-0 + vehicle_id bus-001 + route_id R01 + lat 9.928100 + lon -84.090700 + speed_kmh 45.0 + ts 1703000000000 +``` + +Uso: consulta de las últimas N posiciones de todos los vehículos, ordenadas por tiempo de inserción. Diseñado para alimentar el pipeline del Lab 04 (TimescaleDB via Prefect). + +--- + +## Diseño del Ingestor — Reconexión MQTT + +`rumqttc` gestiona la reconexión automáticamente: cuando `EventLoop::poll()` devuelve `Err`, el loop continúa esperando 5 segundos antes del siguiente intento. El broker NanoMQ retiene los mensajes QoS1 durante la desconexión (dentro del tiempo de persistencia del broker). + +```rust +Err(e) => { + error!(error = %e, "error en el event loop MQTT, reintentando en 5s"); + tokio::time::sleep(Duration::from_secs(5)).await; + // rumqttc reintenta la conexión en el siguiente poll() +} +``` + +--- + +## Conexión con Otros Labs + +| Lab | Relación | +|---|---| +| Lab 04 (Prefect + TimescaleDB) | El stream `transit:positions` puede ser consumido por un flow Prefect para cargar posiciones en TimescaleDB como serie temporal | +| Lab 08 (MQTT + NanoMQ) | Mismo broker NanoMQ; los topics `transit/vehicles/+/position` siguen la misma convención de Lab 08 | +| Lab 07 (Prometheus + Grafana) | El ingestor podría exponer métricas (`messages_ingested_total`, `validation_errors_total`) para scraping | + +--- + +## Redes Docker + +``` +transit-net + ├── nanomq — broker MQTT (puerto 1883) + ├── redis — almacén (puerto 6379) + ├── ingestor — consume MQTT, escribe Redis + └── publisher — genera telemetría simulada +``` + +--- + +## Decisiones de Diseño + +| Decisión | Alternativa | Motivo | +|---|---|---| +| `rumqttc` sobre `paho-mqtt` | paho-mqtt (C bindings) | rumqttc es puro Rust: seguro en memoria, sin dependencias de sistema, reconexión automática integrada | +| Conexión Redis multiplexada | Pool de conexiones | Para el throughput de este lab (≤ 10 msg/s) una sola conexión multiplexada es suficiente y más simple | +| `redis::cmd()` explícito | `AsyncCommands::xadd()` | Más robusto ante cambios de API entre versiones del crate; XADD MAXLEN no tiene wrapper estable en todas las versiones | +| LCG en publisher | Crate `rand` | Evita una dependencia extra para un uso tan acotado; el LCG es suficiente para variación visual | +| `VehiclePosition` compartida entre binarios | DTOs separados | Garantiza que el publisher emite exactamente el esquema que el ingestor espera; reutilización real del crate | +| `#[ignore]` en tests de integración | Siempre correr | Los tests unitarios corren sin Docker; los de integración se activan con `--ignored` cuando el stack está levantado | diff --git a/lab-11-rust-ingestor/src/bin/publisher.rs b/lab-11-rust-ingestor/src/bin/publisher.rs new file mode 100644 index 0000000..6a3079d --- /dev/null +++ b/lab-11-rust-ingestor/src/bin/publisher.rs @@ -0,0 +1,150 @@ +/// Simulador GPS para demostrar el ingestor. +/// +/// Publica posiciones de 3 buses ficticios cada PUBLISH_INTERVAL_MS milisegundos +/// en el topic `transit/vehicles/{vehicle_id}/position`. +/// +/// Variables de entorno: +/// MQTT_HOST — hostname del broker (default: localhost) +/// MQTT_PORT — puerto MQTT (default: 1883) +/// PUBLISH_INTERVAL_MS — intervalo ms (default: 1000) +use rust_ingestor::models::VehiclePosition; +use rumqttc::{AsyncClient, MqttOptions, QoS}; +use std::time::Duration; +use tokio::time::sleep; + +// ── Generador de números pseudo-aleatorios (LCG) ───────────────────────────── +// Evita añadir la dependencia `rand` para un uso tan acotado. + +struct Lcg(u64); + +impl Lcg { + fn new() -> Self { + use std::time::{SystemTime, UNIX_EPOCH}; + let seed = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .subsec_nanos() as u64; + Self(seed ^ 0xdeadbeef_cafebabe) + } + + /// Devuelve un f64 en [0, 1). + fn next_f64(&mut self) -> f64 { + self.0 = self.0 + .wrapping_mul(6_364_136_223_846_793_005) + .wrapping_add(1_442_695_040_888_963_407); + (self.0 >> 33) as f64 / (u32::MAX as f64 + 1.0) + } + + /// Devuelve un f64 en [min, max). + fn range(&mut self, min: f64, max: f64) -> f64 { + min + self.next_f64() * (max - min) + } +} + +// ── Simulador de bus ────────────────────────────────────────────────────────── + +struct Bus { + vehicle_id: &'static str, + route_id: &'static str, + lat: f64, + lon: f64, + bearing: f64, + speed_kmh: f64, +} + +impl Bus { + fn new(vehicle_id: &'static str, route_id: &'static str, lat: f64, lon: f64, bearing: f64) -> Self { + Self { vehicle_id, route_id, lat, lon, bearing, speed_kmh: 40.0 } + } + + /// Avanza ~50 m en la dirección actual con pequeñas variaciones de rumbo y velocidad. + fn step(&mut self, rng: &mut Lcg) { + let delta = 0.0005_f64; // ~55 m por paso + let rad = self.bearing.to_radians(); + self.lat += delta * rad.cos(); + self.lon += delta * rad.sin(); + + self.speed_kmh = (self.speed_kmh + rng.range(-3.0, 3.0)).clamp(5.0, 80.0); + self.bearing = (self.bearing + rng.range(-5.0, 5.0)).rem_euclid(360.0); + } + + fn to_position(&self) -> VehiclePosition { + VehiclePosition { + vehicle_id: self.vehicle_id.into(), + route_id: self.route_id.into(), + lat: (self.lat * 1e6).round() / 1e6, + lon: (self.lon * 1e6).round() / 1e6, + speed_kmh: (self.speed_kmh * 10.0).round() / 10.0, + bearing: Some((self.bearing * 10.0).round() / 10.0), + ts: chrono::Utc::now().timestamp_millis(), + } + } + + fn topic(&self) -> String { + format!("transit/vehicles/{}/position", self.vehicle_id) + } +} + +// ── Entry point ────────────────────────────────────────────────────────────── + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + let mqtt_host = std::env::var("MQTT_HOST").unwrap_or_else(|_| "localhost".into()); + let mqtt_port: u16 = std::env::var("MQTT_PORT") + .unwrap_or_else(|_| "1883".into()) + .parse()?; + let interval_ms: u64 = std::env::var("PUBLISH_INTERVAL_MS") + .unwrap_or_else(|_| "1000".into()) + .parse()?; + + let mut options = MqttOptions::new("rust-publisher", &mqtt_host, mqtt_port); + options.set_keep_alive(Duration::from_secs(30)); + options.set_clean_session(true); + + let (client, mut eventloop) = AsyncClient::new(options, 10); + + // El event loop corre en su propio task; el publisher sólo necesita escribir. + tokio::spawn(async move { + loop { + if let Err(e) = eventloop.poll().await { + eprintln!("[publisher] event loop error: {e}"); + sleep(Duration::from_secs(2)).await; + } + } + }); + + // Esperar a que el broker esté disponible antes de publicar. + sleep(Duration::from_secs(2)).await; + println!("[publisher] conectado a {mqtt_host}:{mqtt_port} — publicando cada {interval_ms}ms"); + + let mut rng = Lcg::new(); + + // Tres buses en distintas rutas del Gran Área Metropolitana de Costa Rica. + let mut buses = vec![ + Bus::new("bus-001", "R01", 9.9281, -84.0907, 45.0), // San José → Cartago + Bus::new("bus-002", "R02", 9.9560, -84.1320, 135.0), // San José → Heredia + Bus::new("bus-003", "R03", 9.9000, -84.1500, 270.0), // San José → Alajuela + ]; + + let mut count: u64 = 0; + + loop { + for bus in &mut buses { + bus.step(&mut rng); + let pos = bus.to_position(); + let payload = serde_json::to_string(&pos)?; + let topic = bus.topic(); + + client + .publish(&topic, QoS::AtLeastOnce, false, payload.as_bytes()) + .await?; + } + + count += 1; + if count % 10 == 0 { + println!("[publisher] {count} ciclos publicados (3 posiciones/ciclo)"); + } + + sleep(Duration::from_millis(interval_ms)).await; + } +} diff --git a/lab-11-rust-ingestor/src/lib.rs b/lab-11-rust-ingestor/src/lib.rs new file mode 100644 index 0000000..d864222 --- /dev/null +++ b/lab-11-rust-ingestor/src/lib.rs @@ -0,0 +1,4 @@ +// Módulos públicos del crate — accesibles desde los binarios y desde los tests de integración. +pub mod models; +pub mod mqtt; +pub mod redis_sink; diff --git a/lab-11-rust-ingestor/src/main.rs b/lab-11-rust-ingestor/src/main.rs new file mode 100644 index 0000000..2cde479 --- /dev/null +++ b/lab-11-rust-ingestor/src/main.rs @@ -0,0 +1,99 @@ +use rust_ingestor::{ + models::VehiclePosition, + mqtt, + redis_sink::RedisSink, +}; +use rumqttc::{Event, Packet}; +use tracing::{error, info, warn}; + +/// Ingestor de telemetría GPS. +/// +/// Flujo: +/// NanoMQ (MQTT) → parse JSON → VehiclePosition.validate() → RedisSink.publish() +/// +/// Variables de entorno: +/// MQTT_HOST — hostname del broker MQTT (default: localhost) +/// MQTT_PORT — puerto MQTT (default: 1883) +/// REDIS_URL — URL de conexión a Redis (default: redis://localhost:6379) +/// RUST_LOG — nivel de logging (default: info) +#[tokio::main] +async fn main() -> anyhow::Result<()> { + // Inicializar logging estructurado; lee RUST_LOG, cae a "info" si no está definido. + tracing_subscriber::fmt() + .with_env_filter( + tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")), + ) + .init(); + + let mqtt_host = std::env::var("MQTT_HOST").unwrap_or_else(|_| "localhost".into()); + let mqtt_port: u16 = std::env::var("MQTT_PORT") + .unwrap_or_else(|_| "1883".into()) + .parse()?; + let redis_url = std::env::var("REDIS_URL") + .unwrap_or_else(|_| "redis://localhost:6379".into()); + + info!( + mqtt_host = %mqtt_host, + mqtt_port, + redis_url = %redis_url, + topic = mqtt::POSITION_TOPIC, + "ingestor iniciando" + ); + + let sink = RedisSink::new(&redis_url)?; + let (client, mut eventloop) = mqtt::create_client(&mqtt_host, mqtt_port, "rust-ingestor"); + + // La suscripción se encola; el event loop la procesa en el primer poll(). + client.subscribe(mqtt::POSITION_TOPIC, mqtt::TOPIC_QOS).await?; + + info!(topic = mqtt::POSITION_TOPIC, "suscrito al topic MQTT"); + + loop { + match eventloop.poll().await { + Ok(Event::Incoming(Packet::ConnAck(_))) => { + info!("conectado al broker MQTT"); + } + + Ok(Event::Incoming(Packet::Publish(msg))) => { + match serde_json::from_slice::(&msg.payload) { + Ok(pos) => match pos.validate() { + Ok(()) => { + if let Err(e) = sink.publish(&pos).await { + error!( + vehicle_id = %pos.vehicle_id, + error = %e, + "error al almacenar posición en Redis" + ); + } + } + Err(e) => { + warn!( + vehicle_id = %pos.vehicle_id, + error = %e, + "posición descartada: validación fallida" + ); + } + }, + Err(e) => { + warn!( + topic = %msg.topic, + error = %e, + "payload inválido: no se pudo deserializar como VehiclePosition" + ); + } + } + } + + Ok(_) => { + // Otros eventos MQTT (SubAck, PingResp, etc.) se ignoran. + } + + Err(e) => { + // rumqttc reintenta la conexión automáticamente al continuar el loop. + error!(error = %e, "error en el event loop MQTT, reintentando en 5s"); + tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; + } + } + } +} diff --git a/lab-11-rust-ingestor/src/models.rs b/lab-11-rust-ingestor/src/models.rs new file mode 100644 index 0000000..24b6230 --- /dev/null +++ b/lab-11-rust-ingestor/src/models.rs @@ -0,0 +1,122 @@ +use serde::{Deserialize, Serialize}; +use std::fmt; + +/// Posición GPS de un vehículo publicada vía MQTT. +/// +/// Payload JSON esperado en el topic `transit/vehicles/{vehicle_id}/position`: +/// ```json +/// { +/// "vehicle_id": "bus-001", +/// "route_id": "R01", +/// "lat": 9.9281, +/// "lon": -84.0907, +/// "speed_kmh": 45.0, +/// "bearing": 270.0, +/// "ts": 1700000000000 +/// } +/// ``` +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct VehiclePosition { + pub vehicle_id: String, + pub route_id: String, + pub lat: f64, + pub lon: f64, + pub speed_kmh: f64, + pub bearing: Option, + /// Timestamp Unix en milisegundos. + pub ts: i64, +} + +/// Errores de validación semántica de una posición GPS. +#[derive(Debug)] +pub enum ValidationError { + EmptyVehicleId, + EmptyRouteId, + InvalidLatitude(f64), + InvalidLongitude(f64), + InvalidSpeed(f64), + InvalidBearing(f64), +} + +impl fmt::Display for ValidationError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::EmptyVehicleId => write!(f, "vehicle_id no puede estar vacío"), + Self::EmptyRouteId => write!(f, "route_id no puede estar vacío"), + Self::InvalidLatitude(v) => write!(f, "latitud {v:.6} fuera de rango [-90, 90]"), + Self::InvalidLongitude(v) => write!(f, "longitud {v:.6} fuera de rango [-180, 180]"), + Self::InvalidSpeed(v) => write!(f, "speed_kmh {v:.1} fuera de rango [0, 200]"), + Self::InvalidBearing(v) => write!(f, "bearing {v:.1} fuera de rango [0, 360)"), + } + } +} + +impl std::error::Error for ValidationError {} + +impl VehiclePosition { + /// Valida que todos los campos estén dentro de rangos aceptables. + /// Velocidad máxima: 200 km/h (límite físico para un bus). + pub fn validate(&self) -> Result<(), ValidationError> { + if self.vehicle_id.trim().is_empty() { + return Err(ValidationError::EmptyVehicleId); + } + if self.route_id.trim().is_empty() { + return Err(ValidationError::EmptyRouteId); + } + if !(-90.0..=90.0).contains(&self.lat) { + return Err(ValidationError::InvalidLatitude(self.lat)); + } + if !(-180.0..=180.0).contains(&self.lon) { + return Err(ValidationError::InvalidLongitude(self.lon)); + } + if !(0.0..=200.0).contains(&self.speed_kmh) { + return Err(ValidationError::InvalidSpeed(self.speed_kmh)); + } + if let Some(b) = self.bearing { + if !(0.0..360.0).contains(&b) { + return Err(ValidationError::InvalidBearing(b)); + } + } + Ok(()) + } + + /// Clave Redis para la última posición conocida del vehículo (TTL 5 min). + pub fn redis_key_latest(&self) -> String { + format!("transit:vehicle:{}:latest", self.vehicle_id) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn valid_pos() -> VehiclePosition { + VehiclePosition { + vehicle_id: "bus-001".into(), + route_id: "R01".into(), + lat: 9.9281, + lon: -84.0907, + speed_kmh: 45.0, + bearing: Some(270.0), + ts: 1_700_000_000_000, + } + } + + #[test] + fn valid_position_passes() { + assert!(valid_pos().validate().is_ok()); + } + + #[test] + fn bearing_none_passes() { + let mut p = valid_pos(); + p.bearing = None; + assert!(p.validate().is_ok()); + } + + #[test] + fn redis_key_format() { + let key = valid_pos().redis_key_latest(); + assert_eq!(key, "transit:vehicle:bus-001:latest"); + } +} diff --git a/lab-11-rust-ingestor/src/mqtt.rs b/lab-11-rust-ingestor/src/mqtt.rs new file mode 100644 index 0000000..a4ad064 --- /dev/null +++ b/lab-11-rust-ingestor/src/mqtt.rs @@ -0,0 +1,22 @@ +use rumqttc::{AsyncClient, EventLoop, MqttOptions, QoS}; +use std::time::Duration; + +/// Topic MQTT al que el ingestor se suscribe. +/// El `+` es un wildcard de nivel único: coincide con cualquier vehicle_id. +pub const POSITION_TOPIC: &str = "transit/vehicles/+/position"; + +/// QoS utilizado tanto en suscripción como en publicación. +/// AtLeastOnce garantiza entrega sin duplicados excesivos para telemetría. +pub const TOPIC_QOS: QoS = QoS::AtLeastOnce; + +/// Crea un cliente MQTT asíncrono y su event loop asociado. +/// +/// El `capacity` del canal interno (10) es suficiente para el throughput +/// de este laboratorio (≤ 10 publicaciones/s). +pub fn create_client(host: &str, port: u16, client_id: &str) -> (AsyncClient, EventLoop) { + let mut options = MqttOptions::new(client_id, host, port); + options.set_keep_alive(Duration::from_secs(30)); + options.set_clean_session(true); + + AsyncClient::new(options, 10) +} diff --git a/lab-11-rust-ingestor/src/redis_sink.rs b/lab-11-rust-ingestor/src/redis_sink.rs new file mode 100644 index 0000000..f99ba93 --- /dev/null +++ b/lab-11-rust-ingestor/src/redis_sink.rs @@ -0,0 +1,63 @@ +use crate::models::VehiclePosition; +use tracing::info; + +/// Escribe posiciones GPS validadas en Redis. +/// +/// Por cada posición se realizan dos operaciones: +/// 1. `SET transit:vehicle:{id}:latest EX 300` — última posición conocida (TTL 5 min) +/// 2. `XADD transit:positions MAXLEN ~ 1000 * ...` — stream acotado a ~1000 entradas +pub struct RedisSink { + client: redis::Client, +} + +impl RedisSink { + pub fn new(url: &str) -> anyhow::Result { + let client = redis::Client::open(url)?; + Ok(Self { client }) + } + + pub async fn publish(&self, pos: &VehiclePosition) -> anyhow::Result<()> { + let mut conn = self.client.get_multiplexed_async_connection().await?; + + let json = serde_json::to_string(pos)?; + + // Última posición conocida: expira en 5 minutos si el vehículo deja de reportar. + let _: () = redis::cmd("SET") + .arg(pos.redis_key_latest()) + .arg(&json) + .arg("EX") + .arg(300i64) + .query_async(&mut conn) + .await?; + + // Stream histórico acotado: permite consultar las últimas N posiciones globales. + let lat_s = pos.lat.to_string(); + let lon_s = pos.lon.to_string(); + let speed_s = pos.speed_kmh.to_string(); + let ts_s = pos.ts.to_string(); + + let _: String = redis::cmd("XADD") + .arg("transit:positions") + .arg("MAXLEN").arg("~").arg(1000i64) + .arg("*") + .arg("vehicle_id").arg(pos.vehicle_id.as_str()) + .arg("route_id") .arg(pos.route_id.as_str()) + .arg("lat") .arg(lat_s.as_str()) + .arg("lon") .arg(lon_s.as_str()) + .arg("speed_kmh") .arg(speed_s.as_str()) + .arg("ts") .arg(ts_s.as_str()) + .query_async(&mut conn) + .await?; + + info!( + vehicle_id = %pos.vehicle_id, + route_id = %pos.route_id, + lat = pos.lat, + lon = pos.lon, + speed_kmh = pos.speed_kmh, + "posición almacenada" + ); + + Ok(()) + } +} diff --git a/lab-11-rust-ingestor/tests/integration_test.rs b/lab-11-rust-ingestor/tests/integration_test.rs new file mode 100644 index 0000000..45df6f3 --- /dev/null +++ b/lab-11-rust-ingestor/tests/integration_test.rs @@ -0,0 +1,215 @@ +/// Tests de integración para el crate rust-ingestor. +/// +/// Los tests de validación no requieren servicios externos — cubren la lógica +/// de negocio pura de `VehiclePosition::validate()`. +/// +/// Los tests que necesitan Redis o NanoMQ están marcados con `#[ignore]` y +/// requieren el stack de Docker levantado: `docker compose up -d` +use rust_ingestor::models::{ValidationError, VehiclePosition}; + +// ── Helpers ────────────────────────────────────────────────────────────────── + +fn base() -> VehiclePosition { + VehiclePosition { + vehicle_id: "bus-001".into(), + route_id: "R01".into(), + lat: 9.9281, + lon: -84.0907, + speed_kmh: 45.0, + bearing: Some(270.0), + ts: 1_700_000_000_000, + } +} + +// ── Casos válidos ───────────────────────────────────────────────────────────── + +#[test] +fn valid_position_passes_validation() { + assert!(base().validate().is_ok()); +} + +#[test] +fn bearing_none_is_valid() { + let mut p = base(); + p.bearing = None; + assert!(p.validate().is_ok()); +} + +#[test] +fn speed_zero_is_valid() { + let mut p = base(); + p.speed_kmh = 0.0; + assert!(p.validate().is_ok()); +} + +#[test] +fn speed_at_max_limit_is_valid() { + let mut p = base(); + p.speed_kmh = 200.0; + assert!(p.validate().is_ok()); +} + +#[test] +fn boundary_lat_plus_90_is_valid() { + let mut p = base(); + p.lat = 90.0; + assert!(p.validate().is_ok()); +} + +#[test] +fn boundary_lat_minus_90_is_valid() { + let mut p = base(); + p.lat = -90.0; + assert!(p.validate().is_ok()); +} + +#[test] +fn bearing_at_zero_is_valid() { + let mut p = base(); + p.bearing = Some(0.0); + assert!(p.validate().is_ok()); +} + +#[test] +fn bearing_at_359_is_valid() { + let mut p = base(); + p.bearing = Some(359.9); + assert!(p.validate().is_ok()); +} + +// ── Casos inválidos ─────────────────────────────────────────────────────────── + +#[test] +fn empty_vehicle_id_fails() { + let mut p = base(); + p.vehicle_id = " ".into(); + assert!(matches!(p.validate(), Err(ValidationError::EmptyVehicleId))); +} + +#[test] +fn empty_route_id_fails() { + let mut p = base(); + p.route_id = String::new(); + assert!(matches!(p.validate(), Err(ValidationError::EmptyRouteId))); +} + +#[test] +fn latitude_above_90_fails() { + let mut p = base(); + p.lat = 90.1; + assert!(matches!(p.validate(), Err(ValidationError::InvalidLatitude(_)))); +} + +#[test] +fn latitude_below_minus_90_fails() { + let mut p = base(); + p.lat = -90.1; + assert!(matches!(p.validate(), Err(ValidationError::InvalidLatitude(_)))); +} + +#[test] +fn longitude_above_180_fails() { + let mut p = base(); + p.lon = 180.1; + assert!(matches!(p.validate(), Err(ValidationError::InvalidLongitude(_)))); +} + +#[test] +fn longitude_below_minus_180_fails() { + let mut p = base(); + p.lon = -180.1; + assert!(matches!(p.validate(), Err(ValidationError::InvalidLongitude(_)))); +} + +#[test] +fn negative_speed_fails() { + let mut p = base(); + p.speed_kmh = -1.0; + assert!(matches!(p.validate(), Err(ValidationError::InvalidSpeed(_)))); +} + +#[test] +fn speed_above_200_fails() { + let mut p = base(); + p.speed_kmh = 200.1; + assert!(matches!(p.validate(), Err(ValidationError::InvalidSpeed(_)))); +} + +#[test] +fn bearing_360_fails() { + let mut p = base(); + p.bearing = Some(360.0); + assert!(matches!(p.validate(), Err(ValidationError::InvalidBearing(_)))); +} + +#[test] +fn negative_bearing_fails() { + let mut p = base(); + p.bearing = Some(-1.0); + assert!(matches!(p.validate(), Err(ValidationError::InvalidBearing(_)))); +} + +// ── Serialización / deserialización ────────────────────────────────────────── + +#[test] +fn json_round_trip() { + let original = base(); + let json = serde_json::to_string(&original).unwrap(); + let decoded: VehiclePosition = serde_json::from_str(&json).unwrap(); + + assert_eq!(decoded.vehicle_id, original.vehicle_id); + assert_eq!(decoded.route_id, original.route_id); + assert!((decoded.lat - original.lat).abs() < 1e-9); + assert!((decoded.lon - original.lon).abs() < 1e-9); + assert!((decoded.speed_kmh - original.speed_kmh).abs() < 1e-9); + assert_eq!(decoded.ts, original.ts); +} + +#[test] +fn deserialization_without_bearing_succeeds() { + let json = r#"{ + "vehicle_id": "bus-001", + "route_id": "R01", + "lat": 9.9281, + "lon": -84.0907, + "speed_kmh": 30.0, + "ts": 1700000000000 + }"#; + let pos: VehiclePosition = serde_json::from_str(json).unwrap(); + assert!(pos.bearing.is_none()); + assert!(pos.validate().is_ok()); +} + +#[test] +fn redis_key_latest_format() { + let key = base().redis_key_latest(); + assert_eq!(key, "transit:vehicle:bus-001:latest"); +} + +// ── Tests de integración con servicios (requieren docker compose up -d) ─────── + +/// Verifica que el ingestor puede almacenar una posición en Redis. +/// Ejecutar con: cargo test -- --ignored +#[tokio::test] +#[ignore = "requiere Redis en localhost:6379"] +async fn redis_sink_stores_position() { + use rust_ingestor::redis_sink::RedisSink; + + let sink = RedisSink::new("redis://localhost:6379").unwrap(); + let pos = base(); + + sink.publish(&pos).await.expect("publish no debe fallar con Redis disponible"); + + // Verificar que la clave existe en Redis + let client = redis::Client::open("redis://localhost:6379").unwrap(); + let mut conn = client.get_multiplexed_async_connection().await.unwrap(); + let value: Option = redis::cmd("GET") + .arg(pos.redis_key_latest()) + .query_async(&mut conn) + .await + .unwrap(); + + assert!(value.is_some(), "la clave latest debe existir tras publish()"); + let stored: VehiclePosition = serde_json::from_str(&value.unwrap()).unwrap(); + assert_eq!(stored.vehicle_id, pos.vehicle_id); +}