#3421 въехал в main параллельно с той же миграцией 288 (deals.doc_type,
параметры region_code/doc_types). Разрешение: 288 — целиком версия main;
наша дельта (FDW-колонки okato/quarter_cad_number/district, выключенный seed
rosreestr_dkp_import_77) переехала в 289. scheduler.py — doc_types из main +
canonical_city-маппинг/raw_payload/per-source чекпоинт. deploy-скрипт —
валидация REGION_CODE и DOC_TYPE (интерполируются в SQL текстом).
Трек 2 подготовки Mera к Москве. import_rosreestr_dkp принимает region_code из
params (default 66 — байт-в-байт прежнее поведение), валидирует его через
app.services.regions.REGIONS. Регион с canonical_city (77 — Москва, Росреестр
отдаёт округ/поселение вместо города) подставляет city/address через одну
SQL-ветку на bind-параметре :canonical_city, а не Python if/else на код региона;
city IS NOT NULL не фильтруется для такого региона (иначе теряется ~10% строк),
исходные city/okato/quarter_cad_number/district уходят в raw_payload.
Чекпоинт курсора (_resume_dkp_cursor) стал per-region: source для поиска
предыдущего прогона строится через _dkp_source_for_region (66 сохраняет
легаси-имя 'rosreestr_dkp_import', остальные — суффикс кода) — иначе прогон по
77 либо никогда не резюмился бы (source-литерал не матчил), либо, при более
наивном фиксе, унёс бы курсор чужого региона.
product_handlers регистрирует wildcard rosreestr_dkp_import_* (по образцу
deactivate_stale_*/avito_city_sweep_*), deploy/import-rosreestr.sh получил
REGION_CODE env (bash-путь не region-generic — city-override только в Python).
Migration 288: deals.doc_type + backfill 'ДКП' для source=rosreestr, foreign
table gendesign_rosreestr_deals расширена okato/quarter_cad_number/district
(проверено live на прод-БД), выключенный seed rosreestr_dkp_import_77.
ПОЧЕМУ: расширение на Москву упирается в два литерала. В источнике за 2024 по региону 77
лежат 30 627 ДДУ с медианой 112 743 против 107 005 ДКП с медианой 256 250 — это цены
котлована, и без различимого признака в deals они развалят любую оценку. При этом тип
сделки терялся при загрузке вовсе (в deals колонки не было), а фильтры region_code = 66
и doc_type = 'ДКП' стояли литералами в scheduler.import_rosreestr_dkp и в двойнике
deploy/import-rosreestr.sh — сменить регион было нельзя, не правя код.
ЧТО:
- миграция 288: deals.doc_type text (idempotent) + бэкфилл 'ДКП' для source='rosreestr'
(корректен, а не эвристика: всё загруженное прошло фильтр ДКП — и в импорте, и в 077)
+ явный region_code=66 в default_params расписания rosreestr_dkp_import вместо неявного
дефолта в коде. Индекс НЕ добавлен: 2-3 значения, живые выборки идут по
region_code/deal_date/geom — заведём частичный, когда появится режущий запрос;
- import_rosreestr_dkp: region_code (default 66) и doc_types (default ['ДКП']) из params,
фильтры через bind-параметры CAST(:region_code AS int) / ANY(CAST(:doc_types AS text[])),
doc_type едет из SELECT в INSERT и в ON CONFLICT DO UPDATE. Дефолты сохраняют текущее
прод-поведение байт-в-байт;
- dedup_hash оставлен как 'ros:dkp:' || id: id уникален в источнике независимо от типа
документа, а смена формы ключа осиротила бы уже загруженные строки (ровно то, что
разгребала миграция 077);
- deploy/import-rosreestr.sh: REGION_CODE / DOC_TYPE как env со старыми дефолтами,
doc_type протащен через staging в deals; шапка про «ЕКБ квартиры» переписана честно —
city-фильтр снят давно, скоуп = весь регион;
- тесты: test_rosreestr_dedup_key переведён с ассертов на литералы на проверку
«параметр + дефолт = скоуп 077»; новый test_3051_* проверяет bind-параметры реальным
вызовом с моком Session, дефолты 66/['ДКП'], doc_type в колонках INSERT и текст 288.
Две живые копии одного модуля с противоположной семантикой counters: app
`mark_done`/`mark_failed`/`mark_banned`/`update_heartbeat` ЗАМЕНЯЛИ
(`counters = CAST(:counters AS jsonb)`), kit — МЕРЖИЛИ
(`COALESCE(counters,'{}') || …`). Расхождение дважды за сутки дало ложные
выводы на ревью (#3388 «отдать только флаг, остальное домержится» — на
replace это стёрло бы измеренное; #3355). Разошлись и другие места: гейт
статуса, `honors_cancel` у mark_cancelled (был только в app), `mark_skipped`
(только в kit), `mark_backfill_finished`/`distinct_sources` (только в app).
Реализация теперь одна — `scraper_kit.orchestration.runs`; в неё перенесены
app-only функции. `app.services.scrape_runs` — алиас kit-модуля через
sys.modules, а не реэкспорт имён: реэкспорт разводит патч-цели
(`patch("app.services.scrape_runs.sentry_sdk")`, `patch.object(runs_mod,
"mark_done")` правили бы глобаль модуля-обёртки, а тело функции читает
глобаль kit'а) — тест остался бы зелёным, не подменив ничего. С алиасом оба
имени ведут в единственную реализацию, и ни один из ~40 вызывающих и ~30
патч-сайтов в тестах не правится.
Победила семантика мержа: у строки прогона несколько писателей (пульс,
финализатор, дрейн), каждый знает лишь свои ключи, и замена теряла чужие —
чекпоинт done_buckets (#930), метку interrupted (#3391), замер из пульса
(#3384). Обратной зависимости («вызывающий рассчитывает, что финализатор
УДАЛИТ ключ заменой») нет: строка создаётся пустой в create_run, резюм читает
counters ПРЕДЫДУЩЕГО прогона по его id.
Тесты по значению на обоих путях импорта (двойник сессии читает SQL: `||`
против CAST, WHERE-гейт из текста): пульс {a:5} + mark_failed {b:1} → {a,b};
пульс/финализатор по финализированной строке — no-op; mark_cancelled
отказывает источнику, который отмену не опрашивает. На main эти тесты
красные для app-пути.
Комментарии в app/services/scheduler.py и kit/pipeline.py, утверждавшие про
живого «перезаписывающего двойника», приведены в соответствие.
Ревью нашло у цианa (в отличие от avito/домклика с живым counters.to_dict()) старый
словарь в общем except: реальные значения присваиваются уже ПОСЛЕ возврата из
backfill_cian_history, а отказ бывает и посреди неё — пул опустел между стадиями, упал
SELECT домов. Тогда поверх измеренного в запись прогона уезжали нули, и SQL-разбор
простоя (#3288/#3367) читал «к площадке не ходили» про прогон, который ходил.
Механизм оказался хуже описанного в ревью: mark_failed мержит counters (`counters ||
:counters`) только в kit-копии, а cian/avito/домклик зовут app.services.scrape_runs, где
UPDATE counters ЗАМЕНЯЕТ (scrape_runs.py:738). Поэтому «отдать только {no_proxy_stop: 1}»
стёрло бы измеренное начисто; вместо этого _heartbeat кладёт свой снимок в те же
counters (nonlocal), и в mark_failed уезжает последнее измеренное + флаг.
Тест по значению: heartbeat записал listings_processed=5, дальше пул пуст → в jsonb-
payload mark_failed должно остаться 5, а не 0 (проверяется сам payload UPDATE'а,
runs_mod настоящий). На HEAD ветки красный: `counters={'listings_processed': 0, ...,
'no_proxy_stop': 1}: нули поверх измеренных 5`.
Стаб пула приведён к проду: RealProxyProvider.acquire при пустом пуле ВОЗВРАЩАЕТ None
(scraper_adapters.py:230), а не поднимает, — исключение из провайдера глотал
`except Exception` в _acquire_lease и приходило к тому же отказу другим путём. Теперь
NoProxyAvailableError рождается там же, где в проде (browser_fetcher.py:712, ветка
`lease is None and use_pool and production`) — проверено прогоном против до-#3384
исходников: все три теста красные, трейс из _acquire_lease.
_prod_pool патчит app.core.config.settings явно + assert, что все три задачи держат тот
же синглтон: раньше патч через chb.settings выглядел настройкой одного циана.
Lease берётся один раз в BrowserFetcher.__aenter__, поэтому на проде с пустым пулом
NoProxyAvailableError вылетает из самого `async with` — ДО первой карточки и мимо
стоп-механики внутри цикла (no_proxy_stop = True; break), которая и пишет
counters.no_proxy_stop. Общий `except Exception` ловил его и делал mark_failed с
нулевыми counters без ключа: прогон, который к площадке не ходил вообще, в SQL-разборе
простоя по counters.no_proxy_stop (#3288/#3367) не находится.
Дыра одинаковая у всех трёх бэкфиллов (avito её тоже не обрабатывал: __aenter__
вызывается напрямую строкой 427, отказ уходит в тот же общий except). Правка — в трёх
уже существующих обработчиках, которые и так зовут mark_failed: ключ no_proxy_stop=1
при caused_by_no_proxy(exc). Оборачивать `async with` в try/except пришлось бы с
переносом ~200 строк тела под новый отступ в каждом файле, и покрывало бы только
падение на входе; здесь ловится любой путь мимо цикла.
Тест — через настоящий BrowserFetcher: подделан только провайдер прокси (его acquire
поднимает NoProxyAvailableError), отказ рождается там же, где в проде. Проверяется
failed + no_proxy_stop=1 + attempted=0 (у циана listings_processed=0) и ноль POST'ов
в сайдкар.
Closes#3384
BrowserFetcher(source="cian") в cian_history_backfill конструировался без
proxy_provider/use_pool/environment — трёх аргументов, которые кладут "proxy" в тело
POST /fetch. Сайдкар брал свой env-прокси (SCRAPER_PROXY_URL): пул из 4 узлов, его
баны и ротация проходили мимо, а прод-отказ «пул пуст → не ходить на env/direct»
(#2616) на этом пути был мёртв, потому что смотрит на environment. Проводка теперь
как у соседей — domclick_detail_backfill и house_imv_backfill.
Ожившему отказу нужен обработчик: NoProxyAvailableError ловился общим except на
объявление, и батч крутил впустую весь список (пул пуст с первого — значит пуст и на
1000-м). Распознаём по цепочке причин, обрываем прогон, counters.no_proxy_stop=1 и
mark_failed вместо mark_banned — отказ нашей стороны не должен записываться как бан
Циана. Дома и оценки после стопа пропускаем: они идут через тот же пул.
caused_by_no_proxy вынесен в scraper_kit.proxy_errors (у avito #3288 и domclick #3283
живут приватные копии — их схлопывание отдельной правкой).
Карточки Циана доставались побочным эффектом cian_history_backfill, а её
выборка ключуется по offer_price_history. Следствия на 30.08: карточка есть
у 4494 из 25222 объявлений (17.8% — последнее место при втором месте по
объёму), 19046 без истории при квоте 100/сутки (190 дней на остаток, тогда
как очередь растёт вдвадцатеро быстрее), и 1697 объявлений с историей и без
карточки, которые исторической выборке недостижимы в принципе.
Фетчер при этом исправен: прогоны 5154/5240/5328 дали 100/100, 99/100,
100/100. Чинить нечего — не выдана мощность.
Добавлен второй режим выборки (listings_pending="detail", по
detail_enriched_at, свежие первыми) и второе расписание поверх ТОГО ЖЕ тела:
машинерия работает, дублировать её новым модулем незачем. Историческая
выборка оставлена побайтово — по ней живёт суточный прогон.
batch_size=400 не на глаз: замеренный темп ~28с на объявление, порог
reap_zombies 6ч по heartbeat, бюджетного сторожа у задачи нет — 400×28с≈3.1ч
проходит, 800 как у Яндекса (≈6.2ч) убивало бы жнецом.
Расписание засеяно enabled=false, как domclick_detail_backfill в миграции
175: это третий круглосуточный добор на общий пул из четырёх узлов, влияние
на соседей надо посмотреть, а не предположить.
Тесты: 13 проверок, ключ выборки / порядок / неизменность прежнего режима /
проводка параметров через посредника / регистрация обоих source. Проверено
мутациями: снятие ORDER BY, молчаливый дефолт вместо ValueError и потеря
listings_pending в посреднике роняют по 2-3 теста каждая. Набор целиком —
5196 passed, 37 skipped.
Сайдкар вообще не читал код ответа page.goto: страница классифицировалась
только по маркерам, снятым с Авито. Домклик отдаёт статическую `403 | Домклик`
на 26 624 байта, где нет ни одного такого маркера (замер прода 28.08.2026) —
она уезжала наверх как валидный HTML, парсер не находил состояние, и прогон
получал блок неизвестной природы. За 14 дней все 14 прогонов домклика легли с
ban_kind='unknown'; у Яндекса счётчика blocked не было вовсе, поэтому ветка
перевода прогона в 'banned' была недостижима по построению — ноль банов.
- browser/server.py: статус целевой навигации сохраняется per-provider и
доезжает в тело /fetch аддитивным ключом "status" (ключ "html" не тронут);
403/429 с маркерами челленджа больше не ждут PoW — ждать нечего, статическая
страница сама себя не перезагрузит. Наверх идёт BanPageDetectedError, а не
заглушка: вернув её контентом, воскресили бы #3045.
- scraper_kit/browser_fetcher.py: BrowserFetcher.last_response_status +
ban_kind_from_status (403/429 → platform, 5xx → infra, прочее → None).
Поток управления не менялся: fetch() по-прежнему отдаёт str.
- domclick: DomClickBlockedError несёт .status — один тип исключения на
маркер-детект и на сбой фетча разводится без размножения типов; прогон
передаёт перепись диагнозов в mark_backfill_finished.
- yandex: появился счётчик blocked, оживляющий ветку бана. Серии блоков и
промахов парсера считаются РАЗДЕЛЬНО: иначе четыре промаха плюс один 403
пятым давали 'banned' с переписью {platform: 1}.
- cian: ban_kinds наполняется только диагностируемым статусом. HTTP 200 с
пустым разбором — дрейф разметки на нашей стороне, а не отказ площадки;
записав его блоком, мы бы штамповали фиктивные баны у здорового источника
(13 done против 1 banned за 14 дней).
Инвариант: непустой ban_kinds ⟺ виден ответ 403/429/5xx. Значения остаются в
пределах CHECK scrape_runs.ban_kind.
Известный пробел: шов providers/domclick/detail.py `blocked.status = status`
тестами не покрыт — существующие домкликовые тесты подают исключение готовым
моком и боевой fetch_detail не исполняют.
last_id жил только в памяти процесса (import_rosreestr_dkp, scheduler.py):
heartbeat писал его в scrape_runs.counters каждый батч (комментарий рядом
прямо называл это чекпоинтом), но при старте last_id всегда инициализировался
литералом 0 — обрыв (деплой/OOM/рестарт хоста) откатывал прогресс и заставлял
пере-сканировать источник с начала.
Разведка: из пяти backfill-циклов issue (avito_detail_backfill,
house_imv_backfill, cian_history_backfill, yandex_detail_backfill,
geocode_missing_listings) ни один не имеет этого дефекта — все устроены как
WHERE ... IS NULL/NOT EXISTS ... LIMIT, естественно резюмируемы без курсора.
Единственный код, буквально описанный в issue (строки/SQL/комментарий),
это шестой, не входящий в таблицу backfill — rosreestr_dkp_import.
Фикс — _resume_dkp_cursor(db, run_id):
- кандидат — последний прогон source='rosreestr_dkp_import';
- резюмится только незавершённый штатно прогон: status running/zombie,
либо done с counters.interrupted=1 (SIGTERM-drain — эта ветка раньше
считала последующий full rescan штатным поведением, теперь помечает
себя как прерванную и резюмится наравне с zombie);
- потолок возраста чекпоинта — 24ч, старше — 'checkpoint_stale', старт с 0;
- чистый 'done' (полный проход) не резюмится — иначе ON CONFLICT DO UPDATE
перестанет ловить правки уже импортированных сделок при следующем проходе.
Вердикт и per-batch чекпоинт пишутся через kit_runs.update_heartbeat (merge
`counters || :counters`) вместо локального runs_mod.update_heartbeat (полная
замена) — иначе resume-вердикт стирался первым же heartbeat'ом батча.
Тесты: tests/test_3168_backfill_cursor_resume.py — резюм с сохранённого
last_id, резюм после SIGTERM-drain, отказ резюмить чистый done, отказ
резюмить протухший (>24ч) чекпоинт, merge не стирает посторонние ключи.
Обратимость проверена вручную (временный откат _resume_dkp_cursor красил
6 из 8 тестов).
A. ДОМ.РФ year backfill: авторитетный commission_year больше не экранируется
невозможным существующим year_built. Валидное существующее значение
выигрывает (COALESCE-семантика #2013 сохранена), но NULL/impossible
(< 1850, > текущий+2, 0) заменяется валидным commission_year → zhkh_year.
parse_int_field получил sanity-гейт [min,max] — мусорный commission_year
не попадает в staging (7 таких строк в текущем staging).
B. propagate_listings_year: добавлен link-consistency guard. Раньше копировал
houses.year_built на listings по house_id_fk без проверки — перепривязанный
FK впрыскивал чужой когортный год. Теперь пропагация только если координаты
объявления в пределах 500м от геометрии дома (ST_DWithin), либо (без коорд.)
консервативный address-фоллбек по short_address. Live: блокирует 1 из 10
текущих кандидатов (>500м mislink).
C. rosreestr_dkp_import: per-row INSERT-ошибки отделены от dedup-skip. Раньше
except инкрементил тот же batch_skipped, что и легитимный ON CONFLICT — сбой
маскировался под дедуп и run рапортовал success. Отдельный rows_errored +
гейт по доле ошибок (> 5% → run FAILED, не silent-green).
D. rosreestr_dkp_import: ON CONFLICT DO NOTHING → DO UPDATE изменяемых сырых
фактов Росреестра (price/area/rooms/floor/year/deal_date/...), чтобы
исправленный/переопубликованный квартал обновлялся. Обогащение
(lat/lon/geom/geocode_tried_at и пр.) не в SET-списке — не затирается.
IS DISTINCT FROM guard сохраняет идемпотентность resume; RETURNING (xmax=0)
отличает insert от update (rows_inserted vs rows_updated).
- scheduler.py import_rosreestr_dkp: снят фильтр city ILIKE '%катеринбург%', address из реального city источника, deals.region_code + новая deals.city заполняются (было: хардкод 'Екатеринбург,' + region_code NULL)
- migration 177: deals.city + индекс + бэкфилл существующих EKB-строк region_code=66/city
- guard city IS NOT NULL → address не NULL для ~15 null-city строк источника
- sync deploy/import-rosreestr.sh (ops-fallback) под тот же oblast-scope
- split dedup-теста: 077 (историческая) хранит EKB-фильтр, живой импорт — нет
Открывает +47183 не-ЕКБ сделок region 66, уже сидящих в источнике, ранее резавшихся на импорте.
Топология подтверждена перед удалением (docker-compose.prod.yml): tradein-backend
(uvicorn app.main:app) — SCHEDULER_ENABLE=false; tradein-scraper (python -m
app.scheduler_main) — SCHEDULER_ENABLE=true + USE_KIT_SCHEDULER=true. Kit-путь
(_run_kit_scheduler → scraper_kit.orchestration.scheduler + product_handlers)
самодостаточен: не импортирует ничего из app.services.scheduler.scheduler_loop
или app.services.scrape_pipeline. Все НЕ-sweep джобы, которые kit-scheduler
диспетчерит через build_product_handlers, идут напрямую в app.tasks.*/
app.services.* (либо lazy-импортят import_rosreestr_dkp/_execute_cian_backfill
из scheduler.py) — мимо удаляемой legacy-машинерии.
app/services/scheduler.py: 2098 → 418 строк. Удалено: scheduler_loop,
get_due_schedules, reap_zombies, _claim_run, _defer_next_run_at, _spawn_tracked/
_drain_inflight/_inflight_tasks, все 27 trigger_*_run-функций, импорт
app.services.scrape_pipeline, константы SCHEDULER_TICK_SEC/ZOMBIE_THRESHOLD_HOURS
(достижимы были только через удалённый scheduler_loop-путь). Оставлено (живые
импортёры вне удалённого): compute_next_run_at + has_running_run (admin.py),
import_rosreestr_dkp + _execute_cian_backfill (lazy-импорты в
product_handlers.py — job-тела kit-handler'ов).
main.py: убран `from app.services.scheduler import scheduler_loop` + lifespan-блок
запуска (`if settings.scheduler_enable: asyncio.create_task(scheduler_loop())`);
прод-backend всегда шёл с SCHEDULER_ENABLE=false, так что это был мёртвый код.
scheduler_main.py: убрана ship-dark развилка #2192 (USE_KIT_SCHEDULER=false →
legacy scheduler_loop fallback) — _run_kit_scheduler() теперь безусловный путь.
Поле settings.use_kit_scheduler оставлено в конфиге (Settings extra="ignore"
защищает от startup-краха на leftover env var), но на ветвление не влияет.
app.services.scrape_pipeline: 0 runtime-импортёров в app/+scripts/+packages/
после этого PR (только тесты, которые Part E удалит вместе с самим файлом) —
подтверждено grep. scrape_pipeline.py не тронут (Part E).
Тесты: удалены test_house_imv_backfill_scheduler.py (100% legacy-триггер,
backfill_house_imv сервис покрыт в test_house_imv_backfill_browser_flag.py /
test_backfill_wave2.py) и test_kit_registry_completeness.py (parity-инвариант
против удалённого dispatch, дублирует test_scraper_kit_scheduler_parity.py).
Точечно вырезаны "Scheduler wiring" секции (trigger_fn_exists/dispatch_branch_
wired/runs_in_executor) из ~10 файлов, тестирующих сами task-функции — сами
task-тесты (SQL-shape, миграции, fake-db поведение) оставлены нетронутыми.
test_scheduler.py: 825 → ~90 строк (остались только compute_next_run_at-тесты).
test_scraper_kit_scheduler_parity.py: убрана golden-parity секция против
удалённого scheduler_loop (SOURCE_TO_OLD_TRIGGER/_drive_old_one_tick/
test_routing_parity_per_source), остальное (claim/reap_zombies/dispatch/
registry-shape тесты kit-модуля) сохранено — источник этих инвариантов не
app.services.scheduler, а сам scraper_kit.orchestration.scheduler.
test_scheduler_main.py: 2 теста, патчившие app.services.scheduler.scheduler_loop,
переведены на монкипатч sm._run_kit_scheduler (единственный путь после этого PR).
test_sweep_imv_phase.py:171-371 (6 прямых импортов run_avito_city_sweep из
scrape_pipeline) намеренно НЕ тронуты — Part E.
Verify: полный pytest 3179 passed / 6 skipped / 1 known-unrelated fail
(test_search_cache_hit, #2208, не связан с этим PR); ruff 0.7.4 чист на всех
изменённых файлах; `python -c "import app.main; import app.scheduler_main"` OK.
Финальный PR issue #2045 (BE-3): GET /api/v1/trade-in/location-coef для
LocationDrawer. FDW foreign table -> локальное зеркало osm_poi_ekb_local
(TRUNCATE+INSERT, тот же паттерн что cad_buildings_local/cadastral_geo_match,
избегает ~1.16s/row FDW round-trip) -> straight-line POI-скоринг, портированный
из Site Finder poi_score.py::compute_poi_weighted_top7 (CATEGORY_WEIGHTS as-is,
радиус 1200м для квартир вместо Ptica 2000м для участков). score->coef -
новая MVP-эвристика (0.95..1.05, не откалибрована на реальных дельтах).
Graceful fallback (не 500, не фабрикуем факторы): пустая/не отрефрешенная
osm_poi_ekb_local или отсутствие lat/lon у оценки -> coef=1.0, factors=[],
geo_source="unavailable".
Scheduler: source=osm_poi_ekb_refresh, daily, зарегистрирован и в боевом
dispatch (scheduler.py), и в kit-registry (product_handlers.py) - иначе
test_kit_registry_completeness падает на ship-dark инварианте (#2192).
Frontend wiring (mappers.ts/LocationDrawer.tsx) - вне scope, отдельная задача
после проверки endpoint'а curl'ом на деплое.
Pre-push review: _await_scheduler ждал только scheduler_loop COORDINATOR, но вся
scrape-работа крутится в detached asyncio.create_task детях (каждый trigger_* делал
`task = create_task(_run())` без join). На SIGTERM coordinator выходил из while True
и завершался → asyncio.run() teardown хард-кансельил ещё бегущих детей mid-await =
ровно #1182 failure mode. Кооперативный checkpoint спасал ребёнка лишь когда его
residual-время случайно перекрывало drain — вероятностно, не гарантированно.
Fix — coordinator теперь дренажит детей перед выходом:
- _spawn_tracked(coro): централизованный detached-spawn, кладёт задачу в module-level
registry _inflight_tasks (strong-ref = RUF006 keep-alive) + done-callback ретривит
exception и убирает из set'а. Заменил 26 одинаковых
`task = create_task(_run()); task.add_done_callback(...)` сайтов.
- _drain_inflight(): один asyncio.wait по живым детям с бюджетом
_CHILD_DRAIN_TIMEOUT_S=80s (< scheduler_main 100s < docker grace 120s). Кооперативные
дети (avito_detail_backfill, rosreestr-executor) дочекивают карточку/батч + mark_done
и резолвятся; некооперативные упираются в timeout и падают на внешний hard-cancel.
- scheduler_loop по выходу из tick-loop (только по SIGTERM) зовёт await _drain_inflight().
NB: raw asyncio.all_tasks()-minus-self здесь НЕЛЬЗЯ — в нашей топологии он захватывает
_run parent-task (блокирован на wait_for(coordinator)) и shutdown_waiter → циклическое
ожидание coordinator↔_run, всегда упирающееся в timeout. Точный registry это исключает.
Tests: tests/test_scheduler.py — detached cooperative child drained-not-cancelled,
non-cooperative child timeout→left for hard-cancel, no-op без детей, scheduler_loop→drain
wiring. Обновил 3 source-inspection теста под новый _spawn_tracked паттерн.
Деплой (docker recreate tradein-scraper) шлёт SIGTERM с stop_grace_period=120s
(Phase 1). Раньше scheduler_main отвечал hard task.cancel() → бегущий scrape-unit
получал CancelledError посреди браузер-карточки → карточка гибла без COMMIT.
Теперь SIGTERM/SIGINT лишь выставляют кооперативный флаг; tick-loop и длинные
задачи опрашивают его на СВОИХ существующих between-unit checkpoint'ах,
докоммичивают текущий unit и выходят сами.
- app/core/shutdown.py: новый standalone-модуль (asyncio.Event + request_shutdown
/ shutdown_requested / wait_for_shutdown), без app-зависимостей → нет циклов.
- scheduler_main._run: SIGTERM → request_shutdown() вместо task.cancel(); ждём
добровольного drain'а, safety-net wait_for(100s < docker grace) с fallback на
hard-cancel для некооперирующей задачи.
- scheduler_loop: проверка флага в начале тика и после каждого dispatch — не
reap/claim новые run'ы во время shutdown, выходим из loop'а.
- avito_detail_backfill: break на границе карточки рядом с budget-guard; snapshot
pending-query идемпотентен → следующий run сам резюмит остаток.
- import_rosreestr_dkp: расширен is_cancelled checkpoint — drain делает mark_done
(partial), не mark_cancelled; user-cancel семантика без изменений.
Без SIGTERM поведение идентично прежнему. reap_zombies не трогает drained-run'ы
(они mark_done, не 'running').
Tests: tests/core/test_shutdown.py, scheduler_main drain+timeout-fallback,
avito_detail_backfill partial-drain. 26 + 29 scheduler passed.
scrape_schedules had no interval column -> every source ran daily. Add
interval_days from default_params (default 1, backward-compatible) so
compute_next_run_at can schedule N days ahead. Set avito_full_load_exhaustive
to interval_days=7 (weekly full pass for last_seen/delisting + silent price
edits); avito_full_load stays daily incremental.
run_avito_full_load gains incremental_days param -> passes since= to the
(merged) incremental SERP engine. avito_full_load schedule flipped to
incremental_days=2 (shallow, date early-stop -> avoids deep-pagination 429
bans). New avito_full_load_exhaustive source runs the full walk weekly to
refresh last_seen (10-day delisting TTL) and catch silent price edits.
Подключает prod-проверенный run_cian_full_load (exhaustive региональный сбор
Cian ЕКБ вторички, room×price партиционирование, incremental on_bucket save) в
in-app scheduler как recurring-источник cian_full_load (окно 20-22 UTC — отделено
от cian_city_sweep/history_backfill на 2-5 чтобы не конкурировать на shared
kf-прокси).
run_cian_city_sweep теперь NB-only (newbuilding_only=True default): SERP-фаза
сохраняет только novostroyki-лоты, вторичку отбрасывает — ею авторитетно владеет
full_load. Убирает дубль secondary SERP-сейва. DETAIL/HOUSES-фазы не затронуты.
Backfill listings.building_cadastral_number (was 0%) from the nearest
cadastral building. gendesign_cad_buildings is a postgres_fdw foreign table
with no geom — a per-listing FDW nearest query is ~1.16s/row (~13h for 43k
listings). Instead materialize the FDW once into a LOCAL cad_buildings_local
table (Point geom + GIST), then run a fast local KNN nearest-neighbour join.
Perf: the distance gate uses geometry-space ST_DWithin(geom, point, deg) (GIST
index, no geography cast) + geom <-> KNN order + ST_DistanceSphere metric
recheck on the single nearest row. A geography-cast ST_DWithin in the WHERE
defeated the index (58s+, full update never finished); the geometry-gate design
does the whole ~41.9k-listing UPDATE in ~5.4s (refresh+match ~7s end-to-end),
matching 16027 listings at 50m across 2229 distinct buildings.
The match is GEO-NEAREST (approximate): a street-level-geocoded listing matches
the nearest building within threshold_m (default 50m), not necessarily its exact
cadastral building. Exact cadastral + parcel-containment deferred (cad_parcels
FDW not exposed). Threshold is logged.
- 124: cad_buildings_local table (empty) + GIST. Migration does not read the FDW
(deploy-independent); populated by the refresh job.
- 125: scrape_schedules seed source=cadastral_geo_match, enabled, next_run_at
tomorrow 09:00 UTC (after geocode_missing).
- tasks/cadastral_geo_match.py: refresh_cad_buildings_local (bulk FDW scan),
match_listings_to_buildings (chunked LATERAL KNN UPDATE), run_cadastral_geo_match
run-lifecycle wrapper.
- scheduler: trigger_cadastral_geo_match_run + dispatch (sync DB-only, executor).
- scheduler.py: исправлен stale-комментарий yandex_detail_backfill
(был: «YandexDetailScraper httpx, no proxy layer»;
теперь: fetch via curl_cffi chrome120 + scraper_proxy_url,
parse via YandexDetailScraper.parse)
- yandex_detail_backfill.py: outer except — сообщение лога изменено на
«save/iteration error for listing_id=%d» чтобы отличать save-ошибку
от fetch/parse-None путей (run_id убран — он менее важен в этом scope)
- New task app/tasks/avito_detail_backfill.py with run_avito_detail_backfill()
* Single snapshot SELECT at start (guarantees termination)
* Same proxy/AsyncSession path as scrape_pipeline.py step 5
* Budget guard (budget_sec), consecutive block abort (mark_done not mark_failed)
* rotate_ip() on every AvitoBlockedError; rollback on generic Exception (#1368)
* start = time.monotonic() initialized before try so except can reference it
- Scheduler wiring: trigger_avito_detail_backfill_run() + elif in scheduler_loop()
- Migration 112: scrape_schedules INSERT window 09-12 UTC, batch_size=800,
budget_sec=3600, request_delay_sec=6, max_consecutive_blocks=5
- 7 unit tests (no pytest-mock, unittest.mock only): all 7 passing,
full CI suite 1809 passed
Регистрирует отдельный in-app scheduler-источник для cian_newbuilding
enrichment-backfill (houses_price_dynamics / house_reliability_checks /
house_reviews), который раньше бежал только инлайн в Cian full-sweep.
- run_newbuilding_enrich() wrapper (heartbeat → backfill → mark_done/failed),
делегирует backfill_newbuilding_enrichment (#972, уже в main).
- trigger_newbuilding_enrich_run() + dispatch-ветка в scheduler_loop
(зеркалит trigger_sber_index_pull_run): _claim_run guard, async task,
zombie-reap наследуется; idempotency наследуется (skip уже обогащённых,
per-house SAVEPOINT, bounded per-fire limit дренирует backlog 306 домов).
- seed-миграция 103: source='newbuilding_enrich', окно 00:00-01:00 UTC
(03:00-04:00 МСК, до 01:00+ UTC sweep-блока), limit=25, next_run_at на завтра.
Засеяно enabled=FALSE (dormant): на prod весь external-HTTP scraping намеренно
на паузе. Расписание+триггер готовы, но не стартуют до намеренного возобновления
(UPDATE ... SET enabled=true). Не активирую одиночный источник при общей паузе.
Closes#973