| |
| 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 |
|
|
| |
| import os |
|
|
| |
|
|
| 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 |
|
|
| |
| session = requests.Session() |
| lock = threading.Lock() |
| GLOBAL_LINK_COUNTER = 0 |
|
|
|
|
| |
| 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): |
| 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() |
|
|
| |
| 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." |
| ) |
|
|
| |
| 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" |
| |
| 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": [], |
| } |
| ) |
|
|
| |
| 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 |
|
|