# scraper.py import json import time import io import threading import http.client import requests from bs4 import BeautifulSoup from PyPDF2 import PdfReader from concurrent.futures import ThreadPoolExecutor # -------- CONFIG ---------- import os # SERPER_API_KEY = os.getenv("SERPER_API_KEY", "") SERPER_API_KEYS = [ os.getenv("SERPER_API_KEY"), os.getenv("SERPER_API_KEY_1"), os.getenv("SERPER_API_KEY_2"), os.getenv("SERPER_API_KEY_3"), os.getenv("SERPER_API_KEY_335"), os.getenv("SERPER_API_KEY_zoe"), os.getenv("SERPER_API_KEY_38732"), os.getenv("SERPER_API_KEY_849"), ] SERPER_API_KEYS = [k for k in SERPER_API_KEYS if k] CURRENT_SERPER_INDEX = 0 SERPER_LOCK = threading.Lock() class SerperExhaustedError(Exception): pass def rotate_serper_key(): global CURRENT_SERPER_INDEX with SERPER_LOCK: CURRENT_SERPER_INDEX += 1 if CURRENT_SERPER_INDEX >= len(SERPER_API_KEYS): raise SerperExhaustedError( "❌ SERPER APIs exhausted. Please contact the developer." ) print(f"🔁 Switching to SERPER API #{CURRENT_SERPER_INDEX}") BAD_KEYWORDS = [ "facebook", "youtube", "tiktok", "instagram", "linkedin", "x.com", "twitter", "pinterest", "snapchat", "blog", "course", ] FAST_MODE = True MAX_WORKERS = 15 SCRAPE_CHAR_LIMIT = 2000 REQUEST_TIMEOUT = 8 PDF_READ_LIMIT = 2000 SERP_MAX_RESULTS = 3 RETRIES = 1 # -------- shared resources ---------- session = requests.Session() lock = threading.Lock() GLOBAL_LINK_COUNTER = 0 # -------- helpers ---------- def is_good_source(url: str) -> bool: if not url: return False u = url.lower() return not any(bad in u for bad in BAD_KEYWORDS) def is_valid_url(url: str) -> bool: if FAST_MODE: return url.startswith("http") try: r = session.head(url, timeout=3, allow_redirects=True) return r.status_code in (200, 301, 302) except: return False def extract_pdf_content(url: str) -> str: try: r = session.get(url, timeout=REQUEST_TIMEOUT) pdf_bytes = r.content reader = PdfReader(io.BytesIO(pdf_bytes)) text = "" for page in reader.pages: text += page.extract_text() or "" if len(text) >= PDF_READ_LIMIT: break return text.strip()[:SCRAPE_CHAR_LIMIT] except: return "" def scrape_page(url, retries=RETRIES, limit=SCRAPE_CHAR_LIMIT): for _ in range(retries): try: if url.lower().endswith(".pdf"): return extract_pdf_content(url)[:limit] r = session.get(url, timeout=REQUEST_TIMEOUT) soup = BeautifulSoup(r.text, "html.parser") text = soup.get_text(separator="\n", strip=True) return text[:limit] except: time.sleep(0.3) return "" # def search_serper(query, max_results=SERP_MAX_RESULTS): # if not SERPER_API_KEY: # # في البروداكشن خليه raise Exception أحسن # return [] # try: # conn = http.client.HTTPSConnection("google.serper.dev") # payload = json.dumps({"q": query, "page": 1}) # headers = {"X-API-KEY": SERPER_API_KEY, "Content-Type": "application/json"} # conn.request("POST", "/search", payload, headers) # res = conn.getresponse() # data = res.read() # response_json = json.loads(data.decode("utf-8")) # except: # return [] # results = [] # seen = set() # if "organic" in response_json: # for r in response_json["organic"]: # url = r.get("link") # if not url or url in seen: # continue # if not is_good_source(url): # continue # if not is_valid_url(url): # continue # seen.add(url) # position = r.get("position", 100) # score = round(1 / max(position, 1), 4) # results.append( # { # "url": url, # "title": r.get("title"), # "content": r.get("snippet"), # "score": score, # } # ) # if len(results) >= max_results: # break # return results def search_serper(query, max_results=SERP_MAX_RESULTS): global CURRENT_SERPER_INDEX attempts = 0 while attempts < len(SERPER_API_KEYS): api_key = SERPER_API_KEYS[CURRENT_SERPER_INDEX] try: conn = http.client.HTTPSConnection( "google.serper.dev", timeout=REQUEST_TIMEOUT ) payload = json.dumps({"q": query, "page": 1}) headers = { "X-API-KEY": api_key, "Content-Type": "application/json", } conn.request("POST", "/search", payload, headers) res = conn.getresponse() status = res.status data = res.read() # ❌ API Key خلصت أو مش صالحة if status in (401, 403): rotate_serper_key() attempts += 1 continue # ❌ أي خطأ تاني if status != 200: return [] response_json = json.loads(data.decode("utf-8")) break except SerperExhaustedError: raise except Exception: rotate_serper_key() attempts += 1 else: raise SerperExhaustedError( "❌ SERPER APIs exhausted. Please contact the developer." ) # -------- normal parsing ---------- results = [] seen = set() for r in response_json.get("organic", []): url = r.get("link") if not url or url in seen: continue if not is_good_source(url): continue if not is_valid_url(url): continue seen.add(url) position = r.get("position", 100) score = round(1 / max(position, 1), 4) results.append( { "url": url, "title": r.get("title"), "content": r.get("snippet"), "score": score, } ) if len(results) >= max_results: break return results def process_link(link_data): url = link_data["url"] idx = link_data.get("counter", -1) print(f"🕷️ [SCRAPE #{idx}] START {url}") t0 = time.perf_counter() page_text = scrape_page(url) duration = round(time.perf_counter() - t0, 2) print(f"🕷️ [SCRAPE #{idx}] DONE in {duration}s — chars: {len(page_text)}") link_data["scraped_content"] = page_text if page_text else "ERROR: EMPTY" # summary_result و raw_result في الكيس دي نفس الحاجة return link_data, dict(link_data) def safe_append(summary_result, raw_result, final_output_summary, final_output_raw): u = summary_result["unit_idx"] t = summary_result["topic_idx"] s = summary_result["sub_idx"] with lock: final_output_summary["units"][u]["topics"][t]["subtopics"][s]["results"].append( summary_result ) final_output_raw["units"][u]["topics"][t]["subtopics"][s]["results"].append( raw_result ) def scrape_course(course: dict): """ course هنا هو نفس الستركشر اللي كان جاي من JSON file ويرجع dicts: (final_output_summary, final_output_raw) """ global GLOBAL_LINK_COUNTER GLOBAL_LINK_COUNTER = 0 final_output_summary = { "course_name": course["course_name"], "course_audience": course["course_audience"], "units": [], } final_output_raw = { "course_name": course["course_name"], "course_audience": course["course_audience"], "units": [], } all_tasks = [] for unit_idx, unit in enumerate(course["expanded_units"]): unit_obj_sum = { "unit_name": unit["unit_name"], "outcome": unit["outcome"], "topics": [], } unit_obj_raw = { "unit_name": unit["unit_name"], "outcome": unit["outcome"], "topics": [], } for topic_idx, topic in enumerate(unit["topics"]): topic_sum = {"title": topic["topic_name"], "subtopics": []} topic_raw = {"title": topic["topic_name"], "subtopics": []} for sub_idx, sub in enumerate(topic["subtopics"]): topic_sum["subtopics"].append( { "title": sub["title"], "description": sub.get("description", ""), "queries": sub["queries"], "results": [], } ) topic_raw["subtopics"].append( { "title": sub["title"], "description": sub.get("description", ""), "queries": sub["queries"], "results": [], } ) # لكل query في الـ subtopic for query_text in sub["queries"]: print( f"\n🔍 [SEARCH] U{unit_idx} T{topic_idx} S{sub_idx} — Query: {query_text}" ) results = search_serper(query_text, max_results=SERP_MAX_RESULTS) for r in results: with lock: GLOBAL_LINK_COUNTER += 1 r["counter"] = GLOBAL_LINK_COUNTER r["query"] = query_text r["unit_idx"] = unit_idx r["topic_idx"] = topic_idx r["sub_idx"] = sub_idx print(f" ➕ [ADD LINK #{r['counter']}] {r['url']}") all_tasks.append(r) unit_obj_sum["topics"].append(topic_sum) unit_obj_raw["topics"].append(topic_raw) final_output_summary["units"].append(unit_obj_sum) final_output_raw["units"].append(unit_obj_raw) print( f"\n🚀 START PARALLEL SCRAPING — Total Links: {len(all_tasks)} — Workers: {MAX_WORKERS}\n" ) def worker(link_data): summary_result, raw_result = process_link(link_data) safe_append(summary_result, raw_result, final_output_summary, final_output_raw) with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: list(executor.map(worker, all_tasks)) return final_output_summary, final_output_raw