421 lines
19 KiB
Python
421 lines
19 KiB
Python
"""
|
||
Pipeline de crawl — fonctionne dans un thread background sans contexte Flask.
|
||
Reçoit un sqlalchemy.engine.Engine + la config, crée sa propre Session native.
|
||
"""
|
||
import json
|
||
import logging
|
||
import requests
|
||
from bs4 import BeautifulSoup
|
||
from datetime import datetime
|
||
from sqlalchemy.orm import Session as SASession
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
_CONDITIONS = [
|
||
'hypertension', 'diabetes_risk', 'cardiovascular_risk', 'anxiety',
|
||
'digestive_issues', 'back_pain', 'osteoporosis_risk', 'metabolic_syndrome',
|
||
'thyroid', 'sleep_apnea',
|
||
]
|
||
_HEALTH_GOALS = [
|
||
'weight_loss', 'muscle_gain', 'longevity', 'energy', 'better_sleep',
|
||
'stress_reduction', 'metabolic_health', 'cardiovascular_health', 'cognitive_health',
|
||
]
|
||
_DIET_TYPES = ['omnivore', 'flexitarian', 'vegetarian', 'vegan', 'paleo', 'keto', 'other']
|
||
_ACTIVITY_LEVELS = ['sedentary', 'light', 'moderate', 'intense', 'athlete']
|
||
|
||
SYSTEM_PROMPT = f"""Tu es un expert en médecine factuelle (evidence-based medicine).
|
||
Analyse le contenu santé fourni et extrait les recommandations médicales/santé qu'il contient.
|
||
|
||
Retourne UNIQUEMENT du JSON valide.
|
||
Si le contenu ne contient pas de recommandation santé pertinente, retourne: {{"relevant": false}}
|
||
|
||
Format attendu:
|
||
{{
|
||
"relevant": true,
|
||
"recommendations": [
|
||
{{
|
||
"domain": "nutrition|fasting|exercise|sleep|stress|supplements|mental_health|other",
|
||
"title": "Titre court (< 60 caractères)",
|
||
"summary": "1-2 phrases résumant l'action recommandée",
|
||
"details": "2-4 phrases avec mécanisme d'action et données quantitatives si disponibles",
|
||
"evidence_level": "A (méta-analyse/RCT solide) | B (observationnel/RCT limité) | C (expert/anecdote)",
|
||
"study_type": "meta-analysis|RCT|observational|expert_opinion|anecdote",
|
||
"intervention": "Action concrète à mettre en oeuvre",
|
||
"effect": "Résultat attendu avec chiffres si disponibles",
|
||
"source_title": "Titre de l'étude ou article source",
|
||
"source_doi": "DOI si disponible sinon null",
|
||
"source_url": "URL de la source",
|
||
"tags": ["tag1", "tag2", "tag3"],
|
||
"contraindications": ["contre-indication 1"],
|
||
"target_age_min": null,
|
||
"target_age_max": null,
|
||
"target_sex": "all|male|female",
|
||
"required_conditions": [],
|
||
"excluded_conditions": [],
|
||
"required_health_goals": [],
|
||
"required_diet_types": [],
|
||
"required_activity_levels": [],
|
||
"not_for_under_18": false,
|
||
"not_for_pregnant": false,
|
||
"priority": 5
|
||
}}
|
||
]
|
||
}}
|
||
|
||
Valeurs acceptées pour required_conditions / excluded_conditions: {_CONDITIONS}
|
||
Valeurs acceptées pour required_health_goals: {_HEALTH_GOALS}
|
||
Valeurs acceptées pour required_diet_types: {_DIET_TYPES}
|
||
Valeurs acceptées pour required_activity_levels: {_ACTIVITY_LEVELS}
|
||
|
||
Pour required_conditions: conditions qui DÉCLENCHENT cette rec.
|
||
Pour excluded_conditions: conditions INCOMPATIBLES avec cette rec.
|
||
Si la recommandation est universelle, laisser toutes les listes vides.
|
||
|
||
RÈGLE ABSOLUE — Citation scientifique:
|
||
- Identifier TOUJOURS l'étude primaire sous-jacente (PubMed, Cochrane, NEJM, Lancet, BMJ…)
|
||
- source_doi: DOI de l'étude originale (ex: 10.1056/NEJMoa1200303), PAS l'URL de l'article source
|
||
- source_url: URL PubMed ou DOI de la source primaire (https://doi.org/...), pas le site référant
|
||
- evidence_level A UNIQUEMENT si méta-analyse ou RCT avec ≥ 100 participants
|
||
- Ne jamais inventer de référence: si incertain, mettre source_doi null et evidence_level C
|
||
- study_type doit refléter le niveau de preuve réel de l'étude citée
|
||
"""
|
||
|
||
|
||
# ── DB helpers (session-aware, no Flask context) ───────────────────────────────
|
||
|
||
def _log(session: SASession, run, msg: str) -> None:
|
||
ts = datetime.utcnow().strftime('%H:%M:%S')
|
||
lines = run.log_lines
|
||
lines.append(f"[{ts}] {msg}")
|
||
run._log_lines = json.dumps(lines[-150:])
|
||
session.add(run)
|
||
session.commit()
|
||
logger.info(msg)
|
||
|
||
|
||
# ── Web fetchers ───────────────────────────────────────────────────────────────
|
||
|
||
_HEADERS = {
|
||
'User-Agent': (
|
||
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) '
|
||
'AppleWebKit/537.36 (KHTML, like Gecko) '
|
||
'Chrome/124.0.0.0 Safari/537.36'
|
||
),
|
||
'Accept-Language': 'en-US,en;q=0.9',
|
||
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
|
||
}
|
||
|
||
|
||
def _fetch_webpage(url: str) -> tuple:
|
||
r = requests.get(url, headers=_HEADERS, timeout=20)
|
||
if r.status_code == 429:
|
||
raise requests.HTTPError(f'429 Too Many Requests – {url}', response=r)
|
||
r.raise_for_status()
|
||
soup = BeautifulSoup(r.content, 'html.parser')
|
||
for tag in soup(['script', 'style', 'nav', 'footer', 'header', 'aside', 'iframe', 'form']):
|
||
tag.decompose()
|
||
title_tag = soup.find('title')
|
||
title = title_tag.get_text(strip=True) if title_tag else url
|
||
main = (
|
||
soup.find('article') or soup.find('main') or
|
||
soup.find(id='content') or soup.find(class_='content') or
|
||
soup.find('body')
|
||
)
|
||
text = main.get_text(separator='\n', strip=True) if main else soup.get_text()
|
||
return title, text[:6000]
|
||
|
||
|
||
def _fetch_rss(url: str) -> list:
|
||
import feedparser
|
||
feed = feedparser.parse(url)
|
||
results = []
|
||
for entry in feed.entries[:30]:
|
||
link = entry.get('link', '')
|
||
title = entry.get('title', '')
|
||
content = entry.get('summary') or entry.get('description') or ''
|
||
if content:
|
||
soup = BeautifulSoup(content, 'html.parser')
|
||
content = soup.get_text(strip=True)
|
||
results.append((link, title, content[:4000]))
|
||
return results
|
||
|
||
|
||
def _fetch_index_page(url: str) -> list:
|
||
"""Crawl an index/listing page and return (article_url, title) pairs found on it."""
|
||
from urllib.parse import urlparse, urljoin
|
||
_SKIP = (
|
||
'/tag/', '/category/', '/author/', '/page/', '/feed', '/wp-content',
|
||
'login', 'search', 'contact', 'about', 'privacy', 'terms',
|
||
'twitter.com', 'facebook.com', 'instagram.com', 'youtube.com',
|
||
)
|
||
r = requests.get(url, headers=_HEADERS, timeout=20)
|
||
if r.status_code == 429:
|
||
raise requests.HTTPError(f'429 Too Many Requests – {url}', response=r)
|
||
r.raise_for_status()
|
||
soup = BeautifulSoup(r.content, 'html.parser')
|
||
base = urlparse(url)
|
||
seen = set()
|
||
links = []
|
||
for a in soup.find_all('a', href=True):
|
||
href = a['href'].strip()
|
||
if not href or href.startswith(('#', 'mailto:', 'tel:', 'javascript:')):
|
||
continue
|
||
full = urljoin(url, href)
|
||
p = urlparse(full)
|
||
if p.netloc != base.netloc:
|
||
continue
|
||
path = p.path.rstrip('/')
|
||
if not path or path == base.path.rstrip('/'):
|
||
continue
|
||
if full in seen:
|
||
continue
|
||
if any(s in full.lower() for s in _SKIP):
|
||
continue
|
||
title = a.get_text(strip=True)
|
||
if len(title) < 10 or len(title) > 200:
|
||
continue
|
||
seen.add(full)
|
||
links.append((full, title))
|
||
return links[:60]
|
||
|
||
|
||
# ── OpenAI analysis ────────────────────────────────────────────────────────────
|
||
|
||
def _analyze(text: str, title: str, api_key: str, model: str = 'gpt-4o') -> list:
|
||
from openai import OpenAI
|
||
client = OpenAI(api_key=api_key)
|
||
user_msg = f"Titre: {title}\n\nContenu:\n{text}"
|
||
resp = client.chat.completions.create(
|
||
model=model,
|
||
response_format={"type": "json_object"},
|
||
messages=[
|
||
{"role": "system", "content": SYSTEM_PROMPT},
|
||
{"role": "user", "content": user_msg},
|
||
],
|
||
temperature=0.2,
|
||
)
|
||
data = json.loads(resp.choices[0].message.content)
|
||
if not data.get('relevant', False):
|
||
return []
|
||
return data.get('recommendations', [])
|
||
|
||
|
||
def _store_recs(session: SASession, pending_model, article_id: int,
|
||
recs: list, fallback_url: str) -> int:
|
||
count = 0
|
||
for rec in recs:
|
||
p = pending_model(
|
||
article_id=article_id,
|
||
domain=rec.get('domain', 'other'),
|
||
title=(rec.get('title') or '')[:256],
|
||
summary=rec.get('summary', ''),
|
||
details=rec.get('details', ''),
|
||
evidence_level=rec.get('evidence_level', 'C'),
|
||
study_type=rec.get('study_type', ''),
|
||
intervention=(rec.get('intervention') or '')[:512],
|
||
effect=(rec.get('effect') or '')[:512],
|
||
source_title=(rec.get('source_title') or '')[:512],
|
||
source_doi=(rec.get('source_doi') or '')[:256],
|
||
source_url=(rec.get('source_url') or fallback_url or '')[:512],
|
||
target_age_min=rec.get('target_age_min'),
|
||
target_age_max=rec.get('target_age_max'),
|
||
target_sex=rec.get('target_sex', 'all'),
|
||
not_for_under_18=bool(rec.get('not_for_under_18', False)),
|
||
not_for_pregnant=bool(rec.get('not_for_pregnant', False)),
|
||
priority=int(rec.get('priority', 5)),
|
||
)
|
||
p.tags = rec.get('tags') or []
|
||
p.contraindications = rec.get('contraindications') or []
|
||
p.required_conditions = rec.get('required_conditions') or []
|
||
p.excluded_conditions = rec.get('excluded_conditions') or []
|
||
p.required_health_goals = rec.get('required_health_goals') or []
|
||
p.required_diet_types = rec.get('required_diet_types') or []
|
||
p.required_activity_levels = rec.get('required_activity_levels') or []
|
||
session.add(p)
|
||
count += 1
|
||
return count
|
||
|
||
|
||
# ── Main entry point ───────────────────────────────────────────────────────────
|
||
|
||
def run_pipeline(engine, source_id: int = None,
|
||
api_key: str = '', model: str = 'gpt-4o', max_articles: int = 20):
|
||
"""
|
||
Runs entirely in a background thread with no Flask context.
|
||
Uses a plain SQLAlchemy session bound to the passed engine.
|
||
"""
|
||
# Import models lazily (avoiding circular imports at module load time)
|
||
from app import CrawlRun, CrawlSource, RawArticle, PendingRecommendation
|
||
|
||
with SASession(engine) as session:
|
||
# Create run record — immediately visible to the polling endpoint
|
||
run = CrawlRun()
|
||
session.add(run)
|
||
session.commit()
|
||
|
||
try:
|
||
if not api_key:
|
||
run.status = 'error'
|
||
_log(session, run, '❌ Clé API OpenAI manquante — configurez-la dans Administration > Configuration')
|
||
return 0
|
||
|
||
query = session.query(CrawlSource).filter_by(active=True)
|
||
if source_id:
|
||
query = query.filter_by(id=source_id)
|
||
sources = query.all()
|
||
|
||
if not sources:
|
||
run.status = 'done'
|
||
run.ended_at = datetime.utcnow()
|
||
_log(session, run, 'ℹ Aucune source active à crawler')
|
||
return 0
|
||
|
||
_log(session, run, f'🚀 Démarrage — {len(sources)} source(s), modèle: {model}')
|
||
|
||
total_articles = 0
|
||
total_recs = 0
|
||
|
||
for source in sources:
|
||
if total_articles >= max_articles:
|
||
_log(session, run, f'⏹ Limite de {max_articles} articles atteinte — arrêt')
|
||
break
|
||
|
||
_log(session, run, f'🔍 {source.name} — {source.url[:70]}…')
|
||
|
||
try:
|
||
if source.source_type == 'rss':
|
||
items = _fetch_rss(source.url)
|
||
_log(session, run, f' → {len(items)} entrée(s) dans le flux RSS')
|
||
src_articles = 0
|
||
src_recs = 0
|
||
|
||
for item_url, item_title, item_content in items:
|
||
if total_articles >= max_articles:
|
||
break
|
||
if not item_url:
|
||
continue
|
||
if session.query(RawArticle).filter_by(url=item_url).first():
|
||
continue
|
||
|
||
article = RawArticle(
|
||
source_id=source.id, url=item_url,
|
||
title=item_title, content=item_content,
|
||
)
|
||
session.add(article)
|
||
session.flush()
|
||
|
||
try:
|
||
_log(session, run, f' 🤖 Analyse : « {(item_title or item_url)[:55]} »')
|
||
recs = _analyze(item_content, item_title, api_key, model)
|
||
n = _store_recs(session, PendingRecommendation, article.id, recs, item_url)
|
||
session.commit()
|
||
article.analyzed = True
|
||
src_articles += 1
|
||
src_recs += n
|
||
total_articles += 1
|
||
total_recs += n
|
||
if recs:
|
||
_log(session, run, f' ✅ {n} recommandation(s) extraite(s)')
|
||
else:
|
||
_log(session, run, f' — Aucune recommandation pertinente')
|
||
except Exception as e:
|
||
_log(session, run, f' ⚠ Erreur analyse IA : {str(e)[:80]}')
|
||
session.rollback()
|
||
|
||
_log(session, run, f'✓ {source.name} — {src_articles} article(s), {src_recs} rec(s)')
|
||
|
||
elif source.source_type == 'index':
|
||
links = _fetch_index_page(source.url)
|
||
_log(session, run, f' → {len(links)} lien(s) repérés sur la page index')
|
||
src_articles = 0
|
||
src_recs = 0
|
||
|
||
for link_url, link_title in links:
|
||
if total_articles >= max_articles:
|
||
break
|
||
if session.query(RawArticle).filter_by(url=link_url).first():
|
||
continue
|
||
try:
|
||
title, text_content = _fetch_webpage(link_url)
|
||
title = title or link_title
|
||
article = RawArticle(
|
||
source_id=source.id, url=link_url,
|
||
title=title, content=text_content,
|
||
)
|
||
session.add(article)
|
||
session.flush()
|
||
_log(session, run, f' 🤖 Analyse : « {(title or link_url)[:55]} »')
|
||
recs = _analyze(text_content, title, api_key, model)
|
||
n = _store_recs(session, PendingRecommendation, article.id, recs, link_url)
|
||
session.commit()
|
||
article.analyzed = True
|
||
src_articles += 1
|
||
src_recs += n
|
||
total_articles += 1
|
||
total_recs += n
|
||
if recs:
|
||
_log(session, run, f' ✅ {n} recommandation(s) extraite(s)')
|
||
else:
|
||
_log(session, run, f' — Aucune recommandation pertinente')
|
||
except Exception as e:
|
||
_log(session, run, f' ⚠ {link_url[:50]}: {str(e)[:60]}')
|
||
session.rollback()
|
||
|
||
_log(session, run, f'✓ {source.name} — {src_articles} article(s), {src_recs} rec(s)')
|
||
|
||
else: # webpage
|
||
if session.query(RawArticle).filter_by(url=source.url).first():
|
||
_log(session, run, f' → Déjà crawlé (ignoré)')
|
||
else:
|
||
title, text = _fetch_webpage(source.url)
|
||
article = RawArticle(
|
||
source_id=source.id, url=source.url,
|
||
title=title, content=text,
|
||
)
|
||
session.add(article)
|
||
session.flush()
|
||
|
||
_log(session, run, f' 🤖 Analyse : « {title[:60]} »')
|
||
try:
|
||
recs = _analyze(text, title, api_key, model)
|
||
n = _store_recs(session, PendingRecommendation, article.id, recs, source.url)
|
||
session.commit()
|
||
article.analyzed = True
|
||
total_articles += 1
|
||
total_recs += n
|
||
if recs:
|
||
_log(session, run, f'✓ {source.name} — {n} recommandation(s) extraite(s)')
|
||
else:
|
||
_log(session, run, f'✓ {source.name} — Aucune recommandation pertinente')
|
||
except Exception as e:
|
||
_log(session, run, f'⚠ {source.name} — Erreur IA : {str(e)[:80]}')
|
||
session.rollback()
|
||
|
||
source.last_crawled_at = datetime.utcnow()
|
||
session.commit()
|
||
|
||
except requests.exceptions.Timeout:
|
||
_log(session, run, f'⚠ {source.name} — Timeout réseau (>20s)')
|
||
session.rollback()
|
||
except requests.exceptions.ConnectionError as e:
|
||
_log(session, run, f'⚠ {source.name} — Erreur connexion : {str(e)[:80]}')
|
||
session.rollback()
|
||
except Exception as e:
|
||
_log(session, run, f'⚠ {source.name} — Erreur : {str(e)[:80]}')
|
||
session.rollback()
|
||
|
||
run.status = 'done'
|
||
run.ended_at = datetime.utcnow()
|
||
run.articles_processed = total_articles
|
||
run.recs_generated = total_recs
|
||
_log(session, run,
|
||
f"🏁 Terminé — {total_articles} article(s), {total_recs} rec(s) en file d'attente")
|
||
return total_articles
|
||
|
||
except Exception as e:
|
||
run.status = 'error'
|
||
run.ended_at = datetime.utcnow()
|
||
_log(session, run, f'💥 Erreur fatale : {str(e)[:120]}')
|
||
logger.exception("Pipeline fatal error")
|
||
return 0
|