diff --git a/mobilede_scraper/worker/planner.py b/mobilede_scraper/worker/planner.py index bef2ea5..fe0599a 100644 --- a/mobilede_scraper/worker/planner.py +++ b/mobilede_scraper/worker/planner.py @@ -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)) diff --git a/mobilede_scraper/worker/refresh_cycle.py b/mobilede_scraper/worker/refresh_cycle.py index 9e27519..d7ee724 100644 --- a/mobilede_scraper/worker/refresh_cycle.py +++ b/mobilede_scraper/worker/refresh_cycle.py @@ -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: diff --git a/mobilede_scraper/worker/self_heal.py b/mobilede_scraper/worker/self_heal.py index 139b1a8..2113f8e 100644 --- a/mobilede_scraper/worker/self_heal.py +++ b/mobilede_scraper/worker/self_heal.py @@ -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 diff --git a/tests/test_self_heal.py b/tests/test_self_heal.py index 3c073ed..772920f 100644 --- a/tests/test_self_heal.py +++ b/tests/test_self_heal.py @@ -60,6 +60,46 @@ class TestSelfHeal(unittest.TestCase): self.assertTrue(has_inflight) self.assertEqual(flags["has_segment_lock"], 1) + @patch("mobilede_scraper.worker.self_heal.time.time", return_value=2_000) + def test_has_inflight_work_ignores_and_deletes_stale_progress(self, _time_mock) -> None: + redis_client = MagicMock() + redis_client.llen.return_value = 0 + redis_client.scan_iter.side_effect = [ + iter([]), + iter(["mobilede:state:task_progress:stale"]), + ] + redis_client.get.return_value = '{"stage":"runtime_segments_planned","ts":1000}' + + has_inflight, flags = self_heal._has_inflight_work( + redis_client, + "mobilede_sync", + progress_max_age_seconds=900, + ) + + self.assertFalse(has_inflight) + self.assertEqual(flags["has_task_progress"], 0) + redis_client.delete.assert_called_once_with("mobilede:state:task_progress:stale") + + @patch("mobilede_scraper.worker.self_heal.time.time", return_value=2_000) + def test_has_inflight_work_accepts_fresh_progress(self, _time_mock) -> None: + redis_client = MagicMock() + redis_client.llen.return_value = 0 + redis_client.scan_iter.side_effect = [ + iter([]), + iter(["mobilede:state:task_progress:fresh"]), + ] + redis_client.get.return_value = '{"stage":"search_page","ts":1500}' + + has_inflight, flags = self_heal._has_inflight_work( + redis_client, + "mobilede_sync", + progress_max_age_seconds=900, + ) + + self.assertTrue(has_inflight) + self.assertEqual(flags["has_task_progress"], 1) + redis_client.delete.assert_not_called() + @patch("mobilede_scraper.worker.self_heal.time.sleep", return_value=None) @patch("mobilede_scraper.worker.self_heal.os.kill") @patch("builtins.open") diff --git a/tests/test_worker_runtime_tasks.py b/tests/test_worker_runtime_tasks.py index 78088fa..6348ddd 100644 --- a/tests/test_worker_runtime_tasks.py +++ b/tests/test_worker_runtime_tasks.py @@ -1,14 +1,40 @@ from __future__ import annotations from datetime import datetime, timezone +import json from types import SimpleNamespace import unittest from unittest.mock import MagicMock, patch -from mobilede_scraper.worker import planner, tasks +from mobilede_scraper.core.runtime_config import RuntimeConfig +from mobilede_scraper.worker import planner, refresh_cycle, search_sync, tasks class TestWorkerRuntimeTaskHelpers(unittest.TestCase): + class _FakeDetailPersistence: + def __init__(self) -> None: + self.enqueued: list[str] = [] + self.released: list[tuple[str, str | None]] = [] + + def enqueue_detail_jobs(self, origin_ids, *, priority=3, reset=False): # noqa: ARG002 + self.enqueued = list(origin_ids) + return {"ready": len(self.enqueued), "deferred": 0, "missing": 0} + + def reserve_detail_jobs(self, *, limit, lease_seconds=900, listing_ids=None): # noqa: ARG002 + return [ + { + "listing_id": listing_id, + "lease_owner": f"lease-{listing_id}", + "attempts": 1, + "priority": 9, + } + for listing_id in list(listing_ids or [])[:limit] + ] + + def release_detail_job(self, listing_id, *, lease_owner=None): + self.released.append((listing_id, lease_owner)) + return True + class _FakeRedis: def __init__(self) -> None: self.store: dict[str, str] = {} @@ -36,6 +62,13 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): bucket.add(member) return 1 if len(bucket) > before else 0 + def srem(self, key: str, member: str): + bucket = self.sets.setdefault(key, set()) + if member not in bucket: + return 0 + bucket.remove(member) + return 1 + def scard(self, key: str): return len(self.sets.get(key, set())) @@ -78,6 +111,33 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(args[2], "lock:key") self.assertEqual(args[3], "owner-token") + def test_refresh_finalization_claim_is_released_after_db_failure(self) -> None: + redis_client = self._FakeRedis() + cycle_id = "cycle-1" + redis_client.set(tasks.MOBILEDE_REFRESH_CYCLE_ID_KEY, cycle_id) + redis_client.set(tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY, "1") + redis_client.set(tasks.MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, "1") + redis_client.set(tasks.MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, datetime.now(timezone.utc).isoformat()) + self.assertTrue(refresh_cycle._mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)) + persistence = MagicMock() + persistence.mark_sold_not_seen_since.side_effect = RuntimeError("db unavailable") + + with self.assertRaisesRegex(RuntimeError, "db unavailable"): + refresh_cycle._mobilede_finalize_refresh_cycle_sold_marking( + redis_client, + cycle_id=cycle_id, + logger=MagicMock(), + get_persistence=lambda: persistence, + origin_prefixes=("mobile.de:",), + post_refresh_probe_enabled=False, + post_refresh_probe_batch_size=0, + schedule_post_refresh_probe=lambda: None, + ) + + claim_key = refresh_cycle._mobilede_refresh_cycle_finalized_key(cycle_id) + self.assertIsNone(redis_client.get(claim_key)) + self.assertTrue(refresh_cycle._mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)) + def test_mobilede_segment_key_is_stable(self) -> None: segment = {"make_id": "BMW", "price_min": 1000, "price_max": 5000} @@ -109,6 +169,178 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(key1, key2) self.assertEqual(len(key1), 16) + def test_runtime_numeric_filters_intersect_planner_segment(self) -> None: + self.assertEqual( + search_sync._intersect_numeric_bounds("10000", "30000", 15000, 25000), + ("15000", "25000", True), + ) + self.assertEqual( + search_sync._intersect_numeric_bounds("10000", "12000", 15000, 25000), + ("10000", "12000", False), + ) + + def test_depth_six_overflow_leaf_can_split_further(self) -> None: + segment = { + "search_url": "https://www.mobile.de/ru/транспортные-средства/поиск.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ml=%3A4687&p=15003%3A", + "listing_url": "https://www.mobile.de/ru/транспортные-средства/поиск.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ml=%3A4687&p=15003%3A", + "label": "BMW hot leaf", + "make_id": "3500", + "mileage_max": "4687", + "price_min": "15003", + "overflow_depth": 6, + "overflow_split": "mileage", + "max_pages": 50, + } + + with patch.object(planner, "MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH", 8): + groups = planner._mobilede_build_overflow_candidate_groups(segment=segment, max_pages=50) + + self.assertTrue(groups) + self.assertTrue(all(int(child["overflow_depth"]) == 7 for _kind, children in groups for child in children)) + self.assertTrue(any(kind == "price" for kind, _children in groups)) + + def test_detail_queue_prioritizes_new_unique_listings(self) -> None: + redis_client = self._FakeRedis() + persistence = self._FakeDetailPersistence() + + with ( + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true", "MOBILEDE_DETAIL_QUEUE_MAX_PENDING": "3"}), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, + ): + result = tasks._queue_mobilede_detail_listing_ids( + redis_client, + ["mobile.de:100", "mobile.de:100", "mobile.de:101"], + priority=9, + ) + + self.assertEqual(result, {"queued": 2, "duplicate": 0, "full": 0}) + self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"100", "101"}) + self.assertEqual(apply_async.call_count, 2) + self.assertTrue(all(call.kwargs["queue"] == tasks.MOBILEDE_IMAGES_QUEUE for call in apply_async.call_args_list)) + self.assertTrue(all(call.kwargs["priority"] == 9 for call in apply_async.call_args_list)) + self.assertEqual( + [call.kwargs["kwargs"] for call in apply_async.call_args_list], + [ + {"listing_id": "100", "lane": "mobile_de_cars", "lease_owner": "lease-100", "attempt": 1}, + {"listing_id": "101", "lane": "mobile_de_cars", "lease_owner": "lease-101", "attempt": 1}, + ], + ) + + def test_detail_queue_respects_pending_limit(self) -> None: + redis_client = self._FakeRedis() + redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY] = {"100", "101"} + persistence = self._FakeDetailPersistence() + + with ( + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true", "MOBILEDE_DETAIL_QUEUE_MAX_PENDING": "2"}), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, + ): + result = tasks._queue_mobilede_detail_listing_ids( + redis_client, + ["mobile.de:102", "mobile.de:103"], + ) + + self.assertEqual(result, {"queued": 0, "duplicate": 0, "full": 2}) + apply_async.assert_not_called() + + def test_detail_queue_skips_recently_completed_listing(self) -> None: + redis_client = self._FakeRedis() + persistence = self._FakeDetailPersistence() + completed_key = tasks.MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id="100") + redis_client.store[completed_key] = "1" + + with ( + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, + ): + result = tasks._queue_mobilede_detail_listing_ids( + redis_client, + ["mobile.de:100", "mobile.de:101"], + ) + + self.assertEqual(result, {"queued": 1, "duplicate": 1, "full": 0}) + self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"101"}) + apply_async.assert_called_once() + + def test_explicit_gallery_refresh_invalidates_completed_listing(self) -> None: + redis_client = self._FakeRedis() + persistence = self._FakeDetailPersistence() + completed_key = tasks.MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id="100") + redis_client.store[completed_key] = "1" + + with ( + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, + ): + result = tasks._queue_mobilede_detail_listing_ids( + redis_client, + ["mobile.de:100"], + invalidate_completed=True, + ) + + self.assertEqual(result, {"queued": 1, "duplicate": 0, "full": 0}) + self.assertNotIn(completed_key, redis_client.store) + self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"100"}) + apply_async.assert_called_once() + + def test_detail_queue_releases_lease_when_dispatch_fails(self) -> None: + redis_client = self._FakeRedis() + persistence = self._FakeDetailPersistence() + + with ( + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.mobilede_sync_detail_task, "apply_async", side_effect=RuntimeError("broker down")), + self.assertRaisesRegex(RuntimeError, "broker down"), + ): + tasks._queue_mobilede_detail_listing_ids(redis_client, ["mobile.de:100"]) + + self.assertEqual(persistence.released, [("100", "lease-100")]) + self.assertNotIn("100", redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY]) + + def test_only_new_segment_with_updates_only_enters_cooldown(self) -> None: + redis_client = self._FakeRedis() + segment = {"make_id": "BMW", "label": "BMW"} + + tasks._mobilede_update_segment_freshness_state( + redis_client, + segment=segment, + only_new=True, + inserted=0, + updated=20, + listings=20, + ) + + cooldown_key = tasks._mobilede_segment_cooldown_key(segment) + hot_key = tasks._mobilede_segment_hot_key(segment) + self.assertEqual(redis_client.store[cooldown_key], "1") + self.assertNotIn(hot_key, redis_client.store) + + def test_detail_scraper_is_reused_within_worker_process(self) -> None: + original = tasks._detail_scraper_instance + tasks._detail_scraper_instance = None + try: + with ( + patch.object(tasks.MobileDeClient, "for_worker") as client_factory, + patch.object(tasks, "_get_persistence") as persistence_factory, + patch.object(tasks, "MobileDeScraper") as scraper_factory, + ): + expected = scraper_factory.return_value + first = tasks._get_detail_scraper() + second = tasks._get_detail_scraper() + + self.assertIs(first, expected) + self.assertIs(second, expected) + client_factory.assert_called_once_with(delay_seconds=0) + persistence_factory.assert_called_once_with() + scraper_factory.assert_called_once() + finally: + tasks._detail_scraper_instance = original + def test_prune_overflow_parent_segments_keeps_only_leaf_segments(self) -> None: parent = { "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:5000&fr=:2004&ml=:100000", @@ -137,6 +369,256 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(pruned, [child_a, child_b]) + def test_compact_runtime_segments_prunes_observed_empty_and_merges_adjacent_ranges(self) -> None: + base_url = "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&fr=2018:2020" + segments = [ + { + "label": "BMW | price=1-100 | year=2018-2020 | total=400", + "search_url": f"{base_url}&p=1:100", + "source_search_url": base_url, + "make_id": "3500", + "price_min": "1", + "price_max": "100", + "year_min": "2018", + "year_max": "2020", + "total_results": 400, + "observed_complete": True, + "max_pages": 50, + }, + { + "label": "BMW | price=101-200 | year=2018-2020 | total=500", + "search_url": f"{base_url}&p=101:200", + "source_search_url": base_url, + "make_id": "3500", + "price_min": "101", + "price_max": "200", + "year_min": "2018", + "year_max": "2020", + "total_results": 500, + "observed_complete": True, + "max_pages": 50, + }, + { + "label": "BMW | price=201-300 | year=2018-2020 | total=0", + "search_url": f"{base_url}&p=201:300", + "source_search_url": base_url, + "make_id": "3500", + "price_min": "201", + "price_max": "300", + "year_min": "2018", + "year_max": "2020", + "total_results": 0, + "observed_complete": True, + "max_pages": 50, + }, + ] + + compacted = tasks._mobilede_compact_runtime_segments(segments) + + self.assertEqual(len(compacted), 1) + self.assertEqual(compacted[0]["price_min"], "1") + self.assertEqual(compacted[0]["price_max"], "200") + self.assertEqual(compacted[0]["total_results"], 900) + self.assertEqual(compacted[0]["max_pages"], 45) + self.assertEqual(compacted[0]["compacted_from"], 2) + + def test_finalize_runtime_plan_compacts_before_interleaving_makes(self) -> None: + base_url = "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car" + segments = [ + { + "label": "BMW | price=1-100 | total=400", + "search_url": f"{base_url}&ms=3500&p=1:100", + "source_search_url": base_url, + "make_id": "3500", + "price_min": "1", + "price_max": "100", + "total_results": 400, + "observed_complete": True, + "max_pages": 20, + }, + { + "label": "BMW | price=101-200 | total=500", + "search_url": f"{base_url}&ms=3500&p=101:200", + "source_search_url": base_url, + "make_id": "3500", + "price_min": "101", + "price_max": "200", + "total_results": 500, + "observed_complete": True, + "max_pages": 25, + }, + { + "label": "Porsche | price=1-5000 | total=100", + "search_url": f"{base_url}&ms=20100&p=1:5000", + "source_search_url": base_url, + "make_id": "20100", + "price_min": "1", + "price_max": "5000", + "total_results": 100, + "observed_complete": True, + "max_pages": 5, + }, + ] + + finalized = planner._mobilede_finalize_runtime_plan(segments) + + self.assertEqual(len(finalized), 2) + self.assertEqual(finalized[0]["make_id"], "3500") + self.assertEqual(finalized[0]["total_results"], 900) + self.assertEqual(finalized[1]["make_id"], "20100") + + def test_complete_runtime_plan_rejects_unknown_and_dense_leaves(self) -> None: + with self.assertRaisesRegex(RuntimeError, "unknown=1 dense=1"): + planner._mobilede_assert_complete_runtime_plan( + [ + {"label": "unknown", "total_results": None}, + {"label": "dense", "total_results": 1001}, + {"label": "valid", "total_results": 950}, + ] + ) + + def test_complete_dense_refine_ignores_soft_probe_limit(self) -> None: + parent = { + "label": "BMW dense", + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:3000", + "make_id": "3500", + "price_min": "1", + "price_max": "3000", + "total_results": 2400, + "max_pages": 50, + } + parent_key = tasks._mobilede_segment_fingerprint(parent) + children = [ + { + **parent, + "label": f"BMW child {index}", + "search_url": f"https://www.mobile.de/ru/search.html?ms=3500&p={start}:{end}", + "price_min": str(start), + "price_max": str(end), + "total_results": None, + "overflow_parent": parent_key, + "overflow_depth": 1, + } + for index, (start, end) in enumerate(((1, 1000), (1001, 2000), (2001, 3000)), start=1) + ] + + with ( + patch.object(planner, "_mobilede_build_overflow_child_segments", return_value=children), + patch.object(planner, "_mobilede_probe_segment_total", side_effect=[800, 800, 800]) as probe, + ): + refined = planner._mobilede_refine_dense_planned_segments([parent]) + + self.assertEqual(probe.call_count, 3) + self.assertEqual([item["total_results"] for item in refined], [800, 800, 800]) + planner._mobilede_assert_complete_runtime_plan(refined) + + def test_dense_refine_recursively_completes_tree_in_one_call(self) -> None: + parent = { + "label": "BMW parent", + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:3000", + "make_id": "3500", + "price_min": "1", + "price_max": "3000", + "total_results": 2400, + "max_pages": 50, + } + child_dense = { + **parent, + "label": "BMW dense child", + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:2000", + "price_max": "2000", + "total_results": None, + "overflow_depth": 1, + } + child_complete = { + **parent, + "label": "BMW complete child", + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=2001:3000", + "price_min": "2001", + "total_results": None, + "overflow_depth": 1, + } + grandchildren = [ + { + **child_dense, + "label": f"BMW grandchild {index}", + "search_url": f"https://www.mobile.de/ru/search.html?ms=3500&p={start}:{end}", + "price_min": str(start), + "price_max": str(end), + "total_results": None, + "overflow_depth": 2, + } + for index, (start, end) in enumerate(((1, 1000), (1001, 2000)), start=1) + ] + + with ( + patch.object( + planner, + "_mobilede_build_overflow_child_segments", + side_effect=[[child_dense, child_complete], grandchildren], + ) as split, + patch.object( + planner, + "_mobilede_probe_segment_total", + side_effect=[1500, 900, 700, 800], + ), + ): + refined = planner._mobilede_refine_dense_planned_segments([parent]) + + self.assertEqual(split.call_count, 2) + self.assertEqual(sorted(item["total_results"] for item in refined), [700, 800, 900]) + planner._mobilede_assert_complete_runtime_plan(refined) + + def test_record_observation_keeps_active_bootstrap_cache_shape(self) -> None: + redis_client = MagicMock() + settings = MagicMock() + first = { + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:100", + "source_search_url": "https://www.mobile.de/ru/search.html?ms=3500", + "make_id": "3500", + "price_min": "1", + "price_max": "100", + "total_results": None, + "max_pages": 50, + } + second = { + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=101:200", + "source_search_url": "https://www.mobile.de/ru/search.html?ms=3500", + "make_id": "3500", + "price_min": "101", + "price_max": "200", + "total_results": 500, + "observed_complete": True, + "max_pages": 25, + } + + with ( + patch.object(tasks, "_acquire_lock", return_value=True), + patch.object(tasks, "_release_lock_if_owner"), + patch.object(tasks, "_get_cached_mobilede_runtime_segments", return_value=[first, second]), + patch.object(planner, "_mobilede_source_segments_from_settings", return_value=[first]), + patch.object(planner, "_mobilede_save_learned_runtime_segments") as save_plan, + ): + removed = tasks._mobilede_record_observed_segment_total( + redis_client, + settings, + segment=first, + total_results=400, + listing_count=400, + unique_count=400, + bootstrap_run=True, + ) + + cached_payload = redis_client.set.call_args.args[1] + cached_segments = json.loads(cached_payload) + self.assertEqual(removed, 1) + self.assertEqual(len(cached_segments), 2) + self.assertEqual(cached_segments[0]["total_results"], 400) + self.assertTrue(cached_segments[0]["observed_complete"]) + learned_segments = save_plan.call_args.args[1] + self.assertEqual(len(learned_segments), 1) + self.assertEqual(learned_segments[0]["total_results"], 900) + def test_runtime_segment_reservation_uses_pending_cache(self) -> None: redis_client = MagicMock() settings = MagicMock() @@ -319,6 +801,57 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(sold_marked, 0) persistence.mark_sold_not_seen_since.assert_not_called() + def test_known_complete_filtered_scope_is_safe_for_reconciliation(self) -> None: + complete_url = "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=25100&ms=8600&ms=15200&ms=3500&ref=dsp" + runtime_config = RuntimeConfig() + settings = MagicMock() + settings.listing.filtered_search_urls = [complete_url] + + with patch.dict( + "os.environ", + {"MOBILEDE_COMPLETE_SCOPE_BRANDS": "BMW,Volvo,Ferrari,Lexus"}, + clear=False, + ): + reasons = tasks._mobilede_complete_scope_unsafe_reasons(runtime_config, settings) + + self.assertEqual(reasons, []) + + def test_partial_filtered_scope_remains_unsafe_for_reconciliation(self) -> None: + runtime_config = RuntimeConfig() + runtime_config.sync.limit = 100 + settings = MagicMock() + settings.listing.filtered_search_urls = ["https://www.mobile.de/ru/search.html?ms=3500"] + + with patch.dict( + "os.environ", + { + "MOBILEDE_COMPLETE_SCOPE_BRANDS": "BMW,Volvo,Ferrari,Lexus", + }, + clear=False, + ): + reasons = tasks._mobilede_complete_scope_unsafe_reasons(runtime_config, settings) + + self.assertIn("sync_limit", reasons) + + def test_effective_full_pass_is_safe_after_only_new_bootstrap(self) -> None: + runtime_config = RuntimeConfig() + runtime_config.sync.only_new = True + settings = MagicMock() + settings.listing.filtered_search_urls = ["https://www.mobile.de/ru/search.html?ms=25100&ms=8600&ms=15200&ms=3500"] + + with patch.dict( + "os.environ", + {"MOBILEDE_COMPLETE_SCOPE_BRANDS": "BMW,Volvo,Ferrari,Lexus"}, + clear=False, + ): + reasons = tasks._mobilede_complete_scope_unsafe_reasons( + runtime_config, + settings, + effective_only_new=False, + ) + + self.assertEqual(reasons, []) + def test_completed_bootstrap_segment_is_skipped_during_active_full_pass(self) -> None: redis_client = MagicMock() redis_client.get.side_effect = lambda key: "1" if str(key).startswith("mobilede:state:bootstrap_segment_done:") else None @@ -485,7 +1018,7 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): "make_id": "3500", } - with patch.object(tasks, "os") as os_mock: + with patch.object(planner, "os") as os_mock: os_mock.getenv.return_value = "true" queued = tasks._queue_mobilede_overflow_child_segments( redis_client, @@ -502,6 +1035,55 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(queued, 0) redis_client.get.assert_not_called() + def test_late_overflow_split_is_saved_for_next_pass_without_extending_bootstrap(self) -> None: + redis_client = MagicMock() + redis_client.sadd.return_value = 1 + settings = MagicMock() + parent = { + "label": "Cars | ms=3500 | price=1-5000", + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:5000", + "make_id": "3500", + "price_min": "1", + "price_max": "5000", + "max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER, + } + parent_key = tasks._mobilede_segment_fingerprint(parent) + child = { + **parent, + "search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:2500", + "price_max": "2500", + "overflow_parent": parent_key, + "overflow_depth": 1, + } + + with ( + patch.object(planner, "_mobilede_build_overflow_child_segments", return_value=[child]), + patch.object(tasks, "_acquire_lock", return_value=True), + patch.object(tasks, "_release_lock_if_owner"), + patch.object(tasks, "_get_cached_mobilede_runtime_segments", return_value=[parent]), + patch.object(planner, "_mobilede_source_segments_from_settings", return_value=[parent]), + patch.object(planner, "_mobilede_save_learned_runtime_segments") as save_plan, + patch.object(planner, "_mobilede_skip_late_overflow_children_during_bootstrap", return_value=True), + ): + added = tasks._mobilede_try_expand_overflow_segment( + redis_client, + settings, + segment=parent, + listing_count=1000, + unique_count=1000, + max_pages=tasks.MOBILEDE_MAX_PAGE_NUMBER, + segment_end_page=tasks.MOBILEDE_MAX_PAGE_NUMBER, + bootstrap_run=True, + ) + + self.assertEqual(added, 0) + self.assertFalse( + any(call.args and call.args[0] == tasks.MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY for call in redis_client.set.call_args_list) + ) + save_plan.assert_called_once() + learned_segments = save_plan.call_args.args[1] + self.assertEqual(learned_segments, [child]) + def test_pre_split_mileage_requires_explicit_flag(self) -> None: original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE try: @@ -532,53 +1114,24 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): finally: tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split - def test_build_runtime_segments_can_fast_start_filtered_search_urls_when_enabled(self) -> None: - import os - - settings = MagicMock() - settings.listing.filtered_search_urls = [ - "https://suchen.mobile.de/fahrzeuge/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=11000&lang=en" - ] - - previous = os.environ.get("MOBILEDE_FILTERED_URL_FAST_START") - try: - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = "true" - segments = tasks._build_mobilede_runtime_segments(settings) - finally: - if previous is None: - os.environ.pop("MOBILEDE_FILTERED_URL_FAST_START", None) - else: - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = previous - - self.assertEqual(len(segments), 2) - self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"3500", "11000"}) - self.assertTrue(all("p=" not in str(segment.get("search_url")) for segment in segments)) - self.assertTrue(all("fr=" not in str(segment.get("search_url")) for segment in segments)) - self.assertTrue(all("ml=" not in str(segment.get("search_url")) for segment in segments)) - def test_build_runtime_segments_splits_multi_make_filtered_search_url(self) -> None: - import os - settings = MagicMock() settings.listing.filtered_search_urls = [ "https://suchen.mobile.de/fahrzeuge/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=11000&lang=en" ] - previous = os.environ.get("MOBILEDE_FILTERED_URL_FAST_START") - try: - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = "false" - with ( - patch.object(planner, "_mobilede_try_keep_root_segment_unsplit", return_value=None), - patch.object(planner, "_expand_mobilede_search_url_segment_by_probe", return_value=None), - patch.object(planner, "_mobilede_probe_total", return_value=None), - patch.object(planner, "_mobilede_preplan_runtime_segments", side_effect=lambda items: items), - ): - segments = tasks._build_mobilede_runtime_segments(settings) - finally: - if previous is None: - os.environ.pop("MOBILEDE_FILTERED_URL_FAST_START", None) - else: - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = previous + with ( + patch.object(planner, "_mobilede_try_keep_root_segment_unsplit", return_value=None), + patch.object(planner, "_expand_mobilede_search_url_segment_by_probe", return_value=None), + patch.object(planner, "_mobilede_probe_total", return_value=None), + patch.object( + planner, + "_mobilede_preplan_runtime_segments", + side_effect=lambda items: [{**item, "total_results": 1} for item in items], + ), + patch.object(planner, "_mobilede_compact_runtime_segments", side_effect=lambda items: items), + ): + segments = tasks._build_mobilede_runtime_segments(settings) self.assertGreater(len(segments), 2) self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"3500", "11000"}) @@ -586,31 +1139,28 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertTrue(all("ms=3500&ms=11000" not in str(segment.get("search_url")) for segment in segments)) def test_full_link_coverage_can_pre_split_every_segment_by_mileage(self) -> None: - import os - settings = MagicMock() settings.listing.filtered_search_urls = [ "https://suchen.mobile.de/fahrzeuge/search.html?isSearchRequest=true&s=Car&vc=Car&ms=111&ms=222&lang=en" ] original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE - previous_fast_start = os.environ.get("MOBILEDE_FILTERED_URL_FAST_START") try: tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = True - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = "false" with ( patch.object(planner, "_mobilede_try_keep_root_segment_unsplit", return_value=None), patch.object(planner, "_expand_mobilede_search_url_segment_by_probe", return_value=None), patch.object(planner, "_mobilede_probe_total", return_value=None), - patch.object(planner, "_mobilede_preplan_runtime_segments", side_effect=lambda items: items), + patch.object( + planner, + "_mobilede_preplan_runtime_segments", + side_effect=lambda items: [{**item, "total_results": 1} for item in items], + ), + patch.object(planner, "_mobilede_compact_runtime_segments", side_effect=lambda items: items), ): segments = tasks._build_mobilede_runtime_segments(settings) finally: tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split - if previous_fast_start is None: - os.environ.pop("MOBILEDE_FILTERED_URL_FAST_START", None) - else: - os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = previous_fast_start self.assertGreater(len(segments), 2) self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"111", "222"})