diff --git a/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py b/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py
new file mode 100644
index 00000000..cc0ba37a
--- /dev/null
+++ b/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py
@@ -0,0 +1,295 @@
+"""#930 добивка: планировщик не подхватывал чекпоинт оборванного прогона.
+
+#930 сделал обе половины механизма — запись точки (`counters.done_buckets`, per-bucket
+heartbeat) и её чтение (`run_*_full_load(resume_run_id=...)`, skip-set в SERP-слое), —
+но единственным входом оставил админку. У avito full-load админского эндпоинта нет
+вовсе, а планировщик передавал `resume_run_id=None` ЛИТЕРАЛОМ (scheduler.py 708/728/799
+на origin/main). То есть боевой путь возобновления не существовал ни одного дня.
+
+Цена на проде (замер 2026-08-12, 90 суток, read-only): 433 корзины в 30 оборванных
+прогонах с ЖИВОЙ незабранной точкой — avito_full_load 242, cian_full_load 134,
+avito_full_load_exhaustive 57. Прогон 3547 (09.08, убит деплоем на третьем часу, 35 из
+84 корзин дерева) лежит до сих пор и будет подхвачен расписанием 139 16.08.
+
+Красный прогон на origin/main:
+ 1. `test_scheduler_hands_checkpoint_to_pipeline` — планировщик отдаёт в пайплайн
+ resume_run_id=None вместо id прошлого прогона (AssertionError на 3 источниках);
+ 2. `test_partial_bucket_is_not_complete` — бакет с выпавшей страницей приезжает в
+ on_bucket неотличимым от целого (у колбэка нет аргумента полноты вообще);
+ 3. `test_pipeline_keeps_partial_bucket_out_of_checkpoint` — TypeError: `_on_bucket`
+ на main принимает два аргумента, признаку полноты некуда приехать.
+Тесты ладдера (`_resume_decision`) на main падают с AttributeError — функции нет.
+"""
+
+from __future__ import annotations
+
+import os
+from types import SimpleNamespace
+from typing import Any
+from unittest.mock import AsyncMock, MagicMock, patch
+
+import pytest
+
+os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
+
+from scraper_kit.orchestration import scheduler as sched
+from scraper_kit.orchestration.pipeline import run_avito_full_load
+from scraper_kit.providers.avito.serp import AvitoScraper
+
+PFX = "scraper_kit.orchestration.pipeline"
+
+# Прод-слепок расписания 139 (avito_full_load_exhaustive) на 2026-08-12: params прогона
+# 3547 совпадают с default_params расписания байт-в-байт — это и есть «то же задание».
+_PARAMS = {
+ "concurrency": 1,
+ "interval_days": 7,
+ "secondary_only": True,
+ "request_delay_sec": 7.0,
+ "price_cap_per_bucket": 1400,
+}
+
+
+def _candidate(**over: Any) -> SimpleNamespace:
+ """Строка-кандидат из _RESUME_CANDIDATE_SQL: прогон 3547 как он лежит на проде."""
+ base = {
+ "prev_id": 3547,
+ "prev_status": "cancelled",
+ "prev_counters": {
+ "unique_fetched": 5496,
+ "done_buckets": [f"room_1_komn:{i}:0" for i in range(35)],
+ },
+ "same_params": True,
+ "age_h": 164.6, # 6.86 суток — столько будет точке 3547 к подхвату 16.08
+ "interval_days": "7",
+ }
+ base.update(over)
+ return SimpleNamespace(**base)
+
+
+class _FakeDb:
+ """Двойник сессии: отдаёт ОДНУ строку-кандидата на любой SELECT, глотает UPDATE."""
+
+ def __init__(self, row: Any) -> None:
+ self.row = row
+ self.written: list[dict[str, Any]] = []
+
+ def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
+ if params and "counters" in params: # update_heartbeat пишет вердикт
+ self.written.append(params)
+ return MagicMock()
+ return MagicMock(fetchone=lambda: self.row)
+
+ def commit(self) -> None:
+ pass
+
+
+# ── 1. Главное: планировщик обязан отдать точку в пайплайн ───────────────────
+
+
+@pytest.mark.parametrize(
+ ("job", "pipeline_fn"),
+ [
+ (sched._job_avito_full_load, "run_avito_full_load"),
+ (sched._job_avito_full_load_exhaustive, "run_avito_full_load"),
+ (sched._job_cian_full_load, "run_cian_full_load"),
+ ],
+)
+async def test_scheduler_hands_checkpoint_to_pipeline(job: Any, pipeline_fn: str) -> None:
+ """Оборванный прогон с валидной точкой → новый прогон продолжает его, а не с нуля.
+
+ Падает на origin/main: планировщик передаёт литеральный None — 433 корзины за 90
+ суток перебирались заново, включая 35 корзин прогона 3547.
+ """
+ db = _FakeDb(_candidate())
+ captured: dict[str, Any] = {}
+
+ async def _spy(*_a: Any, **kw: Any) -> None:
+ captured.update(kw)
+
+ with patch.object(sched, pipeline_fn, _spy):
+ await job(db, 4000, dict(_PARAMS), MagicMock())
+
+ assert captured["resume_run_id"] == 3547
+
+
+async def test_verdict_lands_in_counters_of_new_run() -> None:
+ """Подхватили или нет — видно В СЧЁТЧИКАХ прогона, а не только в docker-логах.
+
+ Логи теряются при редеплое (контейнер tradein-scraper пересоздаётся), поэтому
+ молчаливый отказ подхватить неотличим от отсутствия правки.
+ """
+ db = _FakeDb(_candidate(prev_status="zombie"))
+ with patch.object(sched, "run_avito_full_load", AsyncMock()):
+ await sched._job_avito_full_load(db, 4000, dict(_PARAMS), MagicMock())
+
+ assert db.written, "вердикт о подхвате не записан в counters нового прогона"
+ written = db.written[-1]["counters"]
+ assert '"resume_reason": "status_zombie"' in written
+ assert '"resume_candidate": 3547' in written
+
+
+# ── 2. Ладдер отказов: у каждого нуля своя причина ───────────────────────────
+
+
+@pytest.mark.parametrize(
+ ("row", "reason"),
+ [
+ (None, "no_prev_run"),
+ (_candidate(prev_status="done"), "status_done"),
+ (_candidate(prev_status="zombie"), "status_zombie"),
+ (_candidate(same_params=False), "params_changed"),
+ (_candidate(prev_counters={"unique_fetched": 2977}), "no_checkpoint"),
+ (_candidate(age_h=200.0), "checkpoint_stale"),
+ (_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2}), "chain_limit"),
+ ],
+)
+def test_resume_refusals_are_named(row: Any, reason: str) -> None:
+ """«Не подхватили» — это семь РАЗНЫХ фактов, и в counters они различимы."""
+ resume_id, verdict = sched._resume_decision(row)
+ assert resume_id is None
+ assert verdict["resume_reason"] == reason
+ assert verdict["resume_from"] is None
+
+
+def test_resume_chain_is_bounded() -> None:
+ """Цепочка возобновлений считается и упирается в потолок, а не тянется вечно.
+
+ Потолок выведен из STALE_DIGEST_INTERVAL_FACTOR (см. scheduler.py): полный обход
+ обязан начаться раньше, чем сводка объявит источник просроченным.
+ """
+ assert sched._MAX_RESUME_CHAIN == sched.STALE_DIGEST_INTERVAL_FACTOR - 1
+ _id, first = sched._resume_decision(_candidate())
+ assert first["resume_chain"] == 1
+ _id2, second = sched._resume_decision(
+ _candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 1})
+ )
+ assert second["resume_chain"] == sched._MAX_RESUME_CHAIN
+ third_id, third = sched._resume_decision(
+ _candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2})
+ )
+ assert third_id is None and third["resume_reason"] == "chain_limit"
+
+
+def test_stale_threshold_follows_the_source_tick() -> None:
+ """Срок годности точки считается от такта ИСТОЧНИКА, а не общей константой.
+
+ cian ходит раз в 3 суток, avito — раз в 7; одна и та же точка возрастом 100 ч для
+ первого просрочена, для второго свежая. Плюс сутки — сетка запуска (см.
+ _resume_decision): 164.6 ч прогона 3547 при такте 7 суток обязаны пройти, иначе
+ точку отвергал бы jitter расписания, а пропущенный цикл (13 суток) — нет.
+ """
+ assert sched._resume_decision(_candidate(age_h=100.0, interval_days="3"))[0] is None
+ assert sched._resume_decision(_candidate(age_h=100.0, interval_days="7"))[0] == 3547
+ assert sched._resume_decision(_candidate(age_h=164.6, interval_days="7"))[0] == 3547
+ assert sched._resume_decision(_candidate(age_h=13 * 24.0, interval_days="7"))[0] is None
+
+
+# ── 3. Недособранный бакет не имеет права попасть в чекпоинт ─────────────────
+
+
+def _serp_config() -> SimpleNamespace:
+ return SimpleNamespace(
+ scraper_fetch_mode="curl_cffi",
+ browser_http_endpoint="http://browser.test/fetch",
+ scraper_proxy_url=None,
+ avito_proxy_max_rotations=0,
+ avito_serp_ok_not_banned=True,
+ avito_proxy_rotate_settle_s=0.0,
+ proxy_rotate_attempts=1,
+ proxy_rotate_attempt_timeout_s=1.0,
+ scraper_skip_seen_today=False,
+ )
+
+
+@pytest.mark.parametrize(
+ ("page2_html", "expected_complete"),
+ [(None, False), ("", True)],
+)
+async def test_partial_bucket_is_not_complete(
+ page2_html: str | None, expected_complete: bool
+) -> None:
+ """Страница 2 из 3 выпала → бакет НЕ «сделан»; все три пришли → «сделан».
+
+ Контрольная половина обязательна: реализация «всегда False» тоже прошла бы
+ одностороннюю проверку, но убила бы возобновление целиком.
+
+ Падает на origin/main: `on_bucket` вызывается двумя аргументами, признака полноты
+ в протоколе нет — частичный бакет неотличим от целого и попадает в done_buckets.
+ """
+ scraper = AvitoScraper(_serp_config())
+ scraper.request_delay_sec = 0.0
+ calls: list[tuple[str, bool]] = []
+
+ def _on_bucket(key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
+ calls.append((key, complete))
+
+ async def _fetch_page(_self: Any, _slug: str, page: int, *_a: Any, **_k: Any) -> str | None:
+ return page2_html if page == 2 else f""
+
+ with (
+ patch.object(AvitoScraper, "_fetch_rooms_page_html", _fetch_page),
+ patch.object(
+ AvitoScraper,
+ "_parse_html",
+ lambda _self, html, **_k: [MagicMock(source_id=html, listing_segment="secondary")],
+ ),
+ ):
+ await scraper._paginate_leaf_bucket(
+ room_slug="kvartiry_1_komnatnye",
+ room_label="room_1_komn",
+ lo=0,
+ hi=3999999,
+ html="",
+ max_pages=3,
+ seen={},
+ price_cap_per_bucket=1400,
+ max_pages_per_bucket=100,
+ concurrency=2,
+ secondary_only=False,
+ on_bucket=_on_bucket,
+ skip_buckets=None,
+ expected_total=3 * 50,
+ )
+
+ assert [c[1] for c in calls] == [expected_complete]
+
+
+async def test_pipeline_keeps_partial_bucket_out_of_checkpoint() -> None:
+ """Пайплайн: лоты частичного бакета СОХРАНЕНЫ, но в чекпоинт он не попал.
+
+ Именно здесь «видимая потеря» (перескрап) не превращается в «невидимую»: пропустить
+ частичный бакет на следующем прогоне значит не перечитать его страницы уже никогда.
+ """
+ finals: list[dict[str, Any]] = []
+
+ class _Recorder:
+ def is_cancelled(self, *_a: Any, **_k: Any) -> bool:
+ return False
+
+ def update_heartbeat(self, *_a: Any, **_k: Any) -> None:
+ pass
+
+ def mark_done(self, _db: Any, _rid: int, counters: dict[str, Any]) -> None:
+ finals.append(dict(counters))
+
+ async def _fetch(*_a: Any, on_bucket: Any = None, **_k: Any) -> None:
+ on_bucket("room_1_komn:0:3999999", [MagicMock(source_id="a1")], True)
+ on_bucket("room_1_komn:4000000:4999999", [MagicMock(source_id="a2")], False)
+
+ scraper = MagicMock()
+ scraper.__aenter__ = AsyncMock(return_value=scraper)
+ scraper.__aexit__ = AsyncMock(return_value=None)
+ scraper.fetch_all_secondary = _fetch
+
+ with (
+ patch(f"{PFX}.AvitoScraper", return_value=scraper),
+ patch(f"{PFX}.save_listings", MagicMock(return_value=(1, 0))),
+ patch(f"{PFX}.runs", _Recorder()),
+ ):
+ counters = await run_avito_full_load(
+ MagicMock(), run_id=1, config=_serp_config(), matcher=MagicMock()
+ )
+
+ assert finals[0]["done_buckets"] == ["room_1_komn:0:3999999"]
+ assert finals[0]["partial_buckets"] == 1
+ assert counters.unique_fetched == 2, "лоты частичного бакета обязаны быть сохранены"
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py
index 5c090d75..1f3bdd1a 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py
@@ -2953,6 +2953,11 @@ class CianFullLoadCounters:
detail_enriched: int = 0
detail_failed: int = 0
errors_count: int = 0
+ # Бакеты, отданные SERP-слоем как НЕполные (страница выпала / исключение
+ # проглочено / hard-cap): лоты сохранены, но в done_buckets бакет не попал.
+ # Без этого счётчика «сколько бакетов чекпоинт не покрывает» видно только грепом
+ # логов, которые теряются при редеплое.
+ partial_buckets: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
@@ -3007,8 +3012,32 @@ async def run_cian_full_load(
done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона
- def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
- """Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done."""
+ def _mark_bucket(bucket_key: str, complete: bool) -> None:
+ """В чекпоинт — только ПОЛНОСТЬЮ собранный бакет; частичный лишь считаем.
+
+ Ключ done_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета,
+ поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск
+ частичного бакета на следующем прогоне означал бы, что его непрочитанные
+ страницы не перечитает уже никто, а счётчики покажут успех.
+ """
+ if complete:
+ done.add(bucket_key)
+ return
+ counters.partial_buckets += 1
+ logger.warning(
+ "cian-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, "
+ "следующий прогон перечитает его целиком (partial_buckets=%d)",
+ run_id,
+ bucket_key,
+ counters.partial_buckets,
+ )
+
+ def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
+ """Инкрементальный save после каждого leaf-бакета.
+
+ complete=False → лоты сохраняем, бакет в чекпоинт не пишем
+ (см. _mark_bucket). Дефолт True — для вызывающих без пагинации.
+ """
nonlocal done
if runs.is_cancelled(db, run_id):
logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id)
@@ -3023,7 +3052,7 @@ async def run_cian_full_load(
)
raise RuntimeError("shutdown")
if not lots:
- done.add(bucket_key)
+ _mark_bucket(bucket_key, complete)
runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
@@ -3042,7 +3071,7 @@ async def run_cian_full_load(
counters.saved_inserted += inserted
counters.saved_updated += updated
counters.unique_fetched += len(lots)
- done.add(bucket_key)
+ _mark_bucket(bucket_key, complete)
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d",
@@ -3225,7 +3254,8 @@ async def run_cian_full_load(
except NoProxyAvailableError as exc:
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
- # сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
+ # сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который
+ # его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930).
logger.error("cian-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
@@ -3467,7 +3497,8 @@ async def run_yandex_full_load(
except NoProxyAvailableError as exc:
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
- # сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
+ # сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который
+ # его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930).
logger.error("yandex-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
@@ -3527,6 +3558,8 @@ class AvitoFullLoadCounters:
saved_inserted: int = 0
saved_updated: int = 0
errors_count: int = 0
+ # Бакеты, отданные SERP-слоем как НЕполные — см. CianFullLoadCounters.
+ partial_buckets: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
@@ -3599,8 +3632,32 @@ async def run_avito_full_load(
)
done: set[str] = set(skip_set)
- def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
- """Инкрементальный save после каждого leaf-бакета."""
+ def _mark_bucket(bucket_key: str, complete: bool) -> None:
+ """В чекпоинт — только ПОЛНОСТЬЮ собранный бакет; частичный лишь считаем.
+
+ Ключ done_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета,
+ поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск
+ частичного бакета на следующем прогоне означал бы, что его непрочитанные
+ страницы не перечитает уже никто, а счётчики покажут успех.
+ """
+ if complete:
+ done.add(bucket_key)
+ return
+ counters.partial_buckets += 1
+ logger.warning(
+ "avito-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, "
+ "следующий прогон перечитает его целиком (partial_buckets=%d)",
+ run_id,
+ bucket_key,
+ counters.partial_buckets,
+ )
+
+ def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
+ """Инкрементальный save после каждого leaf-бакета.
+
+ complete=False → лоты сохраняем, бакет в чекпоинт не пишем
+ (см. _mark_bucket). Дефолт True — для вызывающих без пагинации.
+ """
nonlocal done
if runs.is_cancelled(db, run_id):
logger.info("avito-full-load run_id=%d: cancel detected in on_bucket", run_id)
@@ -3613,7 +3670,7 @@ async def run_avito_full_load(
)
raise RuntimeError("shutdown")
if not lots:
- done.add(bucket_key)
+ _mark_bucket(bucket_key, complete)
runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
@@ -3632,7 +3689,7 @@ async def run_avito_full_load(
counters.saved_inserted += inserted
counters.saved_updated += updated
counters.unique_fetched += len(lots)
- done.add(bucket_key)
+ _mark_bucket(bucket_key, complete)
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"avito-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d",
@@ -3699,9 +3756,10 @@ async def run_avito_full_load(
except NoProxyAvailableError as exc:
# #2687: пул опустел mid-run. Ветка стоит ДО generic-RuntimeError намеренно —
# NoProxyAvailableError его подкласс, и без неё отказ уходил в mark_failed,
- # который (в отличие от mark_banned) НЕ пишет done_buckets. То есть чекпоинт
+ # который (в отличие от mark_banned) НЕ передавал done_buckets. То есть чекпоинт
# терялся ровно на НАШЕМ отказе — том исходе, для которого #2686 требовал его
- # сохранять наравне с блокировкой площадкой.
+ # сохранять наравне с блокировкой площадкой. Диагноз ban_kind='infra' ветка
+ # даёт по-прежнему; сам чекпоинт с #930-мержем counters переживает и mark_failed.
logger.error("avito-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py
index c0996e1b..86729701 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py
@@ -100,8 +100,8 @@ def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int:
# РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый
# статус» и «явное поле причины»:
# 1. Побочная функция 'banned' — сохранение done_buckets-чекпоинта (mark_failed
-# его теряет) — нужна ОБОИМ исходам. Оставив статус, получаем её даром; расщепив
-# статус, пришлось бы дублировать её в каждом потребителе.
+# его тогда терял) — нужна ОБОИМ исходам. Оставив статус, получаем её даром;
+# расщепив статус, пришлось бы дублировать её в каждом потребителе.
# 2. Новое значение статуса пришлось бы доучить пяти местам, каждое из которых
# молча даёт неверный ответ, если про него забыть: CHECK-констрейнт схемы,
# IN-списки обоих сторожей (_alert_if_consecutive_failures / _zero_results),
@@ -570,12 +570,21 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
return int(row.id)
-def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
- """UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
+def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None:
+ """UPDATE heartbeat_at + counters (МЕРЖ, не замена) + total_seen/new_count колонки.
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
COALESCE: если ключа нет в counters — старое значение колонки сохраняется.
+
+ `counters || :counters` вместо замены (#930 добивка): чекпоинт `done_buckets`
+ писали ТОЛЬКО сайты, знающие о нём (`_on_bucket`), а heartbeat'ы, которые о нём не
+ знают, целиком перезаписывали объект и СТИРАЛИ точку. На проде это давало
+ немонотонную точку: `_on_progress` (после каждой комнатности) и фоновый heartbeat
+ cian'а (каждые 60 с) отправляли `counters.to_dict()` без ключа — то есть у cian
+ точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15
+ оборванных прогонов с доказанной работой (35 706 + 9 222 fetched) и БЕЗ ключа
+ вообще. Мерж делает точку монотонной для любого писателя, а не только для знающих.
"""
total_seen, new_count = _column_counts(counters)
db.execute(
@@ -583,7 +592,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
"""
UPDATE scrape_runs
SET heartbeat_at = clock_timestamp(),
- counters = CAST(:counters AS jsonb),
+ counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id
@@ -641,7 +650,7 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
UPDATE scrape_runs
SET status = 'done',
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
- counters = CAST(:counters AS jsonb),
+ counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id AND status = 'running'
@@ -681,7 +690,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
UPDATE scrape_runs
SET status = 'failed',
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
- error = :error, counters = CAST(:counters AS jsonb),
+ error = :error,
+ counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id AND status = 'running'
@@ -715,8 +725,9 @@ def mark_banned(
Per migration 015 — 'banned' задокументирован как 'Avito вернул 403/captcha'.
Отличается от 'failed': прогон оборван внешним/блокирующим условием, а не нашим
- багом, и — важно — СОХРАНЯЕТ done_buckets-чекпоинт в counters (mark_failed его
- теряет). Cooldown 2-4 часа.
+ багом. Чекпоинт done_buckets раньше сохранял только этот финализатор — теперь
+ counters мержатся во всех (см. update_heartbeat), и точка переживает любой из них.
+ Cooldown 2-4 часа.
`ban_kind` разводит два исхода, которые раньше схлопывались в один статус:
- BAN_KIND_PLATFORM — площадка нас заблокировала (firewall/403/captcha);
@@ -742,7 +753,8 @@ def mark_banned(
UPDATE scrape_runs
SET status = 'banned',
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
- error = :error, counters = CAST(:counters AS jsonb),
+ error = :error,
+ counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
ban_kind = :ban_kind,
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py
index 06d778a8..8b0b5009 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py
@@ -497,6 +497,132 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
return run_id
+# ── подхват чекпоинта оборванного прогона (#930, вторая половина) ────────────
+# #930 сделал чекпоинт (`counters.done_buckets`) и приёмную сторону
+# (`run_*_full_load(resume_run_id=...)`), но единственным входом оставил админку. У
+# avito full-load её нет вовсе, поэтому 299 из 433 «впустую перебранных» корзин за 90
+# суток не имели НИКАКОГО пути возобновления, даже ручного. Планировщик передавал
+# resume_run_id=None литералом.
+#
+# Точка берётся только когда выполнены ВСЕ условия ниже; иначе прогон честно начинает с
+# нуля, а ПРИЧИНА пишется в его counters (молчаливый отказ неотличим от отсутствия
+# правки — см. _resume_decision).
+
+# 'zombie' НЕ в списке НАМЕРЕННО. reap_zombies снимает пометку 'running', но НЕ убивает
+# процесс (прямо задокументировано в app/tasks/listing_source_snapshot.py) — а
+# has_running_run гейтит claim именно по статусу. То есть после reap'а старый сборщик
+# может продолжать писать в ту же строку: подхват читал бы ДВИЖУЩУЮСЯ точку и запускал
+# второй сборщик на ту же площадку. 147 корзин в 10 zombie-прогонах за 90 суток
+# остаются несобранными сознательно — это цена, а не недосмотр.
+# 'failed' в списке: его чекпоинт больше не стирается финализатором (runs.py, мерж
+# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных
+# корзин — они записаны тем же heartbeat'ом, что и у banned.
+_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"})
+
+# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) —
+# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта =
+# сломан». Цепочка не имеет права отодвинуть полный обход дальше этой же черты,
+# поэтому подряд идущих подхватов допускается на один меньше: прогоны 1 и 2 могут
+# продолжать предшественника, третий обязан пойти с нуля. При такте avito 7 суток это
+# гарантирует попытку полного обхода не реже, чем раз в 21 сутки — ровно в тот момент,
+# когда сводка объявляет источник просроченным.
+_MAX_RESUME_CHAIN = STALE_DIGEST_INTERVAL_FACTOR - 1
+
+# Кандидат — ПОСЛЕДНИЙ прогон источника, а не последний подходящий: если после обрыва
+# уже прошёл полный ('done') прогон, дерево обойдено и возобновлять нечего. Строки
+# 'skipped' — бухгалтерия планировщика, а не прогоны, поэтому не в счёт.
+_RESUME_CANDIDATE_SQL = text("""
+ WITH cur AS (
+ SELECT source, params FROM scrape_runs WHERE id = CAST(:rid AS bigint)
+ )
+ SELECT r.id AS prev_id,
+ r.status AS prev_status,
+ r.counters AS prev_counters,
+ (r.params IS NOT DISTINCT FROM cur.params) AS same_params,
+ EXTRACT(EPOCH FROM (clock_timestamp() - r.heartbeat_at)) / 3600.0 AS age_h,
+ cur.params ->> 'interval_days' AS interval_days
+ FROM scrape_runs r, cur
+ WHERE r.source = cur.source
+ AND r.id <> CAST(:rid AS bigint)
+ AND r.status <> 'skipped'
+ ORDER BY r.started_at DESC
+ LIMIT 1
+""")
+
+
+def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
+ """Решение «подхватывать ли точку» + счётчики-объяснение. Чистая функция.
+
+ Возвращает (resume_run_id | None, counters-заготовка нового прогона). Причина
+ отказа — машиночитаемый слаг в `resume_reason`, по нему «предыдущего прогона не
+ было» отличается от «параметры разъехались» ЗАПРОСОМ, а не грепом логов.
+
+ Срок годности точки — такт источника ПЛЮС сутки (`interval_days` из его же params).
+ Ни одно из слагаемых не выбрано произвольно. Такт — объявленный самим расписанием
+ срок, в течение которого собранное считается свежим; точка старше него пережила
+ цикл, в котором источник обязан был обойти дерево целиком. Сутки — сетка запуска:
+ `compute_next_run_at` выбирает ДЕНЬ (сегодня+interval_days) и случайное время внутри
+ окна, поэтому два соседних запуска отстоят друг от друга на interval_days ± меньше
+ суток. Без этого слагаемого точку отвергал бы jitter расписания, а не устаревание:
+ прогон 3547 убит деплоем 09.08 16:53, расписание 139 подхватит его 16.08 13:37 —
+ 164.6 ч при такте 168 ч, запас 3.4 ч при ширине окна 2 ч. Пропущенный цикл в окно
+ всё равно не влезает: для avito это 13 суток против порога 8.
+ """
+ if row is None:
+ return None, {"resume_from": None, "resume_reason": "no_prev_run", "resume_chain": 0}
+
+ prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {}
+ done_buckets = prev_counters.get("done_buckets")
+ done_n = len(done_buckets) if isinstance(done_buckets, list) else 0
+ chain_raw = prev_counters.get("resume_chain")
+ prev_chain = chain_raw if isinstance(chain_raw, int) else 0
+ verdict: dict[str, Any] = {
+ "resume_from": None,
+ "resume_candidate": int(row.prev_id),
+ "resume_buckets": done_n,
+ "resume_chain": 0,
+ }
+
+ if row.prev_status not in _RESUME_STATUSES:
+ verdict["resume_reason"] = f"status_{row.prev_status}"
+ elif not row.same_params:
+ verdict["resume_reason"] = "params_changed"
+ elif done_n == 0:
+ verdict["resume_reason"] = "no_checkpoint"
+ elif row.age_h is None or float(row.age_h) > 24.0 * (
+ _schedule_interval_days(row.interval_days) + 1
+ ):
+ verdict["resume_reason"] = "checkpoint_stale"
+ elif prev_chain >= _MAX_RESUME_CHAIN:
+ verdict["resume_reason"] = "chain_limit"
+ else:
+ verdict["resume_from"] = int(row.prev_id)
+ verdict["resume_reason"] = "ok"
+ verdict["resume_chain"] = prev_chain + 1
+ return int(row.prev_id), verdict
+ return None, verdict
+
+
+def _pick_resume(db: Session, run_id: int) -> int | None:
+ """Чекпоинт какого прогона наследует `run_id` (или None) + запись вердикта.
+
+ Тождество задания сверяется РОВНО по тем полям, которыми задание задаётся: source
+ (кандидат ищется в пределах одного source) и params целиком, побайтово. Кандидат
+ сравнивается с ТЕКУЩИМ прогоном, а не с расписанием, потому что именно params
+ прогона поехали в pipeline. На проде за 90 суток 89 корзин из 433 (21%) лежат в
+ прогонах, чьи params отличаются от следующего — по ним пропуск был бы неверным:
+ `incremental_days` меняет СМЫСЛ ключа (дочитано до watermark ≠ бакет перебран), а
+ `price_cap_per_bucket` меняет само дерево бисекции, то есть какие ключи существуют.
+ """
+ row = db.execute(_RESUME_CANDIDATE_SQL, {"rid": run_id}).fetchone()
+ resume_run_id, verdict = _resume_decision(row)
+ # Вердикт кладём в counters НОВОГО прогона: `update_heartbeat` мержит jsonb, поэтому
+ # последующие heartbeat'ы пайплайна его не затрут и он доживёт до финализатора.
+ _kit_runs.update_heartbeat(db, run_id, verdict)
+ logger.info("scheduler: resume run_id=%d — %s", run_id, verdict)
+ return resume_run_id
+
+
def _defer_next_run_at(db: Session, schedule_row: dict[str, Any]) -> None:
"""Сдвинуть next_run_at на следующее окно БЕЗ создания run (#1522).
@@ -705,7 +831,7 @@ async def _job_avito_full_load(
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)),
- resume_run_id=None,
+ resume_run_id=_pick_resume(db, run_id),
incremental_days=incremental_days,
)
@@ -725,7 +851,7 @@ async def _job_avito_full_load_exhaustive(
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)),
- resume_run_id=None,
+ resume_run_id=_pick_resume(db, run_id),
incremental_days=None,
)
@@ -796,7 +922,7 @@ async def _job_cian_full_load(
request_delay_sec=float(params.get("request_delay_sec", 4.0)),
enrich_detail=bool(params.get("enrich_detail", False)),
detail_top_n=int(params.get("detail_top_n", 0)),
- resume_run_id=None,
+ resume_run_id=_pick_resume(db, run_id),
)
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py
index 3c3d717a..95e2618a 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py
@@ -1095,7 +1095,9 @@ class AvitoScraper(BaseScraper):
leaf-бакета. Может быть async или sync. Исключение прерывает прогон.
on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket).
skip_buckets: множество ключей «room_label:lo:hi» уже завершённых бакетов —
- пагинация и on_bucket для них пропускаются. Probe-запросы выполняются.
+ пагинация и on_bucket для них пропускаются. В exhaustive-режиме probe-запросы
+ всё равно выполняются (skip проверяется уже в листе, после probe), в
+ инкрементальном probe'а нет — там пропускается весь бакет целиком.
since: если None (default) — EXHAUSTIVE bisection-обход (поведение без
изменений). Если задана date — INCREMENTAL: на каждый (комнатность ×
seed-брекет) последовательная пагинация newest-first с ранней остановкой,
@@ -1141,6 +1143,7 @@ class AvitoScraper(BaseScraper):
max_pages_per_bucket=max_pages_per_bucket,
secondary_only=secondary_only,
on_bucket=on_bucket,
+ skip_buckets=skip_buckets,
)
else:
await self._walk_price_range(
@@ -1390,14 +1393,17 @@ class AvitoScraper(BaseScraper):
# иначе фетчим её как обычную страницу (открытый брекет без probe-html).
sem = asyncio.Semaphore(concurrency)
first_url = self._build_rooms_url(room_slug, 1, _lo_param, _hi_param)
+ dropped_pages = 0 # страницы, не отдавшие карточки по отказу (не по пустоте)
async def _one_page(p: int) -> list[ScrapedLot]:
+ nonlocal dropped_pages
if p == 1 and html is not None:
return self._parse_html(html, source_url_base=first_url)
async with sem:
page_html = await self._fetch_rooms_page_html(room_slug, p, _lo_param, _hi_param)
await asyncio.sleep(self.request_delay_sec)
if page_html is None:
+ dropped_pages += 1
logger.warning(
"avito: page_html=None %s [%d, %s] page=%d — skipping page",
room_label,
@@ -1420,6 +1426,7 @@ class AvitoScraper(BaseScraper):
# Блокировки пробрасываем наверх (mark_banned в pipeline-обёртке).
if isinstance(res, AvitoBlockedError | AvitoRateLimitedError):
raise res
+ dropped_pages += 1
logger.warning(
"avito: page exception %s [%d, %s] page=%d — %r",
room_label,
@@ -1464,8 +1471,25 @@ class AvitoScraper(BaseScraper):
if key:
seen[key] = lot
+ # ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ─────
+ # `bucket_key` — это "room:lo:hi" и больше ничего, поэтому «сделано» едет
+ # отдельным аргументом. Три пути частичности сходятся здесь:
+ # 1. страница молча выпала (page_html=None выше);
+ # 2. исключение страницы проглочено gather'ом (return_exceptions=True);
+ # 3. признанный tail-loss — probe провалился (expected_total=None, открытый
+ # брекет) или страниц нужно больше, чем max_pages.
+ # Ни один из них не оставлял следа в чекпоинте: следующий прогон видел ключ и
+ # пропускал бакет. Пока признак не доехал до done_buckets, включать подхват
+ # нельзя — видимая потеря (перескрап) стала бы невидимой (пропуск страниц).
+ pages_needed = (
+ math.ceil(expected_total / _AVITO_OFFERS_PER_PAGE)
+ if expected_total is not None
+ else None
+ )
+ complete = dropped_pages == 0 and pages_needed is not None and pages_needed <= max_pages
logger.info(
- "avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d",
+ "avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d "
+ "complete=%s dropped_pages=%d",
room_label,
lo,
_hi_repr,
@@ -1473,11 +1497,13 @@ class AvitoScraper(BaseScraper):
collected_this_bucket,
dropped_nb,
len(seen),
+ complete,
+ dropped_pages,
)
# ── on_bucket callback: инкрементальный save ──────────────────────────
if on_bucket is not None and bucket_lots:
- res_cb = on_bucket(bucket_key, bucket_lots)
+ res_cb = on_bucket(bucket_key, bucket_lots, complete)
if inspect.isawaitable(res_cb):
await res_cb
@@ -1493,6 +1519,7 @@ class AvitoScraper(BaseScraper):
max_pages_per_bucket: int,
secondary_only: bool,
on_bucket: Callable[..., Any] | None,
+ skip_buckets: set[str] | None = None,
) -> None:
"""INCREMENTAL пагинация одного (комнатность × seed-брекет) с ранней остановкой.
@@ -1514,16 +1541,37 @@ class AvitoScraper(BaseScraper):
bucket_key, secondary_only-фильтр и дедуп в seen — идентичны _paginate_leaf_bucket.
on_bucket вызывается один раз для собранного бакета (async/sync-aware).
AvitoBlockedError/AvitoRateLimitedError из page-фетчей пробрасываются наверх.
+
+ skip_buckets: ключи, дочитанные ПРЕДЫДУЩИМ прогоном с ТЕМИ ЖЕ params (тождество
+ задания проверяет планировщик, `_pick_resume`). В инкрементальном режиме
+ «сделано» значит «дочитал до watermark `since`», а не «перебрал бакет целиком»,
+ поэтому смешивать такой ключ с exhaustive-ключом нельзя — они дословно совпадают
+ (6 из 11 seed-ключей), но означают разное. Пропуск здесь безопасен по покрытию
+ ровно потому, что окно ретроспективы не уже такта (#2674, гарантируется
+ _job_avito_full_load): бакет, дочитанный до watermark N суток назад, следующий
+ плановый прогон перечитает со своим since = сегодня−N и увидит всё, что успело
+ появиться. Раньше аргумент сюда не передавался вовсе — resume в боевом
+ (инкрементальном) режиме avito_full_load был чистым no-op: 134 из 242 корзин.
+
+ complete=False (см. `_paginate_leaf_bucket`) отдаётся, когда бакет НЕ дочитан до
+ watermark: страница выпала, или страниц не хватило (max_pages), или остановка
+ произошла по grace-эвристике «2 подряд недатированные страницы» — там watermark
+ не доказан, а не достигнут.
"""
_lo_param = lo if lo > 0 else None
_hi_param = hi # None → _build_rooms_url не ставит pmax
_hi_repr = "open" if hi is None else str(hi)
bucket_key = f"{room_label}:{lo}:{_hi_repr}"
+ if skip_buckets and bucket_key in skip_buckets:
+ logger.info("avito: skip bucket %s — already read to watermark (resume)", bucket_key)
+ return
+
bucket_lots: list[ScrapedLot] = []
pages_fetched = 0
not_fresh_streak = 0 # подряд идущие не-свежие (all None-или-старые) страницы
stop_reason = "max-pages" # перетирается ниже на реальную причину
+ complete = False # дочитан ли бакет до watermark; True только на честных стопах
for p in range(1, max_pages_per_bucket + 1):
# Последовательный фетч (НЕ asyncio.gather): early-stop требует читать
@@ -1540,6 +1588,7 @@ class AvitoScraper(BaseScraper):
page_lots = self._parse_html(page_html, source_url_base=page_url)
if not page_lots:
stop_reason = "end-of-pages (empty parse)"
+ complete = True # выдача кончилась — читать в этом брекете больше нечего
break
bucket_lots.extend(page_lots)
@@ -1556,6 +1605,7 @@ class AvitoScraper(BaseScraper):
# Есть даты, но все < since → newest-first гарантирует, что дальше
# только старее → стоп немедленно.
stop_reason = "early-stop (page all older than since)"
+ complete = True # watermark достигнут — ровно то, что значит «сделано»
break
# Все карточки undated (None) → не стопим сразу (могут быть свежие без
# даты), но копим streak; 2 подряд недатированные страницы → grace-стоп.
@@ -1581,7 +1631,7 @@ class AvitoScraper(BaseScraper):
logger.info(
"avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d "
- "unique_total=%d incremental since=%s stop=%s",
+ "unique_total=%d incremental since=%s stop=%s complete=%s",
room_label,
lo,
_hi_repr,
@@ -1591,11 +1641,12 @@ class AvitoScraper(BaseScraper):
len(seen),
since.isoformat(),
stop_reason,
+ complete,
)
# ── on_bucket callback: инкрементальный save ──────────────────────────
if on_bucket is not None and bucket_lots:
- res_cb = on_bucket(bucket_key, bucket_lots)
+ res_cb = on_bucket(bucket_key, bucket_lots, complete)
if inspect.isawaitable(res_cb):
await res_cb
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
index e046ef67..67b5c302 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
@@ -621,8 +621,10 @@ class CianScraper(BaseScraper):
# Страница 1 уже есть (html из probe выше); остальные — параллельно.
sem = asyncio.Semaphore(concurrency)
+ dropped_pages = 0 if html else 1 # пустой probe-html = страница 1 не собрана
async def _one_page(p: int) -> list[ScrapedLot]:
+ nonlocal dropped_pages
if p == 1:
# Используем уже полученный HTML от probe
return self._parse_serp_html(html) if html else []
@@ -630,6 +632,7 @@ class CianScraper(BaseScraper):
page_html = await self._fetch_page_html(rooms, p, _lo_param, hi)
await asyncio.sleep(self.request_delay_sec)
if page_html is None:
+ dropped_pages += 1
logger.warning(
"cian: page_html=None %s [%d, %s] page=%d — skipping page",
room_label,
@@ -648,6 +651,7 @@ class CianScraper(BaseScraper):
bucket_lots: list[ScrapedLot] = []
for p_idx, res in enumerate(page_results, start=1):
if isinstance(res, BaseException):
+ dropped_pages += 1
logger.warning(
"cian: page exception %s [%d, %s] page=%d — %r",
room_label,
@@ -674,8 +678,16 @@ class CianScraper(BaseScraper):
if key:
seen[key] = lot
+ # ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ─────
+ # См. avito/serp.py — те же два пути частичности (выпавшая страница,
+ # проглоченное gather'ом исключение) плюс признанный hard-cap выше. У cian это
+ # тяжелее: в SERP-слое НЕТ класса блок-исключения вообще (ban-детект #2625
+ # агрегатный, на выходе из скрапера), поэтому капча посреди бакета приходит сюда
+ # как page_html=None и раньше давала «сделанный» бакет из уцелевших страниц.
+ complete = dropped_pages == 0 and pages_needed <= max_pages
logger.info(
- "cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d",
+ "cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d "
+ "complete=%s dropped_pages=%d",
room_label,
lo,
_hi_repr,
@@ -683,11 +695,13 @@ class CianScraper(BaseScraper):
collected_this_bucket,
dropped_nb,
len(seen),
+ complete,
+ dropped_pages,
)
# ── on_bucket callback: инкрементальный save ──────────────────────────
if on_bucket is not None and bucket_lots:
- res_cb = on_bucket(bucket_key, bucket_lots)
+ res_cb = on_bucket(bucket_key, bucket_lots, complete)
if inspect.isawaitable(res_cb):
await res_cb