Compare commits
2 Commits
05d1ffb5ac
...
7b7efbf525
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7b7efbf525 | ||
|
|
a8c5f8aa41 |
@@ -40,6 +40,7 @@ x-worker-service: &worker-service
|
|||||||
volumes:
|
volumes:
|
||||||
- ./runtime_config.json:/app/runtime_config.json:ro
|
- ./runtime_config.json:/app/runtime_config.json:ro
|
||||||
- tokens_data:/data
|
- tokens_data:/data
|
||||||
|
- beat_data:/var/lib/celery
|
||||||
|
|
||||||
services:
|
services:
|
||||||
# PostgreSQL
|
# PostgreSQL
|
||||||
@@ -146,11 +147,11 @@ services:
|
|||||||
--pidfile=/tmp/celery-worker.pid
|
--pidfile=/tmp/celery-worker.pid
|
||||||
-Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
|
-Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
|
||||||
healthcheck:
|
healthcheck:
|
||||||
test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"]
|
test: ["CMD-SHELL", "celery -A iaai_scraper.worker.celery_app inspect ping -d celery@$$HOSTNAME -t 15 >/dev/null 2>&1 || exit 1"]
|
||||||
interval: 60s
|
interval: 60s
|
||||||
timeout: 10s
|
timeout: 20s
|
||||||
retries: 3
|
retries: 3
|
||||||
start_period: 90s
|
start_period: 120s
|
||||||
logging:
|
logging:
|
||||||
driver: json-file
|
driver: json-file
|
||||||
options:
|
options:
|
||||||
@@ -172,7 +173,7 @@ services:
|
|||||||
celery -A iaai_scraper.worker.celery_app beat
|
celery -A iaai_scraper.worker.celery_app beat
|
||||||
--loglevel=info
|
--loglevel=info
|
||||||
--pidfile=/tmp/celery-beat.pid
|
--pidfile=/tmp/celery-beat.pid
|
||||||
--schedule=/tmp/celerybeat-schedule
|
--schedule=/var/lib/celery/celerybeat-schedule
|
||||||
healthcheck:
|
healthcheck:
|
||||||
test: ["CMD-SHELL", "test -f /tmp/celery-beat.pid && kill -0 $(cat /tmp/celery-beat.pid)"]
|
test: ["CMD-SHELL", "test -f /tmp/celery-beat.pid && kill -0 $(cat /tmp/celery-beat.pid)"]
|
||||||
interval: 60s
|
interval: 60s
|
||||||
@@ -189,3 +190,4 @@ volumes:
|
|||||||
pgdata:
|
pgdata:
|
||||||
redisdata:
|
redisdata:
|
||||||
tokens_data:
|
tokens_data:
|
||||||
|
beat_data:
|
||||||
|
|||||||
@@ -43,16 +43,37 @@ start_proxy_bridge_if_needed() {
|
|||||||
fi
|
fi
|
||||||
|
|
||||||
echo "[entrypoint] Using browser proxy: ${IAAI_PROXY_SERVER}"
|
echo "[entrypoint] Using browser proxy: ${IAAI_PROXY_SERVER}"
|
||||||
python -m iaai_scraper.proxy_bridge &
|
|
||||||
BRIDGE_PID=$!
|
# Watchdog: автоматически рестартует bridge при падении.
|
||||||
|
# Нужно для долгоиграющих контейнеров (месяцы аптайма): без watchdog
|
||||||
|
# упавший bridge тихо ломает весь скрапинг и IAAI начинает банить реальный IP VPS.
|
||||||
|
(
|
||||||
|
while true; do
|
||||||
|
python -m iaai_scraper.proxy_bridge
|
||||||
|
EXIT_CODE=$?
|
||||||
|
echo "[entrypoint] proxy_bridge exited with code ${EXIT_CODE}; restarting in 3s..."
|
||||||
|
sleep 3
|
||||||
|
done
|
||||||
|
) &
|
||||||
|
BRIDGE_WATCHDOG_PID=$!
|
||||||
sleep 1
|
sleep 1
|
||||||
|
|
||||||
if ! kill -0 "${BRIDGE_PID}" 2>/dev/null; then
|
if ! kill -0 "${BRIDGE_WATCHDOG_PID}" 2>/dev/null; then
|
||||||
echo "[entrypoint] ERROR: iaai_scraper.proxy_bridge failed to start"
|
echo "[entrypoint] ERROR: proxy_bridge watchdog failed to start"
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
echo "[entrypoint] Proxy bridge started (PID ${BRIDGE_PID})"
|
# Проверяем что bridge действительно слушает порт.
|
||||||
|
for i in 1 2 3 4 5; do
|
||||||
|
if python -c "import socket,sys; s=socket.socket(); s.settimeout(2); sys.exit(0 if s.connect_ex(('127.0.0.1', 8899))==0 else 1)" 2>/dev/null; then
|
||||||
|
echo "[entrypoint] Proxy bridge listening on :8899 (watchdog PID ${BRIDGE_WATCHDOG_PID})"
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
sleep 1
|
||||||
|
done
|
||||||
|
|
||||||
|
echo "[entrypoint] ERROR: proxy_bridge did not start listening on :8899"
|
||||||
|
exit 1
|
||||||
}
|
}
|
||||||
|
|
||||||
start_proxy_bridge_if_needed "$@"
|
start_proxy_bridge_if_needed "$@"
|
||||||
|
|||||||
@@ -219,6 +219,8 @@ class CeleryConfig:
|
|||||||
task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600)
|
task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600)
|
||||||
worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4)
|
worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4)
|
||||||
worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5)
|
worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5)
|
||||||
|
# KB; worker перезапустит child при превышении (защита от утечек Chromium/Playwright).
|
||||||
|
worker_max_memory_per_child_kb: int = _env_int("CELERY_WORKER_MAX_MEMORY_PER_CHILD_KB", 500_000)
|
||||||
broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200)
|
broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200)
|
||||||
beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60)
|
beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60)
|
||||||
beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None
|
beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None
|
||||||
|
|||||||
@@ -12,3 +12,7 @@ class SiteStructureChangedError(ScraperError):
|
|||||||
|
|
||||||
class ListingResumeError(ScraperError):
|
class ListingResumeError(ScraperError):
|
||||||
"""Вызывается, когда resume по checkpoint больше недостижим."""
|
"""Вызывается, когда resume по checkpoint больше недостижим."""
|
||||||
|
|
||||||
|
|
||||||
|
class ScraperAbortedError(ScraperError):
|
||||||
|
"""Вызывается при внешнем прерывании (например, потеря Celery lock'а)."""
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import logging
|
import logging
|
||||||
import sys
|
import sys
|
||||||
from contextvars import ContextVar
|
from contextvars import ContextVar
|
||||||
|
from logging.handlers import RotatingFileHandler
|
||||||
|
|
||||||
# Храним trace_id текущего потока/корутины.
|
# Храним trace_id текущего потока/корутины.
|
||||||
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
|
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
|
||||||
@@ -21,7 +22,16 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
|
|||||||
# stderr — Docker и Celery prefork корректно его подхватывают.
|
# stderr — Docker и Celery prefork корректно его подхватывают.
|
||||||
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)]
|
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)]
|
||||||
if log_file:
|
if log_file:
|
||||||
handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
|
# Ротация файлового лога: 50MB × 5 файлов, чтобы диск VPS не переполнился
|
||||||
|
# при многомесячной работе.
|
||||||
|
handlers.append(
|
||||||
|
RotatingFileHandler(
|
||||||
|
log_file,
|
||||||
|
maxBytes=50 * 1024 * 1024,
|
||||||
|
backupCount=5,
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
)
|
||||||
trace_filter = TraceIdFilter()
|
trace_filter = TraceIdFilter()
|
||||||
for handler in handlers:
|
for handler in handlers:
|
||||||
handler.addFilter(trace_filter)
|
handler.addFilter(trace_filter)
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
|
|||||||
|
|
||||||
from .browser import BrowserFactory, HumanPacer, NetworkCapture
|
from .browser import BrowserFactory, HumanPacer, NetworkCapture
|
||||||
from .core.config import Settings, settings, parse_listing_segments
|
from .core.config import Settings, settings, parse_listing_segments
|
||||||
from .core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError
|
from .core.exceptions import AntiBotDetectedError, ListingResumeError, ScraperAbortedError, SiteStructureChangedError
|
||||||
from .core.logs import set_trace_id, setup_logging
|
from .core.logs import set_trace_id, setup_logging
|
||||||
from .core.retry import retryable
|
from .core.retry import retryable
|
||||||
from .core.runtime_config import RuntimeConfig
|
from .core.runtime_config import RuntimeConfig
|
||||||
@@ -120,6 +120,9 @@ class IAAIScraper:
|
|||||||
self.car_mapper = CarMapper()
|
self.car_mapper = CarMapper()
|
||||||
self.persistence = PersistenceService(self.settings)
|
self.persistence = PersistenceService(self.settings)
|
||||||
self._shutdown_requested = False
|
self._shutdown_requested = False
|
||||||
|
# Внешний abort-флаг (например, Celery heartbeat при потере lock'а).
|
||||||
|
# Проверяется в безопасных точках долгих циклов (между сегментами/страницами).
|
||||||
|
self._abort_callback: Callable[[], bool] | None = None
|
||||||
proxy_url = self.settings.proxy.server
|
proxy_url = self.settings.proxy.server
|
||||||
if proxy_url:
|
if proxy_url:
|
||||||
_proxy_kwargs: dict = {
|
_proxy_kwargs: dict = {
|
||||||
@@ -148,6 +151,28 @@ class IAAIScraper:
|
|||||||
set_trace_id(trace_id)
|
set_trace_id(trace_id)
|
||||||
return trace_id
|
return trace_id
|
||||||
|
|
||||||
|
def set_abort_callback(self, callback: Callable[[], bool] | None) -> None:
|
||||||
|
"""Устанавливает внешний abort-флаг (callable → bool).
|
||||||
|
|
||||||
|
Если callback вернёт True, при следующей проверке
|
||||||
|
``_check_abort()`` поднимется ``ScraperAbortedError``.
|
||||||
|
Используется Celery worker'ом, чтобы остановить sync_listing
|
||||||
|
при потере распределённого lock'а.
|
||||||
|
"""
|
||||||
|
self._abort_callback = callback
|
||||||
|
|
||||||
|
def _check_abort(self, where: str) -> None:
|
||||||
|
cb = self._abort_callback
|
||||||
|
if cb is None:
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
should_abort = bool(cb())
|
||||||
|
except Exception:
|
||||||
|
should_abort = False
|
||||||
|
if should_abort:
|
||||||
|
logger.warning("Abort signal received at %s — stopping sync gracefully", where)
|
||||||
|
raise ScraperAbortedError(f"Aborted at {where}")
|
||||||
|
|
||||||
def __enter__(self) -> "IAAIScraper":
|
def __enter__(self) -> "IAAIScraper":
|
||||||
if self.playwright is None:
|
if self.playwright is None:
|
||||||
self.playwright = sync_playwright().start()
|
self.playwright = sync_playwright().start()
|
||||||
@@ -216,7 +241,11 @@ class IAAIScraper:
|
|||||||
except PlaywrightTimeoutError:
|
except PlaywrightTimeoutError:
|
||||||
pass
|
pass
|
||||||
time.sleep(0.3)
|
time.sleep(0.3)
|
||||||
logger.info("Warmup done: %s (title=%s)", page.url, page.title()[:50])
|
try:
|
||||||
|
title = page.title()
|
||||||
|
except Exception:
|
||||||
|
title = ""
|
||||||
|
logger.info("Warmup done: %s (title=%s)", page.url, (title or "")[:50])
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("Warmup visit failed: %s — continuing anyway", e)
|
logger.warning("Warmup visit failed: %s — continuing anyway", e)
|
||||||
time.sleep(1.5)
|
time.sleep(1.5)
|
||||||
@@ -367,80 +396,25 @@ class IAAIScraper:
|
|||||||
logger.warning("Listing resumed at page %d after recovery", target_page_number)
|
logger.warning("Listing resumed at page %d after recovery", target_page_number)
|
||||||
return page, applied_filters
|
return page, applied_filters
|
||||||
|
|
||||||
def _reset_listing_progress_checkpoint(
|
|
||||||
self,
|
|
||||||
progress_callback: Callable[[int], None] | None,
|
|
||||||
*,
|
|
||||||
target_page_number: int,
|
|
||||||
reason: str,
|
|
||||||
) -> None:
|
|
||||||
if progress_callback is None:
|
|
||||||
return
|
|
||||||
try:
|
|
||||||
# page_number=0 — специальный sentinel: следующий запуск должен начать scope с page 1.
|
|
||||||
progress_callback(0)
|
|
||||||
logger.warning(
|
|
||||||
"Reset listing checkpoint after failed resume to page %d: %s",
|
|
||||||
target_page_number,
|
|
||||||
reason,
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
logger.warning(
|
|
||||||
"Failed to reset listing checkpoint after resume failure to page %d",
|
|
||||||
target_page_number,
|
|
||||||
exc_info=True,
|
|
||||||
)
|
|
||||||
|
|
||||||
def _open_listing_for_stream(
|
def _open_listing_for_stream(
|
||||||
self,
|
self,
|
||||||
*,
|
*,
|
||||||
start_page: int,
|
|
||||||
make: str | None,
|
make: str | None,
|
||||||
model: str | None,
|
model: str | None,
|
||||||
progress_callback: Callable[[int], None] | None,
|
|
||||||
listing_url: str | None = None,
|
listing_url: str | None = None,
|
||||||
year_min: int | None = None,
|
year_min: int | None = None,
|
||||||
year_max: int | None = None,
|
year_max: int | None = None,
|
||||||
) -> tuple[Page, dict[str, str | int | None], int]:
|
) -> tuple[Page, dict[str, str | int | None]]:
|
||||||
if start_page <= 1:
|
# Всегда открываем с page 1. Page-level resume убран: эфемерные URL и короткие
|
||||||
page, applied_filters = self._open_listing_page(
|
# сегменты делали pagination-resume хрупким. Bootstrap прогресс сохраняется
|
||||||
make=make,
|
# на уровне сегментов в worker/tasks.py (_save_last_completed_segment).
|
||||||
model=model,
|
return self._open_listing_page(
|
||||||
listing_url=listing_url,
|
make=make,
|
||||||
year_min=year_min,
|
model=model,
|
||||||
year_max=year_max,
|
listing_url=listing_url,
|
||||||
)
|
year_min=year_min,
|
||||||
return page, applied_filters, 1
|
year_max=year_max,
|
||||||
|
)
|
||||||
try:
|
|
||||||
page, applied_filters = self._reopen_listing_and_resume(
|
|
||||||
target_page_number=start_page,
|
|
||||||
make=make,
|
|
||||||
model=model,
|
|
||||||
listing_url=listing_url,
|
|
||||||
year_min=year_min,
|
|
||||||
year_max=year_max,
|
|
||||||
)
|
|
||||||
return page, applied_filters, start_page
|
|
||||||
except ListingResumeError as resume_exc:
|
|
||||||
logger.warning(
|
|
||||||
"Stored checkpoint page %d is unreachable for current listing scope; restarting scope from page 1: %s",
|
|
||||||
start_page,
|
|
||||||
resume_exc,
|
|
||||||
)
|
|
||||||
self._reset_listing_progress_checkpoint(
|
|
||||||
progress_callback,
|
|
||||||
target_page_number=start_page,
|
|
||||||
reason=str(resume_exc),
|
|
||||||
)
|
|
||||||
page, applied_filters = self._open_listing_page(
|
|
||||||
make=make,
|
|
||||||
model=model,
|
|
||||||
listing_url=listing_url,
|
|
||||||
year_min=year_min,
|
|
||||||
year_max=year_max,
|
|
||||||
)
|
|
||||||
return page, applied_filters, 1
|
|
||||||
|
|
||||||
def collect_listing(
|
def collect_listing(
|
||||||
self,
|
self,
|
||||||
@@ -476,8 +450,6 @@ class IAAIScraper:
|
|||||||
limit: int | None,
|
limit: int | None,
|
||||||
effective_only_new: bool,
|
effective_only_new: bool,
|
||||||
started_at: float,
|
started_at: float,
|
||||||
start_page: int = 1,
|
|
||||||
progress_callback: Callable[[int], None] | None = None,
|
|
||||||
listing_url: str | None = None,
|
listing_url: str | None = None,
|
||||||
year_min: int | None = None,
|
year_min: int | None = None,
|
||||||
year_max: int | None = None,
|
year_max: int | None = None,
|
||||||
@@ -619,17 +591,16 @@ class IAAIScraper:
|
|||||||
|
|
||||||
page = None
|
page = None
|
||||||
try:
|
try:
|
||||||
page, applied_filters, effective_start_page = self._open_listing_for_stream(
|
page, applied_filters = self._open_listing_for_stream(
|
||||||
start_page=start_page,
|
|
||||||
make=make,
|
make=make,
|
||||||
model=model,
|
model=model,
|
||||||
progress_callback=progress_callback,
|
|
||||||
listing_url=listing_url,
|
listing_url=listing_url,
|
||||||
year_min=year_min,
|
year_min=year_min,
|
||||||
year_max=year_max,
|
year_max=year_max,
|
||||||
)
|
)
|
||||||
|
|
||||||
for page_number in range(effective_start_page, max(effective_start_page, self.settings.listing.max_pages_per_run) + 1):
|
for page_number in range(1, self.settings.listing.max_pages_per_run + 1):
|
||||||
|
self._check_abort(f"listing page {page_number}")
|
||||||
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
|
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
|
||||||
pages_info.append({
|
pages_info.append({
|
||||||
"page_number": page_result.page_number,
|
"page_number": page_result.page_number,
|
||||||
@@ -673,14 +644,9 @@ class IAAIScraper:
|
|||||||
)
|
)
|
||||||
except ListingResumeError as resume_exc:
|
except ListingResumeError as resume_exc:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Cannot resume at page %d (%s) — resetting checkpoint and stopping pagination for this segment",
|
"Cannot resume at page %d (%s) — stopping pagination for this segment",
|
||||||
page_number, resume_exc,
|
page_number, resume_exc,
|
||||||
)
|
)
|
||||||
self._reset_listing_progress_checkpoint(
|
|
||||||
progress_callback,
|
|
||||||
target_page_number=page_number,
|
|
||||||
reason=str(resume_exc),
|
|
||||||
)
|
|
||||||
page = None
|
page = None
|
||||||
page_urls = []
|
page_urls = []
|
||||||
if not page_urls:
|
if not page_urls:
|
||||||
@@ -750,12 +716,6 @@ class IAAIScraper:
|
|||||||
break
|
break
|
||||||
pending_urls = pending_urls[batch_size:]
|
pending_urls = pending_urls[batch_size:]
|
||||||
|
|
||||||
if progress_callback is not None:
|
|
||||||
try:
|
|
||||||
progress_callback(page_number)
|
|
||||||
except Exception:
|
|
||||||
logger.warning("Failed to persist progress callback for page %d", page_number, exc_info=True)
|
|
||||||
|
|
||||||
if limit is not None and limit > 0 and total >= limit:
|
if limit is not None and limit > 0 and total >= limit:
|
||||||
break
|
break
|
||||||
if (
|
if (
|
||||||
@@ -779,14 +739,9 @@ class IAAIScraper:
|
|||||||
)
|
)
|
||||||
except ListingResumeError as resume_exc:
|
except ListingResumeError as resume_exc:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Cannot resume at page %d (%s) — resetting checkpoint and treating pagination as exhausted",
|
"Cannot resume at page %d (%s) — treating pagination as exhausted",
|
||||||
page_number + 1, resume_exc,
|
page_number + 1, resume_exc,
|
||||||
)
|
)
|
||||||
self._reset_listing_progress_checkpoint(
|
|
||||||
progress_callback,
|
|
||||||
target_page_number=page_number + 1,
|
|
||||||
reason=str(resume_exc),
|
|
||||||
)
|
|
||||||
break
|
break
|
||||||
|
|
||||||
if pending_urls:
|
if pending_urls:
|
||||||
@@ -1559,8 +1514,6 @@ class IAAIScraper:
|
|||||||
lane: str = "iaai_cars",
|
lane: str = "iaai_cars",
|
||||||
limit: int | None = None,
|
limit: int | None = None,
|
||||||
only_new: bool | None = None,
|
only_new: bool | None = None,
|
||||||
start_page: int = 1,
|
|
||||||
progress_callback: Callable[[int], None] | None = None,
|
|
||||||
listing_url: str | None = None,
|
listing_url: str | None = None,
|
||||||
year_min: int | None = None,
|
year_min: int | None = None,
|
||||||
year_max: int | None = None,
|
year_max: int | None = None,
|
||||||
@@ -1597,8 +1550,6 @@ class IAAIScraper:
|
|||||||
limit=limit,
|
limit=limit,
|
||||||
effective_only_new=effective_only_new,
|
effective_only_new=effective_only_new,
|
||||||
started_at=started_at,
|
started_at=started_at,
|
||||||
start_page=max(1, int(start_page)),
|
|
||||||
progress_callback=progress_callback,
|
|
||||||
listing_url=listing_url,
|
listing_url=listing_url,
|
||||||
year_min=year_min,
|
year_min=year_min,
|
||||||
year_max=year_max,
|
year_max=year_max,
|
||||||
@@ -1677,19 +1628,20 @@ class IAAIScraper:
|
|||||||
only_new: bool | None = None,
|
only_new: bool | None = None,
|
||||||
start_segment: int = 0,
|
start_segment: int = 0,
|
||||||
start_page: int = 1,
|
start_page: int = 1,
|
||||||
progress_callback: Callable[[int, int], None] | None = None,
|
progress_callback: Callable[[int], None] | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
"""Итеративный sync_listing по списку сегментов (бренд / бренд+годы).
|
"""Итеративный sync_listing по списку сегментов (бренд / бренд+годы).
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
segments: список dict с ключами make, year_min, year_max.
|
segments: список dict с ключами make, year_min, year_max.
|
||||||
start_segment: индекс сегмента для resume (0-based).
|
start_segment: индекс сегмента для resume (0-based).
|
||||||
start_page: страница внутри start_segment для resume.
|
start_page: не используется (остаётся для обратной совместимости API).
|
||||||
progress_callback: вызывается (segment_index, page_number) после каждой страницы.
|
progress_callback: вызывается (segment_index) после каждого ПОЛНОСТЬЮ пройденного сегмента.
|
||||||
"""
|
"""
|
||||||
trace_id = self._new_trace_id("sync-segmented")
|
trace_id = self._new_trace_id("sync-segmented")
|
||||||
started_at = time.perf_counter()
|
started_at = time.perf_counter()
|
||||||
base_url = self.settings.listing.cars_url
|
base_url = self.settings.listing.cars_url
|
||||||
|
_ = start_page # обратная совместимость — page-level resume удалён
|
||||||
|
|
||||||
total_cars_upserted = 0
|
total_cars_upserted = 0
|
||||||
total_cars_failed = 0
|
total_cars_failed = 0
|
||||||
@@ -1701,39 +1653,33 @@ class IAAIScraper:
|
|||||||
completed_all = True
|
completed_all = True
|
||||||
|
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Starting segmented sync: %d segments, resume from segment=%d page=%d",
|
"Starting segmented sync: %d segments, resume from segment=%d",
|
||||||
len(segments), start_segment, start_page,
|
len(segments), start_segment,
|
||||||
)
|
)
|
||||||
|
|
||||||
for seg_idx in range(start_segment, len(segments)):
|
for seg_idx in range(start_segment, len(segments)):
|
||||||
|
self._check_abort(f"segment {seg_idx + 1}/{len(segments)}")
|
||||||
seg = segments[seg_idx]
|
seg = segments[seg_idx]
|
||||||
seg_make = seg.get("make")
|
seg_make = seg.get("make")
|
||||||
seg_year_min = seg.get("year_min")
|
seg_year_min = seg.get("year_min")
|
||||||
seg_year_max = seg.get("year_max")
|
seg_year_max = seg.get("year_max")
|
||||||
seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None
|
seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None
|
||||||
|
|
||||||
seg_start_page = start_page if seg_idx == start_segment else 1
|
|
||||||
seg_label = f"{seg_make or 'ALL'}"
|
seg_label = f"{seg_make or 'ALL'}"
|
||||||
if seg_year_min is not None or seg_year_max is not None:
|
if seg_year_min is not None or seg_year_max is not None:
|
||||||
seg_label += f" ({seg_year_min}-{seg_year_max})"
|
seg_label += f" ({seg_year_min}-{seg_year_max})"
|
||||||
|
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Segment %d/%d: %s (start_page=%d)",
|
"Segment %d/%d: %s",
|
||||||
seg_idx + 1, len(segments), seg_label, seg_start_page,
|
seg_idx + 1, len(segments), seg_label,
|
||||||
)
|
)
|
||||||
|
|
||||||
def _seg_progress(page_number: int, _si=seg_idx) -> None:
|
|
||||||
if progress_callback:
|
|
||||||
progress_callback(_si, page_number)
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
result = self.sync_listing(
|
result = self.sync_listing(
|
||||||
make=None if seg_url else seg_make,
|
make=None if seg_url else seg_make,
|
||||||
model=None,
|
model=None,
|
||||||
lane=lane,
|
lane=lane,
|
||||||
only_new=only_new,
|
only_new=only_new,
|
||||||
start_page=seg_start_page,
|
|
||||||
progress_callback=_seg_progress,
|
|
||||||
listing_url=seg_url,
|
listing_url=seg_url,
|
||||||
year_min=seg_year_min,
|
year_min=seg_year_min,
|
||||||
year_max=seg_year_max,
|
year_max=seg_year_max,
|
||||||
@@ -1754,9 +1700,20 @@ class IAAIScraper:
|
|||||||
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
|
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
|
||||||
})
|
})
|
||||||
|
|
||||||
if not bool(result.get("full_scan_completed", False)):
|
segment_done = bool(result.get("full_scan_completed", False))
|
||||||
|
if not segment_done:
|
||||||
completed_all = False
|
completed_all = False
|
||||||
|
|
||||||
|
# Сегмент полностью пройден — фиксируем чекпоинт, даже если часть машин failed.
|
||||||
|
if segment_done and progress_callback is not None:
|
||||||
|
try:
|
||||||
|
progress_callback(seg_idx)
|
||||||
|
except Exception:
|
||||||
|
logger.warning(
|
||||||
|
"Failed to persist segment checkpoint for segment=%d",
|
||||||
|
seg_idx, exc_info=True,
|
||||||
|
)
|
||||||
|
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d",
|
"Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d",
|
||||||
seg_idx + 1, len(segments), seg_label,
|
seg_idx + 1, len(segments), seg_label,
|
||||||
@@ -1764,6 +1721,12 @@ class IAAIScraper:
|
|||||||
result.get("cars_failed", 0),
|
result.get("cars_failed", 0),
|
||||||
result.get("listing", {}).get("vehicles_collected", 0),
|
result.get("listing", {}).get("vehicles_collected", 0),
|
||||||
)
|
)
|
||||||
|
except ScraperAbortedError:
|
||||||
|
# Внешний abort — выходим из цикла без отметки как failure.
|
||||||
|
# Незавершённый сегмент НЕ чекпоинтится, следующий запуск его повторит.
|
||||||
|
logger.warning("Segmented sync aborted at segment %d/%d", seg_idx + 1, len(segments))
|
||||||
|
completed_all = False
|
||||||
|
break
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error("Segment %d/%d failed: %s — %s", seg_idx + 1, len(segments), seg_label, exc)
|
logger.error("Segment %d/%d failed: %s — %s", seg_idx + 1, len(segments), seg_label, exc)
|
||||||
all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)})
|
all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)})
|
||||||
|
|||||||
@@ -27,6 +27,12 @@ def _on_worker_process_init(**kwargs):
|
|||||||
level = settings.log_level if settings.log_level else "INFO"
|
level = settings.log_level if settings.log_level else "INFO"
|
||||||
setup_logging(level, None)
|
setup_logging(level, None)
|
||||||
|
|
||||||
|
# Сбрасываем cached engine/redis после fork — иначе child наследует
|
||||||
|
# коннекты родителя, что небезопасно (особенно с psycopg2).
|
||||||
|
from . import tasks as _tasks
|
||||||
|
_tasks._CACHED_PERSISTENCE = None
|
||||||
|
_tasks._CACHED_REDIS = None
|
||||||
|
|
||||||
|
|
||||||
def _broker_url() -> str:
|
def _broker_url() -> str:
|
||||||
return settings.celery.broker_url or settings.redis.url
|
return settings.celery.broker_url or settings.redis.url
|
||||||
@@ -66,15 +72,17 @@ celery_app.conf.update(
|
|||||||
task_acks_late=True,
|
task_acks_late=True,
|
||||||
task_reject_on_worker_lost=True,
|
task_reject_on_worker_lost=True,
|
||||||
task_track_started=True,
|
task_track_started=True,
|
||||||
|
task_ignore_result=True,
|
||||||
worker_concurrency=settings.celery.worker_concurrency,
|
worker_concurrency=settings.celery.worker_concurrency,
|
||||||
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
|
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
|
||||||
|
worker_max_memory_per_child=settings.celery.worker_max_memory_per_child_kb,
|
||||||
worker_pool="prefork",
|
worker_pool="prefork",
|
||||||
worker_prefetch_multiplier=1,
|
worker_prefetch_multiplier=1,
|
||||||
broker_connection_retry_on_startup=True,
|
broker_connection_retry_on_startup=True,
|
||||||
broker_transport_options={
|
broker_transport_options={
|
||||||
"visibility_timeout": settings.celery.broker_visibility_timeout,
|
"visibility_timeout": settings.celery.broker_visibility_timeout,
|
||||||
},
|
},
|
||||||
result_expires=86400,
|
result_expires=3600,
|
||||||
worker_redirect_stdouts=False,
|
worker_redirect_stdouts=False,
|
||||||
worker_hijack_root_logger=False,
|
worker_hijack_root_logger=False,
|
||||||
beat_schedule={
|
beat_schedule={
|
||||||
@@ -107,7 +115,10 @@ def _on_worker_ready(**kwargs):
|
|||||||
health_check_interval=settings.redis.health_check_interval_seconds,
|
health_check_interval=settings.redis.health_check_interval_seconds,
|
||||||
retry_on_timeout=True,
|
retry_on_timeout=True,
|
||||||
)
|
)
|
||||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
# Дедуп TTL = интервал beat (по умолчанию 1ч), чтобы не плодить дубликаты
|
||||||
|
# при flapping-рестартах контейнера.
|
||||||
|
dedupe_ttl = max(600, settings.celery.beat_sync_interval_minutes * 60)
|
||||||
|
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=dedupe_ttl))
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
|
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
||||||
|
|
||||||
from concurrent.futures import ThreadPoolExecutor
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
import json
|
|
||||||
import logging
|
import logging
|
||||||
from threading import Event, Thread
|
from threading import Event, Thread
|
||||||
import time
|
import time
|
||||||
@@ -12,6 +11,7 @@ from celery import shared_task
|
|||||||
from redis import Redis
|
from redis import Redis
|
||||||
|
|
||||||
from ..core.config import Settings, parse_listing_segments
|
from ..core.config import Settings, parse_listing_segments
|
||||||
|
from ..core.exceptions import ScraperAbortedError
|
||||||
from ..scraper import IAAIScraper
|
from ..scraper import IAAIScraper
|
||||||
from ..storage.db import PersistenceService
|
from ..storage.db import PersistenceService
|
||||||
|
|
||||||
@@ -21,13 +21,17 @@ SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
|||||||
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||||
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
||||||
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
|
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
|
||||||
SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT = 2
|
|
||||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
|
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_LIMIT = 3
|
||||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
|
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
|
||||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||||
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
||||||
|
|
||||||
|
# Process-level caches (по одному на prefork-child) — пересоздаются при ошибках
|
||||||
|
# или при старте нового child (max_tasks_per_child).
|
||||||
|
_CACHED_PERSISTENCE: PersistenceService | None = None
|
||||||
|
_CACHED_REDIS: Redis | None = None
|
||||||
|
|
||||||
|
|
||||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||||
last_exc: Exception | None = None
|
last_exc: Exception | None = None
|
||||||
@@ -94,32 +98,75 @@ def _sync_listing_lock_ttl_seconds() -> int:
|
|||||||
|
|
||||||
|
|
||||||
def _get_persistence() -> PersistenceService:
|
def _get_persistence() -> PersistenceService:
|
||||||
settings = Settings()
|
# Кэшируем engine на уровне процесса — иначе каждая task создаёт новый
|
||||||
persistence = PersistenceService(settings)
|
# SQLAlchemy pool (5+10 коннектов), и за сутки worker исчерпает
|
||||||
|
# max_connections Postgres. С max_tasks_per_child=5 пул переиспользуется
|
||||||
|
# для всех 5 тасок, после чего child перезапустится и пул пересоздастся.
|
||||||
|
global _CACHED_PERSISTENCE
|
||||||
|
if _CACHED_PERSISTENCE is None:
|
||||||
|
settings = Settings()
|
||||||
|
_CACHED_PERSISTENCE = PersistenceService(settings)
|
||||||
|
|
||||||
|
persistence = _CACHED_PERSISTENCE
|
||||||
|
|
||||||
def _ping_db() -> None:
|
def _ping_db() -> None:
|
||||||
with persistence.engine.connect() as conn:
|
with persistence.engine.connect() as conn:
|
||||||
conn.exec_driver_sql("SELECT 1")
|
conn.exec_driver_sql("SELECT 1")
|
||||||
|
|
||||||
_retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0)
|
try:
|
||||||
|
_retry_with_backoff(_ping_db, attempts=3, base_delay_s=1.0)
|
||||||
|
except Exception:
|
||||||
|
# Engine мог стать невалидным (PG рестарт) — пересоздаём.
|
||||||
|
logger.warning("Cached DB engine ping failed; recreating engine", exc_info=True)
|
||||||
|
try:
|
||||||
|
_CACHED_PERSISTENCE.engine.dispose()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
_CACHED_PERSISTENCE = PersistenceService(Settings())
|
||||||
|
persistence = _CACHED_PERSISTENCE
|
||||||
|
_retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0)
|
||||||
return persistence
|
return persistence
|
||||||
|
|
||||||
|
|
||||||
def _get_redis() -> Redis:
|
def _get_redis() -> Redis:
|
||||||
settings = Settings()
|
# Кэшируем клиент на уровне процесса — Redis().from_url каждый раз создаёт
|
||||||
redis_client = Redis.from_url(
|
# отдельный connection pool. См. _get_persistence для аналогичного обоснования.
|
||||||
settings.redis.url,
|
global _CACHED_REDIS
|
||||||
decode_responses=True,
|
if _CACHED_REDIS is None:
|
||||||
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
|
settings = Settings()
|
||||||
socket_timeout=settings.redis.socket_timeout_seconds,
|
_CACHED_REDIS = Redis.from_url(
|
||||||
health_check_interval=settings.redis.health_check_interval_seconds,
|
settings.redis.url,
|
||||||
retry_on_timeout=True,
|
decode_responses=True,
|
||||||
)
|
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
|
||||||
|
socket_timeout=settings.redis.socket_timeout_seconds,
|
||||||
|
health_check_interval=settings.redis.health_check_interval_seconds,
|
||||||
|
retry_on_timeout=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
redis_client = _CACHED_REDIS
|
||||||
|
|
||||||
def _ping_redis() -> None:
|
def _ping_redis() -> None:
|
||||||
redis_client.ping()
|
redis_client.ping()
|
||||||
|
|
||||||
_retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0)
|
try:
|
||||||
|
_retry_with_backoff(_ping_redis, attempts=3, base_delay_s=1.0)
|
||||||
|
except Exception:
|
||||||
|
logger.warning("Cached Redis ping failed; recreating client", exc_info=True)
|
||||||
|
try:
|
||||||
|
_CACHED_REDIS.close()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
settings = Settings()
|
||||||
|
_CACHED_REDIS = Redis.from_url(
|
||||||
|
settings.redis.url,
|
||||||
|
decode_responses=True,
|
||||||
|
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
|
||||||
|
socket_timeout=settings.redis.socket_timeout_seconds,
|
||||||
|
health_check_interval=settings.redis.health_check_interval_seconds,
|
||||||
|
retry_on_timeout=True,
|
||||||
|
)
|
||||||
|
redis_client = _CACHED_REDIS
|
||||||
|
_retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0)
|
||||||
return redis_client
|
return redis_client
|
||||||
|
|
||||||
|
|
||||||
@@ -181,7 +228,10 @@ def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> N
|
|||||||
|
|
||||||
def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool:
|
def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool:
|
||||||
try:
|
try:
|
||||||
inspector = celery_app.control.inspect(timeout=1.0)
|
# 10s — под нагрузкой workers могут отвечать с задержкой; 1s давал
|
||||||
|
# ложноположительные "пусто" → удаление валидного lock и параллельный
|
||||||
|
# запуск sync_listing с дублированием запросов к IAAI.
|
||||||
|
inspector = celery_app.control.inspect(timeout=10.0)
|
||||||
snapshots = [
|
snapshots = [
|
||||||
inspector.active() or {},
|
inspector.active() or {},
|
||||||
inspector.reserved() or {},
|
inspector.reserved() or {},
|
||||||
@@ -210,6 +260,7 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_
|
|||||||
owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY)
|
owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY)
|
||||||
if not owner_token:
|
if not owner_token:
|
||||||
return False
|
return False
|
||||||
|
ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to read sync listing lock before cleanup", exc_info=True)
|
logger.warning("Failed to read sync listing lock before cleanup", exc_info=True)
|
||||||
return False
|
return False
|
||||||
@@ -218,8 +269,21 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_
|
|||||||
logger.info("sync_listing lock preserved: active task still detected")
|
logger.info("sync_listing lock preserved: active task still detected")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
# Защита от race: если активный owner_token изменился между inspect и delete —
|
||||||
|
# значит другой worker уже acquire'нул lock; удалять его нельзя.
|
||||||
|
try:
|
||||||
|
owner_token_after = redis_client.get(SYNC_LISTING_LOCK_KEY)
|
||||||
|
if owner_token_after != owner_token:
|
||||||
|
logger.info(
|
||||||
|
"sync_listing lock owner changed during cleanup (%s -> %s); preserving",
|
||||||
|
owner_token, owner_token_after,
|
||||||
|
)
|
||||||
|
return False
|
||||||
|
except Exception:
|
||||||
|
logger.warning("Failed to re-check sync listing lock owner before cleanup", exc_info=True)
|
||||||
|
return False
|
||||||
|
|
||||||
try:
|
try:
|
||||||
ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY)
|
|
||||||
redis_client.delete(SYNC_LISTING_LOCK_KEY)
|
redis_client.delete(SYNC_LISTING_LOCK_KEY)
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Removed orphan sync_listing lock owner=%s ttl=%s after worker restart",
|
"Removed orphan sync_listing lock owner=%s ttl=%s after worker restart",
|
||||||
@@ -248,55 +312,32 @@ def _set_full_scan_done(redis_client: Redis, done: bool) -> None:
|
|||||||
logger.warning("Failed to persist full scan state", exc_info=True)
|
logger.warning("Failed to persist full scan state", exc_info=True)
|
||||||
|
|
||||||
|
|
||||||
def _load_sync_checkpoint(redis_client: Redis) -> dict[str, object] | None:
|
def _load_last_completed_segment(redis_client: Redis) -> int | None:
|
||||||
|
# Segment-only checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||||
try:
|
try:
|
||||||
raw = redis_client.get(SYNC_LISTING_CHECKPOINT_KEY)
|
raw = redis_client.get(SYNC_LISTING_CHECKPOINT_KEY)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to read sync listing checkpoint", exc_info=True)
|
logger.warning("Failed to read sync listing checkpoint", exc_info=True)
|
||||||
return None
|
return None
|
||||||
if not raw:
|
if raw is None:
|
||||||
return None
|
|
||||||
if isinstance(raw, str) and not raw.strip():
|
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
data = json.loads(str(raw))
|
value = int(str(raw).strip())
|
||||||
except Exception:
|
except (TypeError, ValueError):
|
||||||
logger.warning("Failed to decode sync listing checkpoint", exc_info=True)
|
logger.warning("Invalid sync listing checkpoint value %r; clearing", raw)
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
return None
|
return None
|
||||||
if not isinstance(data, dict):
|
if value < 0:
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
return None
|
return None
|
||||||
return data
|
return value
|
||||||
|
|
||||||
|
|
||||||
def _save_sync_checkpoint(
|
def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> None:
|
||||||
redis_client: Redis,
|
|
||||||
*,
|
|
||||||
task_id: str,
|
|
||||||
page_number: int,
|
|
||||||
make: str | None,
|
|
||||||
model: str | None,
|
|
||||||
lane: str,
|
|
||||||
segment_index: int | None = None,
|
|
||||||
) -> None:
|
|
||||||
payload = {
|
|
||||||
"status": "in_progress",
|
|
||||||
"task_id": task_id,
|
|
||||||
# page_number=0 — sentinel: прошлый checkpoint признан stale,
|
|
||||||
# следующий запуск должен начать текущий scope заново с page 1.
|
|
||||||
"last_successful_page": max(0, int(page_number)),
|
|
||||||
"resume_failures": 0,
|
|
||||||
"make": make,
|
|
||||||
"model": model,
|
|
||||||
"lane": lane,
|
|
||||||
"segment_index": segment_index,
|
|
||||||
"updated_at": int(time.time()),
|
|
||||||
}
|
|
||||||
try:
|
try:
|
||||||
redis_client.set(
|
redis_client.set(
|
||||||
SYNC_LISTING_CHECKPOINT_KEY,
|
SYNC_LISTING_CHECKPOINT_KEY,
|
||||||
json.dumps(payload),
|
str(int(segment_index)),
|
||||||
ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS,
|
ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
@@ -310,45 +351,6 @@ def _clear_sync_checkpoint(redis_client: Redis) -> None:
|
|||||||
logger.warning("Failed to clear sync listing checkpoint", exc_info=True)
|
logger.warning("Failed to clear sync listing checkpoint", exc_info=True)
|
||||||
|
|
||||||
|
|
||||||
def _bump_checkpoint_resume_failure(
|
|
||||||
redis_client: Redis,
|
|
||||||
checkpoint: dict[str, object] | None,
|
|
||||||
*,
|
|
||||||
reason: str,
|
|
||||||
) -> tuple[int, bool]:
|
|
||||||
if not checkpoint:
|
|
||||||
return 0, False
|
|
||||||
|
|
||||||
failures = int(checkpoint.get("resume_failures") or 0) + 1
|
|
||||||
if failures >= SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT:
|
|
||||||
logger.warning(
|
|
||||||
"Checkpoint resume failed %d times; deleting checkpoint and restarting from page 1 next run (reason=%s)",
|
|
||||||
failures,
|
|
||||||
reason,
|
|
||||||
)
|
|
||||||
_clear_sync_checkpoint(redis_client)
|
|
||||||
return failures, True
|
|
||||||
|
|
||||||
payload = dict(checkpoint)
|
|
||||||
payload["resume_failures"] = failures
|
|
||||||
payload["updated_at"] = int(time.time())
|
|
||||||
try:
|
|
||||||
redis_client.set(
|
|
||||||
SYNC_LISTING_CHECKPOINT_KEY,
|
|
||||||
json.dumps(payload),
|
|
||||||
ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS,
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
logger.warning("Failed to persist checkpoint resume failure counter", exc_info=True)
|
|
||||||
logger.warning(
|
|
||||||
"Checkpoint resume failure %d/%d recorded (reason=%s)",
|
|
||||||
failures,
|
|
||||||
SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT,
|
|
||||||
reason,
|
|
||||||
)
|
|
||||||
return failures, False
|
|
||||||
|
|
||||||
|
|
||||||
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
|
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
|
||||||
try:
|
try:
|
||||||
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
|
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
|
||||||
@@ -408,20 +410,29 @@ def _start_lock_heartbeat(
|
|||||||
key: str,
|
key: str,
|
||||||
owner_token: str,
|
owner_token: str,
|
||||||
ttl_seconds: int,
|
ttl_seconds: int,
|
||||||
) -> tuple[Event, Thread]:
|
) -> tuple[Event, Event, Thread]:
|
||||||
|
# Возвращает (stop_event, lock_lost_event, thread).
|
||||||
|
# lock_lost_event ставится в True, когда heartbeat обнаружил, что lock
|
||||||
|
# принадлежит другому owner'у (orphan-cleanup сработал ошибочно).
|
||||||
|
# Главная задача должна периодически проверять этот event.
|
||||||
stop_event = Event()
|
stop_event = Event()
|
||||||
|
lock_lost_event = Event()
|
||||||
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
||||||
|
|
||||||
def _heartbeat() -> None:
|
def _heartbeat() -> None:
|
||||||
while not stop_event.wait(interval_seconds):
|
while not stop_event.wait(interval_seconds):
|
||||||
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
||||||
if refreshed is False:
|
if refreshed is False:
|
||||||
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
logger.error(
|
||||||
|
"Lost sync_listing lock ownership for %s — signaling abort",
|
||||||
|
owner_token,
|
||||||
|
)
|
||||||
|
lock_lost_event.set()
|
||||||
return
|
return
|
||||||
|
|
||||||
thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True)
|
thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True)
|
||||||
thread.start()
|
thread.start()
|
||||||
return stop_event, thread
|
return stop_event, lock_lost_event, thread
|
||||||
|
|
||||||
|
|
||||||
@shared_task(
|
@shared_task(
|
||||||
@@ -480,29 +491,40 @@ def sync_listing_task(
|
|||||||
lock_acquired = False
|
lock_acquired = False
|
||||||
lock_ttl = _sync_listing_lock_ttl_seconds()
|
lock_ttl = _sync_listing_lock_ttl_seconds()
|
||||||
heartbeat_stop: Event | None = None
|
heartbeat_stop: Event | None = None
|
||||||
|
heartbeat_lock_lost: Event | None = None
|
||||||
heartbeat_thread: Thread | None = None
|
heartbeat_thread: Thread | None = None
|
||||||
force_bootstrap_full_scan = False
|
force_bootstrap_full_scan = False
|
||||||
|
|
||||||
def _enqueue_bootstrap_followup(
|
def _enqueue_bootstrap_followup(
|
||||||
reason: str,
|
reason: str,
|
||||||
delay_seconds: int = 5,
|
delay_seconds: int | None = None,
|
||||||
*,
|
*,
|
||||||
count_as_failure: bool = False,
|
count_as_failure: bool = False,
|
||||||
) -> None:
|
) -> None:
|
||||||
flag_ttl = max(lock_ttl, delay_seconds + 300)
|
# Всегда растим стрик (даже на soft timeout) — иначе в патологическом сценарии
|
||||||
|
# follow-up уходит каждые N секунд бесконечно. Сбрасывается только полным успехом.
|
||||||
|
streak, should_enqueue = _bump_bootstrap_failure_streak(redis_client, reason=reason)
|
||||||
|
if not should_enqueue:
|
||||||
|
logger.error(
|
||||||
|
"Bootstrap follow-up suppressed by circuit breaker (streak=%d, reason=%s); "
|
||||||
|
"next attempt will go through Celery beat schedule",
|
||||||
|
streak, reason,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
|
# Экспоненциальный backoff: 30s, 60s, 120s, ... до 1 часа.
|
||||||
|
# Для count_as_failure=False (soft timeout с прогрессом) — короче.
|
||||||
|
base = 15 if not count_as_failure else 30
|
||||||
|
computed_delay = min(3600, base * (2 ** max(0, streak - 1)))
|
||||||
|
effective_delay = computed_delay if delay_seconds is None else max(int(delay_seconds), computed_delay)
|
||||||
|
|
||||||
|
flag_ttl = max(lock_ttl, effective_delay + 300)
|
||||||
if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl):
|
if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl):
|
||||||
logger.info(
|
logger.info(
|
||||||
"Bootstrap follow-up already pending; skip enqueue (reason=%s)",
|
"Bootstrap follow-up already pending; skip enqueue (reason=%s)",
|
||||||
reason,
|
reason,
|
||||||
)
|
)
|
||||||
return
|
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)
|
|
||||||
try:
|
try:
|
||||||
self.app.send_task(
|
self.app.send_task(
|
||||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||||
@@ -514,12 +536,13 @@ def sync_listing_task(
|
|||||||
"only_new": only_new,
|
"only_new": only_new,
|
||||||
},
|
},
|
||||||
queue="scraping",
|
queue="scraping",
|
||||||
countdown=max(0, int(delay_seconds)),
|
countdown=effective_delay,
|
||||||
)
|
)
|
||||||
logger.info(
|
logger.info(
|
||||||
"Bootstrap follow-up sync queued in %ss (reason=%s)",
|
"Bootstrap follow-up sync queued in %ss (reason=%s, streak=%d)",
|
||||||
delay_seconds,
|
effective_delay,
|
||||||
reason,
|
reason,
|
||||||
|
streak,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
_clear_followup_pending(redis_client)
|
_clear_followup_pending(redis_client)
|
||||||
@@ -545,36 +568,15 @@ def sync_listing_task(
|
|||||||
force_bootstrap_full_scan = not full_scan_done_before_run
|
force_bootstrap_full_scan = not full_scan_done_before_run
|
||||||
effective_limit = None if force_bootstrap_full_scan else limit
|
effective_limit = None if force_bootstrap_full_scan else limit
|
||||||
effective_only_new = False if force_bootstrap_full_scan else only_new
|
effective_only_new = False if force_bootstrap_full_scan else only_new
|
||||||
checkpoint = _load_sync_checkpoint(redis_client)
|
|
||||||
resume_from_page = 1
|
|
||||||
used_checkpoint_resume = False
|
|
||||||
|
|
||||||
if force_bootstrap_full_scan and checkpoint and str(checkpoint.get("status") or "") == "in_progress":
|
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||||
checkpoint_page = int(checkpoint.get("last_successful_page") or 0)
|
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
||||||
checkpoint_make = checkpoint.get("make")
|
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
|
||||||
checkpoint_model = checkpoint.get("model")
|
last_completed_segment: int | None = None
|
||||||
checkpoint_lane = checkpoint.get("lane")
|
if force_bootstrap_full_scan:
|
||||||
checkpoint_segment = checkpoint.get("segment_index")
|
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||||
same_scope = (
|
else:
|
||||||
checkpoint_make == make
|
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||||
and checkpoint_model == model
|
|
||||||
and checkpoint_lane == lane
|
|
||||||
)
|
|
||||||
# Для сегментированного режима: совпадение по lane + наличие segment_index
|
|
||||||
is_segmented_checkpoint = checkpoint_segment is not None and checkpoint_lane == lane
|
|
||||||
if checkpoint_page > 0 and (same_scope or is_segmented_checkpoint):
|
|
||||||
resume_from_page = checkpoint_page + 1
|
|
||||||
used_checkpoint_resume = True
|
|
||||||
logger.warning(
|
|
||||||
"Resuming sync_listing from page %d (segment=%s) using checkpoint",
|
|
||||||
resume_from_page,
|
|
||||||
checkpoint_segment,
|
|
||||||
)
|
|
||||||
elif checkpoint_page > 0:
|
|
||||||
logger.info("Ignoring stale checkpoint due to different sync parameters")
|
|
||||||
_clear_sync_checkpoint(redis_client)
|
|
||||||
elif checkpoint:
|
|
||||||
logger.info("Ignoring leftover checkpoint because full scan is already complete; next run starts from page 1")
|
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
|
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
@@ -590,7 +592,7 @@ def sync_listing_task(
|
|||||||
|
|
||||||
_clear_followup_pending(redis_client)
|
_clear_followup_pending(redis_client)
|
||||||
|
|
||||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
heartbeat_stop, heartbeat_lock_lost, heartbeat_thread = _start_lock_heartbeat(
|
||||||
redis_client,
|
redis_client,
|
||||||
SYNC_LISTING_LOCK_KEY,
|
SYNC_LISTING_LOCK_KEY,
|
||||||
owner_token,
|
owner_token,
|
||||||
@@ -604,31 +606,37 @@ def sync_listing_task(
|
|||||||
use_segmented = bool(segments) and make is None and model is None
|
use_segmented = bool(segments) and make is None and model is None
|
||||||
|
|
||||||
resume_from_segment = 0
|
resume_from_segment = 0
|
||||||
if force_bootstrap_full_scan and use_segmented and checkpoint and str(checkpoint.get("status") or "") == "in_progress":
|
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
||||||
cp_segment = checkpoint.get("segment_index")
|
resume_from_segment = max(0, last_completed_segment + 1)
|
||||||
if cp_segment is not None and int(cp_segment) >= 0:
|
if resume_from_segment >= len(segments):
|
||||||
resume_from_segment = int(cp_segment)
|
# Все сегменты уже пройдены — чекпоинт устарел, начинаем заново.
|
||||||
# resume_from_page уже вычислен выше
|
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,
|
||||||
|
)
|
||||||
|
|
||||||
def _job():
|
def _job():
|
||||||
with IAAIScraper() as scraper:
|
with IAAIScraper() as scraper:
|
||||||
|
# Прокидываем abort-сигнал: scraper будет проверять его
|
||||||
|
# между сегментами/страницами и поднимет ScraperAbortedError.
|
||||||
|
if heartbeat_lock_lost is not None:
|
||||||
|
scraper.set_abort_callback(heartbeat_lock_lost.is_set)
|
||||||
if use_segmented:
|
if use_segmented:
|
||||||
return scraper.sync_listing_segmented(
|
return scraper.sync_listing_segmented(
|
||||||
segments=segments,
|
segments=segments,
|
||||||
lane=lane,
|
lane=lane,
|
||||||
only_new=effective_only_new,
|
only_new=effective_only_new,
|
||||||
start_segment=resume_from_segment,
|
start_segment=resume_from_segment,
|
||||||
start_page=resume_from_page,
|
start_page=1,
|
||||||
progress_callback=(
|
progress_callback=(
|
||||||
(lambda seg_idx, page_number: _save_sync_checkpoint(
|
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
|
||||||
redis_client,
|
|
||||||
task_id=task_id,
|
|
||||||
page_number=page_number,
|
|
||||||
make=None,
|
|
||||||
model=None,
|
|
||||||
lane=lane,
|
|
||||||
segment_index=seg_idx,
|
|
||||||
))
|
|
||||||
if force_bootstrap_full_scan else None
|
if force_bootstrap_full_scan else None
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
@@ -638,18 +646,6 @@ def sync_listing_task(
|
|||||||
lane=lane,
|
lane=lane,
|
||||||
limit=effective_limit,
|
limit=effective_limit,
|
||||||
only_new=effective_only_new,
|
only_new=effective_only_new,
|
||||||
start_page=resume_from_page,
|
|
||||||
progress_callback=(
|
|
||||||
(lambda page_number: _save_sync_checkpoint(
|
|
||||||
redis_client,
|
|
||||||
task_id=task_id,
|
|
||||||
page_number=page_number,
|
|
||||||
make=make,
|
|
||||||
model=model,
|
|
||||||
lane=lane,
|
|
||||||
))
|
|
||||||
if force_bootstrap_full_scan else None
|
|
||||||
),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
result = _run_browser_job(_job)
|
result = _run_browser_job(_job)
|
||||||
@@ -674,9 +670,7 @@ def sync_listing_task(
|
|||||||
"bootstrap_not_completed",
|
"bootstrap_not_completed",
|
||||||
count_as_failure=count_as_failure,
|
count_as_failure=count_as_failure,
|
||||||
)
|
)
|
||||||
elif result.get("status") in ("success", "partial_success") and bool(result.get("full_scan_completed", False)):
|
else:
|
||||||
_clear_sync_checkpoint(redis_client)
|
|
||||||
elif not force_bootstrap_full_scan:
|
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
|
|
||||||
summary = {
|
summary = {
|
||||||
@@ -715,21 +709,27 @@ def sync_listing_task(
|
|||||||
"note": "partial progress saved to DB; bootstrap continuation queued",
|
"note": "partial progress saved to DB; bootstrap continuation queued",
|
||||||
}
|
}
|
||||||
|
|
||||||
|
except ScraperAbortedError as exc:
|
||||||
|
# Lock потерян (orphan-cleanup в другом worker'е). НЕ освобождаем lock
|
||||||
|
# принудительно — он принадлежит другому owner'у. Не retry'имся, не
|
||||||
|
# планируем follow-up: новый владелец lock'а уже работает.
|
||||||
|
logger.warning("sync_listing_task aborted: %s", exc)
|
||||||
|
return {
|
||||||
|
"status": "aborted",
|
||||||
|
"task_id": task_id,
|
||||||
|
"reason": "lock_lost",
|
||||||
|
}
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
||||||
if force_bootstrap_full_scan and used_checkpoint_resume:
|
|
||||||
_, deleted = _bump_checkpoint_resume_failure(
|
|
||||||
redis_client,
|
|
||||||
checkpoint,
|
|
||||||
reason=str(exc),
|
|
||||||
)
|
|
||||||
if deleted:
|
|
||||||
checkpoint = None
|
|
||||||
# Retry только на не-таймаутные ошибки (сеть, БД, браузер).
|
# Retry только на не-таймаутные ошибки (сеть, БД, браузер).
|
||||||
try:
|
try:
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
_set_full_scan_done(redis_client, False)
|
_set_full_scan_done(redis_client, False)
|
||||||
raise self.retry(exc=exc, countdown=5)
|
# Не делаем мгновенный retry в bootstrap — иначе зацикливание при
|
||||||
|
# стабильно падающем сегменте. Backoff: 60s × 2^attempt.
|
||||||
|
attempt = int(getattr(self.request, "retries", 0) or 0)
|
||||||
|
raise self.retry(exc=exc, countdown=min(600, 60 * (2 ** attempt)))
|
||||||
raise self.retry(exc=exc)
|
raise self.retry(exc=exc)
|
||||||
except self.MaxRetriesExceededError:
|
except self.MaxRetriesExceededError:
|
||||||
logger.error("sync_listing_task max retries exceeded, giving up")
|
logger.error("sync_listing_task max retries exceeded, giving up")
|
||||||
|
|||||||
@@ -5,7 +5,7 @@ from types import SimpleNamespace
|
|||||||
from unittest.mock import MagicMock
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
from iaai_scraper.core.config import Settings
|
from iaai_scraper.core.config import Settings
|
||||||
from iaai_scraper.core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError
|
from iaai_scraper.core.exceptions import AntiBotDetectedError, SiteStructureChangedError
|
||||||
from iaai_scraper.scraper import IAAIScraper
|
from iaai_scraper.scraper import IAAIScraper
|
||||||
from iaai_scraper.storage.schemas import CarRecord
|
from iaai_scraper.storage.schemas import CarRecord
|
||||||
|
|
||||||
@@ -107,16 +107,13 @@ class TestScraperSync(unittest.TestCase):
|
|||||||
{"make": "HONDA", "year_min": None, "year_max": None},
|
{"make": "HONDA", "year_min": None, "year_max": None},
|
||||||
]
|
]
|
||||||
|
|
||||||
# Resume: пропускаем TOYOTA, начинаем с FORD на стр. 5.
|
# Resume: пропускаем TOYOTA, начинаем сразу с FORD.
|
||||||
result = scraper.sync_listing_segmented(
|
result = scraper.sync_listing_segmented(
|
||||||
segments=segments, start_segment=1, start_page=5,
|
segments=segments, start_segment=1,
|
||||||
)
|
)
|
||||||
|
|
||||||
self.assertEqual(result["segments_completed"], 2) # FORD + HONDA
|
self.assertEqual(result["segments_completed"], 2) # FORD + HONDA
|
||||||
self.assertEqual(result["cars_upserted"], 10)
|
self.assertEqual(result["cars_upserted"], 10)
|
||||||
# FORD: start_page=5, HONDA: start_page=1.
|
|
||||||
self.assertEqual(call_args_log[0]["start_page"], 5)
|
|
||||||
self.assertEqual(call_args_log[1]["start_page"], 1)
|
|
||||||
# URL содержит бренд.
|
# URL содержит бренд.
|
||||||
self.assertIn("FORD", call_args_log[0]["listing_url"])
|
self.assertIn("FORD", call_args_log[0]["listing_url"])
|
||||||
self.assertIn("HONDA", call_args_log[1]["listing_url"])
|
self.assertIn("HONDA", call_args_log[1]["listing_url"])
|
||||||
@@ -223,7 +220,7 @@ class TestScraperSync(unittest.TestCase):
|
|||||||
scraper.listing_collector.open_cars_listing.assert_called_once()
|
scraper.listing_collector.open_cars_listing.assert_called_once()
|
||||||
page.reload.assert_not_called()
|
page.reload.assert_not_called()
|
||||||
|
|
||||||
def test_sync_listing_resets_stale_resume_checkpoint_and_restarts_from_page_one(self) -> None:
|
def test_sync_listing_always_starts_from_page_one(self) -> None:
|
||||||
scraper = self._make_scraper()
|
scraper = self._make_scraper()
|
||||||
scraper.settings.celery.batch_size = 1
|
scraper.settings.celery.batch_size = 1
|
||||||
|
|
||||||
@@ -234,10 +231,11 @@ class TestScraperSync(unittest.TestCase):
|
|||||||
next_page_detected=False,
|
next_page_detected=False,
|
||||||
)
|
)
|
||||||
|
|
||||||
progress_calls: list[int] = []
|
|
||||||
|
|
||||||
scraper._get_page_with_warmup = MagicMock(return_value=page)
|
scraper._get_page_with_warmup = MagicMock(return_value=page)
|
||||||
scraper._reopen_listing_and_resume = MagicMock(side_effect=ListingResumeError("checkpoint page is no longer reachable"))
|
# Pagination-resume убран: _reopen_listing_and_resume не должен вызываться.
|
||||||
|
scraper._reopen_listing_and_resume = MagicMock(
|
||||||
|
side_effect=AssertionError("pagination resume must not be used")
|
||||||
|
)
|
||||||
scraper.listing_collector.open_cars_listing = MagicMock()
|
scraper.listing_collector.open_cars_listing = MagicMock()
|
||||||
scraper.listing_collector.apply_filters = MagicMock(return_value={
|
scraper.listing_collector.apply_filters = MagicMock(return_value={
|
||||||
"make": None,
|
"make": None,
|
||||||
@@ -261,13 +259,11 @@ class TestScraperSync(unittest.TestCase):
|
|||||||
limit=None,
|
limit=None,
|
||||||
effective_only_new=False,
|
effective_only_new=False,
|
||||||
started_at=0.0,
|
started_at=0.0,
|
||||||
start_page=50,
|
|
||||||
progress_callback=progress_calls.append,
|
|
||||||
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=EAGLE",
|
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=EAGLE",
|
||||||
)
|
)
|
||||||
|
|
||||||
self.assertEqual(progress_calls, [0, 1])
|
|
||||||
scraper.listing_collector.open_cars_listing.assert_called_once()
|
scraper.listing_collector.open_cars_listing.assert_called_once()
|
||||||
|
scraper._reopen_listing_and_resume.assert_not_called()
|
||||||
scraper.sync_batch.assert_called_once()
|
scraper.sync_batch.assert_called_once()
|
||||||
self.assertEqual(result["cars_upserted"], 1)
|
self.assertEqual(result["cars_upserted"], 1)
|
||||||
self.assertEqual(result["total"], 1)
|
self.assertEqual(result["total"], 1)
|
||||||
|
|||||||
@@ -1,6 +1,5 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
|
||||||
import unittest
|
import unittest
|
||||||
from unittest.mock import MagicMock, patch
|
from unittest.mock import MagicMock, patch
|
||||||
|
|
||||||
@@ -68,7 +67,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
stop_event = MagicMock()
|
||||||
heartbeat_thread = MagicMock()
|
heartbeat_thread = MagicMock()
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
start_heartbeat.return_value = (stop_event, MagicMock(), heartbeat_thread)
|
||||||
|
|
||||||
tasks.sync_listing_task.push_request(id="task-123")
|
tasks.sync_listing_task.push_request(id="task-123")
|
||||||
try:
|
try:
|
||||||
@@ -133,7 +132,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
stop_event = MagicMock()
|
||||||
heartbeat_thread = MagicMock()
|
heartbeat_thread = MagicMock()
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
start_heartbeat.return_value = (stop_event, MagicMock(), heartbeat_thread)
|
||||||
|
|
||||||
tasks.sync_listing_task.push_request(id="task-456")
|
tasks.sync_listing_task.push_request(id="task-456")
|
||||||
try:
|
try:
|
||||||
@@ -146,69 +145,84 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
self.assertEqual(acquire_lock.call_count, 2)
|
self.assertEqual(acquire_lock.call_count, 2)
|
||||||
release_lock.assert_called_once()
|
release_lock.assert_called_once()
|
||||||
|
|
||||||
def test_sync_listing_checkpoint_save_load_and_resume(self) -> None:
|
def test_checkpoint_save_load_roundtrip(self) -> None:
|
||||||
# Roundtrip: save → load.
|
|
||||||
redis_client = MagicMock()
|
|
||||||
storage: dict[str, str] = {}
|
storage: dict[str, str] = {}
|
||||||
|
|
||||||
def _fake_set(key, value, ex=None):
|
def _fake_set(key, value, ex=None):
|
||||||
storage[key] = value
|
storage[key] = value
|
||||||
return True
|
return True
|
||||||
|
|
||||||
|
redis_client = MagicMock()
|
||||||
redis_client.set.side_effect = _fake_set
|
redis_client.set.side_effect = _fake_set
|
||||||
redis_client.get.side_effect = lambda key: storage.get(key)
|
redis_client.get.side_effect = lambda key: storage.get(key)
|
||||||
|
|
||||||
tasks._save_sync_checkpoint(
|
tasks._save_last_completed_segment(redis_client, 12)
|
||||||
redis_client, task_id="task-1", page_number=12,
|
|
||||||
make=None, model=None, lane="iaai_cars",
|
|
||||||
)
|
|
||||||
checkpoint = tasks._load_sync_checkpoint(redis_client)
|
|
||||||
self.assertIsNotNone(checkpoint)
|
|
||||||
self.assertEqual(checkpoint["last_successful_page"], 12)
|
|
||||||
|
|
||||||
# Resume: задача стартует со страницы checkpoint + 1.
|
self.assertEqual(tasks._load_last_completed_segment(redis_client), 12)
|
||||||
|
self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS)
|
||||||
|
|
||||||
|
def test_load_checkpoint_clears_invalid_payload(self) -> None:
|
||||||
|
redis_client = MagicMock()
|
||||||
|
redis_client.get.return_value = "{not-a-number"
|
||||||
|
|
||||||
|
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
|
||||||
|
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY)
|
||||||
|
|
||||||
|
def test_load_checkpoint_empty_when_missing(self) -> None:
|
||||||
|
redis_client = MagicMock()
|
||||||
|
redis_client.get.return_value = None
|
||||||
|
|
||||||
|
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
|
||||||
|
redis_client.delete.assert_not_called()
|
||||||
|
|
||||||
|
def test_sync_listing_resumes_from_next_segment_during_bootstrap(self) -> None:
|
||||||
|
segments = [
|
||||||
|
{"make": "ACURA"},
|
||||||
|
{"make": "AUDI"},
|
||||||
|
{"make": "BMW"},
|
||||||
|
{"make": "EAGLE"},
|
||||||
|
]
|
||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
patch.object(tasks, "_get_redis") as get_redis, \
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
||||||
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
||||||
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
||||||
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
||||||
patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \
|
|
||||||
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
||||||
patch.object(tasks.sync_listing_task, "update_state"), \
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
||||||
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
||||||
get_persistence.return_value = MagicMock()
|
get_persistence.return_value = MagicMock()
|
||||||
redis_client2 = MagicMock()
|
redis_client = MagicMock()
|
||||||
redis_client2.get.side_effect = lambda key: json.dumps({
|
redis_client.get.side_effect = lambda key: (
|
||||||
"status": "in_progress", "last_successful_page": 9,
|
"1" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
||||||
"make": None, "model": None, "lane": "iaai_cars",
|
)
|
||||||
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
get_redis.return_value = redis_client
|
||||||
get_redis.return_value = redis_client2
|
start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock())
|
||||||
stop_event = MagicMock()
|
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
sync_listing_mock = MagicMock(return_value={
|
sync_segmented_mock = MagicMock(return_value={
|
||||||
"run_id": 11, "status": "success", "full_scan_completed": True,
|
"run_id": 11, "status": "success", "full_scan_completed": True,
|
||||||
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
||||||
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
||||||
})
|
})
|
||||||
scraper_ctx = MagicMock()
|
scraper_ctx = MagicMock()
|
||||||
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
|
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
|
||||||
scraper_ctx.__exit__.return_value = None
|
scraper_ctx.__exit__.return_value = None
|
||||||
|
|
||||||
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
||||||
tasks.sync_listing_task.push_request(id="task-789")
|
tasks.sync_listing_task.push_request(id="task-resume-seg")
|
||||||
try:
|
try:
|
||||||
result = tasks.sync_listing_task.run()
|
result = tasks.sync_listing_task.run()
|
||||||
finally:
|
finally:
|
||||||
tasks.sync_listing_task.pop_request()
|
tasks.sync_listing_task.pop_request()
|
||||||
|
|
||||||
self.assertEqual(result["status"], "success")
|
self.assertEqual(result["status"], "success")
|
||||||
self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 10)
|
# last_completed=1 → start_segment=2 (AUDI завершён, возобновляем с BMW).
|
||||||
clear_checkpoint.assert_called()
|
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 2)
|
||||||
|
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_page"], 1)
|
||||||
release_lock.assert_called_once()
|
release_lock.assert_called_once()
|
||||||
|
|
||||||
def test_sync_listing_ignores_checkpoint_after_full_scan_completed(self) -> None:
|
def test_sync_listing_ignores_checkpoint_after_full_scan_completed(self) -> None:
|
||||||
|
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
|
||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
patch.object(tasks, "_get_redis") as get_redis, \
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
||||||
@@ -218,25 +232,23 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \
|
patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \
|
||||||
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
||||||
patch.object(tasks.sync_listing_task, "update_state"), \
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
||||||
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
||||||
get_persistence.return_value = MagicMock()
|
get_persistence.return_value = MagicMock()
|
||||||
redis_client = MagicMock()
|
redis_client = MagicMock()
|
||||||
redis_client.get.side_effect = lambda key: json.dumps({
|
# Оставшийся чекпоинт не должен использоваться.
|
||||||
"status": "in_progress", "last_successful_page": 25,
|
redis_client.get.side_effect = lambda key: (
|
||||||
"make": None, "model": None, "lane": "iaai_cars",
|
"0" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
||||||
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
)
|
||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock())
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
sync_listing_mock = MagicMock(return_value={
|
sync_segmented_mock = MagicMock(return_value={
|
||||||
"run_id": 14, "status": "success", "full_scan_completed": True,
|
"run_id": 14, "status": "success", "full_scan_completed": True,
|
||||||
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
||||||
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
||||||
})
|
})
|
||||||
scraper_ctx = MagicMock()
|
scraper_ctx = MagicMock()
|
||||||
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
|
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
|
||||||
scraper_ctx.__exit__.return_value = None
|
scraper_ctx.__exit__.return_value = None
|
||||||
|
|
||||||
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
||||||
@@ -247,77 +259,51 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
tasks.sync_listing_task.pop_request()
|
tasks.sync_listing_task.pop_request()
|
||||||
|
|
||||||
self.assertEqual(result["status"], "success")
|
self.assertEqual(result["status"], "success")
|
||||||
self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 1)
|
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 0)
|
||||||
self.assertIsNone(sync_listing_mock.call_args.kwargs["progress_callback"])
|
self.assertIsNone(sync_segmented_mock.call_args.kwargs["progress_callback"])
|
||||||
clear_checkpoint.assert_called()
|
clear_checkpoint.assert_called()
|
||||||
release_lock.assert_called_once()
|
release_lock.assert_called_once()
|
||||||
|
|
||||||
def test_sync_listing_checkpoint_zero_page_restarts_from_first_page(self) -> None:
|
def test_sync_listing_checkpoint_beyond_segments_restarts_from_zero(self) -> None:
|
||||||
|
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
|
||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
patch.object(tasks, "_get_redis") as get_redis, \
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
||||||
patch.object(tasks, "_is_full_scan_done", return_value=True), \
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
||||||
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
||||||
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
||||||
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
||||||
patch.object(tasks.sync_listing_task, "update_state"), \
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
||||||
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
||||||
get_persistence.return_value = MagicMock()
|
get_persistence.return_value = MagicMock()
|
||||||
redis_client = MagicMock()
|
redis_client = MagicMock()
|
||||||
redis_client.get.side_effect = lambda key: json.dumps({
|
redis_client.get.side_effect = lambda key: (
|
||||||
"status": "in_progress", "last_successful_page": 0,
|
"99" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
||||||
"make": None, "model": None, "lane": "iaai_cars",
|
)
|
||||||
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
|
||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock())
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
sync_listing_mock = MagicMock(return_value={
|
sync_segmented_mock = MagicMock(return_value={
|
||||||
"run_id": 12, "status": "success", "full_scan_completed": True,
|
"run_id": 15, "status": "success", "full_scan_completed": True,
|
||||||
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
"cars_upserted": 0, "cars_failed": 0, "images_upserted": 0,
|
||||||
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
||||||
})
|
})
|
||||||
scraper_ctx = MagicMock()
|
scraper_ctx = MagicMock()
|
||||||
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
|
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
|
||||||
scraper_ctx.__exit__.return_value = None
|
scraper_ctx.__exit__.return_value = None
|
||||||
|
|
||||||
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
||||||
tasks.sync_listing_task.push_request(id="task-790")
|
tasks.sync_listing_task.push_request(id="task-beyond")
|
||||||
try:
|
try:
|
||||||
result = tasks.sync_listing_task.run()
|
result = tasks.sync_listing_task.run()
|
||||||
finally:
|
finally:
|
||||||
tasks.sync_listing_task.pop_request()
|
tasks.sync_listing_task.pop_request()
|
||||||
|
|
||||||
self.assertEqual(result["status"], "success")
|
self.assertEqual(result["status"], "success")
|
||||||
self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 1)
|
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 0)
|
||||||
release_lock.assert_called_once()
|
release_lock.assert_called_once()
|
||||||
|
|
||||||
def test_save_sync_checkpoint_sets_ttl(self) -> None:
|
def test_sync_listing_task_clears_checkpoint_on_hourly_run(self) -> None:
|
||||||
redis_client = MagicMock()
|
|
||||||
|
|
||||||
tasks._save_sync_checkpoint(
|
|
||||||
redis_client,
|
|
||||||
task_id="task-1",
|
|
||||||
page_number=3,
|
|
||||||
make=None,
|
|
||||||
model=None,
|
|
||||||
lane="iaai_cars",
|
|
||||||
)
|
|
||||||
|
|
||||||
redis_client.set.assert_called_once()
|
|
||||||
self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS)
|
|
||||||
|
|
||||||
def test_load_sync_checkpoint_clears_invalid_payload(self) -> None:
|
|
||||||
redis_client = MagicMock()
|
|
||||||
redis_client.get.return_value = "{broken-json"
|
|
||||||
|
|
||||||
checkpoint = tasks._load_sync_checkpoint(redis_client)
|
|
||||||
|
|
||||||
self.assertIsNone(checkpoint)
|
|
||||||
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY)
|
|
||||||
|
|
||||||
def test_sync_listing_task_clears_checkpoint_on_partial_success_when_full_scan_completed(self) -> None:
|
|
||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
patch.object(tasks, "_get_redis") as get_redis, \
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
||||||
@@ -341,9 +327,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
redis_client = MagicMock()
|
redis_client = MagicMock()
|
||||||
redis_client.get.return_value = None
|
redis_client.get.return_value = None
|
||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock())
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
tasks.sync_listing_task.push_request(id="task-791")
|
tasks.sync_listing_task.push_request(id="task-791")
|
||||||
try:
|
try:
|
||||||
@@ -352,7 +336,9 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
tasks.sync_listing_task.pop_request()
|
tasks.sync_listing_task.pop_request()
|
||||||
|
|
||||||
self.assertEqual(result["status"], "partial_success")
|
self.assertEqual(result["status"], "partial_success")
|
||||||
clear_checkpoint.assert_called_once_with(redis_client)
|
# Hourly-ветка (full_scan_done=True) всегда удаляет оставшийся чекпоинт
|
||||||
|
# до браузерного job + после. Главное — вызов произошёл.
|
||||||
|
clear_checkpoint.assert_called()
|
||||||
release_lock.assert_called_once()
|
release_lock.assert_called_once()
|
||||||
|
|
||||||
def test_try_set_followup_pending_deduplicates(self) -> None:
|
def test_try_set_followup_pending_deduplicates(self) -> None:
|
||||||
@@ -417,9 +403,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
redis_client = MagicMock()
|
redis_client = MagicMock()
|
||||||
redis_client.get.return_value = None
|
redis_client.get.return_value = None
|
||||||
get_redis.return_value = redis_client
|
get_redis.return_value = redis_client
|
||||||
stop_event = MagicMock()
|
start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock())
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app:
|
with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app:
|
||||||
tasks.sync_listing_task.push_request(id="task-900")
|
tasks.sync_listing_task.push_request(id="task-900")
|
||||||
@@ -433,86 +417,6 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
clear_pending.assert_called()
|
clear_pending.assert_called()
|
||||||
task_app.send_task.assert_not_called()
|
task_app.send_task.assert_not_called()
|
||||||
|
|
||||||
def test_bump_checkpoint_resume_failure_deletes_after_second_failure(self) -> None:
|
|
||||||
redis_client = MagicMock()
|
|
||||||
checkpoint = {
|
|
||||||
"status": "in_progress",
|
|
||||||
"last_successful_page": 9,
|
|
||||||
"resume_failures": 1,
|
|
||||||
"lane": "iaai_cars",
|
|
||||||
}
|
|
||||||
|
|
||||||
failures, deleted = tasks._bump_checkpoint_resume_failure(
|
|
||||||
redis_client,
|
|
||||||
checkpoint,
|
|
||||||
reason="resume failed",
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertEqual(failures, 2)
|
|
||||||
self.assertTrue(deleted)
|
|
||||||
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY)
|
|
||||||
|
|
||||||
def test_bump_checkpoint_resume_failure_persists_first_failure(self) -> None:
|
|
||||||
redis_client = MagicMock()
|
|
||||||
checkpoint = {
|
|
||||||
"status": "in_progress",
|
|
||||||
"last_successful_page": 9,
|
|
||||||
"resume_failures": 0,
|
|
||||||
"lane": "iaai_cars",
|
|
||||||
}
|
|
||||||
|
|
||||||
failures, deleted = tasks._bump_checkpoint_resume_failure(
|
|
||||||
redis_client,
|
|
||||||
checkpoint,
|
|
||||||
reason="resume failed",
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertEqual(failures, 1)
|
|
||||||
self.assertFalse(deleted)
|
|
||||||
redis_client.set.assert_called_once()
|
|
||||||
|
|
||||||
def test_sync_listing_task_deletes_checkpoint_after_second_resume_failure(self) -> None:
|
|
||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
|
||||||
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
|
||||||
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
|
||||||
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
|
||||||
patch.object(tasks, "_bump_checkpoint_resume_failure", return_value=(2, True)) as bump_failures, \
|
|
||||||
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
||||||
patch.object(tasks.sync_listing_task, "retry", side_effect=tasks.sync_listing_task.MaxRetriesExceededError()), \
|
|
||||||
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
|
|
||||||
get_persistence.return_value = MagicMock()
|
|
||||||
redis_client = MagicMock()
|
|
||||||
redis_client.get.side_effect = lambda key: json.dumps({
|
|
||||||
"status": "in_progress",
|
|
||||||
"last_successful_page": 9,
|
|
||||||
"resume_failures": 1,
|
|
||||||
"make": None,
|
|
||||||
"model": None,
|
|
||||||
"lane": "iaai_cars",
|
|
||||||
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
|
||||||
get_redis.return_value = redis_client
|
|
||||||
stop_event = MagicMock()
|
|
||||||
heartbeat_thread = MagicMock()
|
|
||||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
|
||||||
|
|
||||||
with patch.object(tasks, "IAAIScraper") as scraper_cls:
|
|
||||||
scraper_ctx = MagicMock()
|
|
||||||
scraper_ctx.__enter__.return_value.sync_listing = MagicMock(side_effect=RuntimeError("Failed to resume listing at page 10"))
|
|
||||||
scraper_ctx.__exit__.return_value = None
|
|
||||||
scraper_cls.return_value = scraper_ctx
|
|
||||||
|
|
||||||
tasks.sync_listing_task.push_request(id="task-793")
|
|
||||||
try:
|
|
||||||
result = tasks.sync_listing_task.run()
|
|
||||||
finally:
|
|
||||||
tasks.sync_listing_task.pop_request()
|
|
||||||
|
|
||||||
self.assertEqual(result["status"], "failed")
|
|
||||||
bump_failures.assert_called_once()
|
|
||||||
release_lock.assert_called_once()
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
Reference in New Issue
Block a user