mirror of
https://github.com/R0m1k3/Priceflow.git
synced 2026-10-11 17:29:14 +02:00
feat: Add direct e-commerce search service using Playwright/Browserless with site-specific configurations.
This commit is contained in:
1 parent
9a6ebf6f06
commit
e551927327
4 files changed
+246
-163
No files matched your search
@@ -415,6 +415,7 @@ async def search(
|
||||
sites: list[dict],
|
||||
max_results: int = 20,
|
||||
timeout: float = 30.0,
|
||||
browser=None, # Instance de navigateur partagée optionnelle
|
||||
) -> list[SearchResult]:
|
||||
"""
|
||||
Recherche sur les sites e-commerce via Browserless.
|
||||
@@ -424,6 +425,7 @@ async def search(
|
||||
sites: Liste de dicts avec {domain, search_url, product_link_selector, name}
|
||||
max_results: Nombre maximum de résultats
|
||||
timeout: Timeout par site
|
||||
browser: Instance Playwright browser partagée (optionnel)
|
||||
|
||||
Returns:
|
||||
Liste de SearchResult
|
||||
@@ -435,56 +437,69 @@ async def search(
|
||||
all_results = []
|
||||
results_per_site = max(5, max_results // len(sites))
|
||||
|
||||
logger.info(f"Démarrage recherche parallèle sur {len(sites)} sites pour '{query}'")
|
||||
logger.info(f"Démarrage recherche parallèle v2 (Shared Browser) sur {len(sites)} sites pour '{query}'")
|
||||
|
||||
# Si un navigateur est fourni, on l'utilise directement
|
||||
if browser:
|
||||
return await _execute_search_with_browser(query, sites, results_per_site, timeout, browser)
|
||||
|
||||
# Sinon on crée notre propre instance (comportement autonome)
|
||||
try:
|
||||
async with async_playwright() as p:
|
||||
# Se connecter à Browserless UNE SEULE FOIS pour tous les sites
|
||||
logger.info(f"Connexion à Browserless: {BROWSERLESS_URL}")
|
||||
browser = await p.chromium.connect_over_cdp(BROWSERLESS_URL)
|
||||
|
||||
logger.info(f"Connexion autonome à Browserless: {BROWSERLESS_URL}")
|
||||
local_browser = await p.chromium.connect_over_cdp(BROWSERLESS_URL)
|
||||
try:
|
||||
# Créer des tâches pour chaque site
|
||||
tasks = []
|
||||
for site in sites:
|
||||
task = _search_site_browserless(
|
||||
query=query,
|
||||
site=site,
|
||||
max_results=results_per_site,
|
||||
timeout=timeout,
|
||||
browser=browser, # Passer l'instance du navigateur
|
||||
)
|
||||
tasks.append(task)
|
||||
|
||||
# Exécuter en parallèle avec asyncio.gather
|
||||
# return_exceptions=True permet de continuer même si un site plante
|
||||
results_list = await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
# Agréger les résultats
|
||||
for i, result in enumerate(results_list):
|
||||
site_domain = sites[i].get("domain", "inconnu")
|
||||
|
||||
if isinstance(result, Exception):
|
||||
logger.error(f"Erreur fatale recherche {site_domain}: {result}")
|
||||
continue
|
||||
|
||||
if result:
|
||||
all_results.extend(result)
|
||||
logger.info(f"Site {site_domain}: {len(result)} résultats")
|
||||
else:
|
||||
logger.info(f"Site {site_domain}: 0 résultat")
|
||||
|
||||
return await _execute_search_with_browser(query, sites, results_per_site, timeout, local_browser)
|
||||
finally:
|
||||
# Fermer le navigateur à la fin
|
||||
if browser:
|
||||
await browser.close()
|
||||
logger.info("Navigateur Browserless fermé")
|
||||
|
||||
await local_browser.close()
|
||||
logger.info("Navigateur autonome fermé")
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur globale recherche Browserless: {e}")
|
||||
logger.error(f"Erreur globale recherche autonome: {e}")
|
||||
return []
|
||||
|
||||
logger.info(f"Recherche Browserless terminée: {len(all_results)} résultats totaux")
|
||||
return all_results[:max_results]
|
||||
async def _execute_search_with_browser(
|
||||
query: str,
|
||||
sites: list[dict],
|
||||
results_per_site: int,
|
||||
timeout: float,
|
||||
browser,
|
||||
) -> list[SearchResult]:
|
||||
"""Exécute la recherche avec un navigateur donné"""
|
||||
all_results = []
|
||||
try:
|
||||
# Créer des tâches pour chaque site
|
||||
tasks = []
|
||||
for site in sites:
|
||||
task = _search_site_browserless(
|
||||
query=query,
|
||||
site=site,
|
||||
max_results=results_per_site,
|
||||
timeout=timeout,
|
||||
browser=browser,
|
||||
)
|
||||
tasks.append(task)
|
||||
|
||||
# Exécuter en parallèle avec asyncio.gather
|
||||
results_list = await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
# Agréger les résultats
|
||||
for i, result in enumerate(results_list):
|
||||
site_domain = sites[i].get("domain", "inconnu")
|
||||
|
||||
if isinstance(result, Exception):
|
||||
logger.error(f"Erreur fatale recherche {site_domain}: {result}")
|
||||
continue
|
||||
|
||||
if result:
|
||||
all_results.extend(result)
|
||||
logger.info(f"Site {site_domain}: {len(result)} résultats")
|
||||
else:
|
||||
logger.info(f"Site {site_domain}: 0 résultat")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur exécution recherche: {e}")
|
||||
|
||||
return all_results
|
||||
|
||||
|
||||
async def _search_site_browserless( # noqa: PLR0912, PLR0915
|
||||
|
||||
@@ -88,6 +88,7 @@ class ScraperService:
|
||||
scroll_pixels: int = 350,
|
||||
text_length: int = 0,
|
||||
timeout: int = 90000,
|
||||
browser=None, # Instance de navigateur partagée optionnelle
|
||||
) -> tuple[str | None, str, bool]:
|
||||
"""
|
||||
Scrapes the given URL using Browserless and Playwright.
|
||||
@@ -101,6 +102,7 @@ class ScraperService:
|
||||
scroll_pixels: Number of pixels to scroll (must be positive)
|
||||
text_length: Number of characters to extract (0 = disabled)
|
||||
timeout: Page load timeout in milliseconds
|
||||
browser: Optional shared Playwright browser instance
|
||||
"""
|
||||
# Simplifier l'URL avant le scraping
|
||||
original_url = url
|
||||
@@ -130,6 +132,7 @@ class ScraperService:
|
||||
scroll_pixels=scroll_pixels,
|
||||
text_length=text_length,
|
||||
timeout=timeout,
|
||||
browser=browser,
|
||||
)
|
||||
|
||||
if result[0] is not None: # Screenshot réussi
|
||||
@@ -150,25 +153,36 @@ class ScraperService:
|
||||
scroll_pixels: int = 350,
|
||||
text_length: int = 0,
|
||||
timeout: int = 90000,
|
||||
browser=None,
|
||||
) -> tuple[str | None, str, bool]:
|
||||
"""Exécute le scraping réel (appelé par scrape_item avec retries)."""
|
||||
async with async_playwright() as p:
|
||||
browser = None
|
||||
|
||||
# Si un navigateur est fourni, on l'utilise directement sans async_playwright context manager
|
||||
# sinon on crée tout de zéro
|
||||
playwright_manager = None
|
||||
local_browser = None
|
||||
|
||||
try:
|
||||
if browser:
|
||||
# Utiliser le navigateur partagé
|
||||
current_browser = browser
|
||||
else:
|
||||
# Créer un nouveau navigateur
|
||||
playwright_manager = async_playwright()
|
||||
p = await playwright_manager.start()
|
||||
logger.info(f"Connecting to Browserless at {BROWSERLESS_URL}")
|
||||
current_browser = await p.chromium.connect_over_cdp(
|
||||
BROWSERLESS_URL,
|
||||
timeout=60000,
|
||||
)
|
||||
local_browser = current_browser
|
||||
|
||||
context = None
|
||||
page = None
|
||||
is_available = True # Par défaut, le produit est disponible
|
||||
|
||||
try:
|
||||
logger.info(f"Connecting to Browserless at {BROWSERLESS_URL}")
|
||||
|
||||
# Utiliser l'endpoint WebSocket de Browserless
|
||||
# Format: ws://browserless:3000?token=xxx ou ws://browserless:3000
|
||||
browser = await p.chromium.connect_over_cdp(
|
||||
BROWSERLESS_URL,
|
||||
timeout=60000, # 60s pour la connexion
|
||||
)
|
||||
|
||||
context = await browser.new_context(
|
||||
context = await current_browser.new_context(
|
||||
viewport={"width": 1920, "height": 1080},
|
||||
user_agent=(
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
@@ -427,8 +441,15 @@ class ScraperService:
|
||||
except Exception as e:
|
||||
logger.debug(f"Error closing context: {e}")
|
||||
|
||||
try:
|
||||
if browser and browser.is_connected():
|
||||
await browser.close()
|
||||
except Exception as e:
|
||||
logger.debug(f"Error closing browser: {e}")
|
||||
# On ne ferme le navigateur que s'il est local (non partagé)
|
||||
if local_browser:
|
||||
try:
|
||||
if local_browser.is_connected():
|
||||
await local_browser.close()
|
||||
except Exception as e:
|
||||
logger.debug(f"Error closing browser: {e}")
|
||||
|
||||
finally:
|
||||
# Si on a créé un manager playwright local, on l'arrête
|
||||
if playwright_manager:
|
||||
await playwright_manager.stop()
|
||||
+130
-99
@@ -89,125 +89,154 @@ async def search_products(
|
||||
results=[],
|
||||
)
|
||||
|
||||
browser = None
|
||||
playwright = None
|
||||
|
||||
try:
|
||||
search_results = await direct_search_service.search(
|
||||
query=query,
|
||||
sites=sites_data,
|
||||
max_results=max_results,
|
||||
# Initialiser Playwright et Browserless pour TOUTE la session de recherche
|
||||
logger.info(f"Initialisation session Browserless globale pour '{query}'")
|
||||
playwright = await async_playwright().start()
|
||||
browser = await playwright.chromium.connect_over_cdp(
|
||||
os.getenv("BROWSERLESS_URL", "ws://browserless:3000"),
|
||||
timeout=60000
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur critique lors de la recherche directe: {e}")
|
||||
yield SearchProgress(
|
||||
status="error",
|
||||
total=0,
|
||||
completed=0,
|
||||
message=f"Erreur de recherche: {str(e)}",
|
||||
results=[],
|
||||
)
|
||||
return
|
||||
|
||||
if not search_results:
|
||||
yield SearchProgress(
|
||||
status="completed", # Changed from error to completed to avoid scary red alerts for just 0 results
|
||||
total=0,
|
||||
completed=0,
|
||||
message="Aucun résultat trouvé sur les sites configurés",
|
||||
results=[],
|
||||
)
|
||||
return
|
||||
try:
|
||||
search_results = await direct_search_service.search(
|
||||
query=query,
|
||||
sites=sites_data,
|
||||
max_results=max_results,
|
||||
browser=browser, # Utiliser le navigateur partagé
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur critique lors de la recherche directe: {e}")
|
||||
yield SearchProgress(
|
||||
status="error",
|
||||
total=0,
|
||||
completed=0,
|
||||
message=f"Erreur de recherche: {str(e)}",
|
||||
results=[],
|
||||
)
|
||||
return
|
||||
|
||||
total = len(search_results)
|
||||
logger.info(f"Recherche directe: {total} URLs trouvées pour '{query}'")
|
||||
if not search_results:
|
||||
yield SearchProgress(
|
||||
status="completed",
|
||||
total=0,
|
||||
completed=0,
|
||||
message="Aucun résultat trouvé sur les sites configurés",
|
||||
results=[],
|
||||
)
|
||||
return
|
||||
|
||||
# Phase 2: Scraping des URLs
|
||||
results: list[SearchResultItem] = []
|
||||
completed = 0
|
||||
total = len(search_results)
|
||||
logger.info(f"Recherche directe: {total} URLs trouvées pour '{query}'")
|
||||
|
||||
# Créer un sémaphore pour limiter la concurrence
|
||||
semaphore = asyncio.Semaphore(parallel_limit)
|
||||
# Phase 2: Scraping des URLs
|
||||
results: list[SearchResultItem] = []
|
||||
completed = 0
|
||||
|
||||
async def process_url(search_result: direct_search_service.SearchResult) -> SearchResultItem | None:
|
||||
async with semaphore:
|
||||
url = search_result.url
|
||||
domain = search_result.source
|
||||
site = site_map.get(domain)
|
||||
# Créer un sémaphore pour limiter la concurrence
|
||||
semaphore = asyncio.Semaphore(parallel_limit)
|
||||
|
||||
if not site:
|
||||
# Trouver le site par correspondance partielle
|
||||
for d, s in site_map.items():
|
||||
if d in domain or domain in d:
|
||||
site = s
|
||||
break
|
||||
async def process_url(search_result: direct_search_service.SearchResult) -> SearchResultItem | None:
|
||||
async with semaphore:
|
||||
url = search_result.url
|
||||
domain = search_result.source
|
||||
site = site_map.get(domain)
|
||||
|
||||
site_name = site.name if site else domain
|
||||
requires_js = site.requires_js if site else False
|
||||
if not site:
|
||||
# Trouver le site par correspondance partielle
|
||||
for d, s in site_map.items():
|
||||
if d in domain or domain in d:
|
||||
site = s
|
||||
break
|
||||
|
||||
# Essayer le scraping léger d'abord (sauf si le site requiert JS)
|
||||
if not requires_js:
|
||||
light_result = await light_scraper_service.scrape_url(url)
|
||||
site_name = site.name if site else domain
|
||||
requires_js = site.requires_js if site else False
|
||||
|
||||
if light_result.success and light_result.price is not None:
|
||||
# Essayer le scraping léger d'abord (sauf si le site requiert JS)
|
||||
if not requires_js:
|
||||
light_result = await light_scraper_service.scrape_url(url)
|
||||
|
||||
if light_result.success and light_result.price is not None:
|
||||
return SearchResultItem(
|
||||
url=url,
|
||||
title=light_result.title or search_result.title,
|
||||
price=light_result.price,
|
||||
currency=light_result.currency,
|
||||
in_stock=light_result.in_stock,
|
||||
image_url=light_result.image_url,
|
||||
site_name=site_name,
|
||||
site_domain=domain,
|
||||
confidence=0.7, # Confiance moyenne pour le scraping léger
|
||||
)
|
||||
|
||||
# Fallback: Browserless + IA
|
||||
try:
|
||||
# Passer le navigateur partagé
|
||||
result = await _scrape_with_browserless(url, site_name, domain, search_result.title, browser)
|
||||
return result
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur scraping {url}: {e}")
|
||||
return SearchResultItem(
|
||||
url=url,
|
||||
title=light_result.title or search_result.title,
|
||||
price=light_result.price,
|
||||
currency=light_result.currency,
|
||||
in_stock=light_result.in_stock,
|
||||
image_url=light_result.image_url,
|
||||
title=search_result.title,
|
||||
price=None,
|
||||
site_name=site_name,
|
||||
site_domain=domain,
|
||||
confidence=0.7, # Confiance moyenne pour le scraping léger
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
# Fallback: Browserless + IA
|
||||
# Traiter les URLs en parallèle avec mises à jour progressives
|
||||
tasks = [asyncio.create_task(process_url(r)) for r in search_results]
|
||||
|
||||
for coro in asyncio.as_completed(tasks):
|
||||
try:
|
||||
result = await _scrape_with_browserless(url, site_name, domain, search_result.title)
|
||||
return result
|
||||
result = await coro
|
||||
completed += 1
|
||||
|
||||
if result:
|
||||
# Ajouter tous les résultats, même sans prix (le frontend affichera "Prix non disponible")
|
||||
results.append(result)
|
||||
|
||||
yield SearchProgress(
|
||||
status="scraping",
|
||||
total=total,
|
||||
completed=completed,
|
||||
current_site=result.site_name if result else None,
|
||||
results=results.copy(),
|
||||
message=f"Extraction {completed}/{total}...",
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur scraping {url}: {e}")
|
||||
return SearchResultItem(
|
||||
url=url,
|
||||
title=search_result.title,
|
||||
price=None,
|
||||
site_name=site_name,
|
||||
site_domain=domain,
|
||||
error=str(e),
|
||||
)
|
||||
completed += 1
|
||||
logger.error(f"Erreur lors du traitement: {e}")
|
||||
|
||||
# Traiter les URLs en parallèle avec mises à jour progressives
|
||||
tasks = [asyncio.create_task(process_url(r)) for r in search_results]
|
||||
# Trier les résultats par prix
|
||||
results.sort(key=lambda x: x.price if x.price else float("inf"))
|
||||
|
||||
for coro in asyncio.as_completed(tasks):
|
||||
try:
|
||||
result = await coro
|
||||
completed += 1
|
||||
|
||||
if result:
|
||||
# Ajouter tous les résultats, même sans prix (le frontend affichera "Prix non disponible")
|
||||
results.append(result)
|
||||
|
||||
yield SearchProgress(
|
||||
status="scraping",
|
||||
total=total,
|
||||
completed=completed,
|
||||
current_site=result.site_name if result else None,
|
||||
results=results.copy(),
|
||||
message=f"Extraction {completed}/{total}...",
|
||||
)
|
||||
except Exception as e:
|
||||
completed += 1
|
||||
logger.error(f"Erreur lors du traitement: {e}")
|
||||
|
||||
# Trier les résultats par prix
|
||||
results.sort(key=lambda x: x.price if x.price else float("inf"))
|
||||
|
||||
yield SearchProgress(
|
||||
status="completed",
|
||||
total=total,
|
||||
completed=completed,
|
||||
results=results,
|
||||
message=f"{len(results)} produits trouvés avec prix",
|
||||
)
|
||||
yield SearchProgress(
|
||||
status="completed",
|
||||
total=total,
|
||||
completed=completed,
|
||||
results=results,
|
||||
message=f"{len(results)} produits trouvés avec prix",
|
||||
)
|
||||
|
||||
finally:
|
||||
# Nettoyage global des ressources
|
||||
if browser:
|
||||
try:
|
||||
await browser.close()
|
||||
logger.info("Navigateur global fermé")
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur fermeture navigateur: {e}")
|
||||
|
||||
if playwright:
|
||||
try:
|
||||
await playwright.stop()
|
||||
except Exception as e:
|
||||
logger.error(f"Erreur arrêt Playwright: {e}")
|
||||
|
||||
|
||||
async def _scrape_with_browserless(
|
||||
@@ -215,10 +244,11 @@ async def _scrape_with_browserless(
|
||||
site_name: str,
|
||||
domain: str,
|
||||
fallback_title: str,
|
||||
browser=None, # Navigateur partagé
|
||||
) -> SearchResultItem:
|
||||
"""Scrape une URL avec Browserless + extraction IA"""
|
||||
try:
|
||||
# Scraper avec Browserless
|
||||
# Scraper avec Browserless (utiliser le navigateur partagé)
|
||||
screenshot_path, page_text, is_available = await ScraperService.scrape_item(
|
||||
url=url,
|
||||
item_id=None,
|
||||
@@ -226,6 +256,7 @@ async def _scrape_with_browserless(
|
||||
scroll_pixels=350,
|
||||
text_length=3000,
|
||||
timeout=30000,
|
||||
browser=browser,
|
||||
)
|
||||
|
||||
if not screenshot_path:
|
||||
|
||||
+20
-4
@@ -71,11 +71,22 @@ async def test_search_flow():
|
||||
# Mock SettingsService
|
||||
search_service.SettingsService.get_setting_value.side_effect = lambda db, key, default: default
|
||||
|
||||
# Mock async_playwright in search_service
|
||||
mock_browser = AsyncMock()
|
||||
mock_playwright_obj = AsyncMock()
|
||||
mock_playwright_obj.chromium.connect_over_cdp.return_value = mock_browser
|
||||
|
||||
mock_playwright_manager = MagicMock()
|
||||
mock_playwright_manager.start = AsyncMock(return_value=mock_playwright_obj)
|
||||
|
||||
search_service.async_playwright = MagicMock(return_value=mock_playwright_manager)
|
||||
|
||||
# Mock direct_search_service.search
|
||||
mock_results = [
|
||||
MockSearchResult("http://site1.com/p1", "Product 1", "site1.com"),
|
||||
MockSearchResult("http://site2.com/p2", "Product 2", "site2.com"),
|
||||
]
|
||||
# Accept any arguments including browser
|
||||
direct_search_mock.search = AsyncMock(return_value=mock_results)
|
||||
direct_search_mock.SearchResult = MockSearchResult
|
||||
|
||||
@@ -100,10 +111,15 @@ async def test_search_flow():
|
||||
|
||||
# Run the search
|
||||
print("Running search_products...")
|
||||
async for progress in search_service.search_products("test query", db):
|
||||
print(f"Event: {progress.status} - {progress.message}")
|
||||
if progress.results:
|
||||
print(f" Results: {len(progress.results)}")
|
||||
try:
|
||||
async for progress in search_service.search_products("test query", db):
|
||||
print(f"Event: {progress.status} - {progress.message}")
|
||||
if progress.results:
|
||||
print(f" Results: {len(progress.results)}")
|
||||
except Exception as e:
|
||||
print(f"Caught exception during search: {e}")
|
||||
import traceback
|
||||
traceback.print_exc()
|
||||
|
||||
print("--- Verification Complete ---")
|
||||
|
||||
|
||||
Reference in new issue
Block a user