diff --git a/.env.example b/.env.example index 6a09d029..d27557c9 100644 --- a/.env.example +++ b/.env.example @@ -17,6 +17,7 @@ OPEN_METEO_DISK_CACHE_PATH=/var/lib/polyweather/open_meteo_cache.json # Windows / macOS can usually keep the defaults. UID=1000 GID=1000 +POLYWEATHER_STATE_STORAGE_MODE=dual ######################################## # 2) Telegram bot minimal diff --git a/artifacts/probability_calibration/rollout_report.json b/artifacts/probability_calibration/rollout_report.json new file mode 100644 index 00000000..8dcb6a73 --- /dev/null +++ b/artifacts/probability_calibration/rollout_report.json @@ -0,0 +1,72 @@ +{ + "evaluation_report_path": "E:\\web\\PolyWeather\\artifacts\\probability_calibration\\evaluation_report.json", + "shadow_report_path": "E:\\web\\PolyWeather\\artifacts\\probability_calibration\\shadow_report.json", + "evaluation_report_exists": true, + "shadow_report_exists": true, + "decision": { + "decision": "hold", + "ready_for_primary": false, + "summary": "当前指标不足以切换 emos_primary,应继续保持 shadow。", + "thresholds": { + "evaluation_min_samples": 80, + "shadow_min_samples": 50, + "max_delta_mae": 0.05, + "min_delta_crps": -0.02, + "min_delta_bucket_hit_rate": 0.0, + "max_delta_bucket_brier_promote": 0.02, + "max_delta_bucket_brier_observe": 0.15 + }, + "evaluation": { + "sample_count": 105, + "delta_crps": -0.093663, + "delta_mae": 0.0, + "delta_bucket_hit_rate": 0.0 + }, + "shadow": { + "sample_count": 103, + "delta_mae": 0.012708, + "delta_bucket_hit_rate": 0.009709, + "delta_bucket_brier": 0.293835 + }, + "blocking_reasons": [ + "shadow bucket brier 退化超限:delta=0.293835" + ], + "worst_shadow_regressions": [ + { + "city": "dallas", + "samples": 4, + "delta_mae": 0.114807, + "delta_bucket_hit_rate": 0.0, + "delta_bucket_brier": 0.778678 + }, + { + "city": "chicago", + "samples": 4, + "delta_mae": 0.075265, + "delta_bucket_hit_rate": 0.0, + "delta_bucket_brier": 0.746156 + }, + { + "city": "seattle", + "samples": 4, + "delta_mae": 0.11262, + "delta_bucket_hit_rate": 0.0, + "delta_bucket_brier": 0.692003 + }, + { + "city": "atlanta", + "samples": 4, + "delta_mae": 0.293028, + "delta_bucket_hit_rate": -0.25, + "delta_bucket_brier": 0.601425 + }, + { + "city": "miami", + "samples": 4, + "delta_mae": 0.241559, + "delta_bucket_hit_rate": -0.5, + "delta_bucket_brier": 0.478245 + } + ] + } +} \ No newline at end of file diff --git a/data/probability_training_snapshots.jsonl b/data/probability_training_snapshots.jsonl index bb34e8e7..96c1b9b4 100644 --- a/data/probability_training_snapshots.jsonl +++ b/data/probability_training_snapshots.jsonl @@ -1 +1,37 @@ {"city": "ankara", "timestamp": "2026-03-20T12:00:00+03:00", "date": "2026-03-20", "temp_symbol": "°C", "raw_mu": 15.2, "raw_sigma": 1.2, "deb_prediction": 15.4, "ensemble": {"p10": 14.8, "median": 15.8, "p90": 17.9}, "multi_model": {"ECMWF": 15.8, "GFS": 14.1, "ICON": 15.9}, "max_so_far": 15.0, "peak_status": "before", "prob_snapshot": [{"v": 15, "p": 0.552}, {"v": 16, "p": 0.377}], "shadow_prob_snapshot": [{"v": 15, "p": 0.324}, {"v": 16, "p": 0.238}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320130245", "calibration_source": "artifacts/probability_calibration/default.json", "calibrated_mu": 15.1, "calibrated_sigma": 1.25} +{"city": "test_city", "timestamp": "2026-03-04 10:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.5625, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 26.0, "peak_status": "before", "prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "shadow_prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.5625} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 33.3, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 33.0, "peak_status": "in_window", "prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "shadow_prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 33.3, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 22:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 16:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.85, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.5, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 29.5, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "shadow_prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.85, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 14:30", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 10:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.5625, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 26.0, "peak_status": "before", "prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "shadow_prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.5625} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 33.3, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 33.0, "peak_status": "in_window", "prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "shadow_prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 33.3, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 22:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 16:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.85, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.5, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 29.5, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "shadow_prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.85, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 14:30", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 10:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.5625, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 26.0, "peak_status": "before", "prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "shadow_prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.5625} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 33.3, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 33.0, "peak_status": "in_window", "prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "shadow_prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 33.3, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 22:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 16:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.85, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.5, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 29.5, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "shadow_prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.85, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 14:30", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 10:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.5625, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 26.0, "peak_status": "before", "prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "shadow_prob_snapshot": [{"v": 30, "p": 0.254}, {"v": 29, "p": 0.234}, {"v": 31, "p": 0.185}, {"v": 28, "p": 0.146}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.5625} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 33.3, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 33.0, "peak_status": "in_window", "prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "shadow_prob_snapshot": [{"v": 33, "p": 0.456}, {"v": 34, "p": 0.391}, {"v": 35, "p": 0.153}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 33.3, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 17:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 22:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": null, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "past", "prob_snapshot": [{"v": 28, "p": 1.0}], "shadow_prob_snapshot": [], "probability_engine": "legacy", "probability_mode": "legacy", "calibration_version": null, "calibration_source": null, "calibrated_mu": null, "calibrated_sigma": null} +{"city": "test_city", "timestamp": "2026-03-04 16:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 23.0, "raw_sigma": 0.46875, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 23.0, "peak_status": "past", "prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "shadow_prob_snapshot": [{"v": 23, "p": 0.834}, {"v": 24, "p": 0.166}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 23.0, "calibrated_sigma": 0.46875} +{"city": "test_city", "timestamp": "2026-03-04 14:00", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.85, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.5, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 29.5, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "shadow_prob_snapshot": [{"v": 30, "p": 0.565}, {"v": 31, "p": 0.341}, {"v": 32, "p": 0.094}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.85, "calibrated_sigma": 1.09375} +{"city": "test_city", "timestamp": "2026-03-04 14:30", "date": "2026-03-04", "temp_symbol": "°C", "raw_mu": 29.7, "raw_sigma": 1.09375, "deb_prediction": null, "ensemble": {"p10": 27.0, "median": 29.0, "p90": 31.0}, "multi_model": {"Open-Meteo": 30.0}, "max_so_far": 28.0, "peak_status": "in_window", "prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "shadow_prob_snapshot": [{"v": 30, "p": 0.35}, {"v": 29, "p": 0.299}, {"v": 31, "p": 0.187}, {"v": 28, "p": 0.117}], "probability_engine": "legacy", "probability_mode": "emos_shadow", "calibration_version": "emos-20260320132525", "calibration_source": "artifacts\\probability_calibration\\default.json", "calibrated_mu": 29.7, "calibrated_sigma": 1.09375} diff --git a/scripts/export_runtime_state_from_sqlite.py b/scripts/export_runtime_state_from_sqlite.py new file mode 100644 index 00000000..fafca072 --- /dev/null +++ b/scripts/export_runtime_state_from_sqlite.py @@ -0,0 +1,53 @@ + +import argparse +import json +import os +import sys + +PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if PROJECT_ROOT not in sys.path: + sys.path.insert(0, PROJECT_ROOT) + +from src.database.runtime_state import ( # noqa: E402 + DailyRecordRepository, + OpenMeteoCacheRepository, + ProbabilitySnapshotRepository, + TelegramAlertStateRepository, +) + + +def main(): + parser = argparse.ArgumentParser(description='Export runtime state from SQLite back to legacy JSON files.') + parser.add_argument('--daily-records', default=os.path.join(PROJECT_ROOT, 'data', 'daily_records.json')) + parser.add_argument('--telegram-state', default=os.path.join(PROJECT_ROOT, 'data', 'telegram_alert_state.json')) + parser.add_argument('--snapshots', default=os.path.join(PROJECT_ROOT, 'data', 'probability_training_snapshots.jsonl')) + parser.add_argument('--open-meteo-cache', default=os.path.join(PROJECT_ROOT, 'data', 'open_meteo_cache.json')) + parser.add_argument('--open-meteo-max-age', type=int, default=int(os.getenv('OPEN_METEO_DISK_CACHE_MAX_AGE_SEC', '86400'))) + args = parser.parse_args() + + daily = DailyRecordRepository().load_all() + telegram = TelegramAlertStateRepository().load_state() + snapshots = ProbabilitySnapshotRepository().load_all_rows() + open_meteo = OpenMeteoCacheRepository().load_payload(args.open_meteo_max_age) + + os.makedirs(os.path.dirname(os.path.abspath(args.daily_records)), exist_ok=True) + with open(args.daily_records, 'w', encoding='utf-8') as fh: + json.dump(daily, fh, ensure_ascii=False, indent=2) + with open(args.telegram_state, 'w', encoding='utf-8') as fh: + json.dump(telegram, fh, ensure_ascii=False, indent=2) + with open(args.snapshots, 'w', encoding='utf-8') as fh: + for row in snapshots: + fh.write(json.dumps(row, ensure_ascii=False) + '\n') + with open(args.open_meteo_cache, 'w', encoding='utf-8') as fh: + json.dump(open_meteo, fh, ensure_ascii=False) + + print(json.dumps({ + 'daily_records_exported': sum(len(v) for v in daily.values()), + 'telegram_state_exported': len((telegram.get('last_by_city') or {})) + len((telegram.get('by_signature') or {})), + 'snapshots_exported': len(snapshots), + 'open_meteo_cache_exported': sum(len((open_meteo.get(k) or {})) for k in ('forecast', 'ensemble', 'multi_model')), + }, ensure_ascii=False, indent=2)) + + +if __name__ == '__main__': + main() diff --git a/scripts/fit_probability_calibration.py b/scripts/fit_probability_calibration.py index 7a27a4c3..9f258bb5 100644 --- a/scripts/fit_probability_calibration.py +++ b/scripts/fit_probability_calibration.py @@ -14,6 +14,11 @@ from src.analysis.probability_calibration import ( # noqa: E402 fit_calibration, ) from src.analysis.deb_algorithm import load_history # noqa: E402 +from src.database.runtime_state import ( # noqa: E402 + ProbabilitySnapshotRepository, + STATE_STORAGE_SQLITE, + get_state_storage_mode, +) def _sf(value): @@ -41,6 +46,8 @@ def _load_history_with_fallback(path): def _load_snapshot_rows(path): + if get_state_storage_mode() == STATE_STORAGE_SQLITE: + return ProbabilitySnapshotRepository().load_all_rows() rows = [] if not path or not os.path.exists(path): return rows diff --git a/scripts/judge_probability_rollout.py b/scripts/judge_probability_rollout.py new file mode 100644 index 00000000..a06862b7 --- /dev/null +++ b/scripts/judge_probability_rollout.py @@ -0,0 +1,56 @@ +import argparse +import json +import os +import sys + +PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if PROJECT_ROOT not in sys.path: + sys.path.insert(0, PROJECT_ROOT) + +from src.analysis.probability_rollout import build_rollout_report # noqa: E402 + + +def main(): + parser = argparse.ArgumentParser(description="Judge whether EMOS is ready for primary rollout.") + parser.add_argument( + "--evaluation-report", + default=os.path.join( + PROJECT_ROOT, + "artifacts", + "probability_calibration", + "evaluation_report.json", + ), + ) + parser.add_argument( + "--shadow-report", + default=os.path.join( + PROJECT_ROOT, + "artifacts", + "probability_calibration", + "shadow_report.json", + ), + ) + parser.add_argument( + "--output", + default=os.path.join( + PROJECT_ROOT, + "artifacts", + "probability_calibration", + "rollout_report.json", + ), + ) + args = parser.parse_args() + + payload = build_rollout_report(args.evaluation_report, args.shadow_report) + output_dir = os.path.dirname(os.path.abspath(args.output)) + if output_dir: + os.makedirs(output_dir, exist_ok=True) + with open(args.output, "w", encoding="utf-8") as fh: + json.dump(payload, fh, ensure_ascii=False, indent=2) + + print(json.dumps(payload["decision"], ensure_ascii=False, indent=2)) + print(f"saved rollout report to {args.output}") + + +if __name__ == "__main__": + main() diff --git a/scripts/migrate_runtime_state_to_sqlite.py b/scripts/migrate_runtime_state_to_sqlite.py new file mode 100644 index 00000000..410dd19e --- /dev/null +++ b/scripts/migrate_runtime_state_to_sqlite.py @@ -0,0 +1,73 @@ + +import argparse +import json +import os +import sys + +PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if PROJECT_ROOT not in sys.path: + sys.path.insert(0, PROJECT_ROOT) + +from src.database.runtime_state import ( # noqa: E402 + DailyRecordRepository, + OpenMeteoCacheRepository, + ProbabilitySnapshotRepository, + TelegramAlertStateRepository, +) + + +def _load_json(path, default): + if not path or not os.path.exists(path): + return default + with open(path, 'r', encoding='utf-8') as fh: + data = json.load(fh) + return data + + +def _load_jsonl(path): + rows = [] + if not path or not os.path.exists(path): + return rows + with open(path, 'r', encoding='utf-8') as fh: + for line in fh: + line = line.strip() + if not line: + continue + try: + row = json.loads(line) + except Exception: + continue + if isinstance(row, dict): + rows.append(row) + return rows + + +def main(): + parser = argparse.ArgumentParser(description='Migrate runtime JSON state into SQLite.') + parser.add_argument('--daily-records', default=os.path.join(PROJECT_ROOT, 'data', 'daily_records.json')) + parser.add_argument('--telegram-state', default=os.path.join(PROJECT_ROOT, 'data', 'telegram_alert_state.json')) + parser.add_argument('--snapshots', default=os.path.join(PROJECT_ROOT, 'data', 'probability_training_snapshots.jsonl')) + parser.add_argument('--open-meteo-cache', default=os.path.join(PROJECT_ROOT, 'data', 'open_meteo_cache.json')) + parser.add_argument('--open-meteo-max-age', type=int, default=int(os.getenv('OPEN_METEO_DISK_CACHE_MAX_AGE_SEC', '86400'))) + args = parser.parse_args() + + daily = _load_json(args.daily_records, {}) + telegram = _load_json(args.telegram_state, {'last_by_city': {}, 'by_signature': {}}) + snapshots = _load_jsonl(args.snapshots) + open_meteo = _load_json(args.open_meteo_cache, {'forecast': {}, 'ensemble': {}, 'multi_model': {}, 'saved_at': 0}) + + daily_count = DailyRecordRepository().replace_all(daily if isinstance(daily, dict) else {}) + telegram_count = TelegramAlertStateRepository().replace_from_state(telegram if isinstance(telegram, dict) else {}) + snapshot_count = ProbabilitySnapshotRepository().replace_all(snapshots) + cache_count = OpenMeteoCacheRepository().replace_payload(open_meteo if isinstance(open_meteo, dict) else {}, args.open_meteo_max_age) + + print(json.dumps({ + 'daily_records_imported': daily_count, + 'telegram_state_imported': telegram_count, + 'snapshots_imported': snapshot_count, + 'open_meteo_cache_imported': cache_count, + }, ensure_ascii=False, indent=2)) + + +if __name__ == '__main__': + main() diff --git a/scripts/verify_runtime_state_storage.py b/scripts/verify_runtime_state_storage.py new file mode 100644 index 00000000..9cb78f4e --- /dev/null +++ b/scripts/verify_runtime_state_storage.py @@ -0,0 +1,101 @@ + +import argparse +import json +import os +import sys + +PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if PROJECT_ROOT not in sys.path: + sys.path.insert(0, PROJECT_ROOT) + +from src.database.runtime_state import ( # noqa: E402 + DailyRecordRepository, + OpenMeteoCacheRepository, + ProbabilitySnapshotRepository, + TelegramAlertStateRepository, +) + + +def _load_json(path, default): + if not path or not os.path.exists(path): + return default + with open(path, 'r', encoding='utf-8') as fh: + return json.load(fh) + + +def _load_jsonl(path): + rows = [] + if not path or not os.path.exists(path): + return rows + with open(path, 'r', encoding='utf-8') as fh: + for line in fh: + line = line.strip() + if not line: + continue + try: + row = json.loads(line) + except Exception: + continue + if isinstance(row, dict): + rows.append(row) + return rows + + +def _norm_json(obj): + return json.loads(json.dumps(obj, ensure_ascii=False, sort_keys=True)) + + +def _norm_open_meteo_payload(payload): + payload = dict(payload or {}) + payload.pop('saved_at', None) + return _norm_json(payload) + + +def main(): + parser = argparse.ArgumentParser(description='Verify SQLite runtime state against legacy JSON files.') + parser.add_argument('--daily-records', default=os.path.join(PROJECT_ROOT, 'data', 'daily_records.json')) + parser.add_argument('--telegram-state', default=os.path.join(PROJECT_ROOT, 'data', 'telegram_alert_state.json')) + parser.add_argument('--snapshots', default=os.path.join(PROJECT_ROOT, 'data', 'probability_training_snapshots.jsonl')) + parser.add_argument('--open-meteo-cache', default=os.path.join(PROJECT_ROOT, 'data', 'open_meteo_cache.json')) + parser.add_argument('--open-meteo-max-age', type=int, default=int(os.getenv('OPEN_METEO_DISK_CACHE_MAX_AGE_SEC', '86400'))) + args = parser.parse_args() + + file_daily = _load_json(args.daily_records, {}) + file_telegram = _load_json(args.telegram_state, {'last_by_city': {}, 'by_signature': {}}) + file_snapshots = _load_jsonl(args.snapshots) + file_cache = _load_json(args.open_meteo_cache, {'forecast': {}, 'ensemble': {}, 'multi_model': {}, 'saved_at': 0}) + + db_daily = DailyRecordRepository().load_all() + db_telegram = TelegramAlertStateRepository().load_state() + db_snapshots = ProbabilitySnapshotRepository().load_all_rows() + db_cache = OpenMeteoCacheRepository().load_payload(args.open_meteo_max_age) + + report = { + 'daily_records': { + 'file_cities': len(file_daily or {}), + 'db_cities': len(db_daily or {}), + 'equal': _norm_json(file_daily or {}) == _norm_json(db_daily or {}), + }, + 'telegram_state': { + 'equal': _norm_json(file_telegram or {}) == _norm_json(db_telegram or {}), + 'file_last_by_city': len((file_telegram or {}).get('last_by_city') or {}), + 'db_last_by_city': len((db_telegram or {}).get('last_by_city') or {}), + }, + 'snapshots': { + 'file_rows': len(file_snapshots), + 'db_rows': len(db_snapshots), + 'equal': _norm_json(file_snapshots) == _norm_json(db_snapshots), + }, + 'open_meteo_cache': { + 'file_forecast': len((file_cache or {}).get('forecast') or {}), + 'db_forecast': len((db_cache or {}).get('forecast') or {}), + 'equal': _norm_open_meteo_payload(file_cache or {}) == _norm_open_meteo_payload(db_cache or {}), + }, + } + report['ok'] = all(section.get('equal') for section in report.values() if isinstance(section, dict)) + print(json.dumps(report, ensure_ascii=False, indent=2)) + raise SystemExit(0 if report['ok'] else 1) + + +if __name__ == '__main__': + main() diff --git a/src/analysis/deb_algorithm.py b/src/analysis/deb_algorithm.py index 853a420a..49814688 100644 --- a/src/analysis/deb_algorithm.py +++ b/src/analysis/deb_algorithm.py @@ -3,6 +3,13 @@ import json from datetime import datetime, timedelta import requests from src.analysis.settlement_rounding import apply_city_settlement +from loguru import logger +from src.database.runtime_state import ( + DailyRecordRepository, + STATE_STORAGE_DUAL, + STATE_STORAGE_SQLITE, + get_state_storage_mode, +) # Cross-platform file locking import sys @@ -36,6 +43,7 @@ else: # Simple memory cache to avoid blasting the disk if queried 10 times a minute _history_cache = {} _history_mtime = 0 +_daily_record_repo = DailyRecordRepository() def _sf(value): @@ -54,7 +62,24 @@ def _is_excluded_model_name(model_name: str) -> bool: def load_history(filepath): global _history_cache, _history_mtime + mode = get_state_storage_mode() + + if mode == STATE_STORAGE_SQLITE: + try: + data = _daily_record_repo.load_all() + _history_cache = data + return data + except Exception as e: + logger.error(f"Error loading daily records from sqlite, fallback to file: {e}") + if not os.path.exists(filepath): + if mode == STATE_STORAGE_DUAL: + try: + data = _daily_record_repo.load_all() + _history_cache = data + return data + except Exception: + return {} return {} try: @@ -80,6 +105,18 @@ def load_history(filepath): def save_history(filepath, data): global _history_cache, _history_mtime _history_cache = data + mode = get_state_storage_mode() + + if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}: + try: + _daily_record_repo.replace_all(data) + except Exception as e: + logger.error(f"Error saving daily records to sqlite: {e}") + if mode == STATE_STORAGE_SQLITE: + return + + if mode == STATE_STORAGE_SQLITE: + return try: with open(filepath, "w", encoding="utf-8") as f: _lock_ex(f) @@ -251,6 +288,7 @@ def update_daily_record( ) history_file = os.path.join(project_root, "data", "daily_records.json") + mode = get_state_storage_mode() data = load_history(history_file) if city_name not in data: data[city_name] = {} @@ -359,7 +397,18 @@ def update_daily_record( for d in old_dates: del data[city][d] - save_history(history_file, data) + if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}: + try: + _daily_record_repo.upsert_record(city_name, date_str, existing) + cutoff = (datetime.now() - timedelta(days=14)).strftime("%Y-%m-%d") + _daily_record_repo.delete_older_than(cutoff) + except Exception as e: + logger.error(f"Error upserting daily record to sqlite city={city_name} date={date_str}: {e}") + if mode == STATE_STORAGE_SQLITE: + raise + + if mode != STATE_STORAGE_SQLITE: + save_history(history_file, data) def calculate_dynamic_weights(city_name, current_forecasts, lookback_days=7): diff --git a/src/analysis/probability_rollout.py b/src/analysis/probability_rollout.py new file mode 100644 index 00000000..aea2919c --- /dev/null +++ b/src/analysis/probability_rollout.py @@ -0,0 +1,217 @@ +from __future__ import annotations + +import json +import os +from typing import Any, Dict, List, Optional + + +def _load_json_file(path: str) -> Dict[str, Any]: + try: + with open(path, "r", encoding="utf-8") as fh: + data = json.load(fh) + return data if isinstance(data, dict) else {} + except Exception: + return {} + + +def _sf(value: Any) -> Optional[float]: + if value is None: + return None + try: + return float(value) + except Exception: + return None + + +def _append_reason(reasons: List[str], condition: bool, message: str) -> None: + if condition: + reasons.append(message) + + +def _top_shadow_regressions(by_city: Dict[str, Any], limit: int = 5) -> List[Dict[str, Any]]: + rows: List[Dict[str, Any]] = [] + for city, metrics in (by_city or {}).items(): + if not isinstance(metrics, dict): + continue + rows.append( + { + "city": city, + "samples": int(metrics.get("samples") or 0), + "delta_mae": _sf(metrics.get("delta_mae")), + "delta_bucket_hit_rate": _sf(metrics.get("delta_bucket_hit_rate")), + "delta_bucket_brier": _sf(metrics.get("delta_bucket_brier")), + } + ) + + rows.sort( + key=lambda row: ( + -(row["delta_bucket_brier"] or 0.0), + row["delta_bucket_hit_rate"] or 0.0, + -(row["delta_mae"] or 0.0), + ) + ) + return rows[:limit] + + +def judge_probability_rollout( + evaluation_report: Dict[str, Any], + shadow_report: Dict[str, Any], +) -> Dict[str, Any]: + thresholds = { + "evaluation_min_samples": 80, + "shadow_min_samples": 50, + "max_delta_mae": 0.05, + "min_delta_crps": -0.02, + "min_delta_bucket_hit_rate": 0.0, + "max_delta_bucket_brier_promote": 0.02, + "max_delta_bucket_brier_observe": 0.15, + } + + eval_summary = (evaluation_report or {}).get("summary") or {} + eval_delta = eval_summary.get("delta") or {} + shadow_summary = (shadow_report or {}).get("summary") or {} + + eval_samples = int(eval_summary.get("sample_count") or 0) + shadow_samples = int(shadow_summary.get("samples") or 0) + delta_crps = _sf(eval_delta.get("crps")) + delta_mae = _sf(eval_delta.get("mae")) + delta_hit = _sf(eval_delta.get("bucket_hit_rate")) + shadow_delta_mae = _sf(shadow_summary.get("delta_mae")) + shadow_delta_hit = _sf(shadow_summary.get("delta_bucket_hit_rate")) + shadow_delta_brier = _sf(shadow_summary.get("delta_bucket_brier")) + + promote_reasons: List[str] = [] + _append_reason( + promote_reasons, + eval_samples < thresholds["evaluation_min_samples"], + f"离线评估样本不足:{eval_samples} < {thresholds['evaluation_min_samples']}", + ) + _append_reason( + promote_reasons, + shadow_samples < thresholds["shadow_min_samples"], + f"shadow 样本不足:{shadow_samples} < {thresholds['shadow_min_samples']}", + ) + _append_reason( + promote_reasons, + delta_crps is None or delta_crps > thresholds["min_delta_crps"], + f"离线 CRPS 改善不足:delta={delta_crps}", + ) + _append_reason( + promote_reasons, + delta_mae is None or delta_mae > thresholds["max_delta_mae"], + f"离线 MAE 退化超限:delta={delta_mae}", + ) + _append_reason( + promote_reasons, + delta_hit is None or delta_hit < thresholds["min_delta_bucket_hit_rate"], + f"离线 bucket 命中率下降:delta={delta_hit}", + ) + _append_reason( + promote_reasons, + shadow_delta_mae is None or shadow_delta_mae > thresholds["max_delta_mae"], + f"shadow MAE 退化超限:delta={shadow_delta_mae}", + ) + _append_reason( + promote_reasons, + shadow_delta_hit is None or shadow_delta_hit < thresholds["min_delta_bucket_hit_rate"], + f"shadow bucket 命中率下降:delta={shadow_delta_hit}", + ) + _append_reason( + promote_reasons, + shadow_delta_brier is None + or shadow_delta_brier > thresholds["max_delta_bucket_brier_promote"], + f"shadow bucket brier 退化超限:delta={shadow_delta_brier}", + ) + + if not promote_reasons: + decision = "promote" + summary = "离线与 shadow 指标均达标,可以考虑切换 emos_primary。" + else: + observe_reasons: List[str] = [] + _append_reason( + observe_reasons, + eval_samples < thresholds["evaluation_min_samples"], + f"离线评估样本不足:{eval_samples}", + ) + _append_reason( + observe_reasons, + delta_crps is None or delta_crps > thresholds["min_delta_crps"], + f"离线 CRPS 改善不足:delta={delta_crps}", + ) + _append_reason( + observe_reasons, + delta_mae is None or delta_mae > thresholds["max_delta_mae"], + f"离线 MAE 退化超限:delta={delta_mae}", + ) + _append_reason( + observe_reasons, + delta_hit is None or delta_hit < thresholds["min_delta_bucket_hit_rate"], + f"离线 bucket 命中率下降:delta={delta_hit}", + ) + _append_reason( + observe_reasons, + shadow_samples < thresholds["shadow_min_samples"], + f"shadow 样本不足:{shadow_samples}", + ) + _append_reason( + observe_reasons, + shadow_delta_mae is None or shadow_delta_mae > thresholds["max_delta_mae"], + f"shadow MAE 退化超限:delta={shadow_delta_mae}", + ) + _append_reason( + observe_reasons, + shadow_delta_hit is None or shadow_delta_hit < thresholds["min_delta_bucket_hit_rate"], + f"shadow bucket 命中率下降:delta={shadow_delta_hit}", + ) + _append_reason( + observe_reasons, + shadow_delta_brier is None + or shadow_delta_brier > thresholds["max_delta_bucket_brier_observe"], + f"shadow bucket brier 退化偏大:delta={shadow_delta_brier}", + ) + + if not observe_reasons: + decision = "observe" + summary = "离线评估达标,但 shadow 仍需继续观察,暂不切主路径。" + else: + decision = "hold" + summary = "当前指标不足以切换 emos_primary,应继续保持 shadow。" + + return { + "decision": decision, + "ready_for_primary": decision == "promote", + "summary": summary, + "thresholds": thresholds, + "evaluation": { + "sample_count": eval_samples, + "delta_crps": delta_crps, + "delta_mae": delta_mae, + "delta_bucket_hit_rate": delta_hit, + }, + "shadow": { + "sample_count": shadow_samples, + "delta_mae": shadow_delta_mae, + "delta_bucket_hit_rate": shadow_delta_hit, + "delta_bucket_brier": shadow_delta_brier, + }, + "blocking_reasons": promote_reasons, + "worst_shadow_regressions": _top_shadow_regressions( + (shadow_report or {}).get("by_city") or {} + ), + } + + +def build_rollout_report( + evaluation_report_path: str, + shadow_report_path: str, +) -> Dict[str, Any]: + evaluation_report = _load_json_file(evaluation_report_path) + shadow_report = _load_json_file(shadow_report_path) + decision = judge_probability_rollout(evaluation_report, shadow_report) + return { + "evaluation_report_path": evaluation_report_path, + "shadow_report_path": shadow_report_path, + "evaluation_report_exists": os.path.exists(evaluation_report_path), + "shadow_report_exists": os.path.exists(shadow_report_path), + "decision": decision, + } diff --git a/src/analysis/probability_snapshot_archive.py b/src/analysis/probability_snapshot_archive.py index c5964d32..02823cf8 100644 --- a/src/analysis/probability_snapshot_archive.py +++ b/src/analysis/probability_snapshot_archive.py @@ -5,10 +5,18 @@ import os from datetime import datetime from typing import Any, Dict, List, Optional +from src.database.runtime_state import ( + ProbabilitySnapshotRepository, + STATE_STORAGE_DUAL, + STATE_STORAGE_SQLITE, + get_state_storage_mode, +) + DEDUP_SCAN_LINES = 200 MU_THRESHOLD = 0.2 SIGMA_THRESHOLD = 0.15 MAX_SO_FAR_THRESHOLD = 0.2 +_snapshot_repo = ProbabilitySnapshotRepository() def _sf(value: Any) -> Optional[float]: @@ -79,7 +87,15 @@ def _load_recent_rows(path: str, max_lines: int = DEDUP_SCAN_LINES) -> List[Dict def _should_skip_append(path: str, payload: Dict[str, Any]) -> bool: - recent_rows = _load_recent_rows(path) + mode = get_state_storage_mode() + if mode == STATE_STORAGE_SQLITE: + recent_rows = _snapshot_repo.load_recent_rows( + str(payload.get("city") or ""), + str(payload.get("date") or ""), + DEDUP_SCAN_LINES, + ) + else: + recent_rows = _load_recent_rows(path) city = payload.get("city") date_str = payload.get("date") if not city or not date_str: @@ -183,5 +199,10 @@ def append_probability_snapshot( if _should_skip_append(path, payload): return - with open(path, "a", encoding="utf-8") as fh: - fh.write(json.dumps(payload, ensure_ascii=False) + "\n") + mode = get_state_storage_mode() + if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}: + _snapshot_repo.append_snapshot(payload) + + if mode != STATE_STORAGE_SQLITE: + with open(path, "a", encoding="utf-8") as fh: + fh.write(json.dumps(payload, ensure_ascii=False) + "\n") diff --git a/src/data_collection/metar_sources.py b/src/data_collection/metar_sources.py index f5912ac8..865f1750 100644 --- a/src/data_collection/metar_sources.py +++ b/src/data_collection/metar_sources.py @@ -8,6 +8,8 @@ from typing import Dict, List, Optional import requests from loguru import logger +from src.utils.metrics import record_source_call + class MetarSourceMixin: def get_icao_code(self, city: str) -> Optional[str]: @@ -24,9 +26,11 @@ class MetarSourceMixin: self, city: str, use_fahrenheit: bool = False, utc_offset: int = 0 ) -> Optional[Dict]: """从 NOAA Aviation Weather Center 获取 METAR 航空气象数据。""" + started = time.perf_counter() icao = self.get_icao_code(city) if not icao: logger.warning(f"未找到城市 {city} 对应的 ICAO 代码") + record_source_call("metar", "current", "missing_icao", (time.perf_counter() - started) * 1000.0) return None cache_key = f"{icao}:{utc_offset}:{use_fahrenheit}" @@ -35,6 +39,7 @@ class MetarSourceMixin: cached = self._metar_cache.get(cache_key) if cached and now_ts - cached["t"] < self.metar_cache_ttl_sec: logger.debug(f"METAR cache hit {icao} age={int(now_ts - cached['t'])}s") + record_source_call("metar", "current", "cache_hit", (time.perf_counter() - started) * 1000.0) return cached["d"] try: @@ -201,6 +206,7 @@ class MetarSourceMixin: ) with self._metar_cache_lock: self._metar_cache[cache_key] = {"d": result, "t": now_ts} + record_source_call("metar", "current", "success", (time.perf_counter() - started) * 1000.0) return result except requests.exceptions.RequestException as exc: @@ -209,10 +215,13 @@ class MetarSourceMixin: stale = self._metar_cache.get(cache_key) if stale: logger.warning(f"METAR {icao} 请求失败,使用缓存回退") + record_source_call("metar", "current", "stale_cache", (time.perf_counter() - started) * 1000.0) return stale["d"] + record_source_call("metar", "current", "error", (time.perf_counter() - started) * 1000.0) return None except (KeyError, IndexError, TypeError) as exc: logger.error(f"METAR 数据解析失败 ({icao}): {exc}") + record_source_call("metar", "current", "parse_error", (time.perf_counter() - started) * 1000.0) return None def fetch_metar_nearby_cluster(self, icaos: List[str], use_fahrenheit: bool = False) -> list: diff --git a/src/data_collection/mgm_sources.py b/src/data_collection/mgm_sources.py index 9a3b7aa0..a43473c7 100644 --- a/src/data_collection/mgm_sources.py +++ b/src/data_collection/mgm_sources.py @@ -5,12 +5,15 @@ from typing import Dict, Optional from loguru import logger +from src.utils.metrics import record_source_call + class MgmSourceMixin: def fetch_from_mgm(self, istno: str) -> Optional[Dict]: """ 从土耳其气象局 (MGM) 获取实时数据和预测 (由用户提供其内部 API) """ + started = datetime.now() base_url = "https://servis.mgm.gov.tr/web" # 必须带 Origin,否则会被反爬拦截 headers = { @@ -202,9 +205,21 @@ class MgmSourceMixin: f"hourly points={len(hourly_rows)}, horizon={horizon_hours:.1f}h" ) + record_source_call( + "mgm", + "station", + "success" if "current" in results else "empty", + (datetime.now() - started).total_seconds() * 1000.0, + ) return results if "current" in results else None except Exception as e: logger.error(f"MGM API 请求失败 ({istno}): {e}") + record_source_call( + "mgm", + "station", + "error", + (datetime.now() - started).total_seconds() * 1000.0, + ) return None def fetch_mgm_nearby_stations(self, province: str, root_ist_no: str = None) -> list: @@ -212,6 +227,7 @@ class MgmSourceMixin: 获取一个土耳其省份内所有气象站的当前温度及经纬度 使用多线程辅助抓取,因为直接通过 il={province} 往往只返回 1 个站。 """ + started = datetime.now() base_url = "https://servis.mgm.gov.tr/web" headers = { "Origin": "https://www.mgm.gov.tr", @@ -243,6 +259,12 @@ class MgmSourceMixin: if not province_ist_nos: logger.warning(f"MGM 找不到省份 {province} 的站点元数据") + record_source_call( + "mgm", + "nearby", + "empty", + (datetime.now() - started).total_seconds() * 1000.0, + ) return [] # 同时确保我们关心的几个核心站一定在里面 @@ -318,8 +340,20 @@ class MgmSourceMixin: }) logger.info(f"📍 MGM 周边测站: 成功并发抓取 {len(results)} 个 {province} 站点的实时气温") + record_source_call( + "mgm", + "nearby", + "success" if results else "empty", + (datetime.now() - started).total_seconds() * 1000.0, + ) return results except Exception as e: logger.error(f"Failed to fetch MGM nearby stations for {province}: {e}") + record_source_call( + "mgm", + "nearby", + "error", + (datetime.now() - started).total_seconds() * 1000.0, + ) return [] diff --git a/src/data_collection/nws_open_meteo_sources.py b/src/data_collection/nws_open_meteo_sources.py index 40983ad7..7d4e8cef 100644 --- a/src/data_collection/nws_open_meteo_sources.py +++ b/src/data_collection/nws_open_meteo_sources.py @@ -6,6 +6,8 @@ from typing import Dict, Optional from loguru import logger +from src.utils.metrics import record_source_call + class NwsOpenMeteoSourceMixin: def fetch_nws(self, lat: float, lon: float) -> Optional[Dict]: @@ -13,6 +15,7 @@ class NwsOpenMeteoSourceMixin: 从 NWS (美国国家气象局) 获取高精度预报 仅适用于美国城市,全球 VPS 均可访问 """ + started = time.perf_counter() try: # 1. 获取网格点 points_url = f"https://api.weather.gov/points/{lat},{lon}" @@ -28,6 +31,7 @@ class NwsOpenMeteoSourceMixin: forecast_url = properties.get("forecast") hourly_url = properties.get("forecastHourly") if not forecast_url: + record_source_call("nws", "forecast", "empty", (time.perf_counter() - started) * 1000.0) return None # 2. 获取预报 @@ -39,6 +43,7 @@ class NwsOpenMeteoSourceMixin: periods = forecast_data.get("properties", {}).get("periods", []) if not periods: + record_source_call("nws", "forecast", "empty", (time.perf_counter() - started) * 1000.0) return None hourly_periods = [] @@ -89,7 +94,7 @@ class NwsOpenMeteoSourceMixin: today_high = p.get("temperature") break - return { + result = { "source": "nws", "today_high": today_high, "unit": "fahrenheit", @@ -124,8 +129,11 @@ class NwsOpenMeteoSourceMixin: ], "active_alerts": active_alerts, } + record_source_call("nws", "forecast", "success", (time.perf_counter() - started) * 1000.0) + return result except Exception as e: logger.warning(f"NWS 请求失败: {e}") + record_source_call("nws", "forecast", "error", (time.perf_counter() - started) * 1000.0) return None def fetch_from_open_meteo( @@ -144,6 +152,7 @@ class NwsOpenMeteoSourceMixin: forecast_days: Number of forecast days to fetch (default 14 to cover all market dates) use_fahrenheit: Whether to return temperatures in Fahrenheit (for US markets) """ + started = time.perf_counter() cache_key = ( f"{round(float(lat), 4)}:{round(float(lon), 4)}:" f"{forecast_days}:{'f' if use_fahrenheit else 'c'}" @@ -156,9 +165,11 @@ class NwsOpenMeteoSourceMixin: remaining = int(self._open_meteo_rate_limit_until - now_ts) logger.debug(f"Open-Meteo 冷却期中,跳过请求,还需 {remaining}s") with self._open_meteo_cache_lock: - stale = self._open_meteo_cache.get(cache_key) - if stale and isinstance(stale.get("data"), dict): - return dict(stale["data"]) + stale = self._open_meteo_cache.get(cache_key) + if stale and isinstance(stale.get("data"), dict): + record_source_call("open_meteo", "forecast", "stale_cache", (time.perf_counter() - started) * 1000.0) + return dict(stale["data"]) + record_source_call("open_meteo", "forecast", "cooldown_skip", (time.perf_counter() - started) * 1000.0) return None with self._open_meteo_cache_lock: cached = self._open_meteo_cache.get(cache_key) @@ -168,6 +179,7 @@ class NwsOpenMeteoSourceMixin: ): cached_data = cached.get("data") if isinstance(cached_data, dict): + record_source_call("open_meteo", "forecast", "cache_hit", (time.perf_counter() - started) * 1000.0) return dict(cached_data) try: url = "https://api.open-meteo.com/v1/forecast" @@ -259,6 +271,7 @@ class NwsOpenMeteoSourceMixin: "data": dict(result), } self._flush_open_meteo_disk_cache() + record_source_call("open_meteo", "forecast", "success", (time.perf_counter() - started) * 1000.0) return result except Exception as e: status_code = getattr(getattr(e, "response", None), "status_code", None) @@ -287,7 +300,9 @@ class NwsOpenMeteoSourceMixin: if stale and isinstance(stale.get("data"), dict): fallback = dict(stale["data"]) fallback["stale_cache"] = True + record_source_call("open_meteo", "forecast", "stale_cache", (time.perf_counter() - started) * 1000.0) return fallback + record_source_call("open_meteo", "forecast", "error", (time.perf_counter() - started) * 1000.0) return None def fetch_ensemble( @@ -300,6 +315,7 @@ class NwsOpenMeteoSourceMixin: 从 Open-Meteo Ensemble API 获取 51 成员集合预报 用于计算预报不确定性范围(散度) """ + started = time.perf_counter() cache_key = ( f"{round(float(lat), 4)}:{round(float(lon), 4)}:" f"{'f' if use_fahrenheit else 'c'}" @@ -314,7 +330,9 @@ class NwsOpenMeteoSourceMixin: with self._ensemble_cache_lock: stale = self._ensemble_cache.get(cache_key) if stale and isinstance(stale.get("data"), dict): + record_source_call("open_meteo", "ensemble", "stale_cache", (time.perf_counter() - started) * 1000.0) return dict(stale["data"]) + record_source_call("open_meteo", "ensemble", "cooldown_skip", (time.perf_counter() - started) * 1000.0) return None with self._ensemble_cache_lock: @@ -326,6 +344,7 @@ class NwsOpenMeteoSourceMixin: ): cached_data = cached.get("data") if isinstance(cached_data, dict): + record_source_call("open_meteo", "ensemble", "cache_hit", (time.perf_counter() - started) * 1000.0) return dict(cached_data) try: url = "https://ensemble-api.open-meteo.com/v1/ensemble" @@ -371,6 +390,7 @@ class NwsOpenMeteoSourceMixin: if len(today_highs) < 3: logger.warning(f"Ensemble 数据不足: 仅获取 {len(today_highs)} 个成员") + record_source_call("open_meteo", "ensemble", "empty", (time.perf_counter() - started) * 1000.0) return None today_highs.sort() @@ -400,6 +420,7 @@ class NwsOpenMeteoSourceMixin: "data": dict(result), } self._flush_open_meteo_disk_cache() + record_source_call("open_meteo", "ensemble", "success", (time.perf_counter() - started) * 1000.0) return result except Exception as e: status_code = getattr(getattr(e, "response", None), "status_code", None) @@ -425,7 +446,9 @@ class NwsOpenMeteoSourceMixin: if stale and isinstance(stale.get("data"), dict): fallback = dict(stale["data"]) fallback["stale_cache"] = True + record_source_call("open_meteo", "ensemble", "stale_cache", (time.perf_counter() - started) * 1000.0) return fallback + record_source_call("open_meteo", "ensemble", "error", (time.perf_counter() - started) * 1000.0) return None def fetch_multi_model( @@ -448,6 +471,7 @@ class NwsOpenMeteoSourceMixin: 返回 3 天的预报数据,支持今日+明日共识分析 """ + started = time.perf_counter() cache_city = str(city or "").strip().lower() cache_key = ( f"{round(float(lat), 4)}:{round(float(lon), 4)}:{cache_city}:" @@ -463,7 +487,9 @@ class NwsOpenMeteoSourceMixin: with self._multi_model_cache_lock: stale = self._multi_model_cache.get(cache_key) if stale and isinstance(stale.get("data"), dict): + record_source_call("open_meteo", "multi_model", "stale_cache", (time.perf_counter() - started) * 1000.0) return dict(stale["data"]) + record_source_call("open_meteo", "multi_model", "cooldown_skip", (time.perf_counter() - started) * 1000.0) return None with self._multi_model_cache_lock: @@ -475,6 +501,7 @@ class NwsOpenMeteoSourceMixin: ): cached_data = cached.get("data") if isinstance(cached_data, dict): + record_source_call("open_meteo", "multi_model", "cache_hit", (time.perf_counter() - started) * 1000.0) return dict(cached_data) try: url = "https://api.open-meteo.com/v1/forecast" @@ -524,6 +551,7 @@ class NwsOpenMeteoSourceMixin: if not daily_forecasts: logger.warning("Multi-model: 无有效模型数据") + record_source_call("open_meteo", "multi_model", "empty", (time.perf_counter() - started) * 1000.0) return None # 今天的预报 (向后兼容) @@ -548,6 +576,7 @@ class NwsOpenMeteoSourceMixin: "data": dict(result), } self._flush_open_meteo_disk_cache() + record_source_call("open_meteo", "multi_model", "success", (time.perf_counter() - started) * 1000.0) return result except Exception as e: status_code = getattr(getattr(e, "response", None), "status_code", None) @@ -573,6 +602,8 @@ class NwsOpenMeteoSourceMixin: if stale and isinstance(stale.get("data"), dict): fallback = dict(stale["data"]) fallback["stale_cache"] = True + record_source_call("open_meteo", "multi_model", "stale_cache", (time.perf_counter() - started) * 1000.0) return fallback + record_source_call("open_meteo", "multi_model", "error", (time.perf_counter() - started) * 1000.0) return None diff --git a/src/data_collection/open_meteo_cache.py b/src/data_collection/open_meteo_cache.py index 87ec1f74..667c8f6b 100644 --- a/src/data_collection/open_meteo_cache.py +++ b/src/data_collection/open_meteo_cache.py @@ -5,11 +5,30 @@ import os import time from loguru import logger +from src.database.runtime_state import ( + OpenMeteoCacheRepository, + STATE_STORAGE_DUAL, + STATE_STORAGE_SQLITE, + get_state_storage_mode, +) + +_open_meteo_cache_repo = OpenMeteoCacheRepository() class OpenMeteoCacheMixin: def _load_open_meteo_disk_cache(self) -> None: """启动时从磁盘加载 Open-Meteo 三类缓存,避免重启后冷启动打爆 API""" + mode = get_state_storage_mode() + if mode == STATE_STORAGE_SQLITE: + try: + saved = _open_meteo_cache_repo.load_payload(self._disk_cache_max_age_sec) + loaded = self._merge_open_meteo_payload(saved) + self._disk_cache_last_mtime = _open_meteo_cache_repo.latest_updated_at() + if loaded: + logger.info(f"✅ 从 SQLite 加载 Open-Meteo 缓存 {loaded} 条") + return + except Exception as exc: + logger.error(f"SQLite Open-Meteo 缓存加载失败,fallback file: {exc}") try: path = self._disk_cache_path if not os.path.exists(path): @@ -58,11 +77,26 @@ class OpenMeteoCacheMixin: self._disk_cache_last_mtime = current_mtime if loaded: logger.info(f"✅ 从磁盘加载 Open-Meteo 缓存 {loaded} 条 ({self._disk_cache_path})") + if mode == STATE_STORAGE_DUAL: + try: + _open_meteo_cache_repo.replace_payload(saved, self._disk_cache_max_age_sec) + except Exception as exc: + logger.warning(f"SQLite Open-Meteo 缓存同步失败: {exc}") except Exception as exc: logger.warning(f"磁盘缓存加载失败(首次启动不影响运行): {exc}") def _maybe_reload_open_meteo_disk_cache(self) -> None: """跨进程共享缓存:当缓存文件有更新时增量重载到当前进程内存""" + mode = get_state_storage_mode() + if mode == STATE_STORAGE_SQLITE: + try: + latest_updated_at = _open_meteo_cache_repo.latest_updated_at() + if latest_updated_at <= self._disk_cache_last_mtime: + return + self._load_open_meteo_disk_cache() + except Exception: + pass + return try: path = self._disk_cache_path if not os.path.exists(path): @@ -76,8 +110,8 @@ class OpenMeteoCacheMixin: def _flush_open_meteo_disk_cache(self) -> None: """将三类 Open-Meteo 内存缓存持久化到磁盘""" + mode = get_state_storage_mode() try: - os.makedirs(os.path.dirname(self._disk_cache_path), exist_ok=True) with self._open_meteo_cache_lock: forecast_snapshot = dict(self._open_meteo_cache) with self._ensemble_cache_lock: @@ -90,15 +124,50 @@ class OpenMeteoCacheMixin: "multi_model": multi_model_snapshot, "saved_at": time.time(), } - with self._disk_cache_lock: - tmp_path = self._disk_cache_path + ".tmp" - with open(tmp_path, "w", encoding="utf-8") as f: - json.dump(payload, f) - os.replace(tmp_path, self._disk_cache_path) - self._disk_cache_last_mtime = os.path.getmtime(self._disk_cache_path) + if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}: + _open_meteo_cache_repo.replace_payload(payload, self._disk_cache_max_age_sec) + self._disk_cache_last_mtime = _open_meteo_cache_repo.latest_updated_at() + if mode != STATE_STORAGE_SQLITE: + os.makedirs(os.path.dirname(self._disk_cache_path), exist_ok=True) + with self._disk_cache_lock: + tmp_path = self._disk_cache_path + ".tmp" + with open(tmp_path, "w", encoding="utf-8") as f: + json.dump(payload, f) + os.replace(tmp_path, self._disk_cache_path) + self._disk_cache_last_mtime = max( + self._disk_cache_last_mtime, + os.path.getmtime(self._disk_cache_path), + ) except Exception as exc: logger.warning(f"磁盘缓存写入失败: {exc}") + def _merge_open_meteo_payload(self, saved: dict) -> int: + now = time.time() + max_age = max(600, self._disk_cache_max_age_sec) + loaded = 0 + with self._open_meteo_cache_lock: + for key, entry in (saved.get("forecast", {}) or {}).items(): + if now - float(entry.get("t", 0)) < max_age: + old = self._open_meteo_cache.get(key) + if old is None or float(entry.get("t", 0)) >= float(old.get("t", 0)): + self._open_meteo_cache[key] = entry + loaded += 1 + with self._ensemble_cache_lock: + for key, entry in (saved.get("ensemble", {}) or {}).items(): + if now - float(entry.get("t", 0)) < max_age: + old = self._ensemble_cache.get(key) + if old is None or float(entry.get("t", 0)) >= float(old.get("t", 0)): + self._ensemble_cache[key] = entry + loaded += 1 + with self._multi_model_cache_lock: + for key, entry in (saved.get("multi_model", {}) or {}).items(): + if now - float(entry.get("t", 0)) < max_age: + old = self._multi_model_cache.get(key) + if old is None or float(entry.get("t", 0)) >= float(old.get("t", 0)): + self._multi_model_cache[key] = entry + loaded += 1 + return loaded + def _wait_open_meteo_slot(self, endpoint: str) -> None: """Simple per-process rate gate for Open-Meteo endpoints.""" min_interval = self._open_meteo_min_interval_sec diff --git a/src/database/runtime_state.py b/src/database/runtime_state.py new file mode 100644 index 00000000..ba4b5db2 --- /dev/null +++ b/src/database/runtime_state.py @@ -0,0 +1,500 @@ + +from __future__ import annotations + +import json +import os +import sqlite3 +import threading +import time +from pathlib import Path +from typing import Any, Dict, List, Optional + +from loguru import logger + +from src.database.db_manager import DBManager + +STATE_STORAGE_FILE = "file" +STATE_STORAGE_DUAL = "dual" +STATE_STORAGE_SQLITE = "sqlite" +VALID_STATE_STORAGE_MODES = { + STATE_STORAGE_FILE, + STATE_STORAGE_DUAL, + STATE_STORAGE_SQLITE, +} + +_LOGGED_MODES: set[str] = set() + + +def get_state_storage_mode() -> str: + raw = str(os.getenv("POLYWEATHER_STATE_STORAGE_MODE") or STATE_STORAGE_DUAL).strip().lower() + if raw not in VALID_STATE_STORAGE_MODES: + logger.warning( + f"invalid POLYWEATHER_STATE_STORAGE_MODE={raw!r}, fallback to {STATE_STORAGE_DUAL}" + ) + raw = STATE_STORAGE_DUAL + if raw not in _LOGGED_MODES: + logger.info(f"runtime state storage mode={raw}") + _LOGGED_MODES.add(raw) + return raw + + +class RuntimeStateDB: + _instance: Optional["RuntimeStateDB"] = None + _instance_lock = threading.Lock() + + def __init__(self, db_path: Optional[str] = None): + self.db_path = DBManager(db_path).db_path + self._init_tables() + + @classmethod + def instance(cls) -> "RuntimeStateDB": + with cls._instance_lock: + if cls._instance is None: + cls._instance = cls() + return cls._instance + + def connect(self) -> sqlite3.Connection: + conn = sqlite3.connect(self.db_path) + conn.row_factory = sqlite3.Row + return conn + + def _init_tables(self) -> None: + with self.connect() as conn: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS daily_records_store ( + city TEXT NOT NULL, + target_date TEXT NOT NULL, + actual_high REAL, + deb_prediction REAL, + mu REAL, + updated_at REAL NOT NULL, + payload_json TEXT NOT NULL, + PRIMARY KEY (city, target_date) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS telegram_alert_last_by_city ( + city TEXT PRIMARY KEY, + signature TEXT, + trigger_key TEXT, + severity TEXT, + ts INTEGER, + active INTEGER DEFAULT 0, + cleared_ts INTEGER, + evidence_json TEXT + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS telegram_alert_signature_state ( + signature TEXT PRIMARY KEY, + ts INTEGER NOT NULL + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS probability_training_snapshots_store ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + city TEXT NOT NULL, + target_date TEXT NOT NULL, + timestamp TEXT NOT NULL, + raw_mu REAL, + raw_sigma REAL, + max_so_far REAL, + peak_status TEXT, + probability_mode TEXT, + legacy_top_bucket INTEGER, + shadow_top_bucket INTEGER, + payload_json TEXT NOT NULL + ) + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_probability_snapshot_city_date ON probability_training_snapshots_store(city, target_date, id DESC)" + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS open_meteo_cache_store ( + source_kind TEXT NOT NULL, + cache_key TEXT NOT NULL, + updated_at REAL NOT NULL, + expires_at REAL, + payload_json TEXT NOT NULL, + PRIMARY KEY (source_kind, cache_key) + ) + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_open_meteo_cache_expires ON open_meteo_cache_store(source_kind, expires_at)" + ) + conn.commit() + + +class DailyRecordRepository: + def __init__(self, db: Optional[RuntimeStateDB] = None): + self.db = db or RuntimeStateDB.instance() + + def load_all(self) -> Dict[str, Dict[str, Dict[str, Any]]]: + out: Dict[str, Dict[str, Dict[str, Any]]] = {} + with self.db.connect() as conn: + rows = conn.execute( + "SELECT city, target_date, payload_json FROM daily_records_store ORDER BY city, target_date" + ).fetchall() + for row in rows: + try: + payload = json.loads(row["payload_json"]) + except Exception: + continue + city = str(row["city"]) + date_str = str(row["target_date"]) + out.setdefault(city, {})[date_str] = payload + return out + + def upsert_record(self, city: str, target_date: str, record: Dict[str, Any]) -> None: + payload_json = json.dumps(record, ensure_ascii=False) + updated_at = time.time() + with self.db.connect() as conn: + conn.execute( + """ + INSERT INTO daily_records_store ( + city, target_date, actual_high, deb_prediction, mu, updated_at, payload_json + ) VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(city, target_date) DO UPDATE SET + actual_high = excluded.actual_high, + deb_prediction = excluded.deb_prediction, + mu = excluded.mu, + updated_at = excluded.updated_at, + payload_json = excluded.payload_json + """, + ( + city, + target_date, + record.get("actual_high"), + record.get("deb_prediction"), + record.get("mu"), + updated_at, + payload_json, + ), + ) + conn.commit() + + def replace_all(self, data: Dict[str, Dict[str, Dict[str, Any]]]) -> int: + count = 0 + with self.db.connect() as conn: + conn.execute("DELETE FROM daily_records_store") + for city, city_rows in (data or {}).items(): + if not isinstance(city_rows, dict): + continue + for target_date, record in city_rows.items(): + payload_json = json.dumps(record, ensure_ascii=False) + conn.execute( + """ + INSERT INTO daily_records_store ( + city, target_date, actual_high, deb_prediction, mu, updated_at, payload_json + ) VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + city, + target_date, + record.get("actual_high"), + record.get("deb_prediction"), + record.get("mu"), + time.time(), + payload_json, + ), + ) + count += 1 + conn.commit() + return count + + def delete_older_than(self, cutoff_date: str) -> int: + with self.db.connect() as conn: + cur = conn.execute( + "DELETE FROM daily_records_store WHERE target_date < ?", + (cutoff_date,), + ) + conn.commit() + return int(cur.rowcount or 0) + + +class TelegramAlertStateRepository: + def __init__(self, db: Optional[RuntimeStateDB] = None): + self.db = db or RuntimeStateDB.instance() + + def load_state(self) -> Dict[str, Any]: + state = {"last_by_city": {}, "by_signature": {}} + with self.db.connect() as conn: + city_rows = conn.execute( + "SELECT city, signature, trigger_key, severity, ts, active, cleared_ts, evidence_json FROM telegram_alert_last_by_city" + ).fetchall() + sig_rows = conn.execute( + "SELECT signature, ts FROM telegram_alert_signature_state" + ).fetchall() + for row in city_rows: + entry = { + "signature": row["signature"], + "trigger_key": row["trigger_key"], + "severity": row["severity"], + "ts": row["ts"], + "active": bool(row["active"]), + } + if row["cleared_ts"] is not None: + entry["cleared_ts"] = row["cleared_ts"] + if row["evidence_json"]: + try: + entry["evidence"] = json.loads(row["evidence_json"]) + except Exception: + pass + state["last_by_city"][str(row["city"])] = entry + for row in sig_rows: + state["by_signature"][str(row["signature"])] = int(row["ts"] or 0) + return state + + def save_state(self, state: Dict[str, Any]) -> None: + last_by_city = state.get("last_by_city") or {} + by_signature = state.get("by_signature") or {} + with self.db.connect() as conn: + conn.execute("DELETE FROM telegram_alert_last_by_city") + conn.execute("DELETE FROM telegram_alert_signature_state") + for city, row in last_by_city.items(): + if not isinstance(row, dict): + continue + conn.execute( + """ + INSERT INTO telegram_alert_last_by_city ( + city, signature, trigger_key, severity, ts, active, cleared_ts, evidence_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + city, + row.get("signature"), + row.get("trigger_key"), + row.get("severity"), + int(row.get("ts") or 0), + 1 if row.get("active") else 0, + row.get("cleared_ts"), + json.dumps(row.get("evidence"), ensure_ascii=False) + if row.get("evidence") is not None + else None, + ), + ) + for signature, ts in by_signature.items(): + conn.execute( + "INSERT INTO telegram_alert_signature_state (signature, ts) VALUES (?, ?)", + (signature, int(ts or 0)), + ) + conn.commit() + + def replace_from_state(self, state: Dict[str, Any]) -> int: + self.save_state(state) + return len((state.get("last_by_city") or {})) + len((state.get("by_signature") or {})) + + +class ProbabilitySnapshotRepository: + def __init__(self, db: Optional[RuntimeStateDB] = None): + self.db = db or RuntimeStateDB.instance() + + def append_snapshot(self, payload: Dict[str, Any]) -> None: + legacy_top = _top_bucket(payload.get("prob_snapshot")) + shadow_top = _top_bucket(payload.get("shadow_prob_snapshot")) + with self.db.connect() as conn: + conn.execute( + """ + INSERT INTO probability_training_snapshots_store ( + city, target_date, timestamp, raw_mu, raw_sigma, max_so_far, + peak_status, probability_mode, legacy_top_bucket, shadow_top_bucket, payload_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + payload.get("city"), + payload.get("date"), + payload.get("timestamp"), + payload.get("raw_mu"), + payload.get("raw_sigma"), + payload.get("max_so_far"), + payload.get("peak_status"), + payload.get("probability_mode"), + legacy_top, + shadow_top, + json.dumps(payload, ensure_ascii=False), + ), + ) + conn.commit() + + def load_recent_rows(self, city: str, target_date: str, limit: int) -> List[Dict[str, Any]]: + with self.db.connect() as conn: + rows = conn.execute( + """ + SELECT payload_json + FROM probability_training_snapshots_store + WHERE city = ? AND target_date = ? + ORDER BY id DESC + LIMIT ? + """, + (city, target_date, int(limit)), + ).fetchall() + out = [] + for row in rows: + try: + out.append(json.loads(row["payload_json"])) + except Exception: + continue + return out + + def load_all_rows(self) -> List[Dict[str, Any]]: + with self.db.connect() as conn: + rows = conn.execute( + "SELECT payload_json FROM probability_training_snapshots_store ORDER BY id" + ).fetchall() + out = [] + for row in rows: + try: + out.append(json.loads(row["payload_json"])) + except Exception: + continue + return out + + def replace_all(self, rows: List[Dict[str, Any]]) -> int: + count = 0 + with self.db.connect() as conn: + conn.execute("DELETE FROM probability_training_snapshots_store") + for payload in rows or []: + if not isinstance(payload, dict): + continue + legacy_top = _top_bucket(payload.get("prob_snapshot")) + shadow_top = _top_bucket(payload.get("shadow_prob_snapshot")) + conn.execute( + """ + INSERT INTO probability_training_snapshots_store ( + city, target_date, timestamp, raw_mu, raw_sigma, max_so_far, + peak_status, probability_mode, legacy_top_bucket, shadow_top_bucket, payload_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + payload.get("city"), + payload.get("date"), + payload.get("timestamp"), + payload.get("raw_mu"), + payload.get("raw_sigma"), + payload.get("max_so_far"), + payload.get("peak_status"), + payload.get("probability_mode"), + legacy_top, + shadow_top, + json.dumps(payload, ensure_ascii=False), + ), + ) + count += 1 + conn.commit() + return count + + +class OpenMeteoCacheRepository: + def __init__(self, db: Optional[RuntimeStateDB] = None): + self.db = db or RuntimeStateDB.instance() + + def replace_payload(self, payload: Dict[str, Any], max_age: int) -> int: + count = 0 + now = time.time() + with self.db.connect() as conn: + conn.execute("DELETE FROM open_meteo_cache_store") + for source_kind in ("forecast", "ensemble", "multi_model"): + bucket = payload.get(source_kind) or {} + if not isinstance(bucket, dict): + continue + for cache_key, entry in bucket.items(): + if not isinstance(entry, dict): + continue + updated_at = float(entry.get("t") or now) + expires_at = updated_at + max_age + conn.execute( + """ + INSERT INTO open_meteo_cache_store ( + source_kind, cache_key, updated_at, expires_at, payload_json + ) VALUES (?, ?, ?, ?, ?) + """, + ( + source_kind, + cache_key, + updated_at, + expires_at, + json.dumps(entry, ensure_ascii=False), + ), + ) + count += 1 + conn.commit() + return count + + def load_payload(self, max_age: int) -> Dict[str, Any]: + now = time.time() + payload: Dict[str, Any] = { + "forecast": {}, + "ensemble": {}, + "multi_model": {}, + "saved_at": now, + } + with self.db.connect() as conn: + rows = conn.execute( + "SELECT source_kind, cache_key, updated_at, payload_json FROM open_meteo_cache_store" + ).fetchall() + for row in rows: + updated_at = float(row["updated_at"] or 0) + if now - updated_at >= max(600, max_age): + continue + try: + entry = json.loads(row["payload_json"]) + except Exception: + continue + payload.setdefault(str(row["source_kind"]), {})[str(row["cache_key"])] = entry + return payload + + def latest_updated_at(self) -> float: + with self.db.connect() as conn: + row = conn.execute( + "SELECT MAX(updated_at) AS max_updated_at FROM open_meteo_cache_store" + ).fetchone() + if not row: + return 0.0 + try: + return float(row["max_updated_at"] or 0.0) + except Exception: + return 0.0 + + +def _top_bucket(snapshot: Optional[List[Dict[str, Any]]]) -> Optional[int]: + best_value = None + best_prob = -1.0 + for row in snapshot or []: + if not isinstance(row, dict): + continue + value = row.get("v") + if value is None: + value = row.get("value") + try: + ivalue = int(value) + except Exception: + continue + prob = row.get("p") + if prob is None: + prob = row.get("probability") + try: + fprob = float(prob) + except Exception: + continue + if fprob > best_prob: + best_prob = fprob + best_value = ivalue + return best_value + + +def get_runtime_data_dir() -> str: + raw = str(os.getenv("POLYWEATHER_RUNTIME_DATA_DIR") or "").strip() + if raw: + return raw + project_root = Path(__file__).resolve().parents[2] + return str(project_root / "data") diff --git a/src/utils/metrics.py b/src/utils/metrics.py new file mode 100644 index 00000000..abf548ce --- /dev/null +++ b/src/utils/metrics.py @@ -0,0 +1,139 @@ +from __future__ import annotations + +import threading +from typing import Dict, Iterable, List, Optional, Tuple + + +LabelTuple = Tuple[Tuple[str, str], ...] + + +class _MetricsRegistry: + def __init__(self) -> None: + self._lock = threading.Lock() + self._counters: Dict[Tuple[str, LabelTuple], float] = {} + self._gauges: Dict[Tuple[str, LabelTuple], float] = {} + self._histograms: Dict[Tuple[str, LabelTuple], Dict[str, float]] = {} + + @staticmethod + def _normalize_labels(labels: Dict[str, object]) -> LabelTuple: + return tuple( + sorted((str(key), str(value)) for key, value in labels.items() if value is not None) + ) + + def inc_counter(self, name: str, amount: float = 1.0, **labels: object) -> None: + key = (name, self._normalize_labels(labels)) + with self._lock: + self._counters[key] = self._counters.get(key, 0.0) + amount + + def set_gauge(self, name: str, value: float, **labels: object) -> None: + key = (name, self._normalize_labels(labels)) + with self._lock: + self._gauges[key] = value + + def observe(self, name: str, value: float, **labels: object) -> None: + key = (name, self._normalize_labels(labels)) + with self._lock: + bucket = self._histograms.setdefault( + key, + {"count": 0.0, "sum": 0.0, "max": 0.0}, + ) + bucket["count"] += 1.0 + bucket["sum"] += value + bucket["max"] = max(bucket["max"], value) + + def snapshot(self) -> Dict[str, object]: + with self._lock: + return { + "counters": dict(self._counters), + "gauges": dict(self._gauges), + "histograms": { + key: dict(value) for key, value in self._histograms.items() + }, + } + + def export_prometheus(self) -> str: + snap = self.snapshot() + lines: List[str] = [] + for name, labels, value in _iter_metrics(snap["counters"]): + lines.append(_prom_line(name, value, labels)) + for name, labels, value in _iter_metrics(snap["gauges"]): + lines.append(_prom_line(name, value, labels)) + for (name, labels), stats in sorted(snap["histograms"].items()): + lines.append(_prom_line(f"{name}_count", stats["count"], labels)) + lines.append(_prom_line(f"{name}_sum", stats["sum"], labels)) + lines.append(_prom_line(f"{name}_max", stats["max"], labels)) + return "\n".join(lines) + ("\n" if lines else "") + + +def _iter_metrics(entries: Dict[Tuple[str, LabelTuple], float]) -> Iterable[Tuple[str, LabelTuple, float]]: + for (name, labels), value in sorted(entries.items()): + yield name, labels, value + + +def _prom_line(name: str, value: float, labels: LabelTuple) -> str: + if labels: + def _escape(label_value: str) -> str: + return label_value.replace("\\", "\\\\").replace('"', '\\"') + + label_str = ",".join( + f'{key}="{_escape(str(val))}"' + for key, val in labels + ) + return f"{name}{{{label_str}}} {value}" + return f"{name} {value}" + + +METRICS = _MetricsRegistry() + + +def counter_inc(name: str, amount: float = 1.0, **labels: object) -> None: + METRICS.inc_counter(name, amount=amount, **labels) + + +def gauge_set(name: str, value: float, **labels: object) -> None: + METRICS.set_gauge(name, value=value, **labels) + + +def histogram_observe(name: str, value: float, **labels: object) -> None: + METRICS.observe(name, value=value, **labels) + + +def record_source_call(source: str, operation: str, outcome: str, duration_ms: Optional[float] = None) -> None: + counter_inc( + "polyweather_source_requests_total", + source=source, + operation=operation, + outcome=outcome, + ) + if duration_ms is not None: + histogram_observe( + "polyweather_source_request_duration_ms", + duration_ms, + source=source, + operation=operation, + outcome=outcome, + ) + + +def build_metrics_summary() -> Dict[str, object]: + snapshot = METRICS.snapshot() + request_total = 0.0 + source_total = 0.0 + source_errors = 0.0 + for (name, labels), value in snapshot["counters"].items(): + if name == "polyweather_http_requests_total": + request_total += value + if name == "polyweather_source_requests_total": + source_total += value + label_map = dict(labels) + if label_map.get("outcome") not in {"success", "cache_hit"}: + source_errors += value + return { + "http_requests_total": int(request_total), + "source_requests_total": int(source_total), + "source_error_total": int(source_errors), + } + + +def export_prometheus_metrics() -> str: + return METRICS.export_prometheus() diff --git a/src/utils/telegram_push.py b/src/utils/telegram_push.py index 1699a309..83d0726e 100644 --- a/src/utils/telegram_push.py +++ b/src/utils/telegram_push.py @@ -9,6 +9,12 @@ from typing import Any, Dict, List, Optional, Tuple from loguru import logger +from src.database.runtime_state import ( + STATE_STORAGE_DUAL, + STATE_STORAGE_SQLITE, + TelegramAlertStateRepository, + get_state_storage_mode, +) from src.data_collection.city_registry import CITY_REGISTRY from src.utils.telegram_chat_ids import get_telegram_chat_ids_from_env @@ -19,6 +25,7 @@ SEVERITY_RANK = { "medium": 2, "high": 3, } +_telegram_state_repo = TelegramAlertStateRepository() def _env_bool(name: str, default: bool) -> bool: @@ -184,7 +191,18 @@ def _state_file() -> str: def _load_state(path: str) -> Dict[str, Any]: + mode = get_state_storage_mode() + if mode == STATE_STORAGE_SQLITE: + try: + return _telegram_state_repo.load_state() + except Exception as exc: + logger.error(f"failed to load telegram push state from sqlite: {exc}") if not os.path.exists(path): + if mode == STATE_STORAGE_DUAL: + try: + return _telegram_state_repo.load_state() + except Exception: + return {"last_by_city": {}, "by_signature": {}} return {"last_by_city": {}, "by_signature": {}} try: with open(path, "r", encoding="utf-8") as fh: @@ -199,6 +217,11 @@ def _load_state(path: str) -> Dict[str, Any]: def _save_state(path: str, state: Dict[str, Any]) -> None: + mode = get_state_storage_mode() + if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}: + _telegram_state_repo.save_state(state) + if mode == STATE_STORAGE_SQLITE: + return os.makedirs(os.path.dirname(path), exist_ok=True) tmp_path = f"{path}.tmp" with open(tmp_path, "w", encoding="utf-8") as fh: diff --git a/tests/test_probability_rollout.py b/tests/test_probability_rollout.py new file mode 100644 index 00000000..73de299b --- /dev/null +++ b/tests/test_probability_rollout.py @@ -0,0 +1,65 @@ +from src.analysis.probability_rollout import judge_probability_rollout + + +def test_judge_probability_rollout_holds_on_shadow_brier_regression(): + evaluation_report = { + "summary": { + "sample_count": 105, + "delta": { + "crps": -0.09, + "mae": 0.0, + "bucket_hit_rate": 0.0, + }, + } + } + shadow_report = { + "summary": { + "samples": 103, + "delta_mae": 0.01, + "delta_bucket_hit_rate": 0.01, + "delta_bucket_brier": 0.29, + }, + "by_city": { + "miami": { + "samples": 4, + "delta_mae": 0.24, + "delta_bucket_hit_rate": -0.5, + "delta_bucket_brier": 0.47, + } + }, + } + + payload = judge_probability_rollout(evaluation_report, shadow_report) + + assert payload["decision"] == "hold" + assert payload["ready_for_primary"] is False + assert payload["blocking_reasons"] + assert payload["worst_shadow_regressions"][0]["city"] == "miami" + + +def test_judge_probability_rollout_promotes_on_clean_metrics(): + evaluation_report = { + "summary": { + "sample_count": 120, + "delta": { + "crps": -0.08, + "mae": 0.0, + "bucket_hit_rate": 0.02, + }, + } + } + shadow_report = { + "summary": { + "samples": 110, + "delta_mae": 0.0, + "delta_bucket_hit_rate": 0.01, + "delta_bucket_brier": 0.01, + }, + "by_city": {}, + } + + payload = judge_probability_rollout(evaluation_report, shadow_report) + + assert payload["decision"] == "promote" + assert payload["ready_for_primary"] is True + assert payload["blocking_reasons"] == [] diff --git a/tests/test_probability_snapshot_archive.py b/tests/test_probability_snapshot_archive.py index 9be1d041..d45c452b 100644 --- a/tests/test_probability_snapshot_archive.py +++ b/tests/test_probability_snapshot_archive.py @@ -1,9 +1,15 @@ import json from pathlib import Path +import pytest from src.analysis.probability_snapshot_archive import append_probability_snapshot +@pytest.fixture(autouse=True) +def _force_file_mode(monkeypatch): + monkeypatch.setenv("POLYWEATHER_STATE_STORAGE_MODE", "file") + + def test_append_probability_snapshot_writes_jsonl(tmp_path: Path): archive_path = tmp_path / "probability_training_snapshots.jsonl" diff --git a/tests/test_runtime_state_storage.py b/tests/test_runtime_state_storage.py new file mode 100644 index 00000000..8a469371 --- /dev/null +++ b/tests/test_runtime_state_storage.py @@ -0,0 +1,74 @@ + +import time + +from src.database.runtime_state import ( + DailyRecordRepository, + OpenMeteoCacheRepository, + ProbabilitySnapshotRepository, + RuntimeStateDB, + TelegramAlertStateRepository, +) + + +def test_daily_record_repository_roundtrip(tmp_path, monkeypatch): + monkeypatch.setenv('POLYWEATHER_DB_PATH', str(tmp_path / 'polyweather.db')) + repo = DailyRecordRepository(RuntimeStateDB(str(tmp_path / 'polyweather.db'))) + repo.upsert_record('ankara', '2026-03-20', {'actual_high': 15.2, 'deb_prediction': 14.8, 'mu': 15.0}) + data = repo.load_all() + assert data['ankara']['2026-03-20']['actual_high'] == 15.2 + + +def test_telegram_alert_state_repository_roundtrip(tmp_path, monkeypatch): + monkeypatch.setenv('POLYWEATHER_DB_PATH', str(tmp_path / 'polyweather.db')) + repo = TelegramAlertStateRepository(RuntimeStateDB(str(tmp_path / 'polyweather.db'))) + state = { + 'last_by_city': { + 'ankara': { + 'signature': 'sig-1', + 'trigger_key': 'mkt:test', + 'severity': 'medium', + 'ts': 123, + 'active': True, + 'evidence': {'x': 1}, + } + }, + 'by_signature': {'sig-1': 123}, + } + repo.save_state(state) + loaded = repo.load_state() + assert loaded == state + + +def test_probability_snapshot_repository_recent_rows(tmp_path, monkeypatch): + monkeypatch.setenv('POLYWEATHER_DB_PATH', str(tmp_path / 'polyweather.db')) + repo = ProbabilitySnapshotRepository(RuntimeStateDB(str(tmp_path / 'polyweather.db'))) + repo.append_snapshot({ + 'city': 'ankara', + 'date': '2026-03-20', + 'timestamp': '2026-03-20T12:00:00Z', + 'raw_mu': 15.2, + 'raw_sigma': 1.1, + 'max_so_far': 14.9, + 'peak_status': 'before', + 'probability_mode': 'emos_shadow', + 'prob_snapshot': [{'v': 15, 'p': 0.6}], + 'shadow_prob_snapshot': [{'v': 15, 'p': 0.4}], + }) + rows = repo.load_recent_rows('ankara', '2026-03-20', 5) + assert len(rows) == 1 + assert rows[0]['raw_mu'] == 15.2 + + +def test_open_meteo_cache_repository_roundtrip(tmp_path, monkeypatch): + monkeypatch.setenv('POLYWEATHER_DB_PATH', str(tmp_path / 'polyweather.db')) + repo = OpenMeteoCacheRepository(RuntimeStateDB(str(tmp_path / 'polyweather.db'))) + payload = { + 'forecast': {'ankara': {'t': time.time(), 'temp': 15}}, + 'ensemble': {'ankara': {'t': time.time(), 'spread': 1.5}}, + 'multi_model': {}, + 'saved_at': 1000, + } + repo.replace_payload(payload, 86400) + loaded = repo.load_payload(86400) + assert loaded['forecast']['ankara']['temp'] == 15 + assert loaded['ensemble']['ankara']['spread'] == 1.5 diff --git a/tests/test_web_observability.py b/tests/test_web_observability.py new file mode 100644 index 00000000..b6837884 --- /dev/null +++ b/tests/test_web_observability.py @@ -0,0 +1,37 @@ + +from fastapi.testclient import TestClient + +from web.app import app + + +client = TestClient(app) + + +def test_healthz_returns_ok_shape(): + response = client.get('/healthz') + assert response.status_code == 200 + payload = response.json() + assert payload['status'] in {'ok', 'degraded'} + assert 'db' in payload + assert 'state_storage_mode' in payload + assert 'cities_count' in payload + + +def test_system_status_returns_summary_shape(): + response = client.get('/api/system/status') + assert response.status_code == 200 + payload = response.json() + assert 'db' in payload + assert 'features' in payload + assert 'integrations' in payload + assert 'cache' in payload + assert 'probability' in payload + assert 'rollout' in payload['probability'] + assert payload['probability']['rollout']['decision']['decision'] in {'hold', 'observe', 'promote'} + assert 'cities_count' in payload + + +def test_metrics_endpoint_returns_prometheus_payload(): + response = client.get('/metrics') + assert response.status_code == 200 + assert 'polyweather_http_requests_total' in response.text diff --git a/web/core.py b/web/core.py index 09eab165..93f084c6 100644 --- a/web/core.py +++ b/web/core.py @@ -3,6 +3,9 @@ PolyWeather Web Core Context """ import os +import sqlite3 +import time +from datetime import datetime, timezone from typing import Dict, Any, Optional from fastapi import FastAPI, HTTPException, Request @@ -16,7 +19,19 @@ from src.data_collection.weather_sources import WeatherDataCollector from src.data_collection.city_risk_profiles import CITY_RISK_PROFILES # noqa: F401 from src.data_collection.polymarket_readonly import PolymarketReadOnlyLayer from src.auth.supabase_entitlement import SUPABASE_ENTITLEMENT, extract_bearer_token +from src.utils.metrics import ( + build_metrics_summary, + counter_inc, + gauge_set, + histogram_observe, +) +from src.analysis.probability_calibration import ( + DEFAULT_CALIBRATION_FILE, + resolve_probability_engine_mode, +) +from src.analysis.probability_rollout import build_rollout_report from src.database.db_manager import DBManager +from src.database.runtime_state import get_state_storage_mode from src.payments import PAYMENT_CHECKOUT, PaymentCheckoutError # noqa: F401 from src.data_collection.city_registry import CITY_REGISTRY @@ -65,6 +80,59 @@ SETTLEMENT_SOURCE_LABELS: Dict[str, str] = { _cache: Dict[str, Dict] = {} CACHE_TTL = 300 CACHE_TTL_ANKARA = 60 +_PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +_PROBABILITY_EVALUATION_REPORT = os.path.join( + _PROJECT_ROOT, + "artifacts", + "probability_calibration", + "evaluation_report.json", +) + + +@app.middleware("http") +async def _metrics_middleware(request: Request, call_next): + started = time.perf_counter() + try: + response = await call_next(request) + except Exception: + duration_ms = (time.perf_counter() - started) * 1000.0 + counter_inc( + "polyweather_http_requests_total", + method=request.method, + path=request.url.path, + status="500", + ) + histogram_observe( + "polyweather_http_request_duration_ms", + duration_ms, + method=request.method, + path=request.url.path, + status="500", + ) + raise + + duration_ms = (time.perf_counter() - started) * 1000.0 + status_code = str(response.status_code) + counter_inc( + "polyweather_http_requests_total", + method=request.method, + path=request.url.path, + status=status_code, + ) + histogram_observe( + "polyweather_http_request_duration_ms", + duration_ms, + method=request.method, + path=request.url.path, + status=status_code, + ) + return response +_PROBABILITY_SHADOW_REPORT = os.path.join( + _PROJECT_ROOT, + "artifacts", + "probability_calibration", + "shadow_report.json", +) def _env_bool(name: str, default: bool = False) -> bool: @@ -264,3 +332,97 @@ def _sf(v) -> Optional[float]: def _is_excluded_model_name(model_name: str) -> bool: normalized = str(model_name or "").strip().lower().replace(" ", "").replace("_", "").replace("-", "") return "meteoblue" in normalized + + +def _sqlite_health() -> Dict[str, Any]: + try: + with sqlite3.connect(_account_db.db_path) as conn: + conn.execute("SELECT 1").fetchone() + return {"ok": True, "db_path": _account_db.db_path} + except Exception as exc: + return {"ok": False, "db_path": _account_db.db_path, "error": str(exc)} + + +def _cache_summary() -> Dict[str, Any]: + gauge_set("polyweather_api_cache_entries", len(_cache)) + gauge_set( + "polyweather_open_meteo_forecast_cache_entries", + len(getattr(_weather, "_open_meteo_cache", {}) or {}), + ) + gauge_set( + "polyweather_open_meteo_ensemble_cache_entries", + len(getattr(_weather, "_ensemble_cache", {}) or {}), + ) + gauge_set( + "polyweather_open_meteo_multi_model_cache_entries", + len(getattr(_weather, "_multi_model_cache", {}) or {}), + ) + return { + "api_cache_entries": len(_cache), + "open_meteo_forecast_entries": len(getattr(_weather, "_open_meteo_cache", {}) or {}), + "open_meteo_ensemble_entries": len(getattr(_weather, "_ensemble_cache", {}) or {}), + "open_meteo_multi_model_entries": len(getattr(_weather, "_multi_model_cache", {}) or {}), + } + + +def _feature_flags_summary() -> Dict[str, Any]: + return { + "auth_enabled": bool(SUPABASE_ENTITLEMENT.enabled), + "auth_required": bool(_SUPABASE_AUTH_REQUIRED), + "entitlement_guard_enabled": bool(_ENTITLEMENT_GUARD_ENABLED), + "payment_enabled": bool(getattr(PAYMENT_CHECKOUT, "enabled", False)), + "state_storage_mode": get_state_storage_mode(), + } + + +def _integration_summary() -> Dict[str, Any]: + weather_cfg = _config.get("weather", {}) if isinstance(_config, dict) else {} + return { + "supabase_configured": bool(SUPABASE_ENTITLEMENT.configured), + "meteoblue_configured": bool(os.getenv("METEOBLUE_API_KEY")), + "telegram_bot_configured": bool((_config.get("telegram", {}) or {}).get("bot_token")), + "walletconnect_configured": bool(os.getenv("NEXT_PUBLIC_WALLETCONNECT_PROJECT_ID")), + "weather_sources": { + "openweather": bool(weather_cfg.get("openweather_api_key")), + "wunderground": bool(weather_cfg.get("wunderground_api_key")), + "visualcrossing": bool(weather_cfg.get("visualcrossing_api_key")), + }, + } + + +def _probability_summary() -> Dict[str, Any]: + rollout = build_rollout_report( + _PROBABILITY_EVALUATION_REPORT, + _PROBABILITY_SHADOW_REPORT, + ) + return { + "engine_mode": resolve_probability_engine_mode(), + "calibration_file": os.getenv("POLYWEATHER_PROBABILITY_CALIBRATION_FILE") + or DEFAULT_CALIBRATION_FILE, + "rollout": rollout, + } + + +def build_health_payload() -> Dict[str, Any]: + db = _sqlite_health() + return { + "status": "ok" if db.get("ok") else "degraded", + "time_utc": datetime.now(timezone.utc).isoformat(), + "db": db, + "state_storage_mode": get_state_storage_mode(), + "cities_count": len(CITIES), + } + + +def build_system_status_payload() -> Dict[str, Any]: + return { + "status": build_health_payload()["status"], + "time_utc": datetime.now(timezone.utc).isoformat(), + "db": _sqlite_health(), + "features": _feature_flags_summary(), + "integrations": _integration_summary(), + "cache": _cache_summary(), + "metrics": build_metrics_summary(), + "probability": _probability_summary(), + "cities_count": len(CITIES), + } diff --git a/web/routes.py b/web/routes.py index 9b78035a..0a9883f2 100644 --- a/web/routes.py +++ b/web/routes.py @@ -4,10 +4,12 @@ import os from typing import Optional from fastapi import APIRouter, HTTPException, Request +from fastapi.responses import PlainTextResponse from loguru import logger from src.analysis.deb_algorithm import load_history from src.data_collection.city_registry import ALIASES +from src.utils.metrics import export_prometheus_metrics from web.analysis_service import ( _analyze, _build_city_detail_payload, @@ -31,6 +33,8 @@ from web.core import ( _SUPABASE_AUTH_REQUIRED, _assert_entitlement, _bind_optional_supabase_identity, + build_health_payload, + build_system_status_payload, _require_supabase_identity, _resolve_auth_points, _resolve_weekly_profile, @@ -49,6 +53,27 @@ def _normalize_city_or_404(name: str) -> str: return city +@router.get("/healthz") +async def healthz(): + payload = build_health_payload() + if payload.get("status") != "ok": + raise HTTPException(status_code=503, detail=payload) + return payload + + +@router.get("/api/system/status") +async def system_status(): + return build_system_status_payload() + + +@router.get("/metrics", response_class=PlainTextResponse) +async def metrics(): + return PlainTextResponse( + export_prometheus_metrics(), + media_type="text/plain; version=0.0.4; charset=utf-8", + ) + + @router.get("/api/cities") async def list_cities(request: Request): _assert_entitlement(request)