Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions lab-11-rust-ingestor/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"] }
60 changes: 60 additions & 0 deletions lab-11-rust-ingestor/Dockerfile
Original file line number Diff line number Diff line change
@@ -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"]
252 changes: 252 additions & 0 deletions lab-11-rust-ingestor/README.md
Original file line number Diff line number Diff line change
@@ -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]`
Loading
Loading