feat(r5.1): socle apps/ai — ingestion anonymisée + recherche sémantique

ADR-004 : embeddings locaux sur CPU (fastembed ONNX,
paraphrase-multilingual-MiniLM-L12-v2, 384 dims — les textes du client
ne quittent jamais le serveur), pgvector dans le Postgres existant
(RagChunk possédé par Prisma, migration r5_ia + état de corpus sur
Document), génération opt-in (mode extractif par défaut : la recette
passe sans clé API), service siop2-ai jamais exposé — joint par l'API
NestJS seule (X-Service-Token).

apps/ai (FastAPI + uv) : pipeline PDF MinIO → texte paginé (pypdf) →
anonymisation D4 (e-mails, téléphones marocains, noms connus de la
base, insensible casse/accents — fonction pure testée) → découpage
avec chevauchement (testé) → embeddings → RagChunk localisé (« p. 42 »,
« bilan du 17/07 »). Bilans codés clôturés ingérés. Exclusion de
corpus (D3) appliquée à l'ingestion ET à la lecture.

14 pytest + ruff, embeddeur déterministe en CI (aucun téléchargement),
job CI ai (uv), deploy en dépend. Vérifié en réel avec le vrai modèle :
corpus seedé réindexé en 7 s (PDF réel → 31 extraits paginés + 3
bilans), recherche sémantique concluante, e-mails → ⟨contact⟩,
0 identité dans les chunks (contrôle SQL).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
pr-daaif
2026-07-17 12:22:15 +01:00
parent 199fce69d0
commit 837dcba1db
22 changed files with 2564 additions and 4 deletions

View File

@@ -0,0 +1 @@
"""SIOP V2 — service IA (R5). L'IA propose, l'humain valide (D1)."""

View File

@@ -0,0 +1,62 @@
"""Anonymisation à l'ingestion (D4, loi 09-08) : identités et coordonnées ne
partent JAMAIS dans les index vectoriels ni dans les prompts.
Fonction pure, testée : e-mails, téléphones (formats marocains et
internationaux), et les noms de personnes CONNUS de la base (utilisateurs,
gardiens, contacts tiers) fournis par l'appelant.
"""
import re
import unicodedata
JETON_PERSONNE = "⟨personne⟩"
JETON_CONTACT = "⟨contact⟩"
_EMAIL = re.compile(r"[\w.+-]+@[\w-]+\.[\w.-]+")
# 06 12 34 56 78 · 0612345678 · +212 6 12 34 56 78 · 05 22-34-56-78…
_TELEPHONE = re.compile(r"(?:\+?\d{1,3}[\s.-]?)?(?:0|\(0\))?\d(?:[\s.-]?\d{2}){4}")
def _sans_accents(texte: str) -> str:
return "".join(
c for c in unicodedata.normalize("NFD", texte) if unicodedata.category(c) != "Mn"
)
def anonymiser(texte: str, noms_connus: list[str] | None = None) -> str:
"""Remplace coordonnées et noms connus par des jetons neutres.
Les noms sont remplacés insensiblement à la casse ET aux accents
(« Idrissi » attrape « idrissi »), prénom seul compris quand il est
assez long pour ne pas mutiler le texte technique.
"""
resultat = _EMAIL.sub(JETON_CONTACT, texte)
resultat = _TELEPHONE.sub(JETON_CONTACT, resultat)
for nom in sorted(noms_connus or [], key=len, reverse=True):
nom = nom.strip()
if len(nom) < 3:
continue
morceaux = [nom] + [m for m in nom.split() if len(m) >= 4]
for morceau in morceaux:
motif = re.compile(
r"\b" + re.escape(_sans_accents(morceau)) + r"\b", re.IGNORECASE
)
# on cherche sur une copie sans accents mais on remplace l'original
copie = _sans_accents(resultat)
sortie: list[str] = []
position = 0
for correspondance in motif.finditer(copie):
sortie.append(resultat[position : correspondance.start()])
sortie.append(JETON_PERSONNE)
position = correspondance.end()
sortie.append(resultat[position:])
resultat = "".join(sortie)
# jetons collés en double (« prénom nom » remplacés séparément)
resultat = re.sub(
rf"{re.escape(JETON_PERSONNE)}(\s+{re.escape(JETON_PERSONNE)})+",
JETON_PERSONNE,
resultat,
)
return resultat

