Files
opencare/pipeline.py
laurentbarontini efd775706b Initial production-ready release of OpenCare v1
- Flask/SQLAlchemy app with profile, dashboard, recommendations
- AI crawl pipeline (GPT-4o) with admin review workflow
- 14 curated health sources (RSS + index crawling)
- Production config: env vars, Gunicorn, systemd, nginx
- deploy/ scripts for VPS setup and updates

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-05 21:38:03 +02:00

408 lines
18 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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 ───────────────────────────────────────────────────────────────
def _fetch_webpage(url: str) -> tuple:
headers = {'User-Agent': 'Mozilla/5.0 (compatible; OpenCare-Bot/1.0)'}
r = requests.get(url, headers=headers, timeout=20)
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',
)
headers = {'User-Agent': 'Mozilla/5.0 (compatible; OpenCare-Bot/1.0)'}
r = requests.get(url, headers=headers, timeout=20)
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