Compare commits
12 commits
fix/canon-
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| baed8f75a3 | |||
| dbac5d4e1c | |||
| f71a438495 | |||
| e5212e52ab | |||
| 98f587add6 | |||
| 7c52622a76 | |||
| a82239b931 | |||
| 1265e69412 | |||
| 8b9b541736 | |||
| a1c7f8ea92 | |||
| d633c0d1c3 | |||
| da5fd24435 |
25 changed files with 1382 additions and 157 deletions
|
|
@ -91,12 +91,16 @@ jobs:
|
||||||
METRICS_TELEGRAM_INFRA_TOPIC_ID: ${{ secrets.METRICS_TELEGRAM_INFRA_TOPIC_ID }}
|
METRICS_TELEGRAM_INFRA_TOPIC_ID: ${{ secrets.METRICS_TELEGRAM_INFRA_TOPIC_ID }}
|
||||||
METRICS_TELEGRAM_ONCALL: ${{ secrets.METRICS_TELEGRAM_ONCALL }}
|
METRICS_TELEGRAM_ONCALL: ${{ secrets.METRICS_TELEGRAM_ONCALL }}
|
||||||
ALERT_ACK_GLITCHTIP_SECRET: ${{ secrets.ALERT_ACK_GLITCHTIP_SECRET }}
|
ALERT_ACK_GLITCHTIP_SECRET: ${{ secrets.ALERT_ACK_GLITCHTIP_SECRET }}
|
||||||
|
# #3589: URL внешнего deadman-приёмника Watchdog (healthchecks.io и
|
||||||
|
# аналоги). Пусто — watchdog-ping остаётся без конфигов, см. блок
|
||||||
|
# METRICS_WATCHDOG_PING_BLOCK ниже.
|
||||||
|
METRICS_WATCHDOG_PING_URL: ${{ secrets.METRICS_WATCHDOG_PING_URL }}
|
||||||
# #3471: секрет ретранслятора Telegram Bot API (tg-relay). Пусто —
|
# #3471: секрет ретранслятора Telegram Bot API (tg-relay). Пусто —
|
||||||
# профиль relay не включаем (см. PROFILES ниже), а не падаем в
|
# профиль relay не включаем (см. PROFILES ниже), а не падаем в
|
||||||
# рестарт-луп: контейнер сам делает SystemExit на пустом секрете.
|
# рестарт-луп: контейнер сам делает SystemExit на пустом секрете.
|
||||||
TG_RELAY_SECRET: ${{ secrets.TG_RELAY_SECRET }}
|
TG_RELAY_SECRET: ${{ secrets.TG_RELAY_SECRET }}
|
||||||
with:
|
with:
|
||||||
envs: METRICS_TELEGRAM_BOT_TOKEN,METRICS_TELEGRAM_CHAT_ID,METRICS_TELEGRAM_TOPIC_ID,METRICS_TELEGRAM_INFRA_TOPIC_ID,METRICS_TELEGRAM_ONCALL,ALERT_ACK_GLITCHTIP_SECRET,TG_RELAY_SECRET
|
envs: METRICS_TELEGRAM_BOT_TOKEN,METRICS_TELEGRAM_CHAT_ID,METRICS_TELEGRAM_TOPIC_ID,METRICS_TELEGRAM_INFRA_TOPIC_ID,METRICS_TELEGRAM_ONCALL,ALERT_ACK_GLITCHTIP_SECRET,TG_RELAY_SECRET,METRICS_WATCHDOG_PING_URL
|
||||||
host: ${{ secrets.INFRA_DEPLOY_HOST || secrets.DEPLOY_HOST }}
|
host: ${{ secrets.INFRA_DEPLOY_HOST || secrets.DEPLOY_HOST }}
|
||||||
username: ${{ secrets.INFRA_DEPLOY_USER || secrets.DEPLOY_USER }}
|
username: ${{ secrets.INFRA_DEPLOY_USER || secrets.DEPLOY_USER }}
|
||||||
key: ${{ secrets.INFRA_DEPLOY_SSH_KEY || secrets.DEPLOY_SSH_KEY }}
|
key: ${{ secrets.INFRA_DEPLOY_SSH_KEY || secrets.DEPLOY_SSH_KEY }}
|
||||||
|
|
@ -193,6 +197,24 @@ jobs:
|
||||||
echo "Инфраструктура: тема ${INFRA_TOPIC_ID} по умолчанию (METRICS_TELEGRAM_INFRA_TOPIC_ID не задана)."
|
echo "Инфраструктура: тема ${INFRA_TOPIC_ID} по умолчанию (METRICS_TELEGRAM_INFRA_TOPIC_ID не задана)."
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
# Watchdog-пинг во внешний deadman-приёмник (healthchecks.io и
|
||||||
|
# аналоги) — заменяет регулярные сообщения в Telegram, см.
|
||||||
|
# комментарий у watchdog-ping в alertmanager.yml.tmpl. Блок
|
||||||
|
# целиком, а не значение — та же причина, что у
|
||||||
|
# METRICS_TELEGRAM_INFRA_TOPIC_LINE: пустой `url:` в
|
||||||
|
# receiver'е не деградирует, а валит amtool check-config
|
||||||
|
# целиком, то есть роняет ВЕСЬ алертинг из-за одного
|
||||||
|
# необязательного получателя. Секрет пока не заведён —
|
||||||
|
# деградация корректна: receiver остаётся без `*_configs` и
|
||||||
|
# молча ничего никуда не шлёт, amtool это пропускает.
|
||||||
|
if [ -n "${METRICS_WATCHDOG_PING_URL:-}" ]; then
|
||||||
|
METRICS_WATCHDOG_PING_BLOCK=$(printf ' webhook_configs:\n - url: "%s"\n send_resolved: false' "${METRICS_WATCHDOG_PING_URL}")
|
||||||
|
echo "Watchdog: внешний deadman-пинг настроен."
|
||||||
|
else
|
||||||
|
METRICS_WATCHDOG_PING_BLOCK=""
|
||||||
|
echo "::warning title=Watchdog без deadman-пинга::METRICS_WATCHDOG_PING_URL пуст — сторож мониторинга никуда не сообщает о своей живости. Заведи аккаунт healthchecks.io (или аналог) и секрет, иначе обрыв канала доставки не заметит никто (#3589)."
|
||||||
|
fi
|
||||||
|
|
||||||
# Резервный приёмник GlitchTip (#3471) отвечает 503 на любой
|
# Резервный приёмник GlitchTip (#3471) отвечает 503 на любой
|
||||||
# запрос, пока секрет пуст: тихо принимать чужие алерты настежь
|
# запрос, пока секрет пуст: тихо принимать чужие алерты настежь
|
||||||
# хуже, чем не принимать вовсе. Молчаливого отказа тут быть не
|
# хуже, чем не принимать вовсе. Молчаливого отказа тут быть не
|
||||||
|
|
@ -221,7 +243,8 @@ jobs:
|
||||||
METRICS_TELEGRAM_CHAT_ID="$METRICS_TELEGRAM_CHAT_ID" \
|
METRICS_TELEGRAM_CHAT_ID="$METRICS_TELEGRAM_CHAT_ID" \
|
||||||
METRICS_TELEGRAM_INFRA_TOPIC_LINE="$METRICS_TELEGRAM_INFRA_TOPIC_LINE" \
|
METRICS_TELEGRAM_INFRA_TOPIC_LINE="$METRICS_TELEGRAM_INFRA_TOPIC_LINE" \
|
||||||
METRICS_TELEGRAM_ONCALL="${METRICS_TELEGRAM_ONCALL:-}" \
|
METRICS_TELEGRAM_ONCALL="${METRICS_TELEGRAM_ONCALL:-}" \
|
||||||
envsubst '${METRICS_TELEGRAM_BOT_TOKEN} ${METRICS_TELEGRAM_CHAT_ID} ${METRICS_TELEGRAM_INFRA_TOPIC_LINE} ${METRICS_TELEGRAM_ONCALL}' \
|
METRICS_WATCHDOG_PING_BLOCK="$METRICS_WATCHDOG_PING_BLOCK" \
|
||||||
|
envsubst '${METRICS_TELEGRAM_BOT_TOKEN} ${METRICS_TELEGRAM_CHAT_ID} ${METRICS_TELEGRAM_INFRA_TOPIC_LINE} ${METRICS_TELEGRAM_ONCALL} ${METRICS_WATCHDOG_PING_BLOCK}' \
|
||||||
< ops/metrics/alertmanager/alertmanager.yml.tmpl \
|
< ops/metrics/alertmanager/alertmanager.yml.tmpl \
|
||||||
> ops/metrics/alertmanager/alertmanager.yml
|
> ops/metrics/alertmanager/alertmanager.yml
|
||||||
chmod 600 ops/metrics/alertmanager/alertmanager.yml
|
chmod 600 ops/metrics/alertmanager/alertmanager.yml
|
||||||
|
|
|
||||||
|
|
@ -166,6 +166,74 @@ _GEO_PRICE_RADIUS_M: float = 3000.0 # 3 км — городской радиу
|
||||||
_GEO_PRICE_MIN_LOTS: int = 10
|
_GEO_PRICE_MIN_LOTS: int = 10
|
||||||
_GEO_PRICE_MIN_COMPLEXES: int = 2
|
_GEO_PRICE_MIN_COMPLEXES: int = 2
|
||||||
|
|
||||||
|
_GEO_RADIUS_PRICE_SQL = text("""
|
||||||
|
-- #1964 physflat-дедуп: objective_lots раздут ~2.91×
|
||||||
|
-- (мульти lot_id на физлот) → n/вес медианы/гейт n≥10
|
||||||
|
-- были по пере-листингам. Дедупим INLINE через DISTINCT ON
|
||||||
|
-- (physflat-ключ, последний снапшот), scope протолкнут В
|
||||||
|
-- CTE через проекты ближних ЖК. НЕ через
|
||||||
|
-- v_objective_lots_latest: view материализует ВСЮ таблицу
|
||||||
|
-- (qual/join не проходят ниже DISTINCT ON) → seq-scan+sort
|
||||||
|
-- 1.76M (~6.4 s на request-path analyze_parcel). Дедуп до price-фильтра:
|
||||||
|
-- цена объективна по физлоту (последний снапшот).
|
||||||
|
-- #3583: дедуп по проекту в LATERAL + premise_kind='квартира' → Index Only Scan
|
||||||
|
-- objective_lots_physflat_covering_v2_idx без сортировки. Общий DISTINCT ON по
|
||||||
|
-- всем снапшотам проектов шёл external merge на диск: прод-EXPLAIN 17.09, 16
|
||||||
|
-- проектов, 690 мс против 62 мс. Апартаменты на проде 17.09 есть у 6 проектов из
|
||||||
|
-- complex_sources, у всех complex без координат → в радиус не попадают, фильтр
|
||||||
|
-- медиану не меняет (замер на 570 участках: 0 отличий).
|
||||||
|
WITH nearby_cx AS (
|
||||||
|
-- #3583: complex → проект Объектива из complex_sources (source='objective',
|
||||||
|
-- 1:1), а НЕ из objective_lots.complex_id. Тот проставлен один раз миграцией 76,
|
||||||
|
-- а еженедельный 70_parse_objective_raw.py UPSERT'ом по objective_lot_id
|
||||||
|
-- переписывает project_name и не трогает complex_id → под id ближнего ЖК
|
||||||
|
-- лежат лоты чужих (прод 17.09: 236 354 из 303 677 строк). Тот же дефект,
|
||||||
|
-- что #2962 в competitors.py, и та же сверка имени: связь fuzzy, у «ЖК VEER
|
||||||
|
-- PARK» стоит 'Clever Park' в 11.7 км, у «ЖК Графит» — 'Гранит'.
|
||||||
|
SELECT c.id, cs.source_id AS project_name
|
||||||
|
FROM complexes c
|
||||||
|
JOIN complex_sources cs
|
||||||
|
ON cs.complex_id = c.id
|
||||||
|
AND cs.source = 'objective'
|
||||||
|
CROSS JOIN LATERAL (
|
||||||
|
SELECT regexp_replace(lower(c.canonical_name), '[^0-9a-zа-яё]', '', 'g') AS cx_key,
|
||||||
|
regexp_replace(lower(cs.source_id), '[^0-9a-zа-яё]', '', 'g') AS project_key
|
||||||
|
) k
|
||||||
|
WHERE c.latitude IS NOT NULL
|
||||||
|
AND c.longitude IS NOT NULL
|
||||||
|
AND ST_DWithin(
|
||||||
|
ST_SetSRID(ST_MakePoint(c.longitude, c.latitude), 4326)::geography,
|
||||||
|
ST_SetSRID(
|
||||||
|
ST_MakePoint(CAST(:lon AS float), CAST(:lat AS float)), 4326
|
||||||
|
)::geography,
|
||||||
|
CAST(:radius_m AS float)
|
||||||
|
)
|
||||||
|
AND ( k.cx_key LIKE '%' || k.project_key || '%'
|
||||||
|
OR k.project_key LIKE '%' || k.cx_key || '%')
|
||||||
|
),
|
||||||
|
latest AS (
|
||||||
|
SELECT l.price_per_m2_rub, nc.id AS complex_id
|
||||||
|
FROM nearby_cx nc
|
||||||
|
CROSS JOIN LATERAL (
|
||||||
|
SELECT DISTINCT ON (ol.corpus_name, ol.section, ol.floor, ol.lot_number)
|
||||||
|
ol.price_per_m2_rub
|
||||||
|
FROM objective_lots ol
|
||||||
|
WHERE ol.project_name = nc.project_name
|
||||||
|
AND ol.premise_kind = 'квартира'
|
||||||
|
ORDER BY ol.corpus_name, ol.section, ol.floor, ol.lot_number,
|
||||||
|
ol.snapshot_date DESC, ol.id DESC
|
||||||
|
) l
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
percentile_cont(0.5) WITHIN GROUP (
|
||||||
|
ORDER BY price_per_m2_rub
|
||||||
|
) AS median,
|
||||||
|
count(*) AS n,
|
||||||
|
count(DISTINCT complex_id) AS n_complexes
|
||||||
|
FROM latest
|
||||||
|
WHERE price_per_m2_rub IS NOT NULL
|
||||||
|
""")
|
||||||
|
|
||||||
# #1960 «Медиана рынка»: минимум сделок квартальной росреестровской MV
|
# #1960 «Медиана рынка»: минимум сделок квартальной росреестровской MV
|
||||||
# (mv_quarter_price_per_m2.deals_count, окно 24 мес), чтобы её медиана вообще
|
# (mv_quarter_price_per_m2.deals_count, окно 24 мес), чтобы её медиана вообще
|
||||||
# могла служить последним fallback'ом для карточки district.median_price_per_m2.
|
# могла служить последним fallback'ом для карточки district.median_price_per_m2.
|
||||||
|
|
@ -3562,7 +3630,7 @@ def analyze_parcel(
|
||||||
# радиусе вокруг центроида участка (ST_DWithin). Закрывает пробел district_reference:
|
# радиусе вокруг центроида участка (ST_DWithin). Закрывает пробел district_reference:
|
||||||
# только 4 из 9 админ-районов ЕКБ матчатся с Objective по имени, остальные 5 без неё
|
# только 4 из 9 админ-районов ЕКБ матчатся с Objective по имени, остальные 5 без неё
|
||||||
# проваливались в class_norm. У objective_lots нет geom — коорды берём через
|
# проваливались в class_norm. У objective_lots нет geom — коорды берём через
|
||||||
# complexes (latitude/longitude), join по ol.complex_id = c.id (паттерн competitors.py).
|
# complexes (latitude/longitude), проект — через complex_sources (#3583, как #2962).
|
||||||
geo_radius_price: dict[str, Any]
|
geo_radius_price: dict[str, Any]
|
||||||
# Вырожденная геометрия → центроид свалился на хардкод-центр ЕКБ. Гео-радиусная
|
# Вырожденная геометрия → центроид свалился на хардкод-центр ЕКБ. Гео-радиусная
|
||||||
# медиана тогда = «3км вокруг центра города», но выдаётся за калиброванную рыночную
|
# медиана тогда = «3км вокруг центра города», но выдаётся за калиброванную рыночную
|
||||||
|
|
@ -3573,60 +3641,7 @@ def analyze_parcel(
|
||||||
with db.begin_nested():
|
with db.begin_nested():
|
||||||
grp_row = (
|
grp_row = (
|
||||||
db.execute(
|
db.execute(
|
||||||
text("""
|
_GEO_RADIUS_PRICE_SQL,
|
||||||
-- #1964 physflat-дедуп: objective_lots раздут ~2.91×
|
|
||||||
-- (мульти lot_id на физлот) → n/вес медианы/гейт n≥10
|
|
||||||
-- были по пере-листингам. Дедупим INLINE через DISTINCT ON
|
|
||||||
-- (physflat-ключ, последний снапшот), scope протолкнут В
|
|
||||||
-- CTE через complex_id ближних ЖК. НЕ через
|
|
||||||
-- v_objective_lots_latest: view материализует ВСЮ таблицу
|
|
||||||
-- (qual/join не проходят ниже DISTINCT ON) → seq-scan+sort
|
|
||||||
-- 1.76M (~6.4 s на request-path analyze_parcel). Inline:
|
|
||||||
-- geo-index по complexes → Nested Loop bitmap
|
|
||||||
-- objective_lots_complex_idx по ~186 ЖК → ~120 ms
|
|
||||||
-- (прод-EXPLAIN deep-review #1964). Дедуп до price-фильтра:
|
|
||||||
-- цена объективна по физлоту (последний снапшот).
|
|
||||||
WITH nearby_cx AS (
|
|
||||||
SELECT c.id
|
|
||||||
FROM complexes c
|
|
||||||
WHERE c.latitude IS NOT NULL
|
|
||||||
AND c.longitude IS NOT NULL
|
|
||||||
AND ST_DWithin(
|
|
||||||
ST_SetSRID(
|
|
||||||
ST_MakePoint(c.longitude, c.latitude),
|
|
||||||
4326
|
|
||||||
)::geography,
|
|
||||||
ST_SetSRID(
|
|
||||||
ST_MakePoint(
|
|
||||||
CAST(:lon AS float),
|
|
||||||
CAST(:lat AS float)
|
|
||||||
), 4326
|
|
||||||
)::geography,
|
|
||||||
CAST(:radius_m AS float)
|
|
||||||
)
|
|
||||||
),
|
|
||||||
latest AS (
|
|
||||||
SELECT DISTINCT ON (
|
|
||||||
ol.project_name, ol.corpus_name, ol.section,
|
|
||||||
ol.floor, ol.lot_number
|
|
||||||
)
|
|
||||||
ol.price_per_m2_rub,
|
|
||||||
ol.complex_id
|
|
||||||
FROM objective_lots ol
|
|
||||||
WHERE ol.complex_id IN (SELECT id FROM nearby_cx)
|
|
||||||
ORDER BY ol.project_name, ol.corpus_name, ol.section,
|
|
||||||
ol.floor, ol.lot_number,
|
|
||||||
ol.snapshot_date DESC, ol.id DESC
|
|
||||||
)
|
|
||||||
SELECT
|
|
||||||
percentile_cont(0.5) WITHIN GROUP (
|
|
||||||
ORDER BY price_per_m2_rub
|
|
||||||
) AS median,
|
|
||||||
count(*) AS n,
|
|
||||||
count(DISTINCT complex_id) AS n_complexes
|
|
||||||
FROM latest
|
|
||||||
WHERE price_per_m2_rub IS NOT NULL
|
|
||||||
"""),
|
|
||||||
{
|
{
|
||||||
"lon": centroid_lon,
|
"lon": centroid_lon,
|
||||||
"lat": centroid_lat,
|
"lat": centroid_lat,
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,123 @@
|
||||||
|
"""Гео-радиусная цена участка берёт лоты проектов ближних ЖК, а не всё под complex_id (#3583).
|
||||||
|
|
||||||
|
`objective_lots.complex_id` проставлен один раз миграцией 76, а еженедельный
|
||||||
|
`70_parse_objective_raw.py` UPSERT'ом по objective_lot_id переписывает project_name и
|
||||||
|
не трогает complex_id: под id ближнего ЖК лежат лоты чужих. Тот же дефект, что #2962
|
||||||
|
в competitors.py.
|
||||||
|
|
||||||
|
Тест герметичный и прогоняет НАСТОЯЩИЙ `_GEO_RADIUS_PRICE_SQL`: временные таблицы
|
||||||
|
затеняют боевые в пределах сессии. Нужен Postgres с PostGIS (ST_DWithin по geography).
|
||||||
|
В CI он есть, и там тест не пропускается: без PostGIS падает с настоящей причиной.
|
||||||
|
Пропуск разрешён только вне CI и объявлен в skip_allowlist.txt.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from sqlalchemy import create_engine, text
|
||||||
|
from sqlalchemy.orm import sessionmaker
|
||||||
|
|
||||||
|
|
||||||
|
def _dsn() -> str:
|
||||||
|
raw = os.environ.get("TEST_DATABASE_URL") or os.environ["DATABASE_URL"]
|
||||||
|
return (
|
||||||
|
raw
|
||||||
|
if raw.startswith("postgresql+")
|
||||||
|
else raw.replace("postgresql://", "postgresql+psycopg://")
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _postgis_reachable() -> tuple[bool, str]:
|
||||||
|
try:
|
||||||
|
eng = create_engine(_dsn(), connect_args={"connect_timeout": 3})
|
||||||
|
with eng.connect() as c:
|
||||||
|
if c.execute(text("SELECT 1 FROM pg_extension WHERE extname = 'postgis'")).first():
|
||||||
|
return True, ""
|
||||||
|
return False, "нет расширения postgis"
|
||||||
|
except Exception as exc:
|
||||||
|
return False, str(exc)
|
||||||
|
|
||||||
|
|
||||||
|
_DB_OK, _DB_ERR = _postgis_reachable()
|
||||||
|
_IN_CI = bool(os.environ.get("GITHUB_ACTIONS") or os.environ.get("CI"))
|
||||||
|
pytestmark = pytest.mark.skipif(
|
||||||
|
not _DB_OK and not _IN_CI, reason=f"Postgres/PostGIS недоступен: {_DB_ERR}"
|
||||||
|
)
|
||||||
|
|
||||||
|
_SCHEMA = [
|
||||||
|
"""CREATE TEMP TABLE complexes (
|
||||||
|
id bigint, canonical_name text, latitude double precision,
|
||||||
|
longitude double precision) ON COMMIT DROP""",
|
||||||
|
"""CREATE TEMP TABLE complex_sources (
|
||||||
|
complex_id bigint, source text, source_id text) ON COMMIT DROP""",
|
||||||
|
"""CREATE TEMP TABLE objective_lots (
|
||||||
|
id bigint, project_name text, corpus_name text, section text, floor int,
|
||||||
|
lot_number text, snapshot_date date, premise_kind text, complex_id bigint,
|
||||||
|
price_per_m2_rub numeric) ON COMMIT DROP""",
|
||||||
|
]
|
||||||
|
|
||||||
|
# Участок в (56.840, 60.600), радиус 3 км. Малахит — в ~10 км, вне радиуса.
|
||||||
|
_DATA = [
|
||||||
|
"""INSERT INTO complexes VALUES
|
||||||
|
(10, 'ЖК Мичуринский', 56.841, 60.601),
|
||||||
|
(20, 'ЖК VEER PARK', 56.845, 60.605),
|
||||||
|
(30, 'СтудияПарк', 56.835, 60.595),
|
||||||
|
(40, 'ЖК Малахит', 56.930, 60.600)""",
|
||||||
|
# 20 → 'Clever Park': неверная fuzzy-связь, как на проде (complexes.id=1493).
|
||||||
|
"""INSERT INTO complex_sources VALUES
|
||||||
|
(10, 'objective', 'Мичуринский'),
|
||||||
|
(20, 'objective', 'Clever Park'),
|
||||||
|
(30, 'objective', 'Студия Парк'),
|
||||||
|
(40, 'objective', 'Малахит')""",
|
||||||
|
# Под complex_id=10 лежат свой лот и три лота чужого «Малахита» с устаревшим
|
||||||
|
# complex_id; у новых лотов complex_id NULL. Лот «Мичуринский/1/1/3/3» в двух
|
||||||
|
# снапшотах: в медиану идёт последний (120 тыс.), а не старый (50 тыс.).
|
||||||
|
"""INSERT INTO objective_lots VALUES
|
||||||
|
(1, 'Мичуринский', '1', '1', 1, '1', '2026-05-10', 'квартира', 10, 100000),
|
||||||
|
(2, 'Мичуринский', '1', '1', 2, '2', '2026-09-15', 'квартира', NULL, 110000),
|
||||||
|
(3, 'Мичуринский', '1', '1', 3, '3', '2026-08-01', 'квартира', NULL, 50000),
|
||||||
|
(4, 'Мичуринский', '1', '1', 3, '3', '2026-09-15', 'квартира', NULL, 120000),
|
||||||
|
(5, 'Малахит', '1', '1', 1, '1', '2026-05-10', 'квартира', 10, 300000),
|
||||||
|
(6, 'Малахит', '1', '1', 2, '2', '2026-05-10', 'квартира', 10, 300000),
|
||||||
|
(7, 'Малахит', '1', '1', 3, '3', '2026-05-10', 'квартира', 10, 300000),
|
||||||
|
(8, 'Clever Park', '1', '1', 1, '1', '2026-09-15', 'квартира', NULL, 500000),
|
||||||
|
(9, 'Clever Park', '1', '1', 2, '2', '2026-09-15', 'квартира', NULL, 500000),
|
||||||
|
(10, 'Clever Park', '1', '1', 3, '3', '2026-09-15', 'квартира', NULL, 500000),
|
||||||
|
(11, 'Студия Парк', '1', '1', 1, '1', '2026-09-15', 'квартира', NULL, 90000),
|
||||||
|
(12, 'Студия Парк', '1', '1', 2, '2', '2026-09-15', 'квартира', NULL, 95000)""",
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(scope="module")
|
||||||
|
def row() -> dict[str, float]:
|
||||||
|
from app.api.v1.parcels import _GEO_PRICE_RADIUS_M, _GEO_RADIUS_PRICE_SQL
|
||||||
|
|
||||||
|
session = sessionmaker(bind=create_engine(_dsn()))()
|
||||||
|
try:
|
||||||
|
for stmt in _SCHEMA + _DATA:
|
||||||
|
session.execute(text(stmt))
|
||||||
|
r = (
|
||||||
|
session.execute(
|
||||||
|
_GEO_RADIUS_PRICE_SQL,
|
||||||
|
{"lon": 60.600, "lat": 56.840, "radius_m": _GEO_PRICE_RADIUS_M},
|
||||||
|
)
|
||||||
|
.mappings()
|
||||||
|
.one()
|
||||||
|
)
|
||||||
|
return {k: float(v) for k, v in r.items()}
|
||||||
|
finally:
|
||||||
|
session.rollback()
|
||||||
|
session.close()
|
||||||
|
|
||||||
|
|
||||||
|
def test_median_counts_own_projects_of_nearby_complexes(row) -> None:
|
||||||
|
"""Лоты 90/95 тыс. («Студия Парк») и 100/110/120 тыс. («Мичуринский») → медиана 100 тыс.
|
||||||
|
|
||||||
|
По устаревшему complex_id было бы 100 + три «Малахита» по 300 тыс. → 300 тыс.;
|
||||||
|
без сверки имени добавились бы три лота «Clever Park» по 500 тыс. → 115 тыс.
|
||||||
|
"""
|
||||||
|
assert row == {"median": 100000.0, "n": 5.0, "n_complexes": 2.0}
|
||||||
|
|
@ -21,7 +21,11 @@
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import importlib.util
|
import importlib.util
|
||||||
|
import logging
|
||||||
|
import socket
|
||||||
import sys
|
import sys
|
||||||
|
import threading
|
||||||
|
from http.server import ThreadingHTTPServer
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from types import ModuleType
|
from types import ModuleType
|
||||||
|
|
||||||
|
|
@ -179,3 +183,48 @@ def test_protuhshiy_token_ne_prinimaetsya(app: ModuleType, monkeypatch: pytest.M
|
||||||
code, _ = app.do_ack(token)
|
code, _ = app.do_ack(token)
|
||||||
assert code == 404
|
assert code == 404
|
||||||
assert app.sent == []
|
assert app.sent == []
|
||||||
|
|
||||||
|
|
||||||
|
# ── Секрет из query не попадает в лог (#3576) ────────────────────────────────
|
||||||
|
# GlitchTip шлёт секрет резервного вебхука только в `?secret=`, а http.server
|
||||||
|
# печатает строку запроса целиком — и в строке доступа, и в тексте ошибки
|
||||||
|
# разбора. Значение ниже выдуманное: проверяется, что его нет ни в одной записи.
|
||||||
|
|
||||||
|
_LEAK = "leak-probe-3576-VALUE"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("request_line", "masked", "lines_with_mask"),
|
||||||
|
[
|
||||||
|
# Строка доступа (log_request), путь резервного вебхука GlitchTip.
|
||||||
|
(f"POST /glitchtip?secret={_LEAK} HTTP/1.1", "/glitchtip?secret=***", 1),
|
||||||
|
# Имя с префиксом и соседний параметр — остальная строка цела.
|
||||||
|
(f"GET /ack/x?a=1&access_token={_LEAK}&b=2 HTTP/1.1", "?a=1&access_token=***&b=2", 1),
|
||||||
|
# Ошибка разбора (log_error): stdlib кладёт строку запроса в текст ошибки,
|
||||||
|
# затем та же строка идёт в строку доступа с кодом 400 — обе записи.
|
||||||
|
(f"POST /glitchtip?secret={_LEAK} junk HTTP/1.1", "/glitchtip?secret=***", 2),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_sekret_iz_query_ne_popadaet_v_log(
|
||||||
|
app: ModuleType,
|
||||||
|
caplog: pytest.LogCaptureFixture,
|
||||||
|
request_line: str,
|
||||||
|
masked: str,
|
||||||
|
lines_with_mask: int,
|
||||||
|
) -> None:
|
||||||
|
caplog.set_level(logging.INFO, logger="alert-ack")
|
||||||
|
srv = ThreadingHTTPServer(("127.0.0.1", 0), app.Handler)
|
||||||
|
threading.Thread(target=srv.serve_forever, daemon=True).start()
|
||||||
|
try:
|
||||||
|
with socket.create_connection(srv.server_address, timeout=5) as sock:
|
||||||
|
sock.sendall(f"{request_line}\r\nHost: x\r\nContent-Length: 0\r\n\r\n".encode())
|
||||||
|
sock.shutdown(socket.SHUT_WR)
|
||||||
|
while sock.recv(4096): # до закрытия: к этому моменту запись лога уже сделана
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
srv.shutdown()
|
||||||
|
srv.server_close()
|
||||||
|
|
||||||
|
messages = [r.getMessage() for r in caplog.records if r.name == "alert-ack"]
|
||||||
|
assert not [m for m in messages if _LEAK in m], f"значение секрета в логе: {messages}"
|
||||||
|
assert sum(masked in m for m in messages) == lines_with_mask, messages
|
||||||
|
|
|
||||||
|
|
@ -21,8 +21,9 @@
|
||||||
переменной — и инфраструктурная тема «метрики» оставалась пустой, а весь трафик,
|
переменной — и инфраструктурная тема «метрики» оставалась пустой, а весь трафик,
|
||||||
и клиентский, и инфраструктурный, копился в теме «алерты». Владелец решил
|
и клиентский, и инфраструктурный, копился в теме «алерты». Владелец решил
|
||||||
развести: инфраструктура — в «метрики», клиентские инциденты — в «алерты».
|
развести: инфраструктура — в «метрики», клиентские инциденты — в «алерты».
|
||||||
Тесты ниже закрепляют именно это: `telegram`/`telegram-heartbeat` получают
|
Тест ниже закрепляет именно это для единственного оставшегося прямого
|
||||||
ИНФРАСТРУКТУРНУЮ тему, а не общую.
|
получателя — `telegram` (Watchdog с 17.09 больше не шлёт в Telegram вовсе,
|
||||||
|
у него теперь внешний webhook-приёмник `watchdog-ping`, см. шаблон).
|
||||||
|
|
||||||
Тесты рендерят шаблон обоими способами и разбирают результат как YAML —
|
Тесты рендерят шаблон обоими способами и разбирают результат как YAML —
|
||||||
проверяется фактический конфиг, а не наличие нужных слов в тексте.
|
проверяется фактический конфиг, а не наличие нужных слов в тексте.
|
||||||
|
|
@ -43,9 +44,56 @@ WORKFLOW = REPO_ROOT / ".forgejo" / "workflows" / "deploy-metrics.yml"
|
||||||
|
|
||||||
INFRA_TOPIC_LINE = " message_thread_id: 245"
|
INFRA_TOPIC_LINE = " message_thread_id: 245"
|
||||||
|
|
||||||
|
# Заглушки для переменных, которые деплой может подставить. Значение
|
||||||
|
# METRICS_WATCHDOG_PING_BLOCK по умолчанию пустое — это реальный дефолт
|
||||||
|
# деплоя, когда секрет не заведён (#3589), а не тестовое упрощение.
|
||||||
|
_DUMMY_VALUES = {
|
||||||
|
"METRICS_TELEGRAM_BOT_TOKEN": "123:ABC",
|
||||||
|
"METRICS_TELEGRAM_CHAT_ID": "-100123",
|
||||||
|
"METRICS_TELEGRAM_ONCALL": "",
|
||||||
|
"METRICS_WATCHDOG_PING_BLOCK": "",
|
||||||
|
}
|
||||||
|
|
||||||
def _render(infra_topic_line: str) -> dict:
|
|
||||||
"""Повторяет подстановку деплоя и разбирает результат как YAML.
|
def _envsubst_allowlist() -> set[str]:
|
||||||
|
"""Реальный список переменных, которые деплой передаёт в envsubst.
|
||||||
|
|
||||||
|
Не хардкодим копию списка — #3589 случился именно так: шаблон завёл
|
||||||
|
`${METRICS_WATCHDOG_PING_URL}`, а список envsubst в деплое не пополнили,
|
||||||
|
и тест этого не заметил, потому что сам подставлял значение мимо деплоя.
|
||||||
|
"""
|
||||||
|
assert WORKFLOW.is_file(), f"нет {WORKFLOW} — воркфлоу переехал, гейт ослеп"
|
||||||
|
text = WORKFLOW.read_text(encoding="utf-8")
|
||||||
|
m = re.search(r"envsubst '([^']+)'", text)
|
||||||
|
assert m, "не нашёл вызов envsubst в деплое"
|
||||||
|
return set(re.findall(r"\$\{(\w+)\}", m.group(1)))
|
||||||
|
|
||||||
|
|
||||||
|
def _template_placeholders() -> set[str]:
|
||||||
|
text = TMPL.read_text(encoding="utf-8")
|
||||||
|
return set(re.findall(r"\$\{(\w+)\}", text))
|
||||||
|
|
||||||
|
|
||||||
|
def test_every_template_placeholder_is_in_envsubst_allowlist() -> None:
|
||||||
|
"""Регресс #3589: переменная шаблона обязана быть в allow-list envsubst.
|
||||||
|
|
||||||
|
Тогда в шаблоне появился `${METRICS_WATCHDOG_PING_URL}`, а список
|
||||||
|
envsubst в деплое не пополнили. envsubst подставляет ТОЛЬКО
|
||||||
|
перечисленные переменные — забытая долетает до `amtool check-config`
|
||||||
|
литералом плейсхолдера и валит проверку (`unsupported scheme ""`), то
|
||||||
|
есть роняет ВЕСЬ Alertmanager, а не только Watchdog.
|
||||||
|
"""
|
||||||
|
missing = _template_placeholders() - _envsubst_allowlist()
|
||||||
|
assert not missing, f"эти переменные шаблона деплой не подставляет: {missing}"
|
||||||
|
|
||||||
|
|
||||||
|
def _render(infra_topic_line: str, watchdog_ping_block: str | None = None) -> dict:
|
||||||
|
"""Повторяет ТОЧНО ТУ ЖЕ подстановку, что делает деплой, и разбирает YAML.
|
||||||
|
|
||||||
|
Подставляются только переменные из реального allow-list envsubst деплоя
|
||||||
|
(`_envsubst_allowlist`) — не весь известный тесту набор. Так регресс
|
||||||
|
#3589 (переменная в шаблоне, забытая в allow-list) ловится именно здесь:
|
||||||
|
`assert "${" not in rendered` ниже упадёт, если что-то не подставилось.
|
||||||
|
|
||||||
Строка темы в шаблоне ровно одна — инфраструктурная (#3163). Тема
|
Строка темы в шаблоне ровно одна — инфраструктурная (#3163). Тема
|
||||||
клиентских инцидентов сюда не подставляется вовсе: маршрут
|
клиентских инцидентов сюда не подставляется вовсе: маршрут
|
||||||
|
|
@ -55,11 +103,18 @@ def _render(infra_topic_line: str) -> dict:
|
||||||
"""
|
"""
|
||||||
assert TMPL.is_file(), f"нет {TMPL} — шаблон переехал, гейт ослеп"
|
assert TMPL.is_file(), f"нет {TMPL} — шаблон переехал, гейт ослеп"
|
||||||
text = TMPL.read_text(encoding="utf-8")
|
text = TMPL.read_text(encoding="utf-8")
|
||||||
rendered = (
|
|
||||||
text.replace("${METRICS_TELEGRAM_BOT_TOKEN}", "123:ABC")
|
values = dict(_DUMMY_VALUES)
|
||||||
.replace("${METRICS_TELEGRAM_CHAT_ID}", "-100123")
|
values["METRICS_TELEGRAM_INFRA_TOPIC_LINE"] = infra_topic_line
|
||||||
.replace("${METRICS_TELEGRAM_INFRA_TOPIC_LINE}", infra_topic_line)
|
if watchdog_ping_block is not None:
|
||||||
)
|
values["METRICS_WATCHDOG_PING_BLOCK"] = watchdog_ping_block
|
||||||
|
|
||||||
|
allowlist = _envsubst_allowlist()
|
||||||
|
rendered = text
|
||||||
|
for name, value in values.items():
|
||||||
|
if name in allowlist:
|
||||||
|
rendered = rendered.replace("${" + name + "}", value)
|
||||||
|
|
||||||
assert "${" not in rendered, (
|
assert "${" not in rendered, (
|
||||||
"в отрендеренном конфиге остался литерал плейсхолдера — "
|
"в отрендеренном конфиге остался литерал плейсхолдера — "
|
||||||
"значит в шаблоне появилась подстановка, о которой тест не знает"
|
"значит в шаблоне появилась подстановка, о которой тест не знает"
|
||||||
|
|
@ -67,6 +122,37 @@ def _render(infra_topic_line: str) -> dict:
|
||||||
return yaml.safe_load(rendered)
|
return yaml.safe_load(rendered)
|
||||||
|
|
||||||
|
|
||||||
|
def test_watchdog_receiver_without_secret_has_no_configs_and_parses() -> None:
|
||||||
|
"""Секрет не заведён (реальное состояние прода сейчас) — конфиг всё равно жив.
|
||||||
|
|
||||||
|
`watchdog-ping` остаётся без единого `*_configs` — валидный receiver,
|
||||||
|
Alertmanager его просто пропускает. Деградация корректна: Watchdog никуда
|
||||||
|
не пингует, но остальной алертинг (`telegram`, `telegram-clients`) цел.
|
||||||
|
"""
|
||||||
|
cfg = _render(INFRA_TOPIC_LINE, watchdog_ping_block="")
|
||||||
|
receivers = {r["name"]: r for r in cfg["receivers"]}
|
||||||
|
assert "watchdog-ping" in receivers, "receiver watchdog-ping пропал из конфига"
|
||||||
|
watchdog = receivers["watchdog-ping"]
|
||||||
|
assert "webhook_configs" not in watchdog, "пустой секрет не должен оставлять webhook_configs"
|
||||||
|
assert "telegram" in receivers and "telegram-clients" in receivers, (
|
||||||
|
"остальной алертинг не должен пострадать из-за пустого watchdog-секрета"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_watchdog_receiver_with_secret_gets_webhook() -> None:
|
||||||
|
"""Секрет задан — Watchdog реально пингует внешний deadman-приёмник."""
|
||||||
|
block = (
|
||||||
|
" webhook_configs:\n"
|
||||||
|
' - url: "https://hc-ping.com/dummy"\n'
|
||||||
|
" send_resolved: false"
|
||||||
|
)
|
||||||
|
cfg = _render(INFRA_TOPIC_LINE, watchdog_ping_block=block)
|
||||||
|
receivers = {r["name"]: r for r in cfg["receivers"]}
|
||||||
|
hooks = receivers["watchdog-ping"].get("webhook_configs") or []
|
||||||
|
assert hooks and hooks[0].get("url") == "https://hc-ping.com/dummy"
|
||||||
|
assert hooks[0].get("send_resolved") is False
|
||||||
|
|
||||||
|
|
||||||
def _telegram_configs(cfg: dict) -> list[dict]:
|
def _telegram_configs(cfg: dict) -> list[dict]:
|
||||||
out = []
|
out = []
|
||||||
for r in cfg.get("receivers", []):
|
for r in cfg.get("receivers", []):
|
||||||
|
|
@ -76,16 +162,16 @@ def _telegram_configs(cfg: dict) -> list[dict]:
|
||||||
|
|
||||||
|
|
||||||
def test_topic_lands_in_every_telegram_receiver() -> None:
|
def test_topic_lands_in_every_telegram_receiver() -> None:
|
||||||
"""Оба прямых получателя адресуют ИНФРАСТРУКТУРНУЮ тему, а не клиентскую (#3163).
|
"""Единственный прямой получатель адресует ИНФРАСТРУКТУРНУЮ тему, а не клиентскую (#3163).
|
||||||
|
|
||||||
Получателей два — `telegram` и `telegram-heartbeat`. До разделения тем оба
|
До разделения тем `telegram` и `telegram-heartbeat` брали топик из одной
|
||||||
брали топик из одной переменной с клиентскими инцидентами, и тема «метрики»
|
переменной с клиентскими инцидентами, и тема «метрики» (245) оставалась
|
||||||
(245) оставалась пустой. Если heartbeat уйдёт не в ту тему, «мониторинг жив»
|
пустой. С 17.09 (устранение шума Watchdog) прямой Telegram-получатель
|
||||||
будет капать мимо, и это заметят не сразу — сюда же попадёт и весь
|
остался один — `telegram`; Watchdog теперь пингует внешний
|
||||||
инфраструктурный шум.
|
deadman-приёмник вебхуком (`watchdog-ping`, без topic вовсе — не Telegram).
|
||||||
"""
|
"""
|
||||||
cfgs = _telegram_configs(_render(INFRA_TOPIC_LINE))
|
cfgs = _telegram_configs(_render(INFRA_TOPIC_LINE))
|
||||||
assert len(cfgs) >= 2, f"ожидалось минимум два получателя telegram, найдено {len(cfgs)}"
|
assert len(cfgs) >= 1, f"ожидался хотя бы один получатель telegram, найдено {len(cfgs)}"
|
||||||
for c in cfgs:
|
for c in cfgs:
|
||||||
assert c.get("message_thread_id") == 245, f"инфраструктурный топик не проставлен: {c}"
|
assert c.get("message_thread_id") == 245, f"инфраструктурный топик не проставлен: {c}"
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -114,7 +114,7 @@ def test_klientskiy_marshrut_idyot_v_servis_knopki() -> None:
|
||||||
"""
|
"""
|
||||||
text = TEMPLATE.read_text(encoding="utf-8")
|
text = TEMPLATE.read_text(encoding="utf-8")
|
||||||
block = text[text.index("- name: telegram-clients") :]
|
block = text[text.index("- name: telegram-clients") :]
|
||||||
block = block[: block.index("- name: telegram-heartbeat")]
|
block = block[: block.index("- name: watchdog-ping")]
|
||||||
assert "webhook_configs" in block, "клиентский приёмник не переключён на сервис"
|
assert "webhook_configs" in block, "клиентский приёмник не переключён на сервис"
|
||||||
assert "alert-ack:8080/alertmanager" in block, "вебхук указывает не на сервис кнопки"
|
assert "alert-ack:8080/alertmanager" in block, "вебхук указывает не на сервис кнопки"
|
||||||
# У прочих приёмников прямой путь сохранён.
|
# У прочих приёмников прямой путь сохранён.
|
||||||
|
|
|
||||||
|
|
@ -161,6 +161,12 @@ tests/services/site_finder/test_2962_competitors_gapfill_bridge.py::test_project
|
||||||
tests/services/site_finder/test_2962_competitors_gapfill_bridge.py::test_nearest_complex_without_lots_does_not_eat_the_match
|
tests/services/site_finder/test_2962_competitors_gapfill_bridge.py::test_nearest_complex_without_lots_does_not_eat_the_match
|
||||||
tests/services/site_finder/test_2962_competitors_gapfill_bridge.py::test_complex_with_two_projects_takes_the_matching_one
|
tests/services/site_finder/test_2962_competitors_gapfill_bridge.py::test_complex_with_two_projects_takes_the_matching_one
|
||||||
|
|
||||||
|
# ── #3583: гео-радиусная цена участка (complex_sources → project_name) ────────
|
||||||
|
# Тот же случай, что #2962: гоняет НАСТОЯЩИЙ parcels._GEO_RADIUS_PRICE_SQL на
|
||||||
|
# временных таблицах, нужен PostGIS. В CI идёт и пропуститься не может (CI=true
|
||||||
|
# выключает skipif). Запись только для машины без базы.
|
||||||
|
tests/api/v1/test_3583_geo_radius_price_project_bridge.py::test_median_counts_own_projects_of_nearby_complexes
|
||||||
|
|
||||||
# ── #2464: backfill act_date (миграция 191) ──────────────────────────────────
|
# ── #2464: backfill act_date (миграция 191) ──────────────────────────────────
|
||||||
# Нужен живой Postgres: тесты создают ВРЕМЕННУЮ копию land_reservation в прод-форме
|
# Нужен живой Postgres: тесты создают ВРЕМЕННУЮ копию land_reservation в прод-форме
|
||||||
# (9+2 строки с датой Генплана + контрольные посторонние) и прогоняют ТЕЛО миграции
|
# (9+2 строки с датой Генплана + контрольные посторонние) и прогоняют ТЕЛО миграции
|
||||||
|
|
|
||||||
93
ops/forgejo/docker-compose.yml
Normal file
93
ops/forgejo/docker-compose.yml
Normal file
|
|
@ -0,0 +1,93 @@
|
||||||
|
# Source-of-truth copy of the Forgejo compose file.
|
||||||
|
#
|
||||||
|
# Forgejo is NOT part of the automated deploy pipeline (deploy.yml only
|
||||||
|
# manages /opt/gendesign via `git reset --hard origin/main` on the main and
|
||||||
|
# obsidian stacks). Forgejo lives separately at /home/gendesign/forgejo on
|
||||||
|
# the VM and is a plain directory there — NOT a git checkout — so changes
|
||||||
|
# here do not auto-apply. Sync manually:
|
||||||
|
#
|
||||||
|
# scp ops/forgejo/docker-compose.yml gendesign:/home/gendesign/forgejo/docker-compose.yml
|
||||||
|
# ssh gendesign "cd /home/gendesign/forgejo && docker compose up -d --force-recreate forgejo"
|
||||||
|
#
|
||||||
|
# `up -d --force-recreate` (not `restart`) is required: Forgejo generates
|
||||||
|
# app.ini from the FORGEJO__* env vars via /usr/local/bin/environment-to-ini
|
||||||
|
# at container start, and `restart` does not re-read `environment:` from a
|
||||||
|
# changed compose file (see docker-compose pitfall #1 in devops CLAUDE.md).
|
||||||
|
#
|
||||||
|
# Config keys verified 2026-09-16 against the running image
|
||||||
|
# (codeberg.org/forgejo/forgejo:10, Forgejo 10.0.3+gitea-1.22.0) by
|
||||||
|
# extracting Go struct tags from the binary (`strings` on
|
||||||
|
# /app/gitea/gitea) and by a dry-run of environment-to-ini in a scratch
|
||||||
|
# dir inside the container (no prod files touched). Do not re-derive these
|
||||||
|
# from memory — a wrong key is silently ignored (empty section) and
|
||||||
|
# creates a false sense of safety.
|
||||||
|
#
|
||||||
|
# Dotted section names ([cron.archive_cleanup]) must be encoded as
|
||||||
|
# `_0x2E_` in the env var per the container's own
|
||||||
|
# `environment-to-ini --help` (confirmed empirically, see PR description).
|
||||||
|
services:
|
||||||
|
forgejo:
|
||||||
|
image: codeberg.org/forgejo/forgejo:10
|
||||||
|
container_name: forgejo
|
||||||
|
restart: unless-stopped
|
||||||
|
environment:
|
||||||
|
USER_UID: 1000
|
||||||
|
USER_GID: 1000
|
||||||
|
FORGEJO__database__DB_TYPE: postgres
|
||||||
|
FORGEJO__database__HOST: infra-postgres:5432
|
||||||
|
FORGEJO__database__NAME: forgejo
|
||||||
|
FORGEJO__database__USER: forgejo
|
||||||
|
FORGEJO__database__PASSWD: ${FORGEJO_DB_PASS}
|
||||||
|
FORGEJO__server__DOMAIN: git.gendsgn.ru
|
||||||
|
FORGEJO__server__ROOT_URL: https://git.gendsgn.ru/
|
||||||
|
FORGEJO__server__SSH_PORT: 2222
|
||||||
|
FORGEJO__server__SSH_LISTEN_PORT: 22
|
||||||
|
FORGEJO__server__START_SSH_SERVER: "false"
|
||||||
|
FORGEJO__service__DISABLE_REGISTRATION: "true"
|
||||||
|
FORGEJO__service__REQUIRE_SIGNIN_VIEW: "false"
|
||||||
|
FORGEJO__actions__ENABLED: "true"
|
||||||
|
FORGEJO__actions__DEFAULT_ACTIONS_URL: "github"
|
||||||
|
# Actions Log/artifact retention — was unset (Forgejo defaults), which
|
||||||
|
# let CI run logs/artifacts accumulate indefinitely. Artifact
|
||||||
|
# retention checked against .forgejo/workflows + .github/workflows on
|
||||||
|
# 2026-09-16: no workflow uploads/downloads artifacts today, so 14d
|
||||||
|
# cannot break a cross-job dependency. Revisit this comment if a
|
||||||
|
# workflow starts using actions/upload-artifact.
|
||||||
|
FORGEJO__actions__LOG_RETENTION_DAYS: "30"
|
||||||
|
FORGEJO__actions__ARTIFACT_RETENTION_DAYS: "14"
|
||||||
|
# Repo-archive cache cleanup — was entirely absent (no
|
||||||
|
# [cron.archive_cleanup] section), so it ran on Forgejo's own default
|
||||||
|
# schedule (once every 24h, deleting archives older than 24h). That
|
||||||
|
# let an external crawler hitting /<owner>/<repo>/archive/<ref>
|
||||||
|
# balloon the cache to ~48GB/145GB disk before the daily sweep caught
|
||||||
|
# up (see PR #3534, which closed the path in Caddy as the primary
|
||||||
|
# fix). This is the second line of defense if that Caddy rule is ever
|
||||||
|
# removed: run hourly, evict anything older than 1h.
|
||||||
|
FORGEJO__CRON_0x2E_ARCHIVE_CLEANUP__ENABLED: "true"
|
||||||
|
FORGEJO__CRON_0x2E_ARCHIVE_CLEANUP__RUN_AT_START: "true"
|
||||||
|
FORGEJO__CRON_0x2E_ARCHIVE_CLEANUP__SCHEDULE: "@every 1h"
|
||||||
|
FORGEJO__CRON_0x2E_ARCHIVE_CLEANUP__OLDER_THAN: "1h"
|
||||||
|
FORGEJO__security__INSTALL_LOCK: "true"
|
||||||
|
volumes:
|
||||||
|
- ./data/forgejo:/data
|
||||||
|
# Container log cap. Measured 2026-09-17: this container's json log had
|
||||||
|
# grown to 3.35 GiB in 23 days (~150 MB/day) with no rotation at all,
|
||||||
|
# the single largest log on the host by two orders of magnitude. Docker's
|
||||||
|
# json-file driver never rotates unless told to, and there is no
|
||||||
|
# /etc/docker/daemon.json on this VM to set a global default. 50m x 3
|
||||||
|
# caps this container at 150 MB; recreating it also drops the old
|
||||||
|
# unbounded file.
|
||||||
|
logging:
|
||||||
|
driver: json-file
|
||||||
|
options:
|
||||||
|
max-size: "50m"
|
||||||
|
max-file: "3"
|
||||||
|
ports:
|
||||||
|
- "2222:22"
|
||||||
|
networks:
|
||||||
|
- gendesign_default
|
||||||
|
|
||||||
|
networks:
|
||||||
|
gendesign_default:
|
||||||
|
external: true
|
||||||
|
name: gendesign_default
|
||||||
|
|
@ -57,6 +57,7 @@ import html
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
|
import re
|
||||||
import secrets
|
import secrets
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
|
|
@ -76,6 +77,16 @@ TTL_SEC = int(os.environ.get("ALERT_ACK_TTL_MIN", "1440")) * 60
|
||||||
GLITCHTIP_SECRET = os.environ.get("ALERT_ACK_GLITCHTIP_SECRET", "")
|
GLITCHTIP_SECRET = os.environ.get("ALERT_ACK_GLITCHTIP_SECRET", "")
|
||||||
API = "https://api.telegram.org/bot{}/{}"
|
API = "https://api.telegram.org/bot{}/{}"
|
||||||
|
|
||||||
|
# Значения секретных query-параметров в логе — `***` (#3576). GlitchTip шлёт
|
||||||
|
# секрет только в `?secret=`, а строка запроса целиком уходит в access-log. То
|
||||||
|
# же выражение, что у бэкенда МЕРЫ (tradein-mvp/backend/app/core/log_scrub.py,
|
||||||
|
# #3154) и у Alloy (#3354); импортировать нельзя — сервис без зависимостей.
|
||||||
|
_SENSITIVE_QUERY = re.compile(
|
||||||
|
r"([?&][\w.-]*(?:secret|token|api[-_]?key|apikey|access[-_]?token|password|signature|sig)=)"
|
||||||
|
r"[^&\s\"'<>]+",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
|
||||||
# token -> {"message_id": int, "title": str, "created": float, "acked_by": str|None}
|
# token -> {"message_id": int, "title": str, "created": float, "acked_by": str|None}
|
||||||
_PENDING: dict[str, dict] = {}
|
_PENDING: dict[str, dict] = {}
|
||||||
_LOCK = threading.Lock()
|
_LOCK = threading.Lock()
|
||||||
|
|
@ -319,7 +330,9 @@ class Handler(BaseHTTPRequestHandler):
|
||||||
protocol_version = "HTTP/1.1"
|
protocol_version = "HTTP/1.1"
|
||||||
|
|
||||||
def log_message(self, fmt: str, *args) -> None: # noqa: A003 — подпись из stdlib
|
def log_message(self, fmt: str, *args) -> None: # noqa: A003 — подпись из stdlib
|
||||||
log.info("%s %s", self.address_string(), fmt % args)
|
# Сюда сходятся и строка доступа (log_request), и ошибки разбора запроса
|
||||||
|
# (log_error: «Bad request syntax ('POST /glitchtip?secret=…')») — маскируем здесь.
|
||||||
|
log.info("%s %s", self.address_string(), _SENSITIVE_QUERY.sub(r"\1***", fmt % args))
|
||||||
|
|
||||||
def _reply(self, code: int, body: bytes, ctype: str = "text/html; charset=utf-8") -> None:
|
def _reply(self, code: int, body: bytes, ctype: str = "text/html; charset=utf-8") -> None:
|
||||||
self.send_response(code)
|
self.send_response(code)
|
||||||
|
|
|
||||||
|
|
@ -15,11 +15,11 @@
|
||||||
# клиентские инциденты МЕРЫ) и тему «метрики» (245, инфраструктурный шум —
|
# клиентские инциденты МЕРЫ) и тему «метрики» (245, инфраструктурный шум —
|
||||||
# диск, память, просевший экспортер). Инфраструктурный шум и клиентский
|
# диск, память, просевший экспортер). Инфраструктурный шум и клиентский
|
||||||
# инцидент не равны по срочности, а смешанные в одной теме они обучают
|
# инцидент не равны по срочности, а смешанные в одной теме они обучают
|
||||||
# пролистывать обе. Поэтому `telegram` и `telegram-heartbeat` ниже адресуют
|
# пролистывать обе. Поэтому `telegram` ниже адресует ИНФРАСТРУКТУРНУЮ тему
|
||||||
# ИНФРАСТРУКТУРНУЮ тему (`METRICS_TELEGRAM_INFRA_TOPIC_LINE`). Получателя
|
# (`METRICS_TELEGRAM_INFRA_TOPIC_LINE`). Получателя `telegram-clients` в этом
|
||||||
# `telegram-clients` в этом списке нет: он не шлёт в Telegram напрямую, а
|
# списке нет: он не шлёт в Telegram напрямую, а вебхуком уходит в alert-ack, и
|
||||||
# вебхуком уходит в alert-ack, и тему адресует сам, своей переменной
|
# тему адресует сам, своей переменной METRICS_TELEGRAM_TOPIC_ID. `watchdog-ping`
|
||||||
# METRICS_TELEGRAM_TOPIC_ID.
|
# тоже не шлёт в Telegram вовсе — см. комментарий у его маршрута ниже.
|
||||||
|
|
||||||
global:
|
global:
|
||||||
resolve_timeout: 5m
|
resolve_timeout: 5m
|
||||||
|
|
@ -35,14 +35,35 @@ route:
|
||||||
repeat_interval: 6h
|
repeat_interval: 6h
|
||||||
|
|
||||||
routes:
|
routes:
|
||||||
# Watchdog не должен смешиваться с настоящими алертами и не должен молчать:
|
# Watchdog — «сторож сторожа», горит ВСЕГДА по построению (`vector(1)`,
|
||||||
# это «сторож сторожа», он горит всегда и подтверждает, что канал доставки жив.
|
# см. infra.yml). Раньше уходил в Telegram раз в 12ч — 14 сообщений в
|
||||||
- receiver: telegram-heartbeat
|
# неделю ни о чём, и именно они приучили пролистывать инфра-тему: 17.09
|
||||||
|
# настоящий DiskWillFillIn24h утонул между Watchdog и вечно горящим
|
||||||
|
# NoActiveCeleryWorkers, диск дошёл до 84% незамеченным.
|
||||||
|
#
|
||||||
|
# ПОЧЕМУ ВНЕШНИЙ DEADMAN-ПРИЁМНИК, А НЕ «РЕЖЕ» И НЕ ОТДЕЛЬНАЯ ТЕМА.
|
||||||
|
# Увеличенный интервал по-прежнему кладёт человеку регулярное сообщение —
|
||||||
|
# просто реже, и его тоже рано или поздно начнут пролистывать. Отдельная
|
||||||
|
# техническая тема — это ещё один chat_id/topic_id и ещё один канал,
|
||||||
|
# за которым НАДО СПЕЦИАЛЬНО следить, то есть тот же человеческий цикл,
|
||||||
|
# сдвинутый в другое место. Внешний deadman-приёмник (healthchecks.io и
|
||||||
|
# аналоги) устроен наоборот: Alertmanager молча шлёт HTTP-пинг на каждый
|
||||||
|
# Watchdog, и пока пинги идут — сервис МОЛЧИТ. Он заговорит (email/свой
|
||||||
|
# alert) только когда пинг ПЕРЕСТАНЕТ приходить, то есть ровно когда
|
||||||
|
# канал доставки умер, — это и есть смысл «сторожа сторожа», без единого
|
||||||
|
# штатного сообщения человеку. Полностью выключать эту проверку нельзя —
|
||||||
|
# remove бы всей ветки Watchdog это и сделал.
|
||||||
|
#
|
||||||
|
# METRICS_WATCHDOG_PING_URL пока НЕ заведён на хосте (нужен аккаунт
|
||||||
|
# healthchecks.io/аналога) — до тех пор webhook будет молча падать по
|
||||||
|
# DNS/сети, Alertmanager это тихо ретраит; человека это не касается ни
|
||||||
|
# раньше, ни теперь.
|
||||||
|
- receiver: watchdog-ping
|
||||||
matchers:
|
matchers:
|
||||||
- alertname = "Watchdog"
|
- alertname = "Watchdog"
|
||||||
group_wait: 0s
|
group_wait: 0s
|
||||||
group_interval: 12h
|
group_interval: 5m
|
||||||
repeat_interval: 12h
|
repeat_interval: 5m
|
||||||
|
|
||||||
# Клиентский инцидент. host="apps" — это продуктовая машина: если на ней
|
# Клиентский инцидент. host="apps" — это продуктовая машина: если на ней
|
||||||
# критично, значит МЕРА и Site Finder недоступны людям, а не «где-то в
|
# критично, значит МЕРА и Site Finder недоступны людям, а не «где-то в
|
||||||
|
|
@ -63,6 +84,20 @@ route:
|
||||||
group_wait: 10s
|
group_wait: 10s
|
||||||
repeat_interval: 30m
|
repeat_interval: 30m
|
||||||
|
|
||||||
|
# Эскалация по длительности (AlertFiringTooLong, prometheus/rules/infra.yml)
|
||||||
|
# — сигнал о том, что какую-то другую тревогу не заметили или на неё
|
||||||
|
# забили дольше 6 часов. Она НЕ про клиентский инцидент, но обязана быть
|
||||||
|
# заметнее обычной инфраструктуры, поэтому уходит в ту же тему, где
|
||||||
|
# владелец бывает чаще, а не смешивается с общим потоком severity=critical
|
||||||
|
# ниже. Матчим по имени, а не по host="apps": исходная тревога может
|
||||||
|
# быть про любой хост, и врать в лейбле не стоит (см. инвариант host
|
||||||
|
# у alert:app в infra.yml).
|
||||||
|
- receiver: telegram-clients
|
||||||
|
matchers:
|
||||||
|
- alertname = "AlertFiringTooLong"
|
||||||
|
group_wait: 10s
|
||||||
|
repeat_interval: 30m
|
||||||
|
|
||||||
# Прочее критичное — инфраструктура, клиенты пока не затронуты.
|
# Прочее критичное — инфраструктура, клиенты пока не затронуты.
|
||||||
- receiver: telegram
|
- receiver: telegram
|
||||||
matchers:
|
matchers:
|
||||||
|
|
@ -117,14 +152,17 @@ ${METRICS_TELEGRAM_INFRA_TOPIC_LINE}
|
||||||
- url: "http://alert-ack:8080/alertmanager"
|
- url: "http://alert-ack:8080/alertmanager"
|
||||||
send_resolved: true
|
send_resolved: true
|
||||||
|
|
||||||
- name: telegram-heartbeat
|
# Внешний deadman-приёмник вместо Telegram — см. комментарий у маршрута
|
||||||
telegram_configs:
|
# Watchdog выше. Блок целиком (не значение) подставляется деплоем в
|
||||||
- bot_token: "${METRICS_TELEGRAM_BOT_TOKEN}"
|
# METRICS_WATCHDOG_PING_BLOCK — тот же приём, что у
|
||||||
chat_id: ${METRICS_TELEGRAM_CHAT_ID}
|
# METRICS_TELEGRAM_INFRA_TOPIC_LINE, и по той же причине: envsubst не умеет
|
||||||
${METRICS_TELEGRAM_INFRA_TOPIC_LINE}
|
# условий. Секрет ещё не заведён на хосте — деплой в этом случае подставит
|
||||||
api_url: "https://api.telegram.org"
|
# ПУСТУЮ строку, и receiver останется без единого `*_configs`. Это валидный
|
||||||
parse_mode: HTML
|
# Alertmanager-конфиг: приёмник без конфигов просто молча отбрасывает
|
||||||
send_resolved: false
|
# уведомление, амtool его пропускает. Одинарная подстановка ЗНАЧЕНИЯ url
|
||||||
message: |
|
# (переменная-URL напрямую внутри готового ключа `url:`) сюда не годится:
|
||||||
⚪ <b>Мониторинг жив</b> — сторож отчитался, канал доставки работает.
|
# непустой ключ с пустым значением или литералом плейсхолдера амtool валит
|
||||||
Если это сообщение перестало приходить дважды подряд, замолчал сам мониторинг.
|
# целиком (`unsupported scheme ""`), а с этим — весь Alertmanager, не
|
||||||
|
# только Watchdog.
|
||||||
|
- name: watchdog-ping
|
||||||
|
${METRICS_WATCHDOG_PING_BLOCK}
|
||||||
|
|
|
||||||
|
|
@ -63,40 +63,85 @@ groups:
|
||||||
- name: host
|
- name: host
|
||||||
interval: 60s
|
interval: 60s
|
||||||
rules:
|
rules:
|
||||||
# Диск. На Beget уже был случай, когда занято 79 % и никто не смотрел;
|
# Диск. Инцидент 16.09: кэш архивов Forgejo ел 2 ГБ/ч на диске 145 ГБ
|
||||||
# порог 85 % даёт запас на реакцию, а не сообщает о свершившемся факте.
|
# (host=infra). Проценты на дисках разного размера значат разное —
|
||||||
|
# host=apps держит /-раздел ~910 ГБ, те же 85 % там это ещё ~135 ГБ
|
||||||
|
# запаса, а на infra (145 ГБ) — почти ничего. Порог переведён в
|
||||||
|
# АБСОЛЮТНЫЕ ГБ свободного места: он одинаково осмыслен на любом диске,
|
||||||
|
# потому что напрямую отвечает на вопрос «сколько времени есть до нуля
|
||||||
|
# при текущей скорости утечки», а не «какая доля занята».
|
||||||
|
#
|
||||||
|
# `mountpoint="/"` — намеренное сужение с прежнего «все fstype кроме
|
||||||
|
# tmpfs/overlay»: /boot и /boot/efi по обоим хостам меньше 1 ГБ
|
||||||
|
# целиком, с порогом 30/15 ГБ они бы горели ПОСТОЯННО (проверено живыми
|
||||||
|
# метриками 16.09 — infra:/boot/efi 98 МБ, apps:/boot 553 МБ). Данные
|
||||||
|
# приложения и Postgres лежат на "/", туда и целится алерт.
|
||||||
- alert: DiskSpaceLow
|
- alert: DiskSpaceLow
|
||||||
expr: |
|
expr: |
|
||||||
(1 - node_filesystem_avail_bytes{fstype!~"tmpfs|overlay"}
|
node_filesystem_avail_bytes{fstype!~"tmpfs|overlay", mountpoint="/"}
|
||||||
/ node_filesystem_size_bytes{fstype!~"tmpfs|overlay"}) > 0.85
|
/ 1073741824 < 30
|
||||||
for: 15m
|
for: 15m
|
||||||
labels:
|
labels:
|
||||||
severity: warning
|
severity: warning
|
||||||
annotations:
|
annotations:
|
||||||
summary: "Диск занят больше 85 %"
|
summary: "Свободно на диске меньше 30 ГБ"
|
||||||
description: "{{ $labels.host }} {{ $labels.mountpoint }}: занято {{ $value | humanizePercentage }}."
|
description: "{{ $labels.host }} {{ $labels.mountpoint }}: свободно {{ printf \"%.1f\" $value }} ГБ."
|
||||||
|
|
||||||
- alert: DiskSpaceCritical
|
- alert: DiskSpaceCritical
|
||||||
expr: |
|
expr: |
|
||||||
(1 - node_filesystem_avail_bytes{fstype!~"tmpfs|overlay"}
|
node_filesystem_avail_bytes{fstype!~"tmpfs|overlay", mountpoint="/"}
|
||||||
/ node_filesystem_size_bytes{fstype!~"tmpfs|overlay"}) > 0.93
|
/ 1073741824 < 15
|
||||||
for: 5m
|
for: 5m
|
||||||
labels:
|
labels:
|
||||||
severity: critical
|
severity: critical
|
||||||
annotations:
|
annotations:
|
||||||
summary: "Диск почти кончился"
|
summary: "Свободно на диске меньше 15 ГБ"
|
||||||
description: "{{ $labels.host }} {{ $labels.mountpoint }}: занято {{ $value | humanizePercentage }}. Postgres при заполнении диска останавливается."
|
description: "{{ $labels.host }} {{ $labels.mountpoint }}: свободно {{ printf \"%.1f\" $value }} ГБ. Postgres при заполнении диска останавливается."
|
||||||
|
|
||||||
# Прогноз важнее порога: он ловит утечку до того, как она упрётся в стену.
|
# Скорость, а не уровень. Ловит именно ту утечку, что была 16.09: пока
|
||||||
|
# DiskSpaceLow/Critical ещё не сработали (места вагон), но оно тает
|
||||||
|
# быстрее ~1 ГБ/ч устойчиво — это уже течь, а не органический рост.
|
||||||
|
# `deriv()` — линейная регрессия по 15-минутному окну (не мгновенная
|
||||||
|
# разница двух точек, ту дёргает шум). Знак минус спереди и деление
|
||||||
|
# переносят результат из "Б/с, отрицательное при убыли" в "ГБ/ч,
|
||||||
|
# положительное когда тает" — чтобы $value в тексте читался нормально
|
||||||
|
# (не "-2.1 ГБ/ч"). `for: 30m` поверх 15-минутного окна регрессии
|
||||||
|
# требует НЕПРЕРЫВНОЙ убыли около часа, чтобы не будить на разовый
|
||||||
|
# скачок (бэкап, ротация логов).
|
||||||
|
- alert: DiskSpaceDepletingFast
|
||||||
|
expr: |
|
||||||
|
-deriv(node_filesystem_avail_bytes{fstype!~"tmpfs|overlay", mountpoint="/"}[15m])
|
||||||
|
* 3600 / 1073741824 > 1
|
||||||
|
for: 30m
|
||||||
|
labels:
|
||||||
|
severity: warning
|
||||||
|
annotations:
|
||||||
|
summary: "Диск пустеет быстрее ГБ в час"
|
||||||
|
description: "{{ $labels.host }} {{ $labels.mountpoint }}: свободное место убывает на {{ printf \"%.1f\" $value }} ГБ/час устойчиво последние ~30 минут — независимо от того, сколько места осталось сейчас."
|
||||||
|
|
||||||
|
# Прогноз важнее порога: он ловит утечку до того, как она упрётся в
|
||||||
|
# стену — но роль другая, чем у DiskSpaceDepletingFast выше. Тот ловит
|
||||||
|
# БЫСТРУЮ (>1 ГБ/ч) утечку рано, за счёт короткого 15-минутного окна.
|
||||||
|
# Этот ловит УМЕРЕННУЮ утечку (может быть медленнее 1 ГБ/ч), которую
|
||||||
|
# короткое окно не поймает, но которая всё равно приведёт к нулю в
|
||||||
|
# течение суток при текущем уровне занятости — окно 6h усредняет шум
|
||||||
|
# ценой более позднего срабатывания. На одной и той же быстрой утечке
|
||||||
|
# оба правила могут сработать (сначала это, следом то) — это
|
||||||
|
# ЗАДУМАННАЯ эскалация двумя разными сигналами (скорость сейчас →
|
||||||
|
# подтверждённый тренд на сутки), а не дублирующее письмо: тексты и
|
||||||
|
# время срабатывания разные. DiskSpaceLow/Critical выше добавляют
|
||||||
|
# третий, независимый от скорости сигнал — «места мало» само по себе,
|
||||||
|
# даже если утечки нет и ничего не тает быстро.
|
||||||
- alert: DiskWillFillIn24h
|
- alert: DiskWillFillIn24h
|
||||||
expr: |
|
expr: |
|
||||||
predict_linear(node_filesystem_avail_bytes{fstype!~"tmpfs|overlay"}[6h], 24*3600) < 0
|
(0 - predict_linear(node_filesystem_avail_bytes{fstype!~"tmpfs|overlay", mountpoint="/"}[6h], 24*3600))
|
||||||
|
/ 1073741824 > 0
|
||||||
for: 30m
|
for: 30m
|
||||||
labels:
|
labels:
|
||||||
severity: warning
|
severity: warning
|
||||||
annotations:
|
annotations:
|
||||||
summary: "По текущему темпу диск кончится за сутки"
|
summary: "По текущему темпу диск кончится за сутки"
|
||||||
description: "{{ $labels.host }} {{ $labels.mountpoint }}: экстраполяция по последним 6 часам."
|
description: "{{ $labels.host }} {{ $labels.mountpoint }}: по экстраполяции последних 6 часов через сутки не хватит ~{{ printf \"%.1f\" $value }} ГБ."
|
||||||
|
|
||||||
- alert: MemoryPressure
|
- alert: MemoryPressure
|
||||||
expr: |
|
expr: |
|
||||||
|
|
@ -442,3 +487,44 @@ groups:
|
||||||
annotations:
|
annotations:
|
||||||
summary: "WAL пишется быстрее 100 МБ/час"
|
summary: "WAL пишется быстрее 100 МБ/час"
|
||||||
description: "{{ $labels.host }} / {{ $labels.db }}: {{ $value | humanize1024 }}B/с. Стоит сверить с реальной пользовательской нагрузкой — расхождение означает лишние записи."
|
description: "{{ $labels.host }} / {{ $labels.db }}: {{ $value | humanize1024 }}B/с. Стоит сверить с реальной пользовательской нагрузкой — расхождение означает лишние записи."
|
||||||
|
|
||||||
|
# ── Эскалация ──────────────────────────────────────────────────────────────
|
||||||
|
# Симптом 17.09: DiskWillFillIn24h пришёл вовремя и утонул между Watchdog
|
||||||
|
# (10080 интервалов firing за 7 суток — горит всегда по построению) и
|
||||||
|
# NoActiveCeleryWorkers (6229 интервалов) в общей ленте; диск дошёл до 84%
|
||||||
|
# незамеченным. Alertmanager сам по длительности не эскалирует — это
|
||||||
|
# правило Prometheus поверх служебной метрики ALERTS_FOR_STATE (unix-время
|
||||||
|
# входа тревоги в pending/firing, см. Robust Perception "The
|
||||||
|
# ALERTS_FOR_STATE metric").
|
||||||
|
- name: escalation
|
||||||
|
interval: 60s
|
||||||
|
rules:
|
||||||
|
# ИСКЛЮЧЕНИЯ В `alertname!~` ОБЯЗАТЕЛЬНЫ, а не для порядка:
|
||||||
|
# - Watchdog горит всегда по построению (`vector(1)` выше) — без
|
||||||
|
# исключения это правило унаследовало бы его вечный firing и стало
|
||||||
|
# ВТОРЫМ таким сигналом, то есть тем самым шумом, который лечим.
|
||||||
|
# - Сама AlertFiringTooLong — иначе, однажды сработав, она бы никогда
|
||||||
|
# не погасла: собственная ALERTS_FOR_STATE тоже старше порога, и
|
||||||
|
# правило продлевало бы себя бесконечно.
|
||||||
|
# Что это НЕ значит: если NoActiveCeleryWorkers (её чинит параллельная
|
||||||
|
# правка, здесь не трогаем) продолжит гореть дольше 6 часов, эта
|
||||||
|
# тревога сработает сразу после мержа — это ожидаемо и верно: она
|
||||||
|
# огонь реального незакрытого инцидента, а не вечная по построению.
|
||||||
|
#
|
||||||
|
# `label_replace(..., "stuck_alertname", "$1", "alertname", "(.+)")`
|
||||||
|
# ОБЯЗАТЕЛЕН, а не косметика: `alertname` — зарезервированный лейбл,
|
||||||
|
# Prometheus молча перезаписывает его именем ЭТОГО правила
|
||||||
|
# (AlertFiringTooLong) на выходе, каким бы ни было значение в expr.
|
||||||
|
# Без копии в `stuck_alertname` текст сообщения называл бы саму себя
|
||||||
|
# виновником, а не исходную тревогу.
|
||||||
|
- alert: AlertFiringTooLong
|
||||||
|
expr: |
|
||||||
|
label_replace(
|
||||||
|
(time() - ALERTS_FOR_STATE{alertname!~"Watchdog|AlertFiringTooLong"}) > 6*3600,
|
||||||
|
"stuck_alertname", "$1", "alertname", "(.+)"
|
||||||
|
)
|
||||||
|
labels:
|
||||||
|
severity: critical
|
||||||
|
annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "{{ $labels.stuck_alertname }}{{ if $labels.host }} ({{ $labels.host }}){{ end }} непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
|
|
||||||
|
|
@ -170,3 +170,85 @@ tests:
|
||||||
exp_annotations:
|
exp_annotations:
|
||||||
summary: "Бэкенд «Меры» не отвечает"
|
summary: "Бэкенд «Меры» не отвечает"
|
||||||
description: "Агент на Poincare 5 минут не может снять /metrics с tradein-backend (up=0), либо цель пропала из скрейпа. Проверь `docker ps` и /health изнутри сети. Лэндинг meraocenka.ru может открываться из кэша и при мёртвом бэкенде — это не признак жизни."
|
description: "Агент на Poincare 5 минут не может снять /metrics с tradein-backend (up=0), либо цель пропала из скрейпа. Проверь `docker ps` и /health изнутри сети. Лэндинг meraocenka.ru может открываться из кэша и при мёртвом бэкенде — это не признак жизни."
|
||||||
|
|
||||||
|
# Эскалация по длительности. `promtool test rules` держит одну общую шкалу
|
||||||
|
# времени и TSDB на весь файл — к 6.5 часам к этому моменту «зависшими»
|
||||||
|
# (input series предыдущих сценариев кончились, но absent()-условия по ним
|
||||||
|
# продолжают гореть) оказываются и другие тестовые тревоги файла, не только
|
||||||
|
# NoActiveCeleryWorkers из этого блока. Список ниже — ровно то, что
|
||||||
|
# реально вернул promtool (проверено запуском, не придумано): 8 тревог,
|
||||||
|
# держащихся дольше 6 часов. ГЛАВНАЯ ПРОВЕРКА в этом списке — то, чего в
|
||||||
|
# нём НЕТ: ни Watchdog (горит вечно с t=0 точно так же, но исключён
|
||||||
|
# матчером), ни сама AlertFiringTooLong (иначе была бы там на восьмое
|
||||||
|
# место и продлевала бы себя бесконечно). В 3 часа — рано, эскалации
|
||||||
|
# ещё быть не должно вовсе.
|
||||||
|
- interval: 1m
|
||||||
|
input_series:
|
||||||
|
- series: 'up{job="celery",host="apps"}'
|
||||||
|
values: '1x420'
|
||||||
|
alert_rule_test:
|
||||||
|
- eval_time: 3h
|
||||||
|
alertname: AlertFiringTooLong
|
||||||
|
exp_alerts: []
|
||||||
|
- eval_time: 6h30m
|
||||||
|
alertname: AlertFiringTooLong
|
||||||
|
exp_alerts:
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: MeraBackendDown
|
||||||
|
host: apps
|
||||||
|
job: app
|
||||||
|
app: mera
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "MeraBackendDown (apps) непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: HostAgentDown
|
||||||
|
host: apps
|
||||||
|
job: node
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "HostAgentDown (apps) непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: RemoteWriteStalled
|
||||||
|
host: apps
|
||||||
|
job: node
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "RemoteWriteStalled (apps) непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: CadvisorDown
|
||||||
|
job: cadvisor
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "CadvisorDown непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: QueueExporterDown
|
||||||
|
job: redis
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "QueueExporterDown непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: TradeInBackgroundContainerMissing
|
||||||
|
name: tradein-scraper
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "TradeInBackgroundContainerMissing непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: TradeInBackgroundContainerMissing
|
||||||
|
name: tradein-tgbot
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "TradeInBackgroundContainerMissing непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
- exp_labels:
|
||||||
|
severity: critical
|
||||||
|
stuck_alertname: NoActiveCeleryWorkers
|
||||||
|
exp_annotations:
|
||||||
|
summary: "Тревога держится дольше 6 часов"
|
||||||
|
description: "NoActiveCeleryWorkers непрерывно firing больше 6 часов — похоже, её не заметили или на неё забили."
|
||||||
|
|
|
||||||
|
|
@ -544,6 +544,9 @@ async def _job_domclick_city_sweep(
|
||||||
region_code=kit_resolve_region_code(params),
|
region_code=kit_resolve_region_code(params),
|
||||||
resume_run_id=kit_pick_resume(db, run_id),
|
resume_run_id=kit_pick_resume(db, run_id),
|
||||||
cookies=cookies,
|
cookies=cookies,
|
||||||
|
watchdog_sec=(
|
||||||
|
int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,56 @@
|
||||||
|
-- 326_domclick_msk_reenable_after_incremental_save.sql
|
||||||
|
-- Вернуть домклик-свипы Москвы (77) и области (50) — причина выключения устранена.
|
||||||
|
--
|
||||||
|
-- Apply after: 325_listing_source_snapshots_change_only_comment.sql
|
||||||
|
--
|
||||||
|
-- WHY:
|
||||||
|
-- Миграция 308 выключила эти строки: свип копил лоты в памяти и сохранял их ОДНИМ
|
||||||
|
-- save_listings после всех шести корзин, поэтому снятие по watchdog теряло всё
|
||||||
|
-- собранное, а чекпоинт при этом помечал корзины пройденными. На выдаче размером с
|
||||||
|
-- Москву (≈23 690 лотов вторички против ≈6 300 у ЕКБ) снятие было гарантировано, и
|
||||||
|
-- итоговый сбор равнялся нулю навсегда.
|
||||||
|
--
|
||||||
|
-- PR #3592 это снял: save_listings зовётся из колбэка on_bucket сразу после КАЖДОЙ
|
||||||
|
-- успешной корзины, туда же переехал чекпоинт — done_buckets теперь означает
|
||||||
|
-- «собрано И сохранено». Снятие по watchdog больше не теряет собранное, а корзина
|
||||||
|
-- без сохранённых строк в чекпоинт не попадает. Тем же колбэком добавлена
|
||||||
|
-- кооперативная отмена по корзинам (раньше is_cancelled проверялся только перед
|
||||||
|
-- SERP-фазой, и повисший свип нельзя было снять три часа).
|
||||||
|
--
|
||||||
|
-- interval_days = 1, а не 3 как было: полный проход Москвы в одно окно watchdog'а
|
||||||
|
-- по-прежнему НЕ помещается (замер run 7344: 2 корзины из 6 за два часа). Механизм
|
||||||
|
-- добора — ротация стартовой корзины (start_bucket_index = run_id % 6) плюс
|
||||||
|
-- skip_buckets из чекпоинта: каждый прогон берёт корзины, которых ещё нет в
|
||||||
|
-- done_buckets. При суточном такте шесть корзин закрываются примерно за трое суток,
|
||||||
|
-- при трёхсуточном — за девять, а корпус живёт 14 суток (LISTINGS_FRESH_DAYS).
|
||||||
|
-- Девять суток на полный оборот не оставляли бы запаса на пропуски из-за
|
||||||
|
-- QRATOR-банов.
|
||||||
|
--
|
||||||
|
-- watchdog_sec НЕ задаётся намеренно. Override в коде есть (PR #3592, читается из
|
||||||
|
-- default_params обоими хендлерами), но поднимать таймаут до замера нечем
|
||||||
|
-- обосновать: с инкрементальным сохранением ранний снос перестал быть потерей, а
|
||||||
|
-- более длинный прогон дольше держит один из ДВУХ узлов provider_affinity='any'
|
||||||
|
-- (id 13 и 14), за которые конкурирует cian. Сначала смотрим реальный выход за
|
||||||
|
-- прогон, потом решаем про таймаут.
|
||||||
|
--
|
||||||
|
-- Окна не меняются: Москва 0-3, область 9-12 (разведены миграцией 307), ЕКБ 3-6.
|
||||||
|
-- Ни одно окно не содержит двух домклик-строк — это условие, за нарушение которого
|
||||||
|
-- 17.09 ЕКБ-свип run 7339 отбился 'banned' за 0 секунд.
|
||||||
|
--
|
||||||
|
-- Чекпоинт прогонов 7333/7344 сбрасывать не нужно: у обоих buckets_completed пуст,
|
||||||
|
-- они были задрейнены деплоем ещё до первой завершённой корзины.
|
||||||
|
--
|
||||||
|
-- ИДЕМПОТЕНТНОСТЬ: UPDATE ... WHERE source IN (...) — повторный прогон пишет те же
|
||||||
|
-- значения. next_run_at не трогаем: планировщик посчитает его сам по окну, а явная
|
||||||
|
-- простановка здесь разъехалась бы с реальным временем применения миграции.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
UPDATE scrape_schedules
|
||||||
|
SET enabled = true,
|
||||||
|
default_params = jsonb_set(default_params, '{interval_days}', '1'::jsonb, true)
|
||||||
|
WHERE source IN ('domclick_city_sweep_moskva', 'domclick_city_sweep_moskovskaya_oblast');
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -173,10 +173,14 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None:
|
||||||
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
|
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
|
||||||
"""Корзина без сохранённых строк не попадает в чекпоинт.
|
"""Корзина без сохранённых строк не попадает в чекпоинт.
|
||||||
|
|
||||||
Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин.
|
С инкрементальным сохранением (fix/domclick-incremental-save) save_listings зовётся из
|
||||||
Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из
|
колбэка on_bucket ПОСЛЕ каждой корзины, а чекпоинт пишется там же и только
|
||||||
корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы,
|
ПОСЛЕ успешного save. Поэтому корзина, чей save упал, в done_buckets не попадает,
|
||||||
что следующий прогон пропустит её через skip_buckets навсегда (миграция 308).
|
хотя скрейпер её фетч прошёл (_s.completed_buckets её содержит). Снятая
|
||||||
|
watchdog'ом фаза — тот же инвариант: до on_bucket она не дошла.
|
||||||
|
|
||||||
|
Отметить такую корзину пройденной значило бы, что следующий прогон пропустит её
|
||||||
|
через skip_buckets навсегда (миграция 308).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
class _Dc(_Scraper):
|
class _Dc(_Scraper):
|
||||||
|
|
@ -185,10 +189,17 @@ async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> Non
|
||||||
buckets_completed, buckets_total = 1, 6
|
buckets_completed, buckets_total = 1, 6
|
||||||
completed_buckets = ["st"] # noqa: RUF012
|
completed_buckets = ["st"] # noqa: RUF012
|
||||||
|
|
||||||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
async def fetch_city(self, **kw: Any) -> list[Any]:
|
||||||
if broken == "fetch_timeout":
|
if broken == "fetch_timeout":
|
||||||
raise TimeoutError
|
raise TimeoutError
|
||||||
return [object()]
|
# Заглушка обязана ВЫЗВАТЬ колбэк: путь сохранения переехал внутрь цикла
|
||||||
|
# по корзинам (serp.py fetch_city), и стаб, который просто возвращает лоты,
|
||||||
|
# проверял бы мёртвую ветку — save_listings не был бы вызван вовсе.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
lots = [object()]
|
||||||
|
if on_bucket is not None:
|
||||||
|
on_bucket("st", lots)
|
||||||
|
return lots
|
||||||
|
|
||||||
saved: list[int] = []
|
saved: list[int] = []
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -186,6 +186,12 @@ class _FakeFetcher:
|
||||||
def report_ban(self, reason: str) -> None:
|
def report_ban(self, reason: str) -> None:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
def request_context_reset(self) -> None:
|
||||||
|
# #3118: свип зовёт сброс тёплого контекста на каждой упавшей корзине —
|
||||||
|
# двойник обязан повторять сигнатуру настоящего фетчера, иначе он проверяет
|
||||||
|
# не поведение свипа, а собственную неполноту.
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def _no_browser(monkeypatch: pytest.MonkeyPatch) -> None:
|
def _no_browser(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
|
|
||||||
132
tradein-mvp/backend/tests/test_3118_domclick_reuse_context.py
Normal file
132
tradein-mvp/backend/tests/test_3118_domclick_reuse_context.py
Normal file
|
|
@ -0,0 +1,132 @@
|
||||||
|
"""Свип ДомКлика ходит в тёплый переиспользуемый браузер-контекст сайдкара (#3118).
|
||||||
|
|
||||||
|
Замер (см. `backend/app/tasks/domclick_detail_backfill.py:394-401`): 26 подряд
|
||||||
|
холодных фетчей = 100% QRATOR-блок, те же карточки в тёплом контексте — 5/5
|
||||||
|
примерно по 2с. Причина — `browser.new_page()` создаёт НОВЫЙ изолированный
|
||||||
|
context сайдкара на каждый `/fetch`, из-за чего живой `qrator_jsid2` (куки,
|
||||||
|
которые сайт ротирует через Set-Cookie) никогда не доживает до следующего
|
||||||
|
запроса, а якорная вкладка не выживает между фетчами (`goto(origin)` валился
|
||||||
|
таймаутом 60с, убивая всю корзину).
|
||||||
|
|
||||||
|
Три проверки:
|
||||||
|
1. `build_browser_fetcher(..., reuse_context=True)` включает флаг на фетчере,
|
||||||
|
дефолт (без параметра) — выключен (не ломаем прочие call-site'ы).
|
||||||
|
2. `DomClickScraper.fetch_city` строит фетчер именно с `reuse_context=True`.
|
||||||
|
3. На QRATOR-блоке `fetcher.request_context_reset()` вызывается ДО
|
||||||
|
`fetcher.report_ban()` — иначе сожжённый блоком тёплый контекст травит
|
||||||
|
остаток прогона тем же `qrator_jsid2`.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
import types
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
|
||||||
|
def _fetcher_config() -> types.SimpleNamespace:
|
||||||
|
return types.SimpleNamespace(
|
||||||
|
browser_http_endpoint="http://sidecar:8080",
|
||||||
|
use_proxy_pool_browser=False,
|
||||||
|
environment="test",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_browser_fetcher_reuse_context_default_is_off() -> None:
|
||||||
|
"""Без явного параметра поведение прежнее — прочие call-site'ы не меняются."""
|
||||||
|
from scraper_kit.providers._base import build_browser_fetcher
|
||||||
|
|
||||||
|
fetcher = build_browser_fetcher(_fetcher_config(), "domclick")
|
||||||
|
|
||||||
|
assert fetcher._reuse_context is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_browser_fetcher_reuse_context_true_sets_flag() -> None:
|
||||||
|
from scraper_kit.providers._base import build_browser_fetcher
|
||||||
|
|
||||||
|
fetcher = build_browser_fetcher(_fetcher_config(), "domclick", reuse_context=True)
|
||||||
|
|
||||||
|
assert fetcher._reuse_context is True
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeFetcher:
|
||||||
|
"""Двойник BrowserFetcher: фиксирует порядок reset/ban."""
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.calls: list[Any] = []
|
||||||
|
|
||||||
|
def request_context_reset(self) -> None:
|
||||||
|
self.calls.append("reset")
|
||||||
|
|
||||||
|
def report_ban(self, reason: str) -> None:
|
||||||
|
self.calls.append(("ban", reason))
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeFetcherCM:
|
||||||
|
def __init__(self, fetcher: _FakeFetcher) -> None:
|
||||||
|
self._fetcher = fetcher
|
||||||
|
|
||||||
|
async def __aenter__(self) -> _FakeFetcher:
|
||||||
|
return self._fetcher
|
||||||
|
|
||||||
|
async def __aexit__(self, *_exc: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_domclick_sweep_builds_fetcher_with_reuse_context() -> None:
|
||||||
|
"""`fetch_city` зовёт фабрику фетчера с `reuse_context=True` (не дефолтом)."""
|
||||||
|
from scraper_kit.domclick_exceptions import DomClickBlockedError
|
||||||
|
from scraper_kit.providers.domclick.serp import DomClickScraper
|
||||||
|
|
||||||
|
fetcher = _FakeFetcher()
|
||||||
|
mock_build = MagicMock(return_value=_FakeFetcherCM(fetcher))
|
||||||
|
|
||||||
|
scraper = DomClickScraper(_fetcher_config())
|
||||||
|
# Обрываем на первом же бакете — детали сбора корзины здесь не проверяются.
|
||||||
|
with (
|
||||||
|
patch("scraper_kit.providers._base.build_browser_fetcher", mock_build),
|
||||||
|
patch.object(
|
||||||
|
scraper, "_sweep_bucket", AsyncMock(side_effect=DomClickBlockedError("qrator"))
|
||||||
|
),
|
||||||
|
):
|
||||||
|
await scraper.fetch_city(city_id=66, pages=1)
|
||||||
|
|
||||||
|
assert mock_build.call_args.kwargs.get("reuse_context") is True, (
|
||||||
|
f"свип построил фетчер без reuse_context=True: {mock_build.call_args}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_domclick_sweep_resets_context_after_failed_bucket() -> None:
|
||||||
|
"""Упавшая корзина сбрасывает тёплый контекст, чтобы он не переполз в следующую.
|
||||||
|
|
||||||
|
Прод 17.09, прогон 7389: после падения 'rooms=3' три корзины подряд легли за две
|
||||||
|
минуты каждая на таймауте goto(origin) — якорная вкладка не поднималась в том же
|
||||||
|
сгоревшем контексте. Без сброса reuse_context=True превращает одну неудачу в
|
||||||
|
цепочку.
|
||||||
|
"""
|
||||||
|
from scraper_kit.providers.domclick.serp import DomClickScraper
|
||||||
|
|
||||||
|
fetcher = _FakeFetcher()
|
||||||
|
mock_build = MagicMock(return_value=_FakeFetcherCM(fetcher))
|
||||||
|
|
||||||
|
scraper = DomClickScraper(_fetcher_config())
|
||||||
|
with (
|
||||||
|
patch("scraper_kit.providers._base.build_browser_fetcher", mock_build),
|
||||||
|
patch.object(scraper, "_sweep_bucket", AsyncMock(side_effect=RuntimeError("boom"))),
|
||||||
|
):
|
||||||
|
await scraper.fetch_city(city_id=66, pages=1)
|
||||||
|
|
||||||
|
assert fetcher.calls.count("reset") == 6, (
|
||||||
|
f"сброс контекста ожидался на каждой из 6 упавших корзин: {fetcher.calls!r}"
|
||||||
|
)
|
||||||
|
assert not any(c != "reset" for c in fetcher.calls), (
|
||||||
|
f"report_ban не должен вызываться на не-QRATOR ошибке: {fetcher.calls!r}"
|
||||||
|
)
|
||||||
|
|
@ -153,7 +153,15 @@ class _CleanSweepScraper:
|
||||||
async def __aexit__(self, *_e: Any) -> None:
|
async def __aexit__(self, *_e: Any) -> None:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
async def fetch_city(self, **kw: Any) -> list[Any]:
|
||||||
|
# Чекпоинт с fix/domclick-incremental-save набирается ИЗ колбэка on_bucket, а не
|
||||||
|
# мержится в конце из completed_buckets — стаб обязан его вызвать, иначе
|
||||||
|
# проверялась бы мёртвая ветка. Лотов нет: корзина пройдена, но пустая —
|
||||||
|
# это законный случай, save_listings для неё не зовётся, чекпоинт пишется.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
if on_bucket is not None:
|
||||||
|
for bucket in self.completed_buckets:
|
||||||
|
on_bucket(bucket, [])
|
||||||
return []
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
234
tradein-mvp/backend/tests/test_domclick_incremental_save.py
Normal file
234
tradein-mvp/backend/tests/test_domclick_incremental_save.py
Normal file
|
|
@ -0,0 +1,234 @@
|
||||||
|
"""fix/domclick-incremental-save: большая выдача (Москва ≈23 690 лотов против ЕКБ
|
||||||
|
≈6 300) теряла ВСЁ собранное при снятии свипа по watchdog — сохранение было ОДНО,
|
||||||
|
в самом конце `run_domclick_city_sweep`. Прод run 7344 (`domclick_city_sweep_moskva`):
|
||||||
|
за 2ч watchdog'а (11 100с) пройдено 2 бакета из 6.
|
||||||
|
|
||||||
|
Три дефекта, три слоя тестов ниже:
|
||||||
|
1. Инкрементальное сохранение по бакетам (`fetch_city.on_bucket`) — секция 1.
|
||||||
|
2. Чекпоинт врал ("пройдено" без "сохранено") — секция 1, тест
|
||||||
|
`test_on_bucket_exception_interrupts_bucket_loop` документирует нюанс, на
|
||||||
|
котором строится фикс в pipeline.py: scraper.completed_buckets отмечает бакет
|
||||||
|
как "фетч прошёл" ДО вызова on_bucket, поэтому пайплайн больше не берёт
|
||||||
|
чекпоинт оттуда — только из факта успешного on_bucket (см. комментарий у
|
||||||
|
"Перенести счётчики" в run_domclick_city_sweep).
|
||||||
|
3. Watchdog не знал о размере выдачи — секция 2, `watchdog_sec` override.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import patch
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from scraper_kit.orchestration.pipeline import run_domclick_city_sweep
|
||||||
|
from scraper_kit.providers.domclick.serp import (
|
||||||
|
ROOM_BUCKETS,
|
||||||
|
DomClickBlockedError,
|
||||||
|
DomClickScraper,
|
||||||
|
)
|
||||||
|
|
||||||
|
PFX = "scraper_kit.orchestration.pipeline"
|
||||||
|
|
||||||
|
|
||||||
|
# ── Секция 1: DomClickScraper.fetch_city(on_bucket=...) ──────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _scraper() -> DomClickScraper:
|
||||||
|
return DomClickScraper(SimpleNamespace(scraper_proxy_url=None))
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeFetcherCtx:
|
||||||
|
async def __aenter__(self) -> SimpleNamespace:
|
||||||
|
# request_context_reset (#3118): вызывается ДО report_ban на QRATOR-блоке —
|
||||||
|
# двойник фетчера обязан его иметь, иначе AttributeError на первом же блоке.
|
||||||
|
return SimpleNamespace(
|
||||||
|
report_ban=lambda *_a, **_k: None,
|
||||||
|
request_context_reset=lambda *_a, **_k: None,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def __aexit__(self, *_exc: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
async def _run_fetch_city(
|
||||||
|
scraper: DomClickScraper,
|
||||||
|
*,
|
||||||
|
sweep_bucket: Any,
|
||||||
|
on_bucket: Any = None,
|
||||||
|
start: int = 0,
|
||||||
|
skip: set[str] | None = None,
|
||||||
|
) -> list[str]:
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_sweep_bucket", sweep_bucket),
|
||||||
|
patch(
|
||||||
|
"scraper_kit.providers._base.build_browser_fetcher",
|
||||||
|
lambda *_a, **_k: _FakeFetcherCtx(),
|
||||||
|
),
|
||||||
|
):
|
||||||
|
return await scraper.fetch_city(
|
||||||
|
city_id=1,
|
||||||
|
pages=1,
|
||||||
|
start_bucket_index=start,
|
||||||
|
skip_buckets=skip,
|
||||||
|
on_bucket=on_bucket,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_called_once_per_bucket_with_only_its_own_lots() -> None:
|
||||||
|
"""on_bucket зовётся по разу на каждый успешный бакет с лотами ИМЕННО его,
|
||||||
|
а не накопленным out_lots (главная регрессия #1 issue — было "одно сохранение
|
||||||
|
в конце")."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-a")
|
||||||
|
out_lots.append(f"{rooms}-b")
|
||||||
|
|
||||||
|
calls: list[tuple[str, list[str]]] = []
|
||||||
|
s = _scraper()
|
||||||
|
result = await _run_fetch_city(
|
||||||
|
s, sweep_bucket=_fake_sweep, on_bucket=lambda n, lots: calls.append((n, list(lots)))
|
||||||
|
)
|
||||||
|
|
||||||
|
assert [c[0] for c in calls] == list(ROOM_BUCKETS)
|
||||||
|
for bucket, lots in calls:
|
||||||
|
assert lots == [f"{bucket}-a", f"{bucket}-b"], (bucket, lots, "получил чужие лоты")
|
||||||
|
assert len(result) == len(ROOM_BUCKETS) * 2
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_skips_failed_and_blocked_buckets() -> None:
|
||||||
|
"""Битый бакет (generic Exception) и бакет с QRATOR-блоком не отдают лоты в
|
||||||
|
on_bucket — там ничего не собрано/не гарантированно собрано."""
|
||||||
|
|
||||||
|
async def _fake_sweep_fail_middle(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
if rooms == "1":
|
||||||
|
raise ValueError("boom")
|
||||||
|
|
||||||
|
calls_a: list[str] = []
|
||||||
|
s_a = _scraper()
|
||||||
|
await _run_fetch_city(
|
||||||
|
s_a, sweep_bucket=_fake_sweep_fail_middle, on_bucket=lambda n, _lots: calls_a.append(n)
|
||||||
|
)
|
||||||
|
assert "1" not in calls_a
|
||||||
|
# continue идёт дальше — остальные бакеты всё равно получают on_bucket.
|
||||||
|
assert calls_a == [b for b in ROOM_BUCKETS if b != "1"], calls_a
|
||||||
|
|
||||||
|
async def _fake_sweep_block(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
if rooms == "1":
|
||||||
|
raise DomClickBlockedError("QRATOR")
|
||||||
|
|
||||||
|
calls_b: list[str] = []
|
||||||
|
s_b = _scraper()
|
||||||
|
await _run_fetch_city(
|
||||||
|
s_b, sweep_bucket=_fake_sweep_block, on_bucket=lambda n, _lots: calls_b.append(n)
|
||||||
|
)
|
||||||
|
# break останавливает обход целиком — после блока ни один бакет не пробуется.
|
||||||
|
assert calls_b == ["st"], calls_b
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_exception_interrupts_bucket_loop() -> None:
|
||||||
|
"""Исключение из on_bucket (канал кооперативной отмены) прерывает цикл по
|
||||||
|
ROOM_BUCKETS целиком — остальные бакеты не идут."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
|
||||||
|
visited: list[str] = []
|
||||||
|
|
||||||
|
def _on_bucket(name: str, lots: list[str]) -> None:
|
||||||
|
visited.append(name)
|
||||||
|
if name == "1":
|
||||||
|
raise RuntimeError("cancelled")
|
||||||
|
|
||||||
|
s = _scraper()
|
||||||
|
with pytest.raises(RuntimeError, match="cancelled"):
|
||||||
|
await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=_on_bucket)
|
||||||
|
|
||||||
|
assert visited == ["st", "1"], visited
|
||||||
|
# Скрейпер-уровень успел зафетчить оба бакета ДО того, как on_bucket поднял
|
||||||
|
# исключение — на этом нюансе строится фикс чекпоинта в pipeline.py: pipeline
|
||||||
|
# больше не берёт done_buckets из scraper.completed_buckets, только из факта
|
||||||
|
# успешного on_bucket (см. run_domclick_city_sweep._on_bucket/_checkpoint).
|
||||||
|
assert s.completed_buckets == ["st", "1"], s.completed_buckets
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_none_preserves_return_value() -> None:
|
||||||
|
"""on_bucket=None → поведение прежнее, байт-в-байт: лоты возвращаются из
|
||||||
|
fetch_city как раньше, никаких промежуточных вызовов."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-only")
|
||||||
|
|
||||||
|
s = _scraper()
|
||||||
|
result = await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=None)
|
||||||
|
|
||||||
|
assert result == [f"{b}-only" for b in ROOM_BUCKETS]
|
||||||
|
assert s.completed_buckets == list(ROOM_BUCKETS)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Секция 2: run_domclick_city_sweep(watchdog_sec=...) ──────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _NoOpRuns:
|
||||||
|
"""Достаточно методов, чтобы SERP-фаза дошла до asyncio.wait_for и честно
|
||||||
|
финализировалась после симулированного TimeoutError — значения не важны,
|
||||||
|
важен ТОЛЬКО timeout, с которым позвали wait_for."""
|
||||||
|
|
||||||
|
def is_cancelled(self, db: Any, run_id: int) -> bool:
|
||||||
|
return False
|
||||||
|
|
||||||
|
def update_heartbeat(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_banned(
|
||||||
|
self, db: Any, run_id: int, error: str, counters: dict[str, Any], **kw: Any
|
||||||
|
) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
async def _drive_and_capture_timeout(**kwargs: Any) -> float:
|
||||||
|
captured: dict[str, float] = {}
|
||||||
|
|
||||||
|
async def _fake_wait_for(coro: Any, timeout: float) -> None:
|
||||||
|
captured["timeout"] = timeout
|
||||||
|
coro.close()
|
||||||
|
raise TimeoutError()
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(f"{PFX}.asyncio.wait_for", _fake_wait_for),
|
||||||
|
patch(f"{PFX}.runs", _NoOpRuns()),
|
||||||
|
):
|
||||||
|
await run_domclick_city_sweep(
|
||||||
|
object(), # type: ignore[arg-type]
|
||||||
|
config=SimpleNamespace(browser_http_endpoint="http://x:9000"),
|
||||||
|
matcher=object(),
|
||||||
|
run_id=1,
|
||||||
|
city_id=4,
|
||||||
|
**kwargs,
|
||||||
|
)
|
||||||
|
return captured["timeout"]
|
||||||
|
|
||||||
|
|
||||||
|
async def test_watchdog_sec_default_formula_matches_prod_11100() -> None:
|
||||||
|
"""watchdog_sec=None (дефолт) → прежняя формула байт-в-байт. pages=100,
|
||||||
|
delay=6.0 — те же параметры, что дали 11 100с в run 7344 (Москва)."""
|
||||||
|
timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0)
|
||||||
|
assert timeout == 11100
|
||||||
|
|
||||||
|
|
||||||
|
async def test_watchdog_sec_override_bypasses_formula() -> None:
|
||||||
|
"""watchdog_sec задан явно → формула не считается вовсе, идёт ровно override
|
||||||
|
— даже с теми же pages/delay, что в тесте формулы выше."""
|
||||||
|
timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0, watchdog_sec=777)
|
||||||
|
assert timeout == 777
|
||||||
|
|
@ -73,7 +73,11 @@ def _make_recorder() -> tuple[type, list[dict[str, Any]]]:
|
||||||
proxy_provider: object | None = None,
|
proxy_provider: object | None = None,
|
||||||
use_pool: bool = False,
|
use_pool: bool = False,
|
||||||
environment: str = "dev",
|
environment: str = "dev",
|
||||||
|
reuse_context: bool = False,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
# reuse_context (#3118) обязан быть в сигнатуре двойника: фабрика
|
||||||
|
# build_browser_fetcher передаёт его ВСЕГДА, и двойник без него падал бы
|
||||||
|
# TypeError на каждом вызывающем, а не проверял то, ради чего написан.
|
||||||
calls.append(
|
calls.append(
|
||||||
{
|
{
|
||||||
"source": source,
|
"source": source,
|
||||||
|
|
@ -81,6 +85,7 @@ def _make_recorder() -> tuple[type, list[dict[str, Any]]]:
|
||||||
"proxy_provider": proxy_provider,
|
"proxy_provider": proxy_provider,
|
||||||
"use_pool": use_pool,
|
"use_pool": use_pool,
|
||||||
"environment": environment,
|
"environment": environment,
|
||||||
|
"reuse_context": reuse_context,
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -411,8 +411,19 @@ async def _drive_domclick(
|
||||||
recorder = _RunsRecorder()
|
recorder = _RunsRecorder()
|
||||||
db = MagicMock()
|
db = MagicMock()
|
||||||
lots = [MagicMock() for _ in range(lots_n)]
|
lots = [MagicMock() for _ in range(lots_n)]
|
||||||
|
|
||||||
|
async def _fetch_city(**kw: Any) -> list[Any]:
|
||||||
|
# fix/domclick-incremental-save: save_listings переехал внутрь цикла по корзинам и
|
||||||
|
# зовётся из колбэка on_bucket. Стаб, который просто возвращает лоты, не
|
||||||
|
# вызвал бы сохранение вовсе — фикстура проверяла бы мёртвую ветку.
|
||||||
|
# Одна корзина со всеми лотами: ровно один save_listings, как и было.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
if on_bucket is not None and lots:
|
||||||
|
on_bucket(ROOM_BUCKETS[0], lots)
|
||||||
|
return lots
|
||||||
|
|
||||||
scraper = _ctx_scraper(
|
scraper = _ctx_scraper(
|
||||||
fetch_city=AsyncMock(return_value=lots),
|
fetch_city=_fetch_city,
|
||||||
blocked=blocked,
|
blocked=blocked,
|
||||||
geo_filtered=0,
|
geo_filtered=0,
|
||||||
fetch_errors=fetch_errors,
|
fetch_errors=fetch_errors,
|
||||||
|
|
|
||||||
|
|
@ -4868,6 +4868,7 @@ async def run_domclick_city_sweep(
|
||||||
region_code: int = DEFAULT_REGION_CODE,
|
region_code: int = DEFAULT_REGION_CODE,
|
||||||
resume_run_id: int | None = None,
|
resume_run_id: int | None = None,
|
||||||
cookies: dict[str, str] | None = None,
|
cookies: dict[str, str] | None = None,
|
||||||
|
watchdog_sec: int | None = None,
|
||||||
) -> DomClickCitySweepCounters:
|
) -> DomClickCitySweepCounters:
|
||||||
"""DomClick citywide sweep через BFF JSON API.
|
"""DomClick citywide sweep через BFF JSON API.
|
||||||
|
|
||||||
|
|
@ -4896,6 +4897,15 @@ async def run_domclick_city_sweep(
|
||||||
зависли на challenge). None (дефолт, сессии в БД нет/протухла) — прежнее
|
зависли на challenge). None (дефолт, сессии в БД нет/протухла) — прежнее
|
||||||
поведение, без инъекции.
|
поведение, без инъекции.
|
||||||
|
|
||||||
|
watchdog_sec (incremental-save): явный override расчётной формулы watchdog'а. Формула
|
||||||
|
ниже (buckets × pages × per_fetch + budget) не знает про бисекцию по цене
|
||||||
|
(ДомКлик режет offset на 2000, каждый лист пагинируется отдельно отдельным
|
||||||
|
деревом сплитов) и на больших городах (Москва ≈23 690 лотов против ЕКБ
|
||||||
|
≈6 300) занижена в разы — прод run 7344: за 2ч watchdog'а (11 100с) пройдено
|
||||||
|
2 бакета из 6. С инкрементальным сохранением (on_bucket, см. ниже) ранний
|
||||||
|
снос по watchdog больше не теряет собранное, поэтому вместо более точной
|
||||||
|
оценки — простой override: None (дефолт) — прежняя формула байт-в-байт.
|
||||||
|
|
||||||
Возвращает DomClickCitySweepCounters.
|
Возвращает DomClickCitySweepCounters.
|
||||||
"""
|
"""
|
||||||
# Гео-скоуп свипа больше не зашит в ЕКБ: его задаёт профиль региона (address
|
# Гео-скоуп свипа больше не зашит в ЕКБ: его задаёт профиль региона (address
|
||||||
|
|
@ -4909,11 +4919,17 @@ async def run_domclick_city_sweep(
|
||||||
_resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0
|
_resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0
|
||||||
counters = DomClickCitySweepCounters()
|
counters = DomClickCitySweepCounters()
|
||||||
|
|
||||||
# Watchdog: 6 buckets × pages × per_fetch + budget.
|
# Watchdog: 6 buckets × pages × per_fetch + budget. watchdog_sec (incremental-save) — явный
|
||||||
|
# override для больших городов, где формула занижена (см. докстринг выше).
|
||||||
_num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages)
|
_num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages)
|
||||||
|
if watchdog_sec is not None:
|
||||||
|
_sweep_timeout = int(watchdog_sec)
|
||||||
|
else:
|
||||||
_sweep_timeout = max(
|
_sweep_timeout = max(
|
||||||
ANCHOR_TIMEOUT_SEC,
|
ANCHOR_TIMEOUT_SEC,
|
||||||
int(_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S),
|
int(
|
||||||
|
_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
||||||
|
|
@ -4991,14 +5007,68 @@ async def run_domclick_city_sweep(
|
||||||
_sweep_timeout,
|
_sweep_timeout,
|
||||||
)
|
)
|
||||||
|
|
||||||
lots: list[ScrapedLot] = []
|
# incremental-save: сохранение инкрементальное — save_listings зовётся из _on_bucket ПОСЛЕ
|
||||||
# #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому
|
# КАЖДОГО room-бакета, а не одним save_listings в самом конце. Раньше снятие фазы
|
||||||
# корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже).
|
# по watchdog'у (asyncio.wait_for TimeoutError) теряло ВСЁ собранное — на большой
|
||||||
_saved = False
|
# выдаче (Москва ≈23 690 лотов vs ЕКБ ≈6 300) формула watchdog'а систематически
|
||||||
|
# не укладывалась в отведённое время (прод run 7344: 2 бакета из 6 за 2ч).
|
||||||
|
_cancel_reason: str | None = None
|
||||||
|
|
||||||
|
def _on_bucket(bucket_key: str, bucket_lots: list[ScrapedLot]) -> None:
|
||||||
|
"""Инкрементальный save сразу после того, как room-бакет отфетчился.
|
||||||
|
|
||||||
|
fetch_city (serp.py) зовёт колбэк ВНЕ try/except конкретного бакета —
|
||||||
|
исключение отсюда прерывает обход ROOM_BUCKETS целиком (канал кооперативной
|
||||||
|
отмены/SIGTERM-дрейна), а не проглатывается generic except'ом бакета.
|
||||||
|
Sentinel-приём (RuntimeError("cancelled")/("shutdown")) — тот же, что в
|
||||||
|
run_cian_full_load._on_bucket (#1182 Phase 3a).
|
||||||
|
"""
|
||||||
|
nonlocal _checkpoint
|
||||||
|
if runs.is_cancelled(db, run_id):
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: cancel detected in on_bucket (%s)",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
)
|
||||||
|
raise RuntimeError("cancelled")
|
||||||
|
elif shutdown_requested():
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: SIGTERM-drain — stopping at bucket %s",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
)
|
||||||
|
raise RuntimeError("shutdown")
|
||||||
|
if bucket_lots:
|
||||||
|
# Имя города — из профиля региона, а не из сравнения с vestigial
|
||||||
|
# city_id: гео-скоп задаёт регион, он же знает, какой город штамповать.
|
||||||
|
# У области (50) city_name=None — одного города нет, угадывать нечего.
|
||||||
|
inserted, updated = save_listings(
|
||||||
|
db,
|
||||||
|
bucket_lots,
|
||||||
|
matcher=matcher,
|
||||||
|
region_code=region_code,
|
||||||
|
run_id=run_id,
|
||||||
|
city=_geo_profile.city_name,
|
||||||
|
)
|
||||||
|
counters.lots_fetched += len(bucket_lots)
|
||||||
|
counters.lots_inserted += inserted
|
||||||
|
counters.lots_updated += updated
|
||||||
|
# #3118, теперь на уровне бакета: done_buckets означает "собрано И
|
||||||
|
# сохранено" — чекпоинт пишем ТОЛЬКО пройдя cancel/shutdown-гейт выше и
|
||||||
|
# save_listings этого бакета, не из scraper.completed_buckets (см.
|
||||||
|
# комментарий у "Перенести счётчики" ниже — там раньше был баг #2 issue).
|
||||||
|
_checkpoint = sorted(set(_checkpoint) | {bucket_key})
|
||||||
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: bucket %s saved lots=%d total=%d",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
len(bucket_lots),
|
||||||
|
counters.lots_fetched,
|
||||||
|
)
|
||||||
|
|
||||||
async def _domclick_phase() -> None:
|
async def _domclick_phase() -> None:
|
||||||
"""Единственная citywide-фаза: fetch_city + save."""
|
"""Единственная citywide-фаза: fetch_city с инкрементальным save по бакетам."""
|
||||||
nonlocal lots, _saved
|
|
||||||
async with DomClickScraper(
|
async with DomClickScraper(
|
||||||
config,
|
config,
|
||||||
proxy_provider=proxy_provider,
|
proxy_provider=proxy_provider,
|
||||||
|
|
@ -5021,36 +5091,21 @@ async def run_domclick_city_sweep(
|
||||||
# источники и растёт неравномерно, — но за 30 суток каждая корзина
|
# источники и растёт неравномерно, — но за 30 суток каждая корзина
|
||||||
# получает порядка пяти стартов, чего достаточно для критерия приёмки
|
# получает порядка пяти стартов, чего достаточно для критерия приёмки
|
||||||
# «объявления с rooms >= 2 появились».
|
# «объявления с rooms >= 2 появились».
|
||||||
lots = await _scraper.fetch_city(
|
await _scraper.fetch_city(
|
||||||
city_id=city_id,
|
city_id=city_id,
|
||||||
rooms=rooms,
|
rooms=rooms,
|
||||||
pages=pages,
|
pages=pages,
|
||||||
start_bucket_index=run_id % len(ROOM_BUCKETS),
|
start_bucket_index=run_id % len(ROOM_BUCKETS),
|
||||||
skip_buckets=skip_buckets or None,
|
skip_buckets=skip_buckets or None,
|
||||||
|
on_bucket=_on_bucket,
|
||||||
)
|
)
|
||||||
counters.lots_fetched += len(lots)
|
|
||||||
if lots:
|
|
||||||
# Имя города — из профиля региона, а не из сравнения с vestigial
|
|
||||||
# city_id: гео-скоп задаёт регион, он же знает, какой город штамповать.
|
|
||||||
# У области (50) city_name=None — одного города нет, угадывать нечего.
|
|
||||||
_dc_city = _geo_profile.city_name
|
|
||||||
inserted, updated = save_listings(
|
|
||||||
db,
|
|
||||||
lots,
|
|
||||||
matcher=matcher,
|
|
||||||
region_code=region_code,
|
|
||||||
run_id=run_id,
|
|
||||||
city=_dc_city,
|
|
||||||
)
|
|
||||||
counters.lots_inserted += inserted
|
|
||||||
counters.lots_updated += updated
|
|
||||||
_saved = True
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
|
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results",
|
"domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results "
|
||||||
|
"(incremental-save: уже сохранены инкрементально, бакет за бакетом)",
|
||||||
run_id,
|
run_id,
|
||||||
_sweep_timeout,
|
_sweep_timeout,
|
||||||
)
|
)
|
||||||
|
|
@ -5080,6 +5135,16 @@ async def run_domclick_city_sweep(
|
||||||
ban_kind=ban_kind_of_exception(exc),
|
ban_kind=ban_kind_of_exception(exc),
|
||||||
)
|
)
|
||||||
return counters
|
return counters
|
||||||
|
except RuntimeError as exc:
|
||||||
|
# on_bucket кидает RuntimeError("cancelled") при кооперативной отмене,
|
||||||
|
# RuntimeError("shutdown") при SIGTERM-дрейне — тот же sentinel-приём, что в
|
||||||
|
# run_cian_full_load._on_bucket (#1182 Phase 3a). Прочие RuntimeError —
|
||||||
|
# обычная поломка фазы, ведём себя как под generic except ниже.
|
||||||
|
if str(exc) in ("cancelled", "shutdown"):
|
||||||
|
_cancel_reason = str(exc)
|
||||||
|
else:
|
||||||
|
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
||||||
|
counters.errors_count += 1
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
||||||
counters.errors_count += 1
|
counters.errors_count += 1
|
||||||
|
|
@ -5095,15 +5160,13 @@ async def run_domclick_city_sweep(
|
||||||
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
|
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
|
||||||
counters.buckets_completed = _s.buckets_completed
|
counters.buckets_completed = _s.buckets_completed
|
||||||
counters.buckets_total = _s.buckets_total
|
counters.buckets_total = _s.buckets_total
|
||||||
# #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне.
|
# incremental-save: _checkpoint СЮДА больше не мержится из _s.completed_buckets — он
|
||||||
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
|
# пишется инкрементально внутри _on_bucket, СРАЗУ после save_listings этого
|
||||||
# затирают, и оборванный болезнью финализации прогон его не теряет.
|
# бакета. _s.completed_buckets на уровне скрейпера отмечает бакет как "фетч
|
||||||
# #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза
|
# прошёл" ДО вызова on_bucket (см. serp.py): если on_bucket поймал
|
||||||
# не дошла до save_listings — у её корзин в БД ноль строк, а отметка
|
# cancel/shutdown ДО save для этого самого бакета, _s.completed_buckets
|
||||||
# «пройдена» заставила бы следующий прогон пропустить их навсегда
|
# включил бы его, а _checkpoint — честно нет (лоты не сохранены). Мердж
|
||||||
# (механизм разобран в миграции 308, из-за него выключены свипы 77/50).
|
# отсюда воспроизвёл бы старый баг — чекпоинт врёт про несохранённые бакеты.
|
||||||
if _saved:
|
|
||||||
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
|
|
||||||
runs.update_heartbeat(db, run_id, _payload())
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
counters.bucket_start_index = _s.bucket_start_index
|
counters.bucket_start_index = _s.bucket_start_index
|
||||||
|
|
||||||
|
|
@ -5111,6 +5174,32 @@ async def run_domclick_city_sweep(
|
||||||
counters.pages_fetched = _num_fetches
|
counters.pages_fetched = _num_fetches
|
||||||
runs.update_heartbeat(db, run_id, _payload())
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
|
|
||||||
|
# incremental-save: кооперативная отмена/SIGTERM-дрейн, пойманные в on_bucket — партиал уже
|
||||||
|
# сохранён инкрементально, финализируем как run_cian_full_load (mark_done partial,
|
||||||
|
# не honest-status ниже: обрыв тут known-signal, а не "прогон не доделал сам").
|
||||||
|
if _cancel_reason == "cancelled":
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: cancelled — partial results lots=%d (ins=%d/upd=%d)",
|
||||||
|
run_id,
|
||||||
|
counters.lots_fetched,
|
||||||
|
counters.lots_inserted,
|
||||||
|
counters.lots_updated,
|
||||||
|
)
|
||||||
|
runs.mark_done(db, run_id, _payload())
|
||||||
|
return counters
|
||||||
|
if _cancel_reason == "shutdown":
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: SIGTERM-drain — partial results lots=%d (ins=%d/upd=%d)",
|
||||||
|
run_id,
|
||||||
|
counters.lots_fetched,
|
||||||
|
counters.lots_inserted,
|
||||||
|
counters.lots_updated,
|
||||||
|
)
|
||||||
|
_drain = {**_payload(), "interrupted": 1}
|
||||||
|
runs.update_heartbeat(db, run_id, _drain)
|
||||||
|
runs.mark_done(db, run_id, _drain)
|
||||||
|
return counters
|
||||||
|
|
||||||
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
||||||
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
||||||
# отличается от cian/yandex (#2625/#2642): там независимые anchor'ы, и
|
# отличается от cian/yandex (#2625/#2642): там независимые anchor'ы, и
|
||||||
|
|
|
||||||
|
|
@ -1270,6 +1270,9 @@ async def _job_domclick_city_sweep(
|
||||||
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
||||||
region_code=_resolve_region_code(params),
|
region_code=_resolve_region_code(params),
|
||||||
resume_run_id=_pick_resume(db, run_id),
|
resume_run_id=_pick_resume(db, run_id),
|
||||||
|
watchdog_sec=(
|
||||||
|
int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -192,6 +192,7 @@ def build_browser_fetcher(
|
||||||
*,
|
*,
|
||||||
proxy_provider: ProxyProvider | None = None,
|
proxy_provider: ProxyProvider | None = None,
|
||||||
fetch_timeout_s: float | None = None,
|
fetch_timeout_s: float | None = None,
|
||||||
|
reuse_context: bool = False,
|
||||||
) -> BrowserFetcher:
|
) -> BrowserFetcher:
|
||||||
"""Собрать `BrowserFetcher` с `config: ScraperConfig` **mandatory**.
|
"""Собрать `BrowserFetcher` с `config: ScraperConfig` **mandatory**.
|
||||||
|
|
||||||
|
|
@ -221,6 +222,11 @@ def build_browser_fetcher(
|
||||||
(120s). Явный таймаут передаёт ровно один call-site — `yandex/serp.py` (30s);
|
(120s). Явный таймаут передаёт ровно один call-site — `yandex/serp.py` (30s);
|
||||||
`yandex/newbuilding.py` идёт на дефолтных 120s.
|
`yandex/newbuilding.py` идёт на дефолтных 120s.
|
||||||
|
|
||||||
|
`reuse_context=False` (дефолт) — сохраняет прежнее поведение всех вызывающих:
|
||||||
|
холодный контекст на каждый /fetch. `reuse_context=True` (#3118, свип домклика)
|
||||||
|
держит один сайдкар-контекст на весь прогон вместо нового камуфокса на каждый
|
||||||
|
фетч.
|
||||||
|
|
||||||
`environment=getattr(config, "environment", "dev")` (#2616 шаг 1) — прокидывается в
|
`environment=getattr(config, "environment", "dev")` (#2616 шаг 1) — прокидывается в
|
||||||
`BrowserFetcher._pool_proxy`: пул пуст/сломан + прод → отказ вместо мёртвого
|
`BrowserFetcher._pool_proxy`: пул пуст/сломан + прод → отказ вместо мёртвого
|
||||||
env-прокси. `getattr` с дефолтом "dev" — минимальные ScraperConfig-заглушки без поля
|
env-прокси. `getattr` с дефолтом "dev" — минимальные ScraperConfig-заглушки без поля
|
||||||
|
|
@ -234,6 +240,7 @@ def build_browser_fetcher(
|
||||||
proxy_provider=proxy_provider,
|
proxy_provider=proxy_provider,
|
||||||
use_pool=config.use_proxy_pool_browser,
|
use_pool=config.use_proxy_pool_browser,
|
||||||
environment=environment,
|
environment=environment,
|
||||||
|
reuse_context=reuse_context,
|
||||||
)
|
)
|
||||||
return BrowserFetcher(
|
return BrowserFetcher(
|
||||||
source=source,
|
source=source,
|
||||||
|
|
@ -242,6 +249,7 @@ def build_browser_fetcher(
|
||||||
proxy_provider=proxy_provider,
|
proxy_provider=proxy_provider,
|
||||||
use_pool=config.use_proxy_pool_browser,
|
use_pool=config.use_proxy_pool_browser,
|
||||||
environment=environment,
|
environment=environment,
|
||||||
|
reuse_context=reuse_context,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -459,6 +459,7 @@ class DomClickScraper(BaseScraper):
|
||||||
pages: int = 100,
|
pages: int = 100,
|
||||||
start_bucket_index: int = 0,
|
start_bucket_index: int = 0,
|
||||||
skip_buckets: set[str] | None = None,
|
skip_buckets: set[str] | None = None,
|
||||||
|
on_bucket: Callable[[str, list[ScrapedLot]], None] | None = None,
|
||||||
) -> list[ScrapedLot]:
|
) -> list[ScrapedLot]:
|
||||||
"""Citywide sweep через BFF JSON API.
|
"""Citywide sweep через BFF JSON API.
|
||||||
|
|
||||||
|
|
@ -492,6 +493,17 @@ class DomClickScraper(BaseScraper):
|
||||||
buckets_total при этом = числу корзин В ЭТОМ прогоне (без
|
buckets_total при этом = числу корзин В ЭТОМ прогоне (без
|
||||||
скипнутых) — иначе honest-status читал бы возобновлённый прогон
|
скипнутых) — иначе honest-status читал бы возобновлённый прогон
|
||||||
как вечно-частичный.
|
как вечно-частичный.
|
||||||
|
on_bucket: колбэк инкрементального сохранения — большая выдача теряла
|
||||||
|
ВСЁ собранное при снятии по watchdog, единственный save был в самом
|
||||||
|
конце (fix/domclick-incremental-save). Зовётся СИНХРОННО сразу после того, как
|
||||||
|
бакет отработал успешно (после DomClickBlockedError/generic
|
||||||
|
Exception — НЕ зовётся), аргументы: имя бакета + лоты ИМЕННО
|
||||||
|
этого бакета (не накопленный out_lots). Вызов стоит ВНЕ
|
||||||
|
try/except этого бакета: исключение из колбэка (кооперативная
|
||||||
|
отмена/SIGTERM-дрейн — см. run_domclick_city_sweep) обязано
|
||||||
|
прервать цикл по ROOM_BUCKETS, а не быть проглоченным generic
|
||||||
|
except'ом. on_bucket=None (дефолт) — поведение прежнее,
|
||||||
|
байт-в-байт.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
Дедуплицированный по source_id список ScrapedLot.
|
Дедуплицированный по source_id список ScrapedLot.
|
||||||
|
|
@ -510,8 +522,17 @@ class DomClickScraper(BaseScraper):
|
||||||
# вживую 09.08: через него 500, через мобильный узел пула — 200 и
|
# вживую 09.08: через него 500, через мобильный узел пула — 200 и
|
||||||
# snippetsCount=678 в бакете 'st'), поэтому свип брал 0 лотов 4 дня подряд.
|
# snippetsCount=678 в бакете 'st'), поэтому свип брал 0 лотов 4 дня подряд.
|
||||||
# proxy_provider=None (тесты/dev) по-прежнему валиден — env-fallback.
|
# proxy_provider=None (тесты/dev) по-прежнему валиден — env-fallback.
|
||||||
|
# reuse_context=True (#3118): без него сайдкар поднимает новый камуфокс на
|
||||||
|
# КАЖДЫЙ /fetch (browser.new_page() создаёт свежий изолированный контекст) —
|
||||||
|
# разовая инъекция cookie выше никогда не видит живой qrator_jsid2, который
|
||||||
|
# сайт ротирует через Set-Cookie (TTL ~2.5ч). Замер: 26 подряд холодных
|
||||||
|
# фетчей = 100% блок, те же карточки в тёплом контексте — 5/5 примерно по 2с
|
||||||
|
# (см. backend/app/tasks/domclick_detail_backfill.py:394-401). Заодно
|
||||||
|
# холодный путь не давал якорной вкладке выжить между фетчами — goto(origin)
|
||||||
|
# валился таймаутом 60с, убивая всю корзину. Тёплый контекст держит один
|
||||||
|
# сайдкар-контекст на весь прогон вместо этого.
|
||||||
async with build_browser_fetcher(
|
async with build_browser_fetcher(
|
||||||
self._config, "domclick", proxy_provider=self._proxy_provider
|
self._config, "domclick", proxy_provider=self._proxy_provider, reuse_context=True
|
||||||
) as fetcher:
|
) as fetcher:
|
||||||
# Циклический сдвиг: состав корзин прежний, меняется только точка входа.
|
# Циклический сдвиг: состав корзин прежний, меняется только точка входа.
|
||||||
# Отрицательный/большой индекс нормализуем — вызывающий передаёт остаток от
|
# Отрицательный/большой индекс нормализуем — вызывающий передаёт остаток от
|
||||||
|
|
@ -548,6 +569,7 @@ class DomClickScraper(BaseScraper):
|
||||||
city_id,
|
city_id,
|
||||||
pages,
|
pages,
|
||||||
)
|
)
|
||||||
|
_bucket_start_len = len(out_lots)
|
||||||
try:
|
try:
|
||||||
await self._sweep_bucket(
|
await self._sweep_bucket(
|
||||||
fetcher=fetcher,
|
fetcher=fetcher,
|
||||||
|
|
@ -569,6 +591,16 @@ class DomClickScraper(BaseScraper):
|
||||||
# прокинут выше) это уже не no-op: узел уходит в
|
# прокинут выше) это уже не no-op: узел уходит в
|
||||||
# scrape_proxy_source_bans и следующий acquire("domclick") его не
|
# scrape_proxy_source_bans и следующий acquire("domclick") его не
|
||||||
# выдаст.
|
# выдаст.
|
||||||
|
# NB: request_context_reset() здесь СОЗНАТЕЛЬНО не зовём. Он лишь
|
||||||
|
# взводит `_context_reset_pending`, а тот уезжает в сайдкар только
|
||||||
|
# ключом `reset_context` СЛЕДУЮЩЕГО fetch() (browser_fetcher.py:565);
|
||||||
|
# `__aexit__` его не сливает. Ниже сразу break — фетчей больше не
|
||||||
|
# будет, флаг умрёт вместе с объектом. Вызов был бы no-op'ом, который
|
||||||
|
# читается как защита.
|
||||||
|
# Дыра остаётся: `_contexts[provider]` в сайдкаре — dict без TTL, так
|
||||||
|
# что сожжённый блоком контекст достанется СЛЕДУЮЩЕМУ прогону. Закрыть
|
||||||
|
# можно только отдельной ручкой сброса в сайдкаре (сейчас там только
|
||||||
|
# /fetch, /fetch-json, /login, /health, /pacing) — отдельной задачей.
|
||||||
fetcher.report_ban(f"domklik QRATOR block during rooms={bucket!r}")
|
fetcher.report_ban(f"domklik QRATOR block during rooms={bucket!r}")
|
||||||
break
|
break
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
|
@ -582,6 +614,14 @@ class DomClickScraper(BaseScraper):
|
||||||
# Exception, а не BaseException — CancelledError (SIGTERM-drain,
|
# Exception, а не BaseException — CancelledError (SIGTERM-drain,
|
||||||
# watchdog asyncio.wait_for) обязан пройти насквозь.
|
# watchdog asyncio.wait_for) обязан пройти насквозь.
|
||||||
self.fetch_errors += 1
|
self.fetch_errors += 1
|
||||||
|
# #3118: корзина упала — тёплый контекст (reuse_context=True) мог
|
||||||
|
# сгореть вместе с ней, и тогда он переползёт в следующую корзину.
|
||||||
|
# Прод 17.09, прогон 7389: после падения 'rooms=3' три корзины подряд
|
||||||
|
# ('5+', 'st', '1') легли за две минуты каждая на таймауте
|
||||||
|
# goto(origin) — якорная вкладка не поднималась в том же контексте.
|
||||||
|
# Сброс стоит РОВНО здесь, а не в ветке DomClickBlockedError: там
|
||||||
|
# сразу break и прогон заканчивается, сбрасывать уже нечего.
|
||||||
|
fetcher.request_context_reset()
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"domklik: bucket rooms=%r failed (%s) — skipping to next bucket",
|
"domklik: bucket rooms=%r failed (%s) — skipping to next bucket",
|
||||||
bucket,
|
bucket,
|
||||||
|
|
@ -591,6 +631,11 @@ class DomClickScraper(BaseScraper):
|
||||||
continue
|
continue
|
||||||
self.buckets_completed += 1
|
self.buckets_completed += 1
|
||||||
self.completed_buckets.append(bucket)
|
self.completed_buckets.append(bucket)
|
||||||
|
# incremental-save: колбэк ВНЕ try/except этого бакета — исключение (канал
|
||||||
|
# кооперативной отмены/SIGTERM-дрейна в pipeline.run_domclick_city_sweep)
|
||||||
|
# обязано прервать цикл, а не попасть в generic except выше.
|
||||||
|
if on_bucket is not None:
|
||||||
|
on_bucket(bucket, out_lots[_bucket_start_len:])
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
|
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue