Guard sold reconciliation by scope

This commit is contained in:
qananasikq
2026-08-18 10:54:29 +03:00
parent ca1fab1034
commit 8196fe50d2
5 changed files with 140 additions and 26 deletions
+1
View File
@@ -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"}
+54 -8
View File
@@ -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)
+20 -1
View File
@@ -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:
+34
View File
@@ -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(
+31 -17
View File
@@ -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()