Files
veye-lalwa/pipeline/db.py
T
Cyber MawonajandClaude Opus 5 c36b8dc631 Phase 2 : parseur de l'annexe sec10 — 52 textes, 12 affaires CC
Briques réutilisables :
- dates_fr : lit les formes réelles du corpus (« 23 avr. 2026 (JO 24 avr.) »,
  « 1er juill. », « 24.07.2026 ») et distingue une date exacte d'une échéance
  prévisionnelle (« Fin août 2026 » → situé au 25, libellé conservé, jamais
  présenté comme un fait)
- markdown : découpage des tableaux (barres échappées, lignes courtes),
  nettoyage, extraction des marqueurs [^dim07-24^] / [^phase5-3.1^]

Parseur sec10 :
- 27 lois promulguées + 21 lignes simples + 2 lignes éclatées en 4 textes = 52
- éclatement justifié : les lignes « accords Kazakhstan et Colombie » et
  « MACF » décrivent chacune deux textes de statuts différents
- migration 002 : regime_libertes stocké à part, faute de savoir quel public
  la cotation globale du rapport vise — dérouler « droits+ » sur les trois
  publics aurait inventé une ventilation absente de la source
- migration 003 : la facette « devant le CC » couvre aussi les saisines sans
  numéro d'affaire (programmation militaire), sinon 6 textes au lieu de 7

Trois pièges du corpus, corrigés et couverts par des tests :
- « eau » matchait comme sous-chaîne dans réseaux / drapeau / Retailleau
- une décision citée sans saisine créait une fausse saisine (loi Philippine)
- la date de saisine était pêchée dans la colonne « échéance », datant la
  saisine de l'aide à mourir du mois suivant

72 tests, 90 % de couverture sur pipeline/.

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

383 lines
15 KiB
Python

