2 Commits

Author SHA1 Message Date
qananasikq
d9c95b1a9f Add tools 2026-04-27 20:50:08 +03:00
qananasikq
00d33d527c Update parser 2026-04-27 20:50:07 +03:00
14 changed files with 399 additions and 21 deletions

4
.gitignore vendored
View File

@@ -19,4 +19,6 @@ build/
artifacts/ artifacts/
celerybeat-schedule* celerybeat-schedule*
tokens_data/ tokens_data/
tokens.json tokens.json
iaai_deploy.tar.gz
.tmp_iaai_fast_ref/

37
check_vps.py Normal file
View File

@@ -0,0 +1,37 @@
from __future__ import annotations
import os
import paramiko
HOST = "2.26.123.84"
USER = "root"
PASSWORD = os.environ["VPS_PASSWORD"]
cmds = [
"cd /root/iaai-parser && docker compose ps",
"docker logs --since 3m iaai-parser-worker-1 2>&1 | tail -n 120",
"docker exec -i iaai-postgres psql -U iaai -d iaai_scraper -c \"SELECT now() AS ts, count(*) AS cars FROM iaai_cars; SELECT count(*) AS runs FROM iaai_sync_runs;\"",
]
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(
HOST,
username=USER,
password=PASSWORD,
look_for_keys=False,
allow_agent=False,
timeout=30,
banner_timeout=30,
auth_timeout=30,
)
try:
for cmd in cmds:
print(f"REMOTE_RUN {cmd}", flush=True)
stdin, stdout, stderr = client.exec_command(cmd, timeout=120)
code = stdout.channel.recv_exit_status()
print(stdout.read().decode("utf-8", "replace"), end="")
print(stderr.read().decode("utf-8", "replace"), end="")
print(f"EXIT {code}", flush=True)
finally:
client.close()

41
clean_vps_db.py Normal file
View File

@@ -0,0 +1,41 @@
from __future__ import annotations
import os
import paramiko
HOST = "2.26.123.84"
USER = "root"
PASSWORD = os.environ["VPS_PASSWORD"]
commands = [
"cd /root/iaai-parser && docker exec -i iaai-postgres psql -U iaai -d iaai_scraper -c 'TRUNCATE TABLE iaai_images, iaai_cars, iaai_sync_runs RESTART IDENTITY CASCADE;'",
"docker exec -i iaai-redis redis-cli FLUSHALL",
"docker exec -i iaai-postgres psql -U iaai -d iaai_scraper -c 'SELECT count(*) AS cars FROM iaai_cars; SELECT count(*) AS runs FROM iaai_sync_runs;'",
"cd /root/iaai-parser && docker compose ps",
"docker logs --since 2m iaai-parser-worker-1 2>&1 | tail -n 100",
]
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(
HOST,
username=USER,
password=PASSWORD,
look_for_keys=False,
allow_agent=False,
timeout=30,
banner_timeout=30,
auth_timeout=30,
)
try:
for command in commands:
print(f"REMOTE_RUN {command}", flush=True)
stdin, stdout, stderr = client.exec_command(command, timeout=180)
exit_code = stdout.channel.recv_exit_status()
print(stdout.read().decode("utf-8", "replace"), end="")
print(stderr.read().decode("utf-8", "replace"), end="")
print(f"EXIT {exit_code}", flush=True)
if exit_code != 0:
raise SystemExit(exit_code)
finally:
client.close()

135
deploy_vps.py Normal file
View File

