Files
set50-system/backend/app/scheduler.py
Kunthawat Greethong d87a1ada39 [verified] Cross-theme surprise normalization + historical factor store (P4 enabler)
A. Cross-theme comparability:
- compute_theme_surprises now weight-normalizes by total |weight| (weighted
  average), so every theme surprise on same [-1,1] scale regardless of factor
  count/weight (retail 0.189->0.145; auto_credit 1.0->0.64).

B. Historical factor store (enables learning macro/demographic factors):
- New factor_history.py: append-only per-factor JSONL, dedupes unchanged
  values, rejects non-finite, records every FACTORS value each scheduler run.
- scheduler.py: jobs carry fetch_module; refresh_all records factor history
  (non-fatal); added bank_npl job.
- GET /api/v1/learning/factors?min_points= reports n_points/learnable per
  factor so users see when P4 learning unlocks (validated query parsing).
- weight_learning: generic learn_factor_series() aggregator (momentum reuses).

Independent review deleg_5dd358e3 passed=true (empty security/logic arrays);
its two robustness suggestions applied (finite guard in record(), clean 400 on
bad min_points). 234 tests pass; Vite build passes.
2026-08-27 07:32:16 +07:00

160 lines
7.2 KiB
Python

"""App-internal background scheduler — refreshes real data on its own schedule.
This app runs on its own server (Docker/EasyPanel), independent of Hermes. So
the periodic data-fetching/analysis must live INSIDE the app. This module runs a
daemon thread inside the Flask process that:
- on an interval (default every 60 min), triggers each collector through the
shared daily cache (so fresh real data lands as the "current" state),
- records each successful refresh as a timestamped snapshot under the data dir
(persistent history the dashboard/sources table can report),
- logs progress; swallows transient failures (a source being down once must
not crash the app or the loop).
Thread-safety: the loop only calls collectors + writes cache/snapshots; the HTTP
layer reads the daily cache. No shared mutable in-process state beyond that.
"""
from __future__ import annotations
import logging
import threading
import time
from pathlib import Path
from typing import Callable, List, Optional
log = logging.getLogger("set50.scheduler")
# collectors returning a .to_dict()/dict, keyed by cache key
# (imported lazily to avoid import cycles at module load)
# `fetch_module` = the FACTORS.fetch module name this job feeds (for history).
_REFRESH_JOBS: List[dict] = [
{"key": "bot_tourism", "label": "ท่องเที่ยว (BOT)", "module": "bot_tourism", "fn": "BotTourismSource().fetch", "fetch_module": "macro_thai"},
{"key": "auto_credit/tourism", "label": "ยอดขายรถ (TradingEconomics)", "module": "auto_credit", "fn": "fetch_auto_credit", "fetch_module": "auto_credit"},
{"key": "auto_npl", "label": "NPL รถยนต์ (BOT)", "module": "auto_npl", "fn": "fetch_auto_npl", "fetch_module": "auto_npl"},
{"key": "energy_thai", "label": "โรงกลั่น TOP", "module": "energy_thai", "fn": "fetch_energy_thai", "fetch_module": "energy_thai"},
{"key": "macro_thai", "label": "ภาพรวมประเทศไทย (BOT)", "module": "macro_thai", "fn": "fetch_macro_thai", "fetch_module": "macro_thai"},
{"key": "bank_npl", "label": "NPL ภาคการเงิน (BOT)", "module": "bank_npl", "fn": "fetch_bank_npl", "fetch_module": "bank_npl"},
]
class AppDataScheduler:
"""Runs periodic refresh of the real data collectors inside the app."""
def __init__(self, cache, data_root: Path, interval_seconds: int = 3600):
self.cache = cache
self.data_root = data_root
self.interval = max(30, int(interval_seconds)) # never faster than 30s
self._stop = threading.Event()
self._thread: Optional[threading.Thread] = None
self._snap_dir = data_root / "scheduler"
self._snap_dir.mkdir(parents=True, exist_ok=True)
# -- public lifecycle ---------------------------------------------------
def start(self) -> None:
if self._thread and self._thread.is_alive():
return
self._stop.clear()
self._thread = threading.Thread(target=self._loop, name="set50-refresh", daemon=True)
self._thread.start()
log.info("set50 data scheduler started (interval=%ss)", self.interval)
def stop(self) -> None:
self._stop.set()
# -- internals ----------------------------------------------------------
def _loop(self) -> None:
# refresh once shortly after boot (so fresh data is live), then on interval
self.refresh_all()
while not self._stop.wait(self.interval):
try:
self.refresh_all()
except Exception: # noqa: BLE001 — loop must survive individual failures
log.exception("set50 scheduler refresh_all failed (will retry)")
def refresh_all(self) -> list[dict]:
"""Run every collector, warm the daily cache, snapshot the state, and
append each factor's value to the historical store (P4 enabler)."""
results: list[dict] = []
fetched_by_module: dict[str, dict] = {}
for job in _REFRESH_JOBS:
res = self._run_job(job)
results.append(res)
fetch_module = job.get("fetch_module")
if res.get("ok") and isinstance(res.get("value"), dict) and fetch_module:
fetched_by_module[fetch_module] = res["value"]
self._record_history(fetched_by_module)
self._write_marker(results)
return results
def _record_history(self, fetched_by_module: dict[str, dict]) -> None:
"""Append current factor values to the historical store (append-only).
`fetched_by_module` maps FACTORS.fetch module name -> collector dict.
Runs after each refresh so vintages accumulate; macro/demographic
factors become learnable (P4) once they have enough history points.
"""
try:
from .factor_history import FactorHistory
fh = FactorHistory(self.data_root / "factor_history")
fh.record_all(fetched_by_module)
except Exception: # noqa: BLE001 — never let history break the refresh loop
log.exception("set50 factor-history record failed (non-fatal)")
def _run_job(self, job: dict) -> dict:
key = job["key"]
label = job["label"]
try:
value = self.cache.fetch_or_stale(key, self._make_fetcher(job))
return {"key": key, "label": label, "ok": True, "at": self._now(),
"value": value}
except Exception as exc: # noqa: BLE001
log.warning("set50 refresh failed for %s: %s", key, exc)
return {"key": key, "label": label, "ok": False, "error": str(exc), "at": self._now()}
def _make_fetcher(self, job: dict) -> Callable[[], dict]:
module_name = job["module"]
fn = job["fn"]
def fetcher() -> dict:
mod = __import__(f"app.{module_name}", fromlist=["*"])
obj = mod
# Support "ClassName().method" (instantiate then call) and plain
# "func_name" and static object attribute chains.
parts = fn.split(".")
for i, part in enumerate(parts):
is_instantiate = part.endswith("()")
attr = part[:-2] if is_instantiate else part
obj = getattr(obj, attr)
if is_instantiate:
obj = obj() # construct instance
if callable(obj):
value = obj()
else:
value = obj
if hasattr(value, "to_dict"):
value = value.to_dict()
return value or {}
return fetcher
def _write_marker(self, results: list[dict]) -> None:
import json
import datetime as _dt
ok = [r for r in results if r["ok"]]
marker = {
"at": self._now(),
"ok": len(ok),
"total": len(results),
"sources": results,
}
path = self._snap_dir / "last_refresh.json"
try:
path.write_text(json.dumps(marker, ensure_ascii=False, indent=2), encoding="utf-8")
except OSError:
log.exception("could not write scheduler marker")
@staticmethod
def _now() -> str:
import datetime as _dt
return _dt.datetime.now(_dt.timezone.utc).isoformat(timespec="seconds")