5 Commits

Author SHA1 Message Date
6b69b8c002 Merge branch 'main' into prod-preparation 2026-04-23 22:24:53 +02:00
qananasikq
e726461e75 fix self heal restart 2026-04-23 21:34:57 +03:00
f221aebef0 Prepared for production 2026-04-23 13:54:53 +03:00
qananasikq
6aba4f0dbd fix resume 2026-04-23 09:56:18 +03:00
qananasikq
09fecdae31 fix watchdog 2026-04-23 09:56:18 +03:00
11 changed files with 322 additions and 506 deletions

View File

@@ -14,6 +14,7 @@ from iaai_scraper.core.config import Settings
from iaai_scraper.storage.models import Base from iaai_scraper.storage.models import Base
config = context.config config = context.config
VERSION_TABLE = "iaai_parser_alembic_version"
# Logging # Logging
if config.config_file_name is not None: if config.config_file_name is not None:
@@ -32,6 +33,7 @@ def run_migrations_offline() -> None:
context.configure( context.configure(
url=url, url=url,
target_metadata=target_metadata, target_metadata=target_metadata,
version_table=VERSION_TABLE,
literal_binds=True, literal_binds=True,
dialect_opts={"paramstyle": "named"}, dialect_opts={"paramstyle": "named"},
) )
@@ -47,7 +49,11 @@ def run_migrations_online() -> None:
poolclass=pool.NullPool, poolclass=pool.NullPool,
) )
with connectable.connect() as connection: with connectable.connect() as connection:
context.configure(connection=connection, target_metadata=target_metadata) context.configure(
connection=connection,
target_metadata=target_metadata,
version_table=VERSION_TABLE,
)
with context.begin_transaction(): with context.begin_transaction():
context.run_migrations() context.run_migrations()

View File

@@ -1,4 +1,4 @@
# Начальная схема: cars, images, sync_runs # Начальная схема: iaai_cars, iaai_images, iaai_sync_runs
from typing import Sequence, Union from typing import Sequence, Union
from alembic import op from alembic import op
@@ -10,70 +10,94 @@ branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None depends_on: Union[str, Sequence[str], None] = None
def _schema_name() -> str | None:
bind = op.get_bind()
return "public" if bind.dialect.name == "postgresql" else None
def _table_exists(table_name: str) -> bool:
inspector = sa.inspect(op.get_bind())
return table_name in inspector.get_table_names(schema=_schema_name())
def _index_exists(table_name: str, index_name: str) -> bool:
if not _table_exists(table_name):
return False
inspector = sa.inspect(op.get_bind())
indexes = inspector.get_indexes(table_name, schema=_schema_name())
return any(index.get("name") == index_name for index in indexes)
def upgrade() -> None: def upgrade() -> None:
# Таблица cars — аукционные машины с IAAI # Таблица iaai_cars — аукционные машины с IAAI
op.create_table( if not _table_exists("iaai_cars"):
"cars", op.create_table(
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), "iaai_cars",
sa.Column("parser_id", sa.String(50), nullable=False, unique=True), sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("brand", sa.String(50), nullable=False), sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
sa.Column("model", sa.String(50), nullable=False), sa.Column("brand", sa.String(50), nullable=False),
sa.Column("year", sa.Integer, nullable=True), sa.Column("model", sa.String(50), nullable=False),
sa.Column("price", sa.BigInteger, nullable=True), sa.Column("year", sa.Integer, nullable=True),
sa.Column("currency", sa.String(10), nullable=False, server_default="USD"), sa.Column("price", sa.BigInteger, nullable=True),
sa.Column("mileage", sa.Integer, nullable=False, server_default="0"), sa.Column("currency", sa.String(10), nullable=False, server_default="USD"),
sa.Column("country", sa.String(10), nullable=False, server_default="NA"), sa.Column("mileage", sa.Integer, nullable=False, server_default="0"),
sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("country", sa.String(10), nullable=False, server_default="NA"),
sa.Column("color", sa.String, nullable=False, server_default="other"), sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("drive", sa.String(10), nullable=True), sa.Column("color", sa.String, nullable=False, server_default="other"),
sa.Column("gearbox", sa.String(10), nullable=True), sa.Column("drive", sa.String(10), nullable=True),
sa.Column("steering_wheel", sa.String(10), nullable=True), sa.Column("gearbox", sa.String(10), nullable=True),
sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"), sa.Column("steering_wheel", sa.String(10), nullable=True),
sa.Column("engine_volume", sa.Integer, nullable=True), sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"),
sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"), sa.Column("engine_volume", sa.Integer, nullable=True),
sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"),
sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("origin", sa.String(20), nullable=False, server_default="NA"), sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("origin_url", sa.String, nullable=False), sa.Column("origin", sa.String(20), nullable=False, server_default="NA"),
sa.Column("origin_id", sa.String, nullable=False, unique=True), sa.Column("origin_url", sa.String, nullable=False),
sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("origin_id", sa.String, nullable=False, unique=True),
sa.Column("evaluation", sa.String, nullable=True), sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")), sa.Column("evaluation", sa.String, nullable=True),
sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")),
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("slug", sa.String, nullable=False), sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), sa.Column("slug", sa.String, nullable=False),
) sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
op.create_index("ix_cars_origin_url", "cars", ["origin_url"]) )
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
# Таблица images — изображения машин if not _index_exists("iaai_cars", "ix_iaai_cars_origin_url"):
op.create_table( op.create_index("ix_iaai_cars_origin_url", "iaai_cars", ["origin_url"])
"images", if not _index_exists("iaai_cars", "ix_iaai_cars_origin_id"):
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), op.create_index("ix_iaai_cars_origin_id", "iaai_cars", ["origin_id"], unique=True)
sa.Column("fullres_image", sa.String, nullable=False),
sa.Column("preview_image", sa.String, nullable=False), # Таблица iaai_images — изображения машин
sa.Column("order_index", sa.Integer, nullable=False), if not _table_exists("iaai_images"):
sa.Column("car_id", sa.Integer, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False), op.create_table(
) "iaai_images",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("fullres_image", sa.String, nullable=False),
sa.Column("preview_image", sa.String, nullable=False),
sa.Column("order_index", sa.Integer, nullable=False),
sa.Column("car_id", sa.Integer, sa.ForeignKey("iaai_cars.id", ondelete="CASCADE"), nullable=False),
)
# Таблица iaai_sync_runs — логирование синхронизаций
if not _table_exists("iaai_sync_runs"):
op.create_table(
"iaai_sync_runs",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("status", sa.Text, nullable=False),
sa.Column("lane", sa.Text, nullable=False),
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("error_summary", sa.Text, nullable=True),
)
# Таблица sync_runs — логирование синхронизаций
op.create_table(
"sync_runs",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("status", sa.Text, nullable=False),
sa.Column("lane", sa.Text, nullable=False),
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("error_summary", sa.Text, nullable=True),
)
def downgrade() -> None: def downgrade() -> None:
op.drop_table("sync_runs") # Compatibility-only migration for shared production databases.
op.drop_table("images") pass
op.drop_table("cars")

