Dette passe 1, lot L5 : boucle MQTT du serveur, livraison de la configuration et journaux #5

Open
Thomas wants to merge 4 commits from pr/dette-l5 into pr/dette-l4
Member

Objet

Lot L5 de la passe 1 de dette : la boucle MQTT du serveur ne dépend plus de la base (files bornées par worker, accusé MQTT après écriture, session persistante), la configuration descendante est livrée jusqu'à son accusé, l'arrêt est propre et les journaux sont structurés, sans payload.

Dettes traitées et preuves

Dette Correction Preuve
Q2, abonnement # remplacé par +, réabonnement à chaque ConnAck ; plus de topic à plusieurs niveaux ni de $SYS aquaserveur/src/mqtt/settings.rs, server.rs
N11, Q9, recommandation 7 set_max_packet_size du client (8 Kio de payload plus le plus long topic) ; broker : message_size_limit 8192, max_connections 1000, max_inflight_messages 20, max_queued_messages 100000 ; reconnexion Backoff 1 s à 60 s ; MQTT_BROKER_URL et MQTT_CLIENT_ID validés au démarrage mqtt/settings.rs, server.rs, configs/mosquitto/mosquitto.conf
N10, boucle et base la boucle interroge le broker, route et dépose dans une file bornée par worker (256 places, réserve QoS 1 de 64) ; INGEST_WORKERS (défaut 4) ; un boîtier a un seul worker (ordre conservé) ; erreurs passagères de la base rejouées (1 s doublé jusqu'à 30 s) sans accuser, erreurs durables journalisées puis accusées aquaserveur/src/server.rs, ingest.rs, worker.rs, database/retry.rs, database/store.rs
Recommandation 10, C17 session persistante (clean_session = false), client id stable, accusé manuel seulement après traitement ; rejeu sans effet grâce à l'empreinte des mesures (L4) mqtt/settings.rs, worker.rs
N7, configuration descendante tâche run_delivery : version publiée, renvoi tant que l'accusé manque (60 s doublé jusqu'à 30 min), message retenu par boîtier sur maj_serv-{id}, état relu au démarrage depuis la table boitier_config_deliveries aquaserveur/src/config_delivery.rs, scripts/db/migrations/002_config_delivery.sql
N15, arrêt SIGTERM et SIGINT : réception suspendue, files vidées par les workers, accusés envoyés, DISCONNECT, fermeture du pool ; vidage borné à 8 s server.rs, main.rs
N12, N13, Q15, journaux tracing_subscriber avec RUST_LOG ; ni payload ni donnée personnelle (topic tronqué à 128 octets, erreur serde réduite à sa nature) ; rotation json-file 10 Mo x 5 dans les deux composes ; Mosquitto sur stdout aquaserveur/src/logging.rs, provisioning/error.rs, composes

Tests et résultats

  • TDD : chaque module vu rouge avant implémentation (logging 7 sur 7, MqttSettings 7 sur 7, ingest 8 sur 8, config_delivery 9 sur 9, boucle serveur contre broker réel, SIGTERM).
  • Contre-épreuves (code neutralisé, test rouge, restauration vérifiée) : clean_session(true), limite client à 10 Kio, boucle bloquée par la base, payload journalisé, workers sans vidage, accusé avant écriture, maj_serv non retenu, renvoi désactivé, suivi non relu au démarrage.
  • Ajouts : 44 tests unitaires ; 11 tests ignorés contre MariaDB et Mosquitto réels (tests/mqtt_server_integration.rs, tests/config_delivery_integration.rs), dont un serveur coupé brutalement au tiers d'un flux de 90 trames puis redémarré : 90 lignes, 90 epochs distincts, rejeu de 15 trames sans nouvelle ligne.
  • cargo fmt --check, cargo clippy ... -D warnings : OK.
  • cargo test --workspace : 332 passés, 0 échec, 68 ignorés (L4 : 288 et 57).
  • Tests ignorés (MariaDB 10.11 jetable avec dump, fixtures et migration 001, deux Mosquitto 2.0.22 jetables) : 68 passés, 0 échec, deux passages.
  • Vérifications avec le binaire : PUBLISH de 20 Kio sans coupure du serveur, SIGKILL pendant un flux de 60 trames puis redémarrage sans perte ni doublon, docker stop traité en 369 ms, livraison de configuration sur quatre démarrages successifs, migration 002 sur MariaDB 10.6.22 jetable (version de la prod).

Migrations et actions de déploiement

  1. Appliquer scripts/db/migrations/002_config_delivery.sql avant le nouveau serveur : scripts/db/apply-migration.sh <conteneur> <fichier_env> scripts/db/migrations/002_config_delivery.sql (mot de passe par MYSQL_PWD). Table dédiée boitier_config_deliveries, clé étrangère vers boitiers(id) en cascade ; boitiers n'est pas modifiée. Sans cette table, le serveur démarre avec un suivi en mémoire et le signale.
  2. Déployer configs/mosquitto/mosquitto.conf et redémarrer Mosquitto.
  3. MQTT_CLIENT_ID vide (défaut aquaserveur) ou valeur stable par environnement ; deux serveurs ne doivent jamais partager un client id. INGEST_WORKERS facultatif.
  4. Au premier démarrage, chaque configuration est publiée une fois (retenue) puis attend son accusé.

Choix faits

  • Filtre + au lieu de data-+ et des autres préfixes, invalides en MQTT 3.1.1 (un + occupe un niveau entier).
  • Message retenu par boîtier plutôt que session persistante du boîtier : aucun firmware à modifier ; en contrepartie, ce message reste lisible par tout client autorisé sur maj_serv-{id} tant que les ACL (L6) ne sont pas posées.
  • Table dédiée au suivi de livraison plutôt qu'une colonne de boitiers (table gérée par Laravel).
  • Ordre des accusés conservé par boîtier, pas entre boîtiers ; INGEST_WORKERS=1 rend l'ordre strict.
  • Erreur durable de la base accusée et journalisée plutôt que conservée.
  • Une trame écrite dont la réponse de la base s'est perdue revient en Duplicate au rejeu, sans alerte (à revoir avec la décision L4-bis « alerte même si l'écriture échoue »).
  • Valeurs : 4 workers, files de 256, réserve 64, inflight 20, file de session 100 000, 1 000 connexions, renvoi de 60 s à 30 min, vidage 8 s, rotation 5 x 10 Mo.

Hors périmètre

Authentification et ACL Mosquitto (L6) ; Dockerfiles, STOPSIGNAL, healthcheck réel N14 (L7) ; session persistante côté boîtier ; volumes mosquitto_log_* devenus inutiles ; documentation générale (L9).

Retour arrière

  • Code : revenir les 4 commits du lot et redéployer l'ancienne image ; remettre l'ancien mosquitto.conf.
  • Base : la table boitier_config_deliveries n'est lue que par ce serveur ; la laisser ou la supprimer (DROP TABLE boitier_config_deliveries).
  • Broker : effacer les messages retenus maj_serv-{id} en publiant un message vide retenu sur chaque topic (mosquitto_pub -r -n -t maj_serv-<id>).

Pile de PR

Base de cette PR : pr/dette-l4. Fusion dans l'ordre de la pile ; apres chaque fusion, vigie.py retarget rebase la PR suivante sur main.

  1. #2 pr/dette-l2 : Dette passe 1, lot L2 : contrats partagés dans aquashared et client MQTT du boîtier
  2. #3 pr/dette-l3 : Dette passe 1, lot L3 : provisioning non destructif et identifiant validé à l'entrée MQTT
  3. #4 pr/dette-l4 : Dette passe 1, lot L4 : mesures, alertes et déduplication
  4. #5 pr/dette-l5 : Dette passe 1, lot L5 : boucle MQTT du serveur, livraison de la configuration et journaux (cette PR)
  5. #6 pr/dette-l4bis : Dette passe 1, lot L4-bis : alertes même si l'écriture échoue, délai unique de 600 s, typage des capteurs, ordre des accusés MQTT
  6. #7 pr/dette-l7 : Dette passe 1, lot L7 : builds reproductibles, dépendances, durcissement des conteneurs
  7. #8 pr/dette-l6 : Dette passe 1, lot L6 : authentification Mosquitto, ACL par compte, healthchecks réels, pile e2e isolée
  8. #9 pr/dette-l8 : Dette passe 1, lot L8 : tests non destructifs, fixtures synthétiques, script de tests sur base jetable

