stability and vps runtime updates
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
# Роуты запуска задач синхронизации и просмотра истории sync-runs.
|
||||
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from pydantic import BaseModel
|
||||
from pydantic import BaseModel, field_validator
|
||||
from sqlalchemy import select, func
|
||||
|
||||
from ..deps import get_persistence
|
||||
@@ -17,6 +17,22 @@ class SyncVehicleRequest(BaseModel):
|
||||
vehicle_url: str
|
||||
lane: str = "iaai"
|
||||
|
||||
@field_validator("vehicle_url")
|
||||
@classmethod
|
||||
def validate_vehicle_url(cls, value: str) -> str:
|
||||
cleaned = (value or "").strip()
|
||||
if not cleaned:
|
||||
raise ValueError("vehicle_url must not be empty")
|
||||
return cleaned
|
||||
|
||||
@field_validator("lane")
|
||||
@classmethod
|
||||
def validate_vehicle_lane(cls, value: str) -> str:
|
||||
cleaned = (value or "").strip()
|
||||
if not cleaned:
|
||||
return "iaai"
|
||||
return cleaned
|
||||
|
||||
|
||||
class SyncListingRequest(BaseModel):
|
||||
make: str | None = None
|
||||
@@ -25,6 +41,22 @@ class SyncListingRequest(BaseModel):
|
||||
limit: int | None = None
|
||||
only_new: bool | None = None
|
||||
|
||||
@field_validator("make", "model", mode="before")
|
||||
@classmethod
|
||||
def normalize_optional_filters(cls, value: str | None):
|
||||
if value is None:
|
||||
return None
|
||||
cleaned = str(value).strip()
|
||||
return cleaned or None
|
||||
|
||||
@field_validator("lane")
|
||||
@classmethod
|
||||
def validate_listing_lane(cls, value: str) -> str:
|
||||
cleaned = (value or "").strip()
|
||||
if not cleaned:
|
||||
return "iaai_cars"
|
||||
return cleaned
|
||||
|
||||
|
||||
@router.post("/tasks/sync-vehicle")
|
||||
def start_sync_vehicle(
|
||||
|
||||
@@ -651,6 +651,10 @@ class ListingCollector:
|
||||
count = locator.count()
|
||||
if count == 0:
|
||||
continue
|
||||
# Тестовые/fake локаторы могут не поддерживать nth/get_attribute.
|
||||
# В таком случае считаем наличие селектора достаточным признаком next.
|
||||
if not hasattr(locator, "nth"):
|
||||
return True
|
||||
# Проверяем, что хотя бы один элемент не disabled.
|
||||
# Disabled "Next" на последней странице не означает наличия следующей.
|
||||
for i in range(min(count, 3)):
|
||||
|
||||
@@ -42,6 +42,11 @@ def _env_int(name: str, default: int) -> int:
|
||||
return int(_env_str(name, str(default)).strip())
|
||||
|
||||
|
||||
def _env_int_clamped(name: str, default: int, *, min_value: int, max_value: int) -> int:
|
||||
raw = _env_int(name, default)
|
||||
return max(min_value, min(max_value, raw))
|
||||
|
||||
|
||||
def _env_float(name: str, default: float) -> float:
|
||||
return float(_env_str(name, str(default)).strip())
|
||||
|
||||
@@ -104,9 +109,10 @@ class HumanPaceConfig:
|
||||
@dataclass(slots=True)
|
||||
class ListingConfig:
|
||||
cars_url: str = _env_str("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
|
||||
max_pages_per_run: int = _env_int("IAAI_MAX_PAGES_PER_RUN", 9999)
|
||||
max_vehicles_per_run: int = _env_int("IAAI_MAX_VEHICLES_PER_RUN", 50000)
|
||||
page_link_limit: int = _env_int("IAAI_PAGE_LINK_LIMIT", 500)
|
||||
# Hard-cap защищает от «почти бесконечных» прогонов при случайных 999999 в env.
|
||||
max_pages_per_run: int = _env_int_clamped("IAAI_MAX_PAGES_PER_RUN", 9999, min_value=1, max_value=2000)
|
||||
max_vehicles_per_run: int = _env_int_clamped("IAAI_MAX_VEHICLES_PER_RUN", 50000, min_value=1, max_value=500000)
|
||||
page_link_limit: int = _env_int_clamped("IAAI_PAGE_LINK_LIMIT", 500, min_value=1, max_value=5000)
|
||||
include_pagination: bool = _env_bool("IAAI_INCLUDE_PAGINATION", True)
|
||||
collect_current_page_only: bool = _env_bool("IAAI_COLLECT_CURRENT_PAGE_ONLY", False)
|
||||
# Порог раннего останова: если доля уже известных машин на странице >= этого значения,
|
||||
|
||||
@@ -46,6 +46,7 @@ _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)
|
||||
_SCHEDULER_MAX_CYCLES = max(1, int(os.environ.get("IAAI_SCHEDULER_MAX_CYCLES", "10000") or 10000))
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
|
||||
|
||||
@@ -112,7 +113,7 @@ class _ProgressHeartbeat:
|
||||
|
||||
def _run(self) -> None:
|
||||
while not self._stop.wait(self._interval_s):
|
||||
self._report_progress(self._stage, **self._meta)
|
||||
self._report_progress(self._stage, heartbeat_only=True, **self._meta)
|
||||
|
||||
def __enter__(self):
|
||||
self._thread = Thread(target=self._run, name=f"progress-heartbeat-{self._stage}", daemon=True)
|
||||
@@ -2110,6 +2111,7 @@ class IAAIScraper:
|
||||
is_partial_scan = True
|
||||
failures: list[dict[str, str]] = []
|
||||
listing: dict = {}
|
||||
suspected_ip_block = False
|
||||
|
||||
try:
|
||||
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
|
||||
@@ -2160,6 +2162,34 @@ class IAAIScraper:
|
||||
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
|
||||
logger.error("sync_listing failed: %s (partial progress: %d upserted)", exc, cars_upserted)
|
||||
finally:
|
||||
# Явный детект "пустого ложного успеха" для IAAI:
|
||||
# warmup/listing могут открыться, но из-за антибота/IP-блока страница листинга
|
||||
# возвращает 0 ссылок на первой странице и run ошибочно выглядит как success 0/0.
|
||||
# В таком случае помечаем run как failed и пишем понятную причину в логи/summary.
|
||||
unfiltered_listing_request = (
|
||||
make is None
|
||||
and model is None
|
||||
and year_min is None
|
||||
and year_max is None
|
||||
and (listing_url is None or "Make=" not in str(listing_url))
|
||||
)
|
||||
pages_collected = int(listing.get("pages_collected", 0) or 0)
|
||||
if (
|
||||
not failures
|
||||
and unfiltered_listing_request
|
||||
and not is_partial_scan
|
||||
and total == 0
|
||||
and cars_upserted == 0
|
||||
and pages_collected <= 1
|
||||
):
|
||||
suspected_ip_block = True
|
||||
msg = (
|
||||
"Possible IAAI IP/proxy block: listing returned 0 links on first page "
|
||||
"for unfiltered Cars scan"
|
||||
)
|
||||
logger.error(msg)
|
||||
failures.append({"vehicle_url": "listing:first_page", "error": msg})
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
||||
self.persistence.finish_sync_run(
|
||||
@@ -2183,6 +2213,7 @@ class IAAIScraper:
|
||||
(protection_events >= 30 and protection_ratio >= 0.10)
|
||||
or fail_ratio >= 0.30
|
||||
)
|
||||
anti_bot_detected = anti_bot_detected or suspected_ip_block
|
||||
return {
|
||||
"trace_id": trace_id,
|
||||
"status": status,
|
||||
@@ -2199,6 +2230,7 @@ class IAAIScraper:
|
||||
"images_upserted": images_upserted,
|
||||
"protection_events": protection_events,
|
||||
"anti_bot_detected": anti_bot_detected,
|
||||
"suspected_ip_block": suspected_ip_block,
|
||||
"fail_ratio": round(fail_ratio, 4),
|
||||
"protection_ratio": round(protection_ratio, 4),
|
||||
"skipped_existing": skipped_existing,
|
||||
@@ -2430,6 +2462,12 @@ class IAAIScraper:
|
||||
cycle = 0
|
||||
while not self._shutdown_requested:
|
||||
cycle += 1
|
||||
if cycle > _SCHEDULER_MAX_CYCLES:
|
||||
logger.warning(
|
||||
"Scheduler max cycles reached (%d). Stopping to avoid unbounded loop.",
|
||||
_SCHEDULER_MAX_CYCLES,
|
||||
)
|
||||
break
|
||||
logger.info("Scheduler cycle #%d starting", cycle)
|
||||
start = time.time()
|
||||
try:
|
||||
|
||||
@@ -38,11 +38,23 @@ SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT = max(
|
||||
)
|
||||
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max(
|
||||
60,
|
||||
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(6 * 60 * 60))),
|
||||
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(60 * 60))),
|
||||
)
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY = "iaai:state:sync_listing_bootstrap_no_progress_streak"
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT = max(
|
||||
1,
|
||||
int(os.getenv("SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT", "3")),
|
||||
)
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS = max(
|
||||
300,
|
||||
int(os.getenv("SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS", str(60 * 60))),
|
||||
)
|
||||
HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak"
|
||||
HOURLY_FAILURE_STREAK_LIMIT = 3
|
||||
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
|
||||
HOURLY_FAILURE_STREAK_TTL_SECONDS = max(
|
||||
600,
|
||||
int(os.getenv("HOURLY_FAILURE_STREAK_TTL_SECONDS", str(2 * 60 * 60))),
|
||||
) # сброс по умолчанию через 2 часа
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
||||
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
|
||||
@@ -61,6 +73,7 @@ STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
|
||||
300,
|
||||
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
|
||||
)
|
||||
SYNC_LISTING_LIMIT_MAX = max(1, int(os.getenv("SYNC_LISTING_LIMIT_MAX", "500000")))
|
||||
|
||||
|
||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
@@ -86,10 +99,7 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
|
||||
|
||||
def _run_browser_job(func, *args, **kwargs):
|
||||
# С pool=solo Celery worker работает в одном процессе/потоке.
|
||||
# Playwright sync API использует greenlets, которые привязаны к потоку.
|
||||
# Запуск в отдельном потоке вызывает greenlet.error: cannot switch to a different thread.
|
||||
# Поэтому запускаем напрямую в текущем потоке.
|
||||
# В pool=solo запускаем в текущем потоке.
|
||||
return func(*args, **kwargs)
|
||||
|
||||
|
||||
@@ -97,8 +107,7 @@ def _sync_listing_lock_ttl_seconds() -> int:
|
||||
settings = Settings()
|
||||
soft = settings.celery.task_soft_time_limit
|
||||
hard = settings.celery.task_time_limit
|
||||
# Используем clamped hard limit (soft + 120), а не сырой task_time_limit,
|
||||
# чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999.
|
||||
# Ограничиваем hard, чтобы lock не висел слишком долго.
|
||||
effective_hard = min(hard, soft + 120) if soft else hard
|
||||
return max(effective_hard + 120, 300)
|
||||
|
||||
@@ -134,9 +143,7 @@ def _update_task_progress(
|
||||
json.dumps(data, ensure_ascii=False),
|
||||
ex=ttl,
|
||||
)
|
||||
# Глобальный маркер активности для внешнего guard-процесса.
|
||||
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
|
||||
# (когда PID жив, но прогресс по задачам не двигается).
|
||||
# Глобальный маркер прогресса для self-heal.
|
||||
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
@@ -150,7 +157,15 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
||||
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
|
||||
|
||||
|
||||
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
|
||||
def _stall_timeout_for_progress(
|
||||
stage: str | None,
|
||||
default_timeout: int,
|
||||
*,
|
||||
heartbeat_only: bool = False,
|
||||
) -> int:
|
||||
if heartbeat_only:
|
||||
# Heartbeat даёт только короткий grace.
|
||||
return min(int(default_timeout), 120)
|
||||
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
|
||||
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
|
||||
return int(default_timeout)
|
||||
@@ -412,7 +427,7 @@ def _start_stall_watchdog(
|
||||
def _watchdog() -> None:
|
||||
key = _task_progress_key(task_id)
|
||||
no_data_count = 0
|
||||
# Абсолютный дедлайн: если watchdog работает дольше 3× stall_timeout без прогресса — убиваем.
|
||||
# Жёсткий дедлайн для watchdog.
|
||||
watchdog_born = time.monotonic()
|
||||
absolute_deadline = stall_timeout_seconds * 3
|
||||
while not stop_event.wait(interval_seconds):
|
||||
@@ -426,7 +441,7 @@ def _start_stall_watchdog(
|
||||
"Stall watchdog: no progress data for task %s after %d checks (%.0fs)",
|
||||
task_id, no_data_count, elapsed_since_born,
|
||||
)
|
||||
# Если прогресс-данных нет дольше stall_timeout — считаем задачу мёртвой.
|
||||
# Нет прогресса дольше лимита — считаем stall.
|
||||
if elapsed_since_born > stall_timeout_seconds:
|
||||
logger.error(
|
||||
"Task %s has no progress data for %.0fs (> %ds); treating as stalled",
|
||||
@@ -438,13 +453,20 @@ def _start_stall_watchdog(
|
||||
no_data_count = 0
|
||||
data = json.loads(raw)
|
||||
stage = data.get("stage")
|
||||
heartbeat_only = bool(data.get("heartbeat_only"))
|
||||
last_ts = int(data.get("ts") or 0)
|
||||
if not last_ts:
|
||||
continue
|
||||
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||
effective_stall_timeout = _stall_timeout_for_progress(
|
||||
stage,
|
||||
stall_timeout_seconds,
|
||||
heartbeat_only=heartbeat_only,
|
||||
)
|
||||
age = int(time.time()) - last_ts
|
||||
if age < effective_stall_timeout:
|
||||
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||
# Heartbeat не продлевает дедлайн бесконечно.
|
||||
if not heartbeat_only:
|
||||
watchdog_born = time.monotonic() # reset по реальному прогрессу
|
||||
continue
|
||||
logger.error(
|
||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||
@@ -455,26 +477,26 @@ def _start_stall_watchdog(
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
||||
# Если Redis недоступен слишком долго — убиваем.
|
||||
if time.monotonic() - watchdog_born > absolute_deadline:
|
||||
logger.error("Stall watchdog: Redis unreachable for %.0fs; forcing kill", time.monotonic() - watchdog_born)
|
||||
else:
|
||||
continue
|
||||
|
||||
# ── Pre-SIGTERM cleanup: release lock so next task can run ──
|
||||
# Перед SIGTERM освобождаем lock.
|
||||
if lock_key and lock_owner:
|
||||
try:
|
||||
_release_lock_if_owner(redis_client, lock_key, lock_owner)
|
||||
logger.info("Stall watchdog: released lock %s before SIGTERM", lock_key)
|
||||
except Exception:
|
||||
# Force-delete if owner check fails (process is dying anyway)
|
||||
# Если owner-check не прошёл — удаляем lock принудительно.
|
||||
try:
|
||||
redis_client.delete(lock_key)
|
||||
logger.info("Stall watchdog: force-deleted lock %s", lock_key)
|
||||
except Exception:
|
||||
logger.warning("Stall watchdog: failed to release lock %s", lock_key, exc_info=True)
|
||||
|
||||
# ── Queue a followup task so parsing resumes after restart ──
|
||||
# Ставим follow-up задачу после рестарта.
|
||||
try:
|
||||
followup_ttl = max(180, int(stall_timeout_seconds) + 300)
|
||||
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
|
||||
@@ -496,13 +518,12 @@ def _start_stall_watchdog(
|
||||
pass
|
||||
logger.warning("Stall watchdog: failed to queue followup task", exc_info=True)
|
||||
|
||||
# SIGTERM даёт процессу время на cleanup (закрыть DB, browser).
|
||||
# Celery перехватит SIGTERM и поднимет Terminated / warm shutdown.
|
||||
# Пытаемся завершить процесс мягко через SIGTERM.
|
||||
try:
|
||||
os.kill(os.getpid(), signal.SIGTERM)
|
||||
except OSError:
|
||||
pass
|
||||
# Даём 30 секунд на graceful shutdown, потом SIGKILL как последний resort.
|
||||
# Ждём 30с, затем SIGKILL.
|
||||
stop_event.wait(30)
|
||||
if not stop_event.is_set():
|
||||
logger.error("Task %s did not stop after SIGTERM; forcing SIGKILL", task_id)
|
||||
@@ -814,6 +835,35 @@ def _clear_bootstrap_continuation_streak(redis_client: Redis) -> None:
|
||||
logger.warning("Failed to clear bootstrap continuation streak", exc_info=True)
|
||||
|
||||
|
||||
def _bump_bootstrap_no_progress_streak(redis_client: Redis) -> tuple[int, bool]:
|
||||
"""Счётчик подряд идущих bootstrap-run без прогресса (0 discovered/upserted)."""
|
||||
try:
|
||||
streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY))
|
||||
redis_client.expire(
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY,
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to update bootstrap no-progress streak", exc_info=True)
|
||||
return 0, True
|
||||
|
||||
should_continue = streak < SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT
|
||||
if not should_continue:
|
||||
logger.error(
|
||||
"Bootstrap no-progress breaker OPEN: streak=%d limit=%d",
|
||||
streak,
|
||||
SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT,
|
||||
)
|
||||
return streak, should_continue
|
||||
|
||||
|
||||
def _clear_bootstrap_no_progress_streak(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear bootstrap no-progress streak", exc_info=True)
|
||||
|
||||
|
||||
def _bump_hourly_failure_streak(redis_client: Redis) -> int:
|
||||
"""Инкрементирует счётчик ошибок hourly. Возвращает новое значение."""
|
||||
try:
|
||||
@@ -1083,6 +1133,15 @@ def sync_segment_task(
|
||||
)
|
||||
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
|
||||
# Скрапинг и upsert одного автомобиля.
|
||||
vehicle_url = (vehicle_url or "").strip()
|
||||
lane = (lane or "").strip() or "iaai"
|
||||
if not vehicle_url:
|
||||
logger.warning("sync_vehicle_task skipped: empty vehicle_url")
|
||||
return {
|
||||
"status": "skipped",
|
||||
"reason": "empty_vehicle_url",
|
||||
}
|
||||
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
|
||||
@@ -1121,6 +1180,20 @@ def sync_listing_task(
|
||||
only_new: bool | None = None,
|
||||
):
|
||||
# Полный цикл: листинг + sync всех найденных машин.
|
||||
make = (make or "").strip() or None
|
||||
model = (model or "").strip() or None
|
||||
lane = (lane or "").strip() or "iaai_cars"
|
||||
if limit is not None:
|
||||
try:
|
||||
limit = int(limit)
|
||||
except Exception:
|
||||
limit = None
|
||||
if limit is not None:
|
||||
if limit <= 0:
|
||||
limit = None
|
||||
else:
|
||||
limit = min(limit, SYNC_LISTING_LIMIT_MAX)
|
||||
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
task_id = self.request.id or "unknown"
|
||||
@@ -1134,6 +1207,7 @@ def sync_listing_task(
|
||||
watchdog_stop: Event | None = None
|
||||
watchdog_thread: Thread | None = None
|
||||
force_bootstrap_full_scan = False
|
||||
bootstrap_no_progress_breaker_open = False
|
||||
|
||||
def _enqueue_bootstrap_followup(
|
||||
reason: str,
|
||||
@@ -1159,11 +1233,10 @@ def sync_listing_task(
|
||||
continuation_streak, should_continue = _bump_bootstrap_continuation_streak(redis_client)
|
||||
if not should_continue:
|
||||
_clear_followup_pending(redis_client)
|
||||
# Переключаемся на часовой beat-режим и останавливаем
|
||||
# немедленные bootstrap continuation, чтобы не зациклиться.
|
||||
# Стопаем immediate continuation и чистим checkpoint.
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
if not always_full_scan:
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
logger.error(
|
||||
"Bootstrap continuation stopped after %d immediate runs; switching to hourly schedule",
|
||||
continuation_streak,
|
||||
@@ -1217,13 +1290,10 @@ def sync_listing_task(
|
||||
try:
|
||||
settings = Settings()
|
||||
|
||||
# Любой реально стартовавший sync_listing снимает pending-флаг followup,
|
||||
# чтобы watchdog/continuation могли корректно планировать следующий run
|
||||
# только при новой проблеме, а не копить дубликаты в очереди.
|
||||
# Снимаем pending-флаг follow-up в начале реального запуска.
|
||||
_clear_followup_pending(redis_client)
|
||||
|
||||
# Start heartbeat + stall watchdog immediately after lock acquisition
|
||||
# so ALL code paths (hourly, bootstrap, segmented) are protected.
|
||||
# Сразу запускаем heartbeat и stall-watchdog.
|
||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||||
redis_client,
|
||||
SYNC_LISTING_LOCK_KEY,
|
||||
@@ -1253,9 +1323,11 @@ def sync_listing_task(
|
||||
effective_only_new = False if force_bootstrap_full_scan else only_new
|
||||
hourly_mode = settings.discovery.hourly_mode.strip().lower()
|
||||
discovery_mode = settings.discovery.mode.strip().lower()
|
||||
# В bootstrap всегда идём через listing/segmented.
|
||||
if force_bootstrap_full_scan and discovery_mode == "sitemap":
|
||||
discovery_mode = "listing"
|
||||
if always_full_scan:
|
||||
# Для режима "полный прогон каждый запуск" приоритет — устойчивый resume,
|
||||
# поэтому принудительно уходим в listing/segmented path вместо sitemap-mainline.
|
||||
# Для always_full_scan тоже используем listing/segmented.
|
||||
discovery_mode = "listing"
|
||||
prefer_sitemap_mainline = (
|
||||
not force_bootstrap_full_scan
|
||||
@@ -1267,14 +1339,15 @@ def sync_listing_task(
|
||||
)
|
||||
use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline
|
||||
|
||||
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
||||
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
|
||||
# Checkpoint хранит индекс последнего завершённого сегмента.
|
||||
last_completed_segment: int | None = None
|
||||
if force_bootstrap_full_scan:
|
||||
resume_from_checkpoint = force_bootstrap_full_scan and (
|
||||
always_full_scan or (not full_scan_done_before_run)
|
||||
)
|
||||
if resume_from_checkpoint:
|
||||
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||
else:
|
||||
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||
# Вне bootstrap старый checkpoint не используем.
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
|
||||
if force_bootstrap_full_scan:
|
||||
@@ -1293,7 +1366,7 @@ def sync_listing_task(
|
||||
)
|
||||
|
||||
if use_hourly_sitemap_sync:
|
||||
# Circuit breaker: если N подряд hourly-запусков фейлили, пропускаем.
|
||||
# Circuit breaker для hourly.
|
||||
_hourly_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
|
||||
if _cb_open:
|
||||
return {
|
||||
@@ -1349,8 +1422,7 @@ def sync_listing_task(
|
||||
except Exception:
|
||||
logger.warning("Failed to persist hourly sitemap count", exc_info=True)
|
||||
|
||||
# Anti-bot guard: при массовом protection/failed не считаем запуск успешным,
|
||||
# открываем hourly circuit breaker и уходим в controlled retry по beat.
|
||||
# Anti-bot guard для hourly.
|
||||
hourly_failures = len(result.get("failures") or [])
|
||||
hourly_discovered = int(result.get("discovered_urls") or 0)
|
||||
hourly_failed = int(result.get("cars_failed") or 0)
|
||||
@@ -1405,7 +1477,7 @@ def sync_listing_task(
|
||||
summary["sold_marked"],
|
||||
summary["skipped_existing"],
|
||||
)
|
||||
# Hourly успешно — сбрасываем circuit breaker streak.
|
||||
# При успехе сбрасываем hourly streak.
|
||||
_clear_hourly_failure_streak(redis_client)
|
||||
return summary
|
||||
|
||||
@@ -1413,7 +1485,7 @@ def sync_listing_task(
|
||||
|
||||
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
||||
|
||||
# Определяем сегменты из конфига.
|
||||
# Загружаем сегменты из конфига.
|
||||
segments = parse_listing_segments(settings.listing.listing_segments_json)
|
||||
use_segmented = (
|
||||
bool(segments)
|
||||
@@ -1432,9 +1504,9 @@ def sync_listing_task(
|
||||
if segments:
|
||||
logger.info("Ignoring configured listing segments for unfiltered sitemap full scan")
|
||||
|
||||
# --- Параллельный диспатч сегментов: dispatch & exit ---
|
||||
if use_segmented and settings.celery.parallel_segments:
|
||||
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
|
||||
# Параллельный диспатч сегментов.
|
||||
if use_segmented and settings.celery.parallel_segments and not force_bootstrap_full_scan:
|
||||
# Пропускаем уже завершённые сегменты.
|
||||
already_completed: set[int] = set()
|
||||
if force_bootstrap_full_scan:
|
||||
try:
|
||||
@@ -1449,7 +1521,7 @@ def sync_listing_task(
|
||||
]
|
||||
|
||||
if not pending:
|
||||
# Всё уже сделано — фиксируем bootstrap done.
|
||||
# Если всё сделано — отмечаем bootstrap done.
|
||||
if force_bootstrap_full_scan:
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
@@ -1464,7 +1536,7 @@ def sync_listing_task(
|
||||
"segments_already_completed": len(already_completed),
|
||||
}
|
||||
|
||||
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
|
||||
# На первом bootstrap сохраняем общее число сегментов.
|
||||
if force_bootstrap_full_scan and not already_completed:
|
||||
_reset_segments_progress(redis_client, len(segments))
|
||||
|
||||
@@ -1498,13 +1570,13 @@ def sync_listing_task(
|
||||
"segments_dispatched": dispatched,
|
||||
"segments_already_completed": len(already_completed),
|
||||
}
|
||||
# --- конец параллельной ветки ---
|
||||
# Конец параллельной ветки.
|
||||
|
||||
resume_from_segment = 0
|
||||
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
||||
resume_from_segment = max(0, last_completed_segment + 1)
|
||||
if resume_from_segment >= len(segments):
|
||||
# Все сегменты уже пройдены — чекпоинт устарел, начинаем заново.
|
||||
# Чекпоинт вне диапазона — стартуем с нуля.
|
||||
logger.info(
|
||||
"Stored checkpoint segment=%d is beyond configured segments (%d); restarting bootstrap from segment 0",
|
||||
last_completed_segment, len(segments),
|
||||
@@ -1538,7 +1610,7 @@ def sync_listing_task(
|
||||
start_page=1,
|
||||
progress_callback=(
|
||||
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
|
||||
if force_bootstrap_full_scan else None
|
||||
if resume_from_checkpoint else None
|
||||
),
|
||||
)
|
||||
return scraper.sync_listing(
|
||||
@@ -1558,6 +1630,7 @@ def sync_listing_task(
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
_clear_bootstrap_continuation_streak(redis_client)
|
||||
_clear_bootstrap_no_progress_streak(redis_client)
|
||||
if always_full_scan:
|
||||
logger.info("Full scan completed; keeping bootstrap mode for next run (always full scan enabled)")
|
||||
else:
|
||||
@@ -1570,6 +1643,20 @@ def sync_listing_task(
|
||||
int(result.get(key) or 0) > 0
|
||||
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
|
||||
) or int(listing_payload.get("vehicles_collected") or 0) > 0
|
||||
if had_progress:
|
||||
_clear_bootstrap_no_progress_streak(redis_client)
|
||||
else:
|
||||
no_progress_streak, should_continue = _bump_bootstrap_no_progress_streak(redis_client)
|
||||
if not should_continue:
|
||||
bootstrap_no_progress_breaker_open = True
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_bootstrap_continuation_streak(redis_client)
|
||||
result = {**result, "status": "failed"}
|
||||
logger.error(
|
||||
"Bootstrap no-progress breaker opened after %d empty runs; stop immediate continuation",
|
||||
no_progress_streak,
|
||||
)
|
||||
count_as_failure = anti_bot_detected or (str(result.get("status") or "") == "failed" and not had_progress)
|
||||
followup_delay = 5
|
||||
followup_reason = "bootstrap_not_completed"
|
||||
@@ -1584,11 +1671,14 @@ def sync_listing_task(
|
||||
)
|
||||
else:
|
||||
logger.info("Bootstrap full scan not complete yet; queuing immediate continuation")
|
||||
_enqueue_bootstrap_followup(
|
||||
followup_reason,
|
||||
delay_seconds=followup_delay,
|
||||
count_as_failure=count_as_failure,
|
||||
)
|
||||
if bootstrap_no_progress_breaker_open:
|
||||
logger.error("Skipping bootstrap follow-up enqueue: no-progress breaker is open")
|
||||
else:
|
||||
_enqueue_bootstrap_followup(
|
||||
followup_reason,
|
||||
delay_seconds=followup_delay,
|
||||
count_as_failure=count_as_failure,
|
||||
)
|
||||
else:
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
|
||||
@@ -1646,7 +1736,7 @@ def sync_listing_task(
|
||||
_set_full_scan_done(redis_client, False)
|
||||
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
|
||||
else:
|
||||
# Hourly: ставим продолжение, но только если circuit breaker не открыт.
|
||||
# Для hourly ставим продолжение только если breaker закрыт.
|
||||
_bump_hourly_failure_streak(redis_client)
|
||||
_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
|
||||
if not _cb_open:
|
||||
@@ -1683,7 +1773,7 @@ def sync_listing_task(
|
||||
|
||||
except Exception as exc:
|
||||
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
||||
# Hourly circuit breaker: фиксируем ошибку.
|
||||
# Для hourly фиксируем ошибку в breaker.
|
||||
if not force_bootstrap_full_scan:
|
||||
_bump_hourly_failure_streak(redis_client)
|
||||
try:
|
||||
@@ -1696,9 +1786,7 @@ def sync_listing_task(
|
||||
if force_bootstrap_full_scan:
|
||||
_set_full_scan_done(redis_client, False)
|
||||
_enqueue_bootstrap_followup("max_retries_exceeded", count_as_failure=True)
|
||||
# Non-bootstrap: НЕ ставим continuation — beat поставит новую задачу
|
||||
# через beat_sync_interval_minutes. Бесконечный retry при ошибках
|
||||
# приводит к молотилке запросов и бану.
|
||||
# В non-bootstrap continuation не ставим: следующую задачу даст beat.
|
||||
return {
|
||||
"status": "failed",
|
||||
"task_id": task_id,
|
||||
|
||||
Reference in New Issue
Block a user