"""Single-day ARR download jobs; the accepted source executor is injected. This module does not implement or certify an OHIP-to-report adapter. The default Web runtime remains unavailable until an accepted executor is explicitly wired. The SQLite queue has one process owner on one host, not a distributed lease. """ from __future__ import annotations import fcntl import hashlib import json import os import re import sqlite3 import threading from contextlib import contextmanager from dataclasses import dataclass from datetime import date, datetime, timedelta, timezone from pathlib import Path from typing import Callable, Protocol from zoneinfo import ZoneInfo from arr_web.contracts import PortalError, validate_job_id ACTIVE = frozenset({"queued", "downloading", "processing"}) MAX_ATTEMPTS = 10 def validate_report_date(value: object) -> str: if isinstance(value, str) and re.fullmatch(r"[0-9]{4}-[0-9]{2}-[0-9]{2}", value): try: if date.fromisoformat(value).isoformat() == value: return value except ValueError: pass raise PortalError("ARR_DOWNLOAD_DATE_INVALID", "请选择一个有效的报表日期") def validate_request_id(value: object) -> str: if not isinstance(value, str) or not re.fullmatch(r"[0-9a-f]{32}", value): raise PortalError("ARR_DOWNLOAD_REQUEST_INVALID", "下载任务编号无效") return value def default_report_date() -> str: return (datetime.now(ZoneInfo("Asia/Bangkok")).date() - timedelta(days=1)).isoformat() def _now() -> str: return datetime.now(timezone.utc).isoformat() @dataclass(frozen=True) class DownloadOutcome: # succeeded means an independently validated Finance commit, not just capture. status: str job_id: str | None = None retryable: bool = False def validate(self) -> None: if self.status not in {"succeeded", "needs_review", "needs_data_review", "failed"}: raise ValueError("invalid download outcome") if self.job_id is not None: validate_job_id(self.job_id) if self.status == "needs_data_review" and (self.job_id is not None or self.retryable): raise ValueError("source review precedes a processing job") if self.status in {"succeeded", "needs_review"} and (not self.job_id or self.retryable): raise ValueError("a processing outcome requires a stable job identity") if type(self.retryable) is not bool: raise ValueError("invalid retryability") class ARRDownloadExecutor(Protocol): def execute( self, *, request_id: str, from_date: date, to_date: date, report_stage: Callable[[str], None], ) -> DownloadOutcome: """Replay the SAME source/package/Finance identity after an uncertain call. Both dates are equal. Publish downloading/processing only. Return success only after authoritative ingestion acknowledgement; return needs_review for a persisted review. A retryable failure must be safe for exact replay. No credentials, raw response, guest fields or arbitrary messages cross here. """ ... class ARRDownloadCoordinator(Protocol): @property def ready(self) -> bool: ... def latest(self) -> dict | None: ... def pending_data_reviews(self) -> list[dict]: ... def create(self, report_date: str, request_id: str) -> dict: ... def get(self, request_id: str) -> dict: ... def retry(self, request_id: str) -> dict: ... def get_data_review(self, request_id: str) -> dict: ... def update_data_review_item(self, request_id, item_id, revision, value, actor) -> dict: ... def finalize_data_review(self, request_id, revision, actor) -> dict: ... class UnavailableARRDownloads: ready = False def latest(self) -> None: return None def pending_data_reviews(self) -> list[dict]: return [] @staticmethod def _unavailable() -> PortalError: return PortalError("ARR_DOWNLOAD_UNAVAILABLE", "自动下载服务暂未就绪", 503) def create(self, report_date: str, request_id: str) -> dict: raise self._unavailable() def get(self, request_id: str) -> dict: raise self._unavailable() def retry(self, request_id: str) -> dict: raise self._unavailable() def get_data_review(self, *args): raise self._unavailable() def update_data_review_item(self, *args): raise self._unavailable() def finalize_data_review(self, *args): raise self._unavailable() class PersistentARRDownloads: """Persist intent before dispatch; response loss never allocates another job.""" def __init__(self, jobs_root: Path, executor: ARRDownloadExecutor) -> None: jobs_root.mkdir(parents=True, exist_ok=True, mode=0o700) if jobs_root.is_symlink(): raise ValueError("download job directory must not be a symlink") os.chmod(jobs_root, 0o700) self._database = jobs_root / "arr-downloads.sqlite3" lock_path = jobs_root / "arr-downloads.lock" if self._database.is_symlink() or lock_path.is_symlink(): raise ValueError("download state must not be a symlink") self._lease = os.open(lock_path, os.O_CREAT | os.O_RDWR, 0o600) try: fcntl.flock(self._lease, fcntl.LOCK_EX | fcntl.LOCK_NB) except BaseException: os.close(self._lease) raise self._executor = executor identity = getattr(executor, "source_identity", None) source_context = identity() if callable(identity) else {"kind": type(executor).__name__} self.context_id = hashlib.sha256(json.dumps({"root": str(jobs_root.resolve()), "source": source_context}, sort_keys=True, separators=(",", ":"), allow_nan=False).encode()).hexdigest() self._condition = threading.Condition(threading.RLock()) self._stopping = False try: with self._connect() as db: db.execute("""CREATE TABLE IF NOT EXISTS downloads ( request_id TEXT PRIMARY KEY, report_date TEXT NOT NULL, status TEXT NOT NULL, job_id TEXT, retryable INTEGER NOT NULL DEFAULT 0, attempts INTEGER NOT NULL DEFAULT 0, created_at TEXT NOT NULL, updated_at TEXT NOT NULL)""") # A process exit is not evidence that a Finance commit failed. db.execute("""UPDATE downloads SET status='interrupted', retryable=1, updated_at=? WHERE status IN ('downloading', 'processing')""", (_now(),)) os.chmod(self._database, 0o600) self._thread = threading.Thread(target=self._work, name="arr-download-worker", daemon=True) self._thread.start() except BaseException: os.close(self._lease) raise @contextmanager def _connect(self): db = sqlite3.connect(self._database, timeout=10) db.row_factory = sqlite3.Row try: db.execute("PRAGMA synchronous=FULL") with db: yield db finally: db.close() @property def ready(self) -> bool: return not self._stopping and self._thread.is_alive() @staticmethod def _public(row: sqlite3.Row) -> dict: return { "request_id": row["request_id"], "report_date": row["report_date"], "from_date": row["report_date"], "to_date": row["report_date"], "status": row["status"], "job_id": row["job_id"], "error_code": ("ARR_SOURCE_FETCH_FAILED" if row["status"] == "failed" else "ARR_DOWNLOAD_INTERRUPTED" if row["status"] == "interrupted" else None) if not row["job_id"] else None, "can_retry": bool(row["retryable"]) and row["attempts"] < MAX_ATTEMPTS, "created_at": row["created_at"], "updated_at": row["updated_at"], } def latest(self) -> dict | None: with self._connect() as db: row = db.execute("SELECT * FROM downloads ORDER BY created_at DESC, rowid DESC LIMIT 1").fetchone() return self._public(row) if row else None def pending_data_reviews(self) -> list[dict]: with self._connect() as db: rows = db.execute("""SELECT * FROM downloads WHERE status='needs_data_review' ORDER BY report_date DESC, created_at, request_id""").fetchall() return [self._public(row) for row in rows] def get(self, request_id: str) -> dict: validate_request_id(request_id) with self._connect() as db: row = db.execute("SELECT * FROM downloads WHERE request_id=?", (request_id,)).fetchone() if row is None: raise PortalError("ARR_DOWNLOAD_NOT_FOUND", "下载任务不存在", 404) # An explicit finalization may have frozen successfully just before the # process/response was lost. Reconciliation resumes that same intent. review_getter = getattr(self._executor, "get_data_review", None) if row["status"] == "needs_data_review" and callable(review_getter) and self.ready: with self._condition: if review_getter(request_id)["status"] == "finalized": with self._connect() as db: db.execute("UPDATE downloads SET status='queued',retryable=0,updated_at=? WHERE request_id=? AND status='needs_data_review'", (_now(), request_id)) row = db.execute("SELECT * FROM downloads WHERE request_id=?", (request_id,)).fetchone() self._condition.notify_all() return self._public(row) def create(self, report_date: str, request_id: str) -> dict: validate_report_date(report_date) validate_request_id(request_id) with self._condition: if not self.ready: raise UnavailableARRDownloads._unavailable() with self._connect() as db: db.execute("BEGIN IMMEDIATE") row = db.execute("SELECT * FROM downloads WHERE request_id=?", (request_id,)).fetchone() if row: if row["report_date"] != report_date: raise PortalError("ARR_DOWNLOAD_DATE_CONFLICT", "原任务日期不能修改", 409) return self._public(row) # Separate tabs cannot start another acquisition for an unresolved day. row = db.execute("""SELECT * FROM downloads WHERE report_date=? AND status IN ('queued','downloading','processing','interrupted','needs_data_review') ORDER BY created_at DESC LIMIT 1""", (report_date,)).fetchone() if row: return self._public(row) now = _now() db.execute("INSERT INTO downloads (request_id,report_date,status,created_at,updated_at) VALUES (?,?,?,?,?)", (request_id, report_date, "queued", now, now)) self._condition.notify_all() return self.get(request_id) def _review_method(self, name): method = getattr(self._executor, name, None) if not callable(method): raise PortalError("ARR_DATA_REVIEW_NOT_FOUND", "该任务没有待完善的数据", 404) return method def get_data_review(self, request_id): self.get(request_id) return self._review_method("get_data_review")(request_id) def update_data_review_item(self, request_id, item_id, revision, value, actor): with self._condition: if self.get(request_id)["status"] != "needs_data_review": raise PortalError("ARR_DATA_REVIEW_FROZEN", "该任务当前不能修改字段", 409) return self._review_method("update_data_review_item")(request_id, item_id, revision, value, actor) def finalize_data_review(self, request_id, revision, actor): with self._condition: if not self.ready: raise UnavailableARRDownloads._unavailable() task = self.get(request_id) review = self.get_data_review(request_id) if task["status"] != "needs_data_review": if review["status"] == "finalized": return task # Reconcile a lost finalization response without another run. raise PortalError("ARR_DATA_REVIEW_FROZEN", "该任务当前不能确认生成", 409) self._review_method("finalize_data_review")(request_id, revision, actor) with self._connect() as db: db.execute("UPDATE downloads SET status='queued',retryable=0,updated_at=? WHERE request_id=? AND status='needs_data_review'", (_now(), request_id)) self._condition.notify_all() return self.get(request_id) def retry(self, request_id: str) -> dict: validate_request_id(request_id) with self._condition: if not self.ready: raise UnavailableARRDownloads._unavailable() task = self.get(request_id) if task["status"] in ACTIVE or task["status"] in {"succeeded", "needs_review", "needs_data_review"}: return task if not task["can_retry"]: raise PortalError("ARR_DOWNLOAD_RETRY_UNAVAILABLE", "该任务不能继续重试,请查看任务日志", 409) with self._connect() as db: other = db.execute("""SELECT request_id FROM downloads WHERE report_date=? AND request_id<>? AND status IN ('queued','downloading','processing','interrupted')""", (task["report_date"], request_id)).fetchone() if other: raise PortalError("ARR_DOWNLOAD_OTHER_ACTIVE", "该日期已有进行中的下载任务", 409) db.execute("UPDATE downloads SET status='queued',retryable=0,updated_at=? WHERE request_id=?", (_now(), request_id)) self._condition.notify_all() return self.get(request_id) def _stage(self, request_id: str, stage: str) -> None: if stage not in {"downloading", "processing"}: raise ValueError("invalid download stage") with self._connect() as db: db.execute("""UPDATE downloads SET status=?,updated_at=? WHERE request_id=? AND status IN ('downloading','processing')""", (stage, _now(), request_id)) def _work(self) -> None: try: while True: with self._condition: if self._stopping: return with self._connect() as db: row = db.execute("SELECT * FROM downloads WHERE status='queued' ORDER BY created_at,rowid LIMIT 1").fetchone() if row: db.execute("""UPDATE downloads SET status='downloading',attempts=attempts+1, updated_at=? WHERE request_id=?""", (_now(), row["request_id"])) if row is None: self._condition.wait() continue request_id = row["request_id"] day = date.fromisoformat(row["report_date"]) try: outcome = self._executor.execute( request_id=request_id, from_date=day, to_date=day, report_stage=lambda stage: self._stage(request_id, stage), ) outcome.validate() status, job_id, retryable = outcome.status, outcome.job_id, outcome.retryable except Exception: # Keep the exact identity for reconciliation; never expose exception text. status, job_id, retryable = "interrupted", row["job_id"], True with self._connect() as db: db.execute("UPDATE downloads SET status=?,job_id=?,retryable=?,updated_at=? WHERE request_id=?", (status, job_id, int(retryable), _now(), request_id)) finally: os.close(self._lease) def close(self, *, wait: bool = False) -> None: """Stop dispatch; owners must wait before closing executor dependencies.""" with self._condition: self._stopping = True self._condition.notify_all() self._thread.join(timeout=None if wait else 2) # The worker retains its lease if an executor is still running.