From 9a33efa89b88bbf2b5680b960e1fa4f983324e86 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Mon, 27 Apr 2026 13:09:34 +0300 Subject: [PATCH] Fix self heal stale progress cleanup --- iaai_scraper/worker/self_heal.py | 45 ++++++++++++++-- tests/test_self_heal.py | 92 ++++++++++++++++++++++++++++++++ 2 files changed, 134 insertions(+), 3 deletions(-) diff --git a/iaai_scraper/worker/self_heal.py b/iaai_scraper/worker/self_heal.py index 762359f..fc027a3 100644 --- a/iaai_scraper/worker/self_heal.py +++ b/iaai_scraper/worker/self_heal.py @@ -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 diff --git a/tests/test_self_heal.py b/tests/test_self_heal.py index 24eb4ce..a18c582 100644 --- a/tests/test_self_heal.py +++ b/tests/test_self_heal.py @@ -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")