diff --git a/mobilede_scraper/worker/search_sync.py b/mobilede_scraper/worker/search_sync.py index ab55437..20e8112 100644 --- a/mobilede_scraper/worker/search_sync.py +++ b/mobilede_scraper/worker/search_sync.py @@ -13,6 +13,27 @@ globals().update( ) +def _intersect_numeric_bounds( + current_min: str | None, + current_max: str | None, + filter_min: int | None, + filter_max: int | None, +) -> tuple[str | None, str | None, bool]: + segment_min = int(current_min) if current_min not in (None, "") else None + segment_max = int(current_max) if current_max not in (None, "") else None + lower_values = [value for value in (segment_min, filter_min) if value is not None] + upper_values = [value for value in (segment_max, filter_max) if value is not None] + lower = max(lower_values) if lower_values else None + upper = min(upper_values) if upper_values else None + if lower is not None and upper is not None and lower > upper: + return current_min, current_max, False + return ( + str(lower) if lower is not None else None, + str(upper) if upper is not None else None, + True, + ) + + def run_mobilede_sync_search_task( self, start_page: int = 1, @@ -387,6 +408,21 @@ def run_mobilede_sync_search_task( if stage == "search_collection_done": window_pages_collected = int(meta.get("pages_collected", 0) or 0) elif stage == "db_upsert_done": + inserted_origin_ids = [ + str(origin_id) + for origin_id in meta.get("inserted_origin_ids", []) + if origin_id + ] + detail_queued = 0 + if inserted_origin_ids: + detail_queue = _queue_mobilede_detail_listing_ids( + redis_client, + inserted_origin_ids, + lane=lane, + priority=9, + ) + detail_queued += detail_queue["queued"] + meta["detail_queued"] = detail_queued if refresh_cycle_id: _refresh_done, total_segments_for_log, _refresh_left = _mobilede_refresh_cycle_progress( redis_client, @@ -397,13 +433,17 @@ def run_mobilede_sync_search_task( total_segments_raw = redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) total_segments_for_log = int(total_segments_raw) if total_segments_raw else None logger.info( - "mobile.de db upsert: segment_no=%s segment_key=%s segment=%s page=%s inserted=%s updated=%s images=%s run_id=%s", + "mobile.de db upsert: segment_no=%s segment_key=%s segment=%s page=%s inserted=%s changed=%s unchanged=%s reappeared=%s updated=%s detail_queued=%s images=%s run_id=%s", _mobilede_segment_position_label(segment_index, total_segments_for_log), _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), meta.get("page_number"), int(meta.get("inserted", 0) or 0), + int(meta.get("changed", 0) or 0), + int(meta.get("unchanged", 0) or 0), + int(meta.get("reappeared", 0) or 0), int(meta.get("updated", 0) or 0), + detail_queued, int(meta.get("images_upserted", 0) or 0), meta.get("run_id"), ) @@ -426,17 +466,26 @@ def run_mobilede_sync_search_task( runtime_settings=settings, ) runtime_filters = runtime_config.filters - if runtime_filters.price.min is not None: - price_min = str(runtime_filters.price.min) - if runtime_filters.price.max is not None: - price_max = str(runtime_filters.price.max) - if runtime_filters.mileage.min is not None: - mileage_min = str(runtime_filters.mileage.min) - if runtime_filters.mileage.max is not None: - mileage_max = str(runtime_filters.mileage.max) + price_min, price_max, price_intersects = _intersect_numeric_bounds( + price_min, price_max, runtime_filters.price.min, runtime_filters.price.max + ) + mileage_min, mileage_max, mileage_intersects = _intersect_numeric_bounds( + mileage_min, mileage_max, runtime_filters.mileage.min, runtime_filters.mileage.max + ) if runtime_filters.include.years: - year_min = str(min(runtime_filters.include.years)) - year_max = str(max(runtime_filters.include.years)) + year_min, year_max, year_intersects = _intersect_numeric_bounds( + year_min, + year_max, + min(runtime_filters.include.years), + max(runtime_filters.include.years), + ) + else: + year_intersects = True + if not all((price_intersects, mileage_intersects, year_intersects)): + logger.info( + "mobile.de runtime filters do not intersect planner segment; retaining segment bounds for safe post-filtering: runtime=%s", + _mobilede_short_segment_label(segment), + ) run_seen_at = _mobilede_refresh_cycle_seen_at(redis_client, cycle_id=refresh_cycle_id) result = scraper.sync_search( start_page=actual_start_page, @@ -492,6 +541,15 @@ def run_mobilede_sync_search_task( ) ) if segment_scan_complete: + _mobilede_record_observed_segment_total( + redis_client, + settings, + segment=segment, + total_results=_mobilede_parse_optional_int(result.get("total_results")), + listing_count=listing_count, + unique_count=unique_count, + bootstrap_run=bootstrap_run_active, + ) overflow_segments_added = _mobilede_try_expand_overflow_segment( redis_client, settings, @@ -500,6 +558,7 @@ def run_mobilede_sync_search_task( unique_count=unique_count, max_pages=max_pages, segment_end_page=actual_end_page, + bootstrap_run=bootstrap_run_active, ) queued_overflow_children = _queue_mobilede_overflow_child_segments( redis_client, @@ -1095,6 +1154,10 @@ def run_mobilede_sync_search_task( max_retries = int(getattr(self, "max_retries", 0) or 0) current_retries = int(getattr(self.request, "retries", 0) or 0) if is_transient_request_error: + network_retry_delay = max( + 5, + int(float(os.getenv("MOBILEDE_SEARCH_NETWORK_RETRY_DELAY_SECONDS", "30"))), + ) logger.warning( "mobile.de network issue: runtime=%s filter=%s pages=%s-%s retry=%s/%s error=%s", _mobilede_segment_label(segment), @@ -1106,7 +1169,10 @@ def run_mobilede_sync_search_task( exc, ) if current_retries < max_retries: - raise self.retry(exc=exc, countdown=max(60, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS)) + raise self.retry( + exc=exc, + countdown=network_retry_delay * (current_retries + 1), + ) if continuous is None: continuous = MOBILEDE_CONTINUOUS_SYNC_ENABLED @@ -1138,7 +1204,7 @@ def run_mobilede_sync_search_task( delayed_retry = ( MOBILEDE_ANTIBOT_BACKOFF_SECONDS if is_antibot_block - else max(300, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS * 4) + else max(60, network_retry_delay * 4) ) if _try_set_mobilede_followup_pending( redis_client,