feat(tradein): чекпоинты #3074 — наследование при claim + yandex_city_sweep #3098

Merged
bot-backend merged 2 commits from fix/3074-checkpoint-survives-claim into main 2026-08-26 07:49:29 +00:00
2 changed files with 74 additions and 0 deletions
Showing only changes of commit e4b3c6cc2b - Show all commits

View file

@ -0,0 +1,63 @@
"""#3074: чекпоинт предшественника наследуется в counters НОВОГО прогона при claim.
Прод-факт (Beget, 23.08): прогон 4707 подхватил чекпоинт 4117 (42 корзины,
resume_reason='ok'), был убит деплоем на 26-й минуте ДО завершения первой
НОВОЙ корзины и не успел ни разу написать heartbeat с done_buckets. Его
собственный чекпоинт остался пуст: следующий кандидат увидел бы no_checkpoint,
и 42 корзины пропали бы, хотя цепочка (resume_chain=2) ещё позволяла подхват.
Фикс: _resume_decision при вердикте 'ok' кладёт done_buckets предшественника в
counters-заготовку нового прогона она пишется в БД прямо при claim
(_pick_resume update_heartbeat), до старта пайплайна. Heartbeat мержит jsonb,
первый настоящий bucket-heartbeat перезапишет ключ надмножеством.
Красный на origin/main по ЗНАЧЕНИЮ: 'done_buckets' отсутствует в вердикте.
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from types import SimpleNamespace
from typing import Any
from scraper_kit.orchestration import scheduler as sched
def _candidate(**over: Any) -> SimpleNamespace:
"""Кандидат из _RESUME_CANDIDATE_SQL: прогон 4117 как он лежал на проде 23.08."""
base = {
"prev_id": 4117,
"prev_status": "banned",
"prev_counters": {
"resume_chain": 1,
"done_buckets": [f"room_1_komn:{i}:0" for i in range(42)],
},
"same_params": True,
"age_h": 164.8,
"interval_days": "7",
}
base.update(over)
return SimpleNamespace(**base)
def test_ok_verdict_carries_predecessor_checkpoint() -> None:
"""Вердикт 'ok' несёт done_buckets предшественника — чекпоинт персистится при
claim и переживает обрыв до первой новой корзины (кейс 4707)."""
resume_id, verdict = sched._resume_decision(_candidate())
assert resume_id == 4117
assert verdict["resume_reason"] == "ok"
assert verdict.get("done_buckets") == sorted(f"room_1_komn:{i}:0" for i in range(42)), (
"counters-заготовка нового прогона не содержит чекпоинт предшественника — "
"обрыв до первой новой корзины снова потеряет всю цепочку"
)
def test_refusal_verdict_carries_no_checkpoint() -> None:
"""Отказ от подхвата чекпоинт не наследует: прогон честно идёт с нуля, и его
counters не должны врать, будто корзины предшественника уже собраны."""
_id, verdict = sched._resume_decision(_candidate(same_params=False))
assert verdict["resume_reason"] == "params_changed"
assert "done_buckets" not in verdict

View file

@ -599,6 +599,17 @@ def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
verdict["resume_from"] = int(row.prev_id)
verdict["resume_reason"] = "ok"
verdict["resume_chain"] = prev_chain + 1
# Наследуем чекпоинт в counters нового прогона ПРЯМО ПРИ CLAIM (#3074).
# Прод-факт (run 4707, 23.08): возобновлённый прогон, убитый деплоем на
# 26-й минуте — ДО завершения первой новой корзины, — не успел ни разу
# написать heartbeat с done_buckets. Его собственный чекпоинт остался
# пуст, кандидат следующей субботы увидел no_checkpoint, и 42 корзины
# предшественника (4117) пропали. Пайплайн сеет `done = set(skip_set)`
# у себя в памяти, но до первого _on_bucket это знание нигде не
# персистится; heartbeat мержит jsonb (`counters || :counters`), так что
# первый же настоящий bucket-heartbeat перезапишет ключ тем же множеством
# плюс новое — двойной записи не возникает.
verdict["done_buckets"] = sorted(done_buckets) if done_buckets else []
return int(row.prev_id), verdict
return None, verdict