View File

@@ -0,0 +1,62 @@
"""Service IA — JAMAIS exposé publiquement (ADR-004 §4) : seule l'API NestJS
le contacte, avec le secret partagé `X-Service-Token`. Les permissions des
utilisateurs restent l'affaire de l'API — ici, un seul appelant de confiance.
"""
from contextlib import asynccontextmanager
from dataclasses import asdict
import asyncpg
from fastapi import Depends, FastAPI, Header, HTTPException
from pydantic import BaseModel, Field
from .config import Reglages, charger_reglages
from .embeddings import construire_embeddeur
from .ingestion import reindexer_tout
from .recherche import chercher
@asynccontextmanager
async def cycle_de_vie(app: FastAPI):
reglages = charger_reglages()
app.state.reglages = reglages
app.state.embeddeur = construire_embeddeur(reglages.ai_embeddings)
app.state.pool = await asyncpg.create_pool(reglages.database_url, min_size=1, max_size=5)
yield
await app.state.pool.close()
app = FastAPI(title="SIOP V2 — service IA (R5)", lifespan=cycle_de_vie)
def verifier_jeton(
x_service_token: str = Header(default=""),
) -> None:
reglages: Reglages = app.state.reglages
if x_service_token != reglages.ai_service_token:
raise HTTPException(status_code=401, detail="Jeton de service invalide")
@app.get("/healthz")
async def sante() -> dict:
"""Sonde interne (compose/Dokploy) — ne révèle rien du corpus."""
return {"status": "ok", "service": "siop2-ai"}
@app.post("/internal/reindex", dependencies=[Depends(verifier_jeton)])
async def reindexer() -> dict:
async with app.state.pool.acquire() as cnx:
resultat = await reindexer_tout(cnx, app.state.reglages, app.state.embeddeur)
return asdict(resultat)
class RequeteRecherche(BaseModel):
question: str = Field(min_length=3, max_length=500)
limite: int = Field(default=5, ge=1, le=10)
@app.post("/internal/search", dependencies=[Depends(verifier_jeton)])
async def rechercher(corps: RequeteRecherche) -> dict:
async with app.state.pool.acquire() as cnx:
extraits = await chercher(cnx, app.state.embeddeur, corps.question, corps.limite)
return {"extraits": [asdict(e) for e in extraits]}

View File

@@ -0,0 +1,32 @@
"""Configuration — validée au démarrage, comme l'API NestJS (même philosophie)."""
from pydantic_settings import BaseSettings
class Reglages(BaseSettings):
# Postgres partagé (schéma possédé par Prisma — apps/api)
database_url: str = "postgresql://siop:siop@localhost:5432/siop"
# MinIO en direct sur le réseau privé (ADR-004 §4 — jamais exposé)
minio_endpoint: str = "localhost"
minio_port: int = 9000
minio_use_ssl: bool = False
minio_access_key: str = "siop"
minio_secret_key: str = "siop-minio"
minio_bucket: str = "siop2"
# Le service n'est JAMAIS public : seul l'API NestJS le contacte,
# porteuse de ce secret partagé (ADR-004 §4).
ai_service_token: str = "dev-only-ai-token"
# Embeddings : « locale » (fastembed ONNX) ou « deterministe » (tests/CI)
ai_embeddings: str = "locale"
# Génération : « off » (mode extractif, défaut honnête) ou « api » (opt-in)
ai_generation: str = "off"
model_config = {"env_prefix": "", "case_sensitive": False}
def charger_reglages() -> Reglages:
reglages = Reglages()
# asyncpg ne comprend pas le paramètre ?schema= de Prisma
if "?" in reglages.database_url:
reglages.database_url = reglages.database_url.split("?")[0]
return reglages

View File

