Compare commits

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

4 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
19 changed files with 969 additions and 88 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,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

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