Compare commits

..

3 Commits

Author SHA1 Message Date
qananasikq
6a46b2f022 add sync coverage tests 2026-05-08 13:31:22 +03:00
qananasikq
66bb930f34 fix sold sync tracking 2026-05-08 13:31:22 +03:00
qananasikq
046e9f125f fix dubizzle listing extraction 2026-05-08 13:31:22 +03:00
5 changed files with 447 additions and 6 deletions

View File

@@ -8,6 +8,11 @@ from ..core.utils import LOT_RE, PRICE_RE, deep_find_all_keys, deep_find_key, fi
logger = logging.getLogger("dubizzle_scraper.parsers")
NEXT_DATA_RE = re.compile(
r'<script[^>]+id=["\']__NEXT_DATA__["\'][^>]*type=["\']application/json["\'][^>]*>(.*?)</script>',
re.IGNORECASE | re.DOTALL,
)
class VehicleParser:
# Парсер страницы авто.
@@ -157,6 +162,11 @@ class VehicleParser:
network_dump = network_dump or {}
responses = network_dump.get("json_responses", [])
payloads = [item.get("payload") for item in responses if isinstance(item.get("payload"), (dict, list))]
embedded = self._extract_embedded_json(page_html)
for item in embedded:
payload = item.get("payload")
if isinstance(payload, (dict, list)):
payloads.append(payload)
dom_kv = self._parse_dom_key_value_pairs(dom_text)
page_title = ""
title_match = re.search(r"<title[^>]*>(.*?)</title>", page_html or "", re.IGNORECASE | re.DOTALL)
@@ -194,11 +204,9 @@ class VehicleParser:
if not summary.get("current_bid") and len(prices) > 1:
summary["current_bid"] = prices[1]
embedded = self._extract_embedded_json(page_html)
for item in embedded:
p = item.get("payload")
if isinstance(p, (dict, list)):
payloads.append(p)
# Доп. проход по JSON.
extra = deep_find_all_keys([p], self.SUMMARY_KEY_MAP)
for field, vals in extra.items():
@@ -299,6 +307,11 @@ class VehicleParser:
@staticmethod
def _guess_currency(summary: dict[str, Any]) -> str:
for key in ["currency", "price", "buy_now", "current_bid", "actual_cash_value", "estimated_repair_cost"]:
value = str(summary.get(key) or "")
upper = value.upper()
if "AED" in upper or "د.إ" in value:
return "AED"
for key in ["buy_now", "current_bid", "actual_cash_value", "estimated_repair_cost"]:
value = str(summary.get(key) or "")
if "$" in value:
@@ -322,9 +335,18 @@ class VehicleParser:
def _extract_embedded_json(html: str) -> list[dict[str, Any]]:
scripts = re.findall(r"<script[^>]*>(.*?)</script>", html or "", flags=re.DOTALL | re.IGNORECASE)
extracted: list[dict[str, Any]] = []
next_match = NEXT_DATA_RE.search(html or "")
if next_match:
try:
next_data = json.loads(html_module.unescape(next_match.group(1).strip()))
extracted.append({"type": "next_data", "payload": next_data})
except Exception:
pass
for script_text in scripts:
if "{" not in script_text and "[" not in script_text:
continue
if "__NEXT_DATA__" in script_text:
continue
# Пропускаем большие блоки.
if len(script_text) > 51_200:
continue

View File

