Merge pull request 'feat(tradein/scheduler): boot-reap — прогоны предыдущего контейнера снимаются на старте (#3122)' (#3123) from feat/3122-boot-reap into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m50s
Deploy Trade-In / build-backend (push) Successful in 1m37s
Deploy Trade-In / deploy (push) Successful in 1m58s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m50s
Deploy Trade-In / build-backend (push) Successful in 1m37s
Deploy Trade-In / deploy (push) Successful in 1m58s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s
This commit is contained in:
commit
dea08a0523
2 changed files with 147 additions and 1 deletions
91
tradein-mvp/backend/tests/test_3122_boot_reap.py
Normal file
91
tradein-mvp/backend/tests/test_3122_boot_reap.py
Normal file
|
|
@ -0,0 +1,91 @@
|
|||
"""#3122: boot-reap — прогоны предыдущего контейнера снимаются на старте, а не через 6ч.
|
||||
|
||||
Прод-факт (27.08, после ночного офлайна #3119): 4 прогона 'running' со стартами
|
||||
03:30–07:08 при старте контейнера 07:57 — процессы доказуемо мертвы, но
|
||||
has_running_run блокировал их источники до 6-часового порогового reap'а: до
|
||||
пяти часов слепоты на источник ровно после простоя.
|
||||
|
||||
Критерий boot-reap — started_at < старт процесса (минус минута на дрейф),
|
||||
пульс не участвует: ложные срабатывания класса #2702 (редкий пульс у длинных
|
||||
прогонов) невозможны по построению. Маркер counters.boot_reaped=true открывает
|
||||
таким зомби подхват чекпоинта (_resume_decision): у порогового zombie процесс
|
||||
может быть жив (движущаяся точка), у boot-зомби — гарантированно мёртв.
|
||||
|
||||
Красные на main по значению: reap-SQL не существует (capability), а
|
||||
resume-вердикт для boot-зомби — status_zombie вместо ok (по значению).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from datetime import UTC, datetime
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
from scraper_kit.orchestration import scheduler as sched
|
||||
|
||||
# ── 1. Сам boot-reap: SQL-критерий по границе старта ─────────────────────────
|
||||
|
||||
|
||||
class _ReapDb:
|
||||
def __init__(self, rows: list[Any]) -> None:
|
||||
self.rows = rows
|
||||
self.params: dict[str, Any] | None = None
|
||||
self.sql: str = ""
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
self.sql = str(stmt)
|
||||
self.params = dict(params or {})
|
||||
return MagicMock(fetchall=lambda: self.rows)
|
||||
|
||||
def commit(self) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def test_boot_reap_criterion_is_started_at_not_heartbeat() -> None:
|
||||
"""Критерий — started_at против границы старта процесса с минутой запаса;
|
||||
heartbeat в SQL не участвует (иначе вернулись бы ложные срабатывания #2702)."""
|
||||
db = _ReapDb([SimpleNamespace(id=5018, source="avito_newbuilding_sweep")])
|
||||
boot = datetime(2026, 8, 27, 7, 57, tzinfo=UTC)
|
||||
n = sched.reap_boot_zombies(db, boot)
|
||||
assert n == 1
|
||||
assert "started_at <" in db.sql and "60 seconds" in db.sql
|
||||
assert "heartbeat" not in db.sql
|
||||
assert "boot_reaped" in db.sql # маркер безопасного подхвата
|
||||
assert db.params and db.params["boot"] == boot
|
||||
|
||||
|
||||
# ── 2. Resume: boot-зомби подхватывается, пороговый — нет ────────────────────
|
||||
|
||||
|
||||
def _candidate(status: str, counters: dict[str, Any]) -> SimpleNamespace:
|
||||
return SimpleNamespace(
|
||||
prev_id=5040,
|
||||
prev_status=status,
|
||||
prev_counters=counters,
|
||||
same_params=True,
|
||||
age_h=2.0,
|
||||
interval_days="1",
|
||||
)
|
||||
|
||||
|
||||
def test_boot_reaped_zombie_is_resumable() -> None:
|
||||
"""Зомби С маркером boot_reaped — процесс гарантированно мёртв, чекпоинт
|
||||
безопасен → подхват. На main: status_zombie (красный по значению)."""
|
||||
rid, verdict = sched._resume_decision(
|
||||
_candidate("zombie", {"boot_reaped": True, "done_buckets": ["a", "b"]})
|
||||
)
|
||||
assert verdict["resume_reason"] == "ok", verdict
|
||||
assert rid == 5040
|
||||
|
||||
|
||||
def test_threshold_zombie_stays_rejected() -> None:
|
||||
"""Пороговый зомби БЕЗ маркера — процесс может быть жив (движущаяся точка,
|
||||
см. _RESUME_STATUSES) → по-прежнему отказ. Инвариант обеих эр."""
|
||||
rid, verdict = sched._resume_decision(_candidate("zombie", {"done_buckets": ["a", "b"]}))
|
||||
assert rid is None
|
||||
assert verdict["resume_reason"] == "status_zombie"
|
||||
|
|
@ -410,6 +410,50 @@ def reap_zombies(db: Session) -> int:
|
|||
return len(rows)
|
||||
|
||||
|
||||
def reap_boot_zombies(db: Session, boot_time: datetime) -> int:
|
||||
"""Однократный boot-reap (#3122): прогон, стартовавший раньше старта СОБСТВЕННОГО
|
||||
процесса планировщика, мёртв по построению — его сборщик жил в предыдущем
|
||||
контейнере и умер вместе с ним. Пульс в критерии не участвует вовсе, поэтому
|
||||
ложные срабатывания класса #2702 (редкий пульс у 5-часовых прогонов)
|
||||
невозможны: живой прогон этого процесса не может быть старше самого процесса.
|
||||
|
||||
Прод-цена без этого (27.08, после ночного офлайна #3119): 4 зомби блокировали
|
||||
свои источники через has_running_run до 6-часового порогового reap'а — до пяти
|
||||
часов слепоты на источник ровно после простоя, когда догон нужнее всего.
|
||||
|
||||
Минута запаса на часовой дрейф между хостом БД и приложением. Маркер
|
||||
`boot_reaped: true` в counters отличает этот исход от порогового 'zombie':
|
||||
у порогового процесс МОЖЕТ быть жив (reap не убивает его — см. комментарий
|
||||
у _RESUME_STATUSES), у boot-зомби — гарантированно мёртв, поэтому его
|
||||
чекпоинт безопасен для подхвата (_resume_decision пускает 'zombie' только
|
||||
с этим маркером).
|
||||
"""
|
||||
result = db.execute(
|
||||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'zombie',
|
||||
finished_at = clock_timestamp(),
|
||||
counters = COALESCE(counters, CAST('{}' AS jsonb))
|
||||
|| CAST('{"boot_reaped": true}' AS jsonb)
|
||||
WHERE status = 'running'
|
||||
AND started_at < CAST(:boot AS timestamptz) - CAST('60 seconds' AS interval)
|
||||
RETURNING id, source
|
||||
"""
|
||||
),
|
||||
{"boot": boot_time},
|
||||
)
|
||||
rows = result.fetchall()
|
||||
db.commit()
|
||||
if rows:
|
||||
logger.warning(
|
||||
"scheduler: boot-reap — %d прогонов предыдущего контейнера сняты с 'running': %s",
|
||||
len(rows),
|
||||
[(r.id, r.source) for r in rows],
|
||||
)
|
||||
return len(rows)
|
||||
|
||||
|
||||
def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> int | None:
|
||||
"""Claim run: INSERT scrape_runs + UPDATE scrape_schedules next_run_at.
|
||||
|
||||
|
|
@ -518,6 +562,10 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
|
|||
# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных
|
||||
# корзин — они записаны тем же heartbeat'ом, что и у banned.
|
||||
_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"})
|
||||
# 'zombie' допускается к подхвату ТОЛЬКО с маркером counters.boot_reaped (#3122):
|
||||
# у boot-зомби процесс гарантированно мёртв (жил в предыдущем контейнере), значит
|
||||
# «движущейся точки» — причины, по которой 'zombie' исключён выше, — быть не может.
|
||||
# Пороговый zombie без маркера по-прежнему отвергается.
|
||||
|
||||
# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) —
|
||||
# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта =
|
||||
|
|
@ -583,7 +631,8 @@ def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
|
|||
"resume_chain": 0,
|
||||
}
|
||||
|
||||
if row.prev_status not in _RESUME_STATUSES:
|
||||
_boot_reaped_zombie = row.prev_status == "zombie" and prev_counters.get("boot_reaped") is True
|
||||
if row.prev_status not in _RESUME_STATUSES and not _boot_reaped_zombie:
|
||||
verdict["resume_reason"] = f"status_{row.prev_status}"
|
||||
elif not row.same_params:
|
||||
verdict["resume_reason"] = "params_changed"
|
||||
|
|
@ -1035,8 +1084,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
|
|||
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
|
||||
"""
|
||||
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
|
||||
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
|
||||
# Initial sleep 30s чтобы дать FastAPI startup завершиться
|
||||
await asyncio.sleep(30)
|
||||
_boot_reap_done = False
|
||||
while True:
|
||||
# #1182 Phase 2: кооперативный SIGTERM-drain. Не reap'аем и не claim'аем
|
||||
# новые run'ы во время shutdown — даём текущему dispatch'у докатиться и выходим.
|
||||
|
|
@ -1046,6 +1097,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
|
|||
try:
|
||||
db = ctx.session_factory()
|
||||
try:
|
||||
# #3122: однократно на старте — снять прогоны предыдущего контейнера.
|
||||
if not _boot_reap_done:
|
||||
reap_boot_zombies(db, boot_time)
|
||||
_boot_reap_done = True
|
||||
# Reap zombies first
|
||||
reap_zombies(db)
|
||||
# #2670: раз в сутки — сводка «что сейчас не собирает». Календарная, а
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue