Build complete runtime segments

This commit is contained in:
qananasikq
2026-08-19 17:11:00 +03:00
parent 8297742a64
commit 7068b7e8fa
5 changed files with 1054 additions and 232 deletions
+357 -166
View File
@@ -18,7 +18,41 @@ from .tasks import (
)
_MOBILEDE_RUNTIME_PLAN_VERSION = 2
_MOBILEDE_RUNTIME_PLAN_VERSION = 3
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]
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)
return {
"page_cap": page_cap,
"unknown": unknown,
"dense": dense,
"advertised_total": advertised_total,
"fetchable_total": fetchable_total,
"unreachable_total": max(0, advertised_total - fetchable_total),
}
def _mobilede_assert_complete_runtime_plan(segments: list[dict[str, object]]) -> None:
issues = _mobilede_runtime_plan_issues(segments)
unknown = issues["unknown"]
dense = issues["dense"]
if not unknown and not dense:
return
dense_preview = ", ".join(
f"{_mobilede_short_segment_label(item)}={_mobilede_segment_total_results(item)}"
for item in dense[:5]
)
raise RuntimeError(
"mobile.de runtime plan is incomplete: "
f"segments={len(segments)} unknown={len(unknown)} dense={len(dense)} "
f"page_cap={issues['page_cap']} unreachable={issues['unreachable_total']} "
f"dense_preview=[{dense_preview}]"
)
def _mobilede_price_ranges() -> list[tuple[int, int | None]]:
@@ -78,10 +112,6 @@ def _mobilede_price_ranges_for_segment(segment: dict[str, object] | None = None)
return _mobilede_price_ranges()
def _mobilede_filtered_url_uses_adaptive_plan() -> bool:
return os.getenv("MOBILEDE_FILTERED_URL_ADAPTIVE_PLAN", "false").strip().lower() in {"1", "true", "yes", "on"}
def _mobilede_skip_late_overflow_children_during_bootstrap() -> bool:
return os.getenv("MOBILEDE_SKIP_LATE_OVERFLOW_CHILDREN_DURING_BOOTSTRAP", "true").strip().lower() in {"1", "true", "yes", "on"}
@@ -294,7 +324,9 @@ def _mobilede_overflow_threshold(max_pages: int) -> int:
def _mobilede_parse_optional_int(value: object) -> int | None:
raw = str(value or "").strip()
if value is None:
return None
raw = str(value).strip()
if not raw:
return None
try:
@@ -338,6 +370,178 @@ def _mobilede_prune_overflow_parent_segments(segments: list[dict[str, object]])
return pruned
_MOBILEDE_COMPACT_RANGE_FIELDS = (
("price_min", "price_max", "p", _mobilede_price_label),
("year_min", "year_max", "fr", _mobilede_year_label),
("mileage_min", "mileage_max", "ml", _mobilede_mileage_label),
)
def _mobilede_runtime_segment_base_url(segment: dict[str, object]) -> str:
search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip()
return _mobilede_strip_query_keys(search_url, {"p", "fr", "ml", "pageNumber", "lang"})
def _mobilede_ranges_are_contiguous(
left: tuple[int | None, int | None],
right: tuple[int | None, int | None],
) -> tuple[int | None, int | None] | None:
left_min, left_max = left
right_min, right_max = right
if left_max is not None and right_min is not None and int(left_max) + 1 == int(right_min):
return left_min, right_max
if right_max is not None and left_min is not None and int(right_max) + 1 == int(left_min):
return right_min, left_max
return None
def _mobilede_compact_segment_pair(
left: dict[str, object],
right: dict[str, object],
*,
target_results: int,
) -> dict[str, object] | None:
left_total = _mobilede_segment_total_results(left)
right_total = _mobilede_segment_total_results(right)
if left_total is None or right_total is None or left_total < 0 or right_total < 0:
return None
combined_total = int(left_total) + int(right_total)
if combined_total > int(target_results):
return None
if _mobilede_runtime_segment_base_url(left) != _mobilede_runtime_segment_base_url(right):
return None
for field in ("make_id", "model_id", "source_search_url", "runtime_brand"):
if str(left.get(field) or "").strip() != str(right.get(field) or "").strip():
return None
differing: list[tuple[str, str, str, object, tuple[int | None, int | None]]] = []
for min_field, max_field, query_key, labeler in _MOBILEDE_COMPACT_RANGE_FIELDS:
left_range = (
_mobilede_parse_optional_int(left.get(min_field)),
_mobilede_parse_optional_int(left.get(max_field)),
)
right_range = (
_mobilede_parse_optional_int(right.get(min_field)),
_mobilede_parse_optional_int(right.get(max_field)),
)
if left_range == right_range:
continue
merged_range = _mobilede_ranges_are_contiguous(left_range, right_range)
if merged_range is None:
return None
differing.append((min_field, max_field, query_key, labeler, merged_range))
if len(differing) != 1:
return None
min_field, max_field, query_key, labeler, merged_range = differing[0]
merged_min, merged_max = merged_range
merged = dict(left)
merged[min_field] = str(merged_min) if merged_min is not None else None
merged[max_field] = str(merged_max) if merged_max is not None else None
range_value = _mobilede_range_value(merged_min, merged_max)
merged_url = _mobilede_make_segment_url(
str(left.get("search_url") or left.get("listing_url") or ""),
**{query_key: range_value},
)
merged["search_url"] = merged_url
merged["listing_url"] = merged_url
merged["start_page"] = 1
merged["max_pages"] = _mobilede_segment_pages_for_total(
combined_total,
max(
int(left.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER),
int(right.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER),
),
)
merged["total_results"] = combined_total
merged["observed_complete"] = True
merged["compacted_from"] = int(left.get("compacted_from") or 1) + int(right.get("compacted_from") or 1)
merged.pop("overflow_parent", None)
merged.pop("overflow_split", None)
merged["overflow_depth"] = min(
int(_mobilede_parse_optional_int(left.get("overflow_depth")) or 0),
int(_mobilede_parse_optional_int(right.get("overflow_depth")) or 0),
)
base_label = str(left.get("label") or right.get("label") or "mobile.de segment")
label_parts = [part.strip() for part in base_label.split(" | ") if part.strip()]
label_parts = [
part
for part in label_parts
if not part.startswith(("price=", "year", "km", "total=", "overflow:d"))
]
range_labels: list[str] = []
for range_min_field, range_max_field, _range_query_key, range_labeler in _MOBILEDE_COMPACT_RANGE_FIELDS:
range_min = _mobilede_parse_optional_int(merged.get(range_min_field))
range_max = _mobilede_parse_optional_int(merged.get(range_max_field))
if range_min is None and range_max is None:
continue
range_labels.append(range_labeler(range_min, range_max))
merged["label"] = " | ".join([*label_parts, *range_labels, f"total={combined_total}"])
return merged
def _mobilede_compact_runtime_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]:
target_results = min(
int(MOBILEDE_SEGMENT_TARGET_RESULTS),
_mobilede_overflow_threshold(MOBILEDE_MAX_PAGE_NUMBER),
)
compacted = [
dict(item)
for item in _mobilede_prune_overflow_parent_segments(segments)
if isinstance(item, dict)
and not (
item.get("observed_complete") is True
and _mobilede_segment_total_results(item) == 0
)
]
before = len(compacted)
changed = True
while changed:
changed = False
for left_index in range(len(compacted)):
for right_index in range(left_index + 1, len(compacted)):
merged = _mobilede_compact_segment_pair(
compacted[left_index],
compacted[right_index],
target_results=target_results,
)
if merged is None:
continue
compacted[left_index] = merged
compacted.pop(right_index)
changed = True
break
if changed:
break
if len(compacted) != before:
logger.info(
"mobile.de runtime segments compacted: before=%s after=%s removed=%s target=%s",
before,
len(compacted),
before - len(compacted),
target_results,
)
return compacted
def _mobilede_finalize_runtime_plan(segments: list[dict[str, object]]) -> list[dict[str, object]]:
"""Keep dispatched leaves compact without merging beyond the page cap."""
compacted = _mobilede_compact_runtime_segments(segments)
if len(compacted) != len(segments):
logger.info(
"mobile.de runtime plan compacted before dispatch: before=%s after=%s",
len(segments),
len(compacted),
)
return _mobilede_interleave_segments_by_make(compacted)
def _mobilede_finalize_complete_runtime_plan(segments: list[dict[str, object]]) -> list[dict[str, object]]:
final_plan = _mobilede_finalize_runtime_plan(segments)
_mobilede_assert_complete_runtime_plan(final_plan)
return final_plan
def _mobilede_learned_segments_file() -> str:
return os.getenv("MOBILEDE_LEARNED_SEGMENTS_FILE", "/data/mobilede_runtime_segments.json").strip()
@@ -377,6 +581,19 @@ def _mobilede_load_learned_runtime_segments(source_segments: list[dict[str, obje
learned = _mobilede_prune_overflow_parent_segments([dict(item) for item in raw_segments if isinstance(item, dict)])
if not learned:
return None
if payload.get("complete") is not True:
logger.info("mobile.de learned segments ignored: file=%s reason=not_marked_complete", path)
return None
issues = _mobilede_runtime_plan_issues(learned)
if issues["unknown"] or issues["dense"]:
logger.warning(
"mobile.de learned segments ignored: file=%s unknown=%s dense=%s unreachable=%s",
path,
len(issues["unknown"]),
len(issues["dense"]),
issues["unreachable_total"],
)
return None
logger.info("mobile.de learned runtime segments loaded: file=%s count=%s", path, len(learned))
return learned
except FileNotFoundError:
@@ -395,11 +612,21 @@ def _mobilede_save_learned_runtime_segments(
return
try:
os.makedirs(os.path.dirname(path) or ".", exist_ok=True)
pruned_segments = _mobilede_prune_overflow_parent_segments(runtime_segments)
pruned_segments = _mobilede_compact_runtime_segments(runtime_segments)
plan_issues = _mobilede_runtime_plan_issues(pruned_segments)
if plan_issues["unknown"] or plan_issues["dense"]:
logger.warning(
"mobile.de incomplete learned plan not saved: unknown=%s dense=%s unreachable=%s",
len(plan_issues["unknown"]),
len(plan_issues["dense"]),
plan_issues["unreachable_total"],
)
return
payload = {
"plan_version": _MOBILEDE_RUNTIME_PLAN_VERSION,
"source_fingerprint": _mobilede_segments_source_fingerprint(source_segments),
"updated_at": datetime.now(timezone.utc).isoformat(),
"complete": True,
"segments": pruned_segments,
}
tmp_path = f"{path}.tmp"
@@ -411,6 +638,80 @@ def _mobilede_save_learned_runtime_segments(
logger.warning("mobile.de learned runtime segments save failed: file=%s", path, exc_info=True)
def _mobilede_record_observed_segment_total(
redis_client: Redis,
settings: Settings,
*,
segment: dict[str, object] | None,
total_results: int | None,
listing_count: int,
unique_count: int,
bootstrap_run: bool,
) -> int:
if not segment:
return 0
observed_total = (
max(0, int(total_results))
if total_results is not None
else max(0, int(listing_count or 0), int(unique_count or 0))
)
segment_fingerprint = _mobilede_segment_fingerprint(segment)
lock_owner = f"observe:{uuid.uuid4().hex}"
if not _acquire_lock(redis_client, MOBILEDE_RUNTIME_SEGMENTS_CACHE_LOCK_KEY, lock_owner, 120):
return 0
try:
cached_segments = _get_cached_mobilede_runtime_segments(redis_client)
if cached_segments is None:
cached_segments = _build_mobilede_runtime_segments(settings)
updated_segments: list[dict[str, object]] = []
matched = False
for item in cached_segments:
if not isinstance(item, dict):
continue
updated = dict(item)
if _mobilede_segment_fingerprint(updated) == segment_fingerprint:
updated["total_results"] = observed_total
updated["observed_complete"] = True
updated["max_pages"] = _mobilede_segment_pages_for_total(
observed_total,
int(updated.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER),
)
matched = True
updated_segments.append(updated)
if not matched:
return 0
learned_segments = _mobilede_compact_runtime_segments(updated_segments)
cache_segments = updated_segments
redis_client.set(
MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY,
json.dumps(cache_segments, ensure_ascii=False),
ex=24 * 60 * 60,
)
_mobilede_save_learned_runtime_segments(
_mobilede_source_segments_from_settings(settings),
learned_segments,
)
logger.info(
"mobile.de complete segment observation saved: segment=%s total=%s runtime_segments=%s learned_segments=%s bootstrap=%s",
_mobilede_short_segment_label(segment),
observed_total,
len(cache_segments),
len(learned_segments),
bootstrap_run,
)
return len(updated_segments) - len(learned_segments)
except Exception:
logger.warning(
"Failed to save mobile.de complete segment observation: segment=%s",
_mobilede_segment_label(segment),
exc_info=True,
)
return 0
finally:
_release_lock_if_owner(redis_client, MOBILEDE_RUNTIME_SEGMENTS_CACHE_LOCK_KEY, lock_owner)
def _mobilede_normalize_make_name(value: str) -> str:
text = unicodedata.normalize("NFKD", str(value or ""))
ascii_text = text.encode("ascii", "ignore").decode("ascii")
@@ -1400,6 +1701,7 @@ def _mobilede_try_expand_overflow_segment(
unique_count: int,
max_pages: int,
segment_end_page: int,
bootstrap_run: bool = False,
) -> int:
if not MOBILEDE_OVERFLOW_SPLIT_ENABLED or not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment:
return 0
@@ -1485,9 +1787,24 @@ def _mobilede_try_expand_overflow_segment(
committed = True
return 0
updated_segments = _mobilede_prune_overflow_parent_segments(
updated_segments = _mobilede_compact_runtime_segments(
[dict(item) for item in cached_segments if isinstance(item, dict)] + appended_segments
)
if bootstrap_run and _mobilede_skip_late_overflow_children_during_bootstrap():
_mobilede_save_learned_runtime_segments(
_mobilede_source_segments_from_settings(settings),
updated_segments,
)
committed = True
logger.info(
"mobile.de overflow split deferred until next pass: parent=%s parent_key=%s observed=%s added=%s learned_segments=%s",
_mobilede_short_segment_label(segment),
_mobilede_short_segment_ref(segment),
observed,
len(appended_segments),
len(updated_segments),
)
return 0
redis_client.set(
MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY,
json.dumps(updated_segments, ensure_ascii=False),
@@ -1638,7 +1955,7 @@ def _mobilede_probe_total(search_url: str, **params: str | int | None) -> int |
page_number=1,
search_url=search_url,
timeout=probe_timeout,
max_retries=0,
max_retries=2,
**params,
)
return int(page.total_results or 0)
@@ -1773,7 +2090,7 @@ def _mobilede_split_search_url_segment_by_make(segment: dict[str, object]) -> li
def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) -> list[dict[str, object]] | None:
if not MOBILEDE_PREPLAN_SEGMENT_PROBES or not _mobilede_filtered_url_uses_adaptive_plan():
if not MOBILEDE_PREPLAN_SEGMENT_PROBES:
return None
search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip()
@@ -1784,34 +2101,22 @@ def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) ->
max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER)
planned: list[dict[str, object]] = []
probes_used = 0
probe_limit = max(1, int(MOBILEDE_PREPLAN_MAX_PROBES))
started_at = time.monotonic()
max_seconds = int(MOBILEDE_ADAPTIVE_URL_MAX_SECONDS)
def _probe(params: dict[str, str | int | None]) -> int | None:
nonlocal probes_used
if probes_used >= probe_limit:
return None
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
return None
probes_used += 1
if probes_used == 1 or probes_used % 25 == 0:
_mobilede_touch_planning_progress("adaptive_url_planning")
if probes_used % 50 == 0:
logger.info(
"mobile.de adaptive planning progress: base=%s probes=%s/%s planned=%s",
"mobile.de adaptive planning progress: base=%s probes=%s planned=%s",
_mobilede_short_segment_label(segment),
probes_used,
probe_limit,
len(planned),
)
return _mobilede_probe_total(search_url, **params)
time_budget_reached = False
for price_min, price_max in _mobilede_price_ranges_for_segment(segment):
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
time_budget_reached = True
break
price_params: dict[str, str | int | None] = {"p": _mobilede_range_value(price_min, price_max)}
price_label = _mobilede_price_label(price_min, price_max)
price_total = _probe(price_params)
@@ -1833,14 +2138,8 @@ def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) ->
continue
for year_min, year_max in _mobilede_year_ranges_for_segment_price(segment, price_min, price_max):
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
time_budget_reached = True
break
year_label = _mobilede_year_label(year_min, year_max)
for refined_price_min, refined_price_max in _mobilede_price_subranges_for_hot_year(price_min, price_max, year_min, year_max):
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
time_budget_reached = True
break
refined_price_label = _mobilede_price_label(refined_price_min, refined_price_max)
year_params = {
"p": _mobilede_range_value(refined_price_min, refined_price_max),
@@ -1866,9 +2165,6 @@ def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) ->
continue
for mileage_min, mileage_max in _mobilede_mileage_ranges():
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
time_budget_reached = True
break
mileage_params = dict(year_params)
mileage_params["ml"] = _mobilede_range_value(mileage_min, mileage_max)
mileage_label = _mobilede_mileage_label(mileage_min, mileage_max)
@@ -1889,27 +2185,11 @@ def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) ->
mileage_range=(mileage_min, mileage_max),
)
)
if time_budget_reached:
break
if time_budget_reached:
break
if time_budget_reached:
logger.warning(
"mobile.de adaptive URL planning time budget reached: base=%s seconds=%s segments=%s probes=%s/%s",
_mobilede_short_segment_label(segment),
max_seconds,
len(planned),
probes_used,
probe_limit,
)
logger.info(
"mobile.de adaptive URL segments planned: base=%s segments=%s probes=%s limit=%s",
"mobile.de adaptive URL segments planned: base=%s segments=%s probes=%s",
_mobilede_short_segment_label(segment),
len(planned),
probes_used,
probe_limit,
)
# Финальную dense-доразбивку делает _build_mobilede_runtime_segments.
# Здесь только возвращаем probe-план с total_results, чтобы не терять хвосты
@@ -1922,27 +2202,12 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
pending: list[dict[str, object]] = [dict(item) for item in segments]
refined: list[dict[str, object]] = []
probes_used = 0
probe_limit = max(1, int(MOBILEDE_PREPLAN_MAX_PROBES))
max_refine_depth = max(MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH, MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH)
started_at = time.monotonic()
max_seconds = int(MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS)
while pending:
if max_seconds > 0 and time.monotonic() - started_at >= max_seconds:
logger.warning(
"mobile.de dense refine time budget reached: seconds=%s refined=%s pending=%s probes=%s/%s",
max_seconds,
len(refined),
len(pending),
probes_used,
probe_limit,
)
refined.extend(_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in pending)
break
current = pending.pop(0)
total_results = current.get("total_results")
if total_results is None and probes_used < probe_limit:
if total_results is None:
total_results = _mobilede_probe_segment_total(current)
current["total_results"] = total_results
probes_used += 1
@@ -1950,9 +2215,8 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
_mobilede_touch_planning_progress("dense_refine")
if probes_used % 50 == 0:
logger.info(
"mobile.de dense refine progress: probes=%s/%s refined=%s pending=%s",
"mobile.de dense refine progress: probes=%s refined=%s pending=%s",
probes_used,
probe_limit,
len(refined),
len(pending),
)
@@ -1962,7 +2226,6 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
total_results is not None
and int(total_results) > threshold
and depth < max_refine_depth
and probes_used < probe_limit
and len(refined) + len(pending) < MOBILEDE_PREPLAN_MAX_SEGMENTS
):
children = _mobilede_build_overflow_child_segments(
@@ -1972,7 +2235,7 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
)
if children:
for child in children:
if probes_used >= probe_limit or len(refined) + len(pending) >= MOBILEDE_PREPLAN_MAX_SEGMENTS:
if len(refined) + len(pending) >= MOBILEDE_PREPLAN_MAX_SEGMENTS:
pending.append(child)
continue
child_total = _mobilede_probe_segment_total(child)
@@ -1982,9 +2245,8 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
_mobilede_touch_planning_progress("dense_refine")
if probes_used % 50 == 0:
logger.info(
"mobile.de dense refine progress: probes=%s/%s refined=%s pending=%s",
"mobile.de dense refine progress: probes=%s refined=%s pending=%s",
probes_used,
probe_limit,
len(refined),
len(pending),
)
@@ -1993,15 +2255,16 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
pending.append(_mobilede_finalize_preplanned_segment(child, child_total))
continue
if total_results is None:
raise RuntimeError(
"mobile.de complete plan could not resolve segment total: "
f"segment={_mobilede_short_segment_label(current)} probes={probes_used}"
)
if total_results is not None and int(total_results) > threshold:
logger.warning(
"mobile.de dense segment remains after refine: segment=%s total=%s depth=%s max_depth=%s probes=%s/%s",
_mobilede_short_segment_label(current),
total_results,
depth,
max_refine_depth,
probes_used,
probe_limit,
raise RuntimeError(
"mobile.de complete plan could not split dense segment: "
f"segment={_mobilede_short_segment_label(current)} total={total_results} "
f"depth={depth}/{max_refine_depth} probes={probes_used}"
)
refined.append(_mobilede_finalize_preplanned_segment(current, int(total_results) if total_results is not None else None))
@@ -2016,37 +2279,6 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -
return refined
def _mobilede_refine_dense_planned_segments_until_stable(segments: list[dict[str, object]]) -> list[dict[str, object]]:
max_passes = max(1, int(os.getenv("MOBILEDE_REFINE_ADAPTIVE_MAX_PASSES", "2")))
refined = [dict(item) for item in segments]
for pass_index in range(1, max_passes + 1):
before = len(refined)
refined = _mobilede_refine_dense_planned_segments(refined)
dense_count = sum(
1
for item in refined
if (_mobilede_segment_total_results(item) or 0) > min(
MOBILEDE_SEGMENT_TARGET_RESULTS,
_mobilede_overflow_threshold(MOBILEDE_MAX_PAGE_NUMBER),
)
)
logger.info(
"mobile.de dense refine pass complete: pass=%s/%s before=%s after=%s dense_left=%s",
pass_index,
max_passes,
before,
len(refined),
dense_count,
)
if dense_count <= 0 or len(refined) >= MOBILEDE_PREPLAN_MAX_SEGMENTS:
break
return refined
def _mobilede_should_refine_adaptive_segments() -> bool:
return os.getenv("MOBILEDE_REFINE_ADAPTIVE_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"}
def _mobilede_try_keep_root_segment_unsplit(
segment: dict[str, object],
*,
@@ -2087,12 +2319,7 @@ def _mobilede_try_keep_root_segment_unsplit(
return [item]
def _expand_mobilede_search_url_segment(
segment: dict[str, object],
*,
allow_adaptive_planning: bool = True,
allow_probe_planning: bool = True,
) -> list[dict[str, object]]:
def _expand_mobilede_search_url_segment(segment: dict[str, object]) -> list[dict[str, object]]:
search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip()
if not search_url:
return [segment]
@@ -2101,13 +2328,7 @@ def _expand_mobilede_search_url_segment(
if len(split_by_make) > 1:
expanded: list[dict[str, object]] = []
for split_segment in split_by_make:
expanded.extend(
_expand_mobilede_search_url_segment(
split_segment,
allow_adaptive_planning=allow_adaptive_planning,
allow_probe_planning=allow_probe_planning,
)
)
expanded.extend(_expand_mobilede_search_url_segment(split_segment))
return expanded
# Готовый URL остаётся главным источником фильтров. Не режем его повторно
@@ -2119,13 +2340,6 @@ def _expand_mobilede_search_url_segment(
if query_keys & existing_range_keys:
return [segment]
if not allow_adaptive_planning:
logger.info(
"mobile.de URL segment fast-start enabled: using search_url as-is for %s",
_mobilede_short_segment_label(segment),
)
return [segment]
max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER)
root_segment = _mobilede_try_keep_root_segment_unsplit(
segment,
@@ -2135,10 +2349,9 @@ def _expand_mobilede_search_url_segment(
if root_segment is not None:
return root_segment
if allow_probe_planning:
probe_planned = _expand_mobilede_search_url_segment_by_probe(segment)
if probe_planned is not None:
return probe_planned
probe_planned = _expand_mobilede_search_url_segment_by_probe(segment)
if probe_planned is not None:
return probe_planned
expanded: list[dict[str, object]] = []
base_label = str(segment.get("label") or "mobile.de segmented URL")
@@ -2248,15 +2461,6 @@ def _expand_mobilede_search_url_segment(
def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, object]]:
env_search_urls = settings.listing.filtered_search_urls
# Готовые фильтрованные URL по умолчанию используют старое
# детерминированное деление по цене/году: оно стартует сразу и его проще
# держать стабильным. Адаптивное планирование на основе probe можно
# включить явно, когда нужен более плотный pre-plan.
fast_start_filtered_urls = (
bool(env_search_urls)
and os.getenv("MOBILEDE_FILTERED_URL_FAST_START", "false").strip().lower() in {"1", "true", "yes", "on"}
)
adaptive_filtered_urls = bool(env_search_urls) and _mobilede_filtered_url_uses_adaptive_plan()
if env_search_urls:
segments = _mobilede_source_segments_from_settings(settings)
if any(str(item.get("runtime_brand") or "").strip() for item in segments):
@@ -2278,36 +2482,23 @@ def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, objec
if segment.get("auto_segment") is False:
expanded.append(segment)
continue
segment_expanded = _expand_mobilede_search_url_segment(
segment,
allow_adaptive_planning=not fast_start_filtered_urls,
allow_probe_planning=not env_search_urls or adaptive_filtered_urls,
)
segment_expanded = _expand_mobilede_search_url_segment(segment)
expanded.extend(segment_expanded)
if any(item.get("total_results") is not None for item in segment_expanded):
expanded_from_adaptive_url = True
if len(expanded) != len(segments):
logger.info("mobile.de segments planned: input=%s total=%s", len(segments), len(expanded))
expanded = _mobilede_interleave_segments_by_make(expanded)
if fast_start_filtered_urls:
logger.info(
"mobile.de fast-start runtime segments ready: input_urls=%s final=%s reason=filtered_search_urls",
len(env_search_urls),
len(expanded),
)
return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in expanded]
if expanded_from_adaptive_url:
if _mobilede_should_refine_adaptive_segments():
refined = _mobilede_refine_dense_planned_segments_until_stable(expanded)
logger.info(
"mobile.de adaptive URL plan refined: before=%s after=%s reason=dense_segments",
len(expanded),
len(refined),
)
return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in refined]
refined = _mobilede_refine_dense_planned_segments(expanded)
logger.info(
"mobile.de preplan skipped after adaptive URL planning: final=%s reason=fast_start",
"mobile.de adaptive URL plan refined: before=%s after=%s reason=dense_segments",
len(expanded),
len(refined),
)
return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in expanded]
return _mobilede_preplan_runtime_segments(expanded)
final_plan = _mobilede_finalize_complete_runtime_plan(
[_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in refined]
)
_mobilede_save_learned_runtime_segments(segments, final_plan)
return final_plan
return _mobilede_finalize_complete_runtime_plan(_mobilede_preplan_runtime_segments(expanded))
+31 -9
View File
@@ -69,17 +69,28 @@ def _mobilede_redis_text(value: object) -> str:
def _mobilede_claim_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> bool:
ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS))
lease_seconds = min(300, max(60, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)))
return bool(
redis_client.set(
_mobilede_refresh_cycle_finalized_key(cycle_id),
"1",
"in_progress",
nx=True,
ex=ttl,
ex=lease_seconds,
)
)
def _mobilede_release_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> None:
key = _mobilede_refresh_cycle_finalized_key(cycle_id)
if _mobilede_redis_text(redis_client.get(key)) == "in_progress":
redis_client.delete(key)
def _mobilede_complete_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> None:
ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS))
redis_client.set(_mobilede_refresh_cycle_finalized_key(cycle_id), "completed", 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 = [
@@ -217,6 +228,7 @@ def _mobilede_finalize_refresh_cycle_sold_marking(
post_refresh_probe_enabled: bool,
post_refresh_probe_batch_size: int,
schedule_post_refresh_probe: Callable[[], None],
scope_brands: tuple[str, ...] = (),
) -> 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)
@@ -231,10 +243,12 @@ def _mobilede_finalize_refresh_cycle_sold_marking(
total_now,
sorted(unsafe_reasons),
)
_mobilede_release_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
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)
_mobilede_release_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
return 0
try:
cutoff = datetime.fromisoformat(_mobilede_redis_text(started_at_raw))
@@ -246,18 +260,26 @@ def _mobilede_finalize_refresh_cycle_sold_marking(
started_at_raw,
cycle_id,
)
_mobilede_release_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
return 0
sold_marked = get_persistence().mark_sold_not_seen_since(
cutoff,
prefix=origin_prefixes,
safety_ratio=0.8,
)
try:
sold_marked = get_persistence().mark_sold_not_seen_since(
cutoff,
prefix=origin_prefixes,
safety_ratio=0.8,
brands=scope_brands,
)
except Exception:
_mobilede_release_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
raise
_mobilede_complete_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
logger.info(
"mobile.de refresh sold marking completed: cycle=%s progress=%s/%s cutoff=%s sold_marked=%s",
"mobile.de refresh sold marking completed: cycle=%s progress=%s/%s cutoff=%s brands=%s sold_marked=%s",
cycle_id,
done_now,
total_now,
cutoff.isoformat(),
list(scope_brands) if scope_brands else "all",
sold_marked,
)
if post_refresh_probe_enabled and post_refresh_probe_batch_size > 0:
+24 -5
View File
@@ -111,7 +111,12 @@ def _reset_bootstrap_checkpoint_for_db_idle(redis_client: Redis) -> None:
pipe.execute()
def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict[str, int]]:
def _has_inflight_work(
redis_client: Redis,
queue_name: str,
*,
progress_max_age_seconds: int = 900,
) -> tuple[bool, dict[str, int]]:
"""Есть ли признаки активной/зависшей работы, даже если очередь пуста."""
queue_len = _safe_int(redis_client.llen(queue_name), 0)
has_segment_lock = 0
@@ -119,9 +124,19 @@ def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict
has_segment_lock = 1
break
has_task_progress = 0
for _ in redis_client.scan_iter(match="mobilede:state:task_progress:*"):
has_task_progress = 1
break
now_ts = int(time.time())
max_age = max(60, int(progress_max_age_seconds))
for key in redis_client.scan_iter(match="mobilede:state:task_progress:*"):
try:
payload = redis_client.get(key)
data = json.loads(payload) if payload else {}
progress_ts = _safe_int(data.get("ts"), 0)
if progress_ts > 0 and now_ts - progress_ts <= max_age:
has_task_progress = 1
break
redis_client.delete(key)
except Exception:
redis_client.delete(key)
flags = {
"queue_len": queue_len,
"has_segment_lock": has_segment_lock,
@@ -193,7 +208,11 @@ def main() -> None:
redis_client = _get_redis()
redis_client.ping()
has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
has_inflight, inflight = _has_inflight_work(
redis_client,
queue_name,
progress_max_age_seconds=stall_seconds,
)
if not has_inflight:
time.sleep(check_interval)
continue