@@ -0,0 +1,52 @@
"""Découpage du texte en extraits indexables — pur et testé.
Paragraphes regroupés jusqu'à ~900 caractères, avec un chevauchement de
queue pour ne pas couper une prescription en deux. Un extrait trop long est
scindé sur les phrases.
"""
import re
TAILLE_CIBLE = 900
CHEVAUCHEMENT = 150
TAILLE_MINIMALE = 40 # en deçà : bruit (titres orphelins, numéros de page)
def _phrases(texte: str) -> list[str]:
return [p.strip() for p in re.split(r"(?<=[.!?;])\s+", texte) if p.strip()]
def decouper(texte: str) -> list[str]:
paragraphes = [p.strip() for p in re.split(r"\n\s*\n", texte) if p.strip()]
extraits: list[str] = []
courant = ""
def pousser() -> None:
nonlocal courant
nettoye = courant.strip()
if len(nettoye) >= TAILLE_MINIMALE:
extraits.append(nettoye)
courant = ""
for paragraphe in paragraphes:
paragraphe = re.sub(r"\s+", " ", paragraphe)
if len(courant) + len(paragraphe) + 1 > TAILLE_CIBLE and courant:
queue = courant[-CHEVAUCHEMENT:]
pousser()
courant = queue + " "
while len(paragraphe) > TAILLE_CIBLE:
phrases = _phrases(paragraphe)
if len(phrases) <= 1:
courant += paragraphe[:TAILLE_CIBLE]
paragraphe = paragraphe[TAILLE_CIBLE - CHEVAUCHEMENT :]
pousser()
continue
morceau = ""
while phrases and len(morceau) + len(phrases[0]) + 1 <= TAILLE_CIBLE:
morceau += phrases.pop(0) + " "
courant += morceau
pousser()
paragraphe = " ".join(phrases)
courant += paragraphe + " "
pousser()
return extraits

View File

@@ -0,0 +1,53 @@
"""Embeddeurs (ADR-004) : le vrai modèle local ONNX, et un déterministe pour
tests/CI — même interface, mêmes 384 dimensions, aucun téléchargement en test.
"""
import hashlib
import math
from typing import Protocol
DIMENSIONS = 384
MODELE_LOCAL = "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2"
class Embeddeur(Protocol):
def encoder(self, textes: list[str]) -> list[list[float]]: ...
class EmbeddeurDeterministe:
"""Sac de tri-grammes haché puis normalisé : stable, sans réseau, et les
textes proches partagent des composantes — assez pour tester le circuit
complet (ingestion → pgvector → similarité)."""
def encoder(self, textes: list[str]) -> list[list[float]]:
return [self._un(t) for t in textes]
def _un(self, texte: str) -> list[float]:
vecteur = [0.0] * DIMENSIONS
mots = texte.lower().split()
grammes = mots + [" ".join(mots[i : i + 3]) for i in range(max(0, len(mots) - 2))]
for gramme in grammes:
empreinte = hashlib.sha256(gramme.encode()).digest()
indice = int.from_bytes(empreinte[:4], "big") % DIMENSIONS
signe = 1.0 if empreinte[4] % 2 == 0 else -1.0
vecteur[indice] += signe
norme = math.sqrt(sum(v * v for v in vecteur)) or 1.0
return [v / norme for v in vecteur]
class EmbeddeurLocal:
"""fastembed (ONNX, CPU) — chargé paresseusement, jamais importé en test."""
def __init__(self) -> None:
from fastembed import TextEmbedding # import différé (dépendance optionnelle)
self._modele = TextEmbedding(model_name=MODELE_LOCAL)
def encoder(self, textes: list[str]) -> list[list[float]]:
return [vecteur.tolist() for vecteur in self._modele.embed(textes)]
def construire_embeddeur(mode: str) -> Embeddeur:
if mode == "deterministe":
return EmbeddeurDeterministe()
return EmbeddeurLocal()

View File

