Files
iaai-parser/iaai_scraper/worker/tasks.py
2026-04-22 21:40:40 +03:00

2050 lines
85 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Задачи Celery.
import json
import logging
import os
import signal
from threading import Event, Thread
import time
import uuid
from typing import Callable
from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task
from redis import Redis
from ..core.config import Settings, parse_listing_segments
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
logger = logging.getLogger("iaai_scraper.worker.tasks")
# Пауза между батчами.
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0)
# Порог ошибок в батче.
_FAIL_RATE_THRESHOLD = min(max(float(os.getenv("IAAI_FAIL_RATE_THRESHOLD", "0.9")), 0.0), 1.0)
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
SYNC_LISTING_RESUME_TARGET_KEY = "iaai:state:sync_listing_resume_target_segment"
SYNC_LISTING_RESUME_TARGET_STREAK_KEY = "iaai:state:sync_listing_resume_target_streak"
SYNC_LISTING_RESUME_TARGET_TTL_SECONDS = 24 * 60 * 60
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT = max(
1,
int(os.getenv("SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT", "2")),
)
SYNC_LISTING_REPEAT_CHECKPOINT_KEY = "iaai:state:sync_listing_repeat_checkpoint"
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY = "iaai:state:sync_listing_repeat_checkpoint_streak"
SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS = 24 * 60 * 60
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT = max(
1,
int(os.getenv("SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT", "2")),
)
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY = "iaai:state:sync_listing_bootstrap_continuation_streak"
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT = max(
1,
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT", "6")),
)
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max(
60,
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(6 * 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 часов.
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}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
SYNC_SEGMENTS_CYCLE_KEY = "iaai:state:sync_segments_cycle_id"
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment"
SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
last_exc: Exception | None = None
for attempt in range(1, attempts + 1):
try:
return func()
except Exception as exc:
last_exc = exc
if attempt >= attempts:
break
delay = base_delay_s * (2 ** (attempt - 1))
logger.warning(
"Operation failed (attempt %d/%d): %s. Retrying in %.1fs",
attempt,
attempts,
exc,
delay,
)
time.sleep(delay)
if last_exc is not None:
raise last_exc
def _run_browser_job(func, *args, **kwargs):
# Запускаем в текущем потоке.
return func(*args, **kwargs)
def _sync_listing_lock_ttl_seconds() -> int:
settings = Settings()
soft = settings.celery.task_soft_time_limit
hard = settings.celery.task_time_limit
# Ограничиваем hard limit.
effective_hard = min(hard, soft + 120) if soft else hard
return max(effective_hard + 120, 300)
def _sync_segment_lock_ttl_seconds() -> int:
return 300
def _task_progress_key(task_id: str) -> str:
return TASK_PROGRESS_KEY_FMT.format(task_id=task_id)
def _update_task_progress(
redis_client: Redis,
*,
task_id: str,
stage: str,
ttl_seconds: int,
**payload,
) -> None:
try:
now_ts = int(time.time())
data = {
"task_id": task_id,
"stage": stage,
"ts": now_ts,
**payload,
}
ttl = max(60, int(ttl_seconds))
pipe = redis_client.pipeline()
pipe.set(
_task_progress_key(task_id),
json.dumps(data, ensure_ascii=False),
ex=ttl,
)
# Глобальный маркер активности для внешнего guard-процесса.
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
# (когда PID жив, но прогресс по задачам не двигается).
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
pipe.execute()
except Exception:
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
try:
redis_client.delete(_task_progress_key(task_id))
except Exception:
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if new_urls:
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
batch_size = settings.celery.batch_size
for batch_start in range(0, len(new_urls), batch_size):
batch_urls = new_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
# Fail-fast: если слишком много ошибок — IAAI блокирует, не тратим ресурсы.
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + batch_size < len(new_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": len(discovered_urls) - len(new_urls),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_diff",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_full_refresh_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
batch_size = settings.celery.batch_size
for batch_start in range(0, len(discovered_urls), batch_size):
batch_urls = discovered_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + batch_size < len(discovered_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": 0,
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_full_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": None,
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_rolling_refresh_sync(
*,
redis_client: Redis,
lane: str,
limit: int | None,
only_new: bool | None,
progress_callback=None,
) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
total_active = persistence.count_active_cars_for_lane("iaai:")
batch_size = max(1, int(settings.discovery.hourly_refresh_batch_size))
try:
offset = int(redis_client.get(SITEMAP_HOURLY_REFRESH_OFFSET_KEY) or 0)
except Exception:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=offset,
limit=batch_size,
)
if not refresh_urls and total_active > 0:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=0,
limit=batch_size,
)
next_offset = 0
if total_active > 0:
next_offset = offset + len(refresh_urls)
if next_offset >= total_active:
next_offset = 0
try:
redis_client.set(SITEMAP_HOURLY_REFRESH_OFFSET_KEY, str(next_offset))
except Exception:
logger.warning("Failed to persist hourly rolling refresh offset", exc_info=True)
seen: set[str] = set()
target_urls: list[str] = []
for url in new_urls + refresh_urls:
if url in seen:
continue
seen.add(url)
target_urls.append(url)
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if target_urls:
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
worker_batch_size = settings.celery.batch_size
for batch_start in range(0, len(target_urls), worker_batch_size):
batch_urls = target_urls[batch_start:batch_start + worker_batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + worker_batch_size < len(target_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": max(0, len(discovered_urls) - len(new_urls)),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"refresh_urls": len(refresh_urls),
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
"refresh_offset": offset,
"refresh_next_offset": next_offset,
"active_total": total_active,
}
def _start_stall_watchdog(
redis_client: Redis,
*,
task_id: str,
stall_timeout_seconds: int,
lock_key: str | None = None,
lock_owner: str | None = None,
on_stall: Callable[[dict[str, object]], None] | None = None,
) -> tuple[Event, Thread]:
stop_event = Event()
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
def _watchdog() -> None:
key = _task_progress_key(task_id)
no_data_count = 0
# Абсолютный дедлайн: если watchdog работает дольше 3× stall_timeout без прогресса — убиваем.
watchdog_born = time.monotonic()
absolute_deadline = stall_timeout_seconds * 3
while not stop_event.wait(interval_seconds):
try:
raw = redis_client.get(key)
if not raw:
no_data_count += 1
elapsed_since_born = time.monotonic() - watchdog_born
if no_data_count % 5 == 0:
logger.warning(
"Stall watchdog: no progress data for task %s after %d checks (%.0fs)",
task_id, no_data_count, elapsed_since_born,
)
# Если прогресс-данных нет дольше stall_timeout — считаем задачу мёртвой.
if elapsed_since_born > stall_timeout_seconds:
logger.error(
"Task %s has no progress data for %.0fs (> %ds); treating as stalled",
task_id, elapsed_since_born, stall_timeout_seconds,
)
else:
continue
else:
no_data_count = 0
data = json.loads(raw)
last_ts = int(data.get("ts") or 0)
if not last_ts:
continue
age = int(time.time()) - last_ts
if age < stall_timeout_seconds:
watchdog_born = time.monotonic() # reset absolute deadline on real progress
continue
logger.error(
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
task_id,
age,
data.get("stage"),
data,
)
if on_stall is not None:
try:
on_stall(data)
except Exception:
logger.warning("Stall watchdog hook failed for task %s", task_id, exc_info=True)
except Exception:
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
# Если 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 ──
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)
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 ──
try:
followup_ttl = max(180, int(stall_timeout_seconds) + 300)
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
from ..worker.celery_app import celery_app
celery_app.send_task(
SYNC_LISTING_TASK_NAME,
kwargs={},
queue="scraping",
countdown=15,
expires=followup_ttl,
)
logger.info("Stall watchdog: queued followup sync_listing_task after stall")
else:
logger.info("Stall watchdog: followup already pending, skip duplicate enqueue")
except Exception:
try:
_clear_followup_pending(redis_client)
except Exception:
pass
logger.warning("Stall watchdog: failed to queue followup task", exc_info=True)
# SIGTERM даёт процессу время на cleanup (закрыть DB, browser).
# Celery перехватит SIGTERM и поднимет Terminated / warm shutdown.
try:
os.kill(os.getpid(), signal.SIGTERM)
except OSError:
pass
# Даём 30 секунд на graceful shutdown, потом SIGKILL как последний resort.
stop_event.wait(30)
if not stop_event.is_set():
logger.error("Task %s did not stop after SIGTERM; forcing SIGKILL", task_id)
os.kill(os.getpid(), signal.SIGKILL)
thread = Thread(target=_watchdog, name=f"task-stall-watchdog-{task_id[:8]}", daemon=True)
thread.start()
return stop_event, thread
_persistence_instance: PersistenceService | None = None
def _get_persistence() -> PersistenceService:
global _persistence_instance
if _persistence_instance is not None:
return _persistence_instance
sett = Settings()
persistence = PersistenceService(sett)
def _ping_db() -> None:
with persistence.engine.connect() as conn:
conn.exec_driver_sql("SELECT 1")
_retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0)
_persistence_instance = persistence
return persistence
_redis_instance: Redis | None = None
def _get_redis() -> Redis:
global _redis_instance
if _redis_instance is not None:
try:
_redis_instance.ping()
return _redis_instance
except Exception:
_redis_instance = None
sett = Settings()
redis_client = Redis.from_url(
sett.redis.url,
decode_responses=True,
socket_connect_timeout=sett.redis.socket_connect_timeout_seconds,
socket_timeout=sett.redis.socket_timeout_seconds,
health_check_interval=sett.redis.health_check_interval_seconds,
retry_on_timeout=True,
)
def _ping_redis() -> None:
redis_client.ping()
_retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0)
_redis_instance = redis_client
return redis_client
def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool:
try:
acquired = bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
if acquired:
return True
# Автовосстановление: если lock завис без TTL, считаем stale и пересоздаём.
ttl = redis_client.ttl(key)
if ttl is not None and ttl < 0:
logger.warning("Detected stale lock without TTL, removing: %s", key)
redis_client.delete(key)
return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
return False
except Exception as exc:
logger.warning("Failed to acquire lock %s", key, exc_info=True)
return False
def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None:
try:
refreshed = redis_client.eval(
"""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
end
return 0
""",
1,
key,
owner_token,
int(ttl_seconds),
)
return bool(refreshed)
except Exception as exc:
logger.warning("Failed to refresh lock %s", key, exc_info=True)
return None
def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None:
try:
redis_client.eval(
"""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
""",
1,
key,
owner_token,
)
except Exception as exc:
logger.warning("Failed to release lock %s", key, exc_info=True)
def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool:
try:
inspector = celery_app.control.inspect(timeout=1.0)
snapshots = [
inspector.active() or {},
inspector.reserved() or {},
inspector.scheduled() or {},
]
except Exception:
logger.warning("Failed to inspect Celery workers for running sync tasks", exc_info=True)
return True
for snapshot in snapshots:
for entries in snapshot.values():
for entry in entries or []:
task_name = str(entry.get("name") or entry.get("request", {}).get("name") or "")
if task_name != SYNC_LISTING_TASK_NAME:
continue
# Исключаем текущую задачу — она не считается "другой запущенной"
entry_id = str(entry.get("id") or entry.get("request", {}).get("id") or "")
if exclude_task_id and entry_id == exclude_task_id:
continue
return True
return False
def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_task_id: str | None = None) -> bool:
try:
owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY)
if not owner_token:
return False
except Exception:
logger.warning("Failed to read sync listing lock before cleanup", exc_info=True)
return False
if _has_running_sync_listing_tasks(celery_app, exclude_task_id=current_task_id):
logger.info("sync_listing lock preserved: active task still detected")
return False
try:
ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY)
redis_client.delete(SYNC_LISTING_LOCK_KEY)
logger.warning(
"Removed orphan sync_listing lock owner=%s ttl=%s after worker restart",
owner_token,
ttl,
)
return True
except Exception:
logger.warning("Failed to clear orphan sync listing lock", exc_info=True)
return False
def _is_full_scan_done(redis_client: Redis) -> bool:
try:
value = redis_client.get(SYNC_FULL_SCAN_DONE_KEY)
except Exception:
logger.warning("Failed to read full scan state", exc_info=True)
return False
return str(value or "").strip() == "1"
def _set_full_scan_done(redis_client: Redis, done: bool) -> None:
try:
redis_client.set(SYNC_FULL_SCAN_DONE_KEY, "1" if done else "0")
except Exception:
logger.warning("Failed to persist full scan state", exc_info=True)
def _load_last_completed_segment(redis_client: Redis) -> int | None:
# Segment-only checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
try:
raw = redis_client.get(SYNC_LISTING_CHECKPOINT_KEY)
except Exception:
logger.warning("Failed to read sync listing checkpoint", exc_info=True)
return None
if raw is None:
return None
try:
value = int(str(raw).strip())
except (TypeError, ValueError):
logger.warning("Invalid sync listing checkpoint value %r; clearing", raw)
_clear_sync_checkpoint(redis_client)
return None
if value < 0:
_clear_sync_checkpoint(redis_client)
return None
return value
def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> None:
try:
redis_client.set(
SYNC_LISTING_CHECKPOINT_KEY,
str(int(segment_index)),
ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to save sync listing checkpoint", exc_info=True)
def _clear_sync_checkpoint(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY)
except Exception:
logger.warning("Failed to clear sync listing checkpoint", exc_info=True)
def _save_stalled_segment(redis_client: Redis, segment_index: int, *, task_id: str, stage: str | None) -> None:
try:
payload = {
"segment_index": int(segment_index),
"task_id": task_id,
"stage": stage,
"ts": int(time.time()),
}
redis_client.set(
SYNC_LISTING_STALLED_SEGMENT_KEY,
json.dumps(payload, ensure_ascii=False),
ex=SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to persist stalled segment info", exc_info=True)
def _infer_stalled_segment_index(redis_client: Redis, payload: dict[str, object]) -> int | None:
"""Пытается определить индекс застрявшего сегмента из payload или checkpoint."""
raw_segment = payload.get("segment_index")
if raw_segment is not None:
try:
return max(0, int(raw_segment))
except Exception:
pass
last_completed = _load_last_completed_segment(redis_client)
if last_completed is None:
return 0
return max(0, int(last_completed) + 1)
def _load_stalled_segment(redis_client: Redis) -> int | None:
try:
raw = redis_client.get(SYNC_LISTING_STALLED_SEGMENT_KEY)
except Exception:
logger.warning("Failed to read stalled segment info", exc_info=True)
return None
if not raw:
return None
try:
data = json.loads(raw)
return int(data.get("segment_index"))
except Exception:
logger.warning("Invalid stalled segment payload %r; clearing", raw)
_clear_stalled_segment(redis_client)
return None
def _clear_stalled_segment(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_STALLED_SEGMENT_KEY)
except Exception:
logger.warning("Failed to clear stalled segment info", exc_info=True)
def _register_bootstrap_resume_target(redis_client: Redis, segment_index: int) -> tuple[int, bool]:
"""Возвращает (streak, should_skip_segment) для защиты от вечного resume на одном сегменте."""
ttl = max(60, int(SYNC_LISTING_RESUME_TARGET_TTL_SECONDS))
try:
raw_target = redis_client.get(SYNC_LISTING_RESUME_TARGET_KEY)
raw_streak = redis_client.get(SYNC_LISTING_RESUME_TARGET_STREAK_KEY)
prev_target = int(str(raw_target).strip()) if raw_target is not None else None
prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0
except Exception:
logger.warning("Failed to read bootstrap resume target guard", exc_info=True)
prev_target = None
prev_streak = 0
streak = (prev_streak + 1) if prev_target == int(segment_index) else 1
try:
pipe = redis_client.pipeline()
pipe.set(SYNC_LISTING_RESUME_TARGET_KEY, str(int(segment_index)), ex=ttl)
pipe.set(SYNC_LISTING_RESUME_TARGET_STREAK_KEY, str(int(streak)), ex=ttl)
pipe.execute()
except Exception:
logger.warning("Failed to persist bootstrap resume target guard", exc_info=True)
# Если bootstrap второй раз подряд возвращается в тот же сегмент,
# считаем его застрявшим и отпускаем checkpoint дальше.
should_skip_segment = streak >= SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT
if should_skip_segment:
logger.error(
"Bootstrap resume loop guard OPEN: segment=%d streak=%d limit=%d",
segment_index,
streak,
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT,
)
elif streak > 1:
logger.warning(
"Bootstrap resume repeated for segment=%d (%d/%d)",
segment_index,
streak,
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT,
)
return streak, should_skip_segment
def _register_repeated_bootstrap_checkpoint(redis_client: Redis, checkpoint_segment: int) -> tuple[int, bool]:
"""Возвращает (streak, should_skip_next_segment), если bootstrap снова стартует с тем же checkpoint."""
ttl = max(60, int(SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS))
try:
raw_checkpoint = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_KEY)
raw_streak = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY)
prev_checkpoint = int(str(raw_checkpoint).strip()) if raw_checkpoint is not None else None
prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0
except Exception:
logger.warning("Failed to read repeated bootstrap checkpoint guard", exc_info=True)
prev_checkpoint = None
prev_streak = 0
streak = (prev_streak + 1) if prev_checkpoint == int(checkpoint_segment) else 1
try:
pipe = redis_client.pipeline()
pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_KEY, str(int(checkpoint_segment)), ex=ttl)
pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY, str(int(streak)), ex=ttl)
pipe.execute()
except Exception:
logger.warning("Failed to persist repeated bootstrap checkpoint guard", exc_info=True)
should_skip_next_segment = streak >= SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT
if should_skip_next_segment:
logger.error(
"Bootstrap checkpoint repeat guard OPEN: checkpoint=%d streak=%d limit=%d",
checkpoint_segment,
streak,
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT,
)
elif streak > 1:
logger.warning(
"Bootstrap repeated with same checkpoint=%d (%d/%d)",
checkpoint_segment,
streak,
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT,
)
return streak, should_skip_next_segment
def _clear_bootstrap_resume_target(redis_client: Redis) -> None:
try:
redis_client.delete(
SYNC_LISTING_RESUME_TARGET_KEY,
SYNC_LISTING_RESUME_TARGET_STREAK_KEY,
SYNC_LISTING_REPEAT_CHECKPOINT_KEY,
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY,
)
except Exception:
logger.warning("Failed to clear bootstrap resume target guard", exc_info=True)
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
try:
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
except Exception:
logger.warning("Failed to set follow-up pending flag", exc_info=True)
return True
def _clear_followup_pending(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
except Exception:
logger.warning("Failed to clear follow-up pending flag", exc_info=True)
def _bump_bootstrap_failure_streak(
redis_client: Redis,
*,
reason: str,
) -> tuple[int, bool]:
try:
streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY))
redis_client.expire(
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY,
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to update bootstrap failure streak", exc_info=True)
return 0, True
should_enqueue = streak < SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT
if should_enqueue:
logger.warning(
"Bootstrap failure streak %d/%d recorded (reason=%s)",
streak,
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT,
reason,
)
else:
logger.error(
"Bootstrap follow-up circuit breaker opened after %d consecutive failures (reason=%s)",
streak,
reason,
)
return streak, should_enqueue
def _clear_bootstrap_failure_streak(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY)
except Exception:
logger.warning("Failed to clear bootstrap failure streak", exc_info=True)
def _bump_bootstrap_continuation_streak(redis_client: Redis) -> tuple[int, bool]:
"""Счётчик подряд идущих bootstrap-followup запусков.
Возвращает (streak, should_continue). Если should_continue=False,
немедленные continuation блокируются до следующего beat-цикла.
"""
try:
streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY))
redis_client.expire(
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY,
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to update bootstrap continuation streak", exc_info=True)
return 0, True
should_continue = streak <= SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT
if not should_continue:
logger.error(
"Bootstrap continuation breaker OPEN: streak=%d limit=%d",
streak,
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT,
)
return streak, should_continue
def _clear_bootstrap_continuation_streak(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY)
except Exception:
logger.warning("Failed to clear bootstrap continuation streak", exc_info=True)
def _bump_hourly_failure_streak(redis_client: Redis) -> int:
"""Инкрементирует счётчик ошибок hourly. Возвращает новое значение."""
try:
streak = int(redis_client.incr(HOURLY_FAILURE_STREAK_KEY))
redis_client.expire(HOURLY_FAILURE_STREAK_KEY, HOURLY_FAILURE_STREAK_TTL_SECONDS)
return streak
except Exception:
logger.warning("Failed to update hourly failure streak", exc_info=True)
return 0
def _check_hourly_circuit_breaker(redis_client: Redis) -> tuple[int, bool]:
"""Проверяет открыт ли circuit breaker. Возвращает (streak, is_open)."""
try:
raw = redis_client.get(HOURLY_FAILURE_STREAK_KEY)
streak = int(raw) if raw else 0
except Exception:
return 0, False
is_open = streak >= HOURLY_FAILURE_STREAK_LIMIT
if is_open:
logger.error("Hourly circuit breaker OPEN (%d/%d failures) — skipping", streak, HOURLY_FAILURE_STREAK_LIMIT)
return streak, is_open
def _clear_hourly_failure_streak(redis_client: Redis) -> None:
try:
redis_client.delete(HOURLY_FAILURE_STREAK_KEY)
except Exception:
pass
def _start_lock_heartbeat(
redis_client: Redis,
key: str,
owner_token: str,
ttl_seconds: int,
) -> tuple[Event, Thread]:
stop_event = Event()
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
def _heartbeat() -> None:
consecutive_failures = 0
while not stop_event.wait(interval_seconds):
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
if refreshed is False:
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
return
if refreshed is None:
consecutive_failures += 1
if consecutive_failures >= 5:
logger.error("Lock heartbeat failed %d times in a row for %s; giving up", consecutive_failures, owner_token)
return
else:
consecutive_failures = 0
thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True)
thread.start()
return stop_event, thread
def _reset_segments_progress(redis_client: Redis, total: int) -> None:
try:
pipe = redis_client.pipeline()
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.execute()
except Exception:
logger.warning("Failed to reset segments progress", exc_info=True)
def _clear_segments_progress_state(redis_client: Redis) -> None:
try:
redis_client.delete(
SYNC_SEGMENTS_PROGRESS_KEY,
SYNC_SEGMENTS_TOTAL_KEY,
SYNC_SEGMENTS_CYCLE_KEY,
)
except Exception:
logger.warning("Failed to clear segments progress state", exc_info=True)
def _ensure_segments_cycle(redis_client: Redis, total: int) -> tuple[str, bool]:
"""Возвращает (cycle_id, created_new_cycle)."""
ttl = max(60, int(SYNC_SEGMENTS_PROGRESS_TTL_SECONDS))
try:
raw_cycle = redis_client.get(SYNC_SEGMENTS_CYCLE_KEY)
cycle_id = str(raw_cycle).strip() if raw_cycle else ""
except Exception:
logger.warning("Failed to read segments cycle id", exc_info=True)
cycle_id = ""
if cycle_id:
try:
pipe = redis_client.pipeline()
pipe.expire(SYNC_SEGMENTS_CYCLE_KEY, ttl)
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, ttl)
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=ttl)
pipe.execute()
except Exception:
logger.warning("Failed to refresh segments cycle ttl", exc_info=True)
return cycle_id, False
cycle_id = uuid.uuid4().hex
try:
_reset_segments_progress(redis_client, total)
redis_client.set(SYNC_SEGMENTS_CYCLE_KEY, cycle_id, ex=ttl)
except Exception:
logger.warning("Failed to initialize segments cycle id", exc_info=True)
return cycle_id, True
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
try:
pipe = redis_client.pipeline()
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.expire(SYNC_SEGMENTS_CYCLE_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
results = pipe.execute()
completed = int(results[3] or 0)
total = int(results[4] or 0) if results[4] else 0
return completed, total
except Exception:
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
return 0, 0
@shared_task(
name="iaai_scraper.worker.tasks.sync_segment_task",
bind=True,
max_retries=2,
default_retry_delay=60,
acks_late=True,
)
def sync_segment_task(
self,
segment_index: int,
segment: dict,
lane: str = "iaai_cars",
only_new: bool | None = None,
is_bootstrap: bool = False,
):
"""Обработка одного сегмента листинга. Запускается параллельно несколькими воркерами."""
persistence = _get_persistence()
persistence.create_tables()
task_id = self.request.id or "unknown"
redis_client = _get_redis()
settings = Settings()
# Per-segment lock — защита от случайного дубля
seg_lock_key = SYNC_SEGMENT_LOCK_KEY_FMT.format(idx=int(segment_index))
owner_token = f"{task_id}:{uuid.uuid4().hex}"
lock_ttl = _sync_segment_lock_ttl_seconds()
lock_acquired = _acquire_lock(redis_client, seg_lock_key, owner_token, lock_ttl)
if not lock_acquired:
logger.info("sync_segment_task[%d] skipped: already running", segment_index)
return {"status": "skipped", "segment_index": segment_index, "reason": "duplicate"}
seg_make = segment.get("make")
seg_year_min = segment.get("year_min")
seg_year_max = segment.get("year_max")
seg_label = f"{seg_make or 'ALL'}"
if seg_year_min is not None or seg_year_max is not None:
seg_label += f" ({seg_year_min}-{seg_year_max})"
heartbeat_stop: Event | None = None
heartbeat_thread: Thread | None = None
watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_started",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(redis_client, seg_lock_key, owner_token, lock_ttl)
watchdog_stop, watchdog_thread = _start_stall_watchdog(
redis_client,
task_id=task_id,
stall_timeout_seconds=stall_timeout,
lock_key=seg_lock_key,
lock_owner=owner_token,
)
try:
def _job():
with IAAIScraper() as scraper:
scraper.set_progress_callback(
lambda stage, meta: _update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
**meta,
)
)
base_url = scraper.settings.listing.cars_url
seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None
return scraper.sync_listing(
make=None if seg_url else seg_make,
model=None,
lane=lane,
only_new=only_new,
listing_url=seg_url,
year_min=seg_year_min,
year_max=seg_year_max,
skip_mark_sold=True,
)
result = _run_browser_job(_job)
segment_done = bool(result.get("full_scan_completed", False))
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_completed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
status=result.get("status", "success"),
cars_upserted=result.get("cars_upserted", 0),
cars_failed=result.get("cars_failed", 0),
)
if is_bootstrap and segment_done:
completed, total = _mark_segment_completed(redis_client, segment_index)
logger.warning(
"Segment %d (%s) bootstrap done: %d/%d completed",
segment_index, seg_label, completed, total,
)
if total > 0 and completed >= total:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("All %d segments completed; bootstrap full scan done", total)
return {
"status": result.get("status", "success"),
"segment_index": segment_index,
"segment_label": seg_label,
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
"full_scan_completed": segment_done,
}
except SoftTimeLimitExceeded:
logger.warning("sync_segment_task[%d] soft timeout — partial progress saved", segment_index)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_soft_timeout",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
return {
"status": "timed_out",
"segment_index": segment_index,
"segment_label": seg_label,
}
except Exception as exc:
logger.error("sync_segment_task[%d] failed: %s%s", segment_index, seg_label, exc, exc_info=True)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_failed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
error=str(exc),
)
try:
raise self.retry(exc=exc)
except self.MaxRetriesExceededError:
return {
"status": "failed",
"segment_index": segment_index,
"segment_label": seg_label,
"error": str(exc),
}
finally:
if watchdog_stop is not None:
watchdog_stop.set()
if watchdog_thread is not None:
watchdog_thread.join(timeout=1)
if heartbeat_stop is not None:
heartbeat_stop.set()
if heartbeat_thread is not None:
heartbeat_thread.join(timeout=1)
_clear_task_progress(redis_client, task_id)
_release_lock_if_owner(redis_client, seg_lock_key, owner_token)
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
# Скрапинг и upsert одного автомобиля.
persistence = _get_persistence()
persistence.create_tables()
try:
def _job():
with IAAIScraper() as scraper:
return scraper.sync_vehicle(vehicle_url, lane=lane)
result = _run_browser_job(_job)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
"vehicle_url": vehicle_url,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
}
except Exception as exc:
logger.error("sync_vehicle_task failed: %s%s", vehicle_url, exc, exc_info=True)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.sync_listing_task",
bind=True,
max_retries=3,
default_retry_delay=120,
acks_late=True,
)
def sync_listing_task(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
# Полный цикл: листинг + sync всех найденных машин.
persistence = _get_persistence()
persistence.create_tables()
task_id = self.request.id or "unknown"
owner_token = f"{task_id}:{uuid.uuid4().hex}"
redis_client = _get_redis()
lock_acquired = False
lock_ttl = _sync_listing_lock_ttl_seconds()
heartbeat_stop: Event | None = None
heartbeat_thread: Thread | None = None
watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None
force_bootstrap_full_scan = False
def _enqueue_bootstrap_followup(
reason: str,
delay_seconds: int = 5,
*,
count_as_failure: bool = False,
) -> None:
flag_ttl = max(lock_ttl, delay_seconds + 300)
if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl):
logger.info(
"Bootstrap follow-up already pending; skip enqueue (reason=%s)",
reason,
)
return
if count_as_failure:
_, should_enqueue = _bump_bootstrap_failure_streak(redis_client, reason=reason)
if not should_enqueue:
_clear_followup_pending(redis_client)
return
else:
_clear_bootstrap_failure_streak(redis_client)
continuation_streak, should_continue = _bump_bootstrap_continuation_streak(redis_client)
if not should_continue:
_clear_followup_pending(redis_client)
# Переключаемся на часовой beat-режим и останавливаем
# немедленные bootstrap continuation, чтобы не зациклиться.
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,
)
else:
logger.error(
"Bootstrap continuation stopped after %d immediate runs (always_full_scan=true)",
continuation_streak,
)
return
try:
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
self.app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task",
kwargs={
"make": make,
"model": model,
"lane": lane,
"limit": limit,
"only_new": only_new,
},
queue="scraping",
countdown=max(0, int(delay_seconds)),
expires=followup_expires,
)
logger.info(
"Bootstrap follow-up sync queued in %ss (reason=%s)",
delay_seconds,
reason,
)
except Exception:
_clear_followup_pending(redis_client)
logger.warning("Failed to enqueue bootstrap follow-up sync", exc_info=True)
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
if not lock_acquired:
orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app, current_task_id=task_id)
if orphan_cleared:
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
if not lock_acquired:
logger.info("sync_listing_task skipped: another sync is already running")
return {
"status": "skipped",
"reason": "sync_already_running",
"task_id": task_id,
}
try:
settings = Settings()
# Любой реально стартовавший sync_listing снимает pending-флаг followup,
# чтобы watchdog/continuation могли корректно планировать следующий run
# только при новой проблеме, а не копить дубликаты в очереди.
_clear_followup_pending(redis_client)
# Start heartbeat + stall watchdog immediately after lock acquisition
# so ALL code paths (hourly, bootstrap, segmented) are protected.
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
redis_client,
SYNC_LISTING_LOCK_KEY,
owner_token,
lock_ttl,
)
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
watchdog_stop, watchdog_thread = _start_stall_watchdog(
redis_client,
task_id=task_id,
stall_timeout_seconds=stall_timeout,
lock_key=SYNC_LISTING_LOCK_KEY,
lock_owner=owner_token,
on_stall=(
lambda payload: (
_save_stalled_segment(
redis_client,
stalled_segment_index,
task_id=task_id,
stage=str(payload.get("stage") or ""),
)
if force_bootstrap_full_scan
and (stalled_segment_index := _infer_stalled_segment_index(redis_client, payload)) is not None
else None
)
),
)
_update_task_progress(
redis_client,
task_id=task_id,
stage="sync_listing_started",
ttl_seconds=progress_ttl,
)
full_scan_done_before_run = _is_full_scan_done(redis_client)
always_full_scan = bool(settings.discovery.always_full_scan)
force_bootstrap_full_scan = always_full_scan or (not full_scan_done_before_run)
effective_limit = None if force_bootstrap_full_scan else limit
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()
if always_full_scan:
# Для режима "полный прогон каждый запуск" приоритет — устойчивый resume,
# поэтому принудительно уходим в listing/segmented path вместо sitemap-mainline.
discovery_mode = "listing"
prefer_sitemap_mainline = (
not force_bootstrap_full_scan
and make is None
and model is None
and effective_limit is None
and not effective_only_new
and not always_full_scan
)
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.
last_completed_segment: int | None = None
if force_bootstrap_full_scan:
last_completed_segment = _load_last_completed_segment(redis_client)
stalled_segment = _load_stalled_segment(redis_client)
else:
# После завершения bootstrap чекпоинт не нужен никогда.
_clear_sync_checkpoint(redis_client)
stalled_segment = None
if force_bootstrap_full_scan:
logger.info(
"Bootstrap mode: forcing full scan (only_new=False, limit=None) until first complete run",
)
if always_full_scan:
logger.info(
"Always full scan mode enabled (IAAI_ALWAYS_FULL_SCAN=true): using listing resume path",
)
logger.info(
"sync_listing options: only_new=%s, limit=%s",
effective_only_new,
effective_limit,
)
if use_hourly_sitemap_sync:
# Circuit breaker: если N подряд hourly-запусков фейлили, пропускаем.
_hourly_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
if _cb_open:
return {
"status": "circuit_breaker_open",
"task_id": task_id,
"hourly_failure_streak": _hourly_streak,
}
_progress_cb = lambda stage, meta: _update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
**meta,
)
try:
if hourly_mode == "diff":
logger.info("Hourly mode: running sitemap diff sync instead of full listing traversal")
result = _hourly_sitemap_diff_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
elif hourly_mode == "full_refresh":
logger.info("Hourly mode: running sitemap full refresh of all active vehicles")
result = _hourly_sitemap_full_refresh_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
else:
logger.info("Hourly mode: running sitemap rolling refresh of active vehicles")
result = _hourly_sitemap_rolling_refresh_sync(
redis_client=redis_client,
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
except SitemapDiscoveryError as exc:
logger.warning(
"Sitemap discovery failed (%s); falling back to listing traversal for hourly sync",
exc,
)
use_hourly_sitemap_sync = False # fall through to listing traversal below
prefer_sitemap_mainline = False # allow segmented fallback
discovery_mode = "listing" # override so use_segmented check passes
if use_hourly_sitemap_sync:
try:
redis_client.set(SITEMAP_HOURLY_LAST_COUNT_KEY, str(result.get("discovered_urls", 0)))
except Exception:
logger.warning("Failed to persist hourly sitemap count", exc_info=True)
# Anti-bot guard: при массовом protection/failed не считаем запуск успешным,
# открываем hourly circuit breaker и уходим в controlled retry по beat.
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)
hourly_protection = int(result.get("protection_events") or 0)
if hourly_discovered > 0:
fail_ratio = hourly_failed / max(1, hourly_discovered)
protection_ratio = hourly_protection / max(1, hourly_discovered)
anti_bot_suspected = (
(hourly_protection >= 30 and protection_ratio >= 0.10)
or fail_ratio >= 0.30
)
if anti_bot_suspected:
_streak = _bump_hourly_failure_streak(redis_client)
logger.error(
"Hourly anti-bot guard triggered: discovered=%d failed=%d protection=%d fail_ratio=%.2f protection_ratio=%.2f streak=%d",
hourly_discovered,
hourly_failed,
hourly_protection,
fail_ratio,
protection_ratio,
_streak,
)
return {
"status": "anti_bot_detected",
"task_id": task_id,
"hourly_mode": result.get("hourly_mode"),
"discovered_urls": hourly_discovered,
"cars_failed": hourly_failed,
"protection_events": hourly_protection,
"hourly_failure_streak": _streak,
}
summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"status": result.get("status", "success"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
"failures_count": len(result.get("failures") or []),
"hourly_mode": result.get("hourly_mode"),
"discovered_urls": result.get("discovered_urls", 0),
"new_urls": result.get("new_urls", 0),
"sold_marked": result.get("sold_marked", 0),
}
logger.info(
"sync_listing_task hourly diff completed: status=%s, new=%d, sold=%d, skipped=%d",
summary["status"],
summary["new_urls"],
summary["sold_marked"],
summary["skipped_existing"],
)
# Hourly успешно — сбрасываем circuit breaker streak.
_clear_hourly_failure_streak(redis_client)
return summary
_clear_followup_pending(redis_client)
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)
and make is None
and model is None
and not prefer_sitemap_mainline
and discovery_mode != "sitemap"
)
if prefer_sitemap_mainline:
if discovery_mode != "sitemap":
logger.warning(
"Unfiltered full scan forcing sitemap discovery despite IAAI_DISCOVERY_MODE=%s",
discovery_mode or "unset",
)
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.
already_completed: set[int] = set()
cycle_id: str | None = None
if force_bootstrap_full_scan:
try:
cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments))
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
if created_new_cycle:
logger.warning("Started new bootstrap cycle: %s", cycle_id)
except Exception:
already_completed = set()
cycle_id = None
pending = [
(idx, seg) for idx, seg in enumerate(segments)
if idx not in already_completed
]
if not pending:
# Всё уже сделано — фиксируем bootstrap done.
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_segments_progress_state(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("Parallel segments: nothing to dispatch (all completed)")
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": 0,
"segments_already_completed": len(already_completed),
}
dispatched = 0
for idx, seg in pending:
try:
self.app.send_task(
"iaai_scraper.worker.tasks.sync_segment_task",
kwargs={
"segment_index": idx,
"segment": seg,
"lane": lane,
"only_new": effective_only_new,
"is_bootstrap": force_bootstrap_full_scan,
},
queue="scraping",
)
dispatched += 1
except Exception:
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
logger.warning(
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)",
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id,
)
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": dispatched,
"segments_already_completed": len(already_completed),
"segments_cycle_id": cycle_id,
}
# --- конец параллельной ветки ---
precomputed_result = None
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),
)
resume_from_segment = 0
_clear_sync_checkpoint(redis_client)
elif resume_from_segment > 0:
logger.warning(
"Resuming segmented bootstrap from segment=%d (last completed=%d)",
resume_from_segment, last_completed_segment,
)
checkpoint_streak, should_skip_from_checkpoint = _register_repeated_bootstrap_checkpoint(
redis_client,
last_completed_segment,
)
if should_skip_from_checkpoint:
skipped_segment = resume_from_segment
logger.error(
"Checkpoint %d repeated (streak=%d); advancing past segment=%d before resume",
last_completed_segment,
checkpoint_streak,
skipped_segment,
)
_save_last_completed_segment(redis_client, skipped_segment)
_update_task_progress(
redis_client,
task_id=task_id,
stage="bootstrap_repeat_checkpoint_segment_skipped",
ttl_seconds=progress_ttl,
checkpoint_segment=last_completed_segment,
segment_index=skipped_segment,
checkpoint_streak=checkpoint_streak,
)
resume_from_segment = skipped_segment + 1
if resume_from_segment >= len(segments):
precomputed_result = {
"status": "partial_success",
"run_id": None,
"full_scan_completed": True,
"cars_upserted": 0,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"elapsed_seconds": 0,
"failures": [{
"vehicle_url": f"segment_{skipped_segment}",
"error": f"Skipped after repeated bootstrap checkpoint loops on segment {skipped_segment}",
}],
}
else:
logger.warning(
"Bootstrap resume will continue from segment=%d after repeated checkpoint skip of segment=%d",
resume_from_segment,
skipped_segment,
)
if (
force_bootstrap_full_scan
and use_segmented
and stalled_segment is not None
and 0 <= stalled_segment < len(segments)
and stalled_segment >= resume_from_segment
):
logger.error(
"Skipping previously stalled segment=%d and advancing checkpoint before resume",
stalled_segment,
)
_save_last_completed_segment(redis_client, stalled_segment)
_clear_stalled_segment(redis_client)
resume_from_segment = stalled_segment + 1
_update_task_progress(
redis_client,
task_id=task_id,
stage="bootstrap_stalled_segment_skipped",
ttl_seconds=progress_ttl,
segment_index=stalled_segment,
)
if force_bootstrap_full_scan and use_segmented and resume_from_segment < len(segments):
resume_streak, should_skip_segment = _register_bootstrap_resume_target(
redis_client,
resume_from_segment,
)
if should_skip_segment:
skipped_segment = resume_from_segment
logger.error(
"Segment %d is stuck on bootstrap resume (streak=%d); skipping it and advancing checkpoint",
skipped_segment,
resume_streak,
)
_save_last_completed_segment(redis_client, skipped_segment)
_update_task_progress(
redis_client,
task_id=task_id,
stage="bootstrap_resume_segment_skipped",
ttl_seconds=progress_ttl,
segment_index=skipped_segment,
resume_streak=resume_streak,
)
resume_from_segment = skipped_segment + 1
if resume_from_segment >= len(segments):
precomputed_result = {
"status": "partial_success",
"run_id": None,
"full_scan_completed": True,
"cars_upserted": 0,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"elapsed_seconds": 0,
"failures": [{
"vehicle_url": f"segment_{skipped_segment}",
"error": f"Skipped after repeated bootstrap resume loops on segment {skipped_segment}",
}],
}
else:
logger.warning(
"Bootstrap resume will continue from segment=%d after skipping segment=%d",
resume_from_segment,
skipped_segment,
)
def _progress_cb_main(stage, meta):
_update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
**meta,
)
def _job():
with IAAIScraper() as scraper:
scraper.set_progress_callback(_progress_cb_main)
if use_segmented:
return scraper.sync_listing_segmented(
segments=segments,
lane=lane,
only_new=effective_only_new,
start_segment=resume_from_segment,
start_page=1,
progress_callback=(
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
if force_bootstrap_full_scan else None
),
)
return scraper.sync_listing(
make=make,
model=model,
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
)
result = precomputed_result if precomputed_result is not None else _run_browser_job(_job)
if force_bootstrap_full_scan:
bootstrap_completed = bool(result.get("full_scan_completed"))
if bootstrap_completed:
_set_full_scan_done(redis_client, False if always_full_scan else True)
_clear_sync_checkpoint(redis_client)
_clear_segments_progress_state(redis_client)
_clear_bootstrap_resume_target(redis_client)
_clear_stalled_segment(redis_client)
_clear_bootstrap_failure_streak(redis_client)
_clear_bootstrap_continuation_streak(redis_client)
if always_full_scan:
logger.info("Full scan completed; keeping bootstrap mode for next run (always full scan enabled)")
else:
logger.info("Bootstrap full scan completed; hourly schedule continues")
else:
_set_full_scan_done(redis_client, False)
listing_payload = result.get("listing") if isinstance(result.get("listing"), dict) else {}
anti_bot_detected = bool(result.get("anti_bot_detected"))
had_progress = any(
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
count_as_failure = anti_bot_detected or (str(result.get("status") or "") == "failed" and not had_progress)
followup_delay = max(3600, int(settings.celery.beat_sync_interval_minutes) * 60)
followup_reason = "bootstrap_not_completed"
if anti_bot_detected:
followup_delay = 180
followup_reason = "bootstrap_anti_bot_detected"
logger.warning(
"Bootstrap anti-bot guard: protection_events=%s fail_ratio=%s protection_ratio=%s; scheduling delayed continuation",
result.get("protection_events"),
result.get("fail_ratio"),
result.get("protection_ratio"),
)
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,
)
else:
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_resume_target(redis_client)
summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"status": result.get("status", "success"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"protection_events": result.get("protection_events", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
"failures_count": len(result.get("failures") or []),
}
if not force_bootstrap_full_scan:
discovered = int(result.get("total_discovered") or result.get("total") or 0)
failed = int(summary["cars_failed"] or 0)
protection = int(summary["protection_events"] or 0)
if discovered > 0:
fail_ratio = failed / max(1, discovered)
protection_ratio = protection / max(1, discovered)
anti_bot_suspected = (
(protection >= 30 and protection_ratio >= 0.10)
or fail_ratio >= 0.30
)
if anti_bot_suspected:
streak = _bump_hourly_failure_streak(redis_client)
logger.error(
"Hourly anti-bot guard triggered (listing fallback): discovered=%d failed=%d protection=%d fail_ratio=%.2f protection_ratio=%.2f streak=%d",
discovered,
failed,
protection,
fail_ratio,
protection_ratio,
streak,
)
summary["status"] = "anti_bot_detected"
summary["hourly_failure_streak"] = streak
return summary
logger.info(
"sync_listing_task completed: status=%s, %d upserted, %d failed, failures=%d",
summary["status"],
summary["cars_upserted"],
summary["cars_failed"],
summary["failures_count"],
)
return summary
except SoftTimeLimitExceeded:
logger.warning(
"sync_listing_task soft timeout exceeded — partial progress already saved to DB"
)
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
else:
# Hourly: ставим продолжение, но только если circuit breaker не открыт.
_bump_hourly_failure_streak(redis_client)
_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
if not _cb_open:
followup_ttl = max(600, int(lock_ttl))
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
try:
self.app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task",
kwargs={
"make": make,
"model": model,
"lane": lane,
"limit": limit,
"only_new": only_new,
},
queue="scraping",
countdown=10,
expires=followup_ttl,
)
logger.info("Queued immediate continuation after soft timeout")
except Exception:
_clear_followup_pending(redis_client)
logger.warning("Failed to queue continuation after soft timeout", exc_info=True)
else:
logger.info("Continuation after soft timeout already pending; skip duplicate enqueue")
else:
logger.warning("Skipping continuation: hourly circuit breaker open (%d failures)", _streak)
return {
"status": "timed_out",
"task_id": task_id,
"reason": "soft_time_limit_exceeded",
"note": "partial progress saved to DB; continuation queued",
}
except Exception as exc:
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
# Hourly circuit breaker: фиксируем ошибку.
if not force_bootstrap_full_scan:
_bump_hourly_failure_streak(redis_client)
try:
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
raise self.retry(exc=exc, countdown=5)
raise self.retry(exc=exc)
except self.MaxRetriesExceededError:
logger.error("sync_listing_task max retries exceeded, giving up")
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 при ошибках
# приводит к молотилке запросов и бану.
return {
"status": "failed",
"task_id": task_id,
"error": str(exc),
}
finally:
if watchdog_stop is not None:
watchdog_stop.set()
if watchdog_thread is not None:
watchdog_thread.join(timeout=5)
if heartbeat_stop is not None:
heartbeat_stop.set()
if heartbeat_thread is not None:
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
if lock_acquired:
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, owner_token)