Compare commits
5 Commits
3299f8f372
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d9c95b1a9f | ||
|
|
00d33d527c | ||
|
|
f278bd1a47 | ||
|
|
0905472991 | ||
|
|
9a33efa89b |
@@ -21,7 +21,7 @@ COPY . .
|
|||||||
RUN pip install --no-cache-dir -e .
|
RUN pip install --no-cache-dir -e .
|
||||||
RUN python -m compileall -q iaai_scraper
|
RUN python -m compileall -q iaai_scraper
|
||||||
RUN chmod +x entrypoint.sh
|
RUN chmod +x entrypoint.sh
|
||||||
RUN chown -R app:app /app
|
RUN mkdir -p /data && chown -R app:app /app /data
|
||||||
|
|
||||||
STOPSIGNAL SIGINT
|
STOPSIGNAL SIGINT
|
||||||
|
|
||||||
|
|||||||
@@ -5,31 +5,43 @@ x-app-env: &app-env
|
|||||||
CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0}
|
CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0}
|
||||||
CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0}
|
CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0}
|
||||||
IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800}
|
IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800}
|
||||||
|
IAAI_DATABASE_POOL_SIZE: ${IAAI_DATABASE_POOL_SIZE:-20}
|
||||||
|
IAAI_DATABASE_MAX_OVERFLOW: ${IAAI_DATABASE_MAX_OVERFLOW:-40}
|
||||||
# Для больших Search-выборок (десятки тысяч лотов) run должен жить дольше одного батча.
|
# Для больших Search-выборок (десятки тысяч лотов) run должен жить дольше одного батча.
|
||||||
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-7200}
|
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-7200}
|
||||||
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-7500}
|
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-7500}
|
||||||
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-10800}
|
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-10800}
|
||||||
|
IAAI_PROFILE: ${IAAI_PROFILE:-fast}
|
||||||
|
IAAI_SCRAPING_PROFILE: ${IAAI_SCRAPING_PROFILE:-${IAAI_PROFILE:-fast}}
|
||||||
|
IAAI_HTTP_FIRST: ${IAAI_HTTP_FIRST:-true}
|
||||||
|
IAAI_BROWSER_FALLBACK_ENABLED: ${IAAI_BROWSER_FALLBACK_ENABLED:-false}
|
||||||
|
IAAI_ANONYMOUS_BOOTSTRAP_ENABLED: ${IAAI_ANONYMOUS_BOOTSTRAP_ENABLED:-false}
|
||||||
|
IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1}
|
||||||
|
IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-1}
|
||||||
|
IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0.03}
|
||||||
|
IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-2}
|
||||||
|
IAAI_LISTING_RETRIES: ${IAAI_LISTING_RETRIES:-1}
|
||||||
# 0 = без лимита.
|
# 0 = без лимита.
|
||||||
CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0}
|
CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0}
|
||||||
# Следующий плановый прогон — не раньше чем через час после предыдущего tick beat.
|
# Следующий плановый прогон — не раньше чем через час после предыдущего tick beat.
|
||||||
CELERY_BEAT_SYNC_INTERVAL_MINUTES: ${CELERY_BEAT_SYNC_INTERVAL_MINUTES:-60}
|
CELERY_BEAT_SYNC_INTERVAL_MINUTES: ${CELERY_BEAT_SYNC_INTERVAL_MINUTES:-60}
|
||||||
# Любой внутренний follow-up после неполного bootstrap тоже не стартует сразу.
|
# Любой внутренний follow-up после неполного bootstrap тоже не стартует сразу.
|
||||||
IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS: ${IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS:-3600}
|
IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS: ${IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS:-3600}
|
||||||
|
CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4}
|
||||||
|
CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork}
|
||||||
CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
|
CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
|
||||||
CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1500}
|
CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1500}
|
||||||
IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0.3}
|
IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0}
|
||||||
# Стабильный fast-профиль: не душим IAAI слишком большим числом HTTPS detail-соединений.
|
# Стабильный fast-профиль: не душим IAAI слишком большим числом HTTPS detail-соединений.
|
||||||
IAAI_SCRAPING_PROFILE: fast
|
|
||||||
IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-96}
|
IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-96}
|
||||||
IAAI_MAX_RETRIES: ${IAAI_MAX_RETRIES:-2}
|
IAAI_MAX_RETRIES: ${IAAI_MAX_RETRIES:-2}
|
||||||
IAAI_RETRY_DELAY_SECONDS: ${IAAI_RETRY_DELAY_SECONDS:-0.35}
|
IAAI_RETRY_DELAY_SECONDS: ${IAAI_RETRY_DELAY_SECONDS:-0.35}
|
||||||
IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1}
|
|
||||||
CELERY_PARALLEL_SEGMENTS: ${CELERY_PARALLEL_SEGMENTS:-false}
|
CELERY_PARALLEL_SEGMENTS: ${CELERY_PARALLEL_SEGMENTS:-false}
|
||||||
CELERY_WORKER_CONCURRENCY: 4
|
|
||||||
CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork}
|
|
||||||
IAAI_FAIL_RATE_THRESHOLD: ${IAAI_FAIL_RATE_THRESHOLD:-0.9}
|
IAAI_FAIL_RATE_THRESHOLD: ${IAAI_FAIL_RATE_THRESHOLD:-0.9}
|
||||||
IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-8}
|
IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-8}
|
||||||
IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true}
|
IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true}
|
||||||
|
IAAI_FILTERED_SEARCH_URL: ${IAAI_FILTERED_SEARCH_URL:-}
|
||||||
|
IAAI_FILTERED_SEARCH_URLS: ${IAAI_FILTERED_SEARCH_URLS:-}
|
||||||
IAAI_FAST_PATH_TIMEOUT_MS: ${IAAI_FAST_PATH_TIMEOUT_MS:-10000}
|
IAAI_FAST_PATH_TIMEOUT_MS: ${IAAI_FAST_PATH_TIMEOUT_MS:-10000}
|
||||||
IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-[]}
|
IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-[]}
|
||||||
IAAI_DISCOVERY_MODE: ${IAAI_DISCOVERY_MODE:-listing}
|
IAAI_DISCOVERY_MODE: ${IAAI_DISCOVERY_MODE:-listing}
|
||||||
@@ -50,6 +62,8 @@ x-app-env: &app-env
|
|||||||
IAAI_SELF_HEAL_STALL_SECONDS: ${IAAI_SELF_HEAL_STALL_SECONDS:-900}
|
IAAI_SELF_HEAL_STALL_SECONDS: ${IAAI_SELF_HEAL_STALL_SECONDS:-900}
|
||||||
IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS: ${IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS:-300}
|
IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS: ${IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS:-300}
|
||||||
IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS: ${IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS:-300}
|
IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS: ${IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS:-300}
|
||||||
|
# Если в Postgres нет записей дольше 60 минут при активной работе — полный restart с segment 1.
|
||||||
|
IAAI_DB_IDLE_RESTART_SECONDS: ${IAAI_DB_IDLE_RESTART_SECONDS:-3600}
|
||||||
TZ: ${TZ:-UTC}
|
TZ: ${TZ:-UTC}
|
||||||
|
|
||||||
x-env-file: &env-file
|
x-env-file: &env-file
|
||||||
|
|||||||
@@ -135,10 +135,11 @@ celery_app.conf.update(
|
|||||||
"task": "iaai.sync_cars_feed",
|
"task": "iaai.sync_cars_feed",
|
||||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||||
"args": (),
|
"args": (),
|
||||||
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": False},
|
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
|
||||||
"options": {
|
"options": {
|
||||||
"queue": IAAI_SYNC_QUEUE,
|
"queue": IAAI_SYNC_QUEUE,
|
||||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||||
|
"headers": {"iaai_beat_task": True},
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -12,7 +12,8 @@ from billiard.exceptions import SoftTimeLimitExceeded
|
|||||||
from celery import shared_task
|
from celery import shared_task
|
||||||
from redis import Redis
|
from redis import Redis
|
||||||
|
|
||||||
from ..core.config import Settings, parse_listing_segments
|
from ..core.config import Settings, build_fast_listing_segments_for_makes, build_listing_segments_for_makes, parse_listing_segments
|
||||||
|
from ..core.runtime_config import RuntimeConfig
|
||||||
from ..scraper import IAAIScraper
|
from ..scraper import IAAIScraper
|
||||||
from ..storage.db import PersistenceService
|
from ..storage.db import PersistenceService
|
||||||
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
|
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
|
||||||
@@ -54,16 +55,35 @@ SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
|
|||||||
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
||||||
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
||||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||||
|
GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts"
|
||||||
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
||||||
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
||||||
STALL_WATCHDOG_NAVIGATION_STAGES = {
|
STALL_WATCHDOG_NAVIGATION_STAGES = {
|
||||||
"listing_next_page_started",
|
"listing_next_page_started",
|
||||||
"listing_resume_progress",
|
"listing_resume_progress",
|
||||||
}
|
}
|
||||||
|
STALL_WATCHDOG_LONG_RUNNING_STAGES = {
|
||||||
|
"fast_listing_collected",
|
||||||
|
"fast_detail_progress",
|
||||||
|
}
|
||||||
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
|
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
|
||||||
300,
|
300,
|
||||||
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
|
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
|
||||||
)
|
)
|
||||||
|
STALL_WATCHDOG_DETAIL_GRACE_SECONDS = max(
|
||||||
|
600,
|
||||||
|
int(os.getenv("STALL_WATCHDOG_DETAIL_GRACE_SECONDS", "1200")),
|
||||||
|
)
|
||||||
|
DB_IDLE_RESTART_SECONDS = max(60, int(os.getenv("IAAI_DB_IDLE_RESTART_SECONDS", "3600")))
|
||||||
|
TERMINAL_PROGRESS_STAGES = {
|
||||||
|
"segment_done",
|
||||||
|
"segment_failed",
|
||||||
|
"segment_task_completed",
|
||||||
|
"segment_task_failed",
|
||||||
|
"segment_task_soft_timeout",
|
||||||
|
"sync_done",
|
||||||
|
"failed",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||||
@@ -124,6 +144,19 @@ def _update_task_progress(
|
|||||||
) -> None:
|
) -> None:
|
||||||
try:
|
try:
|
||||||
now_ts = int(time.time())
|
now_ts = int(time.time())
|
||||||
|
existing_task_started_ts: int | None = None
|
||||||
|
try:
|
||||||
|
existing_raw = redis_client.get(_task_progress_key(task_id))
|
||||||
|
if existing_raw:
|
||||||
|
existing_payload = json.loads(existing_raw)
|
||||||
|
existing_task_started_ts = _safe_int(existing_payload.get("task_started_ts"))
|
||||||
|
except Exception:
|
||||||
|
existing_task_started_ts = None
|
||||||
|
payload.setdefault("task_started_ts", existing_task_started_ts or now_ts)
|
||||||
|
if stage != "fast_db_progress" and "last_db_progress_ts" not in payload:
|
||||||
|
last_db_progress_ts = _safe_int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY))
|
||||||
|
if last_db_progress_ts is not None:
|
||||||
|
payload["last_db_progress_ts"] = last_db_progress_ts
|
||||||
data = {
|
data = {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
"stage": stage,
|
"stage": stage,
|
||||||
@@ -141,6 +174,8 @@ def _update_task_progress(
|
|||||||
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
|
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
|
||||||
# (когда PID жив, но прогресс по задачам не двигается).
|
# (когда PID жив, но прогресс по задачам не двигается).
|
||||||
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||||
|
if stage == "fast_db_progress":
|
||||||
|
pipe.set(GLOBAL_DB_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||||
pipe.execute()
|
pipe.execute()
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
|
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
|
||||||
@@ -156,6 +191,8 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
|||||||
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
|
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
|
||||||
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
|
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
|
||||||
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
|
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
|
||||||
|
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
|
||||||
|
return max(int(default_timeout), STALL_WATCHDOG_DETAIL_GRACE_SECONDS)
|
||||||
return int(default_timeout)
|
return int(default_timeout)
|
||||||
|
|
||||||
|
|
||||||
@@ -408,6 +445,7 @@ def _start_stall_watchdog(
|
|||||||
stall_timeout_seconds: int,
|
stall_timeout_seconds: int,
|
||||||
lock_key: str | None = None,
|
lock_key: str | None = None,
|
||||||
lock_owner: str | None = None,
|
lock_owner: str | None = None,
|
||||||
|
db_idle_restart_seconds: int | None = None,
|
||||||
) -> tuple[Event, Thread]:
|
) -> tuple[Event, Thread]:
|
||||||
stop_event = Event()
|
stop_event = Event()
|
||||||
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
||||||
@@ -419,6 +457,7 @@ def _start_stall_watchdog(
|
|||||||
watchdog_born = time.monotonic()
|
watchdog_born = time.monotonic()
|
||||||
absolute_deadline = stall_timeout_seconds * 3
|
absolute_deadline = stall_timeout_seconds * 3
|
||||||
while not stop_event.wait(interval_seconds):
|
while not stop_event.wait(interval_seconds):
|
||||||
|
db_idle_restart = False
|
||||||
try:
|
try:
|
||||||
raw = redis_client.get(key)
|
raw = redis_client.get(key)
|
||||||
if not raw:
|
if not raw:
|
||||||
@@ -444,18 +483,32 @@ def _start_stall_watchdog(
|
|||||||
last_ts = int(data.get("ts") or 0)
|
last_ts = int(data.get("ts") or 0)
|
||||||
if not last_ts:
|
if not last_ts:
|
||||||
continue
|
continue
|
||||||
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
db_idle_restart = bool(
|
||||||
age = int(time.time()) - last_ts
|
db_idle_restart_seconds
|
||||||
if age < effective_stall_timeout:
|
and _should_restart_for_db_idle(data, db_idle_restart_seconds)
|
||||||
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
|
||||||
continue
|
|
||||||
logger.error(
|
|
||||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
|
||||||
task_id,
|
|
||||||
age,
|
|
||||||
stage,
|
|
||||||
data,
|
|
||||||
)
|
)
|
||||||
|
if db_idle_restart:
|
||||||
|
logger.error(
|
||||||
|
"Task %s has no DB writes for >%ss at segment=%s/%s stage=%s; full restart required",
|
||||||
|
task_id,
|
||||||
|
db_idle_restart_seconds,
|
||||||
|
data.get("segment_index"),
|
||||||
|
data.get("segments_total"),
|
||||||
|
stage,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||||
|
age = int(time.time()) - last_ts
|
||||||
|
if age < effective_stall_timeout:
|
||||||
|
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||||
|
continue
|
||||||
|
logger.error(
|
||||||
|
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||||
|
task_id,
|
||||||
|
age,
|
||||||
|
stage,
|
||||||
|
data,
|
||||||
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||||
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
||||||
@@ -464,6 +517,12 @@ def _start_stall_watchdog(
|
|||||||
else:
|
else:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
if db_idle_restart:
|
||||||
|
_restart_bootstrap_from_first_segment(
|
||||||
|
redis_client,
|
||||||
|
reason=f"no DB writes for >{db_idle_restart_seconds}s",
|
||||||
|
)
|
||||||
|
|
||||||
# ── Pre-SIGTERM cleanup: release lock so next task can run ──
|
# ── Pre-SIGTERM cleanup: release lock so next task can run ──
|
||||||
if lock_key and lock_owner:
|
if lock_key and lock_owner:
|
||||||
try:
|
try:
|
||||||
@@ -516,6 +575,55 @@ def _start_stall_watchdog(
|
|||||||
return stop_event, thread
|
return stop_event, thread
|
||||||
|
|
||||||
|
|
||||||
|
def _should_restart_for_db_idle(progress: dict, db_idle_restart_seconds: int) -> bool:
|
||||||
|
stage = str(progress.get("stage") or "")
|
||||||
|
if stage in TERMINAL_PROGRESS_STAGES:
|
||||||
|
return False
|
||||||
|
# Listing/segment progress means the task is alive even if it writes no new
|
||||||
|
# cars. This is normal for hourly only_new runs when all vehicles are already
|
||||||
|
# present in DB. Do not kill healthy scans just because DB progress is idle.
|
||||||
|
progress_ts = _safe_int(progress.get("ts")) or 0
|
||||||
|
if progress_ts > 0 and int(time.time()) - progress_ts < int(db_idle_restart_seconds):
|
||||||
|
return False
|
||||||
|
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
|
||||||
|
timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds)
|
||||||
|
return progress_ts > 0 and int(time.time()) - progress_ts >= timeout
|
||||||
|
|
||||||
|
segments_total = _safe_int(progress.get("segments_total"))
|
||||||
|
segment_index = _safe_int(progress.get("segment_index"))
|
||||||
|
if segments_total is None or segment_index is None:
|
||||||
|
return False
|
||||||
|
if segments_total <= 0 or segment_index >= segments_total - 1:
|
||||||
|
return False
|
||||||
|
|
||||||
|
now_ts = int(time.time())
|
||||||
|
if progress_ts <= 0:
|
||||||
|
return False
|
||||||
|
|
||||||
|
db_progress_ts = _safe_int(progress.get("last_db_progress_ts"))
|
||||||
|
if db_progress_ts is None:
|
||||||
|
db_progress_ts = _read_global_db_progress_ts()
|
||||||
|
task_started_ts = _safe_int(progress.get("task_started_ts")) or progress_ts
|
||||||
|
last_db_or_start_ts = max(db_progress_ts or 0, task_started_ts)
|
||||||
|
return now_ts - last_db_or_start_ts >= int(db_idle_restart_seconds)
|
||||||
|
|
||||||
|
|
||||||
|
def _safe_int(value) -> int | None:
|
||||||
|
try:
|
||||||
|
return int(value)
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _read_global_db_progress_ts() -> int | None:
|
||||||
|
try:
|
||||||
|
redis_client = _get_redis()
|
||||||
|
raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||||
|
return _safe_int(raw)
|
||||||
|
except Exception:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
_persistence_instance: PersistenceService | None = None
|
_persistence_instance: PersistenceService | None = None
|
||||||
|
|
||||||
|
|
||||||
@@ -723,6 +831,21 @@ def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> Non
|
|||||||
logger.warning("Failed to save sync listing checkpoint", exc_info=True)
|
logger.warning("Failed to save sync listing checkpoint", exc_info=True)
|
||||||
|
|
||||||
|
|
||||||
|
def _restart_bootstrap_from_first_segment(redis_client: Redis, *, reason: str) -> None:
|
||||||
|
"""Clear bootstrap checkpoint so the next full scan starts from segment 1."""
|
||||||
|
try:
|
||||||
|
pipe = redis_client.pipeline()
|
||||||
|
pipe.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
||||||
|
pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
|
||||||
|
pipe.delete(GLOBAL_PROGRESS_TS_KEY)
|
||||||
|
pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||||
|
pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0")
|
||||||
|
pipe.execute()
|
||||||
|
logger.error("Bootstrap restart requested from segment 1: %s", reason)
|
||||||
|
except Exception:
|
||||||
|
logger.warning("Failed to reset bootstrap checkpoint for DB-idle restart", exc_info=True)
|
||||||
|
|
||||||
|
|
||||||
def _clear_sync_checkpoint(redis_client: Redis) -> None:
|
def _clear_sync_checkpoint(redis_client: Redis) -> None:
|
||||||
try:
|
try:
|
||||||
redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
||||||
@@ -752,7 +875,9 @@ def _mark_sync_completed(redis_client: Redis) -> None:
|
|||||||
logger.warning("Failed to mark sync completion timestamp", exc_info=True)
|
logger.warning("Failed to mark sync completion timestamp", exc_info=True)
|
||||||
|
|
||||||
|
|
||||||
def _seconds_until_next_allowed_sync(redis_client: Redis, settings: Settings) -> int:
|
def _seconds_until_next_allowed_sync(redis_client: Redis, settings: Settings, *, is_beat_task: bool = False) -> int:
|
||||||
|
if is_beat_task:
|
||||||
|
return 0
|
||||||
min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60))
|
min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60))
|
||||||
if min_interval <= 0:
|
if min_interval <= 0:
|
||||||
return 0
|
return 0
|
||||||
@@ -878,6 +1003,7 @@ def _start_lock_heartbeat(
|
|||||||
key: str,
|
key: str,
|
||||||
owner_token: str,
|
owner_token: str,
|
||||||
ttl_seconds: int,
|
ttl_seconds: int,
|
||||||
|
stop_on_lost: bool = True,
|
||||||
) -> tuple[Event, Thread]:
|
) -> tuple[Event, Thread]:
|
||||||
stop_event = Event()
|
stop_event = Event()
|
||||||
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
||||||
@@ -888,7 +1014,10 @@ def _start_lock_heartbeat(
|
|||||||
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
||||||
if refreshed is False:
|
if refreshed is False:
|
||||||
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
||||||
return
|
if stop_on_lost:
|
||||||
|
return
|
||||||
|
consecutive_failures += 1
|
||||||
|
continue
|
||||||
if refreshed is None:
|
if refreshed is None:
|
||||||
consecutive_failures += 1
|
consecutive_failures += 1
|
||||||
if consecutive_failures >= 5:
|
if consecutive_failures >= 5:
|
||||||
@@ -929,6 +1058,49 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in
|
|||||||
return 0, 0
|
return 0, 0
|
||||||
|
|
||||||
|
|
||||||
|
def _build_listing_segments(settings: Settings) -> list[dict[str, str | int | None]]:
|
||||||
|
"""Возвращает сегменты листинга с учётом runtime_config.
|
||||||
|
|
||||||
|
При IAAI_LISTING_SEGMENTS=runtime сегменты строятся из filters.brands / filters.include.brands.
|
||||||
|
Остальные значения IAAI_LISTING_SEGMENTS сохраняют прежнее поведение: auto или JSON.
|
||||||
|
"""
|
||||||
|
filtered_urls = settings.listing.filtered_search_urls
|
||||||
|
if len(filtered_urls) > 1:
|
||||||
|
return [
|
||||||
|
{"make": None, "year_min": None, "year_max": None, "listing_url": url}
|
||||||
|
for url in filtered_urls
|
||||||
|
]
|
||||||
|
if len(filtered_urls) == 1:
|
||||||
|
return []
|
||||||
|
|
||||||
|
raw_segments = settings.listing.listing_segments_json.strip()
|
||||||
|
if raw_segments.casefold() != "runtime":
|
||||||
|
return parse_listing_segments(raw_segments)
|
||||||
|
|
||||||
|
parsed_override = parse_listing_segments(raw_segments)
|
||||||
|
if parsed_override:
|
||||||
|
return parsed_override
|
||||||
|
|
||||||
|
runtime_config = RuntimeConfig.from_file(settings.runtime_config_file)
|
||||||
|
brands = runtime_config.filters.include.brands
|
||||||
|
|
||||||
|
# Fast HTTP-first режим не использует UI year-фильтры IAAI.
|
||||||
|
# Значит runtime-сегменты должны быть "один бренд = один сегмент"
|
||||||
|
# без доп. year split, иначе получаем много лишних долгих сегментов,
|
||||||
|
# которые fast-профиль всё равно игнорирует.
|
||||||
|
if settings.scraping_profile.http_first:
|
||||||
|
segments = build_fast_listing_segments_for_makes(list(brands))
|
||||||
|
elif settings.listing.fast_segment_year_splits:
|
||||||
|
segments = build_listing_segments_for_makes(list(brands))
|
||||||
|
else:
|
||||||
|
segments = build_fast_listing_segments_for_makes(list(brands))
|
||||||
|
if not segments:
|
||||||
|
logger.warning(
|
||||||
|
"IAAI_LISTING_SEGMENTS=runtime, but runtime_config filters.brands is empty; segmented listing disabled"
|
||||||
|
)
|
||||||
|
return segments
|
||||||
|
|
||||||
|
|
||||||
@shared_task(
|
@shared_task(
|
||||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||||
queue=IAAI_SYNC_QUEUE,
|
queue=IAAI_SYNC_QUEUE,
|
||||||
@@ -964,7 +1136,10 @@ def sync_segment_task(
|
|||||||
seg_make = segment.get("make")
|
seg_make = segment.get("make")
|
||||||
seg_year_min = segment.get("year_min")
|
seg_year_min = segment.get("year_min")
|
||||||
seg_year_max = segment.get("year_max")
|
seg_year_max = segment.get("year_max")
|
||||||
|
seg_listing_url = segment.get("listing_url")
|
||||||
seg_label = f"{seg_make or 'ALL'}"
|
seg_label = f"{seg_make or 'ALL'}"
|
||||||
|
if seg_listing_url:
|
||||||
|
seg_label = "FILTERED_SEARCH"
|
||||||
if seg_year_min is not None or seg_year_max is not None:
|
if seg_year_min is not None or seg_year_max is not None:
|
||||||
seg_label += f" ({seg_year_min}-{seg_year_max})"
|
seg_label += f" ({seg_year_min}-{seg_year_max})"
|
||||||
|
|
||||||
@@ -1006,14 +1181,12 @@ def sync_segment_task(
|
|||||||
**meta,
|
**meta,
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
base_url = scraper.settings.listing.cars_url
|
|
||||||
seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None
|
|
||||||
return scraper.sync_listing(
|
return scraper.sync_listing(
|
||||||
make=None if seg_url else seg_make,
|
make=seg_make,
|
||||||
model=None,
|
model=None,
|
||||||
lane=lane,
|
lane=lane,
|
||||||
only_new=only_new,
|
only_new=only_new,
|
||||||
listing_url=seg_url,
|
listing_url=str(seg_listing_url) if seg_listing_url else None,
|
||||||
year_min=seg_year_min,
|
year_min=seg_year_min,
|
||||||
year_max=seg_year_max,
|
year_max=seg_year_max,
|
||||||
skip_mark_sold=True,
|
skip_mark_sold=True,
|
||||||
@@ -1170,6 +1343,7 @@ def sync_listing_task(
|
|||||||
0,
|
0,
|
||||||
int(os.getenv("IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS", str(settings.celery.beat_sync_interval_minutes * 60))),
|
int(os.getenv("IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS", str(settings.celery.beat_sync_interval_minutes * 60))),
|
||||||
)
|
)
|
||||||
|
is_beat_task = bool(self.request.headers and self.request.headers.get("iaai_beat_task"))
|
||||||
|
|
||||||
def _enqueue_bootstrap_followup(
|
def _enqueue_bootstrap_followup(
|
||||||
reason: str,
|
reason: str,
|
||||||
@@ -1252,7 +1426,7 @@ def sync_listing_task(
|
|||||||
}
|
}
|
||||||
|
|
||||||
try:
|
try:
|
||||||
next_allowed_delay = _seconds_until_next_allowed_sync(redis_client, settings)
|
next_allowed_delay = _seconds_until_next_allowed_sync(redis_client, settings, is_beat_task=is_beat_task)
|
||||||
if next_allowed_delay > 0:
|
if next_allowed_delay > 0:
|
||||||
logger.info(
|
logger.info(
|
||||||
"sync_listing_task skipped: previous full run finished recently; next run allowed in %ss",
|
"sync_listing_task skipped: previous full run finished recently; next run allowed in %ss",
|
||||||
@@ -1277,16 +1451,10 @@ def sync_listing_task(
|
|||||||
SYNC_LISTING_LOCK_KEY,
|
SYNC_LISTING_LOCK_KEY,
|
||||||
owner_token,
|
owner_token,
|
||||||
lock_ttl,
|
lock_ttl,
|
||||||
|
stop_on_lost=False,
|
||||||
)
|
)
|
||||||
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
|
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
|
||||||
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
|
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
|
||||||
watchdog_stop, watchdog_thread = _start_stall_watchdog(
|
|
||||||
redis_client,
|
|
||||||
task_id=task_id,
|
|
||||||
stall_timeout_seconds=stall_timeout,
|
|
||||||
lock_key=SYNC_LISTING_LOCK_KEY,
|
|
||||||
lock_owner=owner_token,
|
|
||||||
)
|
|
||||||
_update_task_progress(
|
_update_task_progress(
|
||||||
redis_client,
|
redis_client,
|
||||||
task_id=task_id,
|
task_id=task_id,
|
||||||
@@ -1296,7 +1464,17 @@ def sync_listing_task(
|
|||||||
|
|
||||||
full_scan_done_before_run = _is_full_scan_done(redis_client)
|
full_scan_done_before_run = _is_full_scan_done(redis_client)
|
||||||
always_full_scan = bool(settings.discovery.always_full_scan)
|
always_full_scan = bool(settings.discovery.always_full_scan)
|
||||||
force_bootstrap_full_scan = always_full_scan or (not full_scan_done_before_run)
|
explicit_filtered_run = bool(
|
||||||
|
make is not None
|
||||||
|
or model is not None
|
||||||
|
or (limit is not None and int(limit or 0) > 0)
|
||||||
|
or only_new is True
|
||||||
|
)
|
||||||
|
force_bootstrap_full_scan = (always_full_scan or (not full_scan_done_before_run)) and not explicit_filtered_run
|
||||||
|
if explicit_filtered_run and not full_scan_done_before_run:
|
||||||
|
logger.info(
|
||||||
|
"Explicit sync_listing request detected; honoring make/model/limit/only_new before bootstrap full scan is complete",
|
||||||
|
)
|
||||||
effective_limit = None if force_bootstrap_full_scan else limit
|
effective_limit = None if force_bootstrap_full_scan else limit
|
||||||
effective_only_new = False if force_bootstrap_full_scan else only_new
|
effective_only_new = False if force_bootstrap_full_scan else only_new
|
||||||
hourly_mode = settings.discovery.hourly_mode.strip().lower()
|
hourly_mode = settings.discovery.hourly_mode.strip().lower()
|
||||||
@@ -1313,13 +1491,31 @@ def sync_listing_task(
|
|||||||
and not effective_only_new
|
and not effective_only_new
|
||||||
and not always_full_scan
|
and not always_full_scan
|
||||||
)
|
)
|
||||||
use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline
|
# Определяем сегменты из конфига/env или из runtime_config при IAAI_LISTING_SEGMENTS=runtime.
|
||||||
|
filtered_listing_urls = settings.listing.filtered_search_urls
|
||||||
|
filtered_listing_url = filtered_listing_urls[0] if filtered_listing_urls else None
|
||||||
|
segments = _build_listing_segments(settings)
|
||||||
|
|
||||||
|
watchdog_stop, watchdog_thread = _start_stall_watchdog(
|
||||||
|
redis_client,
|
||||||
|
task_id=task_id,
|
||||||
|
stall_timeout_seconds=stall_timeout,
|
||||||
|
lock_key=SYNC_LISTING_LOCK_KEY,
|
||||||
|
lock_owner=owner_token,
|
||||||
|
db_idle_restart_seconds=(DB_IDLE_RESTART_SECONDS if segments else None),
|
||||||
|
)
|
||||||
|
use_hourly_sitemap_sync = (
|
||||||
|
(not always_full_scan)
|
||||||
|
and full_scan_done_before_run
|
||||||
|
and prefer_sitemap_mainline
|
||||||
|
and not segments
|
||||||
|
)
|
||||||
|
|
||||||
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||||
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
||||||
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
|
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
|
||||||
last_completed_segment: int | None = None
|
last_completed_segment: int | None = None
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run):
|
||||||
last_completed_segment = _load_last_completed_segment(redis_client)
|
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||||
else:
|
else:
|
||||||
# После завершения bootstrap чекпоинт не нужен никогда.
|
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||||
@@ -1461,14 +1657,11 @@ def sync_listing_task(
|
|||||||
|
|
||||||
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
||||||
|
|
||||||
# Определяем сегменты из конфига.
|
|
||||||
segments = parse_listing_segments(settings.listing.listing_segments_json)
|
|
||||||
use_segmented = (
|
use_segmented = (
|
||||||
bool(segments)
|
bool(segments)
|
||||||
and make is None
|
and make is None
|
||||||
and model is None
|
and model is None
|
||||||
and not prefer_sitemap_mainline
|
and not use_hourly_sitemap_sync
|
||||||
and discovery_mode != "sitemap"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
if prefer_sitemap_mainline:
|
if prefer_sitemap_mainline:
|
||||||
@@ -1586,7 +1779,7 @@ def sync_listing_task(
|
|||||||
start_page=1,
|
start_page=1,
|
||||||
progress_callback=(
|
progress_callback=(
|
||||||
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
|
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
|
||||||
if force_bootstrap_full_scan else None
|
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run) else None
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
return scraper.sync_listing(
|
return scraper.sync_listing(
|
||||||
@@ -1595,6 +1788,7 @@ def sync_listing_task(
|
|||||||
lane=lane,
|
lane=lane,
|
||||||
limit=effective_limit,
|
limit=effective_limit,
|
||||||
only_new=effective_only_new,
|
only_new=effective_only_new,
|
||||||
|
listing_url=filtered_listing_url,
|
||||||
)
|
)
|
||||||
|
|
||||||
result = _run_browser_job(_job)
|
result = _run_browser_job(_job)
|
||||||
|
|||||||
179
vps_monitor.sh
179
vps_monitor.sh
@@ -1,179 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
set -euo pipefail
|
|
||||||
|
|
||||||
# VPS monitor for iaai-parser. Run on the VPS from any directory:
|
|
||||||
# bash /root/iaai-parser/vps_monitor.sh
|
|
||||||
# Live mode:
|
|
||||||
# watch -n 30 bash /root/iaai-parser/vps_monitor.sh
|
|
||||||
|
|
||||||
APP_DIR="${APP_DIR:-/root/iaai-parser}"
|
|
||||||
cd "$APP_DIR"
|
|
||||||
|
|
||||||
LINE='============================================================'
|
|
||||||
SMALL='------------------------------------------------------------'
|
|
||||||
|
|
||||||
section() {
|
|
||||||
echo
|
|
||||||
echo "$LINE"
|
|
||||||
echo "$1"
|
|
||||||
echo "$LINE"
|
|
||||||
}
|
|
||||||
|
|
||||||
subsection() {
|
|
||||||
echo
|
|
||||||
echo "$SMALL"
|
|
||||||
echo "$1"
|
|
||||||
echo "$SMALL"
|
|
||||||
}
|
|
||||||
|
|
||||||
run_psql() {
|
|
||||||
docker compose exec -T postgres psql -U iaai -d iaai_scraper "$@"
|
|
||||||
}
|
|
||||||
|
|
||||||
run_redis() {
|
|
||||||
docker compose exec -T redis redis-cli "$@"
|
|
||||||
}
|
|
||||||
|
|
||||||
now_utc="$(date -u '+%Y-%m-%d %H:%M:%S UTC')"
|
|
||||||
echo "IAAI VPS monitor | $now_utc | app=$APP_DIR"
|
|
||||||
|
|
||||||
section "1) Services / restarts"
|
|
||||||
docker compose ps
|
|
||||||
|
|
||||||
subsection "Container restart counters"
|
|
||||||
for svc in api worker beat postgres redis; do
|
|
||||||
cid="$(docker compose ps -q "$svc" 2>/dev/null || true)"
|
|
||||||
if [[ -z "$cid" ]]; then
|
|
||||||
printf '%-10s %s\n' "$svc" "not found"
|
|
||||||
continue
|
|
||||||
fi
|
|
||||||
name="$(docker inspect --format '{{.Name}}' "$cid" | sed 's#^/##')"
|
|
||||||
restarts="$(docker inspect --format '{{.RestartCount}}' "$cid")"
|
|
||||||
started="$(docker inspect --format '{{.State.StartedAt}}' "$cid")"
|
|
||||||
status="$(docker inspect --format '{{.State.Status}}' "$cid")"
|
|
||||||
health="$(docker inspect --format '{{if .State.Health}}{{.State.Health.Status}}{{else}}n/a{{end}}' "$cid")"
|
|
||||||
printf '%-22s restarts=%-3s started=%s status=%s health=%s\n' "$name" "$restarts" "$started" "$status" "$health"
|
|
||||||
done
|
|
||||||
|
|
||||||
subsection "Startup / restart log markers"
|
|
||||||
docker compose logs --since '24h' worker beat 2>/dev/null \
|
|
||||||
| grep -E 'Worker ready|startup sync|dispatching initial|Scheduler: Sending due task|warm shutdown|Restoring|mingle:|ready\.|ERROR|CRITICAL' \
|
|
||||||
| tail -80 || true
|
|
||||||
|
|
||||||
section "2) Main DB result"
|
|
||||||
run_psql -c "
|
|
||||||
select
|
|
||||||
count(*) as total_cars,
|
|
||||||
count(*) filter (where is_sold = false) as active_cars,
|
|
||||||
count(*) filter (where is_sold = true) as sold_cars,
|
|
||||||
count(*) filter (where last_seen_at > now() - interval '5 minutes') as touched_5m,
|
|
||||||
count(*) filter (where last_seen_at > now() - interval '15 minutes') as touched_15m,
|
|
||||||
count(*) filter (where last_seen_at > now() - interval '60 minutes') as touched_60m,
|
|
||||||
max(last_seen_at) as last_seen
|
|
||||||
from iaai_cars;
|
|
||||||
"
|
|
||||||
|
|
||||||
subsection "Writes from sync_runs"
|
|
||||||
run_psql -c "
|
|
||||||
select
|
|
||||||
coalesce(sum(cars_upserted) filter (where started_at > now() - interval '1 hour'), 0) as cars_upserted_1h,
|
|
||||||
coalesce(sum(cars_upserted) filter (where started_at > now() - interval '6 hours'), 0) as cars_upserted_6h,
|
|
||||||
coalesce(sum(cars_upserted) filter (where started_at::date = now()::date), 0) as cars_upserted_today,
|
|
||||||
count(*) filter (where status = 'running') as running_runs,
|
|
||||||
count(*) filter (where status = 'success' and started_at > now() - interval '24 hours') as success_runs_24h,
|
|
||||||
count(*) filter (where status = 'failed' and started_at > now() - interval '24 hours') as failed_runs_24h
|
|
||||||
from iaai_sync_runs;
|
|
||||||
"
|
|
||||||
|
|
||||||
section "3) Segments / passes"
|
|
||||||
subsection "Redis checkpoint / active progress"
|
|
||||||
checkpoint="$(run_redis GET iaai:state:sync_listing_checkpoint || true)"
|
|
||||||
segments_done="$(run_redis SCARD iaai:state:sync_segments_done 2>/dev/null || echo 0)"
|
|
||||||
segments_total="$(run_redis GET iaai:state:sync_segments_total || true)"
|
|
||||||
full_scan_done="$(run_redis GET iaai:state:sync_full_scan_done || true)"
|
|
||||||
queue_len="$(run_redis LLEN iaai_sync || true)"
|
|
||||||
lock_value="$(run_redis GET iaai:locks:sync_listing || true)"
|
|
||||||
echo "checkpoint_last_completed_segment=${checkpoint:-empty}"
|
|
||||||
echo "segments_done_set_count=${segments_done:-0}"
|
|
||||||
echo "segments_total=${segments_total:-unknown}"
|
|
||||||
echo "full_scan_done=${full_scan_done:-empty}"
|
|
||||||
echo "queue_len=${queue_len:-unknown}"
|
|
||||||
echo "sync_lock=${lock_value:-empty}"
|
|
||||||
|
|
||||||
echo
|
|
||||||
echo "Active/latest task progress keys:"
|
|
||||||
progress_keys="$(run_redis --scan --pattern 'iaai:state:task_progress:*' || true)"
|
|
||||||
if [[ -z "$progress_keys" ]]; then
|
|
||||||
echo "no task_progress keys"
|
|
||||||
else
|
|
||||||
while IFS= read -r key; do
|
|
||||||
[[ -z "$key" ]] && continue
|
|
||||||
echo "$key"
|
|
||||||
run_redis GET "$key" || true
|
|
||||||
echo
|
|
||||||
done <<< "$progress_keys"
|
|
||||||
fi
|
|
||||||
|
|
||||||
subsection "Completed full passes by worker logs"
|
|
||||||
passes_24h="$(docker compose logs --since '24h' worker 2>/dev/null | grep -Ec 'Segment 65/65 done|Segment 64/65 done: RIMAC|Bootstrap full scan completed|sync_listing_task completed: status=success' || true)"
|
|
||||||
passes_6h="$(docker compose logs --since '6h' worker 2>/dev/null | grep -Ec 'Segment 65/65 done|Segment 64/65 done: RIMAC|Bootstrap full scan completed|sync_listing_task completed: status=success' || true)"
|
|
||||||
echo "estimated_completed_passes_6h=$passes_6h"
|
|
||||||
echo "estimated_completed_passes_24h=$passes_24h"
|
|
||||||
|
|
||||||
subsection "Recent segment events"
|
|
||||||
docker compose logs --tail=220 worker 2>/dev/null \
|
|
||||||
| grep -E 'Segment [0-9]+/[0-9]+|Fast HTTP-first listing collected|Fast HTTP-first DB batch|sync_listing_task completed|Bootstrap full scan completed|Bootstrap continuation stopped|sync_already_running' \
|
|
||||||
| tail -90 || true
|
|
||||||
|
|
||||||
section "4) Errors / warnings"
|
|
||||||
subsection "Failed sync_runs last 24h"
|
|
||||||
run_psql -c "
|
|
||||||
select id, status, cars_upserted, cars_failed, started_at, finished_at, left(coalesce(error_summary, ''), 180) as error_summary
|
|
||||||
from iaai_sync_runs
|
|
||||||
where status <> 'success' and started_at > now() - interval '24 hours'
|
|
||||||
order by id desc
|
|
||||||
limit 20;
|
|
||||||
"
|
|
||||||
|
|
||||||
subsection "Recent ERROR/WARNING logs"
|
|
||||||
docker compose logs --since '6h' worker beat api 2>/dev/null \
|
|
||||||
| grep -E 'ERROR|WARNING|CRITICAL|Traceback|TimeLimitExceeded|SoftTimeLimitExceeded|Read timed out|anti-bot|circuit breaker' \
|
|
||||||
| tail -120 || true
|
|
||||||
|
|
||||||
section "5) Schedule / hourly behavior"
|
|
||||||
subsection "Relevant environment"
|
|
||||||
docker compose exec -T worker sh -lc 'env | grep -E "IAAI_STARTUP_SYNC_ENABLED|CELERY_BEAT_SYNC_INTERVAL_MINUTES|IAAI_LISTING_SEGMENTS|IAAI_ALWAYS_FULL_SCAN|IAAI_FULL_SCAN_ON_STARTUP|IAAI_HOURLY_MODE|IAAI_HOURLY_REFRESH_BATCH_SIZE|IAAI_SELF_HEAL_STALL_SECONDS|IAAI_DB_IDLE_RESTART_SECONDS|CELERY_TASK_SOFT_TIME_LIMIT|CELERY_TASK_TIME_LIMIT|CELERY_BROKER_VISIBILITY_TIMEOUT" | sort' || true
|
|
||||||
|
|
||||||
subsection "Beat sends / skipped duplicates"
|
|
||||||
docker compose logs --since '12h' beat worker 2>/dev/null \
|
|
||||||
| grep -E 'Scheduler: Sending due task periodic-sync-listing|sync_already_running|sync_listing_task hourly|Hourly mode|Worker ready|startup sync dispatch disabled|dispatching initial' \
|
|
||||||
| tail -120 || true
|
|
||||||
|
|
||||||
section "6) Sold cars"
|
|
||||||
run_psql -c "
|
|
||||||
select
|
|
||||||
count(*) filter (where is_sold = true) as sold_total,
|
|
||||||
count(*) filter (where is_sold = true and last_seen_at > now() - interval '24 hours') as sold_seen_24h,
|
|
||||||
min(last_seen_at) filter (where is_sold = true) as oldest_sold_seen,
|
|
||||||
max(last_seen_at) filter (where is_sold = true) as newest_sold_seen
|
|
||||||
from iaai_cars;
|
|
||||||
"
|
|
||||||
|
|
||||||
subsection "Top sold brands"
|
|
||||||
run_psql -c "
|
|
||||||
select brand, count(*) as sold
|
|
||||||
from iaai_cars
|
|
||||||
where is_sold = true
|
|
||||||
group by brand
|
|
||||||
order by sold desc
|
|
||||||
limit 15;
|
|
||||||
"
|
|
||||||
|
|
||||||
subsection "Sold marking logs"
|
|
||||||
docker compose logs --since '24h' worker 2>/dev/null \
|
|
||||||
| grep -E 'Marked [0-9]+ cars as sold|sold_marked|Fast sold reconcile failed|mark_sold safety abort' \
|
|
||||||
| tail -80 || true
|
|
||||||
|
|
||||||
section "7) Quick health hints"
|
|
||||||
echo "OK if: worker/beat healthy, queue_len is 0 or small, one sync_lock exists during active run, checkpoint grows, failed_runs_24h is low."
|
|
||||||
echo "If IAAI_STARTUP_SYNC_ENABLED is empty/true, worker restart can dispatch an immediate sync. Set it to false to rely only on hourly beat."
|
|
||||||
Reference in New Issue
Block a user