from __future__ import annotations import logging import uuid from datetime import datetime, timezone from typing import Callable from redis import Redis from .constants import ( MOBILEDE_REFRESH_CYCLE_DONE_KEY, MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT, MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT, MOBILEDE_REFRESH_CYCLE_ID_KEY, MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, MOBILEDE_REFRESH_CYCLE_TTL_SECONDS, ) def _mobilede_refresh_cycle_done_set_key(cycle_id: str) -> str: return MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT.format(cycle_id=cycle_id) def _mobilede_refresh_cycle_finalized_key(cycle_id: str) -> str: return MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT.format(cycle_id=cycle_id) def _mobilede_redis_text(value: object) -> str: if isinstance(value, bytes): return value.decode("utf-8", errors="ignore") return str(value or "") def _mobilede_claim_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> bool: ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) return bool( redis_client.set( _mobilede_refresh_cycle_finalized_key(cycle_id), "1", nx=True, ex=ttl, ) ) def _mobilede_clear_refresh_cycle(redis_client: Redis) -> None: cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() keys = [ MOBILEDE_REFRESH_CYCLE_ID_KEY, MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, MOBILEDE_REFRESH_CYCLE_DONE_KEY, ] if cycle_id: keys.append(_mobilede_refresh_cycle_done_set_key(cycle_id)) keys.append(_mobilede_refresh_cycle_finalized_key(cycle_id)) redis_client.delete(*keys) def _mobilede_start_refresh_cycle( redis_client: Redis, *, total_segments: int, logger: logging.Logger, ) -> str: cycle_id = uuid.uuid4().hex[:16] started_at = datetime.now(timezone.utc).isoformat() total = max(0, int(total_segments)) ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id) pipe = redis_client.pipeline() pipe.set(MOBILEDE_REFRESH_CYCLE_ID_KEY, cycle_id, ex=ttl) pipe.set(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, started_at, ex=ttl) pipe.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total), ex=ttl) 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.execute() logger.info("mobile.de refresh cycle started: id=%s total=%s started_at=%s", cycle_id, total, started_at) return cycle_id def _mobilede_get_or_start_refresh_cycle( redis_client: Redis, *, total_segments: int, logger: logging.Logger, ) -> str: active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() if active_cycle_id: 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 total_segments or 0) if total_now <= 0 or done_now < total_now: return active_cycle_id return _mobilede_start_refresh_cycle(redis_client, total_segments=total_segments, logger=logger) def _mobilede_refresh_cycle_seen_at(redis_client: Redis, *, cycle_id: str | None, logger: logging.Logger) -> datetime | None: if not cycle_id: return None active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() if active_cycle_id != cycle_id: return None started_at_raw = redis_client.get(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY) if not started_at_raw: return None try: seen_at = datetime.fromisoformat(_mobilede_redis_text(started_at_raw)) if seen_at.tzinfo is None: seen_at = seen_at.replace(tzinfo=timezone.utc) return seen_at except Exception: logger.warning("mobile.de refresh cycle seen_at is invalid: cycle=%s value=%s", cycle_id, started_at_raw) return None def _mobilede_refresh_cycle_progress( redis_client: Redis, *, cycle_id: str | None, total_segments_hint: int = 0, ) -> tuple[int, int, int]: if not cycle_id: return 0, max(0, int(total_segments_hint)), 0 active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() if active_cycle_id != cycle_id: return 0, max(0, int(total_segments_hint)), 0 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 total_segments_hint or 0) if total_now <= 0: total_now = max(done_now, int(total_segments_hint or 0)) left_now = max(0, total_now - min(done_now, total_now)) if total_now > 0 else 0 return done_now, total_now, left_now def _mobilede_track_refresh_cycle_segment( redis_client: Redis, *, cycle_id: str | None, segment: dict[str, object] | None, total_segments_hint: int, segment_fingerprint: Callable[[dict[str, object] | None], str], ) -> tuple[int, int, bool]: if not cycle_id or not segment: return 0, max(0, int(total_segments_hint)), False active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() if active_cycle_id != cycle_id: return 0, 0, False done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id) if int(redis_client.sadd(done_set_key, segment_fingerprint(segment)) or 0) <= 0: 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 total_segments_hint or 0) return done_now, total_now, False ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) redis_client.expire(done_set_key, ttl) done_now = int(redis_client.incr(MOBILEDE_REFRESH_CYCLE_DONE_KEY)) redis_client.expire(MOBILEDE_REFRESH_CYCLE_DONE_KEY, ttl) total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0) if total_now <= 0: total_now = max(done_now, int(total_segments_hint or 0)) redis_client.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total_now), ex=ttl) if total_now > 0 and done_now > total_now: done_now = total_now redis_client.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, str(done_now), ex=ttl) should_finalize = False if total_now > 0 and done_now >= total_now: should_finalize = _mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id) return done_now, total_now, should_finalize def _mobilede_finalize_refresh_cycle_sold_marking( redis_client: Redis, *, cycle_id: str, logger: logging.Logger, get_persistence: Callable[[], object], origin_prefixes: tuple[str, ...], post_refresh_probe_enabled: bool, post_refresh_probe_batch_size: int, schedule_post_refresh_probe: Callable[[], None], ) -> int: 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) return 0 try: cutoff = datetime.fromisoformat(_mobilede_redis_text(started_at_raw)) if cutoff.tzinfo is None: cutoff = cutoff.replace(tzinfo=timezone.utc) except Exception: logger.warning( "mobile.de refresh sold marking skipped: invalid started_at=%s cycle=%s", started_at_raw, cycle_id, ) return 0 sold_marked = get_persistence().mark_sold_not_seen_since( cutoff, 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, done_now, total_now, cutoff.isoformat(), sold_marked, ) if post_refresh_probe_enabled and post_refresh_probe_batch_size > 0: try: schedule_post_refresh_probe() except Exception: logger.warning( "mobile.de post-refresh sold probe scheduling failed: cycle=%s", cycle_id, exc_info=True, ) return sold_marked def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( redis_client: Redis, *, cycle_id: str | None, total_segments_hint: int = 0, logger: logging.Logger, finalize_refresh_cycle_sold_marking: Callable[[Redis], int] | Callable[..., int], ) -> bool: if not cycle_id: return False active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip() if active_cycle_id != cycle_id: return False done_now, total_now, _left_now = _mobilede_refresh_cycle_progress( redis_client, cycle_id=cycle_id, total_segments_hint=total_segments_hint, ) 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", cycle_id, done_now, total_now, ) 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) return True