From 2b8ea9dff47f7d4dbf7b53ff6f839c89c798bdfa Mon Sep 17 00:00:00 2001 From: PrimeThai Automation Date: Wed, 30 Sep 2026 13:54:55 +0700 Subject: [PATCH] Add detailed daily refresh diagnostics --- README.md | 2 + primethai/webui.py | 123 ++++++++++++++++++++++++++++++++++++++------- 2 files changed, 106 insertions(+), 19 deletions(-) diff --git a/README.md b/README.md index da8cf6d..69fd9da 100644 --- a/README.md +++ b/README.md @@ -275,3 +275,5 @@ docker compose up --build -d ``` Compose uses named volumes for `/app/data` and `/app/reports`. On first creation, Docker seeds them from the files included in the image; later restarts and image rebuilds keep the server's updated price archive, portfolio database and reports. Use `docker compose down` to stop the app while keeping its data. The port is bound to localhost so an HTTPS reverse proxy can sit in front of it. + +For automatic refresh diagnostics, view the container log with `docker compose logs -f primethai`. Failed-source entries include the provider, failed stage, exception type and message; retry entries summarize SiamChart, Yahoo Finance, archive reading and signal calculation status. diff --git a/primethai/webui.py b/primethai/webui.py index 0fd8566..4ce1b21 100644 --- a/primethai/webui.py +++ b/primethai/webui.py @@ -231,6 +231,15 @@ def _today_local() -> str: return datetime.now(THAILAND_TZ).date().isoformat() +def _daily_log(event: str, **details: object) -> None: + payload = json.dumps(details, ensure_ascii=False, sort_keys=True, default=str) + print(f"[daily-refresh] {event} {payload}", flush=True) + + +def _selected_state(state: dict[str, object], fields: tuple[str, ...]) -> dict[str, object]: + return {key: state[key] for key in fields if key in state} + + def _last_fetch_attempt_date() -> str | None: state = _read_json(LATEST_FETCH_STATE) value = state.get("last_attempt_date") if state else None @@ -263,8 +272,8 @@ def _fetch_latest_archive(destination: Path) -> str: except (TypeError, ValueError): pass - destination.parent.mkdir(parents=True, exist_ok=True) temporary = destination.with_name(destination.name + ".download") + stage = "prepare_local_storage" state: dict[str, object] = { "last_attempt_date": today, "attempted_at": now.isoformat(), @@ -272,19 +281,26 @@ def _fetch_latest_archive(destination: Path) -> str: "url": LATEST_ARCHIVE_URL, } try: + destination.parent.mkdir(parents=True, exist_ok=True) + stage = "request_siamchart" request = Request( LATEST_ARCHIVE_URL, headers={"User-Agent": "PrimeThai/0.1 daily EOD refresh"}, ) with urlopen(request, timeout=30) as response: + stage = "read_siamchart_response" payload = response.read(MAX_LATEST_ARCHIVE_BYTES + 1) + stage = "validate_download_size" if not payload or len(payload) > MAX_LATEST_ARCHIVE_BYTES: raise ValueError("ไฟล์ที่ดาวน์โหลดมีขนาดไม่ถูกต้อง") + stage = "write_temporary_archive" temporary.write_bytes(payload) incoming: dict[str, bytes] = {} + stage = "validate_zip_archive" with zipfile.ZipFile(temporary) as bundle: if bundle.testzip() is not None: raise ValueError("ไฟล์ ZIP ตรวจสอบความสมบูรณ์ไม่ผ่าน") + stage = "extract_eod_sessions" for name in bundle.namelist(): match = _EOD_DATE.search(Path(name).name) if match: @@ -292,8 +308,10 @@ def _fetch_latest_archive(destination: Path) -> str: if not incoming: raise ValueError("ZIP ไม่มีไฟล์ SiamChart EOD รายวัน") downloaded_date = max(pd.Timestamp(day) for day in incoming) + stage = "read_existing_archive" existing_summary = archive_summary(destination) existing_date = pd.Timestamp(str(existing_summary["last_date"])) if existing_summary.get("last_date") else None + stage = "merge_daily_sessions" stored = merge_daily_entries( destination, incoming, @@ -314,15 +332,25 @@ def _fetch_latest_archive(destination: Path) -> str: return (f"ดึง SiamChart EOD สำเร็จ · เก็บข้อมูลย้อนหลัง {stored['sessions']} วัน " f"ถึง {stored_date.date()}") except Exception as exc: - state.update({"result": "failed", "error": f"{type(exc).__name__}: {exc}"}) + state.update({"result": "failed", "failed_stage": stage, + "error_type": type(exc).__name__, "error": str(exc)}) + _daily_log( + "source_failed", source="SiamChart", date=today, after_close=after_close, + stage=stage, url=LATEST_ARCHIVE_URL, + error_type=type(exc).__name__, error=str(exc), + ) return "เชื่อมต่อ SiamChart ไม่สำเร็จ; คงข้อมูลเดิมไว้และจะลองใหม่อัตโนมัติ" finally: try: temporary.unlink(missing_ok=True) except OSError: pass - LATEST_FETCH_STATE.parent.mkdir(parents=True, exist_ok=True) - LATEST_FETCH_STATE.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8") + try: + LATEST_FETCH_STATE.parent.mkdir(parents=True, exist_ok=True) + LATEST_FETCH_STATE.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8") + except Exception as exc: + _daily_log("state_write_failed", source="SiamChart", path=str(LATEST_FETCH_STATE), + error_type=type(exc).__name__, error=str(exc)) def _required_backup_symbols(archives: tuple[Path, ...], report_dir: Path) -> list[str]: @@ -371,16 +399,22 @@ def _fetch_yahoo_backup(archives: tuple[Path, ...], report_dir: Path, "attempted_at": now.isoformat(), "provider": "Yahoo Finance via yfinance", } + stage = "read_latest_local_archive" try: latest = _latest_archive_date(archives).date().isoformat() + stage = "determine_required_symbols" + required_symbols = _required_backup_symbols(archives, report_dir) + state["required_symbols"] = len(required_symbols) + stage = "fetch_yahoo_history_and_check_coverage" entries, metadata = fetch_yahoo_eod_backup( [str(path) for path in archives], latest_date=latest, end_date=now.date(), - required_symbols=_required_backup_symbols(archives, report_dir), + required_symbols=required_symbols, ) if not entries: raise YahooBackupError("แหล่งสำรองไม่มี session ใหม่") + stage = "merge_yahoo_sessions" stored = merge_daily_entries( destination, entries, @@ -407,28 +441,43 @@ def _fetch_yahoo_backup(archives: tuple[Path, ...], report_dir: Path, message = str(exc) no_new_session = "ยังไม่มีข้อมูล SET หลังวันที่เก็บไว้ล่าสุด" in message state.update({"result": "no_new_session" if no_new_session else "failed", - "error": f"{type(exc).__name__}: {exc}"}) + "failed_stage": stage, "error_type": type(exc).__name__, + "error": message}) + _daily_log( + "source_no_new_session" if no_new_session else "source_failed", + source="Yahoo Finance via yfinance", date=today, stage=stage, + error_type=type(exc).__name__, error=message, + required_symbols=state.get("required_symbols"), + ) if no_new_session: return "ยังไม่พบ session ใหม่จาก SiamChart หรือ Yahoo Finance; คงข้อมูลล่าสุดที่บันทึกไว้" return "SiamChart และ Yahoo Finance ยังใช้ไม่ได้; คงข้อมูลราคาที่บันทึกไว้และจะลองใหม่อัตโนมัติ" finally: - _write_fetch_state(BACKUP_FETCH_STATE, state) + try: + _write_fetch_state(BACKUP_FETCH_STATE, state) + except Exception as exc: + _daily_log("state_write_failed", source="Yahoo Finance via yfinance", + path=str(BACKUP_FETCH_STATE), + error_type=type(exc).__name__, error=str(exc)) def _refresh_worker(archives: tuple[Path, ...], report_dir: Path, start: str, initial_equity: float, fetch_latest: bool = True) -> None: + stage = "inspect_local_archive" try: destination = _daily_archive_destination(archives) archive_mtime_before = (destination.stat().st_mtime_ns if destination is not None and destination.is_file() else None) if fetch_latest and destination is not None: + stage = "refresh_price_sources" download_status = _fetch_latest_archive(destination) primary_state = _read_json(LATEST_FETCH_STATE) or {} now = datetime.now(THAILAND_TZ) + today = now.date().isoformat() scheduled_at = now.replace(hour=AUTO_REFRESH_HOUR, minute=AUTO_REFRESH_MINUTE, second=0, microsecond=0) primary_has_current_close = ( - primary_state.get("result") == "downloaded" + primary_state.get("result") in {"downloaded", "source_older_than_local"} and bool(primary_state.get("after_close")) and str(primary_state.get("stored_eod_date", "")) >= today ) @@ -443,6 +492,7 @@ def _refresh_worker(archives: tuple[Path, ...], report_dir: Path, start: str, download_status = "ชุดไฟล์นี้ไม่มี ZIP EOD รายวันสำหรับอัปเดต; คำนวณจากไฟล์ที่กำหนด" else: download_status = "คำนวณจากไฟล์ EOD ในเครื่อง" + stage = "read_latest_eod_date" latest = _latest_archive_date(archives) archive_mtime_after = (destination.stat().st_mtime_ns if destination is not None and destination.is_file() else None) @@ -450,10 +500,14 @@ def _refresh_worker(archives: tuple[Path, ...], report_dir: Path, start: str, if archive_mtime_before == archive_mtime_after and reports_exist: with _STATE_LOCK: _STATE.update({"running": False, "completed_at": datetime.now(timezone.utc).isoformat(), - "error": None, "data_as_of": latest.date().isoformat(), + "error": None, "failed_stage": None, + "data_as_of": latest.date().isoformat(), "download_status": download_status}) + _daily_log("refresh_completed", data_as_of=latest.date().isoformat(), + reports_rebuilt=False, download_status=download_status) return end = (latest + pd.Timedelta(days=1)).date().isoformat() + stage = "run_model_lab" run_model_lab( [str(path) for path in archives], start=start, end=end, output_dir=report_dir, action_dir=Path("data/raw/actions"), @@ -463,12 +517,16 @@ def _refresh_worker(archives: tuple[Path, ...], report_dir: Path, start: str, ) with _STATE_LOCK: _STATE.update({"running": False, "completed_at": datetime.now(timezone.utc).isoformat(), - "error": None, "data_as_of": latest.date().isoformat(), + "error": None, "failed_stage": None, + "data_as_of": latest.date().isoformat(), "download_status": download_status}) + _daily_log("refresh_completed", data_as_of=latest.date().isoformat(), + reports_rebuilt=True, download_status=download_status) except Exception as exc: # The UI reports local data or calculation errors to the user. with _STATE_LOCK: _STATE.update({"running": False, "completed_at": datetime.now(timezone.utc).isoformat(), - "error": f"{type(exc).__name__}: {exc}"}) + "error": f"{type(exc).__name__}: {exc}", "failed_stage": stage}) + _daily_log("refresh_failed", stage=stage, error_type=type(exc).__name__, error=str(exc)) def _start_refresh(archives: tuple[Path, ...], report_dir: Path, start: str, @@ -477,7 +535,7 @@ def _start_refresh(archives: tuple[Path, ...], report_dir: Path, start: str, if bool(_STATE["running"]): return False, "กำลังคำนวณข้อมูลชุดนี้อยู่แล้ว" _STATE.update({"running": True, "started_at": datetime.now(timezone.utc).isoformat(), - "completed_at": None, "error": None}) + "completed_at": None, "error": None, "failed_stage": None}) worker = threading.Thread( target=_refresh_worker, args=(archives, report_dir, start, initial_equity, fetch_latest), @@ -519,10 +577,11 @@ def _daily_refresh_loop(archives: tuple[Path, ...], report_dir: Path, start: str if (previous.get("last_attempt_date") == run_date.isoformat() and bool(previous.get("after_close")) and previous.get("result") != "failed"): continue + retry_number = 0 while datetime.now(THAILAND_TZ).date() == run_date: started, message = _start_refresh(archives, report_dir, start, initial_equity) if not started: - print("[daily-refresh] another refresh is running; waiting") + _daily_log("waiting_for_active_refresh", date=run_date.isoformat(), reason=message) while True: with _STATE_LOCK: running = bool(_STATE["running"]) @@ -535,7 +594,7 @@ def _daily_refresh_loop(archives: tuple[Path, ...], report_dir: Path, start: str and current_fetch.get("result") != "failed"): break continue - print("[daily-refresh] scheduled refresh started") + _daily_log("refresh_started", date=run_date.isoformat()) while True: with _STATE_LOCK: running = bool(_STATE["running"]) @@ -546,16 +605,24 @@ def _daily_refresh_loop(archives: tuple[Path, ...], report_dir: Path, start: str backup_state = _read_json(BACKUP_FETCH_STATE) or {} with _STATE_LOCK: refresh_error = _STATE["error"] + refresh_failed_stage = _STATE.get("failed_stage") + try: + latest_archive_date = _latest_archive_date(archives).date().isoformat() + latest_archive_error = None + except Exception as exc: + latest_archive_date = None + latest_archive_error = {"type": type(exc).__name__, "error": str(exc)} primary_ok = ( - fetch_state.get("result") == "downloaded" + fetch_state.get("result") in {"downloaded", "source_older_than_local"} and bool(fetch_state.get("after_close")) and str(fetch_state.get("stored_eod_date", "")) >= run_date.isoformat() ) backup_ok = ( - backup_state.get("result") == "downloaded" + latest_archive_date is not None + and backup_state.get("result") == "downloaded" and backup_state.get("last_attempt_date") == run_date.isoformat() and str(backup_state.get("downloaded_through", "")) - >= _latest_archive_date(archives).date().isoformat() + >= latest_archive_date ) both_sources_confirm_no_new_session = ( backup_state.get("result") == "no_new_session" @@ -563,12 +630,30 @@ def _daily_refresh_loop(archives: tuple[Path, ...], report_dir: Path, start: str and fetch_state.get("result") == "downloaded" and bool(fetch_state.get("after_close")) ) - if (primary_ok or backup_ok or both_sources_confirm_no_new_session) and not refresh_error: + if ((primary_ok or backup_ok or both_sources_confirm_no_new_session) + and not refresh_error and latest_archive_error is None): break retry_at = datetime.now(THAILAND_TZ) + timedelta(minutes=FETCH_RETRY_MINUTES) if retry_at.date() != run_date: break - print(f"[daily-refresh] source update failed; retrying around {retry_at:%H:%M} BKK") + retry_number += 1 + _daily_log( + "retry_scheduled", attempt=retry_number, run_date=run_date.isoformat(), + retry_at=retry_at.isoformat(timespec="seconds"), + primary=_selected_state(fetch_state, ( + "result", "last_attempt_date", "attempted_at", "after_close", "url", + "failed_stage", "error_type", "error", "downloaded_eod_date", "stored_eod_date", + )), + backup=_selected_state(backup_state, ( + "result", "last_attempt_date", "attempted_at", "provider", "failed_stage", + "error_type", "error", "downloaded_from", "downloaded_through", + "liquid_universe_coverage", "liquid_universe_present", "liquid_universe_expected", + "required_symbols", + )), + calculation={"stage": refresh_failed_stage, "error": refresh_error}, + local_archive_date=latest_archive_date, + local_archive_error=latest_archive_error, + ) time.sleep(FETCH_RETRY_MINUTES * 60)