From 8196fe50d212c9b0363b12a749d0e8d35a8578d0 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Tue, 18 Aug 2026 10:54:29 +0300 Subject: [PATCH] Guard sold reconciliation by scope --- mobilede_scraper/worker/constants.py | 1 + mobilede_scraper/worker/refresh_cycle.py | 62 +++++++++++++++++++++--- mobilede_scraper/worker/search_sync.py | 21 +++++++- mobilede_scraper/worker/tasks.py | 34 +++++++++++++ tests/test_worker_runtime_tasks.py | 48 +++++++++++------- 5 files changed, 140 insertions(+), 26 deletions(-) diff --git a/mobilede_scraper/worker/constants.py b/mobilede_scraper/worker/constants.py index d9d390d..eff70d0 100644 --- a/mobilede_scraper/worker/constants.py +++ b/mobilede_scraper/worker/constants.py @@ -87,6 +87,7 @@ MOBILEDE_REFRESH_CYCLE_TOTAL_KEY = "mobilede:state:refresh_cycle:total" MOBILEDE_REFRESH_CYCLE_DONE_KEY = "mobilede:state:refresh_cycle:done" MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT = "mobilede:state:refresh_cycle:done_segments:{cycle_id}" MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT = "mobilede:state:refresh_cycle:finalized:{cycle_id}" +MOBILEDE_REFRESH_CYCLE_UNSAFE_REASONS_KEY_FMT = "mobilede:state:refresh_cycle:unsafe_reasons:{cycle_id}" MOBILEDE_REFRESH_CYCLE_TTL_SECONDS = max(60 * 60, int(os.getenv("MOBILEDE_REFRESH_CYCLE_TTL_SECONDS", str(24 * 60 * 60)))) MOBILEDE_OVERFLOW_CHILDREN_QUEUED_KEY_FMT = "mobilede:state:overflow_children_queued:{scope}:{parent_key}" MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED = os.getenv("MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} diff --git a/mobilede_scraper/worker/refresh_cycle.py b/mobilede_scraper/worker/refresh_cycle.py index 80fe03f..9e27519 100644 --- a/mobilede_scraper/worker/refresh_cycle.py +++ b/mobilede_scraper/worker/refresh_cycle.py @@ -15,6 +15,7 @@ from .constants import ( MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, MOBILEDE_REFRESH_CYCLE_TTL_SECONDS, + MOBILEDE_REFRESH_CYCLE_UNSAFE_REASONS_KEY_FMT, ) @@ -26,6 +27,41 @@ def _mobilede_refresh_cycle_finalized_key(cycle_id: str) -> str: return MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT.format(cycle_id=cycle_id) +def _mobilede_refresh_cycle_unsafe_reasons_key(cycle_id: str) -> str: + return MOBILEDE_REFRESH_CYCLE_UNSAFE_REASONS_KEY_FMT.format(cycle_id=cycle_id) + + +def _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client: Redis, + *, + cycle_id: str | None, + reason: str, + logger: logging.Logger, +) -> None: + if not cycle_id or not reason: + return + active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() + if active_cycle_id != cycle_id: + return + key = _mobilede_refresh_cycle_unsafe_reasons_key(cycle_id) + ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) + redis_client.sadd(key, reason) + redis_client.expire(key, ttl) + logger.warning("mobile.de refresh sold reconciliation disabled: cycle=%s reason=%s", cycle_id, reason) + + +def _mobilede_refresh_cycle_unsafe_reasons(redis_client: Redis, *, cycle_id: str) -> set[str]: + key = _mobilede_refresh_cycle_unsafe_reasons_key(cycle_id) + try: + return { + _mobilede_redis_text(value).strip() + for value in redis_client.smembers(key) + if _mobilede_redis_text(value).strip() + } + except Exception: + return set() + + def _mobilede_redis_text(value: object) -> str: if isinstance(value, bytes): return value.decode("utf-8", errors="ignore") @@ -55,6 +91,7 @@ def _mobilede_clear_refresh_cycle(redis_client: Redis) -> None: if cycle_id: keys.append(_mobilede_refresh_cycle_done_set_key(cycle_id)) keys.append(_mobilede_refresh_cycle_finalized_key(cycle_id)) + keys.append(_mobilede_refresh_cycle_unsafe_reasons_key(cycle_id)) redis_client.delete(*keys) @@ -76,6 +113,7 @@ def _mobilede_start_refresh_cycle( pipe.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, "0", ex=ttl) pipe.delete(done_set_key) pipe.delete(_mobilede_refresh_cycle_finalized_key(cycle_id)) + pipe.delete(_mobilede_refresh_cycle_unsafe_reasons_key(cycle_id)) pipe.execute() logger.info("mobile.de refresh cycle started: id=%s total=%s started_at=%s", cycle_id, total, started_at) return cycle_id @@ -180,6 +218,20 @@ def _mobilede_finalize_refresh_cycle_sold_marking( post_refresh_probe_batch_size: int, schedule_post_refresh_probe: Callable[[], None], ) -> int: + active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() + done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0) + total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or 0) + unsafe_reasons = _mobilede_refresh_cycle_unsafe_reasons(redis_client, cycle_id=cycle_id) + if active_cycle_id != cycle_id or total_now <= 0 or done_now < total_now or unsafe_reasons: + logger.warning( + "mobile.de refresh sold marking skipped: cycle=%s active=%s progress=%s/%s unsafe_reasons=%s", + cycle_id, + active_cycle_id, + done_now, + total_now, + sorted(unsafe_reasons), + ) + return 0 started_at_raw = redis_client.get(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY) if not started_at_raw: logger.warning("mobile.de refresh sold marking skipped: missing started_at for cycle=%s", cycle_id) @@ -200,8 +252,6 @@ def _mobilede_finalize_refresh_cycle_sold_marking( prefix=origin_prefixes, safety_ratio=0.8, ) - done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0) - total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or 0) logger.info( "mobile.de refresh sold marking completed: cycle=%s progress=%s/%s cutoff=%s sold_marked=%s", cycle_id, @@ -243,17 +293,13 @@ def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( if total_now <= 0: return False if done_now < total_now: - redis_client.set( - MOBILEDE_REFRESH_CYCLE_DONE_KEY, - str(total_now), - ex=max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)), - ) logger.warning( - "mobile.de refresh finalize fallback: bootstrap completed while refresh progress lagged cycle=%s done=%s/%s", + "mobile.de refresh sold marking skipped: bootstrap completed while refresh progress is incomplete cycle=%s done=%s/%s", cycle_id, done_now, total_now, ) + return False if not _mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id): return False finalize_refresh_cycle_sold_marking(redis_client, cycle_id=cycle_id) diff --git a/mobilede_scraper/worker/search_sync.py b/mobilede_scraper/worker/search_sync.py index 1307601..ab55437 100644 --- a/mobilede_scraper/worker/search_sync.py +++ b/mobilede_scraper/worker/search_sync.py @@ -475,6 +475,13 @@ def run_mobilede_sync_search_task( bootstrap_progress: tuple[int, int, int, int, int, int, int] | None = None overflow_segments_added = 0 queued_overflow_children = 0 + early_stopped = bool(result.get("early_stopped")) + if refresh_cycle_id and early_stopped: + _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client, + cycle_id=refresh_cycle_id, + reason="early_stopped", + ) segment_max_pages = int((segment or {}).get("max_pages") or max_pages or MOBILEDE_MAX_PAGE_NUMBER) segment_scan_complete = bool( (use_cursor and _mobilede_segment_scan_complete(redis_client, segment)) @@ -517,10 +524,16 @@ def run_mobilede_sync_search_task( ex=max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)), ) late_overflow_pending = bool(segment_scan_complete and (overflow_segments_added or queued_overflow_children)) + if refresh_cycle_id and late_overflow_pending: + _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client, + cycle_id=refresh_cycle_id, + reason="plan_changed_or_incomplete", + ) refresh_done = 0 refresh_total = 0 should_finalize_refresh = False - if refresh_cycle_id and segment_scan_complete: + if refresh_cycle_id and segment_scan_complete and not early_stopped: refresh_done, refresh_total, should_finalize_refresh = _mobilede_track_refresh_cycle_segment( redis_client, cycle_id=refresh_cycle_id, @@ -1073,6 +1086,12 @@ def run_mobilede_sync_search_task( status_code = getattr(getattr(exc, "response", None), "status_code", None) is_antibot_block = int(status_code or 0) in {401, 403, 429} is_transient_request_error = _is_mobilede_transient_request_error(exc) + if refresh_cycle_id: + _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client, + cycle_id=refresh_cycle_id, + reason="http_403" if int(status_code or 0) == 403 else "request_error", + ) max_retries = int(getattr(self, "max_retries", 0) or 0) current_retries = int(getattr(self.request, "retries", 0) or 0) if is_transient_request_error: diff --git a/mobilede_scraper/worker/tasks.py b/mobilede_scraper/worker/tasks.py index 9e021c9..0ff4b08 100644 --- a/mobilede_scraper/worker/tasks.py +++ b/mobilede_scraper/worker/tasks.py @@ -39,6 +39,7 @@ from .refresh_cycle import ( _mobilede_clear_refresh_cycle as _refresh_cycle_clear_state, _mobilede_finalize_refresh_cycle_sold_marking as _refresh_cycle_finalize_sold_marking, _mobilede_get_or_start_refresh_cycle as _refresh_cycle_get_or_start, + _mobilede_mark_refresh_cycle_reconciliation_unsafe as _refresh_cycle_mark_unsafe, _mobilede_redis_text as _refresh_cycle_redis_text, _mobilede_refresh_cycle_done_set_key as _refresh_cycle_done_set_key_impl, _mobilede_refresh_cycle_finalized_key as _refresh_cycle_finalized_key_impl, @@ -1435,6 +1436,20 @@ def _mobilede_get_or_start_refresh_cycle(redis_client: Redis, *, total_segments: return _refresh_cycle_get_or_start(redis_client, total_segments=total_segments, logger=logger) +def _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client: Redis, + *, + cycle_id: str | None, + reason: str, +) -> None: + _refresh_cycle_mark_unsafe( + redis_client, + cycle_id=cycle_id, + reason=reason, + logger=logger, + ) + + def _mobilede_refresh_cycle_seen_at(redis_client: Redis, *, cycle_id: str | None) -> datetime | None: return _refresh_cycle_seen_at_impl(redis_client, cycle_id=cycle_id, logger=logger) @@ -2364,6 +2379,25 @@ def mobilede_sync_runtime_segments_task( redis_client, total_segments=len(cached_segments or []), ) + partial_scope_reasons: list[str] = [] + if runtime_config.sync.only_new is True: + partial_scope_reasons.append("only_new") + if runtime_config.sync.limit is not None: + partial_scope_reasons.append("sync_limit") + runtime_filters = getattr(runtime_config, "filters", None) + if runtime_filters is not None and not runtime_filters.is_empty(): + partial_scope_reasons.append("runtime_filters") + if getattr(settings.listing, "filtered_search_urls", None): + partial_scope_reasons.append("filtered_search_urls") + mobilede_config = getattr(runtime_config, "mobilede", None) + if mobilede_config is not None and getattr(mobilede_config, "segments", ()): + partial_scope_reasons.append("runtime_segments") + for reason in partial_scope_reasons: + _mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client, + cycle_id=refresh_cycle_id, + reason=reason, + ) if bootstrap_recovery_pending: redis_client.delete(MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY) segments = _enqueue_mobilede_runtime_segments( diff --git a/tests/test_worker_runtime_tasks.py b/tests/test_worker_runtime_tasks.py index 9f808bb..78088fa 100644 --- a/tests/test_worker_runtime_tasks.py +++ b/tests/test_worker_runtime_tasks.py @@ -39,6 +39,9 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): def scard(self, key: str): return len(self.sets.get(key, set())) + def smembers(self, key: str): + return self.sets.get(key, set()) + def llen(self, key: str): return len(self.lists.get(key, [])) @@ -277,33 +280,44 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertFalse(recovered) self.assertIn(tasks.MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, redis_client.sets) - def test_bootstrap_completion_forces_refresh_sold_finalize_when_progress_lags(self) -> None: + def test_bootstrap_completion_does_not_finalize_incomplete_refresh_cycle(self) -> None: redis_client = self._FakeRedis() started_at = datetime(2026, 5, 6, 10, 0, tzinfo=timezone.utc) redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_ID_KEY] = "cycle-lag" redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY] = started_at.isoformat() redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_TOTAL_KEY] = "377" redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY] = "141" - persistence = MagicMock() - persistence.mark_sold_not_seen_since.return_value = 42 + finalized = tasks._mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( + redis_client, + cycle_id="cycle-lag", + total_segments_hint=377, + ) - with ( - patch.object(tasks, "_get_persistence", return_value=persistence), - patch.object(tasks, "MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED", False), - ): - finalized = tasks._mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( + self.assertFalse(finalized) + self.assertEqual(redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY], "141") + + def test_refresh_cycle_unsafe_reason_skips_sold_marking(self) -> None: + redis_client = self._FakeRedis() + started_at = datetime(2026, 5, 6, 10, 0, tzinfo=timezone.utc) + redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_ID_KEY] = "cycle-unsafe" + redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY] = started_at.isoformat() + redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_TOTAL_KEY] = "1" + redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY] = "1" + persistence = MagicMock() + + tasks._mobilede_mark_refresh_cycle_reconciliation_unsafe( + redis_client, + cycle_id="cycle-unsafe", + reason="only_new", + ) + with patch.object(tasks, "_get_persistence", return_value=persistence): + sold_marked = tasks._mobilede_finalize_refresh_cycle_sold_marking( redis_client, - cycle_id="cycle-lag", - total_segments_hint=377, + cycle_id="cycle-unsafe", ) - self.assertTrue(finalized) - self.assertEqual(redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY], "377") - persistence.mark_sold_not_seen_since.assert_called_once_with( - started_at, - prefix=("mobile.de:", "mobilede:"), - safety_ratio=0.8, - ) + self.assertEqual(sold_marked, 0) + persistence.mark_sold_not_seen_since.assert_not_called() def test_completed_bootstrap_segment_is_skipped_during_active_full_pass(self) -> None: redis_client = MagicMock()