import asyncio import logging from database import DatabaseManager from scrapers import ScraperEngine from validators import validate_and_enrich_lead # Configure logging logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' ) logger = logging.getLogger(__name__) async def main(): logger.info("Initializing Robotics Client Lead Engine...") db_manager = DatabaseManager() scraper_engine = ScraperEngine() query = "robotics automation companies" logger.info(f"Starting asynchronous scraping for query: '{query}'") # Run scrapers concurrently raw_leads = await scraper_engine.run_all_scrapers(query) logger.info(f"Total raw leads gathered: {len(raw_leads)}") processed_count = 0 dropped_count = 0 for raw_lead in raw_leads: # Validate and apply structural gate / email fallback validated_lead = validate_and_enrich_lead(raw_lead) if validated_lead: # Insert into database inserted_id = db_manager.insert_lead(validated_lead) if inserted_id: processed_count += 1 logger.info(f"Successfully ingested lead: {validated_lead['company_name']} (ID: {inserted_id})") else: logger.warning(f"Failed to insert lead: {validated_lead.get('company_name')}") else: dropped_count += 1 logger.debug(f"Lead dropped due to validation failure: {raw_lead.get('company_name', 'Unknown')}") logger.info("--- Ingestion Pipeline Complete ---") logger.info(f"Successfully processed: {processed_count}") logger.info(f"Dropped: {dropped_count}") if __name__ == "__main__": asyncio.run(main())