Commits du lot

  • 5995bdf feat: structured server logs without payloads
  • ceb3b1a feat: bound server MQTT packets, narrow its subscription and back off
  • 6b6736a feat: write MQTT messages through bounded worker queues, ack after write
  • 74df6f6 feat: deliver boitier configuration until acknowledged
## Objet Lot L5 de la passe 1 de dette : la boucle MQTT du serveur ne dépend plus de la base (files bornées par worker, accusé MQTT après écriture, session persistante), la configuration descendante est livrée jusqu'à son accusé, l'arrêt est propre et les journaux sont structurés, sans payload. ## Dettes traitées et preuves | Dette | Correction | Preuve | |---|---|---| | Q2, abonnement | `#` remplacé par `+`, réabonnement à chaque `ConnAck` ; plus de topic à plusieurs niveaux ni de `$SYS` | `aquaserveur/src/mqtt/settings.rs`, `server.rs` | | N11, Q9, recommandation 7 | `set_max_packet_size` du client (8 Kio de payload plus le plus long topic) ; broker : `message_size_limit 8192`, `max_connections 1000`, `max_inflight_messages 20`, `max_queued_messages 100000` ; reconnexion `Backoff` 1 s à 60 s ; `MQTT_BROKER_URL` et `MQTT_CLIENT_ID` validés au démarrage | `mqtt/settings.rs`, `server.rs`, `configs/mosquitto/mosquitto.conf` | | N10, boucle et base | la boucle interroge le broker, route et dépose dans une file bornée par worker (256 places, réserve QoS 1 de 64) ; `INGEST_WORKERS` (défaut 4) ; un boîtier a un seul worker (ordre conservé) ; erreurs passagères de la base rejouées (1 s doublé jusqu'à 30 s) sans accuser, erreurs durables journalisées puis accusées | `aquaserveur/src/server.rs`, `ingest.rs`, `worker.rs`, `database/retry.rs`, `database/store.rs` | | Recommandation 10, C17 | session persistante (`clean_session = false`), client id stable, accusé manuel seulement après traitement ; rejeu sans effet grâce à l'empreinte des mesures (L4) | `mqtt/settings.rs`, `worker.rs` | | N7, configuration descendante | tâche `run_delivery` : version publiée, renvoi tant que l'accusé manque (60 s doublé jusqu'à 30 min), message retenu par boîtier sur `maj_serv-{id}`, état relu au démarrage depuis la table `boitier_config_deliveries` | `aquaserveur/src/config_delivery.rs`, `scripts/db/migrations/002_config_delivery.sql` | | N15, arrêt | SIGTERM et SIGINT : réception suspendue, files vidées par les workers, accusés envoyés, `DISCONNECT`, fermeture du pool ; vidage borné à 8 s | `server.rs`, `main.rs` | | N12, N13, Q15, journaux | `tracing_subscriber` avec `RUST_LOG` ; ni payload ni donnée personnelle (topic tronqué à 128 octets, erreur serde réduite à sa nature) ; rotation `json-file` 10 Mo x 5 dans les deux composes ; Mosquitto sur stdout | `aquaserveur/src/logging.rs`, `provisioning/error.rs`, composes | ## Tests et résultats - TDD : chaque module vu rouge avant implémentation (`logging` 7 sur 7, `MqttSettings` 7 sur 7, `ingest` 8 sur 8, `config_delivery` 9 sur 9, boucle serveur contre broker réel, SIGTERM). - Contre-épreuves (code neutralisé, test rouge, restauration vérifiée) : `clean_session(true)`, limite client à 10 Kio, boucle bloquée par la base, payload journalisé, workers sans vidage, accusé avant écriture, `maj_serv` non retenu, renvoi désactivé, suivi non relu au démarrage. - Ajouts : 44 tests unitaires ; 11 tests ignorés contre MariaDB et Mosquitto réels (`tests/mqtt_server_integration.rs`, `tests/config_delivery_integration.rs`), dont un serveur coupé brutalement au tiers d'un flux de 90 trames puis redémarré : 90 lignes, 90 epochs distincts, rejeu de 15 trames sans nouvelle ligne. - `cargo fmt --check`, `cargo clippy ... -D warnings` : OK. - `cargo test --workspace` : 332 passés, 0 échec, 68 ignorés (L4 : 288 et 57). - Tests ignorés (MariaDB 10.11 jetable avec dump, fixtures et migration 001, deux Mosquitto 2.0.22 jetables) : 68 passés, 0 échec, deux passages. - Vérifications avec le binaire : PUBLISH de 20 Kio sans coupure du serveur, SIGKILL pendant un flux de 60 trames puis redémarrage sans perte ni doublon, `docker stop` traité en 369 ms, livraison de configuration sur quatre démarrages successifs, migration 002 sur MariaDB 10.6.22 jetable (version de la prod). ## Migrations et actions de déploiement 1. Appliquer `scripts/db/migrations/002_config_delivery.sql` avant le nouveau serveur : `scripts/db/apply-migration.sh <conteneur> <fichier_env> scripts/db/migrations/002_config_delivery.sql` (mot de passe par `MYSQL_PWD`). Table dédiée `boitier_config_deliveries`, clé étrangère vers `boitiers(id)` en cascade ; `boitiers` n'est pas modifiée. Sans cette table, le serveur démarre avec un suivi en mémoire et le signale. 2. Déployer `configs/mosquitto/mosquitto.conf` et redémarrer Mosquitto. 3. `MQTT_CLIENT_ID` vide (défaut `aquaserveur`) ou valeur stable par environnement ; deux serveurs ne doivent jamais partager un client id. `INGEST_WORKERS` facultatif. 4. Au premier démarrage, chaque configuration est publiée une fois (retenue) puis attend son accusé. ## Choix faits - Filtre `+` au lieu de `data-+` et des autres préfixes, invalides en MQTT 3.1.1 (un `+` occupe un niveau entier). - Message retenu par boîtier plutôt que session persistante du boîtier : aucun firmware à modifier ; en contrepartie, ce message reste lisible par tout client autorisé sur `maj_serv-{id}` tant que les ACL (L6) ne sont pas posées. - Table dédiée au suivi de livraison plutôt qu'une colonne de `boitiers` (table gérée par Laravel). - Ordre des accusés conservé par boîtier, pas entre boîtiers ; `INGEST_WORKERS=1` rend l'ordre strict. - Erreur durable de la base accusée et journalisée plutôt que conservée. - Une trame écrite dont la réponse de la base s'est perdue revient en `Duplicate` au rejeu, sans alerte (à revoir avec la décision L4-bis « alerte même si l'écriture échoue »). - Valeurs : 4 workers, files de 256, réserve 64, inflight 20, file de session 100 000, 1 000 connexions, renvoi de 60 s à 30 min, vidage 8 s, rotation 5 x 10 Mo. ## Hors périmètre Authentification et ACL Mosquitto (L6) ; Dockerfiles, `STOPSIGNAL`, healthcheck réel N14 (L7) ; session persistante côté boîtier ; volumes `mosquitto_log_*` devenus inutiles ; documentation générale (L9). ## Retour arrière - Code : revenir les 4 commits du lot et redéployer l'ancienne image ; remettre l'ancien `mosquitto.conf`. - Base : la table `boitier_config_deliveries` n'est lue que par ce serveur ; la laisser ou la supprimer (`DROP TABLE boitier_config_deliveries`). - Broker : effacer les messages retenus `maj_serv-{id}` en publiant un message vide retenu sur chaque topic (`mosquitto_pub -r -n -t maj_serv-<id>`). --- ### Pile de PR Base de cette PR : `pr/dette-l4`. Fusion dans l'ordre de la pile ; apres chaque fusion, `vigie.py retarget` rebase la PR suivante sur `main`. 1. #2 `pr/dette-l2` : Dette passe 1, lot L2 : contrats partagés dans aquashared et client MQTT du boîtier 2. #3 `pr/dette-l3` : Dette passe 1, lot L3 : provisioning non destructif et identifiant validé à l'entrée MQTT 3. #4 `pr/dette-l4` : Dette passe 1, lot L4 : mesures, alertes et déduplication 4. #5 `pr/dette-l5` : Dette passe 1, lot L5 : boucle MQTT du serveur, livraison de la configuration et journaux (cette PR) 5. #6 `pr/dette-l4bis` : Dette passe 1, lot L4-bis : alertes même si l'écriture échoue, délai unique de 600 s, typage des capteurs, ordre des accusés MQTT 6. #7 `pr/dette-l7` : Dette passe 1, lot L7 : builds reproductibles, dépendances, durcissement des conteneurs 7. #8 `pr/dette-l6` : Dette passe 1, lot L6 : authentification Mosquitto, ACL par compte, healthchecks réels, pile e2e isolée 8. #9 `pr/dette-l8` : Dette passe 1, lot L8 : tests non destructifs, fixtures synthétiques, script de tests sur base jetable ### Commits du lot - `5995bdf` feat: structured server logs without payloads - `ceb3b1a` feat: bound server MQTT packets, narrow its subscription and back off - `6b6736a` feat: write MQTT messages through bounded worker queues, ack after write - `74df6f6` feat: deliver boitier configuration until acknowledged <!-- vigie:stack -->
Initialise tracing with an EnvFilter read from RUST_LOG (N12) and replace
the println/eprintln output of the server with levelled, timestamped
records carrying named fields (N13). Payloads and client data are never
logged: only topic, boitier id, size and outcome (Q15); firstco parse
errors are reduced to their kind because serde_json quotes payload values.
Docker logs rotate with json-file max-size/max-file on every dev and prod
service, and Mosquitto logs to stdout.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UH3JcaVWsmFLC6kyEsVL18
The server subscribed to '#': any PUBLISH above the 10 KiB rumqttc default
on any topic reset its connection, and it reconnected without delay (N11,
Q9, Q2). It now:
- sets max_packet_size from aquashared::mqtt::MQTT_MAX_PACKET_SIZE;
- subscribes to '+' only. Prefix filters such as 'data-+' are invalid MQTT
  (a wildcard must fill a whole level, MQTT-4.7.1-3), so with single-level
  topic names '+' is the narrowest filter that still receives firstco-{id}
  from unknown boitiers; multi-level and $SYS topics are no longer received;
