294 lines
9.7 KiB
Python
294 lines
9.7 KiB
Python
# Celery tasks для синхронизации автомобилей Encar.
|
||
|
||
import logging
|
||
from threading import Event, Thread
|
||
import uuid
|
||
|
||
from celery import shared_task
|
||
from celery.signals import worker_process_init
|
||
from redis import Redis
|
||
|
||
from ..core.config import Settings
|
||
from ..core.runtime_config import load_runtime_config
|
||
from ..encar import EncarScraper, EncarFilters, expand_allowed_brands
|
||
from ..storage.db import PersistenceService
|
||
|
||
logger = logging.getLogger("encar_scraper.worker.tasks")
|
||
|
||
ENCAR_SYNC_LOCK_KEY = "encar:locks:sync_listing"
|
||
|
||
_persistence: PersistenceService | None = None
|
||
_redis: Redis | None = None
|
||
|
||
|
||
@worker_process_init.connect
|
||
def _reset_globals_after_fork(**kwargs):
|
||
global _persistence, _redis
|
||
_persistence = None
|
||
_redis = None
|
||
|
||
|
||
def _sync_listing_lock_ttl_seconds() -> int:
|
||
settings = Settings()
|
||
return max(settings.celery.task_time_limit + 120, 300)
|
||
|
||
|
||
def _get_persistence() -> PersistenceService:
|
||
global _persistence
|
||
if _persistence is None:
|
||
_persistence = PersistenceService(Settings())
|
||
return _persistence
|
||
|
||
|
||
def _get_redis() -> Redis:
|
||
global _redis
|
||
if _redis is None:
|
||
settings = Settings()
|
||
_redis = Redis.from_url(
|
||
settings.redis.url,
|
||
decode_responses=True,
|
||
socket_timeout=10,
|
||
socket_connect_timeout=5,
|
||
)
|
||
return _redis
|
||
|
||
|
||
def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool:
|
||
try:
|
||
return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
|
||
except Exception as exc:
|
||
logger.warning("Failed to acquire lock %s", key, exc_info=True)
|
||
return False
|
||
|
||
|
||
_REFRESH_LOCK_SCRIPT = """
|
||
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
||
return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
|
||
end
|
||
return 0
|
||
"""
|
||
|
||
_RELEASE_LOCK_SCRIPT = """
|
||
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
||
return redis.call('DEL', KEYS[1])
|
||
end
|
||
return 0
|
||
"""
|
||
|
||
|
||
def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None:
|
||
try:
|
||
script = redis_client.register_script(_REFRESH_LOCK_SCRIPT)
|
||
refreshed = script(keys=[key], args=[owner_token, int(ttl_seconds)])
|
||
return bool(refreshed)
|
||
except Exception as exc:
|
||
logger.warning("Failed to refresh lock %s", key, exc_info=True)
|
||
return None
|
||
|
||
|
||
def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None:
|
||
try:
|
||
script = redis_client.register_script(_RELEASE_LOCK_SCRIPT)
|
||
script(keys=[key], args=[owner_token])
|
||
except Exception as exc:
|
||
logger.warning("Failed to release lock %s", key, exc_info=True)
|
||
|
||
|
||
def _start_lock_heartbeat(
|
||
redis_client: Redis,
|
||
key: str,
|
||
owner_token: str,
|
||
ttl_seconds: int,
|
||
) -> tuple[Event, Thread]:
|
||
stop_event = Event()
|
||
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
||
|
||
def _heartbeat() -> None:
|
||
while not stop_event.wait(interval_seconds):
|
||
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
||
if refreshed is False:
|
||
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
||
return
|
||
|
||
thread = Thread(target=_heartbeat, name="encar-lock-heartbeat", daemon=True)
|
||
thread.start()
|
||
return stop_event, thread
|
||
|
||
|
||
@shared_task(
|
||
name="encar_scraper.worker.tasks.encar_sync_listing_task",
|
||
bind=True,
|
||
max_retries=2,
|
||
default_retry_delay=60,
|
||
acks_late=True,
|
||
)
|
||
def encar_sync_listing_task(
|
||
self,
|
||
car_type: str = "all",
|
||
manufacturer: str | None = None,
|
||
year_from: int | None = None,
|
||
year_to: int | None = None,
|
||
price_max: int | None = None,
|
||
lane: str = "encar",
|
||
limit: int | None = None,
|
||
):
|
||
"""
|
||
Полная синхронизация листинга Encar.
|
||
Собирает ВСЕ авто, обновляет цены/данные, помечает проданные.
|
||
"""
|
||
persistence = _get_persistence()
|
||
persistence.create_tables()
|
||
task_id = self.request.id or "unknown"
|
||
owner_token = f"{task_id}:{uuid.uuid4().hex}"
|
||
redis_client = _get_redis()
|
||
|
||
lock_ttl = _sync_listing_lock_ttl_seconds()
|
||
lock_acquired = _acquire_lock(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token, lock_ttl)
|
||
|
||
if not lock_acquired:
|
||
logger.info("encar_sync_listing_task skipped: another sync is already running")
|
||
return {
|
||
"status": "skipped",
|
||
"reason": "sync_already_running",
|
||
"task_id": task_id,
|
||
}
|
||
|
||
heartbeat_stop: Event | None = None
|
||
heartbeat_thread: Thread | None = None
|
||
run_id: int | None = None
|
||
|
||
try:
|
||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||
redis_client,
|
||
ENCAR_SYNC_LOCK_KEY,
|
||
owner_token,
|
||
lock_ttl,
|
||
)
|
||
self.update_state(state="STARTED", meta={"stage": "encar_sync_started", "task_id": task_id})
|
||
|
||
# --- Загрузка runtime_config ---
|
||
rc = load_runtime_config()
|
||
effective_limit = limit or rc.sync.limit
|
||
effective_lane = lane or rc.sync.lane or "encar"
|
||
|
||
# --- Фиксируем начало sync_run ---
|
||
run_id = persistence.start_sync_run(effective_lane)
|
||
|
||
car_type_map = {"all": "A", "domestic": "Y", "import": "N"}
|
||
filters = EncarFilters(
|
||
car_type=car_type_map.get(car_type, "A"),
|
||
manufacturer=manufacturer,
|
||
year_from=year_from,
|
||
year_to=year_to,
|
||
price_to=price_max,
|
||
)
|
||
|
||
# Набор разрешённых брендов (lowercased) из runtime_config
|
||
# Расширяем английские имена alias'ами (корейский + русский перевод)
|
||
allowed_brands = expand_allowed_brands(
|
||
set(rc.filters.brands) if rc.filters.brands else None
|
||
)
|
||
excluded_brands = expand_allowed_brands(
|
||
set(rc.filters.exclude_brands) if rc.filters.exclude_brands else None
|
||
)
|
||
|
||
scraper = EncarScraper(persistence=persistence)
|
||
result = scraper.sync_listing(
|
||
limit=effective_limit,
|
||
filters=filters,
|
||
lane=effective_lane,
|
||
redis_client=redis_client,
|
||
allowed_brands=allowed_brands,
|
||
excluded_brands=excluded_brands,
|
||
probe_all_photos=rc.sync.probe_all_photos,
|
||
)
|
||
|
||
# --- Запись результата в sync_runs ---
|
||
persistence.finish_sync_run(
|
||
run_id,
|
||
status="success",
|
||
ids_fetched=result.get("items_collected", 0),
|
||
cars_upserted=result.get("cars_synced", 0),
|
||
cars_failed=result.get("cars_failed", 0),
|
||
images_upserted=0,
|
||
)
|
||
|
||
summary = {
|
||
"task_id": task_id,
|
||
"source": "encar",
|
||
"total_available": result.get("total_available", 0),
|
||
"items_collected": result.get("items_collected", 0),
|
||
"cars_synced": result.get("cars_synced", 0),
|
||
"cars_failed": result.get("cars_failed", 0),
|
||
"cars_marked_sold": result.get("cars_marked_sold", 0),
|
||
}
|
||
logger.info(
|
||
"encar_sync_listing_task completed: %d synced, %d failed, %d marked sold",
|
||
summary["cars_synced"], summary["cars_failed"], summary["cars_marked_sold"],
|
||
)
|
||
return {"status": "success", **summary}
|
||
|
||
except Exception as exc:
|
||
logger.error("encar_sync_listing_task failed: %s", exc, exc_info=True)
|
||
try:
|
||
if run_id is not None:
|
||
persistence.finish_sync_run(
|
||
run_id,
|
||
status="failed",
|
||
ids_fetched=0,
|
||
cars_upserted=0,
|
||
cars_failed=0,
|
||
images_upserted=0,
|
||
error_summary=str(exc)[:500],
|
||
)
|
||
except Exception:
|
||
logger.warning("Failed to record sync_run failure", exc_info=True)
|
||
raise self.retry(exc=exc)
|
||
finally:
|
||
if heartbeat_stop is not None:
|
||
heartbeat_stop.set()
|
||
if heartbeat_thread is not None:
|
||
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
|
||
if lock_acquired:
|
||
_release_lock_if_owner(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token)
|
||
try:
|
||
interval_minutes = Settings().celery.encar_beat_interval_minutes
|
||
countdown = max(60, int(interval_minutes) * 60)
|
||
self.apply_async(
|
||
kwargs={"car_type": car_type},
|
||
countdown=countdown,
|
||
queue="encar",
|
||
)
|
||
logger.info(
|
||
"Next encar_sync_listing_task scheduled in %d minutes",
|
||
countdown // 60,
|
||
)
|
||
except Exception:
|
||
logger.warning("Failed to schedule next sync run", exc_info=True)
|
||
|
||
|
||
@shared_task(
|
||
name="encar_scraper.worker.tasks.encar_sync_vehicle_task",
|
||
bind=True,
|
||
max_retries=2,
|
||
default_retry_delay=30,
|
||
acks_late=True,
|
||
)
|
||
def encar_sync_vehicle_task(self, vehicle_url: str, lane: str = "encar"):
|
||
"""Синхронизация одного авто Encar."""
|
||
persistence = _get_persistence()
|
||
persistence.create_tables()
|
||
|
||
try:
|
||
scraper = EncarScraper(persistence=persistence)
|
||
result = scraper.sync_vehicle(vehicle_url, lane=lane)
|
||
logger.info("encar_sync_vehicle_task completed: %s", vehicle_url)
|
||
return {
|
||
"status": "success",
|
||
"vehicle_url": vehicle_url,
|
||
"origin_id": result.get("origin_id"),
|
||
}
|
||
except Exception as exc:
|
||
logger.error("encar_sync_vehicle_task failed: %s — %s", vehicle_url, exc, exc_info=True)
|
||
raise self.retry(exc=exc)
|