@@ -0,0 +1,135 @@
from __future__ import annotations
import os
import posixpath
import stat
import tarfile
import time
import socket
from pathlib import Path
import paramiko
HOST = "2.26.123.84"
USER = "root"
PASSWORD = os.environ["VPS_PASSWORD"]
LOCAL_ROOT = Path(__file__).resolve().parent
ARCHIVE = LOCAL_ROOT / "iaai_deploy.tar.gz"
REMOTE_DIR = "/root/iaai-parser"
REMOTE_ARCHIVE = "/root/iaai_deploy.tar.gz"
EXCLUDE_DIRS = {
".git",
".venv",
"__pycache__",
".pytest_cache",
".mypy_cache",
".ruff_cache",
"iaai_scraper.egg-info",
}
EXCLUDE_FILES = {"iaai_deploy.tar.gz", "deploy_vps.py"}
def _include(path: Path) -> bool:
rel = path.relative_to(LOCAL_ROOT)
parts = set(rel.parts)
if parts & EXCLUDE_DIRS:
return False
if path.name in EXCLUDE_FILES:
return False
if path.suffix in {".pyc", ".pyo"}:
return False
return True
def make_archive() -> None:
if ARCHIVE.exists():
ARCHIVE.unlink()
with tarfile.open(ARCHIVE, "w:gz") as tf:
for path in LOCAL_ROOT.rglob("*"):
if not _include(path):
continue
tf.add(path, arcname=str(path.relative_to(LOCAL_ROOT)))
print(f"ARCHIVE_READY {ARCHIVE} {ARCHIVE.stat().st_size} bytes", flush=True)
def connect() -> paramiko.SSHClient:
last_exc: BaseException | None = None
for attempt in range(1, 6):
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
try:
print(f"SSH_CONNECT attempt={attempt}", flush=True)
client.connect(
HOST,
username=USER,
password=PASSWORD,
look_for_keys=False,
allow_agent=False,
timeout=45,
banner_timeout=60,
auth_timeout=45,
)
return client
except (paramiko.SSHException, socket.timeout, OSError) as exc:
last_exc = exc
client.close()
print(f"SSH_CONNECT_RETRY attempt={attempt} error={type(exc).__name__}: {exc}", flush=True)
time.sleep(3 * attempt)
raise SystemExit(f"SSH_CONNECT_FAILED: {last_exc}")
def run(client: paramiko.SSHClient, cmd: str, timeout: int | None = None) -> None:
print(f"REMOTE_RUN {cmd}", flush=True)
stdin, stdout, stderr = client.exec_command(cmd, timeout=timeout)
exit_code = stdout.channel.recv_exit_status()
out = stdout.read().decode("utf-8", "replace")
err = stderr.read().decode("utf-8", "replace")
if out:
print(out, end="", flush=True)
if err:
print(err, end="", flush=True)
if exit_code != 0:
raise SystemExit(f"REMOTE_FAILED code={exit_code}: {cmd}")
def upload_file(sftp: paramiko.SFTPClient, local: Path, remote: str) -> None:
total = local.stat().st_size
last = 0.0
def cb(done: int, _total: int) -> None:
nonlocal last
now = time.time()
if now - last >= 2 or done == total:
print(f"UPLOAD {done}/{total}", flush=True)
last = now
sftp.put(str(local), remote, callback=cb)
def main() -> None:
make_archive()
client = connect()
try:
run(client, "echo VPS_OK && hostname && docker --version && docker compose version")
sftp = client.open_sftp()
try:
upload_file(sftp, ARCHIVE, REMOTE_ARCHIVE)
finally:
sftp.close()
run(
client,
f"mkdir -p {REMOTE_DIR} && cd {REMOTE_DIR} && docker compose down || true && "
f"find {REMOTE_DIR} -mindepth 1 -maxdepth 1 ! -name '.env' -exec rm -rf {{}} + && "
f"tar -xzf {REMOTE_ARCHIVE} -C {REMOTE_DIR} && rm -f {REMOTE_ARCHIVE}",
timeout=300,
)
run(client, f"cd {REMOTE_DIR} && docker compose up -d --build", timeout=1800)
run(client, f"cd {REMOTE_DIR} && docker compose ps", timeout=120)
run(client, "docker logs --since 2m iaai-parser-worker-1 2>&1 | tail -n 120 || docker compose -f /root/iaai-parser/docker-compose.yml logs --since 2m worker | tail -n 120", timeout=120)
finally:
client.close()
if __name__ == "__main__":
main()

View File