@@ -0,0 +1,175 @@
"""Ingestion du corpus (D3) : bibliothèque R3 (PDF, MinIO) + bilans codés.
Chaque texte passe par l'anonymisation (D4) AVANT découpage et embeddings.
Les chunks vivent dans `RagChunk` (pgvector, schéma possédé par Prisma).
"""
import io
import json
from dataclasses import dataclass
import asyncpg
from minio import Minio
from pypdf import PdfReader
from .anonymisation import anonymiser
from .config import Reglages
from .decoupage import decouper
from .embeddings import Embeddeur
@dataclass
class ResultatIngestion:
documents_indexes: int
documents_ignores: int
bilans_indexes: int
extraits: int
def _vecteur_sql(vecteur: list[float]) -> str:
return "[" + ",".join(f"{v:.6f}" for v in vecteur) + "]"
async def noms_a_anonymiser(cnx: asyncpg.Connection) -> list[str]:
"""Toutes les identités connues de la base (D4) : utilisateurs, gardiens,
contacts tiers, demandeurs du portail."""
lignes = await cnx.fetch(
'''
SELECT "displayName" AS nom FROM "User"
UNION SELECT "guardianName" FROM "Location" WHERE "guardianName" IS NOT NULL
UNION SELECT "contactName" FROM "Partner" WHERE "contactName" IS NOT NULL
'''
)
return [ligne["nom"] for ligne in lignes if ligne["nom"]]
def extraire_texte_pdf(octets: bytes) -> list[tuple[str, str]]:
"""[(texte, localisation)] par page — la citation doit pointer la page."""
lecteur = PdfReader(io.BytesIO(octets))
pages: list[tuple[str, str]] = []
for numero, page in enumerate(lecteur.pages, start=1):
texte = page.extract_text() or ""
if texte.strip():
pages.append((texte, f"p. {numero}"))
return pages
async def indexer_documents(
cnx: asyncpg.Connection,
reglages: Reglages,
embeddeur: Embeddeur,
noms: list[str],
) -> tuple[int, int, int]:
minio = Minio(
f"{reglages.minio_endpoint}:{reglages.minio_port}",
access_key=reglages.minio_access_key,
secret_key=reglages.minio_secret_key,
secure=reglages.minio_use_ssl,
)
documents = await cnx.fetch(
'SELECT id, "fileName", "storageKey", "contentType", "inCorpus" FROM "Document"'
)
indexes, ignores, total_extraits = 0, 0, 0
for doc in documents:
# réindexation idempotente : on repart de zéro pour ce document
await cnx.execute('DELETE FROM "RagChunk" WHERE "documentId" = $1', doc["id"])
if not doc["inCorpus"] or doc["contentType"] != "application/pdf":
await cnx.execute(
'UPDATE "Document" SET "indexedAt" = NULL, "chunkCount" = 0 WHERE id = $1',
doc["id"],
)
ignores += 1
continue
reponse = minio.get_object(reglages.minio_bucket, doc["storageKey"])
try:
octets = reponse.read()
finally:
reponse.close()
reponse.release_conn()
extraits: list[tuple[str, str]] = []
for texte_page, localisation in extraire_texte_pdf(octets):
texte_sur = anonymiser(texte_page, noms)
extraits += [(morceau, localisation) for morceau in decouper(texte_sur)]
if extraits:
vecteurs = embeddeur.encoder([contenu for contenu, _ in extraits])
await cnx.executemany(
'''
INSERT INTO "RagChunk"
(id, "sourceType", "documentId", locator, content, embedding)
VALUES (gen_random_uuid(), 'DOCUMENT', $1, $2, $3, $4::vector)
''',
[
(doc["id"], localisation, contenu, _vecteur_sql(vecteur))
for (contenu, localisation), vecteur in zip(extraits, vecteurs)
],
)
await cnx.execute(
'UPDATE "Document" SET "indexedAt" = now(), "chunkCount" = $2 WHERE id = $1',
doc["id"],
len(extraits),
)
indexes += 1
total_extraits += len(extraits)
return indexes, ignores, total_extraits
async def indexer_bilans(
cnx: asyncpg.Connection, embeddeur: Embeddeur, noms: list[str]
) -> tuple[int, int]:
"""Les bilans codés clôturés — « sur votre parc, ce réglage a déjà… »."""
await cnx.execute('DELETE FROM "RagChunk" WHERE "sourceType" = \'WORK_ORDER\'')
bilans = await cnx.fetch(
'''
SELECT wo.id, wo.reference, wo.title, wo."completedAt",
a.reference AS asset_ref, a.brand, a.model,
(SELECT json_object_agg(rv.field, rv.label)
FROM "InterventionReport" ir2
JOIN "ReferenceValue" rv ON rv.id IN (
ir2."doorStateId", ir2."cabinPositionId", ir2."anomalyId",
ir2."externalCauseId", ir2."actionTakenId", ir2."componentConcernedId")
WHERE ir2."workOrderId" = wo.id) AS bilan
FROM "WorkOrder" wo
JOIN "InterventionReport" ir ON ir."workOrderId" = wo.id
JOIN "Asset" a ON a.id = wo."assetId"
WHERE wo.status = 'DONE'
'''
)
lignes = []
for bilan in bilans:
champs = json.loads(bilan["bilan"]) if bilan["bilan"] else {}
codes = " ; ".join(f"{champ} : {label}" for champ, label in champs.items())
contenu = anonymiser(
f"Intervention {bilan['reference']}{bilan['title']}. "
f"Appareil {bilan['asset_ref']} ({bilan['brand']} {bilan['model'] or ''}). "
f"Bilan codé : {codes}.",
noms,
)
quand = bilan["completedAt"].date().isoformat() if bilan["completedAt"] else "date inconnue"
lignes.append((bilan["id"], f"bilan du {quand}", contenu))
if lignes:
vecteurs = embeddeur.encoder([contenu for _, _, contenu in lignes])
await cnx.executemany(
'''
INSERT INTO "RagChunk"
(id, "sourceType", "workOrderId", locator, content, embedding)
VALUES (gen_random_uuid(), 'WORK_ORDER', $1, $2, $3, $4::vector)
''',
[
(wo_id, localisation, contenu, _vecteur_sql(vecteur))
for (wo_id, localisation, contenu), vecteur in zip(lignes, vecteurs)
],
)
return len(lignes), len(lignes)
async def reindexer_tout(
cnx: asyncpg.Connection, reglages: Reglages, embeddeur: Embeddeur
) -> ResultatIngestion:
noms = await noms_a_anonymiser(cnx)
docs_ok, docs_non, extraits_docs = await indexer_documents(cnx, reglages, embeddeur, noms)
bilans, extraits_bilans = await indexer_bilans(cnx, embeddeur, noms)
return ResultatIngestion(
documents_indexes=docs_ok,
documents_ignores=docs_non,
bilans_indexes=bilans,
extraits=extraits_docs + extraits_bilans,
)