- reconnects with aquashared::backoff::Backoff (1 s doubling to 60 s);
- reads MQTT_BROKER_URL and MQTT_CLIENT_ID through MqttSettings and refuses
  to start on an invalid value.
Mosquitto refuses payloads above 8 KiB (message_size_limit, below the
client limit), caps connections on the listener, bounds QoS 1 inflight
messages and sizes the persistent session queue.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UH3JcaVWsmFLC6kyEsVL18
The MQTT loop awaited every database query of a frame inside
eventloop.poll(): a slow database stopped PINGREQ and the broker cut the
connection (N10). The loop now only routes each PUBLISH into a bounded
per-worker queue and never waits:
- messages of one boitier always go to the same worker, so they are
  processed and acknowledged in order and its in-memory alert state has a
  single owner; INGEST_WORKERS (default 4) sets the worker count;
- when a queue fills up, QoS 0 is dropped first so that a 64-slot reserve
  stays free for QoS 1, which the broker caps at max_inflight_messages (20)
  unacknowledged messages; a QoS 1 refused on a full queue is left
  unacknowledged and redelivered at the next session resume.

Lossless delivery (recommendation 10, C17): persistent session with a
stable configurable client id (MQTT_CLIENT_ID) and manual acks sent only
once a frame is Inserted or Duplicate, or definitively refused. Transient
database errors are retried with a backoff without acknowledging. Replays
are harmless thanks to the L4 measure fingerprint.

Shutdown (N15): SIGTERM and SIGINT stop new processing, workers drain their
queues and the remaining acks are flushed before DISCONNECT; the pool is
closed last.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UH3JcaVWsmFLC6kyEsVL18
feat: deliver boitier configuration until acknowledged
Some checks failed
aquaprocess/revue-statique echec : 1 constat(s) bloquant(s) de la revue
74df6f6fa5
last_push was set when the maj_serv publication was queued, an offline
boitier missed the configuration, AckTimeout was never used and a restart
republished every boitier (N7). The delivery task now:
- tracks per boitier the version sent, its send count and the last
  version acknowledged on maj_boitier_ACK (ok or error), persisted in the
  new boitier_config_deliveries table (migration 002, idempotent, no change
  to the Laravel boitiers table);
- republishes an unacknowledged version after a timeout doubling from 60 s
  to 30 min (ConfigError::AckTimeout), never an acknowledged one, so a
  restart only resumes pending deliveries;
- publishes maj_serv-{id} retained: an offline boitier gets the latest
  configuration as soon as it resubscribes, with both firmwares and no
  per-boitier persistent session; the boitier re-applies an equal version
  idempotently.
firstco of an already provisioned boitier triggers an immediate push
through the same task. The unused send_config_and_wait_ack stub and the
non-retained push_config_to_mqtt are removed.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UH3JcaVWsmFLC6kyEsVL18
Thomas left a comment

Revue statique Vigie : PR #5

Tete 74df6f6fa5, base pr/dette-l4 (2886a1291a), 4 commit(s), 28 fichier(s) ajoutes ou modifies.

Statut aquaprocess/revue-statique : failure (1 constat(s) bloquant(s) de la revue)

Constats : 1 bloquant, 3 majeur, 19 mineur, 1 style. Commentaires en ligne : 2 (mode important, plafond 10).

Verifications automatiques

Verification Resultat Detail
cargo fmt --check success aucune difference
cargo clippy -D warnings success aucun avertissement
cargo test --workspace success 332 passes, 0 echec(s), 68 ignores
tests ignores (base jetable) success 68 passes, 0 echec(s) ; brokers MQTT : mosquitto, mosquitto-nolimit ; dump charge en 29 s ; 1 fichier(s) de fixtures ; migrations du depot : 001_mesure_empreinte.sql, 002_config_delivery.sql
cargo audit attention introduite(s) : aucune ; deja presente(s) sur la base : RUSTSEC-2026-0049 (rustls-webpki), RUSTSEC-2026-0098 (rustls-webpki), RUSTSEC-2026-0099 (rustls-webpki), RUSTSEC-2026-0104 (rustls-webpki), RUSTSEC-2026-0258 (h2), RUSTSEC-2026-0285 (rustls), 2 ignoree(s) par audit.toml
semgrep (regles du depot) success 21 fichier(s) ; lignes ajoutees : 13 resultat(s) hors tests, tests : no-unwrap-in-production x47, sqlx-no-format-in-query x2 ; 39 sur lignes inchangees
gitleaks (plage de commits) success 4 commit(s) analyses, aucune fuite
format des commits success 4 commit(s) conformes
emoji dans le code ajoute success 3584 ligne(s) ajoutee(s), aucune

Analyse qualitative