@@ -7,38 +7,53 @@ x-app-env: &app-env
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_POOL_SIZE: ${IAAI_DATABASE_POOL_SIZE:-20}
IAAI_DATABASE_MAX_OVERFLOW: ${IAAI_DATABASE_MAX_OVERFLOW:-40} IAAI_DATABASE_MAX_OVERFLOW: ${IAAI_DATABASE_MAX_OVERFLOW:-40}
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900} # Для больших Search-выборок (десятки тысяч лотов) run должен жить дольше одного батча.
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200} CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-7200}
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-2400} CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-7500}
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-10800}
IAAI_PROFILE: ${IAAI_PROFILE:-fast} IAAI_PROFILE: ${IAAI_PROFILE:-fast}
IAAI_SCRAPING_PROFILE: ${IAAI_SCRAPING_PROFILE:-${IAAI_PROFILE:-fast}} IAAI_SCRAPING_PROFILE: ${IAAI_SCRAPING_PROFILE:-${IAAI_PROFILE:-fast}}
IAAI_HTTP_FIRST: ${IAAI_HTTP_FIRST:-true} IAAI_HTTP_FIRST: ${IAAI_HTTP_FIRST:-true}
IAAI_BROWSER_FALLBACK_ENABLED: ${IAAI_BROWSER_FALLBACK_ENABLED:-false} IAAI_BROWSER_FALLBACK_ENABLED: ${IAAI_BROWSER_FALLBACK_ENABLED:-false}
IAAI_ANONYMOUS_BOOTSTRAP_ENABLED: ${IAAI_ANONYMOUS_BOOTSTRAP_ENABLED:-false} IAAI_ANONYMOUS_BOOTSTRAP_ENABLED: ${IAAI_ANONYMOUS_BOOTSTRAP_ENABLED:-false}
IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1} IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1}
IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-2} IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-1}
IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0} IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0.03}
IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-1} IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-2}
IAAI_LISTING_RETRIES: ${IAAI_LISTING_RETRIES:-1} 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.
CELERY_BEAT_SYNC_INTERVAL_MINUTES: ${CELERY_BEAT_SYNC_INTERVAL_MINUTES:-60}
# Любой внутренний follow-up после неполного bootstrap тоже не стартует сразу.
IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS: ${IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS:-3600}
CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4} CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4}
CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork} 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:-1000} CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1500}
IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0} IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0}
# Стабильный fast-профиль: не душим IAAI слишком большим числом HTTPS detail-соединений.
IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-96}
IAAI_MAX_RETRIES: ${IAAI_MAX_RETRIES:-2}
IAAI_RETRY_DELAY_SECONDS: ${IAAI_RETRY_DELAY_SECONDS:-0.35}
CELERY_PARALLEL_SEGMENTS: ${CELERY_PARALLEL_SEGMENTS:-false}
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_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-32}
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_URL: ${IAAI_FILTERED_SEARCH_URL:-}
IAAI_FILTERED_SEARCH_URLS: ${IAAI_FILTERED_SEARCH_URLS:-} IAAI_FILTERED_SEARCH_URLS: ${IAAI_FILTERED_SEARCH_URLS:-}
IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-runtime} IAAI_FAST_PATH_TIMEOUT_MS: ${IAAI_FAST_PATH_TIMEOUT_MS:-10000}
IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-[]}
IAAI_DISCOVERY_MODE: ${IAAI_DISCOVERY_MODE:-listing}
IAAI_HOURLY_MODE: ${IAAI_HOURLY_MODE:-rolling_refresh}
IAAI_HOURLY_REFRESH_BATCH_SIZE: ${IAAI_HOURLY_REFRESH_BATCH_SIZE:-500}
IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999} IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999}
IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000} IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000}
IAAI_ALWAYS_FULL_SCAN: ${IAAI_ALWAYS_FULL_SCAN:-true} IAAI_ALWAYS_FULL_SCAN: ${IAAI_ALWAYS_FULL_SCAN:-false}
IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-true} IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-false}
IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/home/app/tokens.json} # Первый sync после clean restart стартует сразу; следующий запуск — только через hourly guard/beat.
IAAI_STARTUP_SYNC_ENABLED: ${IAAI_STARTUP_SYNC_ENABLED:-true}
IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/data/tokens.json}
IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json}
IAAI_BROWSER_ENGINE: chromium IAAI_BROWSER_ENGINE: chromium
# Self-heal для worker. # Self-heal для worker.
@@ -168,7 +183,9 @@ services:
stop_grace_period: 60s stop_grace_period: 60s
command: > command: >
celery -A iaai_scraper.worker.celery_app worker celery -A iaai_scraper.worker.celery_app worker
--loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=${CELERY_WORKER_POOL:-prefork} --loglevel=info
--concurrency=${CELERY_WORKER_CONCURRENCY:-4}
--pool=${CELERY_WORKER_POOL:-prefork}
--pidfile=/tmp/celery-worker.pid --pidfile=/tmp/celery-worker.pid
-Q iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} -Q iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
healthcheck: healthcheck:

