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>
This commit is contained in:
co-authored by
Claude Opus 5
parent
c9fc8f3f67
commit
feb5544599
@@ -0,0 +1,218 @@
|
||||
"""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
|
||||
Reference in New Issue
Block a user