View File

@@ -0,0 +1,57 @@
"""Recherche sémantique dans le corpus (pgvector, distance cosinus).
Ne renvoie QUE des extraits sourcés — la brique de « sourcé ou silencieux ».
"""
from dataclasses import dataclass
import asyncpg
from .embeddings import Embeddeur
from .ingestion import _vecteur_sql
@dataclass
class ExtraitTrouve:
source_type: str
document_id: str | None
work_order_id: str | None
titre: str # nom de fichier ou référence d'OT
locator: str
content: str
score: float # similarité cosinus (0..1)
async def chercher(
cnx: asyncpg.Connection,
embeddeur: Embeddeur,
question: str,
limite: int = 5,
) -> list[ExtraitTrouve]:
vecteur = _vecteur_sql(embeddeur.encoder([question])[0])
lignes = await cnx.fetch(
'''
SELECT c."sourceType", c."documentId", c."workOrderId", c.locator, c.content,
1 - (c.embedding <=> $1::vector) AS score,
COALESCE(d."fileName", wo.reference, '?') AS titre
FROM "RagChunk" c
LEFT JOIN "Document" d ON d.id = c."documentId"
LEFT JOIN "WorkOrder" wo ON wo.id = c."workOrderId"
WHERE c."documentId" IS NULL OR d."inCorpus" -- l'exclusion D3 s'applique aussi à la lecture
ORDER BY c.embedding <=> $1::vector
LIMIT $2
''',
vecteur,
limite,
)
return [
ExtraitTrouve(
source_type=ligne["sourceType"],
document_id=str(ligne["documentId"]) if ligne["documentId"] else None,
work_order_id=str(ligne["workOrderId"]) if ligne["workOrderId"] else None,
titre=ligne["titre"],
locator=ligne["locator"],
content=ligne["content"],
score=float(ligne["score"]),
)
for ligne in lignes
]