fix(scrape_kn): Redis-lock owner-token + resume retry on lock_held (#1216)
Some checks failed
Deploy / deploy (push) Blocked by required conditions
Deploy / changes (push) Successful in 5s
Deploy / build-backend (push) Successful in 1m34s
Deploy / build-frontend (push) Has been skipped
Deploy / build-worker (push) Has been cancelled

Два concurrency-бага в scrape_kn:

1. Lock value="1" без owner-токена + TTL=30мин = заявленной длительности
   sweep'а (нулевой запас). Sweep > TTL → lock истекает, второй sweep
   стартует, первый при завершении безусловным r.delete сносит ЧУЖОЙ
   lock → возможен третий параллельный.

2. SIGKILL/redeploy mid-sweep: finally не выполняется, lock живёт до
   30 мин. worker_ready метит run 'zombie' и enqueue'ит resume_kn_run
   БЕЗ countdown, который ловит lock_held и возвращает skipped без
   retry → run теряется до недельного beat.

Patch:
- _region_lock: value = uuid4().hex; release через Lua check-and-delete
  (`if GET == ARGV[1] then DEL else 0 end`) — не сносим чужой lock.
- _LOCK_TTL_SECONDS 30 → 45 мин (полтора max sweep duration, запас).
- resume_kn_run: при lock_held → raise self.retry(countdown=300,
  max_retries=12) → ~час окно для подхвата вместо silent skip.
- scrape_kn_region (scheduled) намеренно остаётся skipped — beat
  поднимет в следующий weekly tick (другая семантика).

11 новых юнит-тестов (token uniqueness, Lua guard, retry-not-skip,
existing behaviors). 16/16 scrape_kn тестов зелёные. ruff clean.

Closes #1216
This commit is contained in:
Light1YT 2026-06-13 10:36:08 +05:00 committed by bot-backend
parent b00e338667
commit c75e6a4cfa
2 changed files with 373 additions and 15 deletions

View file

@ -6,6 +6,8 @@ import asyncio
import logging import logging
import random import random
import time import time
import uuid
from collections.abc import Iterator
from contextlib import contextmanager from contextlib import contextmanager
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@ -18,10 +20,23 @@ from app.workers.celery_app import celery_app
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
# Lock TTL — full Свердл sweep с extras занимает ~25-30 мин. 30 мин — потолок. # Lock TTL — full Свердл sweep с extras занимает ~25-30 мин. Держим 45 мин —
# Если worker умер с lock'ом — он истечёт за полчаса, а не за 2 часа (раньше). # полтора max sweep duration: запас на хвостовые retry внутри sweep'а и не
# Heartbeat-based autoclean в /queue endpoint всё равно подхватит зомби-runs. # даём lock'у истечь под живым worker'ом (фикс P2 issue #1216 — раньше TTL
_LOCK_TTL_SECONDS = 30 * 60 # равнялся sweep duration → race на границе).
# Если worker умер с lock'ом — он истечёт за 45 мин, а не за 2 часа (раньше),
# heartbeat-based autoclean в /queue endpoint всё равно подхватит зомби-runs.
_LOCK_TTL_SECONDS = 45 * 60
# Lua check-and-delete: освобождаем lock ТОЛЬКО если value совпадает с нашим
# owner-токеном. Защита от P2 race: TTL может истечь под живым sweep'ом, второй
# sweep встанет под тем же key с НОВЫМ токеном; финализация первого без guard
# снесла бы чужой lock → возможен третий параллельный sweep (issue #1216).
_RELEASE_LOCK_LUA = (
'if redis.call("GET", KEYS[1]) == ARGV[1] '
'then return redis.call("DEL", KEYS[1]) '
"else return 0 end"
)
def _lock_key(region_code: int, developers: list[str] | None) -> str: def _lock_key(region_code: int, developers: list[str] | None) -> str:
@ -102,20 +117,31 @@ def force_release_lock(region_code: int, developers: list[str] | None) -> bool:
@contextmanager @contextmanager
def _region_lock(region_code: int, developers: list[str] | None): def _region_lock(region_code: int, developers: list[str] | None) -> Iterator[bool]:
"""Redis singleton lock so beat double-fire or two workers can't run the """Redis singleton lock so beat double-fire or two workers can't run the
same (region, devs) sweep concurrently. SET NX EX is atomic.""" same (region, devs) sweep concurrently. SET NX EX is atomic.
Owner-token (uuid4 hex) кладём в value и снимаем lock через Lua
check-and-delete иначе истекший по TTL lock, перехваченный другим
sweep'ом, был бы снесён первым worker'ом при finally возможен третий
параллельный sweep (P2 race, issue #1216).
"""
key = _lock_key(region_code, developers) key = _lock_key(region_code, developers)
token = uuid.uuid4().hex
r = redis.Redis.from_url(settings.redis_url) r = redis.Redis.from_url(settings.redis_url)
acquired = r.set(key, "1", nx=True, ex=_LOCK_TTL_SECONDS) acquired = r.set(key, token, nx=True, ex=_LOCK_TTL_SECONDS)
try: try:
yield bool(acquired) yield bool(acquired)
finally: finally:
if acquired: if acquired:
try: try:
r.delete(key) # Lua check-and-delete: атомарно, мимо python GIL → нельзя
except Exception: # снести чужой lock даже если наш TTL уже истёк.
pass r.eval(_RELEASE_LOCK_LUA, 1, key, token)
except Exception as e:
# Глотаем намеренно: lock истечёт сам по TTL. Логируем,
# чтобы не молча терять Redis-фейлы (не bare except).
logger.warning("release lock %s failed: %s", key, e)
@celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region", max_retries=2) @celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region", max_retries=2)
@ -188,11 +214,18 @@ def scrape_kn_region(
) )
@celery_app.task(bind=True, name="tasks.scrape_kn.resume_kn_run", max_retries=2) @celery_app.task(bind=True, name="tasks.scrape_kn.resume_kn_run", max_retries=12)
def resume_kn_run(self: Any, run_id: int) -> dict[str, Any]: def resume_kn_run(self: Any, run_id: int) -> dict[str, Any]:
"""Resume a previously interrupted sweep using objects_snapshot + """Resume a previously interrupted sweep using objects_snapshot +
progress_obj_index from kn_scrape_runs. Triggered by worker_ready hook progress_obj_index from kn_scrape_runs. Triggered by worker_ready hook
when a stale 'running' row is found (heartbeat >5 min).""" when a stale 'running' row is found (heartbeat >5 min).
На lock_held делаем self.retry(countdown=300, max_retries=12) окно
в ~час суммарно. Без retry resume пропускался до недельного beat
(P2 issue #1216): worker_ready помечал run 'zombie' и enqueue'ил resume
без countdown; live sweep держал lock resume получал skipped и run
терялся до следующего планового запуска.
"""
from sqlalchemy import text from sqlalchemy import text
from app.core.db import SessionLocal from app.core.db import SessionLocal
@ -234,12 +267,20 @@ def resume_kn_run(self: Any, run_id: int) -> dict[str, Any]:
with _region_lock(region_code, developers) as got_lock: with _region_lock(region_code, developers) as got_lock:
if not got_lock: if not got_lock:
logger.warning( logger.warning(
"resume_kn_run SKIPPED run=%s region=%s — lock held by another sweep", "resume_kn_run lock_held run=%s region=%s — retry in 5min"
" (live sweep still holds the lock)",
run_id, run_id,
region_code, region_code,
) )
_log_task_skipped(region_code, developers, self.request.id, "lock_held") _log_task_skipped(region_code, developers, self.request.id, "lock_held_retry")
return {"skipped": True, "reason": "lock_held", "run_id": run_id} # 5 мин × 12 retries = 60 мин окно. Покрывает full sweep
# duration (~30 мин) + запас. Celery поднимет Retry exception,
# task поставится обратно в очередь с countdown=300.
raise self.retry(
exc=RuntimeError(f"lock_held for run={run_id} region={region_code}"),
countdown=300,
max_retries=12,
)
logger.info("resume_kn_run start run=%s region=%s devs=%s", run_id, region_code, developers) logger.info("resume_kn_run start run=%s region=%s devs=%s", run_id, region_code, developers)
return asyncio.run( return asyncio.run(

View file

@ -0,0 +1,317 @@
"""Тесты для backend/app/workers/tasks/scrape_kn.py — Redis-lock owner-token
и resume_kn_run retry on lock_held (P2 issue #1216).
Покрытие:
- `_region_lock` value содержит уникальный uuid (а не литерал "1");
- два последовательных вызова получают разные токены (uniqueness);
- `r.eval` вызывается с Lua check-and-delete: правильный key + наш token;
- если в Redis под key лежит ЧУЖОЙ token Lua возвращает 0 (mock-симуляция);
- `_LOCK_TTL_SECONDS` поднят до 45 мин (полтора max sweep duration);
- `resume_kn_run` при lock_held `self.retry(countdown=300, max_retries=12)`,
а НЕ silent skip (issue #1216 fix #2).
"""
from __future__ import annotations
from typing import Any
from unittest.mock import MagicMock, patch
# ── _region_lock owner-token ──────────────────────────────────────────────────
def test_region_lock_value_is_uuid_token_not_literal_one() -> None:
"""`r.set` должен получить uuid4-hex (32 hex chars) — НЕ литерал '1'.
Это фикс fix #1 из issue #1216: без owner-токена finally безусловно
делал `r.delete(key)` можно было снести ЧУЖОЙ lock после TTL-expiry.
"""
from app.workers.tasks import scrape_kn as mod
mock_redis = MagicMock()
mock_redis.set.return_value = True # acquired
with patch.object(mod.redis.Redis, "from_url", return_value=mock_redis):
with mod._region_lock(66, None) as got_lock:
assert got_lock is True
# `r.set(key, token, nx=True, ex=...)` — token positional[1]
set_args, _set_kwargs = mock_redis.set.call_args
token = set_args[1]
assert token != "1", "lock value must be unique owner-token, not literal '1'"
assert isinstance(token, str)
assert len(token) == 32, f"uuid4.hex is 32 chars, got len={len(token)}"
int(token, 16) # raises if not hex
def test_region_lock_tokens_are_unique_across_calls() -> None:
"""Два последовательных context-manager call → два РАЗНЫХ токена.
Иначе owner-check (Lua GET == ARGV[1]) станет бессмысленным
оба worker'а получили бы одинаковый token, и второй снёс бы lock первого.
"""
from app.workers.tasks import scrape_kn as mod
mock_redis = MagicMock()
mock_redis.set.return_value = True
captured_tokens: list[str] = []
with patch.object(mod.redis.Redis, "from_url", return_value=mock_redis):
with mod._region_lock(66, None):
captured_tokens.append(mock_redis.set.call_args[0][1])
with mod._region_lock(66, None):
captured_tokens.append(mock_redis.set.call_args[0][1])
assert len(captured_tokens) == 2
assert captured_tokens[0] != captured_tokens[1], "tokens must be unique per call"
def test_region_lock_release_uses_lua_eval_with_own_token() -> None:
"""В finally вызывается `r.eval(LUA, 1, key, token)` — НЕ `r.delete(key)`.
Это фикс fix #1 из issue #1216 (check-and-delete семантика).
"""
from app.workers.tasks import scrape_kn as mod
mock_redis = MagicMock()
mock_redis.set.return_value = True
mock_redis.eval.return_value = 1 # пусть Lua говорит «удалил»
with patch.object(mod.redis.Redis, "from_url", return_value=mock_redis):
with mod._region_lock(66, ["devA"]):
set_token = mock_redis.set.call_args[0][1]
# eval должен быть вызван ровно один раз с правильными аргументами.
mock_redis.eval.assert_called_once()
eval_args = mock_redis.eval.call_args[0]
assert eval_args[0] == mod._RELEASE_LOCK_LUA
assert eval_args[1] == 1 # numkeys
assert eval_args[2] == mod._lock_key(66, ["devA"]) # KEYS[1]
assert eval_args[3] == set_token # ARGV[1] — наш token
# delete (без guard) НЕ должен вызываться — был бы регрессией к старому багу.
mock_redis.delete.assert_not_called()
def test_region_lock_release_skipped_when_not_acquired() -> None:
"""Если acquire не удался — finally НЕ должен делать eval/delete.
Иначе worker, получивший skipped, снёс бы lock владельца.
"""
from app.workers.tasks import scrape_kn as mod
mock_redis = MagicMock()
mock_redis.set.return_value = None # nx=True + key exists → None
with patch.object(mod.redis.Redis, "from_url", return_value=mock_redis):
with mod._region_lock(66, None) as got_lock:
assert got_lock is False
mock_redis.eval.assert_not_called()
mock_redis.delete.assert_not_called()
def test_release_lock_lua_script_contract() -> None:
"""Smoke: Lua script содержит check-and-delete семантику.
Если script отредактируют наугад этот guard упадёт.
"""
from app.workers.tasks.scrape_kn import _RELEASE_LOCK_LUA
assert 'redis.call("GET", KEYS[1])' in _RELEASE_LOCK_LUA
assert "ARGV[1]" in _RELEASE_LOCK_LUA
assert 'redis.call("DEL", KEYS[1])' in _RELEASE_LOCK_LUA
assert "else return 0 end" in _RELEASE_LOCK_LUA
def test_region_lock_uses_full_ttl_45_min() -> None:
"""`r.set` получает `ex=_LOCK_TTL_SECONDS == 45*60 == 2700`.
Issue #1216 fix: было 30 мин (== max sweep duration → нулевой запас);
поднимаем до 45 мин = полтора sweep duration.
"""
from app.workers.tasks import scrape_kn as mod
assert mod._LOCK_TTL_SECONDS == 45 * 60
mock_redis = MagicMock()
mock_redis.set.return_value = True
with patch.object(mod.redis.Redis, "from_url", return_value=mock_redis):
with mod._region_lock(66, None):
pass
_, set_kwargs = mock_redis.set.call_args
assert set_kwargs["ex"] == 2700
assert set_kwargs["nx"] is True
def test_lua_script_simulation_rejects_foreign_token() -> None:
"""Симуляция Lua check-and-delete на python-уровне:
если value под key != нашему token 0 (не удалили), иначе 1.
Гарантирует что контракт скрипта корректен; реальная атомарность
забота Redis (тестируется на integration-уровне).
"""
from app.workers.tasks.scrape_kn import _RELEASE_LOCK_LUA
def simulate(stored: str | None, our_token: str) -> int:
# Грубая модель: parse-free assertion на ключевые ветки скрипта.
assert 'if redis.call("GET", KEYS[1]) == ARGV[1]' in _RELEASE_LOCK_LUA
assert "else return 0 end" in _RELEASE_LOCK_LUA
if stored == our_token:
return 1
return 0
own = "a" * 32
foreign = "b" * 32
assert simulate(stored=own, our_token=own) == 1
assert simulate(stored=foreign, our_token=own) == 0
assert simulate(stored=None, our_token=own) == 0
# ── resume_kn_run retry on lock_held ──────────────────────────────────────────
def _patch_session_local(mock_db: MagicMock) -> Any:
"""Helper: SessionLocal импортируется внутри resume_kn_run из `app.core.db`,
поэтому патчить надо именно там, а не на module-level scrape_kn."""
return patch("app.core.db.SessionLocal", return_value=mock_db)
def test_resume_kn_run_retries_on_lock_held_not_silent_skip() -> None:
"""`resume_kn_run` при lock_held → `self.retry(countdown=300, max_retries=12)`.
Старое поведение (silent skip return {'skipped': True}) теряло run до
недельного beat это и есть P2 issue #1216 fix #2.
Patch'им `mod.resume_kn_run.retry` напрямую: при bind=True Celery передаёт
task-instance как `self`, и self.retry == task.retry. Подменяем его на
side_effect-exception, ловим, проверяем kwargs.
"""
from app.workers.tasks import scrape_kn as mod
row = {"region_codes": [66], "developer_ids": None}
mock_result = MagicMock()
mock_result.mappings.return_value.first.return_value = row
mock_db = MagicMock()
mock_db.execute.return_value = mock_result
mock_redis = MagicMock()
mock_redis.set.return_value = None # lock_held
class FakeRetryError(Exception):
pass
with (
_patch_session_local(mock_db),
patch.object(mod.redis.Redis, "from_url", return_value=mock_redis),
patch.object(mod, "_log_task_received"),
patch.object(mod, "_log_task_skipped"),
patch.object(
mod.resume_kn_run, "retry", side_effect=FakeRetryError("re-queue")
) as mock_retry,
):
raised = False
try:
mod.resume_kn_run.run(42)
except FakeRetryError:
raised = True
assert raised, "resume_kn_run must raise (via self.retry), not return"
mock_retry.assert_called_once()
retry_kwargs = mock_retry.call_args.kwargs
assert retry_kwargs["countdown"] == 300
assert retry_kwargs["max_retries"] == 12
assert "exc" in retry_kwargs
def test_resume_kn_run_lock_held_does_not_silently_return_skipped() -> None:
"""Регрессия-guard: при lock_held НЕ возвращаем dict — поднимаем Retry.
Подменяем self.retry на sentinel-exception; если код пошёл бы по
старой ветке (silent skip return), exception не возник бы.
"""
from app.workers.tasks import scrape_kn as mod
row = {"region_codes": [66], "developer_ids": None}
mock_result = MagicMock()
mock_result.mappings.return_value.first.return_value = row
mock_db = MagicMock()
mock_db.execute.return_value = mock_result
mock_redis = MagicMock()
mock_redis.set.return_value = None # lock_held
class SentinelRetryError(Exception):
pass
result_returned: Any = None
raised = False
with (
_patch_session_local(mock_db),
patch.object(mod.redis.Redis, "from_url", return_value=mock_redis),
patch.object(mod, "_log_task_received"),
patch.object(mod, "_log_task_skipped"),
patch.object(mod.resume_kn_run, "retry", side_effect=SentinelRetryError()),
):
try:
result_returned = mod.resume_kn_run.run(42)
except SentinelRetryError:
raised = True
assert raised, "resume_kn_run must call self.retry (raise Retry), not return"
assert result_returned is None
def test_resume_kn_run_run_not_found_still_returns_skipped() -> None:
"""Если row не найден — возвращаем `{'skipped': True, 'reason': 'run_not_found'}`.
Это легитимный skip (run уже почистили), retry бесполезен. Проверяем что
fix #2 не сломал валидные ранние выходы.
"""
from app.workers.tasks import scrape_kn as mod
mock_result = MagicMock()
mock_result.mappings.return_value.first.return_value = None
mock_db = MagicMock()
mock_db.execute.return_value = mock_result
with (
_patch_session_local(mock_db),
patch.object(mod.resume_kn_run, "retry") as mock_retry,
):
result = mod.resume_kn_run.run(999)
assert result == {"skipped": True, "reason": "run_not_found", "run_id": 999}
mock_retry.assert_not_called()
# ── scrape_kn_region остаётся skipped (а не retry) — это by-design ────────────
def test_scrape_kn_region_lock_held_returns_skipped() -> None:
"""`scrape_kn_region` (scheduled beat) при lock_held → return skipped.
В отличие от resume_kn_run, scheduled run крутится по weekly beat
пропустить разовый запуск ок, beat поднимет в следующий раз. Retry'ить
тут было бы спамом задач.
"""
from app.workers.tasks import scrape_kn as mod
mock_redis = MagicMock()
mock_redis.set.return_value = None # lock_held
with (
patch.object(mod.redis.Redis, "from_url", return_value=mock_redis),
patch.object(mod, "_log_task_received"),
patch.object(mod, "_log_task_skipped"),
patch.object(mod.settings, "scrape_kn_jitter_seconds", 0),
patch.object(mod.settings, "scrape_kn_state_path", ""),
patch.object(mod.scrape_kn_region, "retry") as mock_retry,
):
result = mod.scrape_kn_region.run(66, None, True)
assert result["skipped"] is True
assert result["reason"] == "lock_held"
assert result["region_code"] == 66
mock_retry.assert_not_called()