Fix self heal stale progress cleanup
This commit is contained in:
@@ -19,6 +19,15 @@ SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM)
|
||||
TERMINAL_PROGRESS_STAGES = {
|
||||
"segment_done",
|
||||
"segment_failed",
|
||||
"segment_task_completed",
|
||||
"segment_task_failed",
|
||||
"segment_task_soft_timeout",
|
||||
"sync_done",
|
||||
"failed",
|
||||
}
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
@@ -148,13 +157,43 @@ def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict
|
||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
||||
has_lock = 1 if redis_client.get(SYNC_LISTING_LOCK_KEY) else 0
|
||||
has_task_progress = 0
|
||||
for _ in redis_client.scan_iter(match="iaai:state:task_progress:*"):
|
||||
has_task_progress = 1
|
||||
break
|
||||
progress_keys_seen = 0
|
||||
stale_progress_deleted = 0
|
||||
stale_seconds = max(180, min(_env_int("IAAI_SELF_HEAL_STALL_SECONDS", 720), 1800))
|
||||
now_ts = int(time.time())
|
||||
for key in redis_client.scan_iter(match="iaai:state:task_progress:*"):
|
||||
progress_keys_seen += 1
|
||||
if queue_len > 0 or has_lock == 1:
|
||||
has_task_progress = 1
|
||||
break
|
||||
try:
|
||||
payload = redis_client.get(key)
|
||||
if not payload:
|
||||
continue
|
||||
data = json.loads(payload)
|
||||
stage = str(data.get("stage") or "")
|
||||
ts = _safe_int(data.get("ts"), 0)
|
||||
is_terminal = stage in TERMINAL_PROGRESS_STAGES
|
||||
is_stale = ts <= 0 or now_ts - ts > stale_seconds
|
||||
if is_terminal or is_stale:
|
||||
redis_client.delete(key)
|
||||
stale_progress_deleted += 1
|
||||
continue
|
||||
has_task_progress = 1
|
||||
break
|
||||
except Exception:
|
||||
# Битые progress payload не должны держать worker в вечном false-stall цикле.
|
||||
try:
|
||||
redis_client.delete(key)
|
||||
stale_progress_deleted += 1
|
||||
except Exception:
|
||||
pass
|
||||
flags = {
|
||||
"queue_len": queue_len,
|
||||
"has_lock": has_lock,
|
||||
"has_task_progress": has_task_progress,
|
||||
"progress_keys_seen": progress_keys_seen,
|
||||
"stale_progress_deleted": stale_progress_deleted,
|
||||
}
|
||||
return (queue_len > 0 or has_lock == 1 or has_task_progress == 1), flags
|
||||
|
||||
|
||||
@@ -72,6 +72,98 @@ class TestSelfHeal(unittest.TestCase):
|
||||
|
||||
self.assertTrue(active)
|
||||
|
||||
@patch("iaai_scraper.worker.self_heal.time.time", return_value=5000)
|
||||
def test_has_inflight_work_deletes_stale_terminal_progress_without_work(self, _time_mock) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.llen.return_value = 0
|
||||
redis_client.get.side_effect = lambda key: {
|
||||
self_heal.SYNC_LISTING_LOCK_KEY: None,
|
||||
"iaai:state:task_progress:old": json.dumps(
|
||||
{
|
||||
"task_id": "old",
|
||||
"stage": "segment_done",
|
||||
"ts": 1000,
|
||||
"segment_full_scan_completed": True,
|
||||
}
|
||||
),
|
||||
}.get(key)
|
||||
redis_client.scan_iter.return_value = ["iaai:state:task_progress:old"]
|
||||
|
||||
has_inflight, flags = self_heal._has_inflight_work(redis_client, self_heal.IAAI_SYNC_QUEUE)
|
||||
|
||||
self.assertFalse(has_inflight)
|
||||
self.assertEqual(flags["has_task_progress"], 0)
|
||||
redis_client.delete.assert_called_once_with("iaai:state:task_progress:old")
|
||||
|
||||
@patch("iaai_scraper.worker.self_heal.time.time", return_value=1010)
|
||||
def test_has_inflight_work_keeps_recent_non_terminal_progress(self, _time_mock) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.llen.return_value = 0
|
||||
redis_client.get.side_effect = lambda key: {
|
||||
self_heal.SYNC_LISTING_LOCK_KEY: None,
|
||||
"iaai:state:task_progress:active": json.dumps(
|
||||
{
|
||||
"task_id": "active",
|
||||
"stage": "fast_listing_collected",
|
||||
"ts": 1000,
|
||||
}
|
||||
),
|
||||
}.get(key)
|
||||
redis_client.scan_iter.return_value = ["iaai:state:task_progress:active"]
|
||||
|
||||
has_inflight, flags = self_heal._has_inflight_work(redis_client, self_heal.IAAI_SYNC_QUEUE)
|
||||
|
||||
self.assertTrue(has_inflight)
|
||||
self.assertEqual(flags["has_task_progress"], 1)
|
||||
redis_client.delete.assert_not_called()
|
||||
|
||||
@patch("iaai_scraper.worker.self_heal.time.time", return_value=5000)
|
||||
def test_has_inflight_work_deletes_stale_non_terminal_progress_without_work(self, _time_mock) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.llen.return_value = 0
|
||||
redis_client.get.side_effect = lambda key: {
|
||||
self_heal.SYNC_LISTING_LOCK_KEY: None,
|
||||
"iaai:state:task_progress:stale": json.dumps(
|
||||
{
|
||||
"task_id": "stale",
|
||||
"stage": "segment_started",
|
||||
"ts": 1000,
|
||||
"segment_index": 11,
|
||||
}
|
||||
),
|
||||
}.get(key)
|
||||
redis_client.scan_iter.return_value = ["iaai:state:task_progress:stale"]
|
||||
|
||||
has_inflight, flags = self_heal._has_inflight_work(redis_client, self_heal.IAAI_SYNC_QUEUE)
|
||||
|
||||
self.assertFalse(has_inflight)
|
||||
self.assertEqual(flags["has_task_progress"], 0)
|
||||
self.assertEqual(flags["stale_progress_deleted"], 1)
|
||||
redis_client.delete.assert_called_once_with("iaai:state:task_progress:stale")
|
||||
|
||||
@patch("iaai_scraper.worker.self_heal.time.time", return_value=5000)
|
||||
def test_has_inflight_work_keeps_stale_progress_when_lock_exists(self, _time_mock) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.llen.return_value = 0
|
||||
redis_client.get.side_effect = lambda key: {
|
||||
self_heal.SYNC_LISTING_LOCK_KEY: "owner",
|
||||
"iaai:state:task_progress:stale": json.dumps(
|
||||
{
|
||||
"task_id": "stale",
|
||||
"stage": "segment_started",
|
||||
"ts": 1000,
|
||||
}
|
||||
),
|
||||
}.get(key)
|
||||
redis_client.scan_iter.return_value = ["iaai:state:task_progress:stale"]
|
||||
|
||||
has_inflight, flags = self_heal._has_inflight_work(redis_client, self_heal.IAAI_SYNC_QUEUE)
|
||||
|
||||
self.assertTrue(has_inflight)
|
||||
self.assertEqual(flags["has_lock"], 1)
|
||||
self.assertEqual(flags["has_task_progress"], 1)
|
||||
redis_client.delete.assert_not_called()
|
||||
|
||||
@patch("iaai_scraper.worker.self_heal.time.sleep", return_value=None)
|
||||
@patch("iaai_scraper.worker.self_heal.os.kill")
|
||||
@patch("builtins.open")
|
||||
|
||||
Reference in New Issue
Block a user