View File

@@ -10,20 +10,15 @@ depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None: def upgrade() -> None:
op.create_index("ix_cars_brand", "cars", ["brand"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand ON iaai_cars (brand)")
op.create_index("ix_cars_brand_model", "cars", ["brand", "model"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand_model ON iaai_cars (brand, model)")
op.create_index("ix_cars_year", "cars", ["year"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_year ON iaai_cars (year)")
op.create_index("ix_cars_is_sold", "cars", ["is_sold"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_is_sold ON iaai_cars (is_sold)")
op.create_index("ix_cars_last_seen_at", "cars", ["last_seen_at"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_last_seen_at ON iaai_cars (last_seen_at)")
op.create_index("ix_images_car_id", "images", ["car_id"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id ON iaai_images (car_id)")
op.create_index("ix_sync_runs_status", "sync_runs", ["status"]) op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_status ON iaai_sync_runs (status)")
def downgrade() -> None: def downgrade() -> None:
op.drop_index("ix_sync_runs_status", table_name="sync_runs") # Compatibility-only migration for shared production databases.
op.drop_index("ix_images_car_id", table_name="images") pass
op.drop_index("ix_cars_last_seen_at", table_name="cars")
op.drop_index("ix_cars_is_sold", table_name="cars")
op.drop_index("ix_cars_year", table_name="cars")
op.drop_index("ix_cars_brand_model", table_name="cars")
op.drop_index("ix_cars_brand", table_name="cars")

View File

@@ -13,26 +13,24 @@ def upgrade() -> None:
# Partial index для активных машин IAAI (is_sold = FALSE) # Partial index для активных машин IAAI (is_sold = FALSE)
# Ускоряет поиск при mark_sold операциях # Ускоряет поиск при mark_sold операциях
op.execute( op.execute(
"CREATE INDEX IF NOT EXISTS ix_cars_active_iaai " "CREATE INDEX IF NOT EXISTS ix_iaai_cars_active_iaai "
"ON cars (origin_url) " "ON iaai_cars (origin_url) "
"WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'" "WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'"
) )
# Composite index для проверки дубликатов изображений # Composite index для проверки дубликатов изображений
op.create_index( op.execute(
"ix_images_car_id_fullres", "CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id_fullres "
"images", "ON iaai_images (car_id, fullres_image)"
["car_id", "fullres_image"],
) )
# Index для быстрого поиска existing машин по origin_url # Index для быстрого поиска existing машин по origin_url
op.execute( op.execute(
"CREATE INDEX IF NOT EXISTS ix_cars_origin_url_id " "CREATE INDEX IF NOT EXISTS ix_iaai_cars_origin_url_id "
"ON cars (origin_url, origin_id)" "ON iaai_cars (origin_url, origin_id)"
) )
def downgrade() -> None: def downgrade() -> None:
op.execute("DROP INDEX IF EXISTS ix_cars_origin_url_id") # Compatibility-only migration for shared production databases.
op.drop_index("ix_images_car_id_fullres", table_name="images") pass
op.execute("DROP INDEX IF EXISTS ix_cars_active_iaai")

View File

@@ -60,6 +60,11 @@ start_self_heal_watchdog_if_worker() {
return 0 return 0
fi fi
# После reboot/restart контейнера файловая система контейнера сохраняется,
# поэтому stale pidfile может остаться от прошлого запуска и блокировать старт Celery.
# Удаляем его перед запуском worker-процесса.
rm -f /tmp/celery-worker.pid
if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then
echo "[entrypoint] Self-heal watchdog disabled" echo "[entrypoint] Self-heal watchdog disabled"
return 0 return 0

View File

@@ -11,6 +11,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
from concurrent.futures import TimeoutError as FuturesTimeoutError from concurrent.futures import TimeoutError as FuturesTimeoutError
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from threading import Event, Thread
from typing import Any, Callable from typing import Any, Callable
from urllib.parse import urlsplit, urlunsplit from urllib.parse import urlsplit, urlunsplit
from urllib.request import Request, urlopen from urllib.request import Request, urlopen
@@ -44,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90
_PAGE_NEXT_TIMEOUT_S = 60 _PAGE_NEXT_TIMEOUT_S = 60
# Ротация контекста выключена по умолчанию. # Ротация контекста выключена по умолчанию.
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0) _CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT = int(os.environ.get("IAAI_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT", "3") or 3)
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120 _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2 _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
@@ -97,6 +99,33 @@ class _PageOpWatchdog:
return False return False
class _ProgressHeartbeat:
"""Периодически пульсует progress во время долгой навигации."""
def __init__(self, report_progress: Callable[[str, Any], None], stage: str, interval_s: float = 15.0, **meta: Any) -> None:
self._report_progress = report_progress
self._stage = stage
self._interval_s = max(5.0, float(interval_s))
self._meta = meta
self._stop = Event()
self._thread: Thread | None = None
def _run(self) -> None:
while not self._stop.wait(self._interval_s):
self._report_progress(self._stage, **self._meta)
def __enter__(self):
self._thread = Thread(target=self._run, name=f"progress-heartbeat-{self._stage}", daemon=True)
self._thread.start()
return self
def __exit__(self, exc_type, exc, tb):
self._stop.set()
if self._thread is not None:
self._thread.join(timeout=1.0)
return False
class IAAIScraper: class IAAIScraper:
@staticmethod @staticmethod
@@ -453,16 +482,23 @@ class IAAIScraper:
) )
for expected_page in range(2, target_page_number + 1): for expected_page in range(2, target_page_number + 1):
# Heartbeat для watchdog: при deep-resume (100+ кликов) задача может self._report_progress(
# идти несколько минут без batch_upserted, поэтому пульсуем прогресс. "listing_resume_progress",
if expected_page == 2 or expected_page == target_page_number or expected_page % 5 == 0: current_page=expected_page - 1,
self._report_progress( target_page=target_page_number,
"listing_resume_progress", next_page_number=expected_page,
current_page=expected_page - 1, )
target_page=target_page_number,
)
if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page): with _ProgressHeartbeat(
self._report_progress,
"listing_resume_progress",
current_page=expected_page - 1,
target_page=target_page_number,
next_page_number=expected_page,
):
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=expected_page)
if not next_ok:
self._report_progress( self._report_progress(
"listing_resume_failed", "listing_resume_failed",
current_page=expected_page - 1, current_page=expected_page - 1,
@@ -471,6 +507,13 @@ class IAAIScraper:
page.close() page.close()
raise ListingResumeError(f"Failed to resume listing at page {target_page_number}") raise ListingResumeError(f"Failed to resume listing at page {target_page_number}")
self._report_progress(
"listing_resume_progress",
current_page=expected_page,
target_page=target_page_number,
next_page_number=min(target_page_number, expected_page + 1),
)
self._report_progress( self._report_progress(
"listing_resume_completed", "listing_resume_completed",
current_page=target_page_number, current_page=target_page_number,
@@ -552,6 +595,7 @@ class IAAIScraper:
failures: list[dict[str, str]] = [] failures: list[dict[str, str]] = []
early_stopped = False early_stopped = False
truncated_by_time_budget = False truncated_by_time_budget = False
listing_recovery_attempts = 0
known_origin_ids: set[str] | None = None known_origin_ids: set[str] | None = None
threshold = self.settings.listing.early_stop_threshold threshold = self.settings.listing.early_stop_threshold
@@ -683,6 +727,21 @@ class IAAIScraper:
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)}) failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
return True return True
def _consume_listing_recovery_attempt(*, page_number: int, reason: str) -> None:
nonlocal listing_recovery_attempts
listing_recovery_attempts += 1
logger.warning(
"Listing recovery attempt %d/%d at page %d (%s)",
listing_recovery_attempts,
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT,
page_number,
reason,
)
if listing_recovery_attempts > _MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT:
raise ListingResumeError(
f"Listing recovery exhausted at page {page_number}: {reason}"
)
page = None page = None
try: try:
page, applied_filters = self._open_listing_for_stream( page, applied_filters = self._open_listing_for_stream(
@@ -721,6 +780,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason="context_rotation",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -752,6 +815,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason=f"collect_failed:{type(op_exc).__name__}",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -797,6 +864,10 @@ class IAAIScraper:
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
page.close() page.close()
try: try:
_consume_listing_recovery_attempt(
page_number=page_number,
reason="empty_page_after_retry",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number, target_page_number=page_number,
make=make, make=make,
@@ -973,8 +1044,17 @@ class IAAIScraper:
cars_failed=cars_failed, cars_failed=cars_failed,
) )
try: try:
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"): with _ProgressHeartbeat(
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1) self._report_progress,
"listing_next_page_started",
page_number=page_number,
next_page_number=page_number + 1,
pending_urls=len(pending_urls),
cars_upserted=cars_upserted,
cars_failed=cars_failed,
):
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"):
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1)
except (PageOperationTimeoutError, PlaywrightError) as nav_exc: except (PageOperationTimeoutError, PlaywrightError) as nav_exc:
logger.warning( logger.warning(
"go_to_next_page stalled/failed at page %d (%s) — will reopen", "go_to_next_page stalled/failed at page %d (%s) — will reopen",
@@ -1018,6 +1098,10 @@ class IAAIScraper:
pass pass
page = None page = None
try: try:
_consume_listing_recovery_attempt(
page_number=page_number + 1,
reason=f"next_page_failed:{type(nav_exc).__name__}",
)
page, _ = self._reopen_listing_and_resume( page, _ = self._reopen_listing_and_resume(
target_page_number=page_number + 1, target_page_number=page_number + 1,
make=make, make=make,

View File

@@ -20,6 +20,7 @@ CAR_DB_FIELDS = {
} }
_IN_CHUNK_SIZE = 5000 _IN_CHUNK_SIZE = 5000
CAR_TABLE_NAME = Car.__tablename__
class PersistenceService: class PersistenceService:
@@ -532,7 +533,7 @@ class PersistenceService:
if is_postgres: if is_postgres:
# Считаем кандидатов на sold. # Считаем кандидатов на sold.
total_active = session.execute( total_active = session.execute(
text("SELECT count(*) FROM cars WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%%'") text(f"SELECT count(*) FROM {CAR_TABLE_NAME} WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%%'")
).scalar() or 0 ).scalar() or 0
if total_active == 0: if total_active == 0:
@@ -556,12 +557,12 @@ class PersistenceService:
# Считаем будущие sold. # Считаем будущие sold.
would_mark = session.execute(text(""" would_mark = session.execute(text("""
SELECT count(*) SELECT count(*)
FROM cars c FROM {car_table} c
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
WHERE a.url IS NULL WHERE a.url IS NULL
AND c.is_sold = FALSE AND c.is_sold = FALSE
AND c.origin_id LIKE 'iaai:%%' AND c.origin_id LIKE 'iaai:%%'
""")).scalar() or 0 """.format(car_table=CAR_TABLE_NAME))).scalar() or 0
# Защита от аномалии. # Защита от аномалии.
if total_active > 100 and would_mark > total_active * 0.8: if total_active > 100 and would_mark > total_active * 0.8:
@@ -573,18 +574,18 @@ class PersistenceService:
# Массовая пометка sold. # Массовая пометка sold.
result = session.execute(text(""" result = session.execute(text("""
UPDATE cars UPDATE {car_table}
SET is_sold = TRUE SET is_sold = TRUE
FROM ( FROM (
SELECT c.id SELECT c.id
FROM cars c FROM {car_table} c
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
WHERE a.url IS NULL WHERE a.url IS NULL
AND c.is_sold = FALSE AND c.is_sold = FALSE
AND c.origin_id LIKE 'iaai:%%' AND c.origin_id LIKE 'iaai:%%'
) sub ) sub
WHERE cars.id = sub.id WHERE {car_table}.id = sub.id
""")) """.format(car_table=CAR_TABLE_NAME)))
count = result.rowcount or 0 count = result.rowcount or 0
else: else:
# Упрощённый путь для SQLite. # Упрощённый путь для SQLite.
@@ -672,4 +673,4 @@ class PersistenceService:
Car.is_sold == False, # noqa: E712 Car.is_sold == False, # noqa: E712
).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit) ).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit)
) )
return [str(row[0]) for row in result if row and row[0]] return [str(row[0]) for row in result if row and row[0]]

View File

@@ -20,10 +20,10 @@ class Base(DeclarativeBase):
class Car(Base): class Car(Base):
__tablename__ = "cars" __tablename__ = "iaai_cars"
__table_args__ = ( __table_args__ = (
Index("ix_cars_brand_model", "brand", "model"), Index("ix_iaai_cars_brand_model", "brand", "model"),
Index("ix_cars_origin_id_not_sold", "origin_id", "is_sold"), Index("ix_iaai_cars_origin_id_not_sold", "origin_id", "is_sold"),
) )
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True) parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
@@ -59,17 +59,17 @@ class Car(Base):
class Image(Base): class Image(Base):
__tablename__ = "images" __tablename__ = "iaai_images"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
fullres_image: Mapped[str] = mapped_column(String(), nullable=False) fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
preview_image: Mapped[str] = mapped_column(String(), nullable=False) preview_image: Mapped[str] = mapped_column(String(), nullable=False)
order_index: Mapped[int] = mapped_column(Integer, nullable=False) order_index: Mapped[int] = mapped_column(Integer, nullable=False)
car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False, index=True) car_id: Mapped[int] = mapped_column(Integer, ForeignKey("iaai_cars.id", ondelete="CASCADE"), nullable=False, index=True)
car: Mapped[Car] = relationship("Car", back_populates="images") car: Mapped[Car] = relationship("Car", back_populates="images")
class SyncRun(Base): class SyncRun(Base):
__tablename__ = "sync_runs" __tablename__ = "iaai_sync_runs"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)

