diff --git a/alembic/env.py b/alembic/env.py index f6f4ab0..9f9ae88 100644 --- a/alembic/env.py +++ b/alembic/env.py @@ -14,6 +14,7 @@ from iaai_scraper.core.config import Settings from iaai_scraper.storage.models import Base config = context.config +VERSION_TABLE = "iaai_parser_alembic_version" # Logging if config.config_file_name is not None: @@ -32,6 +33,7 @@ def run_migrations_offline() -> None: context.configure( url=url, target_metadata=target_metadata, + version_table=VERSION_TABLE, literal_binds=True, dialect_opts={"paramstyle": "named"}, ) @@ -47,7 +49,11 @@ def run_migrations_online() -> None: poolclass=pool.NullPool, ) 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(): context.run_migrations() diff --git a/alembic/versions/001_initial.py b/alembic/versions/001_initial.py index 4179d3e..454207b 100644 --- a/alembic/versions/001_initial.py +++ b/alembic/versions/001_initial.py @@ -1,4 +1,4 @@ -# Начальная схема: cars, images, sync_runs +# Начальная схема: iaai_cars, iaai_images, iaai_sync_runs from typing import Sequence, Union from alembic import op @@ -10,70 +10,94 @@ branch_labels: 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: - # Таблица cars — аукционные машины с IAAI - op.create_table( - "cars", - sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), - sa.Column("parser_id", sa.String(50), nullable=False, unique=True), - sa.Column("brand", sa.String(50), nullable=False), - sa.Column("model", sa.String(50), nullable=False), - sa.Column("year", sa.Integer, nullable=True), - sa.Column("price", sa.BigInteger, nullable=True), - sa.Column("currency", sa.String(10), nullable=False, server_default="USD"), - sa.Column("mileage", sa.Integer, nullable=False, server_default="0"), - sa.Column("country", sa.String(10), nullable=False, server_default="NA"), - sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("color", sa.String, nullable=False, server_default="other"), - sa.Column("drive", sa.String(10), nullable=True), - sa.Column("gearbox", sa.String(10), nullable=True), - sa.Column("steering_wheel", sa.String(10), nullable=True), - sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"), - sa.Column("engine_volume", sa.Integer, nullable=True), - sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"), - sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("origin", sa.String(20), nullable=False, server_default="NA"), - sa.Column("origin_url", sa.String, nullable=False), - sa.Column("origin_id", sa.String, nullable=False, unique=True), - sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("evaluation", sa.String, nullable=True), - sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")), - sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")), - sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")), - 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) + # Таблица iaai_cars — аукционные машины с IAAI + if not _table_exists("iaai_cars"): + op.create_table( + "iaai_cars", + sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), + sa.Column("parser_id", sa.String(50), nullable=False, unique=True), + sa.Column("brand", sa.String(50), nullable=False), + sa.Column("model", sa.String(50), nullable=False), + sa.Column("year", sa.Integer, nullable=True), + sa.Column("price", sa.BigInteger, nullable=True), + sa.Column("currency", sa.String(10), nullable=False, server_default="USD"), + sa.Column("mileage", sa.Integer, nullable=False, server_default="0"), + sa.Column("country", sa.String(10), nullable=False, server_default="NA"), + sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("color", sa.String, nullable=False, server_default="other"), + sa.Column("drive", sa.String(10), nullable=True), + sa.Column("gearbox", sa.String(10), nullable=True), + sa.Column("steering_wheel", sa.String(10), nullable=True), + sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"), + sa.Column("engine_volume", sa.Integer, nullable=True), + sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"), + sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("origin", sa.String(20), nullable=False, server_default="NA"), + sa.Column("origin_url", sa.String, nullable=False), + sa.Column("origin_id", sa.String, nullable=False, unique=True), + sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("evaluation", sa.String, nullable=True), + sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")), + sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("slug", sa.String, nullable=False), + sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + ) - # Таблица images — изображения машин - op.create_table( - "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("cars.id", ondelete="CASCADE"), nullable=False), - ) + if not _index_exists("iaai_cars", "ix_iaai_cars_origin_url"): + op.create_index("ix_iaai_cars_origin_url", "iaai_cars", ["origin_url"]) + if not _index_exists("iaai_cars", "ix_iaai_cars_origin_id"): + op.create_index("ix_iaai_cars_origin_id", "iaai_cars", ["origin_id"], unique=True) + + # Таблица iaai_images — изображения машин + if not _table_exists("iaai_images"): + 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: - op.drop_table("sync_runs") - op.drop_table("images") - op.drop_table("cars") + # Compatibility-only migration for shared production databases. + pass diff --git a/alembic/versions/002_add_indexes.py b/alembic/versions/002_add_indexes.py index 2c7b631..6fb5ad5 100644 --- a/alembic/versions/002_add_indexes.py +++ b/alembic/versions/002_add_indexes.py @@ -10,20 +10,15 @@ depends_on: Union[str, Sequence[str], None] = None def upgrade() -> None: - op.create_index("ix_cars_brand", "cars", ["brand"]) - op.create_index("ix_cars_brand_model", "cars", ["brand", "model"]) - op.create_index("ix_cars_year", "cars", ["year"]) - op.create_index("ix_cars_is_sold", "cars", ["is_sold"]) - op.create_index("ix_cars_last_seen_at", "cars", ["last_seen_at"]) - op.create_index("ix_images_car_id", "images", ["car_id"]) - op.create_index("ix_sync_runs_status", "sync_runs", ["status"]) + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand ON iaai_cars (brand)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand_model ON iaai_cars (brand, model)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_year ON iaai_cars (year)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_is_sold ON iaai_cars (is_sold)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_last_seen_at ON iaai_cars (last_seen_at)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id ON iaai_images (car_id)") + op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_status ON iaai_sync_runs (status)") def downgrade() -> None: - op.drop_index("ix_sync_runs_status", table_name="sync_runs") - op.drop_index("ix_images_car_id", table_name="images") - 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") + # Compatibility-only migration for shared production databases. + pass diff --git a/alembic/versions/003_add_composite_indexes.py b/alembic/versions/003_add_composite_indexes.py index 1388d45..393f08c 100644 --- a/alembic/versions/003_add_composite_indexes.py +++ b/alembic/versions/003_add_composite_indexes.py @@ -13,26 +13,24 @@ def upgrade() -> None: # Partial index для активных машин IAAI (is_sold = FALSE) # Ускоряет поиск при mark_sold операциях op.execute( - "CREATE INDEX IF NOT EXISTS ix_cars_active_iaai " - "ON cars (origin_url) " + "CREATE INDEX IF NOT EXISTS ix_iaai_cars_active_iaai " + "ON iaai_cars (origin_url) " "WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'" ) # Composite index для проверки дубликатов изображений - op.create_index( - "ix_images_car_id_fullres", - "images", - ["car_id", "fullres_image"], + op.execute( + "CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id_fullres " + "ON iaai_images (car_id, fullres_image)" ) # Index для быстрого поиска existing машин по origin_url op.execute( - "CREATE INDEX IF NOT EXISTS ix_cars_origin_url_id " - "ON cars (origin_url, origin_id)" + "CREATE INDEX IF NOT EXISTS ix_iaai_cars_origin_url_id " + "ON iaai_cars (origin_url, origin_id)" ) def downgrade() -> None: - op.execute("DROP INDEX IF EXISTS ix_cars_origin_url_id") - op.drop_index("ix_images_car_id_fullres", table_name="images") - op.execute("DROP INDEX IF EXISTS ix_cars_active_iaai") + # Compatibility-only migration for shared production databases. + pass diff --git a/iaai_scraper/storage/db.py b/iaai_scraper/storage/db.py index a1740f7..ffb69ca 100644 --- a/iaai_scraper/storage/db.py +++ b/iaai_scraper/storage/db.py @@ -20,6 +20,7 @@ CAR_DB_FIELDS = { } _IN_CHUNK_SIZE = 5000 +CAR_TABLE_NAME = Car.__tablename__ class PersistenceService: @@ -532,7 +533,7 @@ class PersistenceService: if is_postgres: # Считаем кандидатов на sold. 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 if total_active == 0: @@ -556,12 +557,12 @@ class PersistenceService: # Считаем будущие sold. would_mark = session.execute(text(""" SELECT count(*) - FROM cars c + FROM {car_table} c LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url WHERE a.url IS NULL AND c.is_sold = FALSE 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: @@ -573,18 +574,18 @@ class PersistenceService: # Массовая пометка sold. result = session.execute(text(""" - UPDATE cars + UPDATE {car_table} SET is_sold = TRUE FROM ( SELECT c.id - FROM cars c + FROM {car_table} c LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url WHERE a.url IS NULL AND c.is_sold = FALSE AND c.origin_id LIKE 'iaai:%%' ) sub - WHERE cars.id = sub.id - """)) + WHERE {car_table}.id = sub.id + """.format(car_table=CAR_TABLE_NAME))) count = result.rowcount or 0 else: # Упрощённый путь для SQLite. @@ -672,4 +673,4 @@ class PersistenceService: Car.is_sold == False, # noqa: E712 ).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit) ) - return [str(row[0]) for row in result if row and row[0]] \ No newline at end of file + return [str(row[0]) for row in result if row and row[0]] diff --git a/iaai_scraper/storage/models.py b/iaai_scraper/storage/models.py index 29d5566..1970dd0 100644 --- a/iaai_scraper/storage/models.py +++ b/iaai_scraper/storage/models.py @@ -20,10 +20,10 @@ class Base(DeclarativeBase): class Car(Base): - __tablename__ = "cars" + __tablename__ = "iaai_cars" __table_args__ = ( - Index("ix_cars_brand_model", "brand", "model"), - Index("ix_cars_origin_id_not_sold", "origin_id", "is_sold"), + Index("ix_iaai_cars_brand_model", "brand", "model"), + 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) parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True) @@ -59,17 +59,17 @@ class Car(Base): class Image(Base): - __tablename__ = "images" + __tablename__ = "iaai_images" id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) fullres_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) - 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") 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) 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) diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 66363f2..ca4f320 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -1,6 +1,7 @@ # Инициализация Celery-приложения и периодических задач. import logging +import os from celery import Celery 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") 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 @@ -79,18 +86,19 @@ celery_app.conf.update( worker_hijack_root_logger=False, beat_schedule={ "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, "args": (), "kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True}, "options": { - "queue": "scraping", + "queue": IAAI_SYNC_QUEUE, "expires": settings.celery.beat_sync_interval_minutes * 60.0, }, } }, 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 def _on_worker_ready(**kwargs): """При старте 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 try: 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) try: - queue_len = int(redis_client.llen("scraping") or 0) + queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0) except Exception: 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 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") 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}, - queue="scraping", + queue=IAAI_SYNC_QUEUE, expires=settings.celery.beat_sync_interval_minutes * 60.0, ) diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 5edaf46..9fec704 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -19,6 +19,8 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema 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) # Если доля failed в батче превышает порог — прерываем (IAAI блокирует). @@ -44,7 +46,7 @@ HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak" HOURLY_FAILURE_STREAK_LIMIT = 3 HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" -SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" +SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed" SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}" SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress" SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total" @@ -482,7 +484,7 @@ def _start_stall_watchdog( celery_app.send_task( SYNC_LISTING_TASK_NAME, kwargs={}, - queue="scraping", + queue=IAAI_SYNC_QUEUE, countdown=15, expires=followup_ttl, ) @@ -903,6 +905,7 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in @shared_task( name="iaai_scraper.worker.tasks.sync_segment_task", + queue=IAAI_SYNC_QUEUE, bind=True, max_retries=2, default_retry_delay=60, @@ -1076,6 +1079,7 @@ def sync_segment_task( @shared_task( name="iaai_scraper.worker.tasks.sync_vehicle_task", + queue=IAAI_SYNC_QUEUE, bind=True, max_retries=2, default_retry_delay=30, @@ -1106,7 +1110,8 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): @shared_task( - name="iaai_scraper.worker.tasks.sync_listing_task", + name=SYNC_LISTING_TASK_NAME, + queue=IAAI_SYNC_QUEUE, bind=True, max_retries=3, default_retry_delay=120, @@ -1178,7 +1183,7 @@ def sync_listing_task( try: followup_expires = max(int(lock_ttl), int(delay_seconds) + 300) self.app.send_task( - "iaai_scraper.worker.tasks.sync_listing_task", + SYNC_LISTING_TASK_NAME, kwargs={ "make": make, "model": model, @@ -1186,7 +1191,7 @@ def sync_listing_task( "limit": limit, "only_new": only_new, }, - queue="scraping", + queue=IAAI_SYNC_QUEUE, countdown=max(0, int(delay_seconds)), expires=followup_expires, ) @@ -1480,7 +1485,7 @@ def sync_listing_task( "only_new": effective_only_new, "is_bootstrap": force_bootstrap_full_scan, }, - queue="scraping", + queue=IAAI_SYNC_QUEUE, ) dispatched += 1 except Exception: @@ -1654,7 +1659,7 @@ def sync_listing_task( if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl): try: self.app.send_task( - "iaai_scraper.worker.tasks.sync_listing_task", + SYNC_LISTING_TASK_NAME, kwargs={ "make": make, "model": model, @@ -1662,7 +1667,7 @@ def sync_listing_task( "limit": limit, "only_new": only_new, }, - queue="scraping", + queue=IAAI_SYNC_QUEUE, countdown=10, expires=followup_ttl, )