Files
veye-lalwa/pipeline/collecteurs/http.py
T
Cyber MawonajandClaude Opus 5 feb5544599 Phase 5 : collecteurs officiels et run de veille — dry-run contre sources réelles
Collecteurs, chacun isolé et tolérant à l'échec :
- legifrance : OAuth PISTE validé en production. Deux pièges vérifiés sur le
  compte réel — seul le couple « Client ID / Client secret » ouvre un jeton
  (la clé d'API donne invalid_client), et les identifiants de production ne
  passent pas en bac à sable.
- senat : liste chronologique des lois promulguées. L'URL du §5.1
  (senat.fr/lois/index.html) est morte ; l'équivalent fonctionnel est
  /dossiers-legislatifs/lois-promulguees.html — 31 lois 2026 relevées.
- conseil_constitutionnel : sans le paramètre ?id=32246 la page des affaires en
  instance ne rend que des QPC et aucune affaire DC. Les colonnes DC et QPC
  diffèrent : lecture par en-tête, jamais par position.

Rapprochement à trois issues — numéro officiel, similarité de titre, et une
zone de signalement où le pipeline s'abstient plutôt que de fusionner.

make update-dry contre les sources réelles : aucun faux positif sur les 52
textes du corpus, et trois apports que le rapport n'avait pas — le numéro
d'affaire 2026-907 DC de la programmation militaire, la date de saisine du
16 juillet pour l'aide à mourir, les décisions du 23 juillet sur 908 et 909.

Une asymétrie de sources est documentée par un test plutôt que corrigée en
douce : le Sénat ne qualifie pas la loi 2026-650 d'organique, Légifrance si.
Le collecteur rapporte ce que sa source dit ; le champ  n'est pas comparé.

Fixtures HTML figées : aucun test ne touche au réseau.
156 tests, 91 % de couverture (chaque module au-dessus du seuil de 70 %).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 21:48:56 -04:00

219 lines
7.9 KiB
Python

"""Client HTTP partagé par les collecteurs.
Trois obligations, non négociables vis-à-vis de services publics gratuits :
1. **Se présenter.** Le user-agent porte le nom du projet et une adresse de
contact, configurables dans `.env`. Un scraper anonyme est un scraper qu'on
bloque, à raison.
2. **Ne pas marteler.** Un délai minimal est respecté entre deux requêtes vers
un même hôte, et un cache disque évite de redemander ce qu'on a déjà.
3. **Ne pas insister bêtement.** Trois tentatives au maximum, avec attente
croissante, et seulement sur les erreurs qui peuvent se résoudre d'elles-mêmes
(temporisations, 5xx, 429).
"""
from __future__ import annotations
import hashlib
import json
import os
import time
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any
import httpx
from pipeline import chemins
from pipeline.journal import logger
log = logger("http")
CODES_A_REESSAYER = frozenset({408, 425, 429, 500, 502, 503, 504})
@dataclass(slots=True)
class Reponse:
"""Réponse HTTP, éventuellement servie par le cache."""
url: str
code: int
texte: str
depuis_le_cache: bool = False
@property
def a_reussi(self) -> bool:
return 200 <= self.code < 300
def json(self) -> Any:
return json.loads(self.texte)
class ClientHttp:
"""Client à cache disque et débit maîtrisé."""
def __init__(
self,
*,
cache: Path | None = None,
duree_cache_h: float | None = None,
delai_minimal_s: float | None = None,
user_agent: str | None = None,
hors_ligne: bool = False,
) -> None:
self.cache = cache or chemins.CACHE_HTTP
self.duree_cache = timedelta(
hours=duree_cache_h
if duree_cache_h is not None
else float(os.environ.get("VEILLE_CACHE_TTL_H", "6"))
)
self.delai_minimal = (
delai_minimal_s
if delai_minimal_s is not None
else float(os.environ.get("VEILLE_DELAI_REQUETES", "1.0"))
)
self.user_agent = user_agent or os.environ.get(
"VEILLE_USER_AGENT", "veille-legislative-971/1.0"
)
# Mode hors ligne : seul le cache répond. Indispensable pour que les
# tests ne dépendent jamais du réseau.
self.hors_ligne = hors_ligne
self._dernier_appel: dict[str, float] = {}
self._client = httpx.Client(
timeout=httpx.Timeout(30.0, connect=15.0),
follow_redirects=True,
headers={"User-Agent": self.user_agent, "Accept-Language": "fr-FR,fr;q=0.9"},
)
# ── Cycle de vie ─────────────────────────────────────────────────────────
def fermer(self) -> None:
self._client.close()
def __enter__(self) -> ClientHttp:
return self
def __exit__(self, *_) -> None:
self.fermer()
# ── Requêtes ─────────────────────────────────────────────────────────────
def get(self, url: str, **kwargs) -> Reponse:
return self._requeter("GET", url, **kwargs)
def post(self, url: str, **kwargs) -> Reponse:
return self._requeter("POST", url, **kwargs)
def _requeter(
self,
methode: str,
url: str,
*,
json_corps: Any = None,
donnees: dict | None = None,
entetes: dict | None = None,
utiliser_cache: bool = True,
tentatives: int = 3,
) -> Reponse:
cle = _cle_de_cache(methode, url, json_corps, donnees)
if utiliser_cache and (en_cache := self._lire_cache(cle)):
return Reponse(url=url, code=200, texte=en_cache, depuis_le_cache=True)
if self.hors_ligne:
return Reponse(url=url, code=0, texte="", depuis_le_cache=False)
derniere_erreur: str | None = None
for tentative in range(1, tentatives + 1):
self._respecter_le_debit(url)
try:
reponse = self._client.request(
methode, url, json=json_corps, data=donnees, headers=entetes
)
except httpx.HTTPError as erreur:
derniere_erreur = f"{type(erreur).__name__}: {erreur}"
log.debug("requête en échec", url=url, tentative=tentative, erreur=derniere_erreur)
time.sleep(min(2**tentative, 8))
continue
if reponse.status_code in CODES_A_REESSAYER and tentative < tentatives:
attente = _attente_conseillee(reponse) or min(2**tentative, 8)
log.debug(
"réponse à réessayer",
url=url,
code=reponse.status_code,
attente_s=attente,
)
time.sleep(attente)
continue
if 200 <= reponse.status_code < 300 and utiliser_cache:
self._ecrire_cache(cle, reponse.text)
return Reponse(url=url, code=reponse.status_code, texte=reponse.text)
raise httpx.HTTPError(derniere_erreur or f"échec après {tentatives} tentatives : {url}")
# ── Débit ────────────────────────────────────────────────────────────────
def _respecter_le_debit(self, url: str) -> None:
hote = httpx.URL(url).host or ""
precedent = self._dernier_appel.get(hote)
if precedent is not None:
attente = self.delai_minimal - (time.monotonic() - precedent)
if attente > 0:
time.sleep(attente)
self._dernier_appel[hote] = time.monotonic()
# ── Cache ────────────────────────────────────────────────────────────────
def _chemin_cache(self, cle: str) -> Path:
return self.cache / cle[:2] / f"{cle}.txt"
def _lire_cache(self, cle: str) -> str | None:
fichier = self._chemin_cache(cle)
if not fichier.exists():
return None
age = datetime.now(UTC) - datetime.fromtimestamp(fichier.stat().st_mtime, UTC)
if age > self.duree_cache:
return None
return fichier.read_text(encoding="utf-8")
def _ecrire_cache(self, cle: str, contenu: str) -> None:
fichier = self._chemin_cache(cle)
fichier.parent.mkdir(parents=True, exist_ok=True)
fichier.write_text(contenu, encoding="utf-8")
def vider_le_cache(self) -> int:
"""Supprime les entrées périmées. Retourne le nombre de fichiers effacés."""
efface = 0
if not self.cache.exists():
return 0
limite = datetime.now(UTC) - self.duree_cache
for fichier in self.cache.rglob("*.txt"):
if datetime.fromtimestamp(fichier.stat().st_mtime, UTC) < limite:
fichier.unlink()
efface += 1
return efface
def _cle_de_cache(methode: str, url: str, json_corps: Any, donnees: dict | None) -> str:
empreinte = hashlib.sha256()
empreinte.update(methode.encode())
empreinte.update(url.encode())
if json_corps is not None:
empreinte.update(json.dumps(json_corps, sort_keys=True, ensure_ascii=False).encode())
if donnees:
empreinte.update(json.dumps(donnees, sort_keys=True, ensure_ascii=False).encode())
return empreinte.hexdigest()
def _attente_conseillee(reponse: httpx.Response) -> float | None:
"""Respecte l'en-tête `Retry-After` quand le serveur en envoie un."""
valeur = reponse.headers.get("Retry-After")
if not valeur:
return None
try:
return min(float(valeur), 30.0)
except ValueError:
return None