[verified] Auto-refresh dated dividend ledger in the data scheduler
Automatically keep the real dated dividend ledger fresh inside the app's own refresh loop (this app runs on its own server, independent of Hermes): - backend/app/scheduler.py: AppDataScheduler gained a cooldown-gated _maybe_refresh_dated_dividends() that fetches real dated dividend history (siamchart /stock-info) into data/dividends/ledger.json at most once per dividend_cooldown_seconds (default 6h) — dividend history changes only a few times a year, so we never hammer the source every refresh tick. The fetch is non-fatal: a network failure leaves the previous ledger intact. - backend/app/__init__.py: passes DIVIDEND_REFRESH_COOLDOWN_SECONDS to the scheduler (default 21600s). - tests: cooldown fires once then skips, and refetches after it elapses (2) — full backend 294 passed.
This commit is contained in:
@@ -159,7 +159,10 @@ def create_app(config: dict[str, Any] | None = None) -> Flask:
|
||||
app.extensions["daily_cache"] = cache
|
||||
if not app.config.get("TESTING"):
|
||||
interval_s = int(os.getenv("REFRESH_INTERVAL_SECONDS", "3600"))
|
||||
scheduler = AppDataScheduler(cache, Path(__file__).resolve().parents[1] / "data", interval_seconds=interval_s)
|
||||
div_cooldown_s = int(os.getenv("DIVIDEND_REFRESH_COOLDOWN_SECONDS", str(6 * 3600)))
|
||||
scheduler = AppDataScheduler(cache, Path(__file__).resolve().parents[1] / "data",
|
||||
interval_seconds=interval_s,
|
||||
dividend_cooldown_seconds=div_cooldown_s)
|
||||
scheduler.start()
|
||||
app.extensions["data_scheduler"] = scheduler
|
||||
|
||||
|
||||
@@ -41,14 +41,18 @@ _REFRESH_JOBS: List[dict] = [
|
||||
class AppDataScheduler:
|
||||
"""Runs periodic refresh of the real data collectors inside the app."""
|
||||
|
||||
def __init__(self, cache, data_root: Path, interval_seconds: int = 3600):
|
||||
def __init__(self, cache, data_root: Path, interval_seconds: int = 3600,
|
||||
dividend_cooldown_seconds: int = 6 * 3600):
|
||||
self.cache = cache
|
||||
self.data_root = data_root
|
||||
self.interval = max(30, int(interval_seconds)) # never faster than 30s
|
||||
self.dividend_cooldown = max(300, int(dividend_cooldown_seconds)) # >=5min
|
||||
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)
|
||||
self._div_marker = self._snap_dir / "dividends_last_refresh.json"
|
||||
self._lock = threading.Lock()
|
||||
|
||||
# -- public lifecycle ---------------------------------------------------
|
||||
def start(self) -> None:
|
||||
@@ -84,6 +88,7 @@ class AppDataScheduler:
|
||||
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._maybe_refresh_dated_dividends()
|
||||
self._write_marker(results)
|
||||
return results
|
||||
|
||||
@@ -101,6 +106,54 @@ class AppDataScheduler:
|
||||
except Exception: # noqa: BLE001 — never let history break the refresh loop
|
||||
log.exception("set50 factor-history record failed (non-fatal)")
|
||||
|
||||
def _maybe_refresh_dated_dividends(self) -> None:
|
||||
"""Cooldown-gated fetch of real dated dividend history into the ledger.
|
||||
|
||||
Dividend history changes slowly (a few times a year per name), so we
|
||||
only re-fetch at most once per ``dividend_cooldown`` (default 6h), not
|
||||
every refresh tick. Non-fatal: a network failure leaves the previous
|
||||
ledger intact and just postpones the next update.
|
||||
"""
|
||||
with self._lock:
|
||||
if self._dividends_fresh():
|
||||
return
|
||||
try:
|
||||
from .dividend_ledger import DividendLedger, populate_dated_dividends
|
||||
from .siamchart import fetch_dividend_history
|
||||
from .siamchart_factors import build_factor_view
|
||||
ledger = DividendLedger(self.data_root / "dividends" / "ledger.json")
|
||||
fv = build_factor_view()
|
||||
symbols = [f["symbol"] for f in fv.get("factors", [])]
|
||||
populate_dated_dividends(ledger, symbols, fetch_dividend_history)
|
||||
ledger.save()
|
||||
self._mark_dividends_refreshed()
|
||||
log.info("set50 dated-dividend ledger refreshed (%d symbols)", len(symbols))
|
||||
except Exception: # noqa: BLE001 — non-fatal
|
||||
log.exception("set50 dated-dividend refresh failed (will retry next cooldown)")
|
||||
|
||||
def _dividends_fresh(self) -> bool:
|
||||
import json
|
||||
import datetime as _dt
|
||||
if not self._div_marker.is_file():
|
||||
return False
|
||||
try:
|
||||
payload = json.loads(self._div_marker.read_text(encoding="utf-8"))
|
||||
last = _dt.datetime.fromisoformat(payload["at"])
|
||||
return (_dt.datetime.now().astimezone() - last).total_seconds() < self.dividend_cooldown
|
||||
except (OSError, ValueError, KeyError):
|
||||
return False
|
||||
|
||||
def _mark_dividends_refreshed(self) -> None:
|
||||
import json
|
||||
import datetime as _dt
|
||||
try:
|
||||
self._div_marker.write_text(
|
||||
json.dumps({"at": _dt.datetime.now().astimezone().isoformat(timespec="seconds")}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
except OSError:
|
||||
log.exception("could not write dividend refresh marker")
|
||||
|
||||
def _run_job(self, job: dict) -> dict:
|
||||
key = job["key"]
|
||||
label = job["label"]
|
||||
|
||||
@@ -47,5 +47,54 @@ class SchedulerTest(unittest.TestCase):
|
||||
self.assertIn("at", data)
|
||||
|
||||
|
||||
class DividendCooldownTest(unittest.TestCase):
|
||||
def _sched(self, snap_dir, cooldown):
|
||||
return AppDataScheduler(_FakeCache(), snap_dir, interval_seconds=99999,
|
||||
dividend_cooldown_seconds=cooldown)
|
||||
|
||||
def test_refreshes_once_then_respects_cooldown(self):
|
||||
snap_dir = Path(tempfile.mkdtemp())
|
||||
sched = self._sched(snap_dir, cooldown=3600) # 1h cooldown
|
||||
calls = {"n": 0}
|
||||
|
||||
def fake_fetch(sym):
|
||||
calls["n"] += 1
|
||||
return [("2025-01-01", 1.0)]
|
||||
|
||||
fake_fv = {"factors": [{"symbol": "PTT"}]}
|
||||
with patch("app.siamchart_factors.build_factor_view", return_value=fake_fv), \
|
||||
patch("app.siamchart.fetch_dividend_history", side_effect=fake_fetch):
|
||||
sched._maybe_refresh_dated_dividends() # first -> fetch
|
||||
n_after_first = calls["n"]
|
||||
sched._maybe_refresh_dated_dividends() # second, within cooldown -> skip
|
||||
self.assertEqual(n_after_first, 1) # fetched once
|
||||
self.assertEqual(calls["n"], 1) # second call skipped
|
||||
self.assertTrue((snap_dir / "dividends" / "ledger.json").exists())
|
||||
|
||||
def test_refetches_after_cooldown_elapses(self):
|
||||
snap_dir = Path(tempfile.mkdtemp())
|
||||
sched = self._sched(snap_dir, cooldown=1) # 1s cooldown
|
||||
calls = {"n": 0}
|
||||
|
||||
def fake_fetch(sym):
|
||||
calls["n"] += 1
|
||||
return [("2025-01-01", 1.0)]
|
||||
|
||||
fake_fv = {"factors": [{"symbol": "PTT"}]}
|
||||
with patch("app.siamchart_factors.build_factor_view", return_value=fake_fv), \
|
||||
patch("app.siamchart.fetch_dividend_history", side_effect=fake_fetch):
|
||||
sched._maybe_refresh_dated_dividends()
|
||||
self.assertEqual(calls["n"], 1)
|
||||
# force cooldown to have elapsed by rewriting the marker back in time
|
||||
import json
|
||||
import datetime as _dt
|
||||
sched._div_marker.write_text(
|
||||
json.dumps({"at": (_dt.datetime.now().astimezone() - _dt.timedelta(hours=1)).isoformat()}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
sched._maybe_refresh_dated_dividends()
|
||||
self.assertEqual(calls["n"], 2) # refetched
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
Reference in New Issue
Block a user