Add detailed daily refresh diagnostics

This commit is contained in:
PrimeThai Automation
2026-09-30 13:54:55 +07:00
parent 2e09c56f4d
commit 2b8ea9dff4
2 changed files with 106 additions and 19 deletions

View File

@@ -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.

View File

@@ -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)