API logs (Python / FastAPI)
central/mercure est un microservice Python + FastAPI dédié à l'accès aux logs des
nodes. Il est le seul composant autorisé à dialoguer avec Elasticsearch : l'API centrale
(saturn) et le frontend passent toujours par lui.
:::info Objectif
Découpler l'API centrale d'Elasticsearch. saturn ne connaît plus l'URL ni les
identifiants Elasticsearch : il appelle simplement l'API logs en REST.
:::
Rôle
- Ingestion des logs remontés par les nodes.
- Mise en file Kafka des événements pour absorber les pics de charge.
- Requêtage/filtrage des logs pour la supervision.
- Indexation dans des index journaliers
leukos-node-logs-YYYY.MM.DD.
Pipeline d'ingestion
Mercure accuse réception après mise en file Kafka. L'indexation Elasticsearch est asynchrone via le worker consommateur.
Arborescence
central/mercure/
├── Dockerfile
├── requirements.txt
├── app/
│ ├── api/http.py # application FastAPI + endpoints
│ ├── core/settings.py # settings (env : ES, Kafka, timeouts)
│ ├── domain/log_models.py # modèles Pydantic (LogEvent, réponses)
│ └── infra/
│ ├── elasticsearch_client.py # client httpx async vers Elasticsearch
│ └── kafka_pipeline.py # producer + consumer worker Kafka
└── tests/
└── test_es_client.py
Endpoints REST
| Méthode | Chemin | Description |
|---|---|---|
GET | /health | Santé du service + état des connexions Elasticsearch et Kafka. |
POST | /logs | Enfile un log dans Kafka (202 Accepted). |
GET | /logs | Requête des logs (limit, node_id, level). |
Modèle d'un log (POST /logs)
{
"node_id": "node-01",
"level": "error",
"message": "mqtt: connection lost",
"timestamp": "2026-08-09T10:00:00Z",
"meta": {
"machine_id": "f7f2d7b7f8c54c17b1a2fdb4c0ec8e90",
"mac": "dc:a6:32:11:22:33",
"service": "mosquitto"
}
}
timestamp est optionnel (RFC3339, généré par défaut) ; level est normalisé
(trim + minuscules).
Pour l'observabilité machine, meta.machine_id et meta.mac sont fortement
recommandés sur les logs remontés par les nodes.
Configuration (variables d'environnement)
| Variable | Défaut | Description |
|---|---|---|
ELASTICSEARCH_URL | http://es.data.server:80 | URL du cluster Elasticsearch. |
ES_USERNAME | supervisor | Utilisateur (auth basic). |
ES_PASSWORD | — | Mot de passe (auth basic). |
INDEX_PREFIX | leukos-node-logs | Préfixe des index journaliers. |
PORT | 8140 | Port d'écoute HTTP. |
ES_TIMEOUT | 8 | Timeout des appels Elasticsearch (s). |
LOG_MAX_LIMIT | 500 | Plafond du paramètre limit. |
KAFKA_ENABLED | true | Active l'ingestion via Kafka. |
KAFKA_BOOTSTRAP_SERVERS | central-kafka:9092 | Brokers Kafka. |
KAFKA_TOPIC | leukos-node-logs | Topic d'ingestion logs. |
KAFKA_GROUP_ID | mercure-indexer | Consumer group du worker d'indexation. |
KAFKA_CLIENT_ID | mercure | Préfixe client Kafka producer/consumer. |
Intégration avec l'API centrale
L'API centrale consomme ce service via l'adaptateur Go
central/saturn/internal/adapters/logsapi. L'URL est configurée par
LOGS_API_URL (défaut http://central-mercure:8140).
Lancement
Via la stack centrale (central/docker-compose.central.yml) :
npm run central:up
Vérification :
curl http://localhost:8140/health
Exemple de réponse :
{
"status": "ok",
"service": "mercure",
"storage": {
"elasticsearch": "up",
"kafka": "up"
},
"connections": {
"elasticsearch": {
"status": "up",
"address": "es.data.server",
"port": 80,
"domain": "es.data.server"
},
"kafka": {
"status": "up",
"address": "central-kafka",
"port": 9092,
"domain": null
}
},
"time": "2026-08-10T15:04:11Z"
}
status passe à degraded si Mercure n'arrive plus à joindre Elasticsearch,
ou Kafka quand KAFKA_ENABLED=true.
Flotte réseau
Mercure consomme aussi le topic Kafka leukos-health-state alimenté par toutes
les API du central et expose GET /fleet. Voir Flotte & santé réseau.