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