Files

219 lines
7.9 KiB
Python
Raw Permalink Normal View History

"""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