# QC Élection Forecast — Plateforme de prévision électorale du Québec 2026 # Auteur : Simon-Pierre Boucher # Contact : contact@spboucher.ai # https://www.qc-election.com """Pipeline automatisé : collect → validate → normalize → dedupe → model → simulate → publish. Chaque étape est isolée : l'échec d'une source n'arrête pas le reste (journalisé dans pipeline_logs).""" from __future__ import annotations import logging import traceback from datetime import date, datetime, timezone import numpy as np from sqlalchemy.orm import Session from .config import settings from .db import SessionLocal, init_db from . import models as Mo from .modeling import ensemble as ENS from .modeling import fundamentals as FUND from .modeling import house_effects as HE from .modeling import simulate as SIM from .modeling.forecast import distribution_from, forecast as fc, nowcast as nc from .modeling.trend import fit_trend log = logging.getLogger("pipeline") def _log(db: Session, step: str, status: str, message: str = "") -> None: db.add(Mo.PipelineLog(step=step, status=status, message=message[:2000])) db.commit() def load_polls(db: Session, election_id: int, as_of: date | None = None, region: str = "QC") -> list[dict]: q = (db.query(Mo.Poll).filter(Mo.Poll.election_id == election_id, Mo.Poll.excluded.is_(False), Mo.Poll.region == region)) if as_of: q = q.filter(Mo.Poll.field_end <= as_of) out = [] for poll in q.all(): shares = {r.party: r.normalized_value for r in poll.results} out.append({"id": poll.id, "pollster": poll.pollster.name, "field_end": poll.field_end, "sample_size": poll.sample_size, "mode": poll.mode, "shares": shares}) return out def compute_pollster_profiles(db: Session, exclude_after: date | None = None): """House effects sur les élections passées à résultat connu.""" elections = [{"date": e.election_date, "result": e.actual_result} for e in db.query(Mo.Election).filter(Mo.Election.actual_result.isnot(None)) if not exclude_after or e.election_date <= exclude_after] polls = [] for e in db.query(Mo.Election).all(): polls.extend(load_polls(db, e.id)) profiles = HE.compute_profiles(polls, elections) # publication des cotes for name, prof in profiles.items(): pol = db.query(Mo.Pollster).filter_by(name=name).first() if not pol: continue rating = (db.query(Mo.PollsterRating) .filter_by(pollster_id=pol.id, computed_for="2026-general").first()) if rating is None: rating = Mo.PollsterRating(pollster_id=pol.id, computed_for="2026-general") db.add(rating) rating.n_polls = prof.n_polls rating.mae_pp = prof.mae_pp rating.house_effects = prof.house_effects rating.weight_multiplier = prof.weight_multiplier rating.detail = prof.detail rating.computed_at = datetime.now(timezone.utc) db.commit() return profiles def _district_inputs(db: Session, baseline_national: dict[str, float]) -> SIM.SimulationInput: districts = db.query(Mo.District).order_by(Mo.District.name).all() names, baselines, regions, retire, inc_idx = [], [], [], [], [] for d in districts: names.append(d.name) baselines.append([d.baseline_shares.get(p, 0.0) / 100.0 for p in settings.parties]) regions.append(d.region) retire.append(0 if d.incumbent_running else 1) inc_idx.append(settings.parties.index(d.incumbent_party) if d.incumbent_party in settings.parties else -1) return SIM.SimulationInput( x_mean=np.zeros(5), P=np.eye(5), baseline_national=np.array([baseline_national[p] / 100.0 for p in settings.parties]), district_names=names, district_baselines=np.array(baselines), district_regions=regions, retirement_flags=np.array(retire), incumbent_party_idx=np.array(inc_idx), majority_seats=settings.majority_seats) def run_forecast(db: Session, as_of: date | None = None, label: str | None = None, is_backtest: bool = False, n_sims: int | None = None, save: bool = True) -> Mo.ForecastRun | None: """Ajuste le modèle v2 (sondages + partielles + fondamentaux + médias), simule l'élection cible et persiste un ForecastRun avec la décomposition complète de la contribution de chaque couche.""" as_of = as_of or date.today() target = db.query(Mo.Election).filter_by(is_target=True).one() days_left = (target.election_date - as_of).days profiles = compute_pollster_profiles(db) polls = load_polls(db, target.id, as_of=as_of) # --- couche 2 : votes réels des partielles → observations nationales --- from .ingest.byelections import pseudo_polls bye_polls = pseudo_polls(db, as_of=as_of) trend_a = fit_trend(polls, profiles, as_of) # sondages seuls trend = fit_trend(polls + bye_polls, profiles, as_of) if bye_polls else trend_a if trend is None: raise RuntimeError("Pas assez de sondages pour ajuster le modèle") # --- volatilité de campagne : attention Wikipédia + pouls social (bornée, # variance seulement — v3 §16) --- attention: dict = {"available": False, "drift_multiplier": 1.0} volatility: dict = {"multiplier": 1.0} if not is_backtest: try: from .ingest.wiki_attention import signal as att_signal attention = att_signal(db) except Exception: pass try: from .modeling.signals import volatility as VOLA volatility = VOLA.compute(db, attention) except Exception: volatility = {"multiplier": float(attention.get("drift_multiplier", 1.0))} drift_mult = float(volatility.get("multiplier", 1.0)) # σ_industrie : appris LOEO sur 2007-2022 si le flag l'active (v3.1) industry_sd = None if settings.use_empirical_industry_error: try: from .modeling.national.pollster_error import industry_sigma_loeo industry_sd = industry_sigma_loeo(None) except Exception: industry_sd = None now_d = nc(trend, as_of) fc_b = fc(trend, as_of, target.election_date, drift_multiplier=drift_mult, industry_sd=industry_sd) # sondages + partielles # --- couche 3 : prior de fondamentaux (poids décroissant vers le scrutin) --- result_2022 = (db.query(Mo.Election) .filter(Mo.Election.actual_result.isnot(None)) .order_by(Mo.Election.election_date.desc()).first().actual_result) fund_diag: dict = {"enabled": settings.fundamentals_enabled} x_c, P_c = fc_b.x, fc_b.P if settings.fundamentals_enabled: from .ingest.firecrawl_watch import latest_satisfaction sat, sat_prov = latest_satisfaction(db) prior = FUND.compute_prior(as_of, incumbent="CAQ", terms=2, satisfaction=sat, prev_result=result_2022) x_c, P_c, bl = FUND.blend(fc_b.x, fc_b.P, prior, days_left) fund_diag.update(prior=prior.detail, blend=bl, satisfaction={"value": sat, **sat_prov}) # --- couche 4 : ajustement médias borné --- nudge = ENS.media_nudge(db) x_d = ENS.apply_nudge(x_c, nudge.get("delta_pp", {})) fc_d = distribution_from("forecast", as_of, x_d, P_c) inp = _district_inputs(db, result_2022) inp.x_mean, inp.P = fc_d.x, fc_d.P sim = SIM.run_simulation(inp, n_sims=n_sims) # --- décomposition : contribution de chaque couche (votes et probabilités) --- decomposition = None if not is_backtest and trend_a is not None: def _probe(x, P, seed): inp.x_mean, inp.P = x, P s = SIM.run_simulation(inp, n_sims=8000, rng=np.random.default_rng(seed)) return {p: s.seats[p]["prob_most"] for p in settings.parties} fc_a = fc(trend_a, as_of, target.election_date, drift_multiplier=drift_mult, industry_sd=industry_sd) variants = [ ("sondages", "Sondages seuls (modèle v1)", fc_a.x, fc_a.P), ("partielles", "+ votes réels des partielles", fc_b.x, fc_b.P), ("fondamentaux", "+ prior de fondamentaux", x_c, P_c), ("medias", "+ ajustement médias", x_d, P_c), ] layers, prev_vote, prev_prob = [], None, None for key, lab, x, P in variants: d = distribution_from("probe", as_of, x, P) vote = {p: d.summary[p]["mean"] for p in settings.parties} prob = _probe(x, P, seed=hash(key) % 2 ** 31) layers.append({ "key": key, "label": lab, "vote": vote, "prob_most": prob, "delta_vote": ({p: round(vote[p] - prev_vote[p], 2) for p in settings.parties} if prev_vote else None), "delta_prob": ({p: round(prob[p] - prev_prob[p], 4) for p in settings.parties} if prev_prob else None)}) prev_vote, prev_prob = vote, prob inp.x_mean, inp.P = fc_d.x, fc_d.P # restaure l'entrée finale decomposition = {"layers": layers, "n_sims_probe": 8000} # --- couche 5 : ensemble avec les marchés prédictifs (probabilités) --- market_ens = None if not is_backtest: try: from .ingest.markets import fetch_winner_market mk = fetch_winner_market() market_ens = ENS.market_blend( {p: sim.seats[p]["prob_most"] for p in settings.parties}, (mk or {}).get("implied_prob_most_seats")) if mk: market_ens["market_meta"] = {"event": mk.get("event"), "url": mk.get("url"), "volume_usd": mk.get("volume_usd"), "fetched_at": mk.get("fetched_at")} except Exception: market_ens = None # série lissée compacte pour les graphiques (sous-échantillonnée) step = max(1, len(trend.dates) // 400) idx = list(range(0, len(trend.dates), step)) if idx[-1] != len(trend.dates) - 1: idx.append(len(trend.dates) - 1) series = { "dates": [trend.dates[i].isoformat() for i in idx], "mean": {p: [round(float(trend.share_mean[i, k] * 100), 2) for i in idx] for k, p in enumerate(settings.parties)}, "lo": {p: [round(float(trend.share_lo[i, k] * 100), 2) for i in idx] for k, p in enumerate(settings.parties)}, "hi": {p: [round(float(trend.share_hi[i, k] * 100), 2) for i in idx] for k, p in enumerate(settings.parties)}, } # --- circonscriptions pivots (battlegrounds, v3 §24) --- battlegrounds = None if not is_backtest and settings.enable_tipping_points: try: inp.x_mean, inp.P = fc_d.x, fc_d.P sim_t = SIM.run_simulation(inp, n_sims=settings.tipping_n_sims, rng=np.random.default_rng(777), keep_raw=True) battlegrounds = SIM.tipping_points( sim_t.raw_winners, sim_t.raw_shares, inp.district_names, settings.majority_seats, settings.parties) except Exception: battlegrounds = None # --- distribution prédictive du prochain sondage + surprises (§46-48) --- next_poll, poll_surprises = None, [] if not is_backtest: try: from .modeling.national.next_poll import predictive, surprise next_poll = predictive(now_d.x, now_d.P) # surprise des sondages arrivés depuis le run précédent, mesurée # contre la prédiction STOCKÉE par ce run-là (jamais rétro-ajustée) prev = (db.query(Mo.ForecastRun) .filter(Mo.ForecastRun.is_backtest.is_(False)) .order_by(Mo.ForecastRun.run_at.desc()).first()) prev_pred = ((prev.national.get("beyond") or {}).get("next_poll") if prev else None) if prev and prev_pred: fresh = (db.query(Mo.Poll) .filter(Mo.Poll.election_id == target.id, Mo.Poll.accessed_at > prev.run_at, Mo.Poll.excluded.is_(False)).all()) for poll in fresh[:6]: shares = {r.party: r.normalized_value for r in poll.results} s = surprise(prev_pred, shares) poll_surprises.append({"pollster": poll.pollster.name, "field_end": poll.field_end.isoformat(), **s}) except Exception: next_poll = None beyond = { "decomposition": decomposition, "web_attention": {k: v for k, v in attention.items() if k != "sparkline_dates"}, "volatility": volatility, "next_poll": next_poll, "poll_surprises": poll_surprises, "fundamentals": fund_diag, "byelections_used": [{"pollster": b["pollster"], "field_end": b["field_end"].isoformat(), "implied_national": b["shares"]} for b in bye_polls], "media_adjustment": nudge, "market_ensemble": market_ens, } run = Mo.ForecastRun( as_of=as_of, model_version=settings.model_version, n_polls_used=trend.n_polls, n_simulations=sim.n_sims, is_backtest=is_backtest, label=label, national={"nowcast": now_d.summary, "forecast": fc_d.summary, "trend_series": series, "x_forecast": [float(v) for v in fc_d.x], "P_forecast": [[float(v) for v in row] for row in fc_d.P], "baseline_national": result_2022, "beyond": beyond}, seats={"per_party": sim.seats, "summary": sim.seat_matrix_summary, "sim_national_vote": sim.national_vote, "ensemble": market_ens, "battlegrounds": battlegrounds}, diagnostics={"q_daily_var": trend.q, "loglik": trend.loglik, "days_to_election": days_left, "n_byelections": len(bye_polls), "industry_sd": round(industry_sd, 4) if industry_sd else settings.industry_error_sd, "industry_sd_source": ("loeo-2007-2022" if industry_sd else "config-manuel"), "fundamentals_weight": (fund_diag.get("blend") or {}).get( "precision_share", 0.0)}) for d in sim.districts: run.district_results.append(Mo.ForecastDistrictResult( district_name=d["district"], favorite=d["favorite"], category=d["category"], detail=d)) if save: db.add(run) db.commit() return run def run_pipeline(full_refresh: bool = True) -> dict: """Cycle complet. Tolère l'échec individuel des sources.""" init_db() db = SessionLocal() report = {} try: if full_refresh: try: from .ingest.wikipedia import refresh_all report["collect"] = refresh_all() _log(db, "collect", "ok", str(report["collect"])) except Exception as e: # la suite continue avec les données en cache report["collect"] = f"échec: {e}" _log(db, "collect", "error", traceback.format_exc()) try: from .seed import seed_all report["seed"] = seed_all(db) _log(db, "normalize", "ok", str(report["seed"])) except Exception as e: report["seed"] = f"échec: {e}" _log(db, "normalize", "error", traceback.format_exc()) try: from .ingest.dgeq import refresh_dgeq report["dgeq"] = refresh_dgeq(db) _log(db, "dgeq", "ok", str(report["dgeq"])[:500]) except Exception as e: report["dgeq"] = f"échec: {e}" _log(db, "dgeq", "error", traceback.format_exc()) try: from .ingest.news_rss import ingest_feeds report["sentiment"] = ingest_feeds(db) _log(db, "sentiment", "ok", str(report["sentiment"])) except Exception as e: report["sentiment"] = f"échec: {e}" _log(db, "sentiment", "error", traceback.format_exc()) try: from .ingest.byelections import ensure_byelections, verify_with_firecrawl added = ensure_byelections(db) verif = verify_with_firecrawl(db) if settings.firecrawl_api_key else {} report["byelections"] = {"ajoutées": added, "vérification": verif} _log(db, "byelections", "ok", str(report["byelections"])[:500]) except Exception as e: report["byelections"] = f"échec: {e}" _log(db, "byelections", "error", traceback.format_exc()) try: from .ingest.firecrawl_watch import run_watch report["veille"] = run_watch(db) _log(db, "veille-firecrawl", "ok", str(report["veille"])[:500]) except Exception as e: report["veille"] = f"échec: {e}" _log(db, "veille-firecrawl", "error", traceback.format_exc()) try: from .ingest.pollster_reports import run_primary_watch report["rapports_primaires"] = run_primary_watch(db) _log(db, "rapports-primaires", "ok", str(report["rapports_primaires"])[:500]) except Exception as e: report["rapports_primaires"] = f"échec: {e}" _log(db, "rapports-primaires", "error", traceback.format_exc()) try: # source canonique + découverte de comptes : hebdomadaires marker = (db.query(Mo.Indicator) .filter_by(name="canonical_checks_last") .order_by(Mo.Indicator.as_of.desc()).first()) if marker is None or (date.today() - marker.as_of).days >= 7: from .ingest.elections_quebec import run_canonical_checks from .ingest.social_account_discovery import run_discovery report["dgeq_canonique"] = run_canonical_checks(db) report["découverte_comptes"] = run_discovery(db) # validation hebdomadaire AUTOMATIQUE : replay + ablation # régénérés sans intervention (rapports publics à jour) try: from .modeling.validation.historical_replay import run_replay run_replay() _log(db, "replay-hebdo", "ok", "rapport régénéré") except Exception: _log(db, "replay-hebdo", "error", traceback.format_exc()) try: from .modeling.validation.ablation import run_ablation run_ablation(db, n_sims=6000) _log(db, "ablation-hebdo", "ok", "rapport régénéré") except Exception: _log(db, "ablation-hebdo", "error", traceback.format_exc()) db.add(Mo.Indicator(name="canonical_checks_last", as_of=date.today(), value=1.0, method="scheduler")) db.commit() _log(db, "canonique-hebdo", "ok", str({**report.get("dgeq_canonique", {}), **report.get("découverte_comptes", {})})[:500]) except Exception as e: report["canonique"] = f"échec: {e}" _log(db, "canonique-hebdo", "error", traceback.format_exc()) try: from .ingest.social_pulse import run_pulse report["pouls_social"] = run_pulse(db) _log(db, "pouls-social", "ok", str(report["pouls_social"])[:500]) except Exception as e: report["pouls_social"] = f"échec: {e}" _log(db, "pouls-social", "error", traceback.format_exc()) try: from .ingest.wiki_attention import refresh as att_refresh report["attention"] = att_refresh(db) _log(db, "attention-wiki", "ok", str(report["attention"])[:500]) except Exception as e: report["attention"] = f"échec: {e}" _log(db, "attention-wiki", "error", traceback.format_exc()) try: from .modeling.signals.event_detection import detect_events report["événements"] = detect_events(db) _log(db, "event-engine", "ok", str(report["événements"])) except Exception as e: report["événements"] = f"échec: {e}" _log(db, "event-engine", "error", traceback.format_exc()) try: from .modeling.synthetic_poll import run as synth_run report["sondage_synthétique"] = synth_run(db) _log(db, "sondage-synthetique", "ok", str(report["sondage_synthétique"])[:500]) except Exception as e: report["sondage_synthétique"] = f"échec: {e}" _log(db, "sondage-synthetique", "error", traceback.format_exc()) try: run = run_forecast(db) report["forecast_run_id"] = run.id _log(db, "model+simulate", "ok", f"run {run.id}") except Exception as e: report["forecast"] = f"échec: {e}" _log(db, "model+simulate", "error", traceback.format_exc()) return report finally: db.close() if __name__ == "__main__": import json logging.basicConfig(level=logging.INFO) print(json.dumps(run_pipeline(), default=str, indent=2, ensure_ascii=False))