Refactor worker
This commit is contained in:
260
mobilede_scraper/worker/refresh_cycle.py
Normal file
260
mobilede_scraper/worker/refresh_cycle.py
Normal file
@@ -0,0 +1,260 @@
|
||||
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
|
||||
Reference in New Issue
Block a user