Spaces:
Paused
Paused
| import asyncio | |
| import json | |
| import os | |
| import random | |
| import tempfile | |
| import urllib.parse | |
| import time | |
| import logging | |
| from datetime import datetime, timezone | |
| import concurrent.futures | |
| import re | |
| from typing import Optional | |
| from dotenv import load_dotenv | |
| load_dotenv() # loads .env from project root → SERPER_API_KEY etc. | |
| from openai import OpenAI | |
| import httpx | |
| from fastapi import FastAPI, HTTPException, Query, BackgroundTasks | |
| from pydantic import BaseModel | |
| from starlette.responses import StreamingResponse | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from bs4 import BeautifulSoup | |
| from school_classifier import classify_job, SCHOOLS, PROGRAM_KEYWORDS | |
| from locations import INDIAN_METROS | |
| from database import save_scraped_jobs, get_jobs_by_timeframe, get_jobs_in_timeframe, cleanup_old_jobs, get_binned_jobs, update_company_ratings, get_rated_companies, _connect | |
| # Primary engine: Selenium with undetected-chromedriver | |
| import undetected_chromedriver as uc | |
| from selenium.webdriver.support.ui import WebDriverWait | |
| from selenium.webdriver.support import expected_conditions as EC | |
| from selenium.webdriver.common.by import By | |
| # Fallback engine: Playwright with stealth | |
| from playwright.async_api import async_playwright | |
| from playwright_stealth import Stealth | |
| from apscheduler.schedulers.asyncio import AsyncIOScheduler | |
| from apscheduler.triggers.interval import IntervalTrigger | |
| scheduler = AsyncIOScheduler() | |
| # ── Logging setup ───────────────────────────────────────────── | |
| logging.basicConfig( | |
| level=logging.INFO, | |
| format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", | |
| datefmt="%Y-%m-%d %H:%M:%S", | |
| ) | |
| logger = logging.getLogger("internscrapper") | |
| logger.setLevel(logging.INFO) | |
| app = FastAPI(title="LinkedIn Public Job Scraper") | |
| app.add_middleware( | |
| CORSMiddleware, | |
| allow_origins=["*"], | |
| allow_credentials=False, | |
| allow_methods=["*"], | |
| allow_headers=["*"], | |
| ) | |
| # ── Shared resources ────────────────────────────────────────── | |
| class ScrapeSession: | |
| def __init__(self): | |
| self.history: list[str] = [] | |
| self.queues: list[asyncio.Queue] = [] | |
| self.is_active = False | |
| self.task: Optional[asyncio.Task] = None | |
| def cancel(self): | |
| self.is_active = False | |
| if self.task and not self.task.done(): | |
| self.task.cancel() | |
| global_scrape = ScrapeSession() | |
| executor = concurrent.futures.ThreadPoolExecutor(max_workers=6) | |
| _scrape_lock = asyncio.Lock() | |
| # ── Auto-scrape state (updated live so the admin UI can poll progress) ──── | |
| _auto_scrape_state: dict = { | |
| "is_running": False, | |
| "current_school": None, | |
| "completed": [], | |
| "failed": [], | |
| "started_at": None, | |
| "finished_at": None, | |
| "cancel_requested": False, | |
| } | |
| _auto_scrape_task: Optional[asyncio.Task] = None | |
| # ── Auto-scrape history (in-memory log of all sweep runs) ───── | |
| _auto_scrape_history: list[dict] = [] | |
| # HuggingFace Spaces sets PORT=7860 in its environment. | |
| # Locally, uvicorn defaults to 8000 unless overridden. | |
| _SELF_PORT = int(os.environ.get("PORT", "8000")) | |
| _SELF_BASE_URL = f"http://localhost:{_SELF_PORT}" | |
| se_driver = None # Selenium undetected-chromedriver | |
| pw_browser = None # Playwright browser | |
| pw_stealth_ctx = None # Playwright stealth context manager | |
| # Cookie warm-up cache — browser-derived cookies with a TTL. | |
| _warm_cookies: dict = {} | |
| _warm_cookies_ts: float = 0.0 | |
| _COOKIE_TTL_SECONDS = 600 # 10-minute TTL before re-warming | |
| # ── Lifecycle ───────────────────────────────────────────────── | |
| def _get_chrome_major_version() -> Optional[int]: | |
| """Helper to detect the installed Chrome major version on Windows/Linux to prevent driver mismatch.""" | |
| try: | |
| import winreg | |
| # Check user-level install / BLBeacon first | |
| try: | |
| with winreg.OpenKey(winreg.HKEY_CURRENT_USER, r"Software\Google\Chrome\BLBeacon") as key: | |
| version, _ = winreg.QueryValueEx(key, "version") | |
| return int(version.split(".")[0]) | |
| except Exception: | |
| pass | |
| # Check system-level WOW64 install | |
| try: | |
| with winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE, r"SOFTWARE\Wow6432Node\Microsoft\Windows\CurrentVersion\Uninstall\Google Chrome") as key: | |
| version, _ = winreg.QueryValueEx(key, "DisplayVersion") | |
| return int(version.split(".")[0]) | |
| except Exception: | |
| pass | |
| # Check system-level 64-bit install | |
| try: | |
| with winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE, r"SOFTWARE\Microsoft\Windows\CurrentVersion\Uninstall\Google Chrome") as key: | |
| version, _ = winreg.QueryValueEx(key, "DisplayVersion") | |
| return int(version.split(".")[0]) | |
| except Exception: | |
| pass | |
| except ImportError: | |
| # Linux/macOS | |
| import subprocess | |
| import re | |
| for cmd in ["google-chrome", "google-chrome-stable", "chromium", "chromium-browser"]: | |
| try: | |
| output = subprocess.check_output([cmd, "--version"], stderr=subprocess.DEVNULL) | |
| version_str = output.decode("utf-8").strip() | |
| match = re.search(r"(\d+)\.", version_str) | |
| if match: | |
| return int(match.group(1)) | |
| except Exception: | |
| continue | |
| return None | |
| # ── User-Agent helpers ──────────────────────────────────────── | |
| def _build_user_agent(chrome_version: Optional[int] = None) -> str: | |
| """Build a realistic Chrome User-Agent string using the installed version.""" | |
| v = chrome_version or _get_chrome_major_version() or 131 | |
| return ( | |
| f"Mozilla/5.0 (Windows NT 10.0; Win64; x64) " | |
| f"AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36" | |
| ) | |
| def _build_ua_pool() -> list[str]: | |
| """Build a pool of realistic User-Agent strings for rotation.""" | |
| v = _get_chrome_major_version() or 131 | |
| return [ | |
| # Windows Chrome (primary — matches our Selenium fingerprint) | |
| f"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36", | |
| # Windows Chrome (slightly older minor) | |
| f"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v-1}.0.0.0 Safari/537.36", | |
| # macOS Chrome | |
| f"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36", | |
| # Linux Chrome | |
| f"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36", | |
| # Windows Edge (Chromium-based, same engine) | |
| f"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36 Edg/{v}.0.0.0", | |
| ] | |
| _UA_POOL: list[str] = [] # populated at startup | |
| # ── Response validation helpers ─────────────────────────────── | |
| def _is_empty_stub(html: str) -> bool: | |
| """Detect LinkedIn's known empty/blocked response patterns. | |
| Returns True if the response is the 26-byte empty stub, a login wall, | |
| a CAPTCHA page, or any other pattern that indicates zero real content. | |
| """ | |
| stripped = html.strip() | |
| # 26-byte empty stub: '<!DOCTYPE html> <!----> ' | |
| if len(stripped) < 60 and "<!--" in stripped: | |
| return True | |
| # Login/signup redirect pages | |
| lower = stripped[:2000].lower() | |
| if any(sig in lower for sig in [ | |
| "login", "sign in", "signup", "sign up", | |
| "authwall", "auth_wall", "checkpoint", | |
| ]): | |
| return True | |
| return False | |
| def _create_selenium_driver(): | |
| """Create an undetected Chrome instance (runs in thread pool).""" | |
| options = uc.ChromeOptions() | |
| options.add_argument("--headless=new") | |
| options.add_argument("--no-sandbox") | |
| options.add_argument("--disable-dev-shm-usage") | |
| options.add_argument("--disable-gpu") | |
| options.add_argument("--disable-blink-features=AutomationControlled") | |
| options.add_argument("--window-size=1920,1080") | |
| options.add_argument(f"user-agent={_build_user_agent()}") | |
| version_main = _get_chrome_major_version() | |
| if version_main: | |
| driver = uc.Chrome(options=options, version_main=version_main) | |
| else: | |
| driver = uc.Chrome(options=options) | |
| driver.set_page_load_timeout(30) | |
| return driver | |
| _AUTO_SCRAPE_SCHOOLS = ["socse", "sob", "sodi", "soepp", "solaw", "sofmca", "solas", "soahp"] | |
| async def _auto_scrape_all_schools(): | |
| global _auto_scrape_state, _auto_scrape_history | |
| logger.info("⏰ Auto-scrape: starting school sweep") | |
| _auto_scrape_state = { | |
| "is_running": True, | |
| "current_school": None, | |
| "completed": [], | |
| "failed": [], | |
| "started_at": time.time(), | |
| "finished_at": None, | |
| "cancel_requested": False, | |
| } | |
| # ── Create a history entry for this sweep run ────────────── | |
| sweep_entry = { | |
| "sweep_id": len(_auto_scrape_history) + 1, | |
| "is_running": True, | |
| "current_school": None, | |
| "completed": [], | |
| "failed": [], | |
| "started_at": time.time(), | |
| "finished_at": None, | |
| "cancelled": False, | |
| } | |
| _auto_scrape_history.append(sweep_entry) | |
| all_scraped_companies = set() | |
| for school in _AUTO_SCRAPE_SCHOOLS: | |
| if _auto_scrape_state.get("cancel_requested"): | |
| logger.info("⏰ Auto-scrape: cancelled by user") | |
| sweep_entry["cancelled"] = True | |
| break | |
| _auto_scrape_state["current_school"] = school | |
| sweep_entry["current_school"] = school | |
| try: | |
| logger.info(f"⏰ Auto-scrape: triggering [{school}]") | |
| params = { | |
| "keywords": school, | |
| "freshness": "r86400", | |
| "work_types": ["onsite", "remote", "hybrid"], | |
| "job_types": ["internship"], | |
| } | |
| async with httpx.AsyncClient( | |
| base_url=_SELF_BASE_URL, | |
| timeout=None, | |
| ) as client: | |
| async with client.stream( | |
| "POST", | |
| "/scrape-internships", | |
| params=params, | |
| ) as resp: | |
| total = 0 | |
| async for line in resp.aiter_lines(): | |
| if not line.strip(): | |
| continue | |
| try: | |
| event = json.loads(line) | |
| etype = event.get("type") | |
| if etype == "info": | |
| logger.info(f"⏰ [{school}] {event.get('message')}") | |
| elif etype == "jobs": | |
| jobs_data = event.get("data", []) | |
| total += len(jobs_data) | |
| for j in jobs_data: | |
| if j.get("company"): | |
| all_scraped_companies.add(j["company"]) | |
| logger.info(f"⏰ [{school}] +{len(jobs_data)} jobs streamed") | |
| elif etype == "done": | |
| total = event.get("total", total) | |
| logger.info( | |
| f"⏰ [{school}] done — total={total} " | |
| f"engine={event.get('engine')} " | |
| f"filtered_out={event.get('filtered_out')}" | |
| ) | |
| elif etype == "error": | |
| logger.error(f"⏰ [{school}] error: {event.get('message')}") | |
| except Exception: | |
| pass | |
| _auto_scrape_state["completed"].append({"school": school, "jobs": total}) | |
| sweep_entry["completed"].append({"school": school, "jobs": total}) | |
| await asyncio.sleep(random.uniform(10.0, 20.0)) | |
| except Exception as e: | |
| logger.error(f"⏰ Auto-scrape [{school}]: failed — {e}") | |
| _auto_scrape_state["failed"].append({"school": school, "error": str(e)}) | |
| sweep_entry["failed"].append({"school": school, "error": str(e)}) | |
| continue | |
| _auto_scrape_state["is_running"] = False | |
| _auto_scrape_state["current_school"] = None | |
| _auto_scrape_state["finished_at"] = time.time() | |
| # ── Finalize the history entry ───────────────────────────── | |
| sweep_entry["is_running"] = False | |
| sweep_entry["current_school"] = None | |
| sweep_entry["finished_at"] = time.time() | |
| logger.info("⏰ Auto-scrape: sweep complete") | |
| if all_scraped_companies and not _auto_scrape_state.get("cancel_requested"): | |
| try: | |
| loop = asyncio.get_event_loop() | |
| already_rated = await loop.run_in_executor(executor, get_rated_companies) | |
| unrated_companies = [ | |
| c for c in all_scraped_companies | |
| if c.lower() not in already_rated | |
| ] | |
| logger.info( | |
| f"⏰ Auto-scrape: {len(all_scraped_companies)} companies scraped, " | |
| f"{len(already_rated)} already rated in DB, " | |
| f"rating {len(unrated_companies)} new ones..." | |
| ) | |
| if unrated_companies: | |
| ratings = await process_company_ratings(unrated_companies) | |
| await loop.run_in_executor(executor, _persist_ratings, ratings) | |
| else: | |
| logger.info("⏰ Auto-scrape: all companies already rated — skipping Serper search") | |
| except Exception as e: | |
| logger.error(f"⏰ Auto-scrape: failed to rate companies: {e}") | |
| async def startup_event(): | |
| global se_driver, pw_browser, pw_stealth_ctx, _UA_POOL | |
| loop = asyncio.get_event_loop() | |
| # 0. Build the UA rotation pool | |
| _UA_POOL = _build_ua_pool() | |
| logger.info(f"✓ UA pool built ({len(_UA_POOL)} variants, Chrome v{_get_chrome_major_version() or '?'})") | |
| # 1. Primary: Selenium (best anti-detection for LinkedIn) | |
| try: | |
| se_driver = await loop.run_in_executor(executor, _create_selenium_driver) | |
| logger.info("✓ Selenium undetected-chromedriver ready") | |
| except Exception as e: | |
| logger.warning(f"✗ Selenium init failed (Chrome installed?): {e}") | |
| # 2. Fallback: Playwright + Stealth | |
| try: | |
| stealth = Stealth() | |
| pw_stealth_ctx = stealth.use_async(async_playwright()) | |
| pw = await pw_stealth_ctx.__aenter__() | |
| pw_browser = await pw.chromium.launch( | |
| headless=True, | |
| args=["--disable-blink-features=AutomationControlled"], | |
| ) | |
| logger.info("✓ Playwright stealth browser ready") | |
| except Exception as e: | |
| logger.warning(f"✗ Playwright init failed: {e}") | |
| if not scheduler.running: | |
| scheduler.start() | |
| logger.info("✓ Scheduler started (awaiting manual trigger)") | |
| else: | |
| logger.info("⚠ Scheduler already running, skipping start") | |
| async def shutdown_event(): | |
| global se_driver, pw_browser, pw_stealth_ctx | |
| scheduler.shutdown(wait=False) | |
| loop = asyncio.get_event_loop() | |
| if se_driver: | |
| await loop.run_in_executor(executor, se_driver.quit) | |
| if pw_browser: | |
| await pw_browser.close() | |
| if pw_stealth_ctx: | |
| await pw_stealth_ctx.__aexit__(None, None, None) | |
| # ── Cookie warm-up ──────────────────────────────────────────── | |
| def _selenium_get_cookies(driver, url: str) -> dict: | |
| """Navigate to a URL and extract cookies (runs in executor thread).""" | |
| driver.get(url) | |
| time.sleep(random.uniform(2.0, 4.0)) # let LinkedIn set session cookies | |
| try: | |
| return {c["name"]: c["value"] for c in driver.get_cookies()} | |
| except Exception: | |
| return {} | |
| async def _warm_up_cookies(force: bool = False) -> dict: | |
| """Obtain fresh LinkedIn session cookies via a browser visit. | |
| Uses Selenium (preferred) or Playwright to visit a simple LinkedIn | |
| guest page, allowing LinkedIn to set its session/tracking cookies | |
| (JSESSIONID, bcookie, li_gc, etc.). These cookies are then forwarded | |
| to the httpx HTTP client for API requests. | |
| Results are cached for _COOKIE_TTL_SECONDS (10 min) to avoid | |
| redundant browser visits on consecutive scrapes. | |
| """ | |
| global _warm_cookies, _warm_cookies_ts, se_driver | |
| # Return cached cookies if still fresh | |
| if not force and _warm_cookies and (time.time() - _warm_cookies_ts) < _COOKIE_TTL_SECONDS: | |
| logger.info("cookie warm-up: using cached cookies (still fresh)") | |
| return _warm_cookies | |
| warm_url = "https://www.linkedin.com/jobs/search/?keywords=intern&f_TPR=r86400" | |
| cookies: dict = {} | |
| # Try Selenium first | |
| if se_driver: | |
| try: | |
| loop = asyncio.get_event_loop() | |
| cookies = await loop.run_in_executor( | |
| executor, _selenium_get_cookies, se_driver, warm_url | |
| ) | |
| if cookies: | |
| logger.info(f"cookie warm-up: got {len(cookies)} cookies via Selenium") | |
| except Exception as e: | |
| logger.warning(f"cookie warm-up: Selenium failed: {e}") | |
| err_str = str(e).lower() | |
| if "window" in err_str or "closed" in err_str or "view" in err_str: | |
| logger.info("cookie warm-up: Recreating crashed Selenium driver...") | |
| try: | |
| loop = asyncio.get_event_loop() | |
| try: | |
| await loop.run_in_executor(executor, se_driver.quit) | |
| except Exception: | |
| pass | |
| se_driver = await loop.run_in_executor(executor, _create_selenium_driver) | |
| cookies = await loop.run_in_executor( | |
| executor, _selenium_get_cookies, se_driver, warm_url | |
| ) | |
| if cookies: | |
| logger.info(f"cookie warm-up: got {len(cookies)} cookies via recreated Selenium") | |
| except Exception as e2: | |
| logger.warning(f"cookie warm-up: Selenium recreation failed: {e2}") | |
| # Fall back to Playwright | |
| if not cookies and pw_browser: | |
| try: | |
| ctx = await pw_browser.new_context(user_agent=_build_user_agent()) | |
| page = await ctx.new_page() | |
| await page.goto(warm_url, wait_until="domcontentloaded", timeout=20000) | |
| await asyncio.sleep(random.uniform(2.0, 4.0)) | |
| cookies = {c["name"]: c["value"] for c in await ctx.cookies()} | |
| await page.close() | |
| await ctx.close() | |
| if cookies: | |
| logger.info(f"cookie warm-up: got {len(cookies)} cookies via Playwright") | |
| except Exception as e: | |
| logger.warning(f"cookie warm-up: Playwright failed: {e}") | |
| if cookies: | |
| _warm_cookies = cookies | |
| _warm_cookies_ts = time.time() | |
| else: | |
| logger.warning("cookie warm-up: no cookies obtained from any engine") | |
| return cookies | |
| # ── URL builder ─────────────────────────────────────────────── | |
| # | |
| # Reference: https://www.linkedin.com/jobs/search/ guest-search parameter spec | |
| # sortBy : R = relevance (fixed) | |
| # f_E : 1 = Internship experience level (fixed) | |
| # geoId : LinkedIn region id — narrows location (fixed, Bengaluru region) | |
| # f_PP : comma-joined Primary-Place ids — UI-only narrowing (fixed) | |
| # f_TPR : r<seconds> freshness — variable (dropdown) | |
| # f_WT : comma-joined work-type codes — variable (3-checkbox group) | |
| # f_JT : comma-joined job-type codes — variable (2-checkbox group) | |
| # | |
| # Per project spec, location is fixed (no city picker) and we only expose | |
| # three user-visible filters: work-type, job-type, freshness. | |
| # | |
| # IMPORTANT — geoId vs f_PP split: | |
| # • The full UI page at /jobs/search/ respects BOTH geoId AND f_PP and uses | |
| # f_PP to narrow within a broader geoId. We keep both on the UI URL so the | |
| # "Open in LinkedIn" link matches the user's reference URL. | |
| # • The guest /jobs-guest/jobs/api/seeMoreJobPostings/search endpoint | |
| # SILENTLY returns an empty 26-byte stub when f_PP is present (it does | |
| # not support the f_PP filter), so the HTTP-pagination path strips f_PP | |
| # and relies on geoId only. See _PAGINATION_PARAM_KEYS below. | |
| # Always-on params. geoId 90009633 = "Bengaluru" (covers Bengaluru metro on | |
| # the guest API). f_PP narrows further on the UI URL for the "Open in | |
| # LinkedIn" link; the API silently ignores it so we strip it for pagination. | |
| _FIXED_PARAMS: dict[str, str] = { | |
| "sortBy": "R", | |
| "f_E": "1", | |
| "geoId": "90009633", | |
| } | |
| # Work-type checkbox options (label → LinkedIn f_WT code). | |
| WORK_TYPES: dict[str, str] = { | |
| "onsite": "1", | |
| "remote": "2", | |
| "hybrid": "3", | |
| } | |
| # Job-type checkbox options (label → LinkedIn f_JT code) — spec restricts to | |
| # these two only. | |
| JOB_TYPES: dict[str, str] = { | |
| "internship": "I", | |
| "full_time": "F", | |
| } | |
| # Freshness dropdown presets (label → LinkedIn f_TPR code). | |
| FRESHNESS_PRESETS: dict[str, str] = { | |
| "hour": "r3600", | |
| "day": "r86400", | |
| "week": "r604800", | |
| "month": "r2592000", | |
| } | |
| # ── School-based search term generation ─────────────────────── | |
| # | |
| # Convenient shortcut keywords that resolve to a specific school. | |
| # The user specifically requested that searches for a school (e.g. "SOB") | |
| # should search EXACTLY these keywords on LinkedIn. | |
| KEYWORD_TO_SCHOOL: dict[str, str] = { | |
| # [[SOCSE]] — Computer Science & Engineering | |
| "software": "SOCSE", | |
| "cs": "SOCSE", | |
| "computer science": "SOCSE", | |
| "programming": "SOCSE", | |
| "coding": "SOCSE", | |
| "data": "SOCSE", | |
| "data science": "SOCSE", | |
| "machine learning": "SOCSE", | |
| "ai": "SOCSE", | |
| "cyber": "SOCSE", | |
| "cybersecurity": "SOCSE", | |
| "qa": "SOCSE", | |
| "devops": "SOCSE", | |
| "frontend": "SOCSE", | |
| "backend": "SOCSE", | |
| "fullstack": "SOCSE", | |
| "cloud": "SOCSE", | |
| "developer": "SOCSE", | |
| "engineer": "SOCSE", | |
| "sde": "SOCSE", | |
| "swe": "SOCSE", | |
| "it": "SOCSE", | |
| "web development": "SOCSE", | |
| "app development": "SOCSE", | |
| "mobile development": "SOCSE", | |
| "android": "SOCSE", | |
| "ios": "SOCSE", | |
| "react": "SOCSE", | |
| "node": "SOCSE", | |
| "python": "SOCSE", | |
| "java": "SOCSE", | |
| "c++": "SOCSE", | |
| "golang": "SOCSE", | |
| "flutter": "SOCSE", | |
| "aws": "SOCSE", | |
| "azure": "SOCSE", | |
| "gcp": "SOCSE", | |
| "data analyst": "SOCSE", | |
| "data engineer": "SOCSE", | |
| "nlp": "SOCSE", | |
| "computer vision": "SOCSE", | |
| "quality assurance": "SOCSE", | |
| "testing": "SOCSE", | |
| "blockchain": "SOCSE", | |
| "web3": "SOCSE", | |
| "systems engineer": "SOCSE", | |
| "network engineer": "SOCSE", | |
| "site reliability": "SOCSE", | |
| "sre": "SOCSE", | |
| "hardware": "SOCSE", | |
| "embedded": "SOCSE", | |
| # SODI — Design & Innovation | |
| "designing": "SOCSE,SODI", | |
| "ui ux": "SOCSE,SODI", | |
| "graphic designer": "SOCSE,SODI", | |
| "ui designer": "SOCSE,SODI", | |
| "ux designer": "SOCSE,SODI", | |
| "product designer": "SOCSE,SODI", | |
| "product design": "SODI", | |
| "industrial design":"SODI", | |
| "interaction design":"SODI", | |
| "visual design": "SODI", | |
| # SOB — Business | |
| "business": "SOB", | |
| "marketing": "SOB", | |
| "finance": "SOB", | |
| "hr": "SOB", | |
| "accounting": "SOB", | |
| "commerce": "SOB", | |
| "mba": "SOB", | |
| "sales": "SOB", | |
| "strategy": "SOB", | |
| "operations": "SOB", | |
| "management": "SOB", | |
| "consulting": "SOB", | |
| "analyst": "SOB", | |
| "business analyst": "SOB", | |
| "business development": "SOB", | |
| "bda": "SOB", | |
| "bde": "SOB", | |
| "sales intern": "SOB", | |
| "marketing intern": "SOB", | |
| "hr intern": "SOB", | |
| "finance intern": "SOB", | |
| "investment banking": "SOB", | |
| "venture capital": "SOB", | |
| "private equity": "SOB", | |
| "digital marketing": "SOB", | |
| "seo": "SOB", | |
| "content writer": "SOB", | |
| "social media": "SOB", | |
| "growth": "SOB", | |
| "product manager": "SOB", | |
| "project manager": "SOB", | |
| # SoEPP — Economics & Public Policy | |
| "economics": "SoEPP", | |
| "policy": "SoEPP", | |
| "governance": "SoEPP", | |
| "public policy": "SoEPP", | |
| "research": "SoEPP", | |
| "development studies":"SoEPP", | |
| "economist": "SoEPP", | |
| "policy analyst": "SoEPP", | |
| "public relations": "SoEPP", | |
| "government": "SoEPP", | |
| # SOLaw — Law | |
| "law": "SOLaw", | |
| "legal": "SOLaw", | |
| "paralegal": "SOLaw", | |
| "attorney": "SOLaw", | |
| "counsel": "SOLaw", | |
| "cyber law": "SOLaw", | |
| "law intern": "SOLaw", | |
| "legal intern": "SOLaw", | |
| "litigation": "SOLaw", | |
| "corporate law": "SOLaw", | |
| "ipr": "SOLaw", | |
| # SOFMCA — Film, Media & Creative Arts | |
| "media": "SOFMCA", | |
| "film": "SOFMCA", | |
| "journalism": "SOFMCA", | |
| "animation": "SOFMCA", | |
| "vfx": "SOFMCA", | |
| "gaming": "SOFMCA", | |
| "content creation": "SOFMCA", | |
| "acting": "SOFMCA", | |
| "editor": "SOFMCA", | |
| "video editor": "SOFMCA", | |
| "cinematographer": "SOFMCA", | |
| "copywriter": "SOFMCA", | |
| "media intern": "SOFMCA", | |
| "pr intern": "SOFMCA", | |
| "producer": "SOFMCA", | |
| # SOLAS — Liberal Arts & Sciences | |
| "psychology": "SOLAS", | |
| "environment": "SOLAS", | |
| "liberal arts": "SOLAS", | |
| "sociology": "SOLAS", | |
| "history": "SOLAS", | |
| "behavioral science":"SOLAS", | |
| "psychology intern": "SOLAS", | |
| "counseling": "SOLAS", | |
| "sociologist": "SOLAS", | |
| "research assistant":"SOLAS", | |
| # SOAHP — Allied & Healthcare | |
| "healthcare": "SOAHP", | |
| "medical": "SOAHP", | |
| "clinical": "SOAHP", | |
| "laboratory": "SOAHP", | |
| "nursing": "SOAHP", | |
| "public health": "SOAHP", | |
| "clinical research": "SOAHP", | |
| "public health intern": "SOAHP", | |
| "hospital administration": "SOAHP", | |
| "pharma": "SOAHP", | |
| } | |
| # Build SCHOOL_SEARCH_TERMS directly from KEYWORD_TO_SCHOOL. | |
| # The user explicitly wants us to search *exactly* these keywords. | |
| SCHOOL_SEARCH_TERMS: dict[str, list[str]] = {} | |
| for kw, school_str in KEYWORD_TO_SCHOOL.items(): | |
| for school in school_str.split(","): | |
| school_lower = school.strip().lower() | |
| if school_lower not in SCHOOL_SEARCH_TERMS: | |
| SCHOOL_SEARCH_TERMS[school_lower] = [] | |
| SCHOOL_SEARCH_TERMS[school_lower].append(kw) | |
| def build_search_params( | |
| keywords: str, | |
| freshness: str = "r86400", | |
| work_types: Optional[list[str]] = None, | |
| job_types: Optional[list[str]] = None, | |
| ) -> dict[str, str]: | |
| """Build the LinkedIn guest search query-param dict. | |
| Fixed params (sortBy, f_E, f_PP) are always present. | |
| Variable params reflect user-selected checkboxes / dropdown: | |
| - f_TPR : freshness (raw "r<seconds>" code) | |
| - f_WT : comma-joined work-type codes (only when any box ticked) | |
| - f_JT : comma-joined job-type codes (only when any box ticked) | |
| """ | |
| params: dict[str, str] = { | |
| "keywords": keywords, | |
| **_FIXED_PARAMS, | |
| "f_TPR": freshness, | |
| } | |
| if work_types: | |
| codes = [WORK_TYPES[k] for k in work_types if k in WORK_TYPES] | |
| if codes: | |
| params["f_WT"] = ",".join(codes) | |
| if job_types: | |
| codes = [JOB_TYPES[k] for k in job_types if k in JOB_TYPES] | |
| if codes: | |
| params["f_JT"] = ",".join(codes) | |
| return params | |
| def build_search_url( | |
| keywords: str, | |
| freshness: str = "r86400", | |
| work_types: Optional[list[str]] = None, | |
| job_types: Optional[list[str]] = None, | |
| ) -> str: | |
| """Builds a fully-parameterised LinkedIn guest search URL. | |
| Shape matches LinkedIn's UI-canonical form, e.g.: | |
| https://www.linkedin.com/jobs/search/?keywords=intern&sortBy=R&f_E=1 | |
| &f_PP=105214831,113968072,112565523,119634689&f_TPR=r86400&f_WT=1,2&f_JT=I | |
| """ | |
| params = build_search_params( | |
| keywords=keywords, | |
| freshness=freshness, | |
| work_types=work_types, | |
| job_types=job_types, | |
| ) | |
| # Trailing slash on /jobs/search/ matches LinkedIn's UI-canonical URL. | |
| return f"https://www.linkedin.com/jobs/search/?{urllib.parse.urlencode(params)}" | |
| # ── Shared HTML parser ──────────────────────────────────────── | |
| def _parse_jobs_from_html(html: str) -> list[dict]: | |
| """Extract job listings from raw HTML using BeautifulSoup. | |
| Cascading selectors handle multiple LinkedIn page variants. | |
| NOTE: contact_details (hiring team links) are NOT extracted here. | |
| They live on each job's *detail* page, not the search results list. | |
| Use _fetch_hiring_links_for_jobs() after collection to enrich jobs. | |
| """ | |
| soup = BeautifulSoup(html, "html.parser") | |
| jobs = [] | |
| # Primary selector: public guest search results | |
| cards = soup.select("ul.jobs-search__results-list > li") | |
| # Fallback selectors for alternate page states + seeMoreJobPostings fragments | |
| # (the paginated endpoint returns bare <li> elements with no wrapping <ul>). | |
| if not cards: | |
| cards = soup.select("li.job-search-card, li.result-card") | |
| if not cards: | |
| cards = soup.select("div.base-card.base-search-card") | |
| if not cards: | |
| cards = soup.select("[data-entity-urn*='jobPosting']") | |
| for card in cards: | |
| try: | |
| # Title | |
| title_el = ( | |
| card.select_one(".base-search-card__title") | |
| or card.select_one("h3") | |
| or card.select_one("[class*='title']") | |
| ) | |
| # Company | |
| company_el = ( | |
| card.select_one(".base-search-card__subtitle a") | |
| or card.select_one(".base-search-card__subtitle") | |
| or card.select_one("h4 a") | |
| ) | |
| # Location | |
| loc_el = ( | |
| card.select_one(".job-search-card__location") | |
| or card.select_one("[class*='location']") | |
| ) | |
| # Link | |
| link_el = ( | |
| card.select_one("a.base-card__full-link") | |
| or card.select_one("a[href*='/jobs/view/']") | |
| or card.select_one("a[href*='linkedin.com/jobs']") | |
| ) | |
| # Posted date | |
| date_el = ( | |
| card.select_one("time.job-search-card__listdate") | |
| or card.select_one("time.job-search-card__listdate--new") | |
| or card.select_one("time") | |
| ) | |
| title = title_el.get_text(strip=True) if title_el else None | |
| company = company_el.get_text(strip=True) if company_el else None | |
| location = loc_el.get_text(strip=True) if loc_el else None | |
| raw_link = link_el.get("href", "") if link_el else "" | |
| link = raw_link.split("?")[0] if raw_link else None | |
| posted = date_el.get_text(strip=True) if date_el else None | |
| posted_dt = date_el.get("datetime") if date_el else None | |
| # At least title or company must exist for a valid card | |
| if title or company: | |
| classification = classify_job(title or "", company or "") | |
| jobs.append({ | |
| "title": title, | |
| "company": company, | |
| "location": location, | |
| "link": link, | |
| "posted": posted, | |
| "posted_datetime": posted_dt, | |
| "programs": classification["programs"], | |
| "schools": classification["schools"], | |
| "contact_details": [], # enriched later by _fetch_hiring_links_for_jobs | |
| }) | |
| except Exception: | |
| continue | |
| return jobs | |
| # ── Serper-based hiring manager search ─────────────────────── | |
| _SERPER_API_KEY = os.environ.get("SERPER_API_KEY", "0820e51e9080c289b3849e30189e023774fdad87") | |
| _SERPER_URL = "https://google.serper.dev/search" | |
| # Semaphore: at most 3 concurrent Serper requests | |
| _SERPER_SEM = asyncio.Semaphore(3) | |
| _LINKEDIN_IN_RE = re.compile(r'https?://(?:[\w-]+\.)?linkedin\.com/in/([\w%-]+)', re.IGNORECASE) | |
| def _extract_linkedin_profile_from_serper(data: dict) -> dict | None: | |
| """Parse a Serper JSON response and return the first LinkedIn /in/ profile found. | |
| Priority order: | |
| 1. answerBox.link (Google AI / Featured Snippet) | |
| 2. answerBox.snippet (text contains a linkedin.com/in/ URL) | |
| 3. organic[0].link (first organic result) | |
| 4. organic[].link (first organic result that is a /in/ URL) | |
| Returns: | |
| {"name": "<person_slug_or_extracted_name>", "url": "<linkedin_profile_url>"} | |
| or None if nothing found. | |
| """ | |
| # 1. answerBox direct link | |
| ab = data.get("answerBox", {}) | |
| ab_link = ab.get("link", "") | |
| if ab_link and "linkedin.com/in/" in ab_link: | |
| slug = ab_link.rstrip("/").split("/in/")[-1].split("?")[0] | |
| name = ab.get("title") or slug.replace("-", " ").title() | |
| return {"name": name, "url": ab_link.split("?")[0]} | |
| # 2. answerBox snippet contains a URL | |
| ab_snippet = ab.get("snippet", "") or ab.get("answer", "") | |
| m = _LINKEDIN_IN_RE.search(ab_snippet) | |
| if m: | |
| url = m.group(0).split("?")[0] | |
| slug = m.group(1) | |
| name = ab.get("title") or slug.replace("-", " ").title() | |
| return {"name": name, "url": url} | |
| # 3 & 4. Organic results — prefer linkedin.com/in/ links | |
| for result in data.get("organic", []): | |
| link = result.get("link", "") | |
| if "linkedin.com/in/" in link: | |
| slug = link.rstrip("/").split("/in/")[-1].split("?")[0] | |
| name = result.get("title", slug).split(" - ")[0].split(" | ")[0].strip() | |
| return {"name": name, "url": link.split("?")[0]} | |
| return None | |
| async def _fetch_hiring_contacts_via_serper(jobs: list[dict]) -> None: | |
| """Enrich each job dict in-place with contact_details via Serper Google search. | |
| For each job, searches: | |
| site:linkedin.com/in "{company}" "hiring" OR "recruiter" OR "HR" | |
| contact_details becomes a list with at most 1 entry: | |
| [{"name": "Person Name", "url": "https://linkedin.com/in/slug"}] | |
| Jobs where no profile is found keep contact_details as []. | |
| Requires SERPER_API_KEY in environment. | |
| """ | |
| if not _SERPER_API_KEY: | |
| logger.warning("serper: SERPER_API_KEY not set — skipping hiring contact search") | |
| return | |
| async def _search_one(job: dict) -> None: | |
| company = job.get("company") or "" | |
| if not company: | |
| return | |
| # Build a targeted query: site operator + company + hiring signals | |
| query = f'site:linkedin.com/in "{company}" hiring OR recruiter OR "HR"' | |
| try: | |
| async with _SERPER_SEM: | |
| await asyncio.sleep(random.uniform(0.2, 0.6)) | |
| async with httpx.AsyncClient(timeout=10.0) as client: | |
| resp = await client.post( | |
| _SERPER_URL, | |
| json={"q": query, "num": 5, "gl": "in", "hl": "en"}, | |
| headers={ | |
| "X-API-KEY": _SERPER_API_KEY, | |
| "Content-Type": "application/json", | |
| }, | |
| ) | |
| if resp.status_code != 200: | |
| logger.debug( | |
| f"serper: HTTP {resp.status_code} for company={company!r}" | |
| ) | |
| return | |
| data = resp.json() | |
| profile = _extract_linkedin_profile_from_serper(data) | |
| if profile: | |
| job["contact_details"] = [profile] | |
| logger.info( | |
| f"serper: ✓ {job.get('title','?')} @ {company} → " | |
| f"{profile['name']} ({profile['url']})" | |
| ) | |
| else: | |
| logger.debug(f"serper: no profile found for company={company!r}") | |
| except Exception as e: | |
| logger.debug(f"serper: error searching for {company!r}: {e}") | |
| await asyncio.gather(*[_search_one(j) for j in jobs]) | |
| # ── Paginated HTTP fetcher (seeMoreJobPostings) ─────────────── | |
| _SEE_MORE_URL = ( | |
| "https://www.linkedin.com/jobs-guest/jobs/api/seeMoreJobPostings/search" | |
| ) | |
| # Keys forwarded to the seeMoreJobPostings endpoint. | |
| # NOTE: f_PP is intentionally EXCLUDED here — the guest API returns an empty | |
| # 26-byte stub whenever f_PP is present (verified 2026-05-20). geoId provides | |
| # equivalent narrowing. f_PP stays on the UI URL only. | |
| _PAGINATION_PARAM_KEYS = { | |
| "keywords", "sortBy", "f_E", "geoId", "f_TPR", "f_WT", "f_JT", | |
| } | |
| async def _fetch_more_pages( | |
| base_params: dict, | |
| cookies: dict, | |
| headers: dict, | |
| max_pages: int = 20, | |
| start_offset: int = 25, | |
| ) -> list[dict]: | |
| """Page through LinkedIn's guest seeMoreJobPostings endpoint. | |
| The endpoint returns ~25 job cards as a raw HTML fragment per call. | |
| Stops on: 4xx/5xx (rate limit, with retry on 429), empty body, | |
| empty stub detection, or empty parsed result. Returns whatever was | |
| collected before stopping — never raises. | |
| Retries on transient network errors (DNS, connection timeouts) with | |
| exponential backoff to handle temporary infrastructure issues. | |
| """ | |
| params_clean = { | |
| k: v for k, v in base_params.items() | |
| if k in _PAGINATION_PARAM_KEYS and v not in (None, "") | |
| } | |
| collected: list[dict] = [] | |
| rate_limit_retries_left = 2 # total 429s tolerated across the whole loop | |
| first_response_dumped = False | |
| dump_path = os.path.join(tempfile.gettempdir(), "linkedin_pagination_first_response.html") | |
| try: | |
| async with httpx.AsyncClient( | |
| timeout=20.0, | |
| cookies=cookies, | |
| headers=headers, | |
| follow_redirects=True, | |
| http2=False, | |
| ) as client: | |
| # Brief settle before first request — reduced from 3s for speed. | |
| await asyncio.sleep(random.uniform(1.0, 2.0)) | |
| start = start_offset | |
| for page_idx in range(max_pages): | |
| url = f"{_SEE_MORE_URL}?{urllib.parse.urlencode({**params_clean, 'start': start})}" | |
| logger.info(f"pagination: GET {url}") | |
| # Retry loop for transient network errors | |
| max_network_retries = 3 | |
| for attempt in range(1, max_network_retries + 1): | |
| try: | |
| resp = await client.get(url) | |
| break # Success, exit retry loop | |
| except (httpx.TimeoutException, httpx.NetworkError, OSError) as e: | |
| # Transient network errors: DNS (getaddrinfo), timeout, connection reset | |
| if attempt < max_network_retries: | |
| backoff = min(2 ** attempt, 10.0) * random.uniform(1.0, 1.5) | |
| logger.warning( | |
| f"pagination: transient network error at start={start} " | |
| f"(attempt {attempt}/{max_network_retries}): {type(e).__name__}. " | |
| f"Retrying in {backoff:.1f}s..." | |
| ) | |
| await asyncio.sleep(backoff) | |
| else: | |
| logger.warning( | |
| f"pagination: network error at start={start} after " | |
| f"{max_network_retries} attempts: {e}; stopping with {len(collected)} jobs" | |
| ) | |
| return collected | |
| except httpx.HTTPError as e: | |
| logger.warning( | |
| f"pagination: HTTP error at start={start}: {e}; " | |
| f"stopping with {len(collected)} extra jobs" | |
| ) | |
| return collected | |
| # Retry on 429 with longer backoff — LinkedIn's rate limits | |
| # are usually short-lived and clear after ~10s. | |
| if resp.status_code == 429 and rate_limit_retries_left > 0: | |
| rate_limit_retries_left -= 1 | |
| backoff = random.uniform(15.0, 25.0) | |
| logger.warning( | |
| f"pagination: 429 at start={start}; backing off {backoff:.1f}s " | |
| f"and retrying ({rate_limit_retries_left} retries left)" | |
| ) | |
| await asyncio.sleep(backoff) | |
| try: | |
| resp = await client.get(url) | |
| except (httpx.TimeoutException, httpx.NetworkError, OSError, httpx.HTTPError) as e: | |
| logger.warning(f"pagination: retry HTTP error: {e}; stopping") | |
| return collected | |
| html = resp.text or "" | |
| body_len = len(resp.content) if resp.content is not None else 0 | |
| li_count = html.lower().count("<li") | |
| logger.info( | |
| f"pagination: start={start} status={resp.status_code} " | |
| f"bytes={body_len} <li>={li_count}" | |
| ) | |
| if not first_response_dumped: | |
| try: | |
| with open(dump_path, "w", encoding="utf-8") as f: | |
| f.write(html) | |
| logger.info(f"pagination: dumped first response → {dump_path}") | |
| except Exception as e: | |
| logger.warning(f"pagination: could not dump response: {e}") | |
| first_response_dumped = True | |
| if resp.status_code >= 400: | |
| logger.warning( | |
| f"pagination: status {resp.status_code} at start={start} " | |
| f"(likely rate-limit); stopping with {len(collected)} extra jobs" | |
| ) | |
| break | |
| if not html.strip(): | |
| logger.info( | |
| f"pagination: empty body at start={start}; end of results " | |
| f"(+{len(collected)} extra jobs)" | |
| ) | |
| break | |
| # ── Early stub detection ────────────────────────── | |
| # Detect the 26-byte empty stub and login walls before | |
| # wasting time on HTML parsing. On the very first page | |
| # this signals "HTTP path is blocked" so we can bail fast. | |
| if _is_empty_stub(html): | |
| logger.warning( | |
| f"pagination: empty stub/block detected at start={start} " | |
| f"(bytes={body_len}); stopping with {len(collected)} extra jobs" | |
| ) | |
| break | |
| page_jobs = _parse_jobs_from_html(html) | |
| logger.info( | |
| f"pagination: start={start} parser returned {len(page_jobs)} jobs" | |
| ) | |
| if not page_jobs: | |
| snippet = html[:500].replace("\n", " ") | |
| logger.warning( | |
| f"pagination: parser returned 0 jobs from a non-empty body " | |
| f"(bytes={body_len}, <li>={li_count}). First 500 chars: {snippet!r}" | |
| ) | |
| logger.info( | |
| f"pagination: no jobs parsed at start={start}; end of results " | |
| f"(+{len(collected)} extra jobs)" | |
| ) | |
| break | |
| collected.extend(page_jobs) | |
| logger.info( | |
| f"pagination: start={start} → +{len(page_jobs)} jobs " | |
| f"(running paginated total {len(collected)})" | |
| ) | |
| # Advance by the actual page size LinkedIn returned. | |
| start += len(page_jobs) | |
| # Human-like inter-page jitter — varied enough to avoid | |
| # pattern detection. | |
| await asyncio.sleep(random.uniform(3.5, 6.0)) | |
| except Exception as e: | |
| logger.warning( | |
| f"pagination: unexpected error: {e}; returning {len(collected)} extra jobs" | |
| ) | |
| return collected | |
| # ── Scraping engines ────────────────────────────────────────── | |
| _JOB_CARD_CSS = ( | |
| "ul.jobs-search__results-list > li, " | |
| "div.base-card.base-search-card, " | |
| "[data-entity-urn*='jobPosting']" | |
| ) | |
| def _selenium_scrape(driver, url: str) -> tuple[list[dict], dict]: | |
| """Synchronous Selenium scrape — executed inside executor thread. | |
| Returns (jobs, cookies) so the async caller can paginate over HTTP.""" | |
| driver.get(url) | |
| # Wait up to 15 s for at least one job card — if nothing loads, bail so | |
| # Playwright fallback can try instead. | |
| try: | |
| WebDriverWait(driver, 15).until( | |
| EC.presence_of_element_located((By.CSS_SELECTOR, _JOB_CARD_CSS)) | |
| ) | |
| except Exception: | |
| logger.warning("Selenium: no job cards appeared within 15 s") | |
| return [], {} | |
| # Scroll to bottom repeatedly; break after 3 consecutive unchanged heights. | |
| no_change = 0 | |
| for _ in range(15): | |
| prev_height = driver.execute_script("return document.body.scrollHeight") | |
| driver.execute_script("window.scrollTo(0, document.body.scrollHeight)") | |
| time.sleep(1.5) | |
| # Click "See more jobs" / "Show more results" if present | |
| try: | |
| btn = driver.find_element( | |
| By.XPATH, | |
| "//button[contains(translate(., 'ABCDEFGHIJKLMNOPQRSTUVWXYZ', 'abcdefghijklmnopqrstuvwxyz'), 'see more')" | |
| " or contains(translate(., 'ABCDEFGHIJKLMNOPQRSTUVWXYZ', 'abcdefghijklmnopqrstuvwxyz'), 'show more')]" | |
| ) | |
| btn.click() | |
| time.sleep(1.0) | |
| except Exception: | |
| pass | |
| new_height = driver.execute_script("return document.body.scrollHeight") | |
| if new_height == prev_height: | |
| no_change += 1 | |
| if no_change >= 3: | |
| break | |
| else: | |
| no_change = 0 | |
| jobs = _parse_jobs_from_html(driver.page_source) | |
| try: | |
| cookies = {c["name"]: c["value"] for c in driver.get_cookies()} | |
| except Exception as e: | |
| logger.warning(f"Selenium: could not extract cookies for pagination: {e}") | |
| cookies = {} | |
| return jobs, cookies | |
| async def _playwright_scrape(browser, url: str) -> tuple[list[dict], dict]: | |
| """Async Playwright scrape — fallback engine. | |
| Returns (jobs, cookies) so the caller can paginate over HTTP.""" | |
| context = await browser.new_context( | |
| user_agent=_build_user_agent() | |
| ) | |
| page = await context.new_page() | |
| try: | |
| await page.goto(url, wait_until="domcontentloaded", timeout=30000) | |
| # Wait up to 15 s for at least one job card | |
| try: | |
| await page.wait_for_selector(_JOB_CARD_CSS, timeout=15000) | |
| except Exception: | |
| logger.warning("Playwright: no job cards appeared within 15 s") | |
| return [], {} | |
| # Scroll to bottom repeatedly; break after 3 consecutive unchanged heights. | |
| no_change = 0 | |
| for _ in range(15): | |
| prev_height = await page.evaluate("document.body.scrollHeight") | |
| await page.evaluate("window.scrollTo(0, document.body.scrollHeight)") | |
| await asyncio.sleep(1.5) | |
| # Click "See more jobs" button if present | |
| try: | |
| await page.locator("button:has-text('See more')").click(timeout=1000) | |
| await asyncio.sleep(1.0) | |
| except Exception: | |
| pass | |
| new_height = await page.evaluate("document.body.scrollHeight") | |
| if new_height == prev_height: | |
| no_change += 1 | |
| if no_change >= 3: | |
| break | |
| else: | |
| no_change = 0 | |
| html = await page.content() | |
| jobs = _parse_jobs_from_html(html) | |
| try: | |
| cookies = {c["name"]: c["value"] for c in await context.cookies()} | |
| except Exception as e: | |
| logger.warning(f"Playwright: could not extract cookies for pagination: {e}") | |
| cookies = {} | |
| return jobs, cookies | |
| finally: | |
| await page.close() | |
| await context.close() | |
| # ── Single-URL scrape helper (Selenium → Playwright) ───────── | |
| def _build_pagination_headers(referer_url: str) -> dict: | |
| """Browser-realistic headers for the seeMoreJobPostings endpoint. | |
| Uses a random UA from the rotation pool and matches the installed | |
| Chrome version for Sec-Ch-Ua consistency.""" | |
| v = _get_chrome_major_version() or 131 | |
| ua = random.choice(_UA_POOL) if _UA_POOL else _build_user_agent() | |
| return { | |
| "User-Agent": ua, | |
| "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", | |
| "Accept-Language": "en-US,en;q=0.9", | |
| "Accept-Encoding": "gzip, deflate, br", | |
| "Referer": referer_url, | |
| "Sec-Fetch-Dest": "empty", | |
| "Sec-Fetch-Mode": "cors", | |
| "Sec-Fetch-Site": "same-origin", | |
| "Sec-Ch-Ua": f'"Chromium";v="{v}", "Not_A Brand";v="24"', | |
| "Sec-Ch-Ua-Mobile": "?0", | |
| "Sec-Ch-Ua-Platform": '"Windows"', | |
| "Connection": "keep-alive", | |
| } | |
| async def _scrape_one_url( | |
| url: str, | |
| params: Optional[dict] = None, | |
| paginate: bool = False, | |
| ) -> tuple[list[dict], str]: | |
| """Scrape one URL with cookie warm-up, exponential-backoff retry, | |
| and smarter browser fallback. | |
| Pipeline: | |
| 1. Warm up cookies via browser visit (cached for 10 min). | |
| 2. Try HTTP pagination with warm cookies (up to 3 attempts | |
| with exponential backoff on failure). | |
| 3. If HTTP fails after all retries, fall back to Selenium | |
| then Playwright to scrape the UI page directly. | |
| 4. If a browser engine succeeds, try one more HTTP pagination | |
| pass using the browser's fresh cookies to extend the results. | |
| Returns (jobs, engine_used) where engine_used is one of | |
| 'http' / 'selenium' / 'playwright' / 'none'. | |
| """ | |
| global se_driver | |
| jobs: list[dict] = [] | |
| engine_used = "none" | |
| browser_cookies: dict = {} | |
| # ── Phase 1 — HTTP with warm cookies + retry loop ───────────── | |
| if paginate and params: | |
| max_http_attempts = 3 | |
| for attempt in range(1, max_http_attempts + 1): | |
| # Warm up cookies (uses cache if fresh; force re-warm on retries) | |
| cookies = await _warm_up_cookies(force=(attempt > 1)) | |
| logger.info( | |
| f"HTTP attempt {attempt}/{max_http_attempts} " | |
| f"(cookies={'warm' if cookies else 'cold'})" | |
| ) | |
| http_jobs = await _fetch_more_pages( | |
| base_params=params, | |
| cookies=cookies, | |
| headers=_build_pagination_headers(referer_url=url), | |
| max_pages=40, | |
| start_offset=0, | |
| ) | |
| if http_jobs: | |
| jobs = http_jobs | |
| engine_used = "http" | |
| logger.info( | |
| f"HTTP scraped {len(jobs)} jobs on attempt {attempt}" | |
| ) | |
| break | |
| # Exponential backoff before next attempt | |
| if attempt < max_http_attempts: | |
| backoff = (2 ** attempt) * random.uniform(2.0, 4.0) | |
| logger.warning( | |
| f"HTTP attempt {attempt} returned 0 jobs; " | |
| f"backing off {backoff:.1f}s before retry" | |
| ) | |
| await asyncio.sleep(backoff) | |
| if not jobs: | |
| logger.warning( | |
| f"HTTP pagination failed after {max_http_attempts} attempts " | |
| "— falling back to browser engines" | |
| ) | |
| # ── Phase 2 — Browser fallback (Selenium → Playwright) ──────── | |
| if not jobs and se_driver: | |
| try: | |
| async with _scrape_lock: | |
| loop = asyncio.get_event_loop() | |
| jobs, browser_cookies = await loop.run_in_executor( | |
| executor, _selenium_scrape, se_driver, url | |
| ) | |
| if jobs: | |
| engine_used = "selenium" | |
| logger.info(f"Selenium scraped {len(jobs)} jobs (fallback)") | |
| except Exception as e: | |
| logger.warning(f"Selenium fallback failed: {e}") | |
| err_str = str(e).lower() | |
| if "window" in err_str or "closed" in err_str or "view" in err_str: | |
| logger.info("Selenium fallback: Recreating crashed driver...") | |
| try: | |
| async with _scrape_lock: | |
| loop = asyncio.get_event_loop() | |
| try: | |
| await loop.run_in_executor(executor, se_driver.quit) | |
| except Exception: | |
| pass | |
| se_driver = await loop.run_in_executor(executor, _create_selenium_driver) | |
| jobs, browser_cookies = await loop.run_in_executor( | |
| executor, _selenium_scrape, se_driver, url | |
| ) | |
| if jobs: | |
| engine_used = "selenium" | |
| logger.info(f"Selenium scraped {len(jobs)} jobs via recreated driver") | |
| except Exception as e2: | |
| logger.warning(f"Selenium fallback recreation failed: {e2}") | |
| if not jobs and pw_browser: | |
| try: | |
| jobs, browser_cookies = await _playwright_scrape(pw_browser, url) | |
| if jobs: | |
| engine_used = "playwright" | |
| logger.info(f"Playwright scraped {len(jobs)} jobs (fallback)") | |
| except Exception as e: | |
| logger.warning(f"Playwright fallback failed: {e}") | |
| # ── Phase 3 — Post-browser HTTP extension ───────────────────── | |
| # If a browser engine got results, try HTTP pagination with its | |
| # fresh cookies to collect additional pages beyond what the browser | |
| # initially loaded (the browser typically only sees page 1). | |
| if jobs and browser_cookies and paginate and params: | |
| logger.info( | |
| f"Extending {len(jobs)} browser results via HTTP with " | |
| f"{len(browser_cookies)} fresh browser cookies" | |
| ) | |
| extra_jobs = await _fetch_more_pages( | |
| base_params=params, | |
| cookies=browser_cookies, | |
| headers=_build_pagination_headers(referer_url=url), | |
| max_pages=20, | |
| start_offset=len(jobs), # continue from where the browser left off | |
| ) | |
| if extra_jobs: | |
| jobs.extend(extra_jobs) | |
| logger.info( | |
| f"HTTP extension added {len(extra_jobs)} jobs " | |
| f"(total now {len(jobs)})" | |
| ) | |
| return jobs, engine_used | |
| # ── API endpoints ───────────────────────────────────────────── | |
| async def get_schools(): | |
| """Return the full school registry used for job classification.""" | |
| return SCHOOLS | |
| async def get_locations(): | |
| """Return the city → geoId mapping used for location checkboxes.""" | |
| return INDIAN_METROS | |
| # ── Parallel multi-keyword HTTP scraper ─────────────────────── | |
| # Concurrency limiter — at most 1 parallel HTTP scrape to avoid | |
| # triggering LinkedIn's rate limiter. Sequential is safer and still fast. | |
| _PARALLEL_SEMAPHORE = asyncio.Semaphore(1) | |
| async def _fetch_all_pages_for_keyword( | |
| keyword: str, | |
| base_params: dict, | |
| cookies: dict, | |
| max_pages: int = 10, | |
| ) -> list[dict]: | |
| """Fast HTTP-only scraper for a single keyword. | |
| Designed for parallel fan-out: lightweight, no browser, shorter | |
| inter-page delays. Uses the shared semaphore to throttle concurrency. | |
| Retries on transient network errors with exponential backoff to handle | |
| temporary infrastructure issues. | |
| """ | |
| params = { | |
| k: v for k, v in base_params.items() | |
| if k in _PAGINATION_PARAM_KEYS and v not in (None, "") | |
| } | |
| params["keywords"] = keyword # override with this specific keyword | |
| referer_url = f"https://www.linkedin.com/jobs/search/?{urllib.parse.urlencode(params)}" | |
| headers = _build_pagination_headers(referer_url=referer_url) | |
| collected: list[dict] = [] | |
| async with _PARALLEL_SEMAPHORE: | |
| try: | |
| async with httpx.AsyncClient( | |
| timeout=20.0, | |
| cookies=cookies, | |
| headers=headers, | |
| follow_redirects=True, | |
| http2=False, | |
| ) as client: | |
| # Brief initial settle | |
| await asyncio.sleep(random.uniform(2.0, 4.0)) | |
| start = 0 | |
| consecutive_empty = 0 | |
| for page_idx in range(max_pages): | |
| url = f"{_SEE_MORE_URL}?{urllib.parse.urlencode({**params, 'start': start})}" | |
| logger.info(f"parallel[{keyword[:30]}]: GET start={start}") | |
| # Retry loop for transient network errors | |
| max_network_retries = 3 | |
| resp = None | |
| for attempt in range(1, max_network_retries + 1): | |
| try: | |
| resp = await client.get(url) | |
| break # Success, exit retry loop | |
| except (httpx.TimeoutException, httpx.NetworkError, OSError) as e: | |
| # Transient network errors | |
| if attempt < max_network_retries: | |
| backoff = min(2 ** attempt, 10.0) * random.uniform(1.0, 1.5) | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: transient network error " | |
| f"(attempt {attempt}/{max_network_retries}): {type(e).__name__}. " | |
| f"Retrying in {backoff:.1f}s..." | |
| ) | |
| await asyncio.sleep(backoff) | |
| else: | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: network error at start={start} " | |
| f"after {max_network_retries} attempts: {e}; " | |
| f"stopping with {len(collected)} jobs" | |
| ) | |
| break | |
| except httpx.HTTPError as e: | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: HTTP error at start={start}: {e}; stopping" | |
| ) | |
| break | |
| if resp is None: | |
| break # All retries failed, stop pagination | |
| # Handle rate limits | |
| if resp.status_code == 429: | |
| backoff = random.uniform(15.0, 25.0) | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: 429 at start={start}; " | |
| f"backing off {backoff:.1f}s" | |
| ) | |
| await asyncio.sleep(backoff) | |
| try: | |
| resp = await client.get(url) | |
| except (httpx.TimeoutException, httpx.NetworkError, OSError, httpx.HTTPError): | |
| break | |
| if resp.status_code >= 400: | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: status {resp.status_code}; stopping" | |
| ) | |
| break | |
| html = resp.text or "" | |
| if not html.strip() or _is_empty_stub(html): | |
| consecutive_empty += 1 | |
| if consecutive_empty >= 2: | |
| break | |
| # One empty might be transient — try next offset | |
| start += 25 | |
| await asyncio.sleep(random.uniform(0.5, 1.0)) | |
| continue | |
| consecutive_empty = 0 | |
| page_jobs = _parse_jobs_from_html(html) | |
| logger.info( | |
| f"parallel[{keyword[:30]}]: start={start} → {len(page_jobs)} jobs" | |
| ) | |
| if not page_jobs: | |
| break | |
| collected.extend(page_jobs) | |
| start += len(page_jobs) | |
| # Human-like delays between pages | |
| await asyncio.sleep(random.uniform(4.0, 7.0)) | |
| except Exception as e: | |
| logger.warning( | |
| f"parallel[{keyword[:30]}]: unexpected error: {e}; " | |
| f"returning {len(collected)} jobs" | |
| ) | |
| # Cool-down before releasing semaphore to next keyword | |
| await asyncio.sleep(random.uniform(4.0, 8.0)) | |
| logger.info(f"parallel[{keyword[:30]}]: done — {len(collected)} jobs total") | |
| return collected | |
| def _resolve_search(user_keyword: str) -> tuple[list[str], Optional[str]]: | |
| """Resolve user input into expanded search keywords + optional school filter. | |
| Resolution order: | |
| 1. Exact school code (case-insensitive): "SOCSE" → all SOCSE terms | |
| 2. Shortcut keyword: "software" → SOCSE terms | |
| 3. Direct passthrough: "Data Analyst" → ["Data Analyst"], no filter | |
| Returns: | |
| (search_keywords, school_code_to_filter_by_or_None) | |
| """ | |
| key = user_keyword.strip().lower() | |
| # 1. Exact school code match (SOCSE, solas, SOB, etc.) | |
| school_lower_map = {code.lower(): code for code in SCHOOLS} | |
| if key in school_lower_map: | |
| code = school_lower_map[key] | |
| terms = SCHOOL_SEARCH_TERMS.get(key, []) | |
| logger.info( | |
| f"resolve_search: school code '{code}' → {len(terms)} search terms" | |
| ) | |
| return terms if terms else [user_keyword], code | |
| # 2. Shortcut keyword → school (e.g. "software" → SOCSE) | |
| if key in KEYWORD_TO_SCHOOL: | |
| school_code = KEYWORD_TO_SCHOOL[key] | |
| terms = [] | |
| for school in school_code.split(","): | |
| for t in SCHOOL_SEARCH_TERMS.get(school.strip().lower(), []).copy(): | |
| if t not in terms: | |
| terms.append(t) | |
| # Guarantee the user's explicit shortcut keyword is also searched! | |
| # Sometimes shortcuts are exact job titles (e.g. "marketing", "data"). | |
| if key not in [t.lower() for t in terms]: | |
| # Insert at the beginning so it's prioritized | |
| terms.insert(0, user_keyword.strip()) | |
| logger.info( | |
| f"resolve_search: shortcut '{key}' → school {school_code} " | |
| f"({len(terms)} search terms including shortcut)" | |
| ) | |
| return terms if terms else [user_keyword], school_code | |
| # 3. Direct keyword — no expansion, no school filter | |
| logger.info(f"resolve_search: direct keyword '{user_keyword}' (no expansion)") | |
| return [user_keyword], None | |
| def _deduplicate_jobs(jobs: list[dict]) -> tuple[list[dict], int]: | |
| """Deduplicate jobs by canonical link. Returns (deduped_list, num_removed).""" | |
| seen: dict[str, dict] = {} | |
| untagged: list[dict] = [] | |
| for job in jobs: | |
| key = job.get("link") | |
| if key: | |
| if key not in seen: | |
| seen[key] = job | |
| else: | |
| untagged.append(job) | |
| deduped = list(seen.values()) + untagged | |
| return deduped, len(jobs) - len(deduped) | |
| async def scrape_internships( | |
| keywords: str = Query("Software Engineer", description="Job search keywords"), | |
| freshness: str = Query("r86400", description="Posting age: r3600=1h r86400=24h r604800=week r2592000=month"), | |
| work_types: Optional[list[str]] = Query( | |
| default=None, | |
| description="Work-type checkboxes — repeat per selection: onsite | remote | hybrid", | |
| ), | |
| job_types: Optional[list[str]] = Query( | |
| default=None, | |
| description="Job-type checkboxes — repeat per selection: internship | full_time", | |
| ), | |
| ): | |
| # ── Resolve keywords + optional school filter ────────────── | |
| all_keywords, school_filter = _resolve_search(keywords) | |
| if job_types and any(jt.lower() == "internship" for jt in job_types): | |
| all_keywords = [f"{kw} intern" for kw in all_keywords] | |
| keywords = f"{keywords} intern" | |
| is_multi = len(all_keywords) > 1 | |
| # Build the "primary" URL (used for the "Open in LinkedIn" link) | |
| primary_url = build_search_url( | |
| keywords=keywords, | |
| freshness=freshness, | |
| work_types=work_types, | |
| job_types=job_types, | |
| ) | |
| primary_params = build_search_params( | |
| keywords=keywords, | |
| freshness=freshness, | |
| work_types=work_types, | |
| job_types=job_types, | |
| ) | |
| async def _stream_generator(): | |
| """NDJSON streaming generator. | |
| Emits lines of JSON, one per event: | |
| {"type": "start", ...} | |
| {"type": "info", "message": ...} | |
| {"type": "jobs", "data": [...], ...} | |
| {"type": "done", ...} | |
| {"type": "error", "message": ...} | |
| """ | |
| all_jobs: list[dict] = [] | |
| engine_used = "http" | |
| total_dupes_removed = 0 | |
| total_filtered_out = 0 | |
| try: | |
| # ── Emit start event ────────────────────────────── | |
| yield json.dumps({ | |
| "type": "start", | |
| "engine": "http", | |
| "total_searches": len(all_keywords) if is_multi else 1, | |
| }) + "\n" | |
| if is_multi: | |
| school_label = f" (school: {school_filter})" if school_filter else "" | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": f"Expanding '{keywords}' → {len(all_keywords)} parallel searches{school_label}", | |
| }) + "\n" | |
| # ── Warm up cookies once ────────────────────────── | |
| cookies = await _warm_up_cookies(force=False) | |
| if not cookies: | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": "Cookie warm-up: re-warming...", | |
| }) + "\n" | |
| cookies = await _warm_up_cookies(force=True) | |
| if is_multi: | |
| # ══ PARALLEL MULTI-KEYWORD PATH ══════════════ | |
| # Fan out all keywords as concurrent tasks. | |
| # As each completes, stream its results immediately. | |
| seen_links: set = set() | |
| async def _scrape_and_collect(kw: str, index: int): | |
| """Scrape one keyword with a stagger delay to avoid 429s.""" | |
| # Stagger: each task waits before starting so requests are spread out. | |
| stagger = index * random.uniform(3.0, 6.0) | |
| if stagger > 0: | |
| await asyncio.sleep(stagger) | |
| jobs = await _fetch_all_pages_for_keyword( | |
| keyword=kw, | |
| base_params=primary_params, | |
| cookies=cookies, | |
| max_pages=10, | |
| ) | |
| return kw, jobs | |
| # Create all tasks | |
| tasks = [ | |
| asyncio.create_task(_scrape_and_collect(kw, i)) | |
| for i, kw in enumerate(all_keywords) | |
| ] | |
| completed = 0 | |
| for coro in asyncio.as_completed(tasks): | |
| try: | |
| kw, kw_jobs = await coro | |
| completed += 1 | |
| # Deduplicate against already-seen links | |
| new_jobs = [] | |
| for job in kw_jobs: | |
| link = job.get("link") | |
| if link: | |
| if link not in seen_links: | |
| seen_links.add(link) | |
| new_jobs.append(job) | |
| else: | |
| total_dupes_removed += 1 | |
| else: | |
| new_jobs.append(job) | |
| if school_filter and new_jobs: | |
| pre = len(new_jobs) | |
| filter_schools = [s.strip() for s in school_filter.split(",")] | |
| new_jobs = [ | |
| j for j in new_jobs | |
| if any(sf in j.get("schools", []) for sf in filter_schools) | |
| ] | |
| total_filtered_out += pre - len(new_jobs) | |
| # Strict internship title filter if only internship is selected | |
| if job_types and "internship" in job_types and "full_time" not in job_types: | |
| pre = len(new_jobs) | |
| intern_kw = ["intern", "trainee", "student", "co-op", "apprentice", "fellow"] | |
| new_jobs = [ | |
| j for j in new_jobs | |
| if any(ik in j.get("title", "").lower() for ik in intern_kw) | |
| ] | |
| total_filtered_out += pre - len(new_jobs) | |
| all_jobs.extend(new_jobs) | |
| logger.info( | |
| f"multi-search [{completed}/{len(all_keywords)}]: " | |
| f"'{kw}' → {len(kw_jobs)} raw, {len(new_jobs)} new " | |
| f"(total {len(all_jobs)})" | |
| ) | |
| # Stream the batch | |
| if new_jobs: | |
| yield json.dumps({ | |
| "type": "jobs", | |
| "data": new_jobs, | |
| "keyword": kw, | |
| "engine": "http", | |
| "deduplicated": total_dupes_removed, | |
| }) + "\n" | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": ( | |
| f"[{completed}/{len(all_keywords)}] " | |
| f"'{kw}' → {len(new_jobs)} new jobs " | |
| f"(total: {len(all_jobs)})" | |
| ), | |
| }) + "\n" | |
| except Exception as e: | |
| completed += 1 | |
| logger.warning(f"multi-search: task error: {e}") | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": f"[{completed}/{len(all_keywords)}] Search failed: {e}", | |
| }) + "\n" | |
| engine_used = "http" | |
| else: | |
| # ══ SINGLE-KEYWORD PATH ══════════════════════ | |
| # Original flow: HTTP → Selenium → Playwright | |
| jobs, engine_used = await _scrape_one_url( | |
| primary_url, params=primary_params, paginate=True | |
| ) | |
| if jobs: | |
| jobs, dupes = _deduplicate_jobs(jobs) | |
| total_dupes_removed = dupes | |
| if school_filter: | |
| pre = len(jobs) | |
| filter_schools = [s.strip() for s in school_filter.split(",")] | |
| jobs = [ | |
| j for j in jobs | |
| if any(sf in j.get("schools", []) for sf in filter_schools) | |
| ] | |
| total_filtered_out += pre - len(jobs) | |
| # Strict internship title filter if only internship is selected | |
| if job_types and "internship" in job_types and "full_time" not in job_types: | |
| pre = len(jobs) | |
| intern_kw = ["intern", "trainee", "student", "co-op", "apprentice", "fellow"] | |
| jobs = [ | |
| j for j in jobs | |
| if any(ik in j.get("title", "").lower() for ik in intern_kw) | |
| ] | |
| total_filtered_out += pre - len(jobs) | |
| all_jobs = jobs | |
| yield json.dumps({ | |
| "type": "jobs", | |
| "data": all_jobs, | |
| "engine": engine_used, | |
| "deduplicated": total_dupes_removed, | |
| }) + "\n" | |
| # ── Emit done event ─────────────────────────────── | |
| if not all_jobs: | |
| yield json.dumps({ | |
| "type": "error", | |
| "message": "All scraping engines failed — 0 results", | |
| }) + "\n" | |
| else: | |
| try: | |
| logger.info( | |
| f"serper: Searching hiring contacts for {len(all_jobs)} companies..." | |
| ) | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": f"Searching for hiring contacts via Google for {len(all_jobs)} companies...", | |
| }) + "\n" | |
| await _fetch_hiring_contacts_via_serper(all_jobs) | |
| found_count = sum( | |
| 1 for j in all_jobs if j.get("contact_details") | |
| ) | |
| logger.info( | |
| f"serper: Hiring contact search done — " | |
| f"{found_count}/{len(all_jobs)} jobs have a contact" | |
| ) | |
| # ── VERIFICATION PRINT (visible in terminal logs) ────── | |
| print(f"\n{'='*60}") | |
| print(f"[SERPER HIRING CONTACTS] {found_count}/{len(all_jobs)} jobs enriched") | |
| for j in all_jobs: | |
| contacts = j.get("contact_details") or [] | |
| if contacts: | |
| for c in contacts: | |
| if isinstance(c, dict): | |
| print( | |
| f" ✓ {j.get('title','?')} @ {j.get('company','?')}: " | |
| f"{c.get('name','?')} → {c.get('url','?')}" | |
| ) | |
| else: | |
| print(f" ✓ {j.get('title','?')} @ {j.get('company','?')}: {c}") | |
| print(f"{'='*60}\n") | |
| # ── Stream enriched jobs back to frontend ────────────── | |
| enriched_jobs = [j for j in all_jobs if j.get("contact_details")] | |
| if enriched_jobs: | |
| yield json.dumps({ | |
| "type": "jobs", | |
| "data": enriched_jobs, | |
| "engine": engine_used, | |
| "deduplicated": 0, | |
| "is_contact_update": True, | |
| }) + "\n" | |
| logger.info( | |
| f"serper: Streamed {len(enriched_jobs)} enriched jobs back to frontend" | |
| ) | |
| except Exception as e: | |
| logger.warning(f"serper: Hiring contact search failed: {e}") | |
| # Save jobs to database for student dashboard | |
| search_params = { | |
| 'keywords': keywords, | |
| 'freshness': freshness, | |
| 'work_types': work_types or [], | |
| 'job_types': job_types or [] | |
| } | |
| try: | |
| saved_count = save_scraped_jobs(all_jobs, search_params) | |
| logger.info(f"Saved {saved_count} jobs to database") | |
| except Exception as e: | |
| logger.warning(f"Failed to save jobs to database: {e}") | |
| yield json.dumps({ | |
| "type": "done", | |
| "total": len(all_jobs), | |
| "engine": engine_used, | |
| "url": primary_url, | |
| "deduplicated": total_dupes_removed, | |
| "filtered_out": total_filtered_out, | |
| "school_filter": school_filter, | |
| "searches_completed": len(all_keywords), | |
| }) + "\n" | |
| except Exception as e: | |
| logger.exception(f"scrape-internships stream error: {e}") | |
| # If we have collected any jobs, emit them first with partial results notification | |
| if all_jobs: | |
| yield json.dumps({ | |
| "type": "info", | |
| "message": f"Stream interrupted after collecting {len(all_jobs)} jobs", | |
| }) + "\n" | |
| yield json.dumps({ | |
| "type": "done", | |
| "total": len(all_jobs), | |
| "engine": engine_used, | |
| "partial": True, | |
| "error": str(e), | |
| "url": primary_url, | |
| "deduplicated": total_dupes_removed, | |
| "filtered_out": total_filtered_out, | |
| "school_filter": school_filter, | |
| "searches_completed": len(all_keywords), | |
| }) + "\n" | |
| else: | |
| # Only emit error if we have no jobs | |
| yield json.dumps({ | |
| "type": "error", | |
| "message": str(e), | |
| }) + "\n" | |
| global global_scrape | |
| if global_scrape.is_active: | |
| global_scrape.cancel() | |
| session = ScrapeSession() | |
| session.is_active = True | |
| global_scrape = session | |
| async def _pump(): | |
| search_params_for_db = { | |
| 'keywords': keywords, | |
| 'freshness': freshness, | |
| 'work_types': work_types or [], | |
| 'job_types': job_types or [] | |
| } | |
| try: | |
| async for item in _stream_generator(): | |
| # Intercept and incrementally save jobs to DB so they appear in student portal mid-scrape | |
| try: | |
| data = json.loads(item.strip()) | |
| if data.get("type") == "jobs" and data.get("data"): | |
| save_scraped_jobs(data["data"], search_params_for_db) | |
| except Exception as e: | |
| logger.warning(f"Incremental save failed: {e}") | |
| # Save to history and broadcast to active queues | |
| session.history.append(item) | |
| for q in session.queues: | |
| await q.put(item) | |
| except Exception as e: | |
| logger.error(f"Background scraper error: {e}") | |
| finally: | |
| session.is_active = False | |
| for q in session.queues: | |
| await q.put(None) | |
| session.queues.clear() | |
| # Start the background task so it continues even if the HTTP connection drops | |
| session.task = asyncio.create_task(_pump()) | |
| async def _consumer(): | |
| q = asyncio.Queue() | |
| session.queues.append(q) | |
| try: | |
| while True: | |
| item = await q.get() | |
| if item is None: | |
| break | |
| yield item | |
| except asyncio.CancelledError: | |
| logger.info("Client disconnected, background scraper will continue.") | |
| raise | |
| finally: | |
| if q in session.queues: | |
| session.queues.remove(q) | |
| return StreamingResponse( | |
| _consumer(), | |
| media_type="application/x-ndjson", | |
| headers={ | |
| "Cache-Control": "no-cache", | |
| "X-Accel-Buffering": "no", | |
| }, | |
| ) | |
| async def scrape_status(): | |
| """Return whether a scrape session is currently active or has history.""" | |
| return { | |
| "is_active": global_scrape.is_active, | |
| "has_history": len(global_scrape.history) > 0 | |
| } | |
| async def stream_scrape(): | |
| """Reconnect to the active scrape session or replay the last session's history.""" | |
| async def _reconnect_consumer(): | |
| # Yield all past events instantly | |
| for item in global_scrape.history: | |
| yield item | |
| # If the session is still active, listen for new events | |
| if global_scrape.is_active: | |
| q = asyncio.Queue() | |
| global_scrape.queues.append(q) | |
| try: | |
| while True: | |
| item = await q.get() | |
| if item is None: | |
| break | |
| yield item | |
| except asyncio.CancelledError: | |
| raise | |
| finally: | |
| if q in global_scrape.queues: | |
| global_scrape.queues.remove(q) | |
| return StreamingResponse( | |
| _reconnect_consumer(), | |
| media_type="application/x-ndjson", | |
| headers={ | |
| "Cache-Control": "no-cache", | |
| "X-Accel-Buffering": "no", | |
| }, | |
| ) | |
| # ── Student Dashboard Endpoints ─────────────────────────────── | |
| def get_recent_jobs(hours: int = Query(1, description="Hours ago (1 or 24)")): | |
| """Get jobs scraped within the last N hours for student dashboard""" | |
| try: | |
| if hours == 1: | |
| # Jobs from the last 1 hour | |
| jobs = get_jobs_in_timeframe(1, 0) | |
| elif hours == 24: | |
| # Jobs from 1-24 hours ago (excluding the last hour to avoid overlap) | |
| jobs = get_jobs_in_timeframe(24, 1) | |
| else: | |
| raise HTTPException(status_code=400, detail="Only 1 or 24 hours supported") | |
| return { | |
| "jobs": jobs, | |
| "count": len(jobs), | |
| "timeframe": f"{hours} hour{'s' if hours > 1 else ''} ago" | |
| } | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| def get_all_timeframes(background_tasks: BackgroundTasks): | |
| """Get jobs for all timeframes for student dashboard and auto-cleanup old jobs""" | |
| try: | |
| # Auto-delete jobs older than 30 days in the background | |
| background_tasks.add_task(cleanup_old_jobs, 30) | |
| # Fetch all jobs in one query and bin them to save DB roundtrips | |
| return get_binned_jobs() | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| def cleanup_old_data(days: int = Query(7, description="Remove jobs older than N days")): | |
| """Admin endpoint to cleanup old job data""" | |
| try: | |
| deleted_count = cleanup_old_jobs(days) | |
| return { | |
| "message": f"Cleaned up {deleted_count} jobs older than {days} days", | |
| "deleted_count": deleted_count | |
| } | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| async def trigger_auto_scrape(): | |
| """ | |
| Manually kick off the full school sweep. Returns immediately; | |
| progress can be polled at GET /admin/scrape-status. | |
| If a sweep is already running, returns 409. | |
| """ | |
| global _auto_scrape_task | |
| if _auto_scrape_state.get("is_running"): | |
| return { | |
| "status": "already_running", | |
| "message": "Auto-scrape is already in progress.", | |
| "state": _auto_scrape_state, | |
| } | |
| scheduler.add_job( | |
| _auto_scrape_all_schools, | |
| trigger=IntervalTrigger(hours=1), | |
| id="hourly_school_scrape", | |
| replace_existing=True, | |
| next_run_time=datetime.now(timezone.utc), | |
| misfire_grace_time=300, | |
| ) | |
| return { | |
| "status": "started", | |
| "message": f"Auto-scrape started for {len(_AUTO_SCRAPE_SCHOOLS)} schools (running hourly).", | |
| "schools": _AUTO_SCRAPE_SCHOOLS, | |
| } | |
| async def stop_auto_scrape(): | |
| """ | |
| Stop the currently running auto-scrape sweep and remove the hourly scheduled job. | |
| """ | |
| if scheduler.get_job("hourly_school_scrape"): | |
| scheduler.remove_job("hourly_school_scrape") | |
| _auto_scrape_state["cancel_requested"] = True | |
| global global_scrape | |
| if global_scrape.is_active: | |
| global_scrape.cancel() | |
| return {"status": "stopped", "message": "Auto-scrape has been stopped and scheduled jobs removed."} | |
| async def get_scrape_status(): | |
| """Poll the live progress of the currently-running (or last) auto-scrape sweep.""" | |
| state = dict(_auto_scrape_state) | |
| if state.get("started_at"): | |
| elapsed = (state.get("finished_at") or time.time()) - state["started_at"] | |
| state["elapsed_seconds"] = round(elapsed) | |
| total_done = len(state.get("completed", [])) + len(state.get("failed", [])) | |
| state["progress"] = f"{total_done}/{len(_AUTO_SCRAPE_SCHOOLS)}" | |
| return state | |
| async def get_scrape_history(): | |
| """Return the full history of all auto-scrape sweep runs (newest first). | |
| Each entry has: sweep_id, is_running, current_school, completed, failed, | |
| started_at, finished_at, cancelled, elapsed_seconds, progress. | |
| """ | |
| history = [] | |
| for entry in reversed(_auto_scrape_history): | |
| item = dict(entry) | |
| # Compute elapsed seconds | |
| if item.get("started_at"): | |
| elapsed = (item.get("finished_at") or time.time()) - item["started_at"] | |
| item["elapsed_seconds"] = round(elapsed) | |
| # Compute progress string | |
| total_done = len(item.get("completed", [])) + len(item.get("failed", [])) | |
| item["progress"] = f"{total_done}/{len(_AUTO_SCRAPE_SCHOOLS)}" | |
| item["total_schools"] = len(_AUTO_SCRAPE_SCHOOLS) | |
| item["schools"] = _AUTO_SCRAPE_SCHOOLS | |
| history.append(item) | |
| return {"history": history, "total_runs": len(history)} | |
| async def process_company_ratings(companies: list[str]) -> dict: | |
| if not companies: | |
| return {} | |
| SERPER_API_KEY = os.environ.get("SERPER_API_KEY", "0820e51e9080c289b3849e30189e023774fdad87") | |
| ratings = {} | |
| async def _fetch_company_rating(client: httpx.AsyncClient, company: str, semaphore: asyncio.Semaphore) -> tuple[str, float]: | |
| async with semaphore: | |
| await asyncio.sleep(random.uniform(0.5, 1.5)) | |
| query = f'"{company}" reviews Bangalore site:glassdoor.co.in OR site:ambitionbox.com' | |
| try: | |
| resp = await client.post( | |
| "https://google.serper.dev/search", | |
| json={"q": query}, | |
| headers={"X-API-KEY": SERPER_API_KEY, "Content-Type": "application/json"} | |
| ) | |
| data = resp.json() | |
| for org in data.get("organic", []): | |
| # Try Serper's built-in rating first | |
| if "rating" in org and isinstance(org["rating"], (int, float)): | |
| return company, float(org["rating"]) | |
| # Fallback to regex on snippet | |
| snippet = org.get("snippet", "") | |
| match = re.search(r"Rating:\s*([\d\.]+)", snippet, re.IGNORECASE) | |
| if match: | |
| return company, float(match.group(1)) | |
| except Exception as e: | |
| logger.warning(f"rate-companies: Serper search failed for {company}: {repr(e)}") | |
| return company, 2.5 | |
| try: | |
| # We limit concurrency to 3 to avoid overwhelming the network or API limits | |
| semaphore = asyncio.Semaphore(3) | |
| async with httpx.AsyncClient(timeout=15.0) as client: | |
| tasks = [_fetch_company_rating(client, c, semaphore) for c in companies] | |
| results = await asyncio.gather(*tasks) | |
| for company, rating in results: | |
| # Clamp to [1, 5] just in case | |
| ratings[company] = round(max(1.0, min(5.0, rating)), 1) | |
| # Fill in any companies missed with a neutral score | |
| for company in companies: | |
| if company not in ratings: | |
| ratings[company] = 2.5 | |
| return ratings | |
| except Exception as e: | |
| logger.error(f"rate-companies: process_company_ratings failed: {repr(e)}") | |
| raise e | |
| class RateCompaniesRequest(BaseModel): | |
| companies: list[str] | |
| async def rate_companies(req: RateCompaniesRequest, background_tasks: BackgroundTasks): | |
| """ | |
| Use Serper.dev (Google Search API) to estimate a reputation/internship | |
| quality rating (1.0–5.0) for each submitted company name (Bangalore context). | |
| A concurrent search is made for each company, extracting the rating from | |
| Glassdoor/Ambitionbox search snippets. | |
| """ | |
| try: | |
| ratings = await process_company_ratings(req.companies) | |
| # Persist ratings to DB in the background (non-blocking) | |
| background_tasks.add_task(_persist_ratings, ratings) | |
| logger.info(f"rate-companies: successfully rated {len(ratings)} companies via Search API") | |
| return {"ratings": ratings} | |
| except Exception as e: | |
| logger.error(f"rate-companies: Search API call failed: {repr(e)}") | |
| raise HTTPException(status_code=500, detail=f"Search rating failed: {str(e)}") | |
| def _persist_ratings(ratings: dict): | |
| """Write company ratings back to the DB (called as a background task).""" | |
| try: | |
| updated = update_company_ratings(ratings) | |
| logger.info(f"rate-companies: persisted ratings for {updated} job rows in DB") | |
| except Exception as e: | |
| logger.error(f"rate-companies: DB persist failed: {e}") | |
| # ── Manual Rating Management ───────────────────────────────────── | |
| class SetRatingRequest(BaseModel): | |
| company: str | |
| rating: float | |
| async def set_company_rating(req: SetRatingRequest): | |
| """Manually set a rating for a specific company (for testing/admin purposes).""" | |
| try: | |
| # Validate rating range | |
| rating = max(1.0, min(5.0, req.rating)) | |
| # Update the database | |
| updated = update_company_ratings({req.company: rating}) | |
| return { | |
| "message": f"Updated rating for '{req.company}' to {rating}", | |
| "company": req.company, | |
| "rating": rating, | |
| "rows_updated": updated | |
| } | |
| except Exception as e: | |
| logger.error(f"set-company-rating: failed: {e}") | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| async def get_company_ratings(): | |
| """Get all companies and their current ratings from the database.""" | |
| try: | |
| conn = _connect() | |
| try: | |
| with conn.cursor() as cur: | |
| cur.execute( | |
| """ | |
| SELECT DISTINCT company, company_rating | |
| FROM scraped_jobs | |
| WHERE company IS NOT NULL AND company_rating IS NOT NULL | |
| ORDER BY company | |
| """ | |
| ) | |
| rows = cur.fetchall() | |
| finally: | |
| conn.close() | |
| ratings = {row[0]: row[1] for row in rows} | |
| return {"ratings": ratings} | |
| except Exception as e: | |
| logger.error(f"get-company-ratings: failed: {e}") | |
| raise HTTPException(status_code=500, detail=str(e)) | |