feat: Prepared for production #1

Merged
duuuuuuuden merged 2 commits from prod-preparation into main 2026-04-23 22:25:06 +02:00
2 changed files with 58 additions and 7 deletions
Showing only changes of commit 6b69b8c002 - Show all commits

View File

@@ -45,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90
_PAGE_NEXT_TIMEOUT_S = 60 _PAGE_NEXT_TIMEOUT_S = 60
# Ротация контекста выключена по умолчанию. # Ротация контекста выключена по умолчанию.
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0) _CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT = int(os.environ.get("IAAI_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT", "3") or 3)
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120 _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2 _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
@@ -594,6 +595,7 @@ class IAAIScraper:
failures: list[dict[str, str]] = [] failures: list[dict[str, str]] = []
early_stopped = False early_stopped = False
truncated_by_time_budget = False truncated_by_time_budget = False
listing_recovery_attempts = 0
known_origin_ids: set[str] | None = None known_origin_ids: set[str] | None = None
threshold = self.settings.listing.early_stop_threshold threshold = self.settings.listing.early_stop_threshold
@@ -725,6 +727,21 @@ class IAAIScraper:
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)}) failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
return True return True
def _consume_listing_recovery_attempt(*, page_number: int, reason: str) -> None:
nonlocal listing_recovery_attempts
listing_recovery_attempts += 1
logger.warning(
"Listing recovery attempt %d/%d at page %d (%s)",
listing_recovery_attempts,
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT,
page_number,
reason,
)
if listing_recovery_attempts > _MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT:
raise ListingResumeError(
f"Listing recovery exhausted at page {page_number}: {reason}"
)
page = None page = None
try: try:
page, applied_filters = self._open_listing_for_stream( page, applied_filters = self._open_listing_for_stream(
@@ -763,6 +780,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason="context_rotation",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -794,6 +815,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason=f"collect_failed:{type(op_exc).__name__}",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -839,6 +864,10 @@ class IAAIScraper:
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
page.close() page.close()
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason="empty_page_after_retry",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -1069,6 +1098,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number + 1,
reason=f"next_page_failed:{type(nav_exc).__name__}",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number + 1, target_page_number=page_number + 1,
make=make, make=make,

View File

@@ -12,6 +12,8 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal")
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress" SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress"
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM)
def _env_bool(name: str, default: bool) -> bool: def _env_bool(name: str, default: bool) -> bool:
@@ -76,6 +78,22 @@ def _read_last_progress_ts(redis_client: Redis) -> int | None:
return max_ts or None return max_ts or None
def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict[str, int]]:
"""Есть ли признаки активной/зависшей работы, даже если очередь пуста."""
queue_len = _safe_int(redis_client.llen(queue_name), 0)
has_lock = 1 if redis_client.get(SYNC_LISTING_LOCK_KEY) else 0
has_task_progress = 0
for _ in redis_client.scan_iter(match="iaai:state:task_progress:*"):
has_task_progress = 1
break
flags = {
"queue_len": queue_len,
"has_lock": has_lock,
"has_task_progress": has_task_progress,
}
return (queue_len > 0 or has_lock == 1 or has_task_progress == 1), flags
def _kill_worker_process() -> None: def _kill_worker_process() -> None:
pid_file = "/tmp/celery-worker.pid" pid_file = "/tmp/celery-worker.pid"
pid: int | None = None pid: int | None = None
@@ -101,7 +119,7 @@ def _kill_worker_process() -> None:
# Если процесс ещё жив — принудительно убиваем. # Если процесс ещё жив — принудительно убиваем.
os.kill(pid, 0) os.kill(pid, 0)
logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid) logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid)
os.kill(pid, signal.SIGKILL) os.kill(pid, SIGKILL_FALLBACK)
except ProcessLookupError: except ProcessLookupError:
pass pass
except Exception: except Exception:
@@ -137,8 +155,8 @@ def main() -> None:
redis_client = _get_redis() redis_client = _get_redis()
redis_client.ping() redis_client.ping()
queue_len = _safe_int(redis_client.llen(queue_name), 0) has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
if queue_len <= 0: if not has_inflight:
time.sleep(check_interval) time.sleep(check_interval)
continue continue
@@ -151,8 +169,8 @@ def main() -> None:
time.sleep(check_interval) time.sleep(check_interval)
continue continue
logger.warning( logger.warning(
"Self-heal: queue=%d but no progress timestamp found after startup grace", "Self-heal: inflight=%s but no progress timestamp found after startup grace",
queue_len, inflight,
) )
age = stall_seconds + 1 age = stall_seconds + 1
@@ -174,8 +192,8 @@ def main() -> None:
continue continue
logger.error( logger.error(
"Self-heal: detected global stall (queue=%d, progress_age=%ss > %ss). Restarting worker process...", "Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...",
queue_len, inflight,
age, age,
stall_seconds, stall_seconds,
) )