View File

@@ -519,7 +519,17 @@ class IAAIFastClient:
if max_pages is not None and max_pages > 0: if max_pages is not None and max_pages > 0:
total_pages = min(total_pages, max_pages) total_pages = min(total_pages, max_pages)
gbp_search_query = first_page.gbp_search_query gbp_search_query = first_page.gbp_search_query
logger.info(
"Fast listing first page parsed: scope=%s result_count=%s page_size=%s total_pages=%s first_page_vehicles=%s",
scope_path,
first_page.result_count,
page_size,
total_pages,
len(first_page.vehicles),
)
for page_number in range(2, total_pages + 1): for page_number in range(2, total_pages + 1):
if page_number == 2 or page_number % 25 == 0 or page_number == total_pages:
logger.info("Fast listing page fetch progress: page=%s/%s", page_number, total_pages)
page_html = self._fetch_listing_page(scope_path, gbp_search_query, page_number, page_size) page_html = self._fetch_listing_page(scope_path, gbp_search_query, page_number, page_size)
parsed_page = parse_listing_page(page_html) parsed_page = parse_listing_page(page_html)
gbp_search_query = parsed_page.gbp_search_query gbp_search_query = parsed_page.gbp_search_query

View File

@@ -18,8 +18,9 @@ def set_trace_id(trace_id: str) -> None:
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
# stderr — Docker и Celery prefork корректно его подхватывают. # stdout — основной поток для `docker logs`, Docker Desktop и compose logs.
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)] # stderr в PowerShell часто выглядит как NativeCommandError, хотя это обычные логи.
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)]
if log_file: if log_file:
handlers.append(logging.FileHandler(log_file, encoding="utf-8")) handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
trace_filter = TraceIdFilter() trace_filter = TraceIdFilter()

View File

@@ -183,6 +183,7 @@ class FastSyncEngine:
prepared_rows: list[CarRecord] = [] prepared_rows: list[CarRecord] = []
db_processed = 0 db_processed = 0
retry_candidates: list[FastListingVehicle] = [] retry_candidates: list[FastListingVehicle] = []
transient_failure_log_count = 0
started_details = time.perf_counter() started_details = time.perf_counter()
with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor: with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor:
future_to_vehicle = { future_to_vehicle = {
@@ -232,7 +233,13 @@ class FastSyncEngine:
stats.protection_events += 1 stats.protection_events += 1
if self._is_transient_detail_error(exc): if self._is_transient_detail_error(exc):
retry_candidates.append(vehicle) retry_candidates.append(vehicle)
logger.info("Fast detail transient failure queued for retry inventory_id=%s: %s", vehicle.inventory_id, exc) transient_failure_log_count += 1
if transient_failure_log_count % 100 == 0:
logger.warning(
"Fast detail transient failures queued for retry: count=%d latest_inventory_id=%s",
transient_failure_log_count,
vehicle.inventory_id,
)
else: else:
logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc) logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc)

View File