View File

@@ -1,6 +1,7 @@
# Инициализация Celery-приложения и периодических задач. # Инициализация Celery-приложения и периодических задач.
import logging import logging
import os
from celery import Celery from celery import Celery
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
@@ -11,6 +12,12 @@ from ..core.logs import setup_logging
logger = logging.getLogger("iaai_scraper.worker.celery_app") logger = logging.getLogger("iaai_scraper.worker.celery_app")
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched" STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
IAAI_SYNC_QUEUE = "iaai_sync"
def _env_bool(name: str, default: bool) -> bool:
raw = os.getenv(name, "true" if default else "false").strip().lower()
return raw in {"1", "true", "yes", "on"}
@celery_setup_logging.connect @celery_setup_logging.connect
@@ -79,18 +86,19 @@ celery_app.conf.update(
worker_hijack_root_logger=False, worker_hijack_root_logger=False,
beat_schedule={ beat_schedule={
"periodic-sync-listing": { "periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task", "task": "iaai.sync_cars_feed",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0, "schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (), "args": (),
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True}, "kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
"options": { "options": {
"queue": "scraping", "queue": IAAI_SYNC_QUEUE,
"expires": settings.celery.beat_sync_interval_minutes * 60.0, "expires": settings.celery.beat_sync_interval_minutes * 60.0,
}, },
} }
}, },
task_routes={ task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"} "iaai.sync_cars_feed": {"queue": IAAI_SYNC_QUEUE},
"iaai_scraper.worker.tasks.*": {"queue": IAAI_SYNC_QUEUE},
}, },
) )
@@ -100,6 +108,10 @@ celery_app.autodiscover_tasks(["iaai_scraper.worker"])
@worker_ready.connect @worker_ready.connect
def _on_worker_ready(**kwargs): def _on_worker_ready(**kwargs):
"""При старте worker отправляем первый sync_listing, если очередь пуста.""" """При старте worker отправляем первый sync_listing, если очередь пуста."""
if not _env_bool("IAAI_STARTUP_SYNC_ENABLED", True):
logger.info("Worker ready: startup sync dispatch disabled by IAAI_STARTUP_SYNC_ENABLED")
return
redis_client = None redis_client = None
try: try:
redis_client = Redis.from_url( redis_client = Redis.from_url(
@@ -121,11 +133,11 @@ def _on_worker_ready(**kwargs):
logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True) logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True)
try: try:
queue_len = int(redis_client.llen("scraping") or 0) queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0)
except Exception: except Exception:
queue_len = 0 queue_len = 0
if queue_len > 0: if queue_len > 0:
logger.info("Worker ready: scraping queue already has %d task(s); skip startup dispatch", queue_len) logger.info("Worker ready: iaai_sync queue already has %d task(s); skip startup dispatch", queue_len)
return return
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
@@ -145,8 +157,8 @@ def _on_worker_ready(**kwargs):
logger.info("Worker ready — dispatching initial sync_listing task") logger.info("Worker ready — dispatching initial sync_listing task")
celery_app.send_task( celery_app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task", "iaai.sync_cars_feed",
kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False}, kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False},
queue="scraping", queue=IAAI_SYNC_QUEUE,
expires=settings.celery.beat_sync_interval_minutes * 60.0, expires=settings.celery.beat_sync_interval_minutes * 60.0,
) )

