mirror of
https://github.com/R0m1k3/Priceflow.git
synced 2026-10-11 17:29:14 +02:00
893 lines
35 KiB
Python
893 lines
35 KiB
Python
"""
|
|
Improved Search Service - Using Persistent Browser Connection
|
|
Based on ScraperService pattern for better session management and reliability
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import re
|
|
from typing import AsyncGenerator, Optional
|
|
from urllib.parse import quote_plus, urljoin, urlparse
|
|
|
|
from bs4 import BeautifulSoup
|
|
from playwright.async_api import Browser, BrowserContext, Page, async_playwright
|
|
from sqlalchemy.orm import Session
|
|
|
|
from app.models import SearchSite
|
|
from app.schemas import SearchProgress, SearchResultItem
|
|
|
|
|
|
from app.core.search_config import SITE_CONFIGS, BROWSERLESS_URL
|
|
from app.services.ai_price_extractor import AIPriceExtractor
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Common popup/cookie selectors
|
|
COMMON_POPUP_SELECTORS = [
|
|
"#sp-cc-accept", # Cookie banner
|
|
"#onetrust-accept-btn-handler", # OneTrust
|
|
".cookie-consent-accept",
|
|
"[data-action='accept-cookies']",
|
|
"button[id*='accept']",
|
|
"button[class*='accept']",
|
|
]
|
|
|
|
|
|
class SearchResult:
|
|
"""Search result data class"""
|
|
def __init__(
|
|
self,
|
|
url: str,
|
|
title: str,
|
|
snippet: str,
|
|
source: str,
|
|
price: float | None = None,
|
|
currency: str = "EUR",
|
|
in_stock: bool | None = None,
|
|
image_url: str | None = None,
|
|
):
|
|
self.url = url
|
|
self.title = title
|
|
self.snippet = snippet
|
|
self.source = source
|
|
self.price = price
|
|
self.currency = currency
|
|
self.in_stock = in_stock
|
|
self.image_url = image_url
|
|
|
|
def to_dict(self):
|
|
return {
|
|
"url": self.url,
|
|
"title": self.title,
|
|
"snippet": self.snippet,
|
|
"source": self.source,
|
|
"price": self.price,
|
|
"currency": self.currency,
|
|
"in_stock": self.in_stock,
|
|
"image_url": self.image_url,
|
|
}
|
|
|
|
|
|
class ImprovedSearchService:
|
|
"""Persistent browser service for e-commerce search scraping"""
|
|
|
|
_playwright = None
|
|
_browser: Browser | None = None
|
|
_lock = asyncio.Lock()
|
|
|
|
@classmethod
|
|
async def initialize(cls):
|
|
"""Initialize shared browser (Thread-Safe)"""
|
|
async with cls._lock:
|
|
await cls._initialize()
|
|
|
|
@classmethod
|
|
async def _initialize(cls):
|
|
"""Internal initialization"""
|
|
if cls._browser is None:
|
|
logger.info("Initializing ImprovedSearchService shared browser...")
|
|
cls._playwright = await async_playwright().start()
|
|
cls._browser = await cls._connect_browser(cls._playwright)
|
|
logger.info("ImprovedSearchService initialized.")
|
|
|
|
@classmethod
|
|
async def shutdown(cls):
|
|
"""Shutdown shared browser"""
|
|
async with cls._lock:
|
|
if cls._browser:
|
|
logger.info("Shutting down ImprovedSearchService...")
|
|
await cls._browser.close()
|
|
cls._browser = None
|
|
if cls._playwright:
|
|
await cls._playwright.stop()
|
|
cls._playwright = None
|
|
logger.info("ImprovedSearchService shutdown complete.")
|
|
|
|
@classmethod
|
|
async def _ensure_browser_connected(cls) -> bool:
|
|
"""Ensure browser is connected, reconnect if needed"""
|
|
async with cls._lock:
|
|
try:
|
|
if cls._browser is None:
|
|
logger.warning("Browser not initialized, initializing...")
|
|
await cls._initialize()
|
|
return cls._browser is not None
|
|
|
|
# Test connection
|
|
try:
|
|
test_context = await cls._browser.new_context()
|
|
await test_context.close()
|
|
return True
|
|
except Exception as e:
|
|
logger.error(f"Browser connection test failed: {e}")
|
|
logger.info("Attempting to reconnect...")
|
|
cls._browser = None
|
|
if cls._playwright:
|
|
try:
|
|
await cls._playwright.stop()
|
|
except Exception:
|
|
pass
|
|
cls._playwright = None
|
|
await cls._initialize()
|
|
return cls._browser is not None
|
|
except Exception as e:
|
|
logger.error(f"Failed to ensure browser connection: {e}")
|
|
return False
|
|
|
|
@staticmethod
|
|
async def _connect_browser(p) -> Browser:
|
|
"""Connect to Browserless"""
|
|
logger.info(f"Connecting to Browserless at {BROWSERLESS_URL}")
|
|
return await p.chromium.connect_over_cdp(BROWSERLESS_URL)
|
|
|
|
@staticmethod
|
|
async def _create_context(browser: Browser) -> BrowserContext:
|
|
"""Create browser context with stealth settings"""
|
|
context = await browser.new_context(
|
|
viewport={"width": 1920, "height": 1080},
|
|
user_agent=(
|
|
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
|
"AppleWebKit/537.36 (KHTML, like Gecko) "
|
|
"Chrome/131.0.0.0 Safari/537.36"
|
|
),
|
|
locale="fr-FR",
|
|
timezone_id="Europe/Paris",
|
|
)
|
|
|
|
# Stealth mode
|
|
await context.add_init_script("""
|
|
Object.defineProperty(navigator, 'webdriver', { get: () => undefined });
|
|
window.chrome = { runtime: {} };
|
|
""")
|
|
|
|
await context.route("**/*", lambda route: route.continue_())
|
|
return context
|
|
|
|
@staticmethod
|
|
async def _handle_popups(page: Page):
|
|
"""Close common popups/cookies"""
|
|
logger.debug("Handling popups...")
|
|
for selector in COMMON_POPUP_SELECTORS:
|
|
try:
|
|
if await page.locator(selector).count() > 0:
|
|
logger.debug(f"Found popup: {selector}")
|
|
await page.locator(selector).first.click(timeout=2000)
|
|
await page.wait_for_timeout(500)
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
|
|
@staticmethod
|
|
def _parse_results(html: str, site_key: str, base_url: str, query: str) -> list[SearchResult]:
|
|
"""
|
|
DEPRECATED: Legacy parsing method - now using specialized parsers
|
|
This method is kept for backward compatibility but should not be used
|
|
"""
|
|
config = SITE_CONFIGS[site_key]
|
|
soup = BeautifulSoup(html, "html.parser")
|
|
results = []
|
|
|
|
# Prepare query words for filtering
|
|
query_words = [w.lower() for w in query.split() if len(w) > 2]
|
|
|
|
# Select product links
|
|
logger.debug(f"Parsing content for {site_key} (length: {len(html)}) with selector: {config['product_selector']}")
|
|
links = soup.select(config["product_selector"])
|
|
logger.debug(f"Found {len(links)} raw items for {site_key}")
|
|
|
|
# Deduplicate links
|
|
seen_urls = set()
|
|
|
|
for item in links:
|
|
# Keep reference to original container for image search
|
|
container = item
|
|
if site_key == "centrakor.com":
|
|
logger.debug(f" Processing Centrakor item: {item.name}, classes: {item.get('class')}")
|
|
|
|
# If selector targets the container (div.product-card), we need to find the link inside
|
|
if item.name != 'a':
|
|
# Try to find the link using product_link_selector if available
|
|
if "product_link_selector" in config:
|
|
found_link = item.select_one(config["product_link_selector"])
|
|
if found_link:
|
|
link = found_link
|
|
else:
|
|
# Fallback: find first 'a' tag
|
|
found_link = item.find('a')
|
|
if found_link:
|
|
link = found_link
|
|
else:
|
|
logger.debug(f" ⚠️ No link found in container for {site_key}")
|
|
continue
|
|
else:
|
|
# No link selector specified, try to find first 'a'
|
|
found_link = item.find('a')
|
|
if found_link:
|
|
link = found_link
|
|
else:
|
|
logger.debug(f" ⚠️ No link found in container for {site_key}")
|
|
continue
|
|
else:
|
|
# Item is already a link
|
|
link = item
|
|
|
|
href = link.get("href")
|
|
if not href:
|
|
logger.debug(f" ⚠️ No href in link for {site_key}")
|
|
continue
|
|
|
|
full_url = urljoin(base_url, href)
|
|
|
|
# Basic cleanup
|
|
if full_url in seen_urls:
|
|
continue
|
|
seen_urls.add(full_url)
|
|
|
|
# Extract title
|
|
title = link.get_text(strip=True)
|
|
|
|
# If no text, check title attribute or nested image alt
|
|
if not title:
|
|
if link.get("title"):
|
|
title = link.get("title")
|
|
else:
|
|
img = link.find("img")
|
|
if img and img.get("alt"):
|
|
title = img.get("alt")
|
|
|
|
if not title or len(title) < 3:
|
|
logger.debug(f" ⚠️ Skipped item with empty/short title: '{title}'")
|
|
continue
|
|
|
|
# Extract Image URL - PRIORITIZE CONFIGURED SELECTOR
|
|
image_url = None
|
|
|
|
# PRIORITY 1: Use product_image_selector if configured (site-specific)
|
|
if "product_image_selector" in config:
|
|
# Search in the original container first - get ALL matches
|
|
img_els = container.select(config["product_image_selector"])
|
|
|
|
# Filter and find first valid image
|
|
for img_el in img_els:
|
|
# Try multiple attributes in order of priority
|
|
candidate_url = (
|
|
img_el.get("src") or
|
|
img_el.get("data-src") or
|
|
img_el.get("data-lazy-src") or
|
|
img_el.get("data-original")
|
|
)
|
|
|
|
# Handle srcset (use first URL)
|
|
if not candidate_url and img_el.get("srcset"):
|
|
srcset = img_el.get("srcset")
|
|
candidate_url = srcset.split(",")[0].split()[0]
|
|
|
|
# Skip invalid images (pictos, icons, etc.)
|
|
if candidate_url:
|
|
# SPECIAL: No filtering for Centrakor (debugging)
|
|
if site_key == "centrakor.com":
|
|
# Skip placeholders
|
|
if "placeholder" in candidate_url.lower():
|
|
logger.debug(f" ⏭️ Skipping placeholder: {candidate_url[:50]}")
|
|
continue
|
|
# Filter only tiny pictos
|
|
if 'picto' in candidate_url.lower() and ('width=60' in candidate_url or 'height=80' in candidate_url):
|
|
logger.debug(f" ⏭️ Skipping tiny picto: {candidate_url[:50]}")
|
|
continue
|
|
image_url = candidate_url
|
|
logger.debug(f" 🖼️ Centrakor image: {image_url[:70]}")
|
|
break
|
|
|
|
# Normal filtering for other sites
|
|
# Filter out obvious pictos and small icons
|
|
if any(keyword in candidate_url.lower() for keyword in ['picto', 'icon', 'logo', 'badge']):
|
|
logger.debug(f" ⏭️ Skipping picto/icon: {candidate_url[:50]}")
|
|
continue
|
|
# Filter out VERY small images (less than 100px)
|
|
import re
|
|
width_match = re.search(r'width=(\d+)', candidate_url)
|
|
height_match = re.search(r'height=(\d+)', candidate_url)
|
|
if width_match and height_match:
|
|
width = int(width_match.group(1))
|
|
height = int(height_match.group(1))
|
|
if width < 100 and height < 100:
|
|
logger.debug(f" ⏭️ Skipping small image ({width}x{height}): {candidate_url[:50]}")
|
|
continue
|
|
|
|
# This is a valid product image
|
|
image_url = candidate_url
|
|
logger.debug(f" 🖼️ Image found via product_image_selector: {image_url[:50]}...")
|
|
break
|
|
|
|
# PRIORITY 2: Fallback - Look for any img directly in the link
|
|
if not image_url:
|
|
img = link.find("img")
|
|
if img:
|
|
image_url = (
|
|
img.get("src") or
|
|
img.get("data-src") or
|
|
img.get("data-lazy-src") or
|
|
img.get("data-original")
|
|
)
|
|
|
|
# Handle srcset
|
|
if not image_url and img.get("srcset"):
|
|
srcset = img.get("srcset")
|
|
image_url = srcset.split(",")[0].split()[0]
|
|
|
|
if image_url:
|
|
logger.debug(f" 🖼️ Image found via link.find('img'): {image_url[:50]}...")
|
|
|
|
# PRIORITY 3: Look for picture > source elements
|
|
if not image_url:
|
|
picture = link.find("picture")
|
|
if picture:
|
|
source = picture.find("source")
|
|
if source and source.get("srcset"):
|
|
srcset = source.get("srcset")
|
|
image_url = srcset.split(",")[0].split()[0]
|
|
|
|
if not image_url:
|
|
img_in_picture = picture.find("img")
|
|
if img_in_picture:
|
|
image_url = img_in_picture.get("src") or img_in_picture.get("data-src")
|
|
|
|
if image_url:
|
|
logger.debug(f" 🖼️ Image found via <picture>: {image_url[:50]}...")
|
|
|
|
# Clean up and validate image URL
|
|
if image_url:
|
|
# Remove data URIs, 1x1 pixels, placeholders
|
|
if (image_url.startswith("data:") or
|
|
"1x1" in image_url or
|
|
"placeholder" in image_url.lower() or
|
|
image_url.strip() == ""):
|
|
logger.debug(f" ⏭️ Skipping invalid image: {image_url[:50]}")
|
|
image_url = None
|
|
elif not image_url.startswith("http"):
|
|
image_url = urljoin(base_url, image_url)
|
|
logger.debug(f" 🔗 Made image URL absolute: {image_url[:80]}...")
|
|
|
|
if not image_url:
|
|
logger.warning(f" ⚠️ No image found for: {title[:50]}")
|
|
|
|
# Extract price from search results for sites where product pages are unavailable
|
|
product_price = None
|
|
if site_key == "gifi.fr":
|
|
# Gifi: Extract price from product tile HTML
|
|
import re
|
|
container_html = str(container)
|
|
price_match = re.search(r'(\d+)[,.](\d+)\s*€', container_html)
|
|
if price_match:
|
|
euros = int(price_match.group(1))
|
|
cents = int(price_match.group(2))
|
|
product_price = float(f"{euros}.{cents}")
|
|
logger.debug(f" 💰 Extracted price from search: {product_price}€")
|
|
|
|
# Create result
|
|
results.append(SearchResult(
|
|
url=full_url,
|
|
title=title,
|
|
snippet=f"Product from {config['name']}",
|
|
source=config["name"],
|
|
price=product_price, # Set price if extracted from search
|
|
image_url=image_url
|
|
))
|
|
|
|
logger.debug(f"Parsed {len(results)} results from HTML")
|
|
return results
|
|
|
|
@classmethod
|
|
async def _scrape_item_details(cls, result: SearchResult, context: BrowserContext) -> SearchResult | None:
|
|
"""Scrape price and details for a single item using same context"""
|
|
try:
|
|
# SPECIAL CASE: L'Incroyable - Price is in the title
|
|
if "lincroyable.fr" in result.url:
|
|
import re
|
|
# Extract price from title (e.g., "34€99" or "59€99")
|
|
price_match = re.search(r'(\d+)€(\d+)', result.title)
|
|
if price_match:
|
|
# Convert to float (e.g., "34€99" -> 34.99)
|
|
price_euros = int(price_match.group(1))
|
|
price_cents = int(price_match.group(2))
|
|
result.price = float(f"{price_euros}.{price_cents}")
|
|
|
|
# Clean title by removing price
|
|
result.title = re.sub(r'\d+€\d+', '', result.title).strip()
|
|
logger.debug(f"L'Incroyable - Extracted price {result.price}€ from title")
|
|
|
|
return result
|
|
|
|
# SPECIAL CASE: Gifi - Price extracted from search, no need to visit page
|
|
if "gifi.fr" in result.url and result.price is not None:
|
|
logger.debug(f"Gifi - Price already extracted from search: {result.price}€")
|
|
return result
|
|
|
|
page = await context.new_page()
|
|
try:
|
|
logger.debug(f"Scraping details for: {result.title[:50]}...")
|
|
|
|
# Navigate to product page
|
|
await page.goto(result.url, wait_until="domcontentloaded", timeout=20000)
|
|
|
|
# Wait for network idle
|
|
try:
|
|
await page.wait_for_load_state("networkidle", timeout=5000)
|
|
except PlaywrightTimeoutError:
|
|
pass
|
|
|
|
# Handle popups
|
|
await cls._handle_popups(page)
|
|
|
|
# Wait for content
|
|
await page.wait_for_timeout(1500)
|
|
|
|
# Extract price using multiple selectors
|
|
price = await cls._extract_price(page)
|
|
result.price = price
|
|
|
|
# Extract stock status
|
|
in_stock = await cls._extract_stock_status(page)
|
|
result.in_stock = in_stock
|
|
|
|
# Keep original image URL from search page - don't replace with screenshot
|
|
# (Screenshots would need to be served by FastAPI, and original images are already good)
|
|
|
|
logger.debug(f" ✓ {result.title[:40]}... - {price}€ - Stock: {in_stock}")
|
|
return result
|
|
|
|
finally:
|
|
await page.close()
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error scraping item details {result.url}: {e}")
|
|
return result # Return original result without price
|
|
|
|
@staticmethod
|
|
async def _extract_price(page: Page) -> float | None:
|
|
"""Extract price using Hybrid Strategy: JSON-LD -> AI -> CSS"""
|
|
import re
|
|
|
|
# STRATEGY 1: JSON-LD (Most reliable)
|
|
try:
|
|
json_ld_scripts = await page.query_selector_all('script[type="application/ld+json"]')
|
|
for script in json_ld_scripts:
|
|
try:
|
|
content = await script.inner_text()
|
|
data = json.loads(content)
|
|
|
|
# Handle list of objects
|
|
if isinstance(data, list):
|
|
data = data[0] if data else {}
|
|
|
|
# Check for Product type
|
|
if data.get('@type') == 'Product':
|
|
offers = data.get('offers')
|
|
if isinstance(offers, list):
|
|
offers = offers[0]
|
|
|
|
if offers and 'price' in offers:
|
|
price = float(offers['price'])
|
|
logger.debug(f" ✅ Found price via JSON-LD: {price}€")
|
|
return price
|
|
except:
|
|
continue
|
|
except Exception as e:
|
|
logger.debug(f" JSON-LD extraction failed: {e}")
|
|
|
|
# STRATEGY 2: AI Extraction (Priority over CSS as requested)
|
|
try:
|
|
logger.info(" 🤖 Attempting AI extraction (Priority Strategy)...")
|
|
html_content = await page.content()
|
|
title = await page.title()
|
|
|
|
ai_price = await AIPriceExtractor.extract_price(html_content, title)
|
|
if ai_price:
|
|
logger.info(f" 🤖 AI found price: {ai_price}€")
|
|
return ai_price
|
|
except Exception as e:
|
|
logger.error(f" AI extraction failed: {e}")
|
|
|
|
# STRATEGY 3: CSS Selectors (Fallback)
|
|
|
|
# PRIORITY 1: Selectors for sale/promotional prices (highest priority)
|
|
sale_price_selectors = [
|
|
'.price-current',
|
|
'.prix-actuel',
|
|
'.sale-price',
|
|
'.promo-price',
|
|
'[class*="promo"]',
|
|
'[class*="sale"]',
|
|
'[class*="discount"]',
|
|
]
|
|
|
|
# PRIORITY 2: Standard price selectors
|
|
price_selectors = [
|
|
'.price',
|
|
'[data-testid="price"]',
|
|
'[itemprop="price"]',
|
|
'.product-price',
|
|
'.a-price .a-offscreen',
|
|
'.a-price-whole',
|
|
'span[class*="price"]',
|
|
]
|
|
|
|
# Try sale prices first
|
|
for selector in sale_price_selectors:
|
|
try:
|
|
elements = await page.query_selector_all(selector)
|
|
for elem in elements:
|
|
# Skip if element is strikethrough (old price)
|
|
parent_html = await elem.evaluate('el => el.parentElement.outerHTML')
|
|
if 'text-decoration: line-through' in parent_html or 'text-decoration-line: line-through' in parent_html:
|
|
continue
|
|
|
|
price_text = await elem.inner_text()
|
|
if price_text:
|
|
cleaned = price_text.strip().replace('€', '').replace('EUR', '').strip()
|
|
cleaned = cleaned.replace(' ', '').replace('\xa0', '').replace(',', '.')
|
|
|
|
match = re.search(r'(\d+\.?\d*)', cleaned)
|
|
if match:
|
|
try:
|
|
price_val = float(match.group(1))
|
|
if 0.01 < price_val < 100000:
|
|
logger.debug(f"Found sale price: {price_val}€ from {selector}")
|
|
return price_val
|
|
except ValueError:
|
|
continue
|
|
except Exception:
|
|
continue
|
|
|
|
|
|
# Fallback to standard price selectors
|
|
# Collect ALL valid prices and return the lowest (usually the promotional price)
|
|
all_prices = []
|
|
|
|
for selector in price_selectors:
|
|
try:
|
|
elements = await page.query_selector_all(selector)
|
|
for elem in elements:
|
|
# Skip if element is strikethrough (old price)
|
|
try:
|
|
parent_html = await elem.evaluate('el => el.parentElement.outerHTML')
|
|
# Check for strikethrough in parent
|
|
if 'text-decoration: line-through' in parent_html or 'text-decoration-line: line-through' in parent_html:
|
|
logger.debug(f" Skipping strikethrough price from {selector}")
|
|
continue
|
|
# Also check the element itself
|
|
elem_style = await elem.evaluate('el => window.getComputedStyle(el).textDecoration')
|
|
if 'line-through' in elem_style:
|
|
logger.debug(f" Skipping element with line-through style")
|
|
continue
|
|
except:
|
|
pass
|
|
|
|
price_text = await elem.inner_text()
|
|
if price_text:
|
|
cleaned = price_text.strip().replace('€', '').replace('EUR', '').strip()
|
|
cleaned = cleaned.replace(' ', '').replace('\xa0', '').replace(',', '.')
|
|
|
|
match = re.search(r'(\d+\.?\d*)', cleaned)
|
|
if match:
|
|
try:
|
|
price_val = float(match.group(1))
|
|
if 0.01 < price_val < 100000:
|
|
logger.debug(f" Found candidate price: {price_val}€ from {selector}")
|
|
all_prices.append(price_val)
|
|
except ValueError:
|
|
continue
|
|
except Exception as e:
|
|
logger.debug(f" Error with selector {selector}: {e}")
|
|
continue
|
|
|
|
# Return the LOWEST price found (promotional price is usually lower)
|
|
if all_prices:
|
|
lowest_price = min(all_prices)
|
|
logger.debug(f"Selected lowest price from {len(all_prices)} candidates: {lowest_price}€")
|
|
return lowest_price
|
|
|
|
logger.debug("No price found")
|
|
return None
|
|
|
|
@staticmethod
|
|
async def _extract_stock_status(page: Page) -> bool | None:
|
|
"""Extract stock status from product page"""
|
|
# Check for out of stock indicators
|
|
out_of_stock_texts = [
|
|
"rupture de stock",
|
|
"indisponible",
|
|
"out of stock",
|
|
"unavailable",
|
|
"épuisé",
|
|
"non disponible"
|
|
]
|
|
|
|
try:
|
|
page_text = await page.inner_text("body")
|
|
page_text_lower = page_text.lower()
|
|
|
|
for text in out_of_stock_texts:
|
|
if text in page_text_lower:
|
|
logger.debug(f"Out of stock detected: '{text}'")
|
|
return False
|
|
|
|
# If we find "add to cart" or similar, assume in stock
|
|
add_to_cart_texts = ["ajouter au panier", "add to cart", "acheter", "buy now"]
|
|
for text in add_to_cart_texts:
|
|
if text in page_text_lower:
|
|
logger.debug(f"In stock detected: '{text}'")
|
|
return True
|
|
|
|
except Exception as e:
|
|
logger.debug(f"Stock check error: {e}")
|
|
|
|
return None # Unknown
|
|
|
|
|
|
|
|
|
|
@classmethod
|
|
async def search_site_generator(cls, site_key: str, query: str) -> AsyncGenerator[SearchResult, None]:
|
|
"""Search a single site and yield results as they are scraped"""
|
|
config = SITE_CONFIGS.get(site_key)
|
|
if not config:
|
|
logger.error(f"Unknown site: {site_key}")
|
|
return
|
|
|
|
search_url = config["search_url"].format(query=quote_plus(query))
|
|
logger.info(f"Searching {config['name']} at {search_url}")
|
|
|
|
if not await cls._ensure_browser_connected():
|
|
logger.error("Failed to connect to browser")
|
|
return
|
|
|
|
try:
|
|
context = await cls._create_context(cls._browser)
|
|
|
|
try:
|
|
page = await context.new_page()
|
|
|
|
# Navigate
|
|
await page.goto(search_url, wait_until="domcontentloaded", timeout=30000)
|
|
|
|
# Wait for selector
|
|
wait_selector = config.get("wait_selector") or config.get("product_selector")
|
|
try:
|
|
if wait_selector:
|
|
await page.wait_for_selector(wait_selector, timeout=10000)
|
|
except Exception:
|
|
logger.warning(f"Timeout waiting for selector {wait_selector} on {site_key}")
|
|
|
|
# Get content
|
|
content = await page.content()
|
|
await page.close() # Close search page to free resources
|
|
|
|
# Parse results (Phase 1)
|
|
base_url = search_url.split("/search")[0]
|
|
if "amazon" in site_key:
|
|
base_url = "https://www.amazon.fr"
|
|
|
|
initial_results = cls._parse_results(content, site_key, base_url, query)
|
|
|
|
if not initial_results:
|
|
logger.warning(f"No results found for {site_key}")
|
|
# Dump HTML for debugging
|
|
import os
|
|
dump_dir = "/app/debug_dumps"
|
|
os.makedirs(dump_dir, exist_ok=True)
|
|
dump_path = f"{dump_dir}/{site_key.replace('.', '_')}_no_results.html"
|
|
with open(dump_path, 'w', encoding='utf-8') as f:
|
|
f.write(content)
|
|
logger.warning(f"HTML dumped to {dump_path} for inspection")
|
|
return
|
|
|
|
# Phase 2: Scrape details (Streaming)
|
|
semaphore = asyncio.Semaphore(3)
|
|
|
|
async def scrape_wrapper(res):
|
|
async with semaphore:
|
|
return await cls._scrape_item_details(res, context)
|
|
|
|
tasks = [scrape_wrapper(r) for r in initial_results]
|
|
|
|
for future in asyncio.as_completed(tasks):
|
|
enriched_res = await future
|
|
if enriched_res:
|
|
yield enriched_res
|
|
|
|
finally:
|
|
await context.close()
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error searching {site_key}: {e}")
|
|
|
|
@classmethod
|
|
async def search_all(cls, query: str) -> list[SearchResult]:
|
|
"""Search all configured sites"""
|
|
tasks = []
|
|
for site_key in SITE_CONFIGS.keys():
|
|
tasks.append(cls.search_site(site_key, query))
|
|
|
|
results_list = await asyncio.gather(*tasks)
|
|
all_results = []
|
|
for r in results_list:
|
|
all_results.extend(r)
|
|
return all_results
|
|
|
|
|
|
# ==========================================
|
|
# COMPATIBILITY LAYER FOR API ROUTERS
|
|
# ==========================================
|
|
|
|
async def search_products(
|
|
query: str,
|
|
db: Session,
|
|
site_ids: list[int] | None = None,
|
|
max_results: int | None = None,
|
|
) -> AsyncGenerator[SearchProgress, None]:
|
|
"""
|
|
Compatibility wrapper for search_products using improved service.
|
|
Yields SearchProgress events incrementally.
|
|
"""
|
|
# 1. Get sites to search
|
|
sites = db.query(SearchSite).order_by(SearchSite.priority).all()
|
|
if site_ids:
|
|
sites = [s for s in sites if s.id in site_ids]
|
|
|
|
active_sites = [s for s in sites if s.is_active]
|
|
|
|
# Initial event
|
|
yield SearchProgress(
|
|
status="searching",
|
|
total=len(active_sites),
|
|
completed=0,
|
|
message=f"Démarrage de la recherche sur {len(active_sites)} sites...",
|
|
results=[],
|
|
)
|
|
|
|
# 2. Map DB sites to Config keys with improved matching
|
|
site_keys = []
|
|
for site in active_sites:
|
|
matched_key = None
|
|
|
|
# Normalize domain for comparison (remove www., lowercase, etc.)
|
|
site_domain_normalized = site.domain.lower().replace("www.", "").replace("http://", "").replace("https://", "").strip("/")
|
|
|
|
# Try multiple matching strategies:
|
|
for key in SITE_CONFIGS.keys():
|
|
key_normalized = key.lower().replace("www.", "")
|
|
|
|
# 1. Exact match
|
|
if site_domain_normalized == key_normalized:
|
|
matched_key = key
|
|
logger.info(f"✅ Mapped {site.name} ({site.domain}) → {key} (exact match)")
|
|
break
|
|
|
|
# 2. Contains match (one in the other)
|
|
if key_normalized in site_domain_normalized or site_domain_normalized in key_normalized:
|
|
matched_key = key
|
|
logger.info(f"✅ Mapped {site.name} ({site.domain}) → {key} (contains)")
|
|
break
|
|
|
|
# 3. Normalize punctuation (. vs - vs nothing) and compare
|
|
# e.leclerc → eleclerc, e-leclerc.com → eleclecrcom
|
|
domain_no_punct = site_domain_normalized.replace("-", "").replace(".", "")
|
|
key_no_punct = key_normalized.replace("-", "").replace(".", "")
|
|
|
|
# Exact match without punctuation
|
|
if domain_no_punct == key_no_punct:
|
|
matched_key = key
|
|
logger.info(f"✅ Mapped {site.name} ({site.domain}) → {key} (normalized punctuation)")
|
|
break
|
|
|
|
# Contains match without punctuation (handles .com, .fr suffixes)
|
|
if len(domain_no_punct) > 3 and len(key_no_punct) > 3:
|
|
if domain_no_punct in key_no_punct or key_no_punct in domain_no_punct:
|
|
matched_key = key
|
|
logger.info(f"✅ Mapped {site.name} ({site.domain}) → {key} (normalized contains)")
|
|
break
|
|
|
|
if matched_key:
|
|
site_keys.append(matched_key)
|
|
else:
|
|
logger.warning(f"❌ No config found for {site.name} (domain: {site.domain}, normalized: {site_domain_normalized})")
|
|
|
|
# 3. Execute searches and stream results
|
|
generators = [ImprovedSearchService.search_site_generator(key, query) for key in site_keys]
|
|
|
|
queue = asyncio.Queue()
|
|
active_producers = len(generators)
|
|
|
|
# Limit concurrent sites
|
|
site_semaphore = asyncio.Semaphore(2)
|
|
|
|
async def producer(gen):
|
|
async with site_semaphore:
|
|
try:
|
|
async for item in gen:
|
|
await queue.put(item)
|
|
except Exception as e:
|
|
logger.error(f"Error in search producer: {e}")
|
|
finally:
|
|
await queue.put(None) # Sentinel
|
|
|
|
# Start producers
|
|
for gen in generators:
|
|
asyncio.create_task(producer(gen))
|
|
|
|
# Consumer loop
|
|
results_so_far = []
|
|
completed_sites = 0
|
|
|
|
while active_producers > 0:
|
|
item = await queue.get()
|
|
|
|
if item is None:
|
|
active_producers -= 1
|
|
completed_sites += 1
|
|
yield SearchProgress(
|
|
status="searching",
|
|
total=len(active_sites),
|
|
completed=completed_sites,
|
|
message=f"Recherche en cours... ({completed_sites}/{len(active_sites)} sites terminés)",
|
|
results=results_so_far,
|
|
)
|
|
else:
|
|
# Convert to SearchResultItem
|
|
api_item = SearchResultItem(
|
|
url=item.url,
|
|
title=item.title,
|
|
price=item.price,
|
|
currency=item.currency,
|
|
in_stock=item.in_stock,
|
|
site_name=item.source,
|
|
site_domain=item.source,
|
|
image_url=item.image_url,
|
|
)
|
|
results_so_far.append(api_item)
|
|
|
|
# Yield update with new result
|
|
yield SearchProgress(
|
|
status="searching",
|
|
total=len(active_sites),
|
|
completed=completed_sites,
|
|
message=f"Trouvé: {item.title[:30]}...",
|
|
results=results_so_far,
|
|
)
|
|
|
|
# Final event
|
|
yield SearchProgress(
|
|
status="completed",
|
|
total=len(active_sites),
|
|
completed=len(active_sites),
|
|
message=f"Terminé. {len(results_so_far)} résultats trouvés.",
|
|
results=results_so_far,
|
|
)
|
|
|
|
|
|
# Global instance
|
|
improved_search_service = ImprovedSearchService()
|