@@ -2234,6 +2234,8 @@ class IAAIScraper:
only_new = rc.only_new only_new = rc.only_new
if lane == "iaai_cars" and rc.lane is not None: if lane == "iaai_cars" and rc.lane is not None:
lane = rc.lane lane = rc.lane
if listing_url is None and rc.name:
listing_url = rc.name
trace_id = self._new_trace_id("sync-listing") trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter() started_at = time.perf_counter()
@@ -2260,6 +2262,8 @@ class IAAIScraper:
year_min, year_min,
year_max, year_max,
) )
if listing_url:
logger.info("sync_listing uses explicit listing_url from runtime/args: %s", listing_url)
fast_engine = FastSyncEngine( fast_engine = FastSyncEngine(
client=self.fast_client, client=self.fast_client,
mapper=self.fast_mapper, mapper=self.fast_mapper,

View File

@@ -16,6 +16,9 @@ logger = logging.getLogger("iaai_scraper.worker.celery_app")
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched" STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
IAAI_SYNC_QUEUE = "iaai_sync" IAAI_SYNC_QUEUE = "iaai_sync"
PROGRESS_KEY_PREFIX = "iaai:state:task_progress:" 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: 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 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 @celery_setup_logging.connect
def _configure_logging(loglevel=None, **kwargs): def _configure_logging(loglevel=None, **kwargs):
# Перехватываем логирование Celery и пишем только в stderr (Docker logs). # Перехватываем логирование Celery и пишем только в stderr (Docker logs).
@@ -116,6 +139,7 @@ celery_app.conf.update(
"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},
}, },
} }
}, },
@@ -147,7 +171,15 @@ def _on_worker_ready(**kwargs):
) )
has_fresh_progress = _has_fresh_active_progress(redis_client) 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: try:
ttl = redis_client.ttl(stale_key) ttl = redis_client.ttl(stale_key)
if ttl is not None and ttl != -2 and not has_fresh_progress: if ttl is not None and ttl != -2 and not has_fresh_progress:

View File

@@ -183,6 +183,7 @@ class FastSyncEngine:
prepared_rows: list[CarRecord] = [] prepared_rows: list[CarRecord] = []
db_processed = 0 db_processed = 0
retry_candidates: list[FastListingVehicle] = [] retry_candidates: list[FastListingVehicle] = []
transient_failure_log_count = 0
started_details = time.perf_counter() started_details = time.perf_counter()
with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor: with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor:
future_to_vehicle = { future_to_vehicle = {
@@ -232,7 +233,13 @@ class FastSyncEngine:
stats.protection_events += 1 stats.protection_events += 1
if self._is_transient_detail_error(exc): if self._is_transient_detail_error(exc):
retry_candidates.append(vehicle) retry_candidates.append(vehicle)
logger.info("Fast detail transient failure queued for retry inventory_id=%s: %s", vehicle.inventory_id, exc) transient_failure_log_count += 1
if transient_failure_log_count % 100 == 0:
logger.warning(
"Fast detail transient failures queued for retry: count=%d latest_inventory_id=%s",
transient_failure_log_count,
vehicle.inventory_id,
)
else: else:
logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc) logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc)

View File