@@ -7,6 +7,7 @@ import signal
import time
import uuid
import html as html_module
from collections.abc import Mapping
from concurrent.futures import ThreadPoolExecutor, as_completed
from concurrent.futures import TimeoutError as FuturesTimeoutError
from datetime import datetime, timezone
@@ -44,6 +45,10 @@ VEHICLE_ID_RE = re.compile(
)
HTML_TAG_RE = re.compile(r"<[^>]+>")
SCRIPT_STYLE_RE = re.compile(r"<(script|style)[^>]*>.*?</\1>", re.IGNORECASE | re.DOTALL)
NEXT_DATA_RE = re.compile(
r'<script[^>]+id=["\']__NEXT_DATA__["\'][^>]*type=["\']application/json["\'][^>]*>(.*?)</script>',
re.IGNORECASE | re.DOTALL,
)
# Таймауты операций страницы.
_PAGE_COLLECT_TIMEOUT_S = 90
@@ -797,6 +802,11 @@ class DUBIZZLEScraper:
total = len(records) if effective_only_new else len(all_raw_urls)
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
all_listing_origin_ids = {
origin_id
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_raw_urls)
if origin_id
}
return {
"listing": {
"status": "ok",
@@ -818,6 +828,7 @@ class DUBIZZLEScraper:
"protection_events": 0,
"failures": failures,
"all_listing_origin_urls": all_listing_origin_urls,
"all_listing_origin_ids": all_listing_origin_ids,
}
def _sync_listing_streaming(
@@ -952,6 +963,11 @@ class DUBIZZLEScraper:
all_raw_urls = list(algolia_result.vehicle_urls)
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
all_listing_origin_ids = {
str(origin_id).strip()
for origin_id in algolia_result.origin_ids_by_url.values()
if str(origin_id).strip()
}
pages_info = list(algolia_result.pages)
effective_urls = list(all_raw_urls)
@@ -1065,6 +1081,7 @@ class DUBIZZLEScraper:
"protection_events": protection_events,
"failures": failures,
"all_listing_origin_urls": all_listing_origin_urls,
"all_listing_origin_ids": all_listing_origin_ids,
}
def _sync_listing_streaming_sitemap_fallback(
@@ -1107,6 +1124,11 @@ class DUBIZZLEScraper:
failures.append({"vehicle_url": f"sitemap_batch_{i}", "error": str(exc)})
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_urls if url}
all_listing_origin_ids = {
origin_id
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_urls)
if origin_id
}
return {
"listing": {
"status": "ok",
@@ -1128,6 +1150,7 @@ class DUBIZZLEScraper:
"protection_events": 0,
"failures": failures,
"all_listing_origin_urls": all_listing_origin_urls,
"all_listing_origin_ids": all_listing_origin_ids,
}
def _sync_listing_streaming_browser_fallback(
@@ -1275,6 +1298,11 @@ class DUBIZZLEScraper:
page.close()
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
all_listing_origin_ids = {
origin_id
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_raw_urls)
if origin_id
}
total = len(all_raw_urls) if not effective_only_new else max(0, len(all_raw_urls) - skipped_existing)
return {
"listing": {
@@ -1296,6 +1324,7 @@ class DUBIZZLEScraper:
"protection_events": 0,
"failures": failures,
"all_listing_origin_urls": all_listing_origin_urls,
"all_listing_origin_ids": all_listing_origin_ids,
}
def _get_page(self) -> Page:
@@ -1314,6 +1343,48 @@ class DUBIZZLEScraper:
_JS_EXTRACT = """
() => {
try {
const nextData = window.__NEXT_DATA__;
const listing = nextData?.props?.pageProps?.reduxWrapperActionsGIPP
?.find((entry) => entry && entry.payload && entry.payload.listing)
?.payload?.listing;
if (listing && listing.details) {
const detailSections = listing.details;
const normalizeSection = (items) => Array.isArray(items)
? items.map((item) => ({
label: item?.label || '',
value: item?.value ?? '',
slug: item?.slug || '',
}))
: [];
return {
ok: true,
source: 'next_data',
listing: {
name: listing.name || '',
description: listing.description || listing.long_description || '',
absolute_url: listing.absolute_url || {},
short_url: listing.short_url || '',
price: listing.price || {},
location: listing.location || {},
posted_timestamp: listing.posted_timestamp || null,
tracking: listing.tracking || {},
categories: Array.isArray(listing.categories) ? listing.categories : [],
details: {
make_model_trim: normalizeSection(detailSections.make_model_trim),
primary: normalizeSection(detailSections.primary),
secondary: normalizeSection(detailSections.secondary),
rental_details: normalizeSection(detailSections.rental_details),
requirements: normalizeSection(detailSections.requirements),
},
photos: Array.isArray(listing.photos_combined)
? listing.photos_combined.map((photo) => photo?.url || photo?.large || photo?.medium || photo?.small || '').filter(Boolean)
: Array.isArray(listing.photos)
? listing.photos.map((photo) => photo?.url || photo?.large || photo?.medium || photo?.small || '').filter(Boolean)
: [],
},
};
}
const scripts = document.querySelectorAll('script:not([src])');
for (const s of scripts) {
const t = s.textContent || '';
@@ -1383,6 +1454,142 @@ class DUBIZZLEScraper:
if not js_data or not js_data.get("ok"):
return {}
if js_data.get("source") == "next_data":
listing = js_data.get("listing") or {}
details = listing.get("details") or {}
categories = listing.get("categories") or []
absolute_url = listing.get("absolute_url") or {}
location = listing.get("location") or {}
tracking = listing.get("tracking") or {}
def detail_value_by_slug(target_slug: str) -> Any:
target = target_slug.strip().lower()
for section_items in details.values():
if not isinstance(section_items, list):
continue
for item in section_items:
if not isinstance(item, Mapping):
continue
if str(item.get("slug") or "").strip().lower() == target:
return item.get("value")
return None
def pick_photo_url(photo: Any) -> str:
if isinstance(photo, str):
return photo.strip()
if not isinstance(photo, Mapping):
return ""
for key in ("url", "main", "large", "medium", "small", "micro"):
value = photo.get(key)
if value:
return str(value).strip()
return ""
detail_sections: dict[str, list[dict[str, Any]]] = {}
flat_details: dict[str, dict[str, dict[str, Any]]] = {}
for section_name, items in details.items():
if not isinstance(items, list):
continue
normalized_items: list[dict[str, Any]] = []
for item in items:
if not isinstance(item, Mapping):
continue
label = str(item.get("label") or "").strip()
value = item.get("value")
slug = str(item.get("slug") or "").strip()
normalized_items.append({"label": {"en": label}, "value": {"en": value}, "slug": slug})
if label:
flat_details[label] = {"en": {"label": label, "value": value}}
if normalized_items:
detail_sections[section_name] = normalized_items
category_v2 = None
category_make = None
category_model = None
if categories:
names_en = [str(item.get("name") or "").strip() for item in categories if isinstance(item, Mapping)]
slug_paths = [str(item.get("full_slug") or item.get("slug") or "").strip() for item in categories if isinstance(item, Mapping)]
ids = [item.get("legacy_id") for item in categories if isinstance(item, Mapping)]
category_v2 = {
"names_en": names_en,
"slug_paths": slug_paths,
"ids": ids,
}
if len(names_en) >= 4:
category_make = names_en[2] or None
category_model = names_en[3] or None
photos: list[str] = []
for photo in listing.get("photos_combined") or []:
picked = pick_photo_url(photo)
if picked and picked not in photos:
photos.append(picked)
for photo in listing.get("photos") or []:
picked = pick_photo_url(photo)
if picked and picked not in photos:
photos.append(picked)
posted_timestamp = listing.get("posted_timestamp")
posted_at_iso = None
if isinstance(posted_timestamp, (int, float)) and posted_timestamp > 0:
posted_at_iso = datetime.fromtimestamp(posted_timestamp, tz=timezone.utc).isoformat()
price = listing.get("price") if isinstance(listing.get("price"), Mapping) else {}
price_raw = price.get("raw") or price.get("formatted")
neighborhood_name = tracking.get("neighbourhood", {}).get("name") if isinstance(tracking.get("neighbourhood"), Mapping) else None
year = detail_value_by_slug("year")
body_type = detail_value_by_slug("body_type")
kilometers = detail_value_by_slug("kilometers")
engine_capacity = detail_value_by_slug("engine_capacity_cc")
transmission_type = detail_value_by_slug("transmission_type")
steering_side = detail_value_by_slug("steering_side")
exterior_color = detail_value_by_slug("exterior_color")
seller_type = detail_value_by_slug("seller_type")
trim = detail_value_by_slug("motors_trim")
summary: dict[str, Any] = {
"source_url": vehicle_url,
"name": {"en": listing.get("name")},
"title": listing.get("name"),
"description": listing.get("description") or listing.get("long_description"),
"make": category_make,
"model": category_model,
"trim": trim,
"year": year,
"body_type": body_type,
"odometer": kilometers,
"kilometers": kilometers,
"engine": engine_capacity,
"engine_volume": engine_capacity,
"transmission_type": transmission_type,
"gearbox": transmission_type,
"steering_side": steering_side,
"steering_wheel": steering_side,
"exterior_color": exterior_color,
"color": exterior_color,
"seller": seller_type,
"price": price_raw,
"buy_now": price_raw,
"currency": price.get("currency") or "AED",
"details_v2": detail_sections,
"details": flat_details,
"image_urls": photos,
"photo_mains": photos,
"category_v2": category_v2,
"location": location.get("name"),
"location_name": location.get("name"),
"site": {"en": "UAE"},
"posted_at": posted_at_iso,
"permalink": listing.get("short_url") or vehicle_url,
"absolute_url": absolute_url,
"id": tracking.get("legacy_id") or listing.get("listing_id") or listing.get("object_id"),
"objectID": listing.get("encoded_object_id") or listing.get("object_id"),
"uuid": listing.get("listing_uuid") or listing.get("uuid"),
}
if neighborhood_name:
summary["neighbourhood"] = {"en": neighborhood_name}
return summary
brnch = str(js_data.get("BranchNumber", "") or "").strip()
img_keys = js_data.get("imageKeys") or []
image_urls = [
@@ -1421,16 +1628,18 @@ class DUBIZZLEScraper:
@staticmethod
def _build_payload_insights(vehicle_summary: dict) -> dict:
currency = vehicle_summary.get("currency") or "USD"
effective_price = vehicle_summary.get("buy_now") or vehicle_summary.get("price")
return {
"vehicle_core": vehicle_summary,
"pricing": {
"buy_now": vehicle_summary.get("buy_now"),
"buy_now": effective_price,
"current_bid": vehicle_summary.get("current_bid"),
"actual_cash_value": vehicle_summary.get("actual_cash_value"),
"estimated_repair_cost": vehicle_summary.get("estimated_repair_cost"),
"currency": "USD",
"currency": currency,
},
"bids": {"amount": vehicle_summary.get("current_bid"), "currency": "USD"},
"bids": {"amount": vehicle_summary.get("current_bid"), "currency": currency},
"damage": {"primary": vehicle_summary.get("primary_damage"), "secondary": vehicle_summary.get("secondary_damage")},
"auction": {"auction_date": vehicle_summary.get("auction_date"), "branch": vehicle_summary.get("location")},
"images": {"count": len(vehicle_summary.get("image_urls") or []), "urls": vehicle_summary.get("image_urls") or []},
@@ -1530,6 +1739,27 @@ class DUBIZZLEScraper:
Использует json.JSONDecoder.raw_decode для быстрого поиска JSON
вместо посимвольного сканирования скобок.
"""
next_match = NEXT_DATA_RE.search(html)
if next_match:
try:
next_data = json.loads(html_module.unescape(next_match.group(1)))
actions = (
next_data.get("props", {})
.get("pageProps", {})
.get("reduxWrapperActionsGIPP", [])
)
for entry in actions:
payload = entry.get("payload") if isinstance(entry, dict) else None
listing = payload.get("listing") if isinstance(payload, dict) else None
if isinstance(listing, dict) and isinstance(listing.get("details"), dict):
return {
"ok": True,
"source": "next_data",
"listing": listing,
}
except Exception:
pass
search = "inventoryView"
pos = html.find(search)
if pos < 0:
@@ -2259,6 +2489,7 @@ class DUBIZZLEScraper:
protection_events = int(stream_result.get("protection_events", 0))
failures.extend(stream_result["failures"])
all_listing_origin_urls = stream_result["all_listing_origin_urls"]
all_listing_origin_ids = stream_result.get("all_listing_origin_ids", set())
logger.info("Streaming sync processed %d vehicles", total)
@@ -2271,6 +2502,13 @@ class DUBIZZLEScraper:
)
if skip_mark_sold:
logger.debug("Skipping mark_sold: caller requested")
elif all_listing_origin_ids and not is_partial_scan:
try:
sold_count = self.persistence.mark_sold_not_in_listing(all_listing_origin_ids)
if sold_count:
logger.info("Marked %d cars as sold by origin_id", sold_count)
except Exception as exc:
logger.warning("Failed to mark sold cars by origin_id: %s", exc)
elif all_listing_origin_urls and not is_partial_scan:
try:
sold_count = self.persistence.mark_sold_not_in_listing_by_urls(all_listing_origin_urls)
@@ -2363,6 +2601,8 @@ class DUBIZZLEScraper:
segment_results: list[dict[str, Any]] = []
completed_all = True
skipped_segments = 0
aggregated_active_urls: set[str] = set()
aggregated_active_ids: set[str] = set()
# В segmented-режиме полный прогон должен проходить ВСЕ сегменты.
# Runtime-фильтры применяются позже (на уровне конкретных карточек),
@@ -2425,6 +2665,7 @@ class DUBIZZLEScraper:
listing_url=seg_url,
year_min=seg_year_min,
year_max=seg_year_max,
skip_mark_sold=True,
)
total_cars_upserted += result.get("cars_upserted", 0)
total_cars_failed += result.get("cars_failed", 0)
@@ -2432,6 +2673,8 @@ class DUBIZZLEScraper:
total_skipped += result.get("skipped_existing", 0)
total_discovered += result.get("listing", {}).get("vehicles_collected", 0)
all_failures.extend(result.get("failures", []))
aggregated_active_urls.update(result.get("all_listing_origin_urls", set()) or set())
aggregated_active_ids.update(result.get("all_listing_origin_ids", set()) or set())
segment_results.append({
"segment": seg,
"segment_index": seg_idx,
@@ -2514,6 +2757,15 @@ class DUBIZZLEScraper:
finally:
self.set_progress_callback(outer_progress_callback)
if completed_all:
try:
if aggregated_active_ids:
self.persistence.mark_sold_not_in_listing(aggregated_active_ids)
elif aggregated_active_urls:
self.persistence.mark_sold_not_in_listing_by_urls(aggregated_active_urls)
except Exception:
logger.warning("Segmented sync final sold-mark failed", exc_info=True)
elapsed = round(time.perf_counter() - started_at, 3)
status = "success" if not all_failures else "partial_success" if total_cars_upserted else "failed"
@@ -2537,6 +2789,8 @@ class DUBIZZLEScraper:
"images_upserted": total_images_upserted,
"skipped_existing": total_skipped,
"total_discovered": total_discovered,
"all_listing_origin_urls": aggregated_active_urls,
"all_listing_origin_ids": aggregated_active_ids,
"elapsed_seconds": elapsed,
"failures": all_failures,
"segment_results": segment_results,

View File

@@ -64,6 +64,7 @@ SYNC_SEGMENTS_PROGRESS_KEY = "dubizzle:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "dubizzle:state:sync_segments_total"
SYNC_SEGMENTED_SCAN_ACTIVE_KEY = "dubizzle:state:sync_segmented_scan_active"
SYNC_SEGMENTED_ACTIVE_URLS_KEY = "dubizzle:state:sync_segmented_active_urls"
SYNC_SEGMENTED_ACTIVE_IDS_KEY = "dubizzle:state:sync_segmented_active_ids"
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "dubizzle:state:task_progress:{task_id}"
GLOBAL_PROGRESS_TS_KEY = "dubizzle:state:last_progress_ts"
@@ -91,6 +92,7 @@ SYNC_SEGMENTS_PROGRESS_KEY = _runtime_key(SYNC_SEGMENTS_PROGRESS_KEY)
SYNC_SEGMENTS_TOTAL_KEY = _runtime_key(SYNC_SEGMENTS_TOTAL_KEY)
SYNC_SEGMENTED_SCAN_ACTIVE_KEY = _runtime_key(SYNC_SEGMENTED_SCAN_ACTIVE_KEY)
SYNC_SEGMENTED_ACTIVE_URLS_KEY = _runtime_key(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
SYNC_SEGMENTED_ACTIVE_IDS_KEY = _runtime_key(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
TASK_PROGRESS_KEY_FMT = _runtime_key(TASK_PROGRESS_KEY_FMT)
GLOBAL_PROGRESS_TS_KEY = _runtime_key(GLOBAL_PROGRESS_TS_KEY)
SITEMAP_HOURLY_LAST_COUNT_KEY = _runtime_key(SITEMAP_HOURLY_LAST_COUNT_KEY)
@@ -759,6 +761,27 @@ def _add_segmented_active_urls(redis_client: Redis, urls: list[str] | set[str])
return 0
def _add_segmented_active_ids(redis_client: Redis, origin_ids: list[str] | set[str]) -> int:
normalized_ids = [str(origin_id).strip() for origin_id in origin_ids if str(origin_id).strip()]
if not normalized_ids:
return 0
try:
added = 0
pipe = redis_client.pipeline()
chunk_size = 5000
for start in range(0, len(normalized_ids), chunk_size):
chunk = normalized_ids[start:start + chunk_size]
pipe.sadd(SYNC_SEGMENTED_ACTIVE_IDS_KEY, *chunk)
pipe.expire(SYNC_SEGMENTED_ACTIVE_IDS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
results = pipe.execute()
added += int(results[0] or 0)
return added
except Exception:
logger.warning("Failed to persist segmented active origin_ids", exc_info=True)
return 0
def _load_segmented_active_urls(redis_client: Redis) -> set[str]:
try:
raw = redis_client.smembers(SYNC_SEGMENTED_ACTIVE_URLS_KEY) or set()
@@ -768,6 +791,15 @@ def _load_segmented_active_urls(redis_client: Redis) -> set[str]:
return {str(url).strip() for url in raw if str(url).strip()}
def _load_segmented_active_ids(redis_client: Redis) -> set[str]:
try:
raw = redis_client.smembers(SYNC_SEGMENTED_ACTIVE_IDS_KEY) or set()
except Exception:
logger.warning("Failed to read segmented active origin_ids", exc_info=True)
return set()
return {str(origin_id).strip() for origin_id in raw if str(origin_id).strip()}
def _clear_segmented_active_urls(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
@@ -775,6 +807,13 @@ def _clear_segmented_active_urls(redis_client: Redis) -> None:
logger.warning("Failed to clear segmented active URLs", exc_info=True)
def _clear_segmented_active_ids(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
except Exception:
logger.warning("Failed to clear segmented active origin_ids", exc_info=True)
def _is_segmented_scan_in_progress(redis_client: Redis, celery_app) -> bool:
try:
marker = redis_client.get(SYNC_SEGMENTED_SCAN_ACTIVE_KEY)
@@ -1036,6 +1075,7 @@ def _reset_segments_progress(redis_client: Redis, total: int) -> None:
pipe = redis_client.pipeline()
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.delete(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
pipe.delete(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.execute()
except Exception:
@@ -1152,9 +1192,12 @@ def sync_segment_task(
segment_done = bool(result.get("full_scan_completed", False))
listing_payload = result.get("listing") if isinstance(result.get("listing"), dict) else {}
active_urls = listing_payload.get("vehicle_urls") or []
active_ids = result.get("all_listing_origin_ids") or []
active_urls_count = len(active_urls)
if active_urls:
_add_segmented_active_urls(redis_client, active_urls)
if active_ids:
_add_segmented_active_ids(redis_client, active_ids)
_update_task_progress(
redis_client,
@@ -1176,8 +1219,22 @@ def sync_segment_task(
)
if total > 0 and completed >= total:
sold_count = 0
active_ids_for_scan = _load_segmented_active_ids(redis_client)
active_urls_for_scan = _load_segmented_active_urls(redis_client)
if active_urls_for_scan:
if active_ids_for_scan:
try:
sold_count = persistence.mark_sold_not_in_listing(
active_ids_for_scan,
lane=_origin_prefix().rstrip(":"),
)
logger.info(
"Segmented full scan sold-mark completed by origin_id: sold=%d active_ids=%d",
sold_count,
len(active_ids_for_scan),
)
except Exception:
logger.warning("Segmented full scan sold-mark by origin_id failed", exc_info=True)
elif active_urls_for_scan:
try:
sold_count = persistence.mark_sold_not_in_listing_by_urls(
active_urls_for_scan,
@@ -1198,6 +1255,7 @@ def sync_segment_task(
_clear_sync_checkpoint(redis_client)
_clear_segmented_scan_active(redis_client)
_clear_segmented_active_urls(redis_client)
_clear_segmented_active_ids(redis_client)
_clear_bootstrap_failure_streak(redis_client)
_clear_bootstrap_continuation_streak(redis_client)
if always_full_scan:

View File

@@ -122,6 +122,50 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(len(cars), 1)
self.assertEqual(cars[0].origin_id, "NEW-ID")
def test_mark_sold_by_urls_marks_missing_records(self) -> None:
first = self._record("111")
second = self._record("222")
first.origin_id = "dubizzle:111"
second.origin_id = "dubizzle:222"
self.persistence.upsert_car(first)
self.persistence.upsert_car(second)
marked = self.persistence.mark_sold_not_in_listing_by_urls({first.origin_url})
self.assertEqual(marked, 1)
with self.persistence.session_scope() as session:
cars = {
car.origin_id: car
for car in session.execute(select(Car)).scalars().all()
}
self.assertFalse(cars["dubizzle:111"].is_sold)
self.assertTrue(cars["dubizzle:222"].is_sold)
def test_upsert_reactivates_previously_sold_car(self) -> None:
record = self._record("333", price=1000)
record.origin_id = "dubizzle:333"
self.persistence.upsert_car(record)
self.persistence.mark_sold_not_in_listing_by_urls({"https://www.dubizzle.com/VehicleDetail/other~US"})
with self.persistence.session_scope() as session:
sold_car = session.execute(
select(Car).where(Car.origin_id == "dubizzle:333")
).scalar_one()
self.assertTrue(sold_car.is_sold)
revived = self._record("333", price=1500)
revived.origin_id = "dubizzle:333"
self.persistence.upsert_car(revived)
with self.persistence.session_scope() as session:
active_car = session.execute(
select(Car).where(Car.origin_id == "dubizzle:333")
).scalar_one()
self.assertFalse(active_car.is_sold)
self.assertEqual(active_car.price, 1500)
if __name__ == "__main__":
unittest.main()

View File

@@ -269,6 +269,69 @@ class TestScraperSync(unittest.TestCase):
self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(result["total"], 1)
def test_sync_listing_marks_sold_after_full_scan(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=7)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.mark_sold_not_in_listing = MagicMock(return_value=1)
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock(return_value=0)
scraper._sync_listing_streaming = MagicMock(return_value={
"listing": {
"vehicles_collected": 2,
"early_stopped": False,
"truncated_by_time_budget": False,
},
"total": 2,
"skipped_existing": 0,
"cars_upserted": 2,
"cars_failed": 0,
"images_upserted": 4,
"failures": [],
"all_listing_origin_urls": {
"https://www.dubizzle.com/VehicleDetail/111~US",
"https://www.dubizzle.com/VehicleDetail/222~US",
},
"all_listing_origin_ids": {
"dubizzle:111",
"dubizzle:222",
},
})
result = scraper.sync_listing(only_new=False, limit=None)
self.assertEqual(result["status"], "success")
scraper.persistence.mark_sold_not_in_listing.assert_called_once()
scraper.persistence.mark_sold_not_in_listing_by_urls.assert_not_called()
def test_sync_listing_skips_mark_sold_for_partial_scan(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=8)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock()
scraper._sync_listing_streaming = MagicMock(return_value={
"listing": {
"vehicles_collected": 1,
"early_stopped": False,
"truncated_by_time_budget": False,
},
"total": 1,
"skipped_existing": 0,
"cars_upserted": 1,
"cars_failed": 0,
"images_upserted": 2,
"failures": [],
"all_listing_origin_urls": {
"https://www.dubizzle.com/VehicleDetail/111~US",
},
})
result = scraper.sync_listing(only_new=True)
self.assertEqual(result["status"], "success")
scraper.persistence.mark_sold_not_in_listing_by_urls.assert_not_called()
def test_sync_batch_timeout_does_not_duplicate_already_processed_urls(self) -> None:
scraper = self._make_scraper()
scraper.persistence.upsert_cars_batch = MagicMock(return_value={