Verdict : lot solide et bien documenté, 1 bloquant (fragment de payload d'une trame refusée recopié dans les journaux, contraire au contrat annoncé), 3 majeurs à traiter avant la prod, le reste en améliorations.

La séparation boucle MQTT, workers et tâche de livraison est propre, l'accusé après écriture est correctement placé, aucun unwrap ni expect hors tests, et les références fichier:ligne du compte rendu sont exactes. Les réserves portent sur des chemins de saturation (files, canaux, inflight) qui ne se rétablissent pas seuls, et sur un trou dans le nettoyage des journaux.

Sécurité

  • Journaux : la troncature du topic et ProvisioningError::log_summary sont bien faites, mais une trame data- malformée journalise le message de parse_custom_frame, qui recopie le jeton fautif (Invalid ext sensor value: '...', Invalid epoch: '...') sans échappement ni borne (aquaserveur/src/worker.rs:183). Avec un broker encore anonyme, n'importe quel client peut écrire des lignes arbitraires dans les journaux. Même famille, en moins grave : l'accusé de configuration invalide (worker.rs:277).
  • Broker : max_queued_messages 100000 vaut pour toute session persistante, combiné à allow_anonymous true et à l'absence de persistent_client_expiration et de max_queued_bytes (configs/mosquitto/mosquitto.conf:27). Amplification de déni de service à fermer avant ou avec le lot L6.
  • Bornes : taille de paquet client et message_size_limit cohérentes (8 Kio de payload, limite client supérieure, contrôle compilé dans aquashared), max_connections 1000, files de 256 par worker, canal de requêtes de 1000. Abonnement réduit à + : justification MQTT-4.7.1-3 correcte, $SYS et topics à plusieurs niveaux exclus, tests à l'appui.
  • Messages retenus maj_serv-{id} : lisibles par tout client tant que les ACL manquent, conservés après suppression du boîtier ; reconnu par le compte rendu (renvoi L6).
  • Migration 002 : idempotente, sans modification de boitiers. La clé étrangère vers une table gérée par Laravel est un couplage à signaler (002_config_delivery.sql:35).
  • apply-migration.sh : mot de passe par MYSQL_PWD et docker exec -e sans valeur, ni en argument ni à l'écran, noms SQL validés par expression régulière. Seule la lecture du fichier d'environnement est fragile (apply-migration.sh:26).

Correction

  • Saturation sans retour : un QoS 1 refusé ou un try_ack en échec garde une place inflight que seule une reconnexion libère, et rien ne la provoque (server.rs:282). Les rechargements de règles contournent la réserve QoS 1 (ingest.rs:193), ce qui rend ce cas atteignable lors d'une modification groupée des configurations, surtout avec INGEST_WORKERS=1.
  • Interblocage possible entre workers et tâche de livraison : notify attend une place dans le canal de livraison pendant que la tâche de livraison attend une place dans la file d'un worker (worker.rs:300 et config_delivery.rs:355). La perte d'événement étant déjà admise par conception, try_send suffit.
  • Abonnement non retenté si le canal de requêtes est plein au ConnAck (server.rs:199).
  • Rejeu des erreurs passagères sans plafond ni alerte, avec Protocol et Tls classés passagers (worker.rs:99).
  • Arrêt : la séquence est bonne (réception suspendue, vidage, accusés, DISCONNECT, fermeture du pool), mais la tâche de livraison garde un clone des files, donc le vidage attend souvent les 8 s (server.rs:177), et le pire cas total dépasse les 10 s de docker stop (server.rs:48).
  • Livraison de configuration : la logique decide / on_ack est juste et bien testée. Le message retenu fait réappliquer une configuration refusée à chaque reconnexion et réinsère une alerte 15 (config_delivery.rs:291).
  • Point d'attention pour L4-bis (non bloquant ici) : le chemin Duplicate et les erreurs d'insertion d'alerte accusées sans rejeu reposent sur « alertes seulement sur Inserted » (worker.rs:172), que la décision de l'utilisateur (alerte même si l'écriture échoue) va changer.

Cohérence avec le compte rendu

Plus de 30 références vérifiées à la tête 74df6f6, toutes exactes : server.rs (run, dispatch, abonnement au ConnAck, délai de reconnexion, disconnect, signaux, délai de vidage), ingest.rs (réserve 64 avec contrôle compilé, admission, répartition, bornes 1 à 8, clé firstco), worker.rs (accusé après traitement, PushNow, transmission de l'accusé), mqtt/settings.rs, database/retry.rs, store.rs, data_handler.rs:181, provisioning/error.rs:50, logging.rs, config_delivery.rs, migration, mosquitto.conf, composes et main.rs. Plus aucun println! ni eprintln! dans aquaserveur/src hors --health-check.

Écarts :

  • « Les valeurs de la trame ne sont plus journalisées » et « ni payload » : faux pour les trames refusées (worker.rs:183).
  • « La réserve garde la place aux QoS 1 » : contournée par les rechargements (ingest.rs:193).
  • « Renvoyé par le broker à la reconnexion » : aucune reconnexion n'est provoquée (server.rs:282).
  • « Accusé de configuration idempotent » : un accusé error rejoué réinsère une alerte 15 (config_delivery.rs:291).
  • Comportement de rumqttc à la reconnexion (PUBACK en canal jetés) : la version 0.24 lue les replace dans pending ; le verrou cite 0.25.1, à vérifier (worker.rs:121).
  • Le test d'intégration « redémarrage » simule le redémarrage par injection d'état dans la table (config_delivery_integration.rs:426).

Tests

Bonne couverture unitaire des fonctions pures (admission, répartition, échéances, classification, journaux) et tests d'intégration réels convaincants pour l'absence de perte après coupure brutale, les PINGREQ pendant une base bloquée, le vidage à l'arrêt et le paquet de 20 Kio. Limites :

  • le test des journaux ne publie ni trame malformée ni accusé invalide (mqtt_server_integration.rs:389) ;
  • aucun test du chemin PushNow ni des saturations (canal de livraison plein, refus QoS 1, échec d'abonnement) ;
  • la tâche de livraison tourne dans tous les tests serveur et partage l'état entre binaires (tests/common/mqtt.rs:38) ;
  • plusieurs attentes fixes (tests/common/mqtt.rs:99 et suivants) rendent certains tests sensibles à la charge.

Questions ouvertes

  • Les boîtiers publient-ils en QoS 0 ou en QoS 1 ? En QoS 1 seul, le refus en file pleine devient improbable, sauf rafale de rechargements.
  • Comportement exact de EventLoop::clean() en 0.25.1 pour les accusés restés en canal.
  • Faut-il effacer le message retenu après un accusé error, ou le garder pour qu'un boîtier réessaie à chaque connexion (avec alerte 15 dédupliquée) ?
  • Le déploiement de mosquitto.conf est-il prévu avant L6 ? Si oui, persistent_client_expiration et max_queued_bytes devraient l'accompagner.

Autres constats (non publies en ligne)

Gravite Emplacement Source Constat
bloquant aquaserveur/src/worker.rs:183 revue Une trame refusée est journalisée avec error = %e : pour DataError::Parse, le message vient de aquashared::parse_custom_frame, qui recopie le jeton fautif (Invalid ext sensor value: '{t}', Invalid epoch: '{epoch_str}'), et EpochError recopie l'epoch. Le payload, choisi par l'émetteur (broker anonyme), atteint donc les journaux, sans échappement (% = Display, retours à la ligne possibles, donc fausses lignes de journal) et sans borne de taille (jusqu'à 8 Kio par ligne). C'est contraire au contrat annoncé (« ni payload », « valeurs de la trame plus journalisées ») : journaliser seulement la nature de l'erreur (comme ProvisioningError::log_summary) et ajouter une trame malformée au test des journaux. (non publie en ligne : Forgejo l'ancrerait a la ligne 180 du commit 6b6736a, decalee dans la vue de la PR)
majeur aquaserveur/src/server.rs:282 revue Un QoS 1 refusé (file pleine) n'est jamais accusé, et un try_ack qui échoue dans ack_now (canal de requêtes plein, l.298) non plus : chacun garde une place inflight du broker. Mosquitto ne renvoie qu'à la reconnexion, et rien ne la provoque ici. Après 20 cas, le broker cesse d'envoyer tout QoS 1 au serveur, sans erreur visible autre que ces journaux. Le commentaire « renvoyé par le broker à la reconnexion » suppose une reconnexion qui n'arrive pas. Piste : forcer une reconnexion (déconnexion du client) après un refus, ou garder les accusés en échec pour les renvoyer. (non publie en ligne : Forgejo l'ancrerait a la ligne 280 du commit 6b6736a, decalee dans la vue de la PR)
mineur aquaserveur/src/config_delivery.rs:291 revue Avec retain, chaque reconnexion d'un boîtier lui redonne la configuration ; s'il échoue à l'appliquer, il renvoie un accusé error et handle_config_ack insère une nouvelle alerte 15 à chaque fois (seulement bornée par admit_insert). Cela contredit « un accusé error : pas de renvoi de cette version » et l'idempotence de l'accusé de configuration invoquée pour l'accusé après écriture (worker.rs:14). Piste : ne pas réinsérer l'alerte 15 si (version, error) est déjà enregistré, ou effacer le message retenu après un accusé error.
mineur aquaserveur/src/data_handler.rs:360 semgrep semgrep no-unwrap-in-production (WARNING) : .unwrap() ou .expect() peut paniquer en production. 1 occurrence(s) sur lignes ajoutees (360), y compris d'eventuels modules de test internes.
mineur aquaserveur/src/ingest.rs:193 revue reload dépose un Job::Reload par send().await, sans tenir compte de QOS1_RESERVE : une rafale de rechargements (modification groupée de boitiers.updated_at, environ 113 par worker pour 450 boîtiers, 450 avec INGEST_WORKERS=1) peut occuper les places réservées au QoS 1. L'invariant de l'en-tête du module (« la réserve leur garde toujours de la place ») ne tient donc pas, et le refus qui en résulte bloque une place inflight (voir server.rs:282). Piste : borner les rechargements (try_send au-delà de la réserve, ou déduplication par boîtier).
mineur aquaserveur/src/ingest.rs:204 semgrep semgrep no-unwrap-in-production (WARNING) : .unwrap() ou .expect() peut paniquer en production. 2 occurrence(s) sur lignes ajoutees (204, 208), y compris d'eventuels modules de test internes.
mineur aquaserveur/src/logging.rs:91 semgrep semgrep no-unwrap-in-production (WARNING) : .unwrap() ou .expect() peut paniquer en production. 4 occurrence(s) sur lignes ajoutees (91, 108, 108, 155), y compris d'eventuels modules de test internes.
mineur aquaserveur/src/mqtt/settings.rs:231 semgrep semgrep no-unwrap-in-production (WARNING) : .unwrap() ou .expect() peut paniquer en production. 5 occurrence(s) sur lignes ajoutees (231, 242, 250, 255, 261), y compris d'eventuels modules de test internes.
mineur aquaserveur/src/provisioning/error.rs:97 semgrep semgrep no-unwrap-in-production (WARNING) : .unwrap() ou .expect() peut paniquer en production. 1 occurrence(s) sur lignes ajoutees (97), y compris d'eventuels modules de test internes.
mineur aquaserveur/src/server.rs:48 revue Le budget d'arrêt dépasse les 10 s de docker stop dans le pire cas : 8 s de vidage, puis jusqu'à 2 s d'attente du DISCONNECT (l.318), puis background.shutdown() et pool.close() (main.rs:80). Un SIGKILL peut donc tomber pendant la fermeture. Piste : réduire DRAIN_TIMEOUT (6 s par exemple) ou fixer stop_grace_period dans les composes.
mineur aquaserveur/src/server.rs:177 revue queue = None ne ferme pas les files : la tâche de livraison garde un clone d'IngestQueue et ne s'arrête qu'en revenant à son select!. Si elle est au milieu d'un passage (base lente, reload en attente), les workers ne voient jamais leur canal fermé et l'arrêt attend toujours les 8 s de DRAIN_TIMEOUT. Piste : donner à la tâche de livraison un émetteur faible (WeakSender) ou fermer les files explicitement.
mineur aquaserveur/src/server.rs:199 revue Si try_subscribe_many échoue (canal de requêtes de 1000 places plein au moment du ConnAck, par exemple après une coupure pendant laquelle publications de livraison et accusés se sont accumulés), l'erreur est journalisée et l'abonnement n'est jamais retenté. Avec session_present = false (première connexion, broker redémarré sans persistance), le serveur reste connecté sans rien recevoir. Piste : retenter à l'itération suivante, ou vérifier la réception du SUBACK et se reconnecter sinon.
mineur aquaserveur/src/worker.rs:99 revue Le rejeu d'une erreur passagère n'a pas de plafond : tant qu'elle revient, le message reste en tête de la file et tous les boîtiers de ce worker sont suspendus, avec un simple avertissement toutes les 30 s. database/retry.rs:22 classe aussi Protocol et Tls comme passagères ; si l'une d'elles était déterministe pour une requête donnée, ce message bloquerait le worker indéfiniment. Piste : compteur exposé et alerte au-delà de N essais, et classification plus restrictive.
mineur aquaserveur/src/worker.rs:121 revue Après une reconnexion, un message déjà en file est accusé avec le pkid de l'ancienne connexion, puis sa copie renvoyée par le broker (même pkid) est retraitée en Duplicate et accusée une seconde fois. Si le broker a entre-temps réattribué ce pkid, ce second PUBACK accuserait un autre message avant son écriture. Le risque semble faible avec Mosquitto (pkid croissants), mais le compte rendu affirme que rumqttc jette les PUBACK en canal à la reconnexion, alors que dans la version 0.24 consultée EventLoop::clean() les déplace dans pending pour les renvoyer : à vérifier pour 0.25.1.
mineur aquaserveur/src/worker.rs:172 revue Point d'attention pour L4-bis (non bloquant ici) : l'accusé après écriture repose sur « alertes seulement sur Inserted ». Une trame écrite dont la réponse s'est perdue, ou un arrêt entre l'écriture de la mesure et celle des alertes, revient en Duplicate et ses alertes ne partent jamais. De plus, une erreur passagère pendant l'insertion d'alertes est seulement journalisée (data_handler.rs, « alertes légionellose non insérées ») et le message est accusé. La décision « alerte même si l'écriture échoue » et le délai unique de 600 s (server.rs:54) devront être raccordés à Step::Retry sans doubler les alertes au rejeu.
mineur aquaserveur/src/worker.rs:277 revue Un accusé de configuration refusé est journalisé avec error = %e ; pour ConfigAckError::InvalidPayload, ce message est celui de serde_json et peut recopier une valeur du payload (échappée). Même raisonnement que pour firstco (log_summary) : réduire à la nature de l'erreur.
mineur aquaserveur/tests/common/mqtt.rs:38 revue Tous les tests de la boucle serveur démarrent aussi la tâche de livraison avec les réglages par défaut : ils publient des maj_serv-{id} retenus sur le broker partagé et écrivent dans boitier_config_deliveries si la table existe. Les deux binaires de test partagent donc un état que le verrou SERIAL (par binaire) ne protège pas. Piste : un réglage qui désactive la livraison dans server_config, ou un nettoyage explicite.
mineur aquaserveur/tests/config_delivery_integration.rs:426 revue « a_restart_republishes_only_unacknowledged_configurations » ne fait aucun redémarrage : l'état est injecté par save_state, puis un seul serveur démarre. Le chemin réel (suivi écrit par une instance, relu par la suivante, envoi non accusé repris à l'échéance et non au démarrage) n'est couvert que par le test unitaire de decide. Piste : démarrer une première instance, l'arrêter avant l'accusé, puis vérifier l'absence de publication immédiate au redémarrage.
mineur aquaserveur/tests/mqtt_server_integration.rs:389 revue Le test « payloads_and_client_data_never_reach_the_logs » ne publie qu'une trame valide, un firstco refusé et un firstco illisible. Il ne couvre ni une trame data- malformée ni un maj_boitier_ACK- invalide, qui sont justement les chemins qui recopient le payload (worker.rs:183, :277). Ajouter ces deux cas avec un marqueur dans le payload.
mineur scripts/db/apply-migration.sh:26 revue Avec set -euo pipefail, une clé absente du fichier d'environnement fait échouer grep dans la substitution : le script s'arrête sans message, et le contrôle « DB_PASSWORD vide » (l.37) n'est atteint que pour une valeur vide. Une valeur entre guillemets (DB_PASSWORD="...", courant dans un .env) est transmise avec ses guillemets. Le passage par MYSQL_PWD hors argument est correct.
mineur scripts/db/migrations/002_config_delivery.sql:35 revue La clé étrangère vers boitiers(id) lie une table du serveur à une table possédée par Laravel : un TRUNCATE boitiers, une migration Laravel qui modifie le type de boitiers.id ou recrée la table échouera désormais. Le choix se défend (nettoyage en cascade, refus d'un id inconnu), mais il mérite d'être signalé à l'équipe Laravel dans la documentation de déploiement.
style aquaserveur/tests/common/mqtt.rs:99 revue Plusieurs tests reposent sur des attentes fixes (500 ms après le premier ConnAck ici, 3 s dans le test de redémarrage, 6 s dans le test de renvoi, 1 s avant le compte des trames en file) plutôt que sur un événement observé (SUBACK, compteur). Ils peuvent devenir instables sur une machine chargée ; préférer wait_until sur une condition.

Limites

Revue publiee en COMMENT : l'auteur des PR et le relecteur sont le meme compte, Forgejo refuse APPROVE et REQUEST_CHANGES. Le blocage s'exprime par le statut de commit aquaprocess/revue-statique. Analyse statique et tests automatises seulement, sans fusion ni deploiement.

## Revue statique Vigie : PR #5 Tete `74df6f6fa5`, base `pr/dette-l4` (`2886a1291a`), 4 commit(s), 28 fichier(s) ajoutes ou modifies. **Statut `aquaprocess/revue-statique` : failure** (1 constat(s) bloquant(s) de la revue) Constats : 1 bloquant, 3 majeur, 19 mineur, 1 style. Commentaires en ligne : 2 (mode `important`, plafond 10). ### Verifications automatiques | Verification | Resultat | Detail | |---|---|---| | cargo fmt --check | success | aucune difference | | cargo clippy -D warnings | success | aucun avertissement | | cargo test --workspace | success | 332 passes, 0 echec(s), 68 ignores | | tests ignores (base jetable) | success | 68 passes, 0 echec(s) ; brokers MQTT : mosquitto, mosquitto-nolimit ; dump charge en 29 s ; 1 fichier(s) de fixtures ; migrations du depot : 001_mesure_empreinte.sql, 002_config_delivery.sql | | cargo audit | attention | introduite(s) : aucune ; deja presente(s) sur la base : RUSTSEC-2026-0049 (rustls-webpki), RUSTSEC-2026-0098 (rustls-webpki), RUSTSEC-2026-0099 (rustls-webpki), RUSTSEC-2026-0104 (rustls-webpki), RUSTSEC-2026-0258 (h2), RUSTSEC-2026-0285 (rustls), 2 ignoree(s) par audit.toml | | semgrep (regles du depot) | success | 21 fichier(s) ; lignes ajoutees : 13 resultat(s) hors tests, tests : no-unwrap-in-production x47, sqlx-no-format-in-query x2 ; 39 sur lignes inchangees | | gitleaks (plage de commits) | success | 4 commit(s) analyses, aucune fuite | | format des commits | success | 4 commit(s) conformes | | emoji dans le code ajoute | success | 3584 ligne(s) ajoutee(s), aucune | ### Analyse qualitative Verdict : lot solide et bien documenté, 1 bloquant (fragment de payload d'une trame refusée recopié dans les journaux, contraire au contrat annoncé), 3 majeurs à traiter avant la prod, le reste en améliorations. La séparation boucle MQTT, workers et tâche de livraison est propre, l'accusé après écriture est correctement placé, aucun `unwrap` ni `expect` hors tests, et les références fichier:ligne du compte rendu sont exactes. Les réserves portent sur des chemins de saturation (files, canaux, inflight) qui ne se rétablissent pas seuls, et sur un trou dans le nettoyage des journaux. ### Sécurité - Journaux : la troncature du topic et `ProvisioningError::log_summary` sont bien faites, mais une trame `data-` malformée journalise le message de `parse_custom_frame`, qui recopie le jeton fautif (`Invalid ext sensor value: '...'`, `Invalid epoch: '...'`) sans échappement ni borne (`aquaserveur/src/worker.rs:183`). Avec un broker encore anonyme, n'importe quel client peut écrire des lignes arbitraires dans les journaux. Même famille, en moins grave : l'accusé de configuration invalide (`worker.rs:277`). - Broker : `max_queued_messages 100000` vaut pour toute session persistante, combiné à `allow_anonymous true` et à l'absence de `persistent_client_expiration` et de `max_queued_bytes` (`configs/mosquitto/mosquitto.conf:27`). Amplification de déni de service à fermer avant ou avec le lot L6. - Bornes : taille de paquet client et `message_size_limit` cohérentes (8 Kio de payload, limite client supérieure, contrôle compilé dans aquashared), `max_connections 1000`, files de 256 par worker, canal de requêtes de 1000. Abonnement réduit à `+` : justification MQTT-4.7.1-3 correcte, `$SYS` et topics à plusieurs niveaux exclus, tests à l'appui. - Messages retenus `maj_serv-{id}` : lisibles par tout client tant que les ACL manquent, conservés après suppression du boîtier ; reconnu par le compte rendu (renvoi L6). - Migration 002 : idempotente, sans modification de `boitiers`. La clé étrangère vers une table gérée par Laravel est un couplage à signaler (`002_config_delivery.sql:35`). - `apply-migration.sh` : mot de passe par `MYSQL_PWD` et `docker exec -e` sans valeur, ni en argument ni à l'écran, noms SQL validés par expression régulière. Seule la lecture du fichier d'environnement est fragile (`apply-migration.sh:26`). ### Correction - Saturation sans retour : un QoS 1 refusé ou un `try_ack` en échec garde une place inflight que seule une reconnexion libère, et rien ne la provoque (`server.rs:282`). Les rechargements de règles contournent la réserve QoS 1 (`ingest.rs:193`), ce qui rend ce cas atteignable lors d'une modification groupée des configurations, surtout avec `INGEST_WORKERS=1`. - Interblocage possible entre workers et tâche de livraison : `notify` attend une place dans le canal de livraison pendant que la tâche de livraison attend une place dans la file d'un worker (`worker.rs:300` et `config_delivery.rs:355`). La perte d'événement étant déjà admise par conception, `try_send` suffit. - Abonnement non retenté si le canal de requêtes est plein au `ConnAck` (`server.rs:199`). - Rejeu des erreurs passagères sans plafond ni alerte, avec `Protocol` et `Tls` classés passagers (`worker.rs:99`). - Arrêt : la séquence est bonne (réception suspendue, vidage, accusés, DISCONNECT, fermeture du pool), mais la tâche de livraison garde un clone des files, donc le vidage attend souvent les 8 s (`server.rs:177`), et le pire cas total dépasse les 10 s de `docker stop` (`server.rs:48`). - Livraison de configuration : la logique `decide` / `on_ack` est juste et bien testée. Le message retenu fait réappliquer une configuration refusée à chaque reconnexion et réinsère une alerte 15 (`config_delivery.rs:291`). - Point d'attention pour L4-bis (non bloquant ici) : le chemin `Duplicate` et les erreurs d'insertion d'alerte accusées sans rejeu reposent sur « alertes seulement sur Inserted » (`worker.rs:172`), que la décision de l'utilisateur (alerte même si l'écriture échoue) va changer. ### Cohérence avec le compte rendu Plus de 30 références vérifiées à la tête `74df6f6`, toutes exactes : `server.rs` (`run`, `dispatch`, abonnement au `ConnAck`, délai de reconnexion, `disconnect`, signaux, délai de vidage), `ingest.rs` (réserve 64 avec contrôle compilé, admission, répartition, bornes 1 à 8, clé `firstco`), `worker.rs` (accusé après traitement, `PushNow`, transmission de l'accusé), `mqtt/settings.rs`, `database/retry.rs`, `store.rs`, `data_handler.rs:181`, `provisioning/error.rs:50`, `logging.rs`, `config_delivery.rs`, migration, `mosquitto.conf`, composes et `main.rs`. Plus aucun `println!` ni `eprintln!` dans `aquaserveur/src` hors `--health-check`. Écarts : - « Les valeurs de la trame ne sont plus journalisées » et « ni payload » : faux pour les trames refusées (`worker.rs:183`). - « La réserve garde la place aux QoS 1 » : contournée par les rechargements (`ingest.rs:193`). - « Renvoyé par le broker à la reconnexion » : aucune reconnexion n'est provoquée (`server.rs:282`). - « Accusé de configuration idempotent » : un accusé `error` rejoué réinsère une alerte 15 (`config_delivery.rs:291`). - Comportement de `rumqttc` à la reconnexion (PUBACK en canal jetés) : la version 0.24 lue les replace dans `pending` ; le verrou cite 0.25.1, à vérifier (`worker.rs:121`). - Le test d'intégration « redémarrage » simule le redémarrage par injection d'état dans la table (`config_delivery_integration.rs:426`). ### Tests Bonne couverture unitaire des fonctions pures (admission, répartition, échéances, classification, journaux) et tests d'intégration réels convaincants pour l'absence de perte après coupure brutale, les PINGREQ pendant une base bloquée, le vidage à l'arrêt et le paquet de 20 Kio. Limites : - le test des journaux ne publie ni trame malformée ni accusé invalide (`mqtt_server_integration.rs:389`) ; - aucun test du chemin `PushNow` ni des saturations (canal de livraison plein, refus QoS 1, échec d'abonnement) ; - la tâche de livraison tourne dans tous les tests serveur et partage l'état entre binaires (`tests/common/mqtt.rs:38`) ; - plusieurs attentes fixes (`tests/common/mqtt.rs:99` et suivants) rendent certains tests sensibles à la charge. ### Questions ouvertes - Les boîtiers publient-ils en QoS 0 ou en QoS 1 ? En QoS 1 seul, le refus en file pleine devient improbable, sauf rafale de rechargements. - Comportement exact de `EventLoop::clean()` en 0.25.1 pour les accusés restés en canal. - Faut-il effacer le message retenu après un accusé `error`, ou le garder pour qu'un boîtier réessaie à chaque connexion (avec alerte 15 dédupliquée) ? - Le déploiement de `mosquitto.conf` est-il prévu avant L6 ? Si oui, `persistent_client_expiration` et `max_queued_bytes` devraient l'accompagner. ### Autres constats (non publies en ligne) | Gravite | Emplacement | Source | Constat | |---|---|---|---| | bloquant | `aquaserveur/src/worker.rs:183` | revue | Une trame refusée est journalisée avec `error = %e` : pour `DataError::Parse`, le message vient de `aquashared::parse_custom_frame`, qui recopie le jeton fautif (`Invalid ext sensor value: '{t}'`, `Invalid epoch: '{epoch_str}'`), et `EpochError` recopie l'epoch. Le payload, choisi par l'émetteur (broker anonyme), atteint donc les journaux, sans échappement (`%` = Display, retours à la ligne possibles, donc fausses lignes de journal) et sans borne de taille (jusqu'à 8 Kio par ligne). C'est contraire au contrat annoncé (« ni payload », « valeurs de la trame plus journalisées ») : journaliser seulement la nature de l'erreur (comme `ProvisioningError::log_summary`) et ajouter une trame malformée au test des journaux. (non publie en ligne : Forgejo l'ancrerait a la ligne 180 du commit 6b6736a, decalee dans la vue de la PR) | | majeur | `aquaserveur/src/server.rs:282` | revue | Un QoS 1 refusé (file pleine) n'est jamais accusé, et un `try_ack` qui échoue dans `ack_now` (canal de requêtes plein, l.298) non plus : chacun garde une place inflight du broker. Mosquitto ne renvoie qu'à la reconnexion, et rien ne la provoque ici. Après 20 cas, le broker cesse d'envoyer tout QoS 1 au serveur, sans erreur visible autre que ces journaux. Le commentaire « renvoyé par le broker à la reconnexion » suppose une reconnexion qui n'arrive pas. Piste : forcer une reconnexion (déconnexion du client) après un refus, ou garder les accusés en échec pour les renvoyer. (non publie en ligne : Forgejo l'ancrerait a la ligne 280 du commit 6b6736a, decalee dans la vue de la PR) | | mineur | `aquaserveur/src/config_delivery.rs:291` | revue | Avec `retain`, chaque reconnexion d'un boîtier lui redonne la configuration ; s'il échoue à l'appliquer, il renvoie un accusé `error` et `handle_config_ack` insère une nouvelle alerte 15 à chaque fois (seulement bornée par `admit_insert`). Cela contredit « un accusé error : pas de renvoi de cette version » et l'idempotence de l'accusé de configuration invoquée pour l'accusé après écriture (worker.rs:14). Piste : ne pas réinsérer l'alerte 15 si `(version, error)` est déjà enregistré, ou effacer le message retenu après un accusé `error`. | | mineur | `aquaserveur/src/data_handler.rs:360` | semgrep | semgrep `no-unwrap-in-production` (WARNING) : `.unwrap()` ou `.expect()` peut paniquer en production. 1 occurrence(s) sur lignes ajoutees (360), y compris d'eventuels modules de test internes. | | mineur | `aquaserveur/src/ingest.rs:193` | revue | `reload` dépose un `Job::Reload` par `send().await`, sans tenir compte de `QOS1_RESERVE` : une rafale de rechargements (modification groupée de `boitiers.updated_at`, environ 113 par worker pour 450 boîtiers, 450 avec `INGEST_WORKERS=1`) peut occuper les places réservées au QoS 1. L'invariant de l'en-tête du module (« la réserve leur garde toujours de la place ») ne tient donc pas, et le refus qui en résulte bloque une place inflight (voir server.rs:282). Piste : borner les rechargements (`try_send` au-delà de la réserve, ou déduplication par boîtier). | | mineur | `aquaserveur/src/ingest.rs:204` | semgrep | semgrep `no-unwrap-in-production` (WARNING) : `.unwrap()` ou `.expect()` peut paniquer en production. 2 occurrence(s) sur lignes ajoutees (204, 208), y compris d'eventuels modules de test internes. | | mineur | `aquaserveur/src/logging.rs:91` | semgrep | semgrep `no-unwrap-in-production` (WARNING) : `.unwrap()` ou `.expect()` peut paniquer en production. 4 occurrence(s) sur lignes ajoutees (91, 108, 108, 155), y compris d'eventuels modules de test internes. | | mineur | `aquaserveur/src/mqtt/settings.rs:231` | semgrep | semgrep `no-unwrap-in-production` (WARNING) : `.unwrap()` ou `.expect()` peut paniquer en production. 5 occurrence(s) sur lignes ajoutees (231, 242, 250, 255, 261), y compris d'eventuels modules de test internes. | | mineur | `aquaserveur/src/provisioning/error.rs:97` | semgrep | semgrep `no-unwrap-in-production` (WARNING) : `.unwrap()` ou `.expect()` peut paniquer en production. 1 occurrence(s) sur lignes ajoutees (97), y compris d'eventuels modules de test internes. | | mineur | `aquaserveur/src/server.rs:48` | revue | Le budget d'arrêt dépasse les 10 s de `docker stop` dans le pire cas : 8 s de vidage, puis jusqu'à 2 s d'attente du DISCONNECT (l.318), puis `background.shutdown()` et `pool.close()` (main.rs:80). Un SIGKILL peut donc tomber pendant la fermeture. Piste : réduire `DRAIN_TIMEOUT` (6 s par exemple) ou fixer `stop_grace_period` dans les composes. | | mineur | `aquaserveur/src/server.rs:177` | revue | `queue = None` ne ferme pas les files : la tâche de livraison garde un clone d'`IngestQueue` et ne s'arrête qu'en revenant à son `select!`. Si elle est au milieu d'un passage (base lente, `reload` en attente), les workers ne voient jamais leur canal fermé et l'arrêt attend toujours les 8 s de `DRAIN_TIMEOUT`. Piste : donner à la tâche de livraison un émetteur faible (`WeakSender`) ou fermer les files explicitement. | | mineur | `aquaserveur/src/server.rs:199` | revue | Si `try_subscribe_many` échoue (canal de requêtes de 1000 places plein au moment du `ConnAck`, par exemple après une coupure pendant laquelle publications de livraison et accusés se sont accumulés), l'erreur est journalisée et l'abonnement n'est jamais retenté. Avec `session_present = false` (première connexion, broker redémarré sans persistance), le serveur reste connecté sans rien recevoir. Piste : retenter à l'itération suivante, ou vérifier la réception du SUBACK et se reconnecter sinon. | | mineur | `aquaserveur/src/worker.rs:99` | revue | Le rejeu d'une erreur passagère n'a pas de plafond : tant qu'elle revient, le message reste en tête de la file et tous les boîtiers de ce worker sont suspendus, avec un simple avertissement toutes les 30 s. `database/retry.rs:22` classe aussi `Protocol` et `Tls` comme passagères ; si l'une d'elles était déterministe pour une requête donnée, ce message bloquerait le worker indéfiniment. Piste : compteur exposé et alerte au-delà de N essais, et classification plus restrictive. | | mineur | `aquaserveur/src/worker.rs:121` | revue | Après une reconnexion, un message déjà en file est accusé avec le pkid de l'ancienne connexion, puis sa copie renvoyée par le broker (même pkid) est retraitée en `Duplicate` et accusée une seconde fois. Si le broker a entre-temps réattribué ce pkid, ce second PUBACK accuserait un autre message avant son écriture. Le risque semble faible avec Mosquitto (pkid croissants), mais le compte rendu affirme que `rumqttc` jette les PUBACK en canal à la reconnexion, alors que dans la version 0.24 consultée `EventLoop::clean()` les déplace dans `pending` pour les renvoyer : à vérifier pour 0.25.1. | | mineur | `aquaserveur/src/worker.rs:172` | revue | Point d'attention pour L4-bis (non bloquant ici) : l'accusé après écriture repose sur « alertes seulement sur `Inserted` ». Une trame écrite dont la réponse s'est perdue, ou un arrêt entre l'écriture de la mesure et celle des alertes, revient en `Duplicate` et ses alertes ne partent jamais. De plus, une erreur passagère pendant l'insertion d'alertes est seulement journalisée (data_handler.rs, « alertes légionellose non insérées ») et le message est accusé. La décision « alerte même si l'écriture échoue » et le délai unique de 600 s (server.rs:54) devront être raccordés à `Step::Retry` sans doubler les alertes au rejeu. | | mineur | `aquaserveur/src/worker.rs:277` | revue | Un accusé de configuration refusé est journalisé avec `error = %e` ; pour `ConfigAckError::InvalidPayload`, ce message est celui de `serde_json` et peut recopier une valeur du payload (échappée). Même raisonnement que pour `firstco` (`log_summary`) : réduire à la nature de l'erreur. | | mineur | `aquaserveur/tests/common/mqtt.rs:38` | revue | Tous les tests de la boucle serveur démarrent aussi la tâche de livraison avec les réglages par défaut : ils publient des `maj_serv-{id}` retenus sur le broker partagé et écrivent dans `boitier_config_deliveries` si la table existe. Les deux binaires de test partagent donc un état que le verrou `SERIAL` (par binaire) ne protège pas. Piste : un réglage qui désactive la livraison dans `server_config`, ou un nettoyage explicite. | | mineur | `aquaserveur/tests/config_delivery_integration.rs:426` | revue | « a_restart_republishes_only_unacknowledged_configurations » ne fait aucun redémarrage : l'état est injecté par `save_state`, puis un seul serveur démarre. Le chemin réel (suivi écrit par une instance, relu par la suivante, envoi non accusé repris à l'échéance et non au démarrage) n'est couvert que par le test unitaire de `decide`. Piste : démarrer une première instance, l'arrêter avant l'accusé, puis vérifier l'absence de publication immédiate au redémarrage. | | mineur | `aquaserveur/tests/mqtt_server_integration.rs:389` | revue | Le test « payloads_and_client_data_never_reach_the_logs » ne publie qu'une trame valide, un firstco refusé et un firstco illisible. Il ne couvre ni une trame `data-` malformée ni un `maj_boitier_ACK-` invalide, qui sont justement les chemins qui recopient le payload (worker.rs:183, :277). Ajouter ces deux cas avec un marqueur dans le payload. | | mineur | `scripts/db/apply-migration.sh:26` | revue | Avec `set -euo pipefail`, une clé absente du fichier d'environnement fait échouer `grep` dans la substitution : le script s'arrête sans message, et le contrôle « DB_PASSWORD vide » (l.37) n'est atteint que pour une valeur vide. Une valeur entre guillemets (`DB_PASSWORD="..."`, courant dans un `.env`) est transmise avec ses guillemets. Le passage par `MYSQL_PWD` hors argument est correct. | | mineur | `scripts/db/migrations/002_config_delivery.sql:35` | revue | La clé étrangère vers `boitiers(id)` lie une table du serveur à une table possédée par Laravel : un `TRUNCATE boitiers`, une migration Laravel qui modifie le type de `boitiers.id` ou recrée la table échouera désormais. Le choix se défend (nettoyage en cascade, refus d'un id inconnu), mais il mérite d'être signalé à l'équipe Laravel dans la documentation de déploiement. | | style | `aquaserveur/tests/common/mqtt.rs:99` | revue | Plusieurs tests reposent sur des attentes fixes (500 ms après le premier ConnAck ici, 3 s dans le test de redémarrage, 6 s dans le test de renvoi, 1 s avant le compte des trames en file) plutôt que sur un événement observé (SUBACK, compteur). Ils peuvent devenir instables sur une machine chargée ; préférer `wait_until` sur une condition. | ### Limites Revue publiee en `COMMENT` : l'auteur des PR et le relecteur sont le meme compte, Forgejo refuse APPROVE et REQUEST_CHANGES. Le blocage s'exprime par le statut de commit `aquaprocess/revue-statique`. Analyse statique et tests automatises seulement, sans fusion ni deploiement. <!-- vigie:review -->
@ -0,0 +297,4 @@
/// l'arrêt du serveur : l'événement est alors perdu, la livraison reprend
/// sur l'état enregistré au démarrage suivant).
async fn notify(deps: &WorkerDeps, event: DeliveryEvent) {
if deps.delivery.send(event).await.is_err() {
Author
Member

[majeur] notify attend une place dans le canal de livraison (send().await, 256 places) alors que la tâche de livraison, pendant un passage, attend elle-même une place dans la file d'un worker (queue.reload(id).await, config_delivery.rs:355) sans lire ses événements. Si le canal de livraison est plein pendant qu'une file de worker est pleine, le worker bloqué ne vide plus sa file et la tâche bloquée ne vide plus le canal : interblocage permanent, sans coupure ni journal, et l'ingestion de ce worker s'arrête (puis des autres, bloqués à leur tour sur notify). Même sans interblocage, un passage long (450 publications et écritures séquentielles, interval en mode rafale) suspend tous les workers qui reçoivent un accusé de configuration. La perte d'un événement étant déjà admise (commentaire l.296), utiliser try_send et journaliser, ou découpler les rechargements dans une tâche dédiée.

Source : revue.

**[majeur]** `notify` attend une place dans le canal de livraison (`send().await`, 256 places) alors que la tâche de livraison, pendant un passage, attend elle-même une place dans la file d'un worker (`queue.reload(id).await`, config_delivery.rs:355) sans lire ses événements. Si le canal de livraison est plein pendant qu'une file de worker est pleine, le worker bloqué ne vide plus sa file et la tâche bloquée ne vide plus le canal : interblocage permanent, sans coupure ni journal, et l'ingestion de ce worker s'arrête (puis des autres, bloqués à leur tour sur `notify`). Même sans interblocage, un passage long (450 publications et écritures séquentielles, `interval` en mode rafale) suspend tous les workers qui reçoivent un accusé de configuration. La perte d'un événement étant déjà admise (commentaire l.296), utiliser `try_send` et journaliser, ou découpler les rechargements dans une tâche dédiée. _Source : revue._
@ -7,0 +24,4 @@
# File d'une session persistante pendant une coupure du serveur. Ordre de
# grandeur : 450 boîtiers à une trame par minute, soit environ 3 h 40 avant
# que le broker ne commence à écarter des messages.
max_queued_messages 100000
Author
Member

[majeur] max_queued_messages 100000 s'applique à toutes les sessions persistantes, pas seulement à celle du serveur (défaut Mosquitto : 1000). Avec allow_anonymous true (l.1, inchangé), sans max_queued_bytes ni persistent_client_expiration, n'importe quel client peut ouvrir des sessions persistantes sous des client id différents, s'abonner puis se déconnecter : chacune accumule jusqu'à 100 000 messages en mémoire et dans la persistance, sans expiration (max_connections ne borne pas les sessions déconnectées). Ajouter max_queued_bytes et persistent_client_expiration, et ne pas déployer cette valeur avant les ACL du lot L6.

Source : revue.

**[majeur]** `max_queued_messages 100000` s'applique à toutes les sessions persistantes, pas seulement à celle du serveur (défaut Mosquitto : 1000). Avec `allow_anonymous true` (l.1, inchangé), sans `max_queued_bytes` ni `persistent_client_expiration`, n'importe quel client peut ouvrir des sessions persistantes sous des client id différents, s'abonner puis se déconnecter : chacune accumule jusqu'à 100 000 messages en mémoire et dans la persistance, sans expiration (`max_connections` ne borne pas les sessions déconnectées). Ajouter `max_queued_bytes` et `persistent_client_expiration`, et ne pas déployer cette valeur avant les ACL du lot L6. _Source : revue._
Some checks failed
aquaprocess/revue-statique echec : 1 constat(s) bloquant(s) de la revue
This pull request can be merged automatically.
You are not authorized to merge this pull request.
View command line instructions

Checkout

From your project repository, check out a new branch and test the changes.
git fetch -u origin pr/dette-l5:pr/dette-l5
git switch pr/dette-l5

Merge

Merge the changes and update on Forgejo.

Warning: The "Autodetect manual merge" setting is not enabled for this repository, you will have to mark this pull request as manually merged afterwards.

git switch pr/dette-l4
git merge --no-ff pr/dette-l5
git switch pr/dette-l5
git rebase pr/dette-l4
git switch pr/dette-l4
git merge --ff-only pr/dette-l5
git switch pr/dette-l5
git rebase pr/dette-l4
git switch pr/dette-l4
git merge --no-ff pr/dette-l5
git switch pr/dette-l4
git merge --squash pr/dette-l5
git switch pr/dette-l4
git merge --ff-only pr/dette-l5
git switch pr/dette-l4
git merge pr/dette-l5
git push origin pr/dette-l4
Sign in to join this conversation.
No reviewers
No labels
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference
AcadeNice/Acquarefactoring-Thomas!5
No description provided.