Update model translations

This commit is contained in:
qananasikq
2026-07-01 14:35:49 +03:00
commit 3904cd003e
53 changed files with 7629 additions and 0 deletions

View File

@@ -0,0 +1,3 @@
from .celery_app import celery_app
__all__ = ["celery_app"]

View File

@@ -0,0 +1,111 @@
# Инициализация Celery-приложения и периодических задач для Encar.
import logging
from celery import Celery
from celery.signals import beat_init
from redis import Redis
from sqlalchemy import func, select
from ..core.config import settings
logger = logging.getLogger("encar_scraper.worker.celery_app")
_INITIAL_SYNC_BOOTSTRAP_KEY = "encar:bootstrap:initial_sync_enqueued"
def _broker_url() -> str:
return settings.celery.broker_url or settings.redis.url
def _result_backend() -> str:
return settings.celery.result_backend or settings.redis.url
celery_app = Celery(
"encar_scraper",
broker=_broker_url(),
backend=_result_backend(),
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
worker_concurrency=settings.celery.worker_concurrency,
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
worker_pool="prefork",
worker_prefetch_multiplier=1,
broker_connection_retry_on_startup=True,
broker_transport_options={
"visibility_timeout": settings.celery.broker_visibility_timeout,
},
result_expires=86400,
worker_hijack_root_logger=False,
# Следующий запуск планирует сама задача в finally через apply_async(countdown=...),
# чтобы час отсчитывался от завершения прогона, а не от старта.
beat_schedule={},
task_routes={
"encar_scraper.worker.tasks.*": {"queue": "encar"},
},
)
celery_app.autodiscover_tasks(["encar_scraper.worker"])
@beat_init.connect
def _schedule_initial_sync_on_first_start(**kwargs):
"""Ставит первый sync сразу после первого запуска проекта на пустой БД."""
from ..storage.db import PersistenceService
from ..storage.models import SyncRun
from .tasks import encar_sync_listing_task
persistence = PersistenceService(settings)
redis_client: Redis | None = None
try:
with persistence.session_scope() as session:
total_runs = session.execute(select(func.count(SyncRun.id))).scalar() or 0
if total_runs:
return
redis_client = Redis.from_url(
_broker_url(),
decode_responses=True,
socket_timeout=10,
socket_connect_timeout=5,
)
marked = redis_client.set(_INITIAL_SYNC_BOOTSTRAP_KEY, "1", nx=True, ex=600)
if not marked:
logger.info("Initial Encar sync already bootstrapped on this startup")
return
encar_sync_listing_task.apply_async(kwargs={"car_type": "all"}, queue="encar")
logger.info("Scheduled initial Encar sync because sync_runs table is empty")
except Exception:
logger.warning("Failed to schedule initial Encar sync", exc_info=True)
finally:
persistence.engine.dispose()
if redis_client is not None:
try:
redis_client.close()
except Exception:
pass
from celery.signals import setup_logging as celery_setup_logging
@celery_setup_logging.connect
def _configure_logging(loglevel, logfile, format, colorize, **kwargs):
"""Перехватываем настройку логирования Celery, чтобы использовать свой формат."""
from ..core.logs import setup_logging
import logging
level_name = logging.getLevelName(loglevel) if isinstance(loglevel, int) else str(loglevel)
setup_logging(level=level_name, log_file=logfile)

View File

@@ -0,0 +1,319 @@
# Celery tasks для синхронизации автомобилей Encar.
import logging
from threading import Event, Thread
import uuid
from celery import shared_task
from celery.exceptions import MaxRetriesExceededError, Retry
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
# После fork закрываем унаследованные TCP-сокеты родителя (иначе дети
# делят одни connection pool'ы → повреждение данных).
if _persistence is not None:
try:
_persistence.engine.dispose()
except Exception:
logger.debug("engine.dispose after fork failed", exc_info=True)
if _redis is not None:
try:
_redis.close()
except Exception:
logger.debug("redis.close after fork failed", exc_info=True)
_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
# Планировать следующий прогон только когда текущий прогон завершён
# (success или final failure после исчерпания retries). Retry сам
# перепланирует задачу — не дублируем apply_async в этом случае.
reschedule_next = True
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,
runtime_filters=rc.filters,
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)
try:
# Retry сам перепланирует задачу — не планируем следующий прогон
reschedule_next = False
raise self.retry(exc=exc)
except MaxRetriesExceededError:
# Retries исчерпаны — запускаем следующий прогон по расписанию
reschedule_next = True
raise
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)
if reschedule_next:
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)