Compare commits

...
Sign in to create a new pull request.

8 commits

Author SHA1 Message Date
baed8f75a3 Свип ДомКлика ходит в тёплом браузерном контексте, а не поднимает камуфокс на каждый фетч (#3595)
All checks were successful
Deploy Trade-In / changes (push) Successful in 16s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m55s
Deploy Trade-In / build-backend (push) Successful in 1m47s
Deploy Trade-In / deploy (push) Successful in 3m28s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m42s
2026-09-17 18:10:25 +00:00
dbac5d4e1c Вернуть домклик-свипы Москвы и области (#3594)
All checks were successful
Deploy Trade-In / changes (push) Successful in 14s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m53s
Deploy Trade-In / build-backend (push) Successful in 34s
Deploy Trade-In / deploy (push) Successful in 1m27s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m41s
2026-09-17 14:03:00 +00:00
f71a438495 fix(ops): Watchdog уходит из Telegram, появляется эскалация по длительности (#3589)
All checks were successful
Deploy / changes (push) Successful in 14s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy Metrics / server (push) Successful in 50s
Deploy / build-worker (push) Successful in 49s
Deploy / build-backend (push) Successful in 51s
Deploy Metrics / agent-apps (push) Successful in 29s
Deploy Metrics / agent-infra (push) Successful in 29s
Deploy / deploy (push) Successful in 1m14s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 1m41s
Watchdog пингует внешний deadman-приёмник вместо человеческой ленты. Новое правило AlertFiringTooLong при firing дольше 6 часов. Пустой секрет пинга оставляет приёмник без конфигов, а не роняет Alertmanager.
2026-09-17 13:54:02 +00:00
e5212e52ab Свип ДомКлика сохраняет лоты по корзинам, а не одним махом в конце (#3592)
All checks were successful
Deploy Trade-In / changes (push) Successful in 17s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m4s
Deploy Trade-In / build-backend (push) Successful in 2m23s
Deploy Trade-In / deploy (push) Successful in 8m6s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m43s
2026-09-17 13:37:47 +00:00
98f587add6 Merge pull request 'ПТИЦА: гео-радиусная цена участка считается по лотам проектов ближних ЖК, а не по чужим лотам под устаревшим complex_id (#3583)' (#3591) from fix/3583-radius-price-complex into main
All checks were successful
Deploy / changes (push) Successful in 12s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m38s
Deploy / build-worker (push) Successful in 3m50s
Deploy / deploy (push) Successful in 1m31s
Deploy / deploy-status (push) Successful in 2s
Deploy / perimeter-smoke (push) Successful in 1m43s
2026-09-17 13:32:55 +00:00
a82239b931 Merge pull request 'alert-ack: значение секрета вебхука GlitchTip больше не пишется в лог — строка доступа и ошибки разбора дают secret=***' (#3588) from fix/3576-alert-ack-redact-secret into main
All checks were successful
Deploy / changes (push) Successful in 13s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy Metrics / server (push) Successful in 49s
Deploy / build-backend (push) Successful in 47s
Deploy / build-worker (push) Successful in 49s
Deploy Metrics / agent-apps (push) Successful in 28s
Deploy Metrics / agent-infra (push) Successful in 28s
Deploy / deploy (push) Successful in 1m16s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 1m41s
2026-09-17 13:15:01 +00:00
1265e69412 chore(forgejo): ограничить размер json-лога контейнера (#3587)
Лог forgejo дорос до 3.35 ГиБ за 23 дня без ротации. 50m x 3 = потолок 150 МБ.
Co-authored-by: lekss361 <lekss361@gendsgn.local>
Co-committed-by: lekss361 <lekss361@gendsgn.local>
2026-09-17 13:08:15 +00:00
8b9b541736 alert-ack: значение секрета из query больше не пишется в лог (#3576)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 19s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 23s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 2m57s
CI / backend-tests (pull_request) Successful in 7m35s
GlitchTip шлёт секрет резервного вебхука только в `?secret=`, а
BaseHTTPRequestHandler печатает строку запроса целиком: в строке доступа
(log_request) и в тексте ошибки разбора (log_error, «Bad request syntax
('POST /glitchtip?secret=…')»). Оба пути сходятся в log_message — маскируем
там одним выражением, тем же, что у бэкенда МЕРЫ (#3154, log_scrub.py) и у
Alloy (#3354). Импортировать его нельзя: сервис намеренно без зависимостей.

Тест в backend/tests/ops/test_3078_alert_ack.py (его гоняет CI по ops/**):
настоящий сокет, три строки запроса — доступ, имя с префиксом, ошибка
разбора; значения в записях нет, `=***` стоит в ожидаемом числе записей.

Ротация секрета — за владельцем, здесь не делается.

Refs #3576

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-17 18:05:16 +05:00
22 changed files with 1044 additions and 89 deletions

View file

@ -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

View file

@ -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

View file

@ -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}"

View file

@ -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, "вебхук указывает не на сервис кнопки"
# У прочих приёмников прямой путь сохранён. # У прочих приёмников прямой путь сохранён.

View file

@ -70,6 +70,18 @@ services:
FORGEJO__security__INSTALL_LOCK: "true" FORGEJO__security__INSTALL_LOCK: "true"
volumes: volumes:
- ./data/forgejo:/data - ./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: ports:
- "2222:22" - "2222:22"
networks: networks:

View file

@ -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)

View file

@ -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}

View file

@ -487,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 часов — похоже, её не заметили или на неё забили."

View file

@ -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 часов — похоже, её не заметили или на неё забили."

View file

@ -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
),
) )

View file

@ -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;

View file

@ -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] = []

View file

@ -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:

View 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}"
)

View file

@ -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 []

View 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

View file

@ -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,
} }
) )

View file

@ -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,

View file

@ -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'ы, и

View file

@ -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
),
) )

View file

@ -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,
) )

View file

@ -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 "