diff --git a/.gitignore b/.gitignore index 86adc10..1331060 100644 --- a/.gitignore +++ b/.gitignore @@ -2,6 +2,8 @@ __pycache__/ .pytest_cache/ .coverage coverage.xml +.env +.env.* .venv/ venv/ *.pyc @@ -16,4 +18,5 @@ dist/ build/ artifacts/ celerybeat-schedule* -tokens_data/ \ No newline at end of file +tokens_data/ +tokens.json \ No newline at end of file diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index 32a312d..f1c324b 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -18,15 +18,22 @@ def set_trace_id(trace_id: str) -> None: def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: - handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)] + # stderr — Docker и Celery prefork корректно его подхватывают. + handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)] if log_file: handlers.append(logging.FileHandler(log_file, encoding="utf-8")) trace_filter = TraceIdFilter() for handler in handlers: handler.addFilter(trace_filter) - logging.basicConfig( - level=getattr(logging, level.upper(), logging.INFO), - format="%(asctime)s | %(levelname)s | %(name)s | trace=%(trace_id)s | %(message)s", - handlers=handlers, - force=True, + handler.setLevel(getattr(logging, level.upper(), logging.INFO)) + root = logging.getLogger() + root.setLevel(getattr(logging, level.upper(), logging.INFO)) + # Убираем старые хендлеры, чтобы не дублировать после fork. + root.handlers.clear() + for handler in handlers: + root.addHandler(handler) + fmt = logging.Formatter( + "%(asctime)s | %(levelname)s | %(name)s | trace=%(trace_id)s | %(message)s" ) + for handler in root.handlers: + handler.setFormatter(fmt) diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 4f5682b..b1c4530 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -1,8 +1,25 @@ # Инициализация Celery-приложения и периодических задач. from celery import Celery +from celery.signals import worker_process_init, setup_logging as celery_setup_logging from ..core.config import settings +from ..core.logs import setup_logging + + +@celery_setup_logging.connect +def _configure_logging(loglevel=None, **kwargs): + # Перехватываем логирование Celery и пишем только в stderr (Docker logs). + level = settings.log_level if settings.log_level else "INFO" + setup_logging(level, None) + + +@worker_process_init.connect +def _on_worker_process_init(**kwargs): + # Повторно настраиваем логирование в каждом дочернем prefork-процессе, + # чтобы StreamHandler(stderr) корректно работал после fork. + level = settings.log_level if settings.log_level else "INFO" + setup_logging(level, None) def _broker_url() -> str: @@ -39,6 +56,8 @@ celery_app.conf.update( "visibility_timeout": settings.celery.broker_visibility_timeout, }, result_expires=86400, + worker_redirect_stdouts=False, + worker_hijack_root_logger=False, beat_schedule={ "periodic-sync-listing": { "task": "iaai_scraper.worker.tasks.sync_listing_task", diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index e244f3a..77fa051 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -18,6 +18,28 @@ logger = logging.getLogger("iaai_scraper.worker.tasks") SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" +def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): + last_exc: Exception | None = None + for attempt in range(1, attempts + 1): + try: + return func() + except Exception as exc: + last_exc = exc + if attempt >= attempts: + break + delay = base_delay_s * (2 ** (attempt - 1)) + logger.warning( + "Operation failed (attempt %d/%d): %s. Retrying in %.1fs", + attempt, + attempts, + exc, + delay, + ) + time.sleep(delay) + if last_exc is not None: + raise last_exc + + def _run_browser_job(func, *args, **kwargs): # Браузерный код запускаем в отдельном потоке без активного loop. with ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") as executor: @@ -32,17 +54,42 @@ def _sync_listing_lock_ttl_seconds() -> int: def _get_persistence() -> PersistenceService: - return PersistenceService(Settings()) + settings = Settings() + persistence = PersistenceService(settings) + + def _ping_db() -> None: + with persistence.engine.connect() as conn: + conn.exec_driver_sql("SELECT 1") + + _retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0) + return persistence def _get_redis() -> Redis: settings = Settings() - return Redis.from_url(settings.redis.url, decode_responses=True) + redis_client = Redis.from_url(settings.redis.url, decode_responses=True) + + def _ping_redis() -> None: + redis_client.ping() + + _retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0) + return redis_client 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)) + acquired = bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds)) + if acquired: + return True + + # Автовосстановление: если lock завис без TTL, считаем stale и пересоздаём. + ttl = redis_client.ttl(key) + if ttl is not None and ttl < 0: + logger.warning("Detected stale lock without TTL, removing: %s", key) + redis_client.delete(key) + return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds)) + + return False except Exception as exc: logger.warning("Failed to acquire lock %s", key, exc_info=True) return False