Update parser
This commit is contained in:
@@ -16,6 +16,9 @@ logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
PROGRESS_KEY_PREFIX = "iaai:state:task_progress:"
|
||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||
SYNC_LAST_COMPLETED_AT_KEY = "iaai:state:sync_listing_last_completed_at"
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
@@ -43,6 +46,26 @@ def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 18
|
||||
return False
|
||||
|
||||
|
||||
def _seconds_until_next_allowed_sync(redis_client: Redis) -> int:
|
||||
"""Запрещает новый автозапуск раньше чем через интервал beat после завершения полного run."""
|
||||
min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60))
|
||||
if min_interval <= 0:
|
||||
return 0
|
||||
try:
|
||||
full_done = str(redis_client.get(SYNC_FULL_SCAN_DONE_KEY) or "").strip().lower() in {"1", "true", "yes", "on"}
|
||||
if not full_done:
|
||||
return 0
|
||||
completed_raw = redis_client.get(SYNC_LAST_COMPLETED_AT_KEY)
|
||||
if not completed_raw:
|
||||
return 0
|
||||
completed_at = int(float(completed_raw))
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect last sync completion timestamp", exc_info=True)
|
||||
return 0
|
||||
elapsed = int(time.time()) - completed_at
|
||||
return max(0, min_interval - elapsed)
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
def _configure_logging(loglevel=None, **kwargs):
|
||||
# Перехватываем логирование Celery и пишем только в stderr (Docker logs).
|
||||
@@ -116,6 +139,7 @@ celery_app.conf.update(
|
||||
"options": {
|
||||
"queue": IAAI_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"headers": {"iaai_beat_task": True},
|
||||
},
|
||||
}
|
||||
},
|
||||
@@ -147,7 +171,15 @@ def _on_worker_ready(**kwargs):
|
||||
)
|
||||
|
||||
has_fresh_progress = _has_fresh_active_progress(redis_client)
|
||||
for stale_key in ("iaai:locks:sync_listing",):
|
||||
next_allowed_delay = _seconds_until_next_allowed_sync(redis_client)
|
||||
if next_allowed_delay > 0:
|
||||
logger.info(
|
||||
"Worker ready: last full sync finished recently; next auto sync allowed in %ss, skip startup dispatch",
|
||||
next_allowed_delay,
|
||||
)
|
||||
return
|
||||
|
||||
for stale_key in (SYNC_LISTING_LOCK_KEY,):
|
||||
try:
|
||||
ttl = redis_client.ttl(stale_key)
|
||||
if ttl is not None and ttl != -2 and not has_fresh_progress:
|
||||
|
||||
Reference in New Issue
Block a user