"""Accès SQLite : migrations versionnées, écriture des textes, requêtes de veille.
Une seule règle : rien n'entre en base sans passer par un modèle pydantic
validé. Les fonctions de ce module reçoivent des `Texte`, jamais des
dictionnaires bruts venus d'un parseur.
"""
from __future__ import annotations
import json
import sqlite3
from collections.abc import Iterable, Iterator
from contextlib import contextmanager
from pathlib import Path
from typing import Any
from pipeline import chemins
from pipeline.journal import logger
from pipeline.modeles import DecisionCC, Evenement, Insight, Source, Texte
log = logger("db")
# ─────────────────────────────────────────────────────────────────────────────
# Connexion et migrations
# ─────────────────────────────────────────────────────────────────────────────
def connecter(base: Path | None = None, *, lecture_seule: bool = False) -> sqlite3.Connection:
"""Ouvre la base avec les réglages attendus par le reste du pipeline."""
base = base or chemins.BASE_SQLITE
base.parent.mkdir(parents=True, exist_ok=True)
if lecture_seule:
cx = sqlite3.connect(f"file:{base}?mode=ro", uri=True)
else:
cx = sqlite3.connect(base)
cx.row_factory = sqlite3.Row
cx.execute("PRAGMA foreign_keys = ON")
if not lecture_seule:
cx.execute("PRAGMA journal_mode = WAL")
cx.execute("PRAGMA synchronous = NORMAL")
return cx
@contextmanager
def transaction(cx: sqlite3.Connection) -> Iterator[sqlite3.Connection]:
"""Valide à la sortie, annule intégralement à la moindre exception."""
try:
yield cx
except Exception:
cx.rollback()
raise
else:
cx.commit()
def _version_schema(cx: sqlite3.Connection) -> int:
return int(cx.execute("PRAGMA user_version").fetchone()[0])
def migrer(cx: sqlite3.Connection, dossier: Path | None = None) -> list[str]:
"""Applique les migrations non encore jouées, dans l'ordre des noms.
Retourne la liste des migrations appliquées lors de cet appel.
"""
dossier = dossier or chemins.MIGRATIONS
fichiers = sorted(dossier.glob("*.sql"))
version = _version_schema(cx)
appliquees: list[str] = []
for fichier in fichiers:
numero = int(fichier.name.split("_", 1)[0])
if numero <= version:
continue
log.info("migration", fichier=fichier.name, numero=numero)
cx.executescript(fichier.read_text(encoding="utf-8"))
cx.execute(f"PRAGMA user_version = {numero}")
cx.commit()
appliquees.append(fichier.name)
return appliquees
def initialiser(base: Path | None = None, *, repartir_de_zero: bool = False) -> sqlite3.Connection:
"""Ouvre la base et s'assure que le schéma est à jour."""
base = base or chemins.BASE_SQLITE
if repartir_de_zero and base.exists():
for suffixe in ("", "-wal", "-shm"):
Path(str(base) + suffixe).unlink(missing_ok=True)
log.info("base supprimée avant réinitialisation", base=str(base))
cx = connecter(base)
migrer(cx)
return cx
# ─────────────────────────────────────────────────────────────────────────────
# Écriture
# ─────────────────────────────────────────────────────────────────────────────
_CHAMPS_TEXTE = (
"id", "numero_officiel", "type", "titre_court", "titre_officiel",
"statut", "statut_date", "date_depot", "date_adoption", "date_promulgation",
"date_entree_vigueur", "prochaine_echeance", "prochaine_echeance_label",
"themes", "impacts", "regime_libertes", "regime_libertes_note",
"guadeloupe_pertinence", "guadeloupe_note",
"resume", "points_cles", "confiance", "source_seed",
"a_verifier", "motif_verification",
)
def enregistrer_texte(cx: sqlite3.Connection, texte: Texte) -> None:
"""Insère ou met à jour un texte, avec ses sources, événements et décisions."""
ligne = texte.en_ligne_sql()
colonnes = ", ".join(_CHAMPS_TEXTE)
valeurs = ", ".join(f":{c}" for c in _CHAMPS_TEXTE)
maj = ", ".join(f"{c} = excluded.{c}" for c in _CHAMPS_TEXTE if c != "id")
cx.execute(
f"INSERT INTO textes ({colonnes}) VALUES ({valeurs}) "
f"ON CONFLICT (id) DO UPDATE SET {maj}, maj_le = datetime('now')",
ligne,
)
for source in texte.sources:
enregistrer_source(cx, texte.id, source)
for evenement in texte.evenements:
enregistrer_evenement(cx, texte.id, evenement)
for decision in texte.decisions_cc:
enregistrer_decision_cc(cx, decision, texte_id=texte.id)
def enregistrer_source(cx: sqlite3.Connection, texte_id: str, source: Source) -> None:
cx.execute(
"""
INSERT INTO sources (texte_id, url, titre, editeur, date_publication, tier,
extrait_verbatim, contexte, confiance, marqueur, fichier_origine)
VALUES (:texte_id, :url, :titre, :editeur, :date_publication, :tier,
:extrait_verbatim, :contexte, :confiance, :marqueur, :fichier_origine)
ON CONFLICT (texte_id, url, COALESCE(marqueur, '')) DO UPDATE SET
titre = COALESCE(excluded.titre, titre),
editeur = COALESCE(excluded.editeur, editeur),
date_publication = COALESCE(excluded.date_publication, date_publication),
tier = excluded.tier,
extrait_verbatim = COALESCE(excluded.extrait_verbatim, extrait_verbatim),
contexte = COALESCE(excluded.contexte, contexte),
confiance = COALESCE(excluded.confiance, confiance)
""",
{
"texte_id": texte_id,
"url": source.url,
"titre": source.titre,
"editeur": source.editeur,
"date_publication": source.date_publication.isoformat()
if source.date_publication
else None,
"tier": str(source.tier),
"extrait_verbatim": source.extrait_verbatim,
"contexte": source.contexte,
"confiance": str(source.confiance) if source.confiance else None,
"marqueur": source.marqueur,
"fichier_origine": source.fichier_origine,
},
)
def enregistrer_evenement(cx: sqlite3.Connection, texte_id: str, evenement: Evenement) -> None:
cx.execute(
"""
INSERT INTO evenements (texte_id, date_evenement, type_etape, description,
source_url, previsionnel)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT (texte_id, date_evenement, type_etape, description) DO UPDATE SET
source_url = COALESCE(excluded.source_url, source_url),
previsionnel = excluded.previsionnel
""",
(
texte_id,
evenement.date_evenement.isoformat(),
str(evenement.type_etape),
evenement.description,
evenement.source_url,
int(evenement.previsionnel),
),
)
def enregistrer_decision_cc(
cx: sqlite3.Connection, decision: DecisionCC, *, texte_id: str | None = None
) -> None:
cx.execute(
"""
INSERT INTO decisions_cc (texte_id, numero_affaire, date_saisine, date_decision,
resultat, saisissants, resume, url, date_decision_attendue)
VALUES (:texte_id, :numero_affaire, :date_saisine, :date_decision,
:resultat, :saisissants, :resume, :url, :date_decision_attendue)
ON CONFLICT (numero_affaire) DO UPDATE SET
texte_id = COALESCE(excluded.texte_id, texte_id),
date_saisine = COALESCE(excluded.date_saisine, date_saisine),
date_decision = COALESCE(excluded.date_decision, date_decision),
resultat = COALESCE(excluded.resultat, resultat),
saisissants = COALESCE(excluded.saisissants, saisissants),
resume = COALESCE(excluded.resume, resume),
url = COALESCE(excluded.url, url),
date_decision_attendue = COALESCE(excluded.date_decision_attendue,
date_decision_attendue),
maj_le = datetime('now')
""",
{
"texte_id": texte_id,
"numero_affaire": decision.numero_affaire,
"date_saisine": decision.date_saisine.isoformat() if decision.date_saisine else None,
"date_decision": decision.date_decision.isoformat()
if decision.date_decision
else None,
"resultat": str(decision.resultat) if decision.resultat else None,
"saisissants": decision.saisissants,
"resume": decision.resume,
"url": decision.url,
"date_decision_attendue": decision.date_decision_attendue.isoformat()
if decision.date_decision_attendue
else None,
},
)
def enregistrer_insight(cx: sqlite3.Connection, insight: Insight) -> None:
cx.execute(
"""
INSERT INTO insights (numero, titre, corps, implications, confiance,
derive_de, fichier_origine)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (numero) DO UPDATE SET
titre = excluded.titre, corps = excluded.corps,
implications = excluded.implications, confiance = excluded.confiance,
derive_de = excluded.derive_de, fichier_origine = excluded.fichier_origine
""",
(
insight.numero,
insight.titre,
insight.corps,
insight.implications,
insight.confiance,
json.dumps(insight.derive_de, ensure_ascii=False),
insight.fichier_origine,
),
)
def enregistrer_textes(cx: sqlite3.Connection, textes: Iterable[Texte]) -> int:
"""Écrit un lot de textes dans une transaction unique."""
total = 0
with transaction(cx):
for texte in textes:
enregistrer_texte(cx, texte)
total += 1
return total
# ─────────────────────────────────────────────────────────────────────────────
# Journal des exécutions
# ─────────────────────────────────────────────────────────────────────────────
def ouvrir_run(cx: sqlite3.Connection, mode: str) -> int:
curseur = cx.execute("INSERT INTO veille_log (mode) VALUES (?)", (mode,))
cx.commit()
return int(curseur.lastrowid)
def cloturer_run(
cx: sqlite3.Connection,
run_id: int,
*,
duree_s: float,
ajouts: int = 0,
modifications: int = 0,
echecs: int = 0,
alertes: list[str] | None = None,
rapport: str | None = None,
) -> None:
cx.execute(
"""
UPDATE veille_log
SET duree_s = ?, ajouts = ?, modifications = ?, echecs = ?,
alertes = ?, rapport = ?
WHERE id = ?
""",
(
duree_s,
ajouts,
modifications,
echecs,
json.dumps(alertes or [], ensure_ascii=False),
rapport,
run_id,
),
)
cx.commit()
def enregistrer_changement(
cx: sqlite3.Connection,
run_id: int,
*,
texte_id: str | None,
nature: str,
champ: str | None = None,
ancienne_valeur: str | None = None,
nouvelle_valeur: str | None = None,
description: str | None = None,
) -> None:
cx.execute(
"""
INSERT INTO veille_changements (run_id, texte_id, nature, champ,
ancienne_valeur, nouvelle_valeur, description)
VALUES (?, ?, ?, ?, ?, ?, ?)
""",
(run_id, texte_id, nature, champ, ancienne_valeur, nouvelle_valeur, description),
)
# ─────────────────────────────────────────────────────────────────────────────
# Lecture
# ─────────────────────────────────────────────────────────────────────────────
def statistiques(cx: sqlite3.Connection) -> dict[str, Any]:
"""Compteurs affichés en fin de seed et sur le tableau de bord."""
def un(requete: str) -> int:
return int(cx.execute(requete).fetchone()[0])
return {
"textes": un("SELECT COUNT(*) FROM textes"),
"sources": un("SELECT COUNT(*) FROM sources"),
"evenements": un("SELECT COUNT(*) FROM evenements"),
"decisions_cc": un("SELECT COUNT(*) FROM decisions_cc"),
"insights": un("SELECT COUNT(*) FROM insights"),
"textes_sans_source": un(
"SELECT COUNT(*) FROM textes t "
"WHERE NOT EXISTS (SELECT 1 FROM sources s WHERE s.texte_id = t.id)"
),
"a_verifier": un("SELECT COUNT(*) FROM textes WHERE a_verifier = 1"),
"par_statut": {
r["statut"]: r["n"]
for r in cx.execute(
"SELECT statut, COUNT(*) AS n FROM textes GROUP BY statut ORDER BY n DESC"
)
},
"par_pertinence_gpe": {
(r["guadeloupe_pertinence"] or "non_cotee"): r["n"]
for r in cx.execute(
"SELECT guadeloupe_pertinence, COUNT(*) AS n FROM textes "
"GROUP BY guadeloupe_pertinence"
)
},
"devant_cc": un("SELECT COUNT(*) FROM v_textes WHERE devant_cc = 1"),
}
def main() -> None:
"""`make db-init` — crée la base si besoin et applique les migrations."""
import argparse
from pipeline.journal import configurer
analyseur = argparse.ArgumentParser(description="Initialise ou migre la base SQLite.")
analyseur.add_argument(
"--repartir-de-zero",
action="store_true",
help="supprime la base existante avant de recréer le schéma",
)
arguments = analyseur.parse_args()
configurer()
cx = initialiser(repartir_de_zero=arguments.repartir_de_zero)
log.info(
"base prête",
chemin=str(chemins.BASE_SQLITE),
version_schema=_version_schema(cx),
**statistiques(cx) | {"par_statut": None, "par_pertinence_gpe": None},
)
if __name__ == "__main__":
main()