View File

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

View File

@@ -1,4 +1,4 @@
# Задачи Celery. # Задачи Celery для синхронизации автомобилей и листинга IAAI.
import json import json
import logging import logging
@@ -7,7 +7,6 @@ import signal
from threading import Event, Thread from threading import Event, Thread
import time import time
import uuid import uuid
from typing import Callable
from billiard.exceptions import SoftTimeLimitExceeded from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task from celery import shared_task
@@ -20,29 +19,17 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema
logger = logging.getLogger("iaai_scraper.worker.tasks") logger = logging.getLogger("iaai_scraper.worker.tasks")
# Пауза между батчами. IAAI_SYNC_QUEUE = "iaai_sync"
# Минимальная пауза между батчами (секунды) — не давит IAAI.
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0) _INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0)
# Порог ошибок в батче. # Если доля failed в батче превышает порог — прерываем (IAAI блокирует).
_FAIL_RATE_THRESHOLD = min(max(float(os.getenv("IAAI_FAIL_RATE_THRESHOLD", "0.9")), 0.0), 1.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_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_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_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
@@ -57,20 +44,25 @@ SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max(
) )
HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak" HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak"
HOURLY_FAILURE_STREAK_LIMIT = 3 HOURLY_FAILURE_STREAK_LIMIT = 3
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # Сброс через 6 часов. HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
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.sync_cars_feed"
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}" SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress" SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total" 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 SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}" TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count" SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset" SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment" STALL_WATCHDOG_NAVIGATION_STAGES = {
SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60 "listing_next_page_started",
"listing_resume_progress",
}
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
300,
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
)
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):
@@ -96,7 +88,10 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
def _run_browser_job(func, *args, **kwargs): def _run_browser_job(func, *args, **kwargs):
# Запускаем в текущем потоке. # С pool=solo Celery worker работает в одном процессе/потоке.
# Playwright sync API использует greenlets, которые привязаны к потоку.
# Запуск в отдельном потоке вызывает greenlet.error: cannot switch to a different thread.
# Поэтому запускаем напрямую в текущем потоке.
return func(*args, **kwargs) return func(*args, **kwargs)
@@ -104,7 +99,8 @@ def _sync_listing_lock_ttl_seconds() -> int:
settings = Settings() settings = Settings()
soft = settings.celery.task_soft_time_limit soft = settings.celery.task_soft_time_limit
hard = settings.celery.task_time_limit hard = settings.celery.task_time_limit
# Ограничиваем hard limit. # Используем clamped hard limit (soft + 120), а не сырой task_time_limit,
# чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999.
effective_hard = min(hard, soft + 120) if soft else hard effective_hard = min(hard, soft + 120) if soft else hard
return max(effective_hard + 120, 300) return max(effective_hard + 120, 300)
@@ -156,6 +152,12 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True) logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
return int(default_timeout)
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]: def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new del limit, only_new
settings = Settings() settings = Settings()
@@ -405,7 +407,6 @@ def _start_stall_watchdog(
stall_timeout_seconds: int, stall_timeout_seconds: int,
lock_key: str | None = None, lock_key: str | None = None,
lock_owner: str | None = None, lock_owner: str | None = None,
on_stall: Callable[[dict[str, object]], None] | None = None,
) -> tuple[Event, Thread]: ) -> tuple[Event, Thread]:
stop_event = Event() stop_event = Event()
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3)) interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
@@ -438,25 +439,22 @@ def _start_stall_watchdog(
else: else:
no_data_count = 0 no_data_count = 0
data = json.loads(raw) data = json.loads(raw)
stage = data.get("stage")
last_ts = int(data.get("ts") or 0) last_ts = int(data.get("ts") or 0)
if not last_ts: if not last_ts:
continue continue
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
age = int(time.time()) - last_ts age = int(time.time()) - last_ts
if age < stall_timeout_seconds: if age < effective_stall_timeout:
watchdog_born = time.monotonic() # reset absolute deadline on real progress watchdog_born = time.monotonic() # reset absolute deadline on real progress
continue continue
logger.error( logger.error(
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
task_id, task_id,
age, age,
data.get("stage"), stage,
data, 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: except Exception:
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
# Если Redis тоже не отвечает дольше дедлайна — убиваем. # Если Redis тоже не отвечает дольше дедлайна — убиваем.
@@ -486,7 +484,7 @@ def _start_stall_watchdog(
celery_app.send_task( celery_app.send_task(
SYNC_LISTING_TASK_NAME, SYNC_LISTING_TASK_NAME,
kwargs={}, kwargs={},
queue="scraping", queue=IAAI_SYNC_QUEUE,
countdown=15, countdown=15,
expires=followup_ttl, expires=followup_ttl,
) )
@@ -731,158 +729,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 _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: 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))))
@@ -1040,59 +886,17 @@ def _reset_segments_progress(redis_client: Redis, total: int) -> None:
logger.warning("Failed to reset segments progress", exc_info=True) 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]: def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
"""Помечает сегмент завершённым. Возвращает (completed_count, total).""" """Помечает сегмент завершённым. Возвращает (completed_count, total)."""
try: try:
pipe = redis_client.pipeline() pipe = redis_client.pipeline()
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index))) 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_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.scard(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.get(SYNC_SEGMENTS_TOTAL_KEY) pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
results = pipe.execute() results = pipe.execute()
completed = int(results[3] or 0) completed = int(results[2] or 0)
total = int(results[4] or 0) if results[4] else 0 total = int(results[3] or 0) if results[3] else 0
return completed, total return completed, total
except Exception: except Exception:
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True) logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
@@ -1101,6 +905,7 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in
@shared_task( @shared_task(
name="iaai_scraper.worker.tasks.sync_segment_task", name="iaai_scraper.worker.tasks.sync_segment_task",
queue=IAAI_SYNC_QUEUE,
bind=True, bind=True,
max_retries=2, max_retries=2,
default_retry_delay=60, default_retry_delay=60,
@@ -1274,6 +1079,7 @@ def sync_segment_task(
@shared_task( @shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task", name="iaai_scraper.worker.tasks.sync_vehicle_task",
queue=IAAI_SYNC_QUEUE,
bind=True, bind=True,
max_retries=2, max_retries=2,
default_retry_delay=30, default_retry_delay=30,
@@ -1304,7 +1110,8 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
@shared_task( @shared_task(
name="iaai_scraper.worker.tasks.sync_listing_task", name=SYNC_LISTING_TASK_NAME,
queue=IAAI_SYNC_QUEUE,
bind=True, bind=True,
max_retries=3, max_retries=3,
default_retry_delay=120, default_retry_delay=120,
@@ -1376,7 +1183,7 @@ def sync_listing_task(
try: try:
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300) followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
self.app.send_task( self.app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task", SYNC_LISTING_TASK_NAME,
kwargs={ kwargs={
"make": make, "make": make,
"model": model, "model": model,
@@ -1384,7 +1191,7 @@ def sync_listing_task(
"limit": limit, "limit": limit,
"only_new": only_new, "only_new": only_new,
}, },
queue="scraping", queue=IAAI_SYNC_QUEUE,
countdown=max(0, int(delay_seconds)), countdown=max(0, int(delay_seconds)),
expires=followup_expires, expires=followup_expires,
) )
@@ -1436,19 +1243,6 @@ def sync_listing_task(
stall_timeout_seconds=stall_timeout, stall_timeout_seconds=stall_timeout,
lock_key=SYNC_LISTING_LOCK_KEY, lock_key=SYNC_LISTING_LOCK_KEY,
lock_owner=owner_token, 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( _update_task_progress(
redis_client, redis_client,
@@ -1484,11 +1278,9 @@ def sync_listing_task(
last_completed_segment: int | None = None last_completed_segment: int | None = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
last_completed_segment = _load_last_completed_segment(redis_client) last_completed_segment = _load_last_completed_segment(redis_client)
stalled_segment = _load_stalled_segment(redis_client)
else: else:
# После завершения bootstrap чекпоинт не нужен никогда. # После завершения bootstrap чекпоинт не нужен никогда.
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
stalled_segment = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
logger.info( logger.info(
@@ -1649,17 +1441,12 @@ def sync_listing_task(
if use_segmented and settings.celery.parallel_segments: if use_segmented and settings.celery.parallel_segments:
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET. # Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
already_completed: set[int] = set() already_completed: set[int] = set()
cycle_id: str | None = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
try: try:
cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments))
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set() raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()} 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: except Exception:
already_completed = set() already_completed = set()
cycle_id = None
pending = [ pending = [
(idx, seg) for idx, seg in enumerate(segments) (idx, seg) for idx, seg in enumerate(segments)
@@ -1671,7 +1458,6 @@ def sync_listing_task(
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, True) _set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
_clear_segments_progress_state(redis_client)
_clear_bootstrap_failure_streak(redis_client) _clear_bootstrap_failure_streak(redis_client)
logger.warning("Parallel segments: nothing to dispatch (all completed)") logger.warning("Parallel segments: nothing to dispatch (all completed)")
return { return {
@@ -1683,6 +1469,10 @@ def sync_listing_task(
"segments_already_completed": len(already_completed), "segments_already_completed": len(already_completed),
} }
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
if force_bootstrap_full_scan and not already_completed:
_reset_segments_progress(redis_client, len(segments))
dispatched = 0 dispatched = 0
for idx, seg in pending: for idx, seg in pending:
try: try:
@@ -1695,15 +1485,15 @@ def sync_listing_task(
"only_new": effective_only_new, "only_new": effective_only_new,
"is_bootstrap": force_bootstrap_full_scan, "is_bootstrap": force_bootstrap_full_scan,
}, },
queue="scraping", queue=IAAI_SYNC_QUEUE,
) )
dispatched += 1 dispatched += 1
except Exception: except Exception:
logger.warning("Failed to dispatch segment %d", idx, exc_info=True) logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
logger.warning( logger.warning(
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)", "Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id, dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
) )
return { return {
"status": "success", "status": "success",
@@ -1712,11 +1502,9 @@ def sync_listing_task(
"segments_total": len(segments), "segments_total": len(segments),
"segments_dispatched": dispatched, "segments_dispatched": dispatched,
"segments_already_completed": len(already_completed), "segments_already_completed": len(already_completed),
"segments_cycle_id": cycle_id,
} }
# --- конец параллельной ветки --- # --- конец параллельной ветки ---
precomputed_result = None
resume_from_segment = 0 resume_from_segment = 0
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None: if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
resume_from_segment = max(0, last_completed_segment + 1) resume_from_segment = max(0, last_completed_segment + 1)
@@ -1734,117 +1522,6 @@ def sync_listing_task(
resume_from_segment, last_completed_segment, 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): def _progress_cb_main(stage, meta):
_update_task_progress( _update_task_progress(
redis_client, redis_client,
@@ -1877,16 +1554,13 @@ def sync_listing_task(
only_new=effective_only_new, only_new=effective_only_new,
) )
result = precomputed_result if precomputed_result is not None else _run_browser_job(_job) result = _run_browser_job(_job)
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
bootstrap_completed = bool(result.get("full_scan_completed")) bootstrap_completed = bool(result.get("full_scan_completed"))
if bootstrap_completed: if bootstrap_completed:
_set_full_scan_done(redis_client, False if always_full_scan else True) _set_full_scan_done(redis_client, False if always_full_scan else True)
_clear_sync_checkpoint(redis_client) _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_failure_streak(redis_client)
_clear_bootstrap_continuation_streak(redis_client) _clear_bootstrap_continuation_streak(redis_client)
if always_full_scan: if always_full_scan:
@@ -1902,7 +1576,7 @@ def sync_listing_task(
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered") for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
) or int(listing_payload.get("vehicles_collected") or 0) > 0 ) 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) 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_delay = 5
followup_reason = "bootstrap_not_completed" followup_reason = "bootstrap_not_completed"
if anti_bot_detected: if anti_bot_detected:
followup_delay = 180 followup_delay = 180
@@ -1922,7 +1596,6 @@ def sync_listing_task(
) )
else: else:
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
_clear_bootstrap_resume_target(redis_client)
summary = { summary = {
"task_id": task_id, "task_id": task_id,
@@ -1986,7 +1659,7 @@ def sync_listing_task(
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl): if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
try: try:
self.app.send_task( self.app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task", SYNC_LISTING_TASK_NAME,
kwargs={ kwargs={
"make": make, "make": make,
"model": model, "model": model,
@@ -1994,7 +1667,7 @@ def sync_listing_task(
"limit": limit, "limit": limit,
"only_new": only_new, "only_new": only_new,
}, },
queue="scraping", queue=IAAI_SYNC_QUEUE,
countdown=10, countdown=10,
expires=followup_ttl, expires=followup_ttl,
) )