@@ -48,6 +48,7 @@ HOURLY_FAILURE_STREAK_LIMIT = 3
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed" SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed"
SYNC_LAST_COMPLETED_AT_KEY = "iaai:state:sync_listing_last_completed_at"
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}" SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress" SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total" SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
@@ -867,6 +868,33 @@ def _clear_followup_pending(redis_client: Redis) -> None:
logger.warning("Failed to clear follow-up pending flag", exc_info=True) logger.warning("Failed to clear follow-up pending flag", exc_info=True)
def _mark_sync_completed(redis_client: Redis) -> None:
try:
redis_client.set(SYNC_LAST_COMPLETED_AT_KEY, str(int(time.time())), ex=7 * 24 * 60 * 60)
except Exception:
logger.warning("Failed to mark sync completion timestamp", exc_info=True)
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))
if min_interval <= 0:
return 0
try:
if not _is_full_scan_done(redis_client):
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 next allowed sync time", exc_info=True)
return 0
elapsed = int(time.time()) - completed_at
return max(0, min_interval - elapsed)
def _bump_bootstrap_failure_streak( def _bump_bootstrap_failure_streak(
redis_client: Redis, redis_client: Redis,
*, *,
@@ -1310,6 +1338,12 @@ def sync_listing_task(
watchdog_stop: Event | None = None watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None watchdog_thread: Thread | None = None
force_bootstrap_full_scan = False force_bootstrap_full_scan = False
settings = Settings()
followup_min_delay_seconds = max(
0,
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,
@@ -1317,6 +1351,7 @@ def sync_listing_task(
*, *,
count_as_failure: bool = False, count_as_failure: bool = False,
) -> None: ) -> None:
delay_seconds = max(int(delay_seconds), followup_min_delay_seconds)
flag_ttl = max(lock_ttl, delay_seconds + 300) flag_ttl = max(lock_ttl, delay_seconds + 300)
if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl): if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl):
logger.info( logger.info(
@@ -1391,7 +1426,18 @@ def sync_listing_task(
} }
try: try:
settings = Settings() next_allowed_delay = _seconds_until_next_allowed_sync(redis_client, settings, is_beat_task=is_beat_task)
if next_allowed_delay > 0:
logger.info(
"sync_listing_task skipped: previous full run finished recently; next run allowed in %ss",
next_allowed_delay,
)
return {
"status": "skipped",
"reason": "next_sync_not_due_yet",
"retry_after_seconds": next_allowed_delay,
"task_id": task_id,
}
# Любой реально стартовавший sync_listing снимает pending-флаг followup, # Любой реально стартовавший sync_listing снимает pending-флаг followup,
# чтобы watchdog/continuation могли корректно планировать следующий run # чтобы watchdog/continuation могли корректно планировать следующий run
@@ -1832,6 +1878,7 @@ def sync_listing_task(
summary["cars_failed"], summary["cars_failed"],
summary["failures_count"], summary["failures_count"],
) )
_mark_sync_completed(redis_client)
return summary return summary
except SoftTimeLimitExceeded: except SoftTimeLimitExceeded:

View File

@@ -1,6 +1,6 @@
{ {
"sync": { "sync": {
"name": null, "name": "https://www.iaai.com/Search?url=B%2bU066dM8%2flZtRvnzTwIEi8Pib9%2fM2fLpCTQDIgQRm4%3d",
"ids_initial_size": null, "ids_initial_size": null,
"ids_next_size": null, "ids_next_size": null,
"ids_max_pages": null, "ids_max_pages": null,

38
status_vps.py Normal file
View File

@@ -0,0 +1,38 @@
from __future__ import annotations
import os
import paramiko
HOST = "2.26.123.84"
USER = "root"
PASSWORD = os.environ["VPS_PASSWORD"]
commands = [
"cd /root/iaai-parser && docker compose ps",
"cd /root/iaai-parser && docker compose exec -T worker celery -A iaai_scraper.worker.celery_app inspect active reserved scheduled",
"docker exec -i iaai-postgres psql -U iaai -d iaai_scraper -c \"SELECT now() AS ts, (SELECT count(*) FROM iaai_cars) AS cars, (SELECT count(*) FROM iaai_sync_runs) AS runs, (SELECT status FROM iaai_sync_runs ORDER BY id DESC LIMIT 1) AS last_status, (SELECT cars_upserted FROM iaai_sync_runs ORDER BY id DESC LIMIT 1) AS last_upserted, (SELECT cars_failed FROM iaai_sync_runs ORDER BY id DESC LIMIT 1) AS last_failed, (SELECT started_at FROM iaai_sync_runs ORDER BY id DESC LIMIT 1) AS last_started;\"",
"docker logs --since 10m iaai-parser-worker-1 2>&1 | tail -n 160",
]
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(
HOST,
username=USER,
password=PASSWORD,
look_for_keys=False,
allow_agent=False,
timeout=30,
banner_timeout=30,
auth_timeout=30,
)
try:
for command in commands:
print(f"\n=== REMOTE_RUN {command} ===", flush=True)
stdin, stdout, stderr = client.exec_command(command, timeout=180)
exit_code = stdout.channel.recv_exit_status()
print(stdout.read().decode("utf-8", "replace"), end="")
print(stderr.read().decode("utf-8", "replace"), end="")
print(f"\n=== EXIT {exit_code} ===", flush=True)
finally:
client.close()