Build complete search segments
This commit is contained in:
@@ -2,6 +2,8 @@ from __future__ import annotations
|
||||
|
||||
# Планирование сегментов вынесено сюда, чтобы `tasks.py` оставался короче.
|
||||
# Модуль использует хелперы из `tasks.py`.
|
||||
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
|
||||
|
||||
from .tasks import * # noqa: F401,F403
|
||||
from .tasks import (
|
||||
_mobilede_make_segment_url,
|
||||
@@ -18,19 +20,32 @@ from .tasks import (
|
||||
)
|
||||
|
||||
|
||||
_MOBILEDE_RUNTIME_PLAN_VERSION = 3
|
||||
_MOBILEDE_RUNTIME_PLAN_VERSION = 4
|
||||
|
||||
|
||||
def _mobilede_runtime_plan_issues(segments: list[dict[str, object]]) -> dict[str, object]:
|
||||
page_cap = max(MOBILEDE_RESULTS_PER_PAGE, MOBILEDE_MAX_PAGE_NUMBER * MOBILEDE_RESULTS_PER_PAGE)
|
||||
unknown = [item for item in segments if _mobilede_segment_total_results(item) is None]
|
||||
dense = [item for item in segments if (_mobilede_segment_total_results(item) or 0) > page_cap]
|
||||
under_capacity = [
|
||||
item
|
||||
for item in segments
|
||||
if (_mobilede_segment_total_results(item) or 0)
|
||||
> max(1, int(item.get("max_pages") or 0)) * MOBILEDE_RESULTS_PER_PAGE
|
||||
]
|
||||
advertised_total = sum(_mobilede_segment_total_results(item) or 0 for item in segments)
|
||||
fetchable_total = sum(min(_mobilede_segment_total_results(item) or 0, page_cap) for item in segments)
|
||||
fetchable_total = sum(
|
||||
min(
|
||||
_mobilede_segment_total_results(item) or 0,
|
||||
max(1, int(item.get("max_pages") or 0)) * MOBILEDE_RESULTS_PER_PAGE,
|
||||
)
|
||||
for item in segments
|
||||
)
|
||||
return {
|
||||
"page_cap": page_cap,
|
||||
"unknown": unknown,
|
||||
"dense": dense,
|
||||
"under_capacity": under_capacity,
|
||||
"advertised_total": advertised_total,
|
||||
"fetchable_total": fetchable_total,
|
||||
"unreachable_total": max(0, advertised_total - fetchable_total),
|
||||
@@ -41,7 +56,8 @@ def _mobilede_assert_complete_runtime_plan(segments: list[dict[str, object]]) ->
|
||||
issues = _mobilede_runtime_plan_issues(segments)
|
||||
unknown = issues["unknown"]
|
||||
dense = issues["dense"]
|
||||
if not unknown and not dense:
|
||||
under_capacity = issues["under_capacity"]
|
||||
if not unknown and not dense and not under_capacity:
|
||||
return
|
||||
dense_preview = ", ".join(
|
||||
f"{_mobilede_short_segment_label(item)}={_mobilede_segment_total_results(item)}"
|
||||
@@ -50,6 +66,7 @@ def _mobilede_assert_complete_runtime_plan(segments: list[dict[str, object]]) ->
|
||||
raise RuntimeError(
|
||||
"mobile.de runtime plan is incomplete: "
|
||||
f"segments={len(segments)} unknown={len(unknown)} dense={len(dense)} "
|
||||
f"under_capacity={len(under_capacity)} "
|
||||
f"page_cap={issues['page_cap']} unreachable={issues['unreachable_total']} "
|
||||
f"dense_preview=[{dense_preview}]"
|
||||
)
|
||||
@@ -1430,11 +1447,14 @@ def _mobilede_overflow_candidate_group_is_useful(
|
||||
|
||||
target, lower_target, tiny_threshold = _mobilede_target_results_band()
|
||||
useful_floor = max(MOBILEDE_RESULTS_PER_PAGE, int(target * MOBILEDE_OVERFLOW_MIN_USEFUL_CHILD_RATIO))
|
||||
probed_children = children[:MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN]
|
||||
if any(_mobilede_parse_optional_int(child.get("total_results")) is None for child in probed_children):
|
||||
return True
|
||||
child_totals = [
|
||||
int(total)
|
||||
for total in (
|
||||
_mobilede_parse_optional_int(child.get("total_results"))
|
||||
for child in children[:MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN]
|
||||
for child in probed_children
|
||||
)
|
||||
if total is not None
|
||||
]
|
||||
@@ -1968,7 +1988,7 @@ def _mobilede_segment_pages_for_total(total_results: int | None, fallback_max_pa
|
||||
if total_results is None or total_results <= 0:
|
||||
return min(fallback_max_pages, MOBILEDE_MAX_PAGE_NUMBER)
|
||||
pages = max(1, min(MOBILEDE_MAX_PAGE_NUMBER, (int(total_results) + MOBILEDE_RESULTS_PER_PAGE - 1) // MOBILEDE_RESULTS_PER_PAGE))
|
||||
return min(fallback_max_pages, pages)
|
||||
return pages
|
||||
|
||||
|
||||
def _mobilede_make_expanded_segment(
|
||||
|
||||
@@ -82,6 +82,16 @@ def run_mobilede_sync_search_task(
|
||||
if not lock_acquired:
|
||||
logger.info("mobile.de sync skipped: segment already running key=%s", segment_runtime_key)
|
||||
return {"status": "skipped", "reason": "segment_already_running", "segment_key": segment_runtime_key}
|
||||
blocked_for_ms = max(0, int(redis_client.pttl(MOBILEDE_SEARCH_ANTIBOT_BLOCK_KEY) or 0))
|
||||
if blocked_for_ms:
|
||||
retry_in_seconds = max(1, (blocked_for_ms + 999) // 1000)
|
||||
_release_lock_if_owner(redis_client, segment_lock_key, lock_owner)
|
||||
logger.warning(
|
||||
"mobile.de search circuit breaker active: runtime=%s retry_in=%ss",
|
||||
_mobilede_short_segment_label(segment),
|
||||
retry_in_seconds,
|
||||
)
|
||||
raise self.retry(countdown=retry_in_seconds)
|
||||
heartbeat_stop: Event | None = None
|
||||
heartbeat_thread: Thread | None = None
|
||||
watchdog_stop: Event | None = None
|
||||
@@ -1145,6 +1155,15 @@ 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 is_antibot_block:
|
||||
cooldown_seconds = MOBILEDE_ANTIBOT_BACKOFF_SECONDS
|
||||
redis_client.set(MOBILEDE_SEARCH_ANTIBOT_BLOCK_KEY, "1", ex=cooldown_seconds)
|
||||
logger.warning(
|
||||
"mobile.de search circuit breaker opened: status=%s cooldown=%ss runtime=%s",
|
||||
status_code,
|
||||
cooldown_seconds,
|
||||
_mobilede_short_segment_label(segment),
|
||||
)
|
||||
if refresh_cycle_id:
|
||||
_mobilede_mark_refresh_cycle_reconciliation_unsafe(
|
||||
redis_client,
|
||||
|
||||
Reference in New Issue
Block a user