diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 14f5099..6c65c6f 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -45,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90 _PAGE_NEXT_TIMEOUT_S = 60 # Ротация контекста выключена по умолчанию. _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_PAGE = 2 @@ -594,6 +595,7 @@ class IAAIScraper: failures: list[dict[str, str]] = [] early_stopped = False truncated_by_time_budget = False + listing_recovery_attempts = 0 known_origin_ids: set[str] | None = None 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)}) 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 try: page, applied_filters = self._open_listing_for_stream( @@ -763,6 +780,10 @@ class IAAIScraper: pass page = None try: + _consume_listing_recovery_attempt( + page_number=page_number, + reason="context_rotation", + ) page, _ = self._reopen_listing_and_resume( target_page_number=page_number, make=make, @@ -794,6 +815,10 @@ class IAAIScraper: pass page = None try: + _consume_listing_recovery_attempt( + page_number=page_number, + reason=f"collect_failed:{type(op_exc).__name__}", + ) page, _ = self._reopen_listing_and_resume( target_page_number=page_number, make=make, @@ -839,6 +864,10 @@ class IAAIScraper: logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) page.close() try: + _consume_listing_recovery_attempt( + page_number=page_number, + reason="empty_page_after_retry", + ) page, _ = self._reopen_listing_and_resume( target_page_number=page_number, make=make, @@ -1069,6 +1098,10 @@ class IAAIScraper: pass page = None 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( target_page_number=page_number + 1, make=make, diff --git a/iaai_scraper/worker/self_heal.py b/iaai_scraper/worker/self_heal.py index d337032..882a086 100644 --- a/iaai_scraper/worker/self_heal.py +++ b/iaai_scraper/worker/self_heal.py @@ -12,6 +12,8 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal") GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" 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: @@ -76,6 +78,22 @@ def _read_last_progress_ts(redis_client: Redis) -> int | 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: pid_file = "/tmp/celery-worker.pid" pid: int | None = None @@ -101,7 +119,7 @@ def _kill_worker_process() -> None: # Если процесс ещё жив — принудительно убиваем. os.kill(pid, 0) 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: pass except Exception: @@ -137,8 +155,8 @@ def main() -> None: redis_client = _get_redis() redis_client.ping() - queue_len = _safe_int(redis_client.llen(queue_name), 0) - if queue_len <= 0: + has_inflight, inflight = _has_inflight_work(redis_client, queue_name) + if not has_inflight: time.sleep(check_interval) continue @@ -151,8 +169,8 @@ def main() -> None: time.sleep(check_interval) continue logger.warning( - "Self-heal: queue=%d but no progress timestamp found after startup grace", - queue_len, + "Self-heal: inflight=%s but no progress timestamp found after startup grace", + inflight, ) age = stall_seconds + 1 @@ -174,8 +192,8 @@ def main() -> None: continue logger.error( - "Self-heal: detected global stall (queue=%d, progress_age=%ss > %ss). Restarting worker process...", - queue_len, + "Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...", + inflight, age, stall_seconds, )