diff --git a/.forgejo/workflows/deploy-metrics.yml b/.forgejo/workflows/deploy-metrics.yml index 7b101edc..559784aa 100644 --- a/.forgejo/workflows/deploy-metrics.yml +++ b/.forgejo/workflows/deploy-metrics.yml @@ -419,11 +419,22 @@ jobs: docker compose -p gendesign-metrics-agent \ -f docker-compose.metrics-agent.yml up -d + # Конфиг Alloy — бинд-маунт ОДНОГО файла, а `git reset --hard` выше пишет + # его новым инодом: `up -d` изменения не видит, контейнер держит старый. + METRICS_ROLE=apps \ + METRICS_ALLOY_CONFIG=alloy-apps.alloy \ + COMPOSE_PROFILES="$EXPORTER_PROFILE" \ + docker compose -p gendesign-metrics-agent \ + -f docker-compose.metrics-agent.yml up -d --force-recreate alloy + sleep 10 METRICS_ROLE=apps METRICS_ALLOY_CONFIG=alloy-apps.alloy COMPOSE_PROFILES="$EXPORTER_PROFILE" \ docker compose -p gendesign-metrics-agent \ -f docker-compose.metrics-agent.yml ps + [ "$(stat -c %i ops/metrics/alloy/alloy-apps.alloy)" = "$(docker exec gendesign-alloy stat -c %i /etc/alloy/config.alloy)" ] \ + || { echo "::error::alloy читает старый инод конфига"; exit 1; } + agent-infra: runs-on: ubuntu-latest needs: server @@ -478,7 +489,18 @@ jobs: docker compose -p gendesign-metrics-agent \ -f docker-compose.metrics-agent.yml up -d + # Конфиг Alloy — бинд-маунт ОДНОГО файла, а `git reset --hard` (джоба server) + # пишет его новым инодом: `up -d` изменения не видит, контейнер держит старый. + METRICS_ROLE=infra \ + METRICS_ALLOY_CONFIG=alloy-infra.alloy \ + COMPOSE_PROFILES="$EXPORTER_PROFILE" \ + docker compose -p gendesign-metrics-agent \ + -f docker-compose.metrics-agent.yml up -d --force-recreate alloy + sleep 10 METRICS_ROLE=infra METRICS_ALLOY_CONFIG=alloy-infra.alloy COMPOSE_PROFILES="$EXPORTER_PROFILE" \ docker compose -p gendesign-metrics-agent \ -f docker-compose.metrics-agent.yml ps + + [ "$(stat -c %i ops/metrics/alloy/alloy-infra.alloy)" = "$(docker exec gendesign-alloy stat -c %i /etc/alloy/config.alloy)" ] \ + || { echo "::error::alloy читает старый инод конфига"; exit 1; } diff --git a/.forgejo/workflows/deploy.yml b/.forgejo/workflows/deploy.yml index 8b8edc7a..9b2c1866 100644 --- a/.forgejo/workflows/deploy.yml +++ b/.forgejo/workflows/deploy.yml @@ -596,6 +596,11 @@ jobs: # #3029: подлинность хоста. Секрет НЕ задан → пустая строка → easyssh-proxy # оставляет ssh.InsecureIgnoreHostKey(), то есть сегодняшнее поведение. fingerprint: ${{ secrets.DEPLOY_SSH_FINGERPRINT }} + # #3324: дефолт appleboy/ssh-action — command_timeout 10m, а worst-case + # гейта в конце скрипта ~17 мин (4 сервиса × 240s ожидания healthy + + # фронт + diagnose). Сессию убило бы посреди печати диагноза, и авария + # выглядела бы обрывом связи, а не мёртвым контейнером. + command_timeout: 30m envs: IMAGE_TAG,SENTRY_RELEASE_VAL,GHCR_PAT,GLITCHTIP_BACKEND_DSN,OBJECTIVE_API_KEY,OPENAI_API_KEY,LLM_ENABLED,OWN_DEVELOPER_IDS,WORKER_RECREATE_GUARD,WORKER_GUARD_MAX_SKIP_H script: | set -euo pipefail @@ -1037,6 +1042,171 @@ jobs: fi echo "→ backend healthy на /health." + # ── Гейт деплоя (#3324) ──────────────────────────────────────────── + # ДО этого блока весь гейт ПТИЦЫ = один `curl backend /health` выше: + # worker/beat не проверялись вообще (у них и healthcheck'а в compose не + # было), у frontend была только TCP-проба внутри контейнера, которую + # деплой не читал. То есть crash-loop воркера, вставший beat и фронт, + # отдающий 500, уезжали ЗЕЛЁНЫМ деплоем. Дисциплина перенесена из + # deploy-tradein.yml (health каждого сервиса + HTTP фронта + сверка + # образов), сюда добавлено чтение docker-health, потому что у ПТИЦЫ + # пробы теперь описаны в compose. + # `up -d --wait` НЕ используется намеренно: подъём здесь разбит на + # несколько `up` (bulk без worker'а → guard #3029 → caddy → forwarder), + # и общий --wait ждал бы ещё и профильные/инфраструктурные сервисы, + # ломая порядок «миграции до подъёма кода». Читаем состояние явно. + # Все проверки выполняются ДО выхода (не падаем на первой) — один + # прогон обязан показать ВСЕ поломанные сервисы, а не первый по списку. + # tail -n1: у сервиса может остаться залежавшийся exited-контейнер, и + # тогда `ps -aq` вернёт НЕСКОЛЬКО id через \n — `docker inspect` с таким + # аргументом падает, статус приходит пустым и гейт краснеет на ровном + # месте. Последний id — самый свежий контейнер сервиса. + cid() { docker compose -p gendesign -f docker-compose.prod.yml ps -aq "$1" 2>/dev/null | tail -n1 || true; } + + diagnose() { # $1 сервис, $2 id контейнера (может быть пустым) + local svc="$1" c="$2" + echo "── ДИАГНОЗ $svc ──" + if [ -z "$c" ]; then + echo " контейнера нет вообще (docker compose ps -aq $svc пусто)" + return 0 + fi + docker inspect -f ' state={{.State.Status}} health={{if .State.Health}}{{.State.Health.Status}}{{else}}{{end}} restarts={{.RestartCount}} exit_code={{.State.ExitCode}} image={{.Image}}' "$c" || true + echo " healthcheck log:" + docker inspect -f '{{json .State.Health}}' "$c" 2>/dev/null | head -c 2000 || true + echo "" + echo " последние 40 строк логов $svc:" + docker logs --tail 40 "$c" 2>&1 | sed 's/^/ /' || true + } + + wait_healthy() { # $1 сервис, $2 таймаут, с + local svc="$1" deadline="$2" c="" status="" alive="" r0="" r1="" waited=0 + while [ "$waited" -lt "$deadline" ]; do + c="$(cid "$svc")" + if [ -n "$c" ]; then + status="$(docker inspect -f '{{if .State.Health}}{{.State.Health.Status}}{{else}}none:{{.State.Status}}{{end}}' "$c" 2>/dev/null || echo '')" + case "$status" in + healthy) + echo "→ $svc healthy (за ${waited}s)" + return 0 + ;; + none:running) + # Контейнер без healthcheck-конфига = создан ДО этой правки + # compose и в этом прогоне не пересоздавался (штатный случай — + # worker, пропущенный guard'ом #3029). Валить деплой за это + # нельзя, но и молчать нельзя: падаем на «стабильный running» + # (двойное чтение, как tgbot/scraper в deploy-tradein.yml). + # Одного `running` дважды НЕДОСТАТОЧНО: crash-loop с временем + # жизни больше паузы читается как «стабилен» — контейнер оба + # раза running, просто это разные его жизни. Поэтому вместе со + # статусом сверяем RestartCount: изменился за окно = именно + # тот дефект, ради которого этот гейт и писался. + r0="$(docker inspect -f '{{.RestartCount}}' "$c" 2>/dev/null || echo '')" + sleep 15 + alive="$(docker inspect -f '{{.State.Status}}' "$c" 2>/dev/null || echo unknown)" + r1="$(docker inspect -f '{{.RestartCount}}' "$c" 2>/dev/null || echo '')" + if [ "$alive" = "running" ] && [ -n "$r0" ] && [ "$r0" = "$r1" ]; then + echo "→ $svc: healthcheck не сконфигурирован (контейнер не пересоздавался), running стабилен (RestartCount=$r0 не изменился за 15s)" + return 0 + fi + # waited растёт на длину ЭТОЙ паузы тоже — иначе таймаут + # 240s превратился бы в ~24 минуты реального ожидания и упёрся + # бы в command_timeout SSH-сессии. + waited=$((waited + 15)) + echo " $svc: running нестабилен — status='$alive', RestartCount ${r0:-<нет>}→${r1:-<нет>} (контейнер перезапускался внутри окна наблюдения); продолжаю ждать" + ;; + esac + fi + waited=$((waited + 3)) + sleep 3 + done + echo "ERROR (#3324): $svc не стал healthy за ${deadline}s (последний статус: '${status:-<контейнера нет>}') — деплой FAILED" + diagnose "$svc" "$c" + return 1 + } + + health_rc=0 + for gate_svc in backend worker beat frontend; do + wait_healthy "$gate_svc" 240 || health_rc=$? + done + + # Фронт: HTTP-СТАТУС, а не только «порт слушает». Compose-проба фронта + # намеренно TCP-only (в node:alpine нет ни curl, ни wget), и она не + # отличает живой Next.js от процесса, отдающего 500 на каждый запрос. + # Тянем с хоста через опубликованный 127.0.0.1:3000. basePath у ПТИЦЫ + # нет (frontend/next.config.*), корень — настоящий маршрут приложения. + # Годным считаем 2xx/3xx: редирект middleware'а на логин — это живой + # роутинг, а не поломка (тот же критерий, что `curl -f` в tradein). + fe_rc=0 + fe_code=000 + for i in $(seq 1 5); do + fe_code="$(curl -s -o /dev/null -w '%{http_code}' --max-time 10 http://localhost:3000/ || echo 000)" + case "$fe_code" in + 2*|3*) break ;; + esac + sleep 3 + done + case "$fe_code" in + 2*|3*) echo "→ frontend отвечает HTTP $fe_code на /." ;; + *) + echo "ERROR (#3324): frontend на http://localhost:3000/ вернул '$fe_code' (000 = соединения нет) — деплой FAILED" + diagnose frontend "$(cid frontend)" + fe_rc=1 + ;; + esac + + # Сверка образов (приём #2679 из deploy-tradein.yml, адаптирован под + # ПТИЦУ). Здесь не одно «backend-семейство»: backend и beat бегут один + # образ gendesign-backend, worker и frontend — свои. Поэтому эталон не + # «образ backend'а», а то, что реально лежит локально под тегом + # $IMAGE_TAG после pull'а: контейнер, оставшийся на другом id, работает + # на старом коде при зелёном деплое. + # Гард свежести самого :latest в registry — отдельный шаг выше + # (scripts/check-latest-image-revision.sh, #2950); здесь проверяется + # следующее звено: доехал ли уже скачанный образ до контейнера. + check_image() { # $1 сервис, $2 репозиторий образа + local svc="$1" repo="$2" want run c + want="$(docker image inspect -f '{{.Id}}' "$repo:$IMAGE_TAG" 2>/dev/null || echo '')" + c="$(cid "$svc")" + run="$(docker inspect -f '{{.Image}}' "$c" 2>/dev/null || echo '')" + if [ -z "$want" ]; then + echo "ERROR (#3324): локально нет образа $repo:$IMAGE_TAG — сверять не с чем (pull не отработал?)" + return 1 + fi + if [ -z "$run" ]; then + echo "ERROR (#3324): контейнера сервиса $svc НЕТ — это не «отставший образ», а неполный стек" + return 1 + fi + if [ "$want" != "$run" ]; then + echo "ERROR (#3324): $svc ОТСТАЛ: работает на $run, а $repo:$IMAGE_TAG — это $want" + echo " лечение: docker compose -p gendesign -f docker-compose.prod.yml up -d --force-recreate --no-deps $svc" + diagnose "$svc" "$c" + return 1 + fi + echo "→ $svc на свежем $repo:$IMAGE_TAG ($run)" + } + + image_rc=0 + check_image backend ghcr.io/lekss361/gendesign-backend || image_rc=$? + check_image beat ghcr.io/lekss361/gendesign-backend || image_rc=$? + check_image frontend ghcr.io/lekss361/gendesign-frontend || image_rc=$? + # worker сверяем ТОЛЬКО если этот прогон его пересоздавал: guard #3029 + # намеренно оставляет worker на старом образе, пока идёт живой прогон + # скрейпа, и это уже отражено WARNING'ом выше. Падать здесь означало бы + # красить деплой за штатное поведение guard'а. + case " $WORKER_SERVICES " in + *" worker "*) check_image worker ghcr.io/lekss361/gendesign-worker || image_rc=$? ;; + *) echo "→ сверка образа worker'а пропущена: guard #3029 не пересоздавал его в этом прогоне (см. WARNING выше)" ;; + esac + + # Явный rc: «зелёная сводка» ниже печатается ДО выхода, поэтому итог + # обязан быть числом в логе, а не выводом из отсутствия ERROR-строк. + gate_rc=0 + [ "$health_rc" = 0 ] || gate_rc=1 + [ "$fe_rc" = 0 ] || gate_rc=1 + [ "$image_rc" = 0 ] || gate_rc=1 + echo "Гейт деплоя (#3324): health_rc=$health_rc frontend_rc=$fe_rc image_rc=$image_rc → rc=$gate_rc" + exit "$gate_rc" + # Честный итог прогона (#2841). ПРОБЛЕМА: `deploy` пропускается своим `if:` # молча (result=skipped), когда build падает (например, битый blob в # buildcache роняет `docker/build-push-action` — до ретрая выше, #2841). diff --git a/_wt-mskcol/tradein-mvp/scripts/local-avito-msk/collect.py b/_wt-mskcol/tradein-mvp/scripts/local-avito-msk/collect.py new file mode 100644 index 00000000..5b0f2461 --- /dev/null +++ b/_wt-mskcol/tradein-mvp/scripts/local-avito-msk/collect.py @@ -0,0 +1,785 @@ +#!/usr/bin/env python3 +"""Локальный ручной сборщик SERP Авито по Москве и МО (эпик #2989, трек 1). + +Запускается ВРУЧНУЮ с машины владельца. Прод-скрейпер, его расписания и +прокси-пул не задействованы вообще: браузер — уже открытый Chrome владельца +(подключение по CDP), парсер — импорт из scraper-kit, заливка — поток в psql +через ssh. Скрипт ничего не устанавливает и своего профиля не поднимает. + +Дефолтный режим — --measure 100 (замер): полный проход только по явному --full. +""" + +from __future__ import annotations + +import argparse +import asyncio +import csv +import io +import json +import math +import os +import random +import re +import subprocess +import sys +import time +from dataclasses import dataclass, field +from datetime import datetime, timezone +from pathlib import Path +from types import SimpleNamespace +from typing import Any, Iterable +from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit + +# --- импорт парсера из scraper-kit без установки backend ------------------- +_KIT_SRC = Path(__file__).resolve().parents[2] / "packages" / "scraper-kit" / "src" +if str(_KIT_SRC) not in sys.path: + sys.path.insert(0, str(_KIT_SRC)) + +from scraper_kit.providers.avito.serp import ( # noqa: E402 + AvitoScraper, + _is_firewall_page, +) + +# Вкладка, открытая у владельца: вторичка, Москва + МО. +DEFAULT_BASE_URL = ( + "https://www.avito.ru/moskva_i_mo/kvartiry/prodam/vtorichka-ASgBAgICAkSSA8YQ5geMUg" + "?f=ASgBAgICA0SSA8YQ5geMUuDC3c0CgOkP" +) + +# Замерено живым проходом (не из документации): выдача Москва+МО отдаёт 50 карточек +# на страницу. Пока здесь стояло 60, planned_pages считал count/60 и не запрашивал +# последние ~17% каждого коридора — 34 753 по счётчику против 28 352 собранных. +# Молчаливое усечение читается как «покрыто всё», поэтому число проверяется живьём. +PAGE_SIZE = 50 # карточек на странице выдачи +MAX_PAGES = 30 # потолок пагинации Авито → 30*50 = 1500 на один запрос +HARD_CAP = PAGE_SIZE * MAX_PAGES + +PRICE_FLOOR = 500_000 # нижняя граница осмысленного коридора, ₽ +PRICE_PROBE_START = 8_000_000 # старт удвоения при поиске верхней границы +PRICE_CEIL = 2_000_000_000 +MIN_WIDTH_RATIO = 1.05 # уже этого коридор не делим (геометрическая ширина) +MAX_DEPTH = 12 + +_BATCH_ID_RE = re.compile(r"^[A-Za-z0-9._-]+$") +_POW_MARKERS = ("startpow", "доступ ограничен: проверка безопасности") + + +# Заполняется при входе в Loader.__aenter__ (playwright импортируется лениво, +# чтобы --help работал без установленного пакета). +PlaywrightTimeoutError: type[BaseException] = TimeoutError + + +# Через столько загрузок вкладка сборщика пересоздаётся (см. _recycle_if_needed). +PAGE_RECYCLE_EVERY = 75 + + +class Blocked(Exception): + """Первый признак блока. Ретраев нет — только немедленный стоп.""" + + def __init__(self, reason: str, detail: str = "") -> None: + super().__init__(f"{reason}: {detail}" if detail else reason) + self.reason = reason + self.detail = detail + + +class BudgetExhausted(Exception): + """Потолок --measure выбран: штатный выход, не ошибка.""" + + +# --- план коридоров -------------------------------------------------------- + + +@dataclass +class Corridor: + lo: int | None + hi: int | None + count: int | None = None + truncated: bool = False + pages_done: int = 0 + status: str = "pending" # pending | done + missed: int = 0 # заведомо недобрано (count - HARD_CAP), если truncated + + def label(self) -> str: + lo = "-" if self.lo is None else f"{self.lo:_}" + hi = "-" if self.hi is None else f"{self.hi:_}" + return f"[{lo} .. {hi}]" + + def to_json(self) -> dict[str, Any]: + return { + "lo": self.lo, "hi": self.hi, "count": self.count, + "truncated": self.truncated, "pages_done": self.pages_done, + "status": self.status, "missed": self.missed, + } + + @staticmethod + def from_json(d: dict[str, Any]) -> "Corridor": + return Corridor( + lo=d.get("lo"), hi=d.get("hi"), count=d.get("count"), + truncated=bool(d.get("truncated")), + pages_done=int(d.get("pages_done") or 0), + status=d.get("status") or "pending", missed=int(d.get("missed") or 0), + ) + + def planned_pages(self) -> int: + if not self.count: + return 1 + return max(1, min(MAX_PAGES, math.ceil(self.count / PAGE_SIZE))) + + +@dataclass +class Plan: + base_url: str + target: int + batch_id: str + corridors: list[Corridor] = field(default_factory=list) + created_at: str = "" + + def save(self, path: Path) -> None: + path.write_text( + json.dumps( + { + "version": 1, "base_url": self.base_url, "target": self.target, + "batch_id": self.batch_id, "created_at": self.created_at, + "corridors": [c.to_json() for c in self.corridors], + }, + ensure_ascii=False, indent=1, + ), + encoding="utf-8", + ) + + @staticmethod + def load(path: Path) -> "Plan": + d = json.loads(path.read_text(encoding="utf-8")) + return Plan( + base_url=d["base_url"], target=int(d["target"]), batch_id=d["batch_id"], + created_at=d.get("created_at", ""), + corridors=[Corridor.from_json(c) for c in d.get("corridors", [])], + ) + + +def build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str: + """URL коридора: pmin/pmax + пагинация. Гео-параметры не добавляем — #3043.""" + parts = urlsplit(base_url) + q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True) + if k not in {"p", "pmin", "pmax"}] + if lo is not None: + q.append(("pmin", str(int(lo)))) + if hi is not None: + q.append(("pmax", str(int(hi)))) + if page > 1: + q.append(("p", str(page))) + return urlunsplit( + (parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment) + ) + + +def geometric_mid(lo: int | None, hi: int) -> int: + """Геометрическая середина коридора. + + Цены логнормальны: арифметическая середина 1 млн..100 млн (≈50 млн) + отрезает вырожденно-пустую верхнюю половину. sqrt(lo*hi) делит выборку + заметно ровнее. + """ + low = max(int(lo or PRICE_FLOOR), 1) + mid = int(math.sqrt(low * float(hi))) + return max(low + 1, min(hi - 1, mid)) + + +def width_ratio(lo: int | None, hi: int | None) -> float: + if hi is None: + return float("inf") + return float(hi) / max(float(lo or PRICE_FLOOR), 1.0) + + +# --- загрузка страницы ----------------------------------------------------- + + +def _guard(html: str, status: int | None) -> None: + """Порядок проверок фиксирован заданием; первое срабатывание = стоп.""" + if status in (403, 439): + raise Blocked("platform", f"HTTP {status}") + if status == 429: + raise Blocked("ratelimit", "HTTP 429") + if _is_firewall_page(html): + raise Blocked("firewall", "firewall-страница на HTTP 200") + head = html[:4096].lower() + if any(m in head for m in _POW_MARKERS): + raise Blocked("challenge", "PoW / проверка безопасности") + + +class Loader: + """Одна СВОЯ вкладка в уже открытом Chrome владельца (CDP). + + Ни браузер, ни контекст, ни чужие вкладки не закрываются и не трогаются: + это рабочий Chrome с залогиненным техаккаунтом. + """ + + def __init__(self, delay: float, page_budget: int | None) -> None: + self._delay = delay + self._budget = page_budget + self.loads = 0 + self._page: Any = None + self._pw: Any = None + self._browser: Any = None + self._ctx: Any = None + self._last_load = 0.0 + self._loads_on_page = 0 + + async def __aenter__(self) -> "Loader": + from playwright.async_api import async_playwright + from playwright.async_api import TimeoutError as _PwTimeout + + global PlaywrightTimeoutError + PlaywrightTimeoutError = _PwTimeout + + endpoint = os.environ.get("AVITO_CDP", "http://localhost:9222") + self._pw = await async_playwright().start() + try: + self._browser = await self._pw.chromium.connect_over_cdp(endpoint) + except Exception as exc: # noqa: BLE001 — подсказка важнее типа + await self._pw.stop() + raise SystemExit( + f"Не удалось подключиться по CDP к {endpoint}: {exc}\n" + "Запусти Chrome с залогиненным техаккаунтом Авито и ключом " + "--remote-debugging-port=9222, либо укажи адрес в AVITO_CDP." + ) from exc + if not self._browser.contexts: + await self._pw.stop() + raise SystemExit( + "В подключённом Chrome нет ни одного контекста. Открой обычное окно " + "Chrome, запущенное с --remote-debugging-port=9222." + ) + self._ctx = self._browser.contexts[0] + self._page = await self._ctx.new_page() + return self + + async def __aexit__(self, *exc: object) -> None: + if self._page is not None: + try: + await self._page.close() # ТОЛЬКО своя вкладка + except Exception: # noqa: BLE001 + pass + if self._pw is not None: + try: + await self._pw.stop() + except Exception: # noqa: BLE001 + pass + + def budget_left(self) -> bool: + return self._budget is None or self.loads < self._budget + + async def _pause(self) -> None: + if self._last_load == 0.0: + return + jitter = self._delay * random.uniform(-0.2, 0.2) + wait = max(0.0, self._delay + jitter - (time.monotonic() - self._last_load)) + if wait > 0: + print(f" пауза {wait:.1f} с", flush=True) + await asyncio.sleep(wait) + + async def fetch(self, url: str) -> tuple[str, int | None]: + if not self.budget_left(): + raise BudgetExhausted() + await self._recycle_if_needed() + await self._pause() + resp = await self._goto(url) + self._loads_on_page += 1 + self.loads += 1 + self._last_load = time.monotonic() + status = resp.status if resp is not None else None + try: + await self._page.wait_for_selector('[data-marker="item"]', timeout=7_000) + except Exception: # noqa: BLE001 — пустая/блочная страница разбирается ниже + pass + html = await self._read_content() + _guard(html, status) + return html, status + + async def _recycle_if_needed(self) -> None: + """Каждые PAGE_RECYCLE_EVERY загрузок пересоздаём свою вкладку. + + Рендерер Chrome накапливает память по всем навигациям вкладки, а проход + по Москве — под тысячу страниц в одной. Наблюдалось живьём: вкладка + падала с «Опаньки… Код ошибки: Out of Memory» при 34 ГБ свободных в + системе, то есть упирался именно рендерер, а не машина. Свежая вкладка + стоит одну навигацию и обнуляет счёт. + + Закрывается ТОЛЬКО своя вкладка; контекст и чужие вкладки владельца не + трогаются — это его рабочий Chrome. + """ + if self._loads_on_page < PAGE_RECYCLE_EVERY: + return + print(f" вкладка пересоздаётся после {self._loads_on_page} загрузок " + "(память рендерера)", flush=True) + old = self._page + self._page = await self._ctx.new_page() + self._loads_on_page = 0 + try: + await old.close() + except Exception: # noqa: BLE001 — старая вкладка могла уже умереть + pass + + async def _goto(self, url: str, attempts: int = 3): + """goto с ограниченным ретраем на таймаут навигации. + + Авито изредка держит соединение до упора и goto падает по timeout. Это + НЕ признак отказа: в наблюдавшемся случае вкладка показывала нормальную + выдачу, а маркеров фаервола/PoW не было. Но и молча ретраить бесконечно + нельзя — тихий отказ выглядит ровно так же. Поэтому: перед каждым + повтором пробуем прочитать то, что в документе, и прогнать через _guard, + чтобы настоящий блок остановил прогон с правильной причиной; исчерпали + попытки — жёсткий стоп с причиной nav_timeout. + """ + for i in range(attempts): + try: + return await self._page.goto(url, wait_until="domcontentloaded", + timeout=90_000) + except PlaywrightTimeoutError: + try: + partial = await self._page.content() + except Exception: # noqa: BLE001 — документа может не быть вовсе + partial = "" + if partial: + _guard(partial, None) # настоящий блок остановит прогон здесь + if i == attempts - 1: + raise Blocked("nav_timeout") from None + print(f" таймаут навигации, попытка {i + 2}/{attempts}", flush=True) + await asyncio.sleep(10.0) + + async def _read_content(self, attempts: int = 4) -> str: + """page.content() с узким ретраем на гонку клиентской перенавигации. + + Авито дорисовывает выдачу после domcontentloaded, и content() иногда + попадает ровно в момент смены документа: "Unable to retrieve content + because the page is navigating and changing the content". Это НЕ отказ + площадки — гвардов не касается, поэтому ретраим только эту ошибку и + только её, а любую другую поднимаем как есть. + """ + last: Exception | None = None + for i in range(attempts): + try: + return await self._page.content() + except Exception as exc: # noqa: BLE001 — сузили проверкой текста ниже + if "page is navigating" not in str(exc): + raise + last = exc + print(f" content() поймал перенавигацию, попытка {i + 2}/{attempts}", + flush=True) + await asyncio.sleep(1.5) + raise RuntimeError(f"page.content() не отдал документ за {attempts} попыток") from last + + +# --- заливка в msk_raw ----------------------------------------------------- + + +def _sql_str(value: str) -> str: + return "'" + value.replace("'", "''") + "'" + + +def _csv_rows(rows: Iterable[dict[str, Any]]) -> str: + buf = io.StringIO() + writer = csv.writer(buf, lineterminator="\n") + for r in rows: + writer.writerow([ + r["source_id"], r["observed_at"], r["batch_id"], r["kind"], + r["url"], r["price"], r["payload"], + ]) + return buf.getvalue() + + +def build_sql(batch_id: str, query: str, rows: list[dict[str, Any]], + started_at: str, kind: str = "serp") -> str: + """Один поток на `psql -f -`: batch (FK!) → TEMP staging → \\copy → INSERT. + + Одиночный `psql -c` через ssh ломается на квотинге скобок и кавычек, поэтому + только поток. rows_new = разница count(*) по batch_id до и после вставки. + """ + bid = _sql_str(batch_id) + return ( + "BEGIN;\n" + "INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n" + f"VALUES ({bid}, {_sql_str(kind)}, {_sql_str(query)}, " + f"CAST({_sql_str(started_at)} AS timestamptz))\n" + "ON CONFLICT (batch_id) DO NOTHING;\n" + "CREATE TEMP TABLE _stg (LIKE msk_raw.avito_cards INCLUDING DEFAULTS) " + "ON COMMIT DROP;\n" + "CREATE TEMP TABLE _before ON COMMIT DROP AS\n" + f" SELECT count(*) AS n FROM msk_raw.avito_cards WHERE batch_id = {bid};\n" + "\\copy _stg (source_id,observed_at,batch_id,kind,url,price,payload) " + "FROM STDIN WITH (FORMAT csv)\n" + + _csv_rows(rows) + + "\\.\n" + "INSERT INTO msk_raw.avito_cards " + "(source_id,observed_at,batch_id,kind,url,price,payload)\n" + "SELECT source_id,observed_at,batch_id,kind,url,price,payload FROM _stg\n" + "ON CONFLICT (source_id,batch_id,kind) DO NOTHING;\n" + "UPDATE msk_raw.batches b SET\n" + " rows_sent = coalesce(b.rows_sent,0) + (SELECT count(*) FROM _stg),\n" + " rows_new = coalesce(b.rows_new,0) +\n" + f" ((SELECT count(*) FROM msk_raw.avito_cards WHERE batch_id = {bid})\n" + " - (SELECT n FROM _before))\n" + f"WHERE b.batch_id = {bid};\n" + "COMMIT;\n" + ) + + +def build_finalize_sql(batch_id: str, query: str, notes: str) -> str: + bid = _sql_str(batch_id) + return ( + "INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n" + f"VALUES ({bid}, 'serp', {_sql_str(query)}, now())\n" + "ON CONFLICT (batch_id) DO NOTHING;\n" + f"UPDATE msk_raw.batches SET finished_at = now(), notes = {_sql_str(notes)}\n" + f"WHERE batch_id = {bid};\n" + ) + + +def run_psql(sql: str, ssh_host: str, container: str, db_user: str, db_name: str, + attempts: int = 4) -> None: + """Заливка батча через ssh с ретраем на обрыв транспорта. + + Прогон длится часами, и ssh рвётся: живьём поймано «Connection reset by peer» + (ssh возвращает 255) прямо посреди заливки — весь прогон умирал, а несброшенный + батч терялся. Ретраить безопасно: SQL идемпотентен (batch через ON CONFLICT DO + NOTHING, карточки через ON CONFLICT (source_id,batch_id,kind) DO NOTHING). + + Ретраится ТОЛЬКО транспорт (ssh 255). Ошибка самого psql (ON_ERROR_STOP, любой + другой код) — это дефект данных или SQL, её повтор не лечит: поднимаем сразу. + """ + cmd = [ + "ssh", ssh_host, + f"docker exec -i {container} psql -U {db_user} -d {db_name} " + "-v ON_ERROR_STOP=1 -f -", + ] + for i in range(attempts): + proc = subprocess.run(cmd, input=sql.encode("utf-8"), capture_output=True) + out = (proc.stdout + proc.stderr).decode("utf-8", "replace").strip() + if proc.returncode == 0: + if out: + print(f" psql: {out}", flush=True) + return + if proc.returncode != 255 or i == attempts - 1: + raise RuntimeError(f"psql через ssh вернул {proc.returncode}:\n{out}") + tail = out.splitlines()[-1] if out else "без вывода" + print(f" ssh оборвался ({tail}), повтор заливки {i + 2}/{attempts}", flush=True) + time.sleep(15.0 * (i + 1)) + + +# --- накопитель карточек --------------------------------------------------- + + +@dataclass +class Sink: + """Батчами на прод (ssh+psql) или в локальный CSV при --dry-run.""" + + batch_id: str + started_at: str + query: str + batch_size: int + dry_run: bool + csv_path: Path + ssh_host: str + container: str + db_user: str + db_name: str + buffer: list[dict[str, Any]] = field(default_factory=list) + sent: int = 0 + skipped_non_numeric: int = 0 + + def add(self, lot: Any) -> None: + raw_id = str(getattr(lot, "source_id", "") or "") + try: + source_id = int(raw_id) # в БД bigint, у ScrapedLot — строка + except (TypeError, ValueError): + self.skipped_non_numeric += 1 + return + payload = lot.model_dump(mode="json") + self.buffer.append({ + "source_id": source_id, + "observed_at": datetime.now(timezone.utc).isoformat(), + "batch_id": self.batch_id, + "kind": "serp", + "url": payload.get("source_url"), + "price": payload.get("price_rub"), + "payload": json.dumps(payload, ensure_ascii=False), + }) + + def maybe_flush(self) -> None: + if len(self.buffer) >= self.batch_size: + self.flush() + + def flush(self) -> None: + if not self.buffer: + return + rows, self.buffer = self.buffer, [] + if self.dry_run: + fresh = not self.csv_path.exists() + with self.csv_path.open("a", encoding="utf-8", newline="") as fh: + if fresh: + fh.write("source_id,observed_at,batch_id,kind,url,price,payload\n") + fh.write(_csv_rows(rows)) + print(f" [dry-run] {len(rows)} строк → {self.csv_path}", flush=True) + else: + run_psql( + build_sql(self.batch_id, self.query, rows, self.started_at), + self.ssh_host, self.container, self.db_user, self.db_name, + ) + print(f" залито {len(rows)} строк в msk_raw.avito_cards", flush=True) + self.sent += len(rows) + + def finalize(self, notes: str) -> None: + self.flush() + if self.dry_run: + print(f" [dry-run] finalize: {notes}", flush=True) + return + run_psql( + build_finalize_sql(self.batch_id, self.query, notes), + self.ssh_host, self.container, self.db_user, self.db_name, + ) + + +# --- сбор ------------------------------------------------------------------ + + +def parse_page(scraper: AvitoScraper, html: str, url: str) -> tuple[int | None, list[Any]]: + count = scraper._extract_total_count(html) + lots = scraper._parse_html(html, "https://www.avito.ru") + if not lots and count: + raise Blocked("empty_page", f"0 карточек при счётчике {count}: {url}") + return count, lots + + +async def probe(loader: Loader, scraper: AvitoScraper, base_url: str, + lo: int | None, hi: int | None) -> tuple[int | None, list[Any]]: + url = build_url(base_url, 1, lo, hi) + html, _ = await loader.fetch(url) + return parse_page(scraper, html, url) + + +async def build_plan(loader: Loader, scraper: AvitoScraper, base_url: str, target: int, + cache: dict[tuple[int | None, int | None], list[Any]] + ) -> list[Corridor]: + """Адаптивная бисекция по цене; страница 1 каждого коридора кэшируется.""" + corridors: list[Corridor] = [] + + def emit(lo: int | None, hi: int | None, count: int | None, + truncated: bool, lots: list[Any]) -> None: + missed = max(0, (count or 0) - HARD_CAP) if truncated else 0 + c = Corridor(lo=lo, hi=hi, count=count, truncated=truncated, missed=missed) + corridors.append(c) + cache[(lo, hi)] = lots + flag = " TRUNCATED" if truncated else "" + print(f" коридор {c.label()} count={count} " + f"страниц={c.planned_pages()}{flag}", flush=True) + if truncated: + print(f" ВНИМАНИЕ: коридор {c.label()} не влезает в потолок " + f"{HARD_CAP}; заведомо не добрано ~{missed} объявлений", flush=True) + + async def find_upper(lo: int | None) -> int: + """Верхнюю границу открытого коридора ищем удвоением от разумного старта.""" + cand = max(int(lo or PRICE_FLOOR) * 2, PRICE_PROBE_START) + while cand < PRICE_CEIL: + cnt, _ = await probe(loader, scraper, base_url, cand, None) + print(f" проба хвоста pmin={cand:_} count={cnt}", flush=True) + if cnt is not None and cnt <= target: + return cand + cand *= 2 + return cand + + async def split(lo: int | None, hi: int | None, depth: int, + count: int | None, lots: list[Any]) -> None: + if count is None: + url = build_url(base_url, 1, lo, hi) + raise Blocked("empty_page", f"счётчик не прочитался: {url}") + if count <= target: + emit(lo, hi, count, False, lots) + return + if depth >= MAX_DEPTH or width_ratio(lo, hi) <= MIN_WIDTH_RATIO: + # Предохранитель: не молчим — помечаем truncated и считаем недобор. + emit(lo, hi, count, count > HARD_CAP, lots) + return + upper = hi if hi is not None else await find_upper(lo) + if hi is None: + tail_cnt, tail_lots = await probe(loader, scraper, base_url, upper, None) + emit(upper, None, tail_cnt, + bool(tail_cnt and tail_cnt > HARD_CAP), tail_lots) + mid = geometric_mid(lo, upper) + for sub_lo, sub_hi in ((lo, mid), (mid, upper)): + sub_cnt, sub_lots = await probe(loader, scraper, base_url, sub_lo, sub_hi) + print(f" проба {sub_lo or '-'}..{sub_hi} count={sub_cnt}", flush=True) + await split(sub_lo, sub_hi, depth + 1, sub_cnt, sub_lots) + + root_cnt, root_lots = await probe(loader, scraper, base_url, None, None) + print(f"Всего по базовому запросу: {root_cnt}", flush=True) + await split(None, None, 0, root_cnt, root_lots) + return corridors + + +async def collect(args: argparse.Namespace) -> int: + # avito_serp_ekb_only=False обязателен: с True парсер выбрасывает всё, где в + # URL нет /ekaterinburg/ — то есть все подмосковные слаги (serp.py:2154). + scraper = AvitoScraper( + SimpleNamespace(avito_serp_ekb_only=False), # type: ignore[arg-type] + target_city_slug="moskva", + ) + out_dir = Path(args.out_dir).resolve() + out_dir.mkdir(parents=True, exist_ok=True) + plan_path = out_dir / f"plan-{args.batch_id}.json" + csv_path = out_dir / f"cards-{args.batch_id}.csv" + started_at = datetime.now(timezone.utc).isoformat() + + plan: Plan | None = None + if args.resume: + if not plan_path.exists(): + print(f"--resume: плана нет — {plan_path}", file=sys.stderr) + return 1 + plan = Plan.load(plan_path) + done = sum(1 for c in plan.corridors if c.status == "done") + print(f"Resume по {plan_path}: коридоров {len(plan.corridors)}, " + f"готово {done}", flush=True) + + # Resume: URL берём из сохранённого плана, а не из CLI — коридоры посчитаны + # именно под него. Расхождение = молчаливая заливка чужой выдачи под тем же + # batch_id, поэтому это ошибка, а не тихий приоритет одного из двух. + if plan is not None and plan.base_url != args.base_url: + raise SystemExit( + "--resume: план построен для другого URL." + f" В плане {plan.base_url}, в аргументах {args.base_url}." + " Убери --base-url (возьмётся из плана) либо начни новый batch_id." + ) + base_url = plan.base_url if plan is not None else args.base_url + + page_budget = None if args.full else args.measure + mode = "FULL" if args.full else f"MEASURE<={page_budget}" + print(f"Режим: {mode}; batch_id={args.batch_id}; delay={args.delay}s; " + f"target={args.target_count}; dry_run={args.dry_run}", flush=True) + + sink = Sink( + batch_id=args.batch_id, started_at=started_at, query=base_url, + batch_size=args.batch_size, dry_run=args.dry_run, csv_path=csv_path, + ssh_host=args.ssh_host, container=args.container, + db_user=args.db_user, db_name=args.db_name, + ) + cache: dict[tuple[int | None, int | None], list[Any]] = {} + total = 0 + stop_reason = "" + rc = 0 + loads = 0 + + async with Loader(args.delay, page_budget) as loader: + try: + if plan is None: + print("Строю план коридоров...", flush=True) + corridors = await build_plan(loader, scraper, base_url, + args.target_count, cache) + plan = Plan(base_url=base_url, target=args.target_count, + batch_id=args.batch_id, corridors=corridors, + created_at=started_at) + plan.save(plan_path) + print(f"План сохранён: {plan_path} ({len(corridors)} коридоров)", + flush=True) + + for corridor in plan.corridors: + if corridor.status == "done": + continue + pages = corridor.planned_pages() + print(f"Коридор {corridor.label()} count={corridor.count} " + f"страниц={pages} (с {corridor.pages_done + 1})", flush=True) + for page in range(corridor.pages_done + 1, pages + 1): + key = (corridor.lo, corridor.hi) + if page == 1 and key in cache: + lots = cache.pop(key) # страница 1 уже скачана при планировании + else: + url = build_url(base_url, page, corridor.lo, corridor.hi) + html, _ = await loader.fetch(url) + _, lots = parse_page(scraper, html, url) + for lot in lots: + sink.add(lot) + total += len(lots) + corridor.pages_done = page + print(f" стр.{page}/{pages}: карточек {len(lots)}, " + f"итого {total}", flush=True) + sink.maybe_flush() + plan.save(plan_path) + if not lots: + print(" пустая страница — конец коридора", flush=True) + break + corridor.status = "done" + plan.save(plan_path) + except BudgetExhausted: + stop_reason = "потолок --measure исчерпан" + print(f"Стоп: {stop_reason}", flush=True) + except Blocked as exc: + stop_reason = f"BLOCKED/{exc.reason}: {exc.detail}" + print(f"СТОП: {stop_reason}", file=sys.stderr, flush=True) + rc = 2 + finally: + loads = loader.loads + if plan is not None: + plan.save(plan_path) + + truncated = [c for c in (plan.corridors if plan else []) if c.truncated] + missed = sum(c.missed for c in truncated) + notes = "; ".join(x for x in [ + f"mode={mode}", f"loads={loads}", f"cards={total}", + f"skipped_non_numeric={sink.skipped_non_numeric}", + (f"truncated_corridors={len(truncated)} missed~{missed}" if truncated else ""), + stop_reason, + ] if x) + try: + sink.finalize(notes) + except Exception as exc: # noqa: BLE001 — не прятать исходную причину стопа + print(f"finalize провалился: {exc}", file=sys.stderr) + rc = rc or 1 + print(f"Готово. Загрузок: {loads}; карточек: {total}; отправлено: {sink.sent}; " + f"пропущено нечисловых source_id: {sink.skipped_non_numeric}; " + f"notes: {notes}", flush=True) + return rc + + +def parse_args(argv: list[str] | None = None) -> argparse.Namespace: + p = argparse.ArgumentParser( + prog="collect.py", + description="Ручной сбор SERP Авито (вторичка, Москва+МО) в прод-схему msk_raw.", + ) + p.add_argument("--base-url", default=DEFAULT_BASE_URL, + help="базовый URL выдачи (дефолт — вкладка владельца)") + p.add_argument("--measure", type=int, default=100, metavar="N", + help="режим замера: не больше N загрузок страниц (дефолт 100)") + p.add_argument("--full", action="store_true", + help="полный проход без потолка страниц (включается только явно)") + p.add_argument("--dry-run", action="store_true", + help="ничего не слать на прод, писать CSV локально") + p.add_argument("--resume", action="store_true", + help="продолжить по сохранённому плану коридоров") + p.add_argument("--delay", type=float, default=8.0, + help="пауза между загрузками, с (±20%% джиттер, дефолт 8.0)") + p.add_argument("--batch-size", type=int, default=1000, + help="карточек в одной заливке (дефолт 1000)") + p.add_argument("--target-count", type=int, default=1500, + help="целевой размер коридора; больше — делим (дефолт 1500)") + p.add_argument("--batch-id", default=None, + help="batch_id в msk_raw.batches (дефолт msk-serp-)") + p.add_argument("--out-dir", default=str(Path(__file__).resolve().parent / "runs"), + help="каталог плана/CSV") + p.add_argument("--ssh-host", default="selectel", help="ssh-хост прода") + p.add_argument("--container", default="tradein-postgres", + help="имя контейнера Postgres на проде") + p.add_argument("--db-user", default="tradein") + p.add_argument("--db-name", default="tradein") + args = p.parse_args(argv) + if args.batch_id is None: + args.batch_id = "msk-serp-" + datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + if not _BATCH_ID_RE.match(args.batch_id): + p.error("--batch-id: допустимы только символы [A-Za-z0-9._-]") + if args.measure < 1: + p.error("--measure должен быть >= 1") + return args + + +def main(argv: list[str] | None = None) -> int: + return asyncio.run(collect(parse_args(argv))) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/auth/roles.yaml b/auth/roles.yaml index 1273eb38..d2f1b028 100644 --- a/auth/roles.yaml +++ b/auth/roles.yaml @@ -133,6 +133,11 @@ users: # продукта; ранее expired с 2026-06-27). Безлимитная квота оценок # выдана через account_quota_overrides.unlimited (migration 191), # не через код — см. app.services.account_quota.is_unlimited. + buyer1: pilot # Тестовый доступ потенциального покупателя — заведён 2026-09-02 по + # просьбе владельца. Квота 50 оценок/мес через + # account_quota_overrides.monthly_limit (не unlimited). DB-роль + # manager (как praktika/kopylov — самостоятельный внешний аккаунт, + # не employee под чьим-то manager_id). admintest: admin # temp QA 2026-05-26 pilottest: pilot # temp QA 2026-05-26 analysttest: analyst # temp QA 2026-06-07 (#962) diff --git a/backend/app/services/site_finder/rosseti_wfs_loader.py b/backend/app/services/site_finder/rosseti_wfs_loader.py index 5c7d27f4..d586a4b3 100644 --- a/backend/app/services/site_finder/rosseti_wfs_loader.py +++ b/backend/app/services/site_finder/rosseti_wfs_loader.py @@ -19,6 +19,7 @@ per-row SAVEPOINT при UPSERT (битая фича не валит weekly-sync import hashlib import json import logging +import math import re import httpx @@ -101,20 +102,47 @@ def parse_voltage_class(name: str | None) -> str | None: return m.group(1).replace(".", "/") -def _stable_external_id(feature: dict, props: dict) -> str: - """Стабильный external_id фичи: feature['id'] или хэш ключевых полей. +def _coord_e5(value: float | None) -> str: + """Координата → целое в единицах 1e-5 градуса (~1 м), полукруглением от нуля. - WFS обычно отдаёт стабильный ``feature['id']``; если его нет — детерминированный - sha1 по (sc_name, координаты) чтобы UPSERT оставался идемпотентным. + Целое, а не форматированный float: ключ обязан совпадать байт-в-байт с SQL- + бэкфиллом (99c), а текстовое представление double в питоне и в PG разное. + Округление ДВОИЧНОЕ (по значению double, не по десятичному представлению): + 64.423605*1e5 == 6442360.499999999 → 6442360, хотя «по десятичному» было бы + 6442361. В 99c та же семантика: floor/abs/sign над float8, БЕЗ каста в numeric + (каст округляет по кратчайшему десятичному repr и расходится в 0.19% координат). """ - fid = feature.get("id") - if fid: - return str(fid) - geom = feature.get("geometry") or {} - coords = geom.get("coordinates") - seed = f"{props.get('sc_name', '')}|{coords}" - # sha1 здесь — стабильный дедуп-id фичи, не криптография. - return "h:" + hashlib.sha1(seed.encode("utf-8")).hexdigest()[:16] + if value is None: + return "" + n = math.floor(abs(value) * 100000 + 0.5) + return str(-n if value < 0 else n) + + +def _stable_external_id(feature: dict, props: dict) -> str: + """Стабильный external_id фичи — хэш атрибутов. ``feature['id']`` ИГНОРИРУЕТСЯ. + + GeoServer отдаёт СЕССИОННЫЙ fid (``sc_points_fullview.fid--``), новый на + каждый GetFeature → ON CONFLICT (source, external_id) не срабатывал ни разу и + таблица росла ×10 (4880 строк на 481 ЦП, #3322). Ключ считаем только по стабильным + атрибутам: нормализованное имя | класс напряжения | координаты в 1e-5 градуса. + + ФОРМУЛА ПРОДУБЛИРОВАНА в ``data/sql/99c_power_supply_centers_dedup.sql`` (бэкфилл + существующих строк) — менять только синхронно. sha256, а не sha1: sha256 встроен + в PG16, sha1 требует pgcrypto. Префикс ``h:`` отличает новый ключ от старого fid. + """ + geom_pair = _point_geom_sql(feature) + coords = geom_pair[1] if geom_pair else {} + sc_name = props.get("sc_name") + seed = "|".join( + ( + normalize_sc_name(sc_name), + parse_voltage_class(sc_name) or "", + _coord_e5(coords.get("lon")), + _coord_e5(coords.get("lat")), + ) + ) + # sha256 здесь — стабильный дедуп-id фичи, не криптография. + return "h:" + hashlib.sha256(seed.encode("utf-8")).hexdigest()[:16] def _map_load_index(props: dict) -> str | None: diff --git a/backend/tests/test_connection_capacity_loaders.py b/backend/tests/test_connection_capacity_loaders.py index ff550bae..ecc35091 100644 --- a/backend/tests/test_connection_capacity_loaders.py +++ b/backend/tests/test_connection_capacity_loaders.py @@ -18,6 +18,7 @@ Pure / mock-based — без реальной сети и БД. Покрывае from __future__ import annotations +import hashlib import io from contextlib import contextmanager from typing import Any @@ -96,6 +97,64 @@ def test_map_load_index_unknown_and_missing() -> None: assert rw._map_load_index({"sc_indexload_id": "мусор"}) is None +# ── _stable_external_id (#3322: сессионный fid раздувал таблицу ×10) ─────────── + + +def _wfs_feature(fid: str, name: str, lon: float, lat: float) -> dict[str, Any]: + return { + "id": fid, + "geometry": {"type": "Point", "coordinates": [lon, lat]}, + "properties": {"sc_name": name}, + } + + +def test_stable_external_id_ignores_session_fid() -> None: + """Разные сессионные fid + одинаковые атрибуты → ОДИН ключ (регрессия #3322).""" + a = _wfs_feature("sc_points_fullview.fid--1a2b3c", "ПС 110/10 Уктус", 60.47123, 56.77456) + b = _wfs_feature("sc_points_fullview.fid--9f8e7d", "ПС 110/10 Уктус", 60.47123, 56.77456) + + key_a = rw._stable_external_id(a, a["properties"]) + assert key_a == rw._stable_external_id(b, b["properties"]) + # Проверка ПО ЗНАЧЕНИЮ: ключ = sha256 по «имя|напряжение|lon_e5|lat_e5», + # не fid. Тот же seed повторён в data/sql/99c_power_supply_centers_dedup.sql. + expected = "h:" + hashlib.sha256("уктус|110/10|6047123|5677456".encode()).hexdigest()[:16] + assert key_a == expected == "h:844e54f0152d2799" + + +def test_stable_external_id_differs_on_attributes() -> None: + """Одинаковый fid, разные атрибуты (координата / имя) → РАЗНЫЕ ключи.""" + base = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Уктус", 60.47123, 56.77456) + moved = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Уктус", 60.47124, 56.77456) + renamed = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Северная", 60.47123, 56.77456) + + keys = {rw._stable_external_id(f, f["properties"]) for f in (base, moved, renamed)} + assert len(keys) == 3 + assert rw._stable_external_id(moved, moved["properties"]) == "h:8c5b5fa5f35c651b" + assert rw._stable_external_id(renamed, renamed["properties"]) == "h:297e00f6e0d4f3b8" + + +def test_stable_external_id_no_geometry() -> None: + """Фича без геометрии: координатные компоненты пустые, ключ всё равно стабилен.""" + f: dict[str, Any] = {"id": "fid--x", "properties": {"sc_name": "ПС 110/10 Уктус"}} + assert rw._stable_external_id(f, f["properties"]) == "h:eb91917f35aff23f" + + +def test_coord_e5_rounds_on_binary_double_not_decimal() -> None: + """Округление по ДВОИЧНОМУ double, не по десятичному представлению. + + 64.423605*1e5 == 6442360.499999999 → 6442360; «по десятичному» вышло бы 6442361 + (так считал бы round(ST_X(geom)::numeric*100000) — расхождение на 0.19% реальных + координат). 99c обязана давать те же цифры, поэтому семантика закреплена тестом. + """ + assert rw._coord_e5(64.423605) == "6442360" + assert rw._coord_e5(-64.423605) == "-6442360" + # 60.123455*1e5 == ровно 6012345.5 → полукругление ОТ нуля, симметрично знаку. + assert rw._coord_e5(60.123455) == "6012346" + assert rw._coord_e5(-60.123455) == "-6012346" + assert rw._coord_e5(60.6) == "6060000" + assert rw._coord_e5(None) == "" + + # ── sanitize_tp_capacity_mva (кВА-санитайз) ─────────────────────────────────── diff --git a/caddy/deploy-window.caddy.snippet b/caddy/deploy-window.caddy.snippet new file mode 100644 index 00000000..971464de --- /dev/null +++ b/caddy/deploy-window.caddy.snippet @@ -0,0 +1,54 @@ +# ═══════════════════════════════════════════════════════════════════════════ +# caddy/deploy-window.caddy.snippet — ответ на окно деплоя (#3274) +# +# Импортируется ВНУТРЬ `handle_errors 502 503 504 { ... }` (см. apps.caddy): +# сам по себе снипет ничего не перехватывает, он только решает, ЧТО отдать, +# когда апстрим не отвечает. +# +# ЧТО ЭТО ЛЕЧИТ, А ЧТО НЕТ. Каждый деплой tradein-frontend/tradein-backend +# оставляет окно 30–90 с, в котором контейнера просто нет: Caddy набирает +# новый апстрим сразу, тот ещё не слушает (замер по access-логам, #3274 — +# все 502 кластеризуются на окнах мержа, duration < 2 мс = мгновенный отказ +# соединения). Снипет НЕ УБИРАЕТ окно — он меняет то, что видит человек и +# клиент внутри окна. Настоящее лечение (готовность нового контейнера до +# переключения) — п.1 issue, решение владельца, здесь его нет. +# +# ПОЧЕМУ 503, А НЕ 502. 502 значит «апстрим ответил мусором» — постоянная +# поломка; поисковик по нему выкидывает страницу из индекса, клиентские +# библиотеки не ретраят. 503 + `Retry-After: 30` — стандартный код «временно +# недоступен, приходи через 30 секунд»: Googlebot держит страницу в индексе, +# HTTP-клиенты понимают, что повтор осмыслен. +# +# ПОЧЕМУ ДВА ТЕЛА. `/trade-in/api/*` вызывают из JS и внешних клиентов — они +# парсят JSON, и HTML-страница у них превращается в ошибку разбора вместо +# читаемого статуса. Всё остальное открывает человек браузером. +# +# ВНЕШНИХ РЕСУРСОВ В СТРАНИЦЕ НЕТ ВООБЩЕ — ни шрифта, ни CSS-файла, ни +# картинки. В окне деплоя они пришли бы с того же мёртвого апстрима, и +# страница-заглушка отрисовалась бы голым текстом. Отсюда же инлайновые +# `style=` вместо блока `