merge main (#2702 clock_timestamp) в ветку #2670
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 3m5s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 3m5s
This commit is contained in:
commit
a5ce832069
12 changed files with 865 additions and 28 deletions
|
|
@ -2,6 +2,31 @@
|
||||||
|
|
||||||
Таблица scrape_runs создана в 015_scrape_runs.sql.
|
Таблица scrape_runs создана в 015_scrape_runs.sql.
|
||||||
Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled.
|
Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled.
|
||||||
|
|
||||||
|
ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL —
|
||||||
|
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции и не двигается,
|
||||||
|
сколько бы та ни жила. Финализаторы (mark_done/mark_failed/mark_banned) выполняются
|
||||||
|
ТОЙ ЖЕ сессией, что и работа задачи, — и если рабочая транзакция всё это время
|
||||||
|
оставалась открытой (задача ничего не коммитила: нечего было сохранять, батч читающий,
|
||||||
|
сохранение шло чужой сессией), их UPDATE попадал ВНУТРЬ неё, и `finished_at` получал
|
||||||
|
время НАЧАЛА работы, а не её конца.
|
||||||
|
|
||||||
|
Замер на проде 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик
|
||||||
|
counters.duration_sec): у 153 заявленная длительность превышала собственное окно
|
||||||
|
finished_at − started_at более чем в 1.5 раза, у 133 окно было меньше секунды при
|
||||||
|
работе дольше 10 с. 126 из этих 133 окон лежат в диапазоне 9-64 мс — это не разброс,
|
||||||
|
а подпись механизма: столько проходит от коммита claim'а до первого запроса рабочей
|
||||||
|
транзакции. Крайний случай — прогон 346 (cian_history_backfill): 18230 с работы,
|
||||||
|
окно 32 мс.
|
||||||
|
|
||||||
|
Дефект был не сплошной ровно потому, что зависел от того, коммитила ли задача перед
|
||||||
|
финалом: cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят
|
||||||
|
поштучно, у них окно совпадало с работой; yandex_address_backfill (45 из 50 прогонов),
|
||||||
|
newbuilding_enrich, cian_history_backfill — нет.
|
||||||
|
|
||||||
|
Побочно это чинит и `heartbeat_at`: он писался тем же `now()` и по той же причине
|
||||||
|
отставал от реальности на возраст открытой транзакции, а на нём стоит поиск зависших
|
||||||
|
прогонов (reap_zombies, порог 6 ч).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -323,7 +348,10 @@ def _alert_on_run_id(
|
||||||
|
|
||||||
|
|
||||||
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
"""INSERT scrape_runs(source, status='running', params, started_at=NOW()).
|
"""INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()).
|
||||||
|
|
||||||
|
started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей
|
||||||
|
транзакции задачи его уже не достаёт (#2702).
|
||||||
|
|
||||||
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
||||||
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
||||||
|
|
@ -334,7 +362,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
||||||
VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW())
|
VALUES (
|
||||||
|
:source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp()
|
||||||
|
)
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
|
|
@ -346,7 +376,7 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
|
|
||||||
|
|
||||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
"""UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
|
||||||
|
|
||||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||||
|
|
@ -357,7 +387,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET heartbeat_at = NOW(),
|
SET heartbeat_at = clock_timestamp(),
|
||||||
counters = CAST(:counters AS jsonb),
|
counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -405,7 +435,7 @@ def is_cancelled(db: Session, run_id: int) -> bool:
|
||||||
|
|
||||||
|
|
||||||
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
"""Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count.
|
"""Финализация run: status='done', finished_at + counters + total_seen/new_count.
|
||||||
|
|
||||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
||||||
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
||||||
|
|
@ -415,7 +445,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'done',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
counters = CAST(:counters AS jsonb),
|
counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -454,7 +485,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'failed',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
error = :error, counters = CAST(:counters AS jsonb),
|
error = :error, counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -513,7 +545,8 @@ def mark_banned(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'banned',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
error = :error, counters = CAST(:counters AS jsonb),
|
error = :error, counters = CAST(:counters AS jsonb),
|
||||||
ban_kind = :ban_kind,
|
ban_kind = :ban_kind,
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
|
|
@ -624,7 +657,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'cancelled', finished_at = NOW()
|
SET status = 'cancelled', finished_at = clock_timestamp()
|
||||||
WHERE id = :run_id AND status = 'running'
|
WHERE id = :run_id AND status = 'running'
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,71 @@
|
||||||
|
-- 223_scrape_runs_time_columns_meaning.sql
|
||||||
|
-- Purpose (#2702): зафиксировать в схеме, что отметки времени прогона до этой
|
||||||
|
-- правки не охватывали его работу, и куда смотреть аналитике вместо разности.
|
||||||
|
--
|
||||||
|
-- Dependencies: 015_scrape_runs.sql (создала started_at/finished_at/heartbeat_at),
|
||||||
|
-- 051_scrape_runs_extend.sql (finished_at/counters).
|
||||||
|
-- Apply after: 221_backfill_house_suggestions_image_link.sql
|
||||||
|
-- Идемпотентно: только COMMENT ON COLUMN (перезаписывает сам себя), данных не трогает.
|
||||||
|
--
|
||||||
|
-- ── ЧТО БЫЛО СЛОМАНО ─────────────────────────────────────────────────────────
|
||||||
|
-- Финализаторы (mark_done / mark_failed / mark_banned) писали finished_at и
|
||||||
|
-- heartbeat_at через now(). В PostgreSQL now() == transaction_timestamp(): он
|
||||||
|
-- замерзает на СТАРТЕ транзакции. Финализатор выполняется той же сессией, что и
|
||||||
|
-- работа задачи; если рабочая транзакция всё это время оставалась открытой (задаче
|
||||||
|
-- нечего было коммитить — читающий батч, ноль сохранений, сохранение чужой сессией),
|
||||||
|
-- UPDATE финализатора попадал ВНУТРЬ неё и получал время НАЧАЛА работы.
|
||||||
|
--
|
||||||
|
-- Прод-замер 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик
|
||||||
|
-- counters.duration_sec):
|
||||||
|
-- 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза;
|
||||||
|
-- 133 — окно меньше секунды при работе дольше 10 с.
|
||||||
|
-- 126 из этих 133 окон лежат в 9-64 мс: это не разброс, а подпись механизма —
|
||||||
|
-- столько проходит от коммита claim'а до первого запроса рабочей транзакции.
|
||||||
|
-- Крайние: прогон 346 (cian_history_backfill) — 18 230 с работы при окне 32 мс;
|
||||||
|
-- 497 (newbuilding_enrich) — 6 124 с при 21 мс; 341 (yandex_address_backfill) —
|
||||||
|
-- 1 460 с при 19 мс.
|
||||||
|
--
|
||||||
|
-- Дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом.
|
||||||
|
-- Средние окно/duration_sec по источникам на том же замере:
|
||||||
|
-- cian_history_backfill 2554 / 4222 ← окно короче работы
|
||||||
|
-- newbuilding_enrich 1302 / 2078 ← короче
|
||||||
|
-- yandex_address_backfill 107 / 1022 ← короче в 10 раз
|
||||||
|
-- avito_detail_backfill 2829 / 2016 ← длиннее (норма)
|
||||||
|
-- house_imv_backfill 1173 / 1173 ← совпадает
|
||||||
|
-- cadastral_geo_match 11 / 10 ← совпадает
|
||||||
|
--
|
||||||
|
-- ── ПОЧЕМУ ИСТОРИЮ НЕ ЧИНИМ ──────────────────────────────────────────────────
|
||||||
|
-- Восстановить настоящий finished_at по строке нельзя: реальное время конца нигде
|
||||||
|
-- не сохранилось. Но counters.duration_sec измерялся монотонными часами процесса
|
||||||
|
-- (time.monotonic / time.time в самих задачах) и транзакцией не затронут — он у
|
||||||
|
-- этих строк верный. Поэтому история не переписывается, а помечается: аналитика
|
||||||
|
-- обязана брать длительность из счётчика, а не из разности отметок.
|
||||||
|
--
|
||||||
|
-- Правка кода (clock_timestamp() вместо now() во всех финализаторах и в heartbeat)
|
||||||
|
-- живёт в app/services/scrape_runs.py + packages/scraper-kit/.../orchestration/runs.py.
|
||||||
|
|
||||||
|
COMMENT ON COLUMN scrape_runs.started_at IS
|
||||||
|
'Момент claim''а прогона. Пишется create_run своей транзакцией (commit сразу '
|
||||||
|
'после INSERT), поэтому откат рабочей транзакции задачи его не затрагивает. '
|
||||||
|
'С #2702 — clock_timestamp(); до него now() (= старт транзакции тика планировщика), '
|
||||||
|
'что давало сдвиг в пределах тика.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN scrape_runs.finished_at IS
|
||||||
|
'Момент финализации прогона. ВНИМАНИЕ: у строк ДО #2702 (2026-08-06) значение '
|
||||||
|
'недостоверно — писалось now() (= transaction_timestamp) внутри рабочей транзакции '
|
||||||
|
'задачи, поэтому у прогонов, ничего не коммитивших по ходу работы, равно времени '
|
||||||
|
'её НАЧАЛА. На проде так вышло у 133 из 487 прогонов со счётчиком длительности '
|
||||||
|
'(окно < 1 с при работе > 10 с). Длительность таких прогонов брать из '
|
||||||
|
'counters->>''duration_sec'' (монотонные часы процесса, транзакцией не затронуты), '
|
||||||
|
'а НЕ из finished_at − started_at.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN scrape_runs.heartbeat_at IS
|
||||||
|
'Последний признак жизни прогона; на нём стоит поиск зависших (reap_zombies, порог '
|
||||||
|
'6 ч). У строк ДО #2702 отставал от реальности на возраст открытой рабочей '
|
||||||
|
'транзакции по той же причине, что finished_at, — то есть критерий «завис» решал '
|
||||||
|
'по замороженной отметке. С #2702 пишется clock_timestamp().';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN scrape_runs.counters IS
|
||||||
|
'Счётчики прогона (jsonb). Ключ duration_sec, где он есть, измерен монотонными '
|
||||||
|
'часами процесса и остаётся единственным достоверным источником длительности для '
|
||||||
|
'строк до #2702 (см. комментарий к finished_at).';
|
||||||
83
tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql
Normal file
83
tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql
Normal file
|
|
@ -0,0 +1,83 @@
|
||||||
|
-- 224_houses_house_type_canon.sql
|
||||||
|
-- Issue #2678 (хвост #2675/#2674): twin миграции 141 для таблицы ДОМОВ.
|
||||||
|
--
|
||||||
|
-- Миграция 141 привела camelCase-вокабуляр Циана к канону только в listings.
|
||||||
|
-- В houses он остался — и каждый читатель типа дома чинил его у себя (#2675
|
||||||
|
-- починил домовую оценку; поштучный путь и подбор аналогов продолжали сравнивать
|
||||||
|
-- 'monolithBrick' с 'monolith_brick' и не совпадать).
|
||||||
|
--
|
||||||
|
-- ЗАМЕР ПРОДА 2026-08-06 (SELECT source, house_type, count(*) FROM houses GROUP BY 1,2):
|
||||||
|
-- канон: brick 388+5+4+1 · panel 325+3+2+1 · monolith 118+37+6+1 ·
|
||||||
|
-- block 90+1+1 · monolith_brick 2+1 · wood 2
|
||||||
|
-- camelCase: monolithBrick 48 (derived) + 8 (cian_newbuilding) = 56 ·
|
||||||
|
-- gasSilicateBlock 1 · aerocreteBlock 1
|
||||||
|
-- прочее: stalin 3 · other 18 · wireframe 1
|
||||||
|
-- NULL: 8440 из 8880 строк (тип дома вообще неизвестен — не наш случай)
|
||||||
|
--
|
||||||
|
-- ЖИВОГО ПИСАТЕЛЯ camelCase В houses НЕТ: у всех 80 неканоничных строк
|
||||||
|
-- last_scraped_at = 2026-05-24 14:04:20.012209 — одна и та же метка, т.е.
|
||||||
|
-- единственный прогон backfill'а 063 (промоут типов из listings ДО миграции 141).
|
||||||
|
-- Единственный живой писатель houses.house_type — avito-каталог домов
|
||||||
|
-- (providers/avito/houses.py), он пишет русские подписи через свою карту; в том
|
||||||
|
-- же PR он переведён на общий normalize_house_type, чтобы неизвестное значение
|
||||||
|
-- шло как NULL, а не как 'other' (его единственный источник неканона).
|
||||||
|
--
|
||||||
|
-- ПРОВЕРКА СМЫСЛА ПЕРЕД СКЛЕЙКОЙ (требование #2674 — не слепить разное):
|
||||||
|
-- контрольная группа в своих же данных. Для каждой неканоничной строки взяты
|
||||||
|
-- типы её ЖЕ объявлений (listings.house_id_fk), уже нормализованных 141:
|
||||||
|
-- monolithBrick 56 домов — monolith_brick присутствует у ВСЕХ 56 → одно и то же
|
||||||
|
-- stalin 3 дома — brick (совпадает с решением 141: «сталинка» = кирпич)
|
||||||
|
-- aerocreteBlock 1 дом — block
|
||||||
|
-- gasSilicateBlock 1 дом — block
|
||||||
|
-- other 18 домов — разброс monolith/monolith_brick/brick/block, т.е.
|
||||||
|
-- 'other' = «неизвестно», а не отдельный материал
|
||||||
|
-- wireframe 1 дом — wireframe и у объявлений (само-согласовано)
|
||||||
|
-- Вывод: склейка безопасна ТОЛЬКО для четырёх camelCase-токенов + stalin.
|
||||||
|
--
|
||||||
|
-- ЧТО НАМЕРЕННО НЕ ТРОГАЕМ:
|
||||||
|
-- 'other' (18) и 'wireframe' (1) — честного соответствия в каноне нет
|
||||||
|
-- (см. #2675: normalize_house_type схлопывает их в None на чтении, и это
|
||||||
|
-- правильный ответ — NULL нейтрален для soft-penalty эстиматора, а выдуманный
|
||||||
|
-- материал был бы враньём). Стирать их здесь тоже не будем: это единственный
|
||||||
|
-- след того, что источник что-то про дом сказал.
|
||||||
|
--
|
||||||
|
-- BACKFILL (счётчики сняты на проде ДО применения, 2026-08-06):
|
||||||
|
-- monolithBrick -> monolith_brick : 56 строк
|
||||||
|
-- stalin -> brick : 3 строки
|
||||||
|
-- aerocreteBlock -> block : 1 строка
|
||||||
|
-- gasSilicateBlock -> block : 1 строка
|
||||||
|
-- foamConcreteBlock-> block : 0 строк (в houses не встречается,
|
||||||
|
-- оставлен для паритета с картой 141)
|
||||||
|
-- ИТОГО ожидаемо тронуто: 61 строка.
|
||||||
|
--
|
||||||
|
-- Idempotent: WHERE перечисляет только мапимые токены → повторный прогон 0 строк.
|
||||||
|
-- Маппинг тождественен house_type_normalizer._RAW_TO_CANON и миграции 141 —
|
||||||
|
-- третьего словаря не заводим.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
UPDATE houses
|
||||||
|
SET house_type = CASE house_type
|
||||||
|
WHEN 'monolithBrick' THEN 'monolith_brick'
|
||||||
|
WHEN 'gasSilicateBlock' THEN 'block'
|
||||||
|
WHEN 'aerocreteBlock' THEN 'block'
|
||||||
|
WHEN 'foamConcreteBlock' THEN 'block'
|
||||||
|
WHEN 'stalin' THEN 'brick'
|
||||||
|
ELSE house_type
|
||||||
|
END
|
||||||
|
WHERE house_type IN (
|
||||||
|
'monolithBrick', 'gasSilicateBlock', 'aerocreteBlock',
|
||||||
|
'foamConcreteBlock', 'stalin'
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON COLUMN houses.house_type IS
|
||||||
|
'Материал/тип дома, канон: panel/brick/monolith/monolith_brick/block/wood '
|
||||||
|
'(тот же enum, что listings.house_type и scraper_kit.house_type_normalizer). '
|
||||||
|
'Писать сюда только через normalize_house_type — сырые вокабуляры источников '
|
||||||
|
'(cian camelCase monolithBrick/gasSilicateBlock/stalin, yandex SCREAMING '
|
||||||
|
'MONOLIT_BRICK, русские подписи Авито) приводятся ДО записи, миграция 224 '
|
||||||
|
'вычистила исторические. Вне канона осталось намеренно: other (источник сказал '
|
||||||
|
'«другое») и wireframe (каркас — материала в каноне нет). Неизвестный тип = '
|
||||||
|
'NULL, а не панель и не other: NULL нейтрален для soft-penalty эстиматора.';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
128
tradein-mvp/backend/tests/test_2678_house_type_canon.py
Normal file
128
tradein-mvp/backend/tests/test_2678_house_type_canon.py
Normal file
|
|
@ -0,0 +1,128 @@
|
||||||
|
"""#2678: тип дома приводится к канону У ИСТОЧНИКА, и словарь ровно один.
|
||||||
|
|
||||||
|
Кейсы взяты не из головы, а из фактического замера прода 2026-08-06
|
||||||
|
(`SELECT source, house_type, count(*) FROM houses GROUP BY 1,2`):
|
||||||
|
|
||||||
|
monolithBrick 56 · other 18 · stalin 3 · aerocreteBlock 1 ·
|
||||||
|
gasSilicateBlock 1 · wireframe 1 · плюс канон (brick/panel/monolith/
|
||||||
|
monolith_brick/block/wood) и 8440 NULL.
|
||||||
|
|
||||||
|
Проверяется три вещи:
|
||||||
|
1. каждый фактический вариант → канон (или честный None);
|
||||||
|
2. живой писатель houses.house_type (avito-каталог) больше не изобретает
|
||||||
|
'other' и ходит через общий нормализатор;
|
||||||
|
3. миграция 224 не заводит третий словарь — её CASE совпадает с картой кода.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import re
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from scraper_kit.house_type_normalizer import _RAW_TO_CANON, normalize_house_type
|
||||||
|
from scraper_kit.providers.avito.houses import _normalize_house_type as avito_house_type
|
||||||
|
|
||||||
|
_MIGRATION_224 = (
|
||||||
|
Path(__file__).resolve().parents[1] / "data" / "sql" / "224_houses_house_type_canon.sql"
|
||||||
|
)
|
||||||
|
|
||||||
|
# Фактический словарь houses.house_type на проде 2026-08-06 → чем он обязан стать.
|
||||||
|
# None = «честно неизвестно» (NULL нейтрален для soft-penalty эстиматора, в отличие
|
||||||
|
# от выдуманного материала).
|
||||||
|
_PROD_VALUES: list[tuple[str, str | None]] = [
|
||||||
|
("monolithBrick", "monolith_brick"), # 56 строк
|
||||||
|
("other", None), # 18 строк — источник сказал «другое», материала нет
|
||||||
|
("stalin", "brick"), # 3 строки — «сталинка» = кирпич (решение миграции 141)
|
||||||
|
("aerocreteBlock", "block"), # 1 строка
|
||||||
|
("gasSilicateBlock", "block"), # 1 строка
|
||||||
|
("wireframe", None), # 1 строка — каркас, в каноне такого материала нет
|
||||||
|
("brick", "brick"),
|
||||||
|
("panel", "panel"),
|
||||||
|
("monolith", "monolith"),
|
||||||
|
("monolith_brick", "monolith_brick"),
|
||||||
|
("block", "block"),
|
||||||
|
("wood", "wood"),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(("raw", "expected"), _PROD_VALUES)
|
||||||
|
def test_prod_value_maps_to_canon(raw: str, expected: str | None) -> None:
|
||||||
|
assert normalize_house_type(raw) == expected
|
||||||
|
|
||||||
|
|
||||||
|
def test_canon_survives_uppercase_including_monolith_brick() -> None:
|
||||||
|
"""#2678 п.6: сквозной проброс канона был регистрозависим — кроме monolith_brick.
|
||||||
|
|
||||||
|
Значений в верхнем регистре в базе сегодня ноль; это страховка на новый источник,
|
||||||
|
который отдаст канон «как в документации».
|
||||||
|
"""
|
||||||
|
assert normalize_house_type("MONOLITH_BRICK") == "monolith_brick"
|
||||||
|
assert normalize_house_type("Monolith_Brick") == "monolith_brick"
|
||||||
|
# Остальной канон и раньше переживал регистр — фиксируем, что не сломали.
|
||||||
|
for token in ("BRICK", "Panel", "MONOLITH", "Block", "WOOD"):
|
||||||
|
assert normalize_house_type(token) == token.lower()
|
||||||
|
|
||||||
|
|
||||||
|
# ── живой писатель houses.house_type: avito-каталог домов ────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("label", "expected"),
|
||||||
|
[
|
||||||
|
("Монолитно-кирпичный", "monolith_brick"),
|
||||||
|
("Панельный", "panel"),
|
||||||
|
("КИРПИЧНЫЙ", "brick"),
|
||||||
|
(" Блочный ", "block"),
|
||||||
|
("Деревянный", "wood"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_avito_house_label_maps_to_canon(label: str, expected: str) -> None:
|
||||||
|
assert avito_house_type(label) == expected
|
||||||
|
|
||||||
|
|
||||||
|
def test_avito_unknown_label_is_null_not_other() -> None:
|
||||||
|
"""Незнакомая подпись → NULL. До #2678 здесь появлялось 'other'.
|
||||||
|
|
||||||
|
'other' всегда != канону, т.е. читатель получал не «неизвестно», а гарантированное
|
||||||
|
несовпадение: ложный штраф при подборе аналогов и пропуск оценки.
|
||||||
|
"""
|
||||||
|
assert avito_house_type("Саманный") is None
|
||||||
|
assert avito_house_type("") is None
|
||||||
|
assert avito_house_type(None) is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_avito_writer_handles_foreign_vocabulary() -> None:
|
||||||
|
"""Писатель ходит через общий нормализатор, а не только через свою карту."""
|
||||||
|
assert avito_house_type("monolithBrick") == "monolith_brick"
|
||||||
|
assert avito_house_type("MONOLIT_BRICK") == "monolith_brick"
|
||||||
|
|
||||||
|
|
||||||
|
# ── миграция 224: тот же словарь, что в коде ─────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _migration_case_pairs() -> dict[str, str]:
|
||||||
|
"""WHEN 'x' THEN 'y' из исполняемой части миграции (без `--`-комментариев)."""
|
||||||
|
code = "\n".join(
|
||||||
|
line.split("--", 1)[0] for line in _MIGRATION_224.read_text(encoding="utf-8").splitlines()
|
||||||
|
)
|
||||||
|
return dict(re.findall(r"WHEN\s+'([^']+)'\s+THEN\s+'([^']+)'", code))
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_224_mapping_matches_code() -> None:
|
||||||
|
"""Миграция не заводит третий словарь — каждая пара есть в _RAW_TO_CANON."""
|
||||||
|
pairs = _migration_case_pairs()
|
||||||
|
assert pairs, "в миграции 224 не нашлось ни одного WHEN ... THEN"
|
||||||
|
for raw, canon in pairs.items():
|
||||||
|
assert _RAW_TO_CANON.get(raw) == canon, f"{raw!r} расходится с house_type_normalizer"
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_224_touches_only_mapped_tokens() -> None:
|
||||||
|
"""WHERE ограничен теми же токенами → 'other'/'wireframe'/канон не трогаются."""
|
||||||
|
code = "\n".join(
|
||||||
|
line.split("--", 1)[0] for line in _MIGRATION_224.read_text(encoding="utf-8").splitlines()
|
||||||
|
)
|
||||||
|
where_tokens = set(re.findall(r"'([A-Za-z]+)'", code.split("WHERE", 1)[1].split(";", 1)[0]))
|
||||||
|
assert where_tokens == set(_migration_case_pairs())
|
||||||
|
assert "other" not in where_tokens
|
||||||
|
assert "wireframe" not in where_tokens
|
||||||
237
tradein-mvp/backend/tests/test_2702_run_timestamps.py
Normal file
237
tradein-mvp/backend/tests/test_2702_run_timestamps.py
Normal file
|
|
@ -0,0 +1,237 @@
|
||||||
|
"""#2702: отметки времени прогона не охватывают его работу.
|
||||||
|
|
||||||
|
Что было. Финализаторы писали `finished_at`/`heartbeat_at` через `now()`, а `now()`
|
||||||
|
в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции.
|
||||||
|
Выполняются финализаторы той же сессией, что и работа задачи, поэтому если рабочая
|
||||||
|
транзакция всё это время оставалась открытой, их UPDATE попадал ВНУТРЬ неё и получал
|
||||||
|
время НАЧАЛА работы.
|
||||||
|
|
||||||
|
Прод-замер 2026-08-06 (487 прогонов с finished_at и counters.duration_sec):
|
||||||
|
* 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза;
|
||||||
|
* 133 — окно меньше секунды при работе дольше 10 с, причём 126 из них укладываются
|
||||||
|
в 9-64 мс: столько проходит от коммита claim'а до первого запроса рабочей
|
||||||
|
транзакции — подпись механизма, а не разброс.
|
||||||
|
Крайний случай, воспроизведённый ниже дословно: прогон 346 (cian_history_backfill) —
|
||||||
|
18 230 с работы, окно 32 мс.
|
||||||
|
|
||||||
|
Почему дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом.
|
||||||
|
`cadastral_geo_match` / `house_imv_backfill` / `avito_detail_backfill` коммитят
|
||||||
|
поштучно — у них окно совпадало с работой; `yandex_address_backfill` (45 из 50
|
||||||
|
прогонов), `newbuilding_enrich`, `cian_history_backfill` — нет.
|
||||||
|
|
||||||
|
Фальсификация: на старом коде (`now()`) тесты 1 и 3 дают другой ответ.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import inspect
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||||
|
|
||||||
|
from scraper_kit.orchestration import runs as kit_runs
|
||||||
|
from scraper_kit.orchestration import scheduler as kit_scheduler
|
||||||
|
|
||||||
|
from app.services import scrape_runs as app_runs
|
||||||
|
|
||||||
|
_MODULES = {"kit": kit_runs, "app": app_runs}
|
||||||
|
|
||||||
|
# Колонки, которые обязаны нести НАСТОЯЩЕЕ время, а не время старта транзакции.
|
||||||
|
_TS_COLS = ("started_at", "finished_at", "heartbeat_at")
|
||||||
|
|
||||||
|
|
||||||
|
def _split_top(items: str) -> list[str]:
|
||||||
|
"""Разбить список SQL-элементов по запятым ВЕРХНЕГО уровня (CAST(:x AS t) — один)."""
|
||||||
|
out: list[str] = []
|
||||||
|
depth = 0
|
||||||
|
cur = ""
|
||||||
|
for ch in items:
|
||||||
|
if ch == "," and depth == 0:
|
||||||
|
out.append(cur.strip())
|
||||||
|
cur = ""
|
||||||
|
continue
|
||||||
|
depth += (ch == "(") - (ch == ")")
|
||||||
|
cur += ch
|
||||||
|
out.append(cur.strip())
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
class _Row:
|
||||||
|
"""Строка ответа: id для create_run, source для alert-хука."""
|
||||||
|
|
||||||
|
id = 1
|
||||||
|
source = "src"
|
||||||
|
status = "done"
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeResult:
|
||||||
|
def fetchone(self) -> _Row:
|
||||||
|
return _Row()
|
||||||
|
|
||||||
|
def first(self) -> _Row:
|
||||||
|
return _Row()
|
||||||
|
|
||||||
|
def fetchall(self) -> list[_Row]:
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
class _FakePg:
|
||||||
|
"""Мини-модель PostgreSQL на две функции времени и ленивую транзакцию.
|
||||||
|
|
||||||
|
`now()` == `transaction_timestamp()` — замерзает на старте транзакции;
|
||||||
|
`clock_timestamp()` — настоящие часы. Транзакция открывается лениво на первом
|
||||||
|
execute (autobegin SQLAlchemy) и закрывается commit/rollback. Больше модель
|
||||||
|
ничего не умеет — этого достаточно, чтобы отличить одно от другого.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.wall: float = 0.0 # «стенные часы» теста
|
||||||
|
self.tx_start: float | None = None
|
||||||
|
self.row: dict[str, float] = {} # что осело в scrape_runs
|
||||||
|
self.calls: list[str] = []
|
||||||
|
|
||||||
|
def _stamp(self, col: str, func: str) -> None:
|
||||||
|
assert self.tx_start is not None
|
||||||
|
self.row[col] = self.tx_start if func.lower() == "now" else self.wall
|
||||||
|
|
||||||
|
def execute(self, stmt: Any, params: Any = None) -> _FakeResult:
|
||||||
|
if self.tx_start is None:
|
||||||
|
self.tx_start = self.wall
|
||||||
|
self.calls.append("execute")
|
||||||
|
sql = str(stmt)
|
||||||
|
for col in _TS_COLS: # UPDATE ... SET <col> = <func>()
|
||||||
|
m = re.search(rf"\b{col}\s*=\s*(now|clock_timestamp)\s*\(\s*\)", sql, re.I)
|
||||||
|
if m is not None:
|
||||||
|
self._stamp(col, m.group(1))
|
||||||
|
ins = re.search(
|
||||||
|
r"INSERT INTO scrape_runs\s*\((.*?)\).*?VALUES\s*\((.*)\)", sql, re.S | re.I
|
||||||
|
)
|
||||||
|
if ins is not None: # INSERT — отметки времени позиционные, в VALUES
|
||||||
|
for col, val in zip(_split_top(ins.group(1)), _split_top(ins.group(2)), strict=False):
|
||||||
|
f = re.fullmatch(r"(now|clock_timestamp)\s*\(\s*\)", val, re.I)
|
||||||
|
if col in _TS_COLS and f is not None:
|
||||||
|
self._stamp(col, f.group(1))
|
||||||
|
return _FakeResult()
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self.calls.append("commit")
|
||||||
|
self.tx_start = None
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
self.calls.append("rollback")
|
||||||
|
self.tx_start = None
|
||||||
|
|
||||||
|
|
||||||
|
# ── 1. Финал прогона датируется концом работы, а не стартом транзакции ────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("name", list(_MODULES))
|
||||||
|
def test_finished_at_covers_the_work_not_the_transaction_start(name: str) -> None:
|
||||||
|
"""Прогон 346 дословно: 18 230 с работы в одной незакоммиченной транзакции.
|
||||||
|
|
||||||
|
На старом коде finished_at = 0.032 (старт рабочей транзакции) → окно 32 мс при
|
||||||
|
пяти часах работы. Это и есть та строка, ради которой заведена задача.
|
||||||
|
"""
|
||||||
|
mod = _MODULES[name]
|
||||||
|
db = _FakePg()
|
||||||
|
db.row["started_at"] = 0.0 # create_run уже закоммитил claim
|
||||||
|
|
||||||
|
db.wall = 0.020
|
||||||
|
mod.update_heartbeat(db, 1, {"listings_processed": 0}) # heartbeat + commit
|
||||||
|
|
||||||
|
db.wall = 0.032
|
||||||
|
db.execute("SELECT id FROM listings WHERE history IS NULL") # рабочая транзакция
|
||||||
|
|
||||||
|
db.wall = 18230.0 # пять часов работы, ни одного коммита
|
||||||
|
mod.mark_done(db, 1, {"listings_processed": 1200, "duration_sec": 18230})
|
||||||
|
|
||||||
|
assert db.row["finished_at"] == pytest.approx(18230.0)
|
||||||
|
assert db.row["heartbeat_at"] == pytest.approx(18230.0)
|
||||||
|
window = db.row["finished_at"] - db.row["started_at"]
|
||||||
|
assert window == pytest.approx(18230.0), "окно прогона обязано охватывать его работу"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("name", list(_MODULES))
|
||||||
|
@pytest.mark.parametrize("finalizer", ["mark_failed", "mark_banned"])
|
||||||
|
def test_failed_and_banned_finals_are_wall_clock_too(name: str, finalizer: str) -> None:
|
||||||
|
"""Тот же инвариант для неуспешных финалов.
|
||||||
|
|
||||||
|
Их спасал defensive-rollback в начале (он закрывал рабочую транзакцию), но
|
||||||
|
полагаться на побочный эффект чужой защиты нельзя — проверяем явно.
|
||||||
|
"""
|
||||||
|
mod = _MODULES[name]
|
||||||
|
db = _FakePg()
|
||||||
|
db.wall = 0.019
|
||||||
|
db.execute("SELECT 1") # рабочая транзакция открыта
|
||||||
|
db.wall = 1460.0
|
||||||
|
getattr(mod, finalizer)(db, 1, "boom", {"checked": 0, "duration_sec": 1460})
|
||||||
|
assert db.row["finished_at"] == pytest.approx(1460.0)
|
||||||
|
|
||||||
|
|
||||||
|
# ── 2. started_at переживает откат рабочей транзакции ────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("name", list(_MODULES))
|
||||||
|
def test_started_at_survives_rolled_back_work_transaction(name: str) -> None:
|
||||||
|
"""Требование #2702 п.1: отметка старта живёт в СВОЕЙ закоммиченной транзакции.
|
||||||
|
|
||||||
|
create_run коммитит INSERT до возврата run_id, а ни один финализатор прогона
|
||||||
|
started_at не переписывает — поэтому откат рабочей транзакции его не достаёт.
|
||||||
|
Финал при этом обязан быть датирован концом работы (это и падает на старом коде).
|
||||||
|
"""
|
||||||
|
mod = _MODULES[name]
|
||||||
|
db = _FakePg()
|
||||||
|
|
||||||
|
run_id = mod.create_run(db, source="cian_history_backfill", params={})
|
||||||
|
assert run_id == 1
|
||||||
|
assert db.calls == ["execute", "commit"], "INSERT прогона обязан коммититься сразу"
|
||||||
|
|
||||||
|
db.wall = 0.030
|
||||||
|
db.execute("UPDATE listings SET address = 'x'") # рабочая транзакция
|
||||||
|
db.wall = 100.0
|
||||||
|
db.rollback() # работа упала и откатилась
|
||||||
|
|
||||||
|
db.wall = 100.5
|
||||||
|
db.execute("SELECT count(*) FROM listings") # новая рабочая транзакция
|
||||||
|
db.wall = 1460.0
|
||||||
|
mod.mark_done(db, 1, {"checked": 0, "duration_sec": 1460})
|
||||||
|
|
||||||
|
assert db.row["started_at"] == pytest.approx(0.0), "started_at не должен сдвигаться"
|
||||||
|
assert db.row["finished_at"] == pytest.approx(1460.0)
|
||||||
|
|
||||||
|
|
||||||
|
# ── 3. Инвариант источника: никакая отметка времени не пишется now() ─────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("name", list(_MODULES))
|
||||||
|
def test_no_run_timestamp_is_written_with_now(name: str) -> None:
|
||||||
|
"""`now()` в этих модулях не имеет корректного применения — его быть не должно.
|
||||||
|
|
||||||
|
Проверяем весь исходник, а не отдельные запросы: INSERT в create_run пишет
|
||||||
|
started_at/heartbeat_at позиционно (в VALUES), и построчная проверка его бы
|
||||||
|
пропустила — ровно так дефект и дожил до 3 300 прогонов.
|
||||||
|
"""
|
||||||
|
src = inspect.getsource(_MODULES[name])
|
||||||
|
# Регистрозависимо: SQL в этих модулях пишется в верхнем регистре, а строчное
|
||||||
|
# `now()` встречается в объяснительной прозе docstring'ов — ловим SQL, не текст.
|
||||||
|
assert re.search(r"\bNOW\s*\(\s*\)", src) is None
|
||||||
|
assert "clock_timestamp()" in src
|
||||||
|
|
||||||
|
|
||||||
|
def test_zombie_criterion_compares_real_clocks() -> None:
|
||||||
|
"""#2702 п.2: поиск зависших сравнивает записанный heartbeat со «сейчас».
|
||||||
|
|
||||||
|
Обе стороны сравнения обязаны быть настоящим временем: на проде у всех 6
|
||||||
|
прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и
|
||||||
|
остался на отметке старта (max advance 0.0 с) — критерий решал по замороженной
|
||||||
|
отметке, хотя нормальный прогон этого источника длится до 5.06 ч.
|
||||||
|
"""
|
||||||
|
src = inspect.getsource(kit_scheduler.reap_zombies)
|
||||||
|
stmt = re.search(r"UPDATE scrape_runs.*?RETURNING id", src, re.S)
|
||||||
|
assert stmt is not None
|
||||||
|
assert re.search(r"\bNOW\s*\(\s*\)", stmt.group(0)) is None
|
||||||
|
assert stmt.group(0).count("clock_timestamp()") == 2 # finished_at + порог сравнения
|
||||||
|
|
@ -129,7 +129,9 @@ def test_mark_skipped_collapse_refreshes_detail_and_started_at() -> None:
|
||||||
update = db.sql_of("UPDATE scrape_runs r")
|
update = db.sql_of("UPDATE scrape_runs r")
|
||||||
assert update is not None
|
assert update is not None
|
||||||
sql, params = update
|
sql, params = update
|
||||||
assert "started_at = NOW()" in sql
|
# clock_timestamp(), а не NOW(): отметки времени прогона пишутся настоящими
|
||||||
|
# часами, иначе внутри долгой открытой транзакции они замерзают (#2702).
|
||||||
|
assert "started_at = clock_timestamp()" in sql
|
||||||
assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска"
|
assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска"
|
||||||
assert "first_skip_at" in sql, "начало стрика потеряно"
|
assert "first_skip_at" in sql, "начало стрика потеряно"
|
||||||
assert params["details"] == "37 дн. назад"
|
assert params["details"] == "37 дн. назад"
|
||||||
|
|
|
||||||
|
|
@ -125,6 +125,11 @@ FETCH_JSON_SETTLE_MS: int = int(os.environ.get("FETCH_JSON_SETTLE_MS", "1200"))
|
||||||
FETCH_JSON_INPAGE_RETRIES: int = int(os.environ.get("FETCH_JSON_INPAGE_RETRIES", "1"))
|
FETCH_JSON_INPAGE_RETRIES: int = int(os.environ.get("FETCH_JSON_INPAGE_RETRIES", "1"))
|
||||||
# Пауза между in-page попытками fetch(), мс.
|
# Пауза между in-page попытками fetch(), мс.
|
||||||
FETCH_JSON_RETRY_DELAY_MS: int = int(os.environ.get("FETCH_JSON_RETRY_DELAY_MS", "800"))
|
FETCH_JSON_RETRY_DELAY_MS: int = int(os.environ.get("FETCH_JSON_RETRY_DELAY_MS", "800"))
|
||||||
|
# Сколько ждать события `load` на ПОВТОРЕ после гонки «execution context destroyed»
|
||||||
|
# (#2676). Только на повторе: happy-path остаётся на дешёвом FETCH_JSON_SETTLE_MS,
|
||||||
|
# иначе бесконечно дозагружающаяся страница удлиняла бы КАЖДЫЙ запрос. Ожидание
|
||||||
|
# best-effort — по таймауту всё равно пробуем evaluate.
|
||||||
|
FETCH_JSON_LOAD_WAIT_MS: int = int(os.environ.get("FETCH_JSON_LOAD_WAIT_MS", "15000"))
|
||||||
|
|
||||||
# Известные поставщики. "generic" — фолбэк для всех прочих хостов (один общий
|
# Известные поставщики. "generic" — фолбэк для всех прочих хостов (один общий
|
||||||
# инстанс на неузнанные домены). Порядок задаёт детерминированный health-вывод.
|
# инстанс на неузнанные домены). Порядок задаёт детерминированный health-вывод.
|
||||||
|
|
@ -1000,6 +1005,27 @@ async def _do_fetch_json(
|
||||||
return await _fetch_json_once(
|
return await _fetch_json_once(
|
||||||
provider, url, method=method, headers=headers, body=body, origin=origin
|
provider, url, method=method, headers=headers, body=body, origin=origin
|
||||||
)
|
)
|
||||||
|
if _is_page_context_lost(exc):
|
||||||
|
# #2676: браузер жив, умерла ОДНА страница — relaunch не нужен (стоил бы
|
||||||
|
# ~10-20с и тёплые cookies инстанса). Повторяем на свежей странице, но с
|
||||||
|
# ожиданием `load`: без него повтор попадает в то же окно клиентской
|
||||||
|
# навигации, и «транзиентная» ошибка воспроизводится детерминированно.
|
||||||
|
logger.warning(
|
||||||
|
"tradein-browser[%s]: страница ушла в навигацию (%s), retry fetch-json "
|
||||||
|
"с ожиданием load: %s",
|
||||||
|
provider,
|
||||||
|
type(exc).__name__,
|
||||||
|
url,
|
||||||
|
)
|
||||||
|
return await _fetch_json_once(
|
||||||
|
provider,
|
||||||
|
url,
|
||||||
|
method=method,
|
||||||
|
headers=headers,
|
||||||
|
body=body,
|
||||||
|
origin=origin,
|
||||||
|
wait_for_load=True,
|
||||||
|
)
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -1011,6 +1037,7 @@ async def _fetch_json_once(
|
||||||
headers: dict,
|
headers: dict,
|
||||||
body: object,
|
body: object,
|
||||||
origin: str,
|
origin: str,
|
||||||
|
wait_for_load: bool = False,
|
||||||
) -> dict:
|
) -> dict:
|
||||||
"""Переходит на origin (same-origin якорь) и выполняет in-page fetch(url).
|
"""Переходит на origin (same-origin якорь) и выполняет in-page fetch(url).
|
||||||
|
|
||||||
|
|
@ -1029,6 +1056,22 @@ async def _fetch_json_once(
|
||||||
await _apply_resource_block(page)
|
await _apply_resource_block(page)
|
||||||
await _pace_provider(provider)
|
await _pace_provider(provider)
|
||||||
await page.goto(origin, timeout=BROWSER_NAV_TIMEOUT_MS, wait_until="domcontentloaded") # type: ignore[attr-defined]
|
await page.goto(origin, timeout=BROWSER_NAV_TIMEOUT_MS, wait_until="domcontentloaded") # type: ignore[attr-defined]
|
||||||
|
if wait_for_load:
|
||||||
|
# Только ретрай после #2676: даём клиентской навигации доиграть до `load`,
|
||||||
|
# иначе повтор попадает ровно в то же окно и падает так же.
|
||||||
|
try:
|
||||||
|
await page.wait_for_load_state( # type: ignore[attr-defined]
|
||||||
|
"load", timeout=FETCH_JSON_LOAD_WAIT_MS
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
# Best-effort: страница может дозагружаться бесконечно (реклама/трекеры).
|
||||||
|
# Не пробрасываем — evaluate ниже сам скажет, готова страница или нет;
|
||||||
|
# молчать нельзя, поэтому пишем в лог.
|
||||||
|
logger.info(
|
||||||
|
"tradein-browser[%s]: load не дождались (%s), пробуем evaluate как есть",
|
||||||
|
provider,
|
||||||
|
type(exc).__name__,
|
||||||
|
)
|
||||||
# БЕЗ полного BROWSER_WAIT_MS: нам нужен лишь origin-контекст (cookies +
|
# БЕЗ полного BROWSER_WAIT_MS: нам нужен лишь origin-контекст (cookies +
|
||||||
# same-origin scope для fetch), а не отрендеренные listings. Settle-паузы
|
# same-origin scope для fetch), а не отрендеренные listings. Settle-паузы
|
||||||
# (#1917, FETCH_JSON_SETTLE_MS) хватает, чтобы страница инициализировалась
|
# (#1917, FETCH_JSON_SETTLE_MS) хватает, чтобы страница инициализировалась
|
||||||
|
|
@ -1106,6 +1149,23 @@ def _is_browser_crash(exc: BaseException) -> bool:
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# Литерал playwright, а не строка из нашего лога: driver 1.60.0 (в образе сайдкара)
|
||||||
|
# бросает ровно «Execution context was destroyed» / «... , most likely because of a
|
||||||
|
# navigation.» — обе формы начинаются одинаково, поэтому хватает одного маркера.
|
||||||
|
# Проверено grep'ом по playwright/driver/package/lib/coreBundle.js в живом контейнере.
|
||||||
|
_PAGE_CONTEXT_LOST_MARKER = "execution context was destroyed"
|
||||||
|
|
||||||
|
|
||||||
|
def _is_page_context_lost(exc: BaseException) -> bool:
|
||||||
|
"""Страница потеряла JS-контекст (ушла в навигацию между goto и evaluate), #2676.
|
||||||
|
|
||||||
|
НЕ краш браузера: инстанс жив, потеряна одна страница. Поэтому обрабатывается
|
||||||
|
отдельно от _is_browser_crash — relaunch здесь стоил бы ~10-20с и тёплый профиль
|
||||||
|
(cookies/фингерпринт инстанса) ради браузера, с которым всё в порядке.
|
||||||
|
"""
|
||||||
|
return _PAGE_CONTEXT_LOST_MARKER in str(exc).lower()
|
||||||
|
|
||||||
|
|
||||||
# ── login handler ──────────────────────────────────────────────────────────────
|
# ── login handler ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -67,6 +67,7 @@ class _FakePage:
|
||||||
def __init__(self, evaluate_result: dict[str, Any]) -> None:
|
def __init__(self, evaluate_result: dict[str, Any]) -> None:
|
||||||
self.goto_urls: list[str] = []
|
self.goto_urls: list[str] = []
|
||||||
self.waits: list[int] = [] # записанные wait_for_timeout(ms) — settle-проверка #1917
|
self.waits: list[int] = [] # записанные wait_for_timeout(ms) — settle-проверка #1917
|
||||||
|
self.load_waits: list[int] = [] # wait_for_load_state("load", timeout=) — #2676
|
||||||
self.closed = 0
|
self.closed = 0
|
||||||
# evaluate — AsyncMock, чтобы проверять как сам результат, так и аргументы.
|
# evaluate — AsyncMock, чтобы проверять как сам результат, так и аргументы.
|
||||||
self.evaluate = AsyncMock(return_value=evaluate_result)
|
self.evaluate = AsyncMock(return_value=evaluate_result)
|
||||||
|
|
@ -80,6 +81,9 @@ class _FakePage:
|
||||||
async def wait_for_timeout(self, ms: int) -> None:
|
async def wait_for_timeout(self, ms: int) -> None:
|
||||||
self.waits.append(ms)
|
self.waits.append(ms)
|
||||||
|
|
||||||
|
async def wait_for_load_state(self, state: str, timeout: int = 0) -> None:
|
||||||
|
self.load_waits.append(timeout)
|
||||||
|
|
||||||
async def close(self) -> None:
|
async def close(self) -> None:
|
||||||
self.closed += 1
|
self.closed += 1
|
||||||
|
|
||||||
|
|
@ -365,3 +369,175 @@ def test_do_fetch_json_relaunch_on_browser_crash(monkeypatch: pytest.MonkeyPatch
|
||||||
healthy_page.evaluate.assert_awaited_once()
|
healthy_page.evaluate.assert_awaited_once()
|
||||||
assert crashing_page.closed == 1
|
assert crashing_page.closed == 1
|
||||||
assert healthy_page.closed == 1
|
assert healthy_page.closed == 1
|
||||||
|
|
||||||
|
|
||||||
|
# ── гонка «execution context was destroyed» (#2676) ───────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeSequenceBrowser:
|
||||||
|
"""Отдаёт страницы по очереди: первая попытка ≠ вторая (retry на СВЕЖЕЙ странице)."""
|
||||||
|
|
||||||
|
def __init__(self, pages: list[_FakePage]) -> None:
|
||||||
|
self._pages = list(pages)
|
||||||
|
self.opened = 0
|
||||||
|
|
||||||
|
async def new_page(self) -> _FakePage:
|
||||||
|
page = self._pages[min(self.opened, len(self._pages) - 1)]
|
||||||
|
self.opened += 1
|
||||||
|
return page
|
||||||
|
|
||||||
|
|
||||||
|
def _no_relaunch(monkeypatch: pytest.MonkeyPatch) -> list[str]:
|
||||||
|
"""Подменяет _relaunch_browser счётчиком — тест падает, если его всё-таки позвали."""
|
||||||
|
calls: list[str] = []
|
||||||
|
|
||||||
|
async def _fake(provider: str) -> None:
|
||||||
|
calls.append(provider)
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_relaunch_browser", _fake)
|
||||||
|
return calls
|
||||||
|
|
||||||
|
|
||||||
|
# Обе формы, которые бросает playwright 1.60 (driver coreBundle.js) — короткая и полная.
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"message",
|
||||||
|
[
|
||||||
|
"Page.evaluate: Execution context was destroyed, most likely because of a navigation.",
|
||||||
|
"Execution context was destroyed",
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_do_fetch_json_retries_on_destroyed_context(
|
||||||
|
monkeypatch: pytest.MonkeyPatch, message: str
|
||||||
|
) -> None:
|
||||||
|
"""#2676: страница ушла в навигацию → повтор на свежей странице, БЕЗ relaunch.
|
||||||
|
|
||||||
|
До правки такая ошибка не попадала ни в одну ветку восстановления (_is_browser_crash
|
||||||
|
матчит только закрытие цели/браузера/соединения) и уезжала наверх как 500.
|
||||||
|
"""
|
||||||
|
racing_page = _FakePage({"status": 0, "body": ""})
|
||||||
|
racing_page.evaluate = AsyncMock(side_effect=RuntimeError(message))
|
||||||
|
settled_page = _FakePage({"status": 200, "body": '{"recovered": true}'})
|
||||||
|
|
||||||
|
server._browsers["avito"] = _FakeSequenceBrowser([racing_page, settled_page])
|
||||||
|
monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000)
|
||||||
|
monkeypatch.setattr(server, "FETCH_JSON_LOAD_WAIT_MS", 4242)
|
||||||
|
relaunched = _no_relaunch(monkeypatch)
|
||||||
|
|
||||||
|
result = asyncio.run(
|
||||||
|
server._do_fetch_json(
|
||||||
|
"avito",
|
||||||
|
"https://www.avito.ru/api/x",
|
||||||
|
method="GET",
|
||||||
|
headers={},
|
||||||
|
body=None,
|
||||||
|
origin="https://www.avito.ru/",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result == {"status": 200, "body": '{"recovered": true}'}
|
||||||
|
# Браузер живой — перезапускать его нельзя (тёплые cookies + ~10-20с).
|
||||||
|
assert relaunched == []
|
||||||
|
# Первая попытка НЕ ждала load (happy-path не удлиняется), повтор — ждал.
|
||||||
|
assert racing_page.load_waits == []
|
||||||
|
assert settled_page.load_waits == [4242]
|
||||||
|
assert racing_page.closed == 1 and settled_page.closed == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_do_fetch_json_retry_survives_load_timeout(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Ожидание load на повторе — best-effort: таймаут не отменяет саму попытку."""
|
||||||
|
racing_page = _FakePage({"status": 0, "body": ""})
|
||||||
|
racing_page.evaluate = AsyncMock(side_effect=RuntimeError("Execution context was destroyed"))
|
||||||
|
settled_page = _FakePage({"status": 200, "body": "ok"})
|
||||||
|
settled_page.wait_for_load_state = AsyncMock( # type: ignore[method-assign]
|
||||||
|
side_effect=TimeoutError("Timeout 15000ms exceeded")
|
||||||
|
)
|
||||||
|
|
||||||
|
server._browsers["avito"] = _FakeSequenceBrowser([racing_page, settled_page])
|
||||||
|
monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000)
|
||||||
|
_no_relaunch(monkeypatch)
|
||||||
|
|
||||||
|
result = asyncio.run(
|
||||||
|
server._do_fetch_json(
|
||||||
|
"avito",
|
||||||
|
"https://www.avito.ru/api/x",
|
||||||
|
method="GET",
|
||||||
|
headers={},
|
||||||
|
body=None,
|
||||||
|
origin="https://www.avito.ru/",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert result == {"status": 200, "body": "ok"}
|
||||||
|
|
||||||
|
|
||||||
|
def test_do_fetch_json_gives_up_after_one_context_retry(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Повтор ровно один: вторая та же ошибка уезжает наверх, а не крутит цикл."""
|
||||||
|
message = "Page.evaluate: Execution context was destroyed"
|
||||||
|
first = _FakePage({"status": 0, "body": ""})
|
||||||
|
first.evaluate = AsyncMock(side_effect=RuntimeError(message))
|
||||||
|
second = _FakePage({"status": 0, "body": ""})
|
||||||
|
second.evaluate = AsyncMock(side_effect=RuntimeError(message))
|
||||||
|
|
||||||
|
server._browsers["avito"] = _FakeSequenceBrowser([first, second])
|
||||||
|
monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000)
|
||||||
|
_no_relaunch(monkeypatch)
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="Execution context was destroyed"):
|
||||||
|
asyncio.run(
|
||||||
|
server._do_fetch_json(
|
||||||
|
"avito",
|
||||||
|
"https://www.avito.ru/api/x",
|
||||||
|
method="GET",
|
||||||
|
headers={},
|
||||||
|
body=None,
|
||||||
|
origin="https://www.avito.ru/",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
first.evaluate.assert_awaited_once()
|
||||||
|
second.evaluate.assert_awaited_once()
|
||||||
|
|
||||||
|
|
||||||
|
def test_do_fetch_json_does_not_retry_unrelated_error(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Чужая ошибка НЕ ретраится — повтор невосстановимого жжёт бюджет прогона."""
|
||||||
|
page = _FakePage({"status": 0, "body": ""})
|
||||||
|
page.evaluate = AsyncMock(side_effect=RuntimeError("net::ERR_PROXY_CONNECTION_FAILED"))
|
||||||
|
server._browsers["avito"] = _FakeSequenceBrowser([page])
|
||||||
|
monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000)
|
||||||
|
_no_relaunch(monkeypatch)
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="ERR_PROXY_CONNECTION_FAILED"):
|
||||||
|
asyncio.run(
|
||||||
|
server._do_fetch_json(
|
||||||
|
"avito",
|
||||||
|
"https://www.avito.ru/api/x",
|
||||||
|
method="GET",
|
||||||
|
headers={},
|
||||||
|
body=None,
|
||||||
|
origin="https://www.avito.ru/",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
page.evaluate.assert_awaited_once()
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_json_handler_500_carries_reason_in_body(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Причина отказа остаётся в теле 500 — её читает _raise_for_sidecar_status (#2708)."""
|
||||||
|
message = "Execution context was destroyed"
|
||||||
|
first = _FakePage({"status": 0, "body": ""})
|
||||||
|
first.evaluate = AsyncMock(side_effect=RuntimeError(message))
|
||||||
|
second = _FakePage({"status": 0, "body": ""})
|
||||||
|
second.evaluate = AsyncMock(side_effect=RuntimeError(message))
|
||||||
|
server._browsers["avito"] = _FakeSequenceBrowser([first, second])
|
||||||
|
monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000)
|
||||||
|
_no_relaunch(monkeypatch)
|
||||||
|
|
||||||
|
async def _ensure(provider: str, proxy_override: str | None = None) -> bool:
|
||||||
|
return True
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_ensure_browser", _ensure)
|
||||||
|
|
||||||
|
response = asyncio.run(
|
||||||
|
server.fetch_json_handler(
|
||||||
|
_make_request({"url": "https://www.avito.ru/api/x", "source": "avito"})
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert response.status == 500
|
||||||
|
assert "Execution context was destroyed" in _json_body(response)["error"]
|
||||||
|
|
|
||||||
|
|
@ -47,6 +47,11 @@ _RAW_TO_CANON: dict[str, str] = {
|
||||||
"panel": "panel",
|
"panel": "panel",
|
||||||
"block": "block",
|
"block": "block",
|
||||||
"wood": "wood",
|
"wood": "wood",
|
||||||
|
# #2678 п.6: канон целиком, включая monolith_brick. Pass-through по _CANON
|
||||||
|
# регистрозависим, и без этого ключа 'MONOLITH_BRICK'/'Monolith_Brick' (форма,
|
||||||
|
# в которой канон может прийти от нового источника) уезжали бы в None —
|
||||||
|
# единственный канонический токен без такой страховки.
|
||||||
|
"monolith_brick": "monolith_brick",
|
||||||
"monolithBrick": "monolith_brick",
|
"monolithBrick": "monolith_brick",
|
||||||
"gasSilicateBlock": "block",
|
"gasSilicateBlock": "block",
|
||||||
"aerocreteBlock": "block",
|
"aerocreteBlock": "block",
|
||||||
|
|
|
||||||
|
|
@ -9,6 +9,15 @@ app-копии:
|
||||||
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
|
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
|
||||||
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
|
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
|
||||||
app-копии эта функция не нужна.
|
app-копии эта функция не нужна.
|
||||||
|
|
||||||
|
ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL —
|
||||||
|
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы
|
||||||
|
выполняются той же сессией, что и работа задачи, и если её транзакция оставалась
|
||||||
|
открытой всё время работы (нечего было коммитить), их UPDATE попадал ВНУТРЬ неё —
|
||||||
|
`finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с
|
||||||
|
окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в
|
||||||
|
диапазоне 9-64 мс (= от коммита claim'а до первого запроса рабочей транзакции).
|
||||||
|
Полный разбор — в docstring app-копии `app/services/scrape_runs.py`.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -336,7 +345,10 @@ def _alert_on_run_id(
|
||||||
|
|
||||||
|
|
||||||
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
"""INSERT scrape_runs(source, status='running', params, started_at=NOW()).
|
"""INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()).
|
||||||
|
|
||||||
|
started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей
|
||||||
|
транзакции задачи его уже не достаёт (#2702).
|
||||||
|
|
||||||
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
||||||
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
||||||
|
|
@ -347,7 +359,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
||||||
VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW())
|
VALUES (
|
||||||
|
:source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp()
|
||||||
|
)
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
|
|
@ -396,9 +410,9 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
)
|
)
|
||||||
UPDATE scrape_runs r
|
UPDATE scrape_runs r
|
||||||
SET heartbeat_at = NOW(),
|
SET heartbeat_at = clock_timestamp(),
|
||||||
started_at = NOW(),
|
started_at = clock_timestamp(),
|
||||||
finished_at = NOW(),
|
finished_at = clock_timestamp(),
|
||||||
counters = COALESCE(r.counters, '{}'::jsonb) || jsonb_build_object(
|
counters = COALESCE(r.counters, '{}'::jsonb) || jsonb_build_object(
|
||||||
'skips', COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1,
|
'skips', COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1,
|
||||||
'detail', CAST(:details AS text),
|
'detail', CAST(:details AS text),
|
||||||
|
|
@ -424,7 +438,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
||||||
source, status, error, counters, started_at, heartbeat_at, finished_at
|
source, status, error, counters, started_at, heartbeat_at, finished_at
|
||||||
)
|
)
|
||||||
VALUES (
|
VALUES (
|
||||||
:source, 'skipped', :reason, CAST(:counters AS jsonb), NOW(), NOW(), NOW()
|
:source, 'skipped', :reason, CAST(:counters AS jsonb), clock_timestamp(), clock_timestamp(), clock_timestamp()
|
||||||
)
|
)
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
|
|
@ -448,7 +462,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
||||||
|
|
||||||
|
|
||||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
"""UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
|
||||||
|
|
||||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||||
|
|
@ -459,7 +473,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET heartbeat_at = NOW(),
|
SET heartbeat_at = clock_timestamp(),
|
||||||
counters = CAST(:counters AS jsonb),
|
counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -486,7 +500,7 @@ def is_cancelled(db: Session, run_id: int) -> bool:
|
||||||
|
|
||||||
|
|
||||||
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
"""Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count.
|
"""Финализация run: status='done', finished_at + counters + total_seen/new_count.
|
||||||
|
|
||||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
||||||
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
||||||
|
|
@ -496,7 +510,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'done',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
counters = CAST(:counters AS jsonb),
|
counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -535,7 +550,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'failed',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
error = :error, counters = CAST(:counters AS jsonb),
|
error = :error, counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||||
|
|
@ -594,7 +610,8 @@ def mark_banned(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(),
|
SET status = 'banned',
|
||||||
|
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||||
error = :error, counters = CAST(:counters AS jsonb),
|
error = :error, counters = CAST(:counters AS jsonb),
|
||||||
ban_kind = :ban_kind,
|
ban_kind = :ban_kind,
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
|
|
@ -625,7 +642,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'cancelled', finished_at = NOW()
|
SET status = 'cancelled', finished_at = clock_timestamp()
|
||||||
WHERE id = :run_id AND status = 'running'
|
WHERE id = :run_id AND status = 'running'
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
|
|
@ -235,16 +235,26 @@ def has_running_run(db: Session, source: str) -> bool:
|
||||||
|
|
||||||
|
|
||||||
def reap_zombies(db: Session) -> int:
|
def reap_zombies(db: Session) -> int:
|
||||||
"""Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours."""
|
"""Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours.
|
||||||
|
|
||||||
|
clock_timestamp(), а не now() (#2702): критерий сравнивает ЗАПИСАННЫЙ heartbeat со
|
||||||
|
временем «сейчас», и обе стороны сравнения должны быть настоящим временем. `now()`
|
||||||
|
замерзает на старте транзакции — у писавшей heartbeat стороны это давало отставание
|
||||||
|
на весь возраст открытой рабочей транзакции (см. docstring runs.py), у читающей
|
||||||
|
стороны — на возраст тика. На проде это уже стоило ложных срабатываний: у всех 6
|
||||||
|
прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и
|
||||||
|
остался на отметке старта (max advance 0.0 с) — при том что нормальный прогон этого
|
||||||
|
источника длится до 5.06 ч (прогон 346) и обязан был двигать heartbeat.
|
||||||
|
"""
|
||||||
zombie_interval = f"{ZOMBIE_THRESHOLD_HOURS} hours"
|
zombie_interval = f"{ZOMBIE_THRESHOLD_HOURS} hours"
|
||||||
result = db.execute(
|
result = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
UPDATE scrape_runs
|
||||||
SET status = 'zombie', finished_at = NOW()
|
SET status = 'zombie', finished_at = clock_timestamp()
|
||||||
WHERE status = 'running'
|
WHERE status = 'running'
|
||||||
AND (heartbeat_at IS NULL
|
AND (heartbeat_at IS NULL
|
||||||
OR heartbeat_at < NOW() - CAST(:interval AS interval))
|
OR heartbeat_at < clock_timestamp() - CAST(:interval AS interval))
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
|
|
|
||||||
|
|
@ -39,6 +39,7 @@ from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from scraper_kit.avito_exceptions import AvitoBlockedError, AvitoRateLimitedError
|
from scraper_kit.avito_exceptions import AvitoBlockedError, AvitoRateLimitedError
|
||||||
|
from scraper_kit.house_type_normalizer import normalize_house_type
|
||||||
from scraper_kit.providers.avito.serp import _is_firewall_page
|
from scraper_kit.providers.avito.serp import _is_firewall_page
|
||||||
from scraper_kit.providers.avito.shared import RUS_MONTHS, _unix_to_date
|
from scraper_kit.providers.avito.shared import RUS_MONTHS, _unix_to_date
|
||||||
|
|
||||||
|
|
@ -258,10 +259,24 @@ def _strip_price(price_str: str | None) -> int | None:
|
||||||
|
|
||||||
|
|
||||||
def _normalize_house_type(raw: str | None) -> str | None:
|
def _normalize_house_type(raw: str | None) -> str | None:
|
||||||
"""Нормализует тип дома: "Монолитный" → "monolith"."""
|
"""Нормализует тип дома: "Монолитный" → "monolith". Незнакомое → None.
|
||||||
|
|
||||||
|
Единственный живой писатель houses.house_type (см. save-функцию ниже), поэтому
|
||||||
|
канон обязан приводиться ЗДЕСЬ, а не у каждого читателя (#2678).
|
||||||
|
|
||||||
|
Русская подпись Авито снимается локальной HOUSE_TYPE_MAP, результат прогоняется
|
||||||
|
через общий scraper_kit.house_type_normalizer: он знает и канон, и чужие
|
||||||
|
вокабуляры (cian camelCase, yandex SCREAMING) — на случай, если карточка дома
|
||||||
|
однажды придёт с чужим токеном.
|
||||||
|
|
||||||
|
#2678: раньше незнакомое значение становилось 'other' — единственный источник
|
||||||
|
неканоничных значений в houses среди живых писателей. 'other' всегда != канон,
|
||||||
|
т.е. для soft-penalty эстиматора это ложный штраф, а NULL нейтрален (тот же
|
||||||
|
довод, что в docstring house_type_normalizer).
|
||||||
|
"""
|
||||||
if not raw:
|
if not raw:
|
||||||
return None
|
return None
|
||||||
return HOUSE_TYPE_MAP.get(raw.lower(), "other")
|
return normalize_house_type(HOUSE_TYPE_MAP.get(raw.strip().lower(), raw))
|
||||||
|
|
||||||
|
|
||||||
def _normalize_house_class(raw: str | None) -> str | None:
|
def _normalize_house_class(raw: str | None) -> str | None:
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue