424 lines
20 KiB
Python
424 lines
20 KiB
Python
"""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
|
|
from integrations.ohip import acquisition_progress
|
|
|
|
|
|
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 pending_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 []
|
|
|
|
def pending_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()
|
|
|
|
def _acquisition_progress(self, row):
|
|
observation = None
|
|
getter = getattr(self._executor, "get_acquisition_progress", None)
|
|
if callable(getter):
|
|
try:
|
|
value = getter(request_id=row["request_id"], report_date=row["report_date"])
|
|
if value is not None:
|
|
observation = acquisition_progress.validate(value)
|
|
except Exception:
|
|
pass
|
|
# Acquisition completion remains true when later report processing fails.
|
|
# Other observations cannot override the queue's restart/retry state.
|
|
if observation and observation["phase"] == "completed":
|
|
return observation
|
|
status = row["status"]
|
|
if status == "queued":
|
|
return acquisition_progress.snapshot("queued", updated_at=row["updated_at"])
|
|
if status in {"failed", "interrupted"}:
|
|
return acquisition_progress.snapshot(status, observation["completed"] if observation else 0,
|
|
observation["total"] if observation else None, updated_at=row["updated_at"])
|
|
if status == "downloading" and (observation is None
|
|
or datetime.fromisoformat(observation["updated_at"]) < datetime.fromisoformat(row["updated_at"])):
|
|
return acquisition_progress.snapshot("searching", updated_at=row["updated_at"])
|
|
return observation
|
|
|
|
def _public(self, 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"],
|
|
"acquisition_progress": self._acquisition_progress(row),
|
|
}
|
|
|
|
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 pending_reviews(self) -> list[dict]:
|
|
# The portal reconciles price-review jobs against authoritative state;
|
|
# queue entries retain their acquisition history after review completion.
|
|
with self._connect() as db:
|
|
rows = db.execute("""SELECT * FROM downloads
|
|
WHERE status IN ('needs_data_review','needs_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:
|
|
self._recover_source_reviews()
|
|
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 _recover_source_reviews(self):
|
|
"""Recover saved review continuations at startup, without browser reads.
|
|
|
|
Individual corrupt/conflicting reviews stay pending. A verified freeze
|
|
survives a crash before this queue update and is reused on next startup.
|
|
"""
|
|
recover = getattr(self._executor, "recover_data_review", None)
|
|
if not callable(recover):
|
|
return
|
|
with self._connect() as db:
|
|
rows = db.execute("""SELECT request_id FROM downloads WHERE status='needs_data_review'
|
|
ORDER BY created_at, request_id""").fetchall()
|
|
for row in rows:
|
|
with self._condition:
|
|
if self._stopping:
|
|
return
|
|
try:
|
|
if recover(row["request_id"]) is not True:
|
|
continue
|
|
with self._condition:
|
|
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(), row["request_id"]))
|
|
self._condition.notify_all()
|
|
except Exception:
|
|
# Recovery is local, optional maintenance, not evidence of a
|
|
# failed acquisition or permission to skip an unresolved review.
|
|
continue
|
|
|
|
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.
|