| 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()
|
|
|
| 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
|
|
|
|
|
| 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
|
|
|
|
|
| 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.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=["*"],
|
| )
|
|
|
|
|
| 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: 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: list[dict] = []
|
|
|
|
|
|
|
| _SELF_PORT = int(os.environ.get("PORT", "8000"))
|
| _SELF_BASE_URL = f"http://localhost:{_SELF_PORT}"
|
|
|
| se_driver = None
|
| pw_browser = None
|
| pw_stealth_ctx = None
|
|
|
|
|
| _warm_cookies: dict = {}
|
| _warm_cookies_ts: float = 0.0
|
| _COOKIE_TTL_SECONDS = 600
|
|
|
|
|
|
|
|
|
| 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
|
|
|
| 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
|
|
|
| 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
|
|
|
| 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:
|
|
|
| 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
|
|
|
|
|
|
|
|
|
| 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 [
|
|
|
| f"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36",
|
|
|
| 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",
|
|
|
| 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",
|
|
|
| f"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/{v}.0.0.0 Safari/537.36",
|
|
|
| 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] = []
|
|
|
|
|
|
|
|
|
| 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()
|
|
|
| if len(stripped) < 60 and "<!--" in stripped:
|
| return True
|
|
|
| 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,
|
| }
|
|
|
|
|
| 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()
|
|
|
|
|
| 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}")
|
|
|
|
|
| @app.on_event("startup")
|
| async def startup_event():
|
| global se_driver, pw_browser, pw_stealth_ctx, _UA_POOL
|
| loop = asyncio.get_event_loop()
|
|
|
|
|
| _UA_POOL = _build_ua_pool()
|
| logger.info(f"β UA pool built ({len(_UA_POOL)} variants, Chrome v{_get_chrome_major_version() or '?'})")
|
|
|
|
|
| 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}")
|
|
|
|
|
| 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")
|
|
|
|
|
|
|
| @app.on_event("shutdown")
|
| 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)
|
|
|
|
|
|
|
|
|
| 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))
|
| 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
|
|
|
|
|
| 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 = {}
|
|
|
|
|
| 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}")
|
|
|
|
|
| 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
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| _FIXED_PARAMS: dict[str, str] = {
|
| "sortBy": "R",
|
| "f_E": "1",
|
| "geoId": "90009633",
|
| }
|
|
|
|
|
| WORK_TYPES: dict[str, str] = {
|
| "onsite": "1",
|
| "remote": "2",
|
| "hybrid": "3",
|
| }
|
|
|
|
|
|
|
| JOB_TYPES: dict[str, str] = {
|
| "internship": "I",
|
| "full_time": "F",
|
| }
|
|
|
|
|
| FRESHNESS_PRESETS: dict[str, str] = {
|
| "hour": "r3600",
|
| "day": "r86400",
|
| "week": "r604800",
|
| "month": "r2592000",
|
| }
|
|
|
|
|
|
|
|
|
|
|
|
|
| KEYWORD_TO_SCHOOL: dict[str, str] = {
|
|
|
| "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",
|
|
|
| "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",
|
|
|
| "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",
|
|
|
| "economics": "SoEPP",
|
| "policy": "SoEPP",
|
| "governance": "SoEPP",
|
| "public policy": "SoEPP",
|
| "research": "SoEPP",
|
| "development studies":"SoEPP",
|
| "economist": "SoEPP",
|
| "policy analyst": "SoEPP",
|
| "public relations": "SoEPP",
|
| "government": "SoEPP",
|
|
|
| "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",
|
|
|
| "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",
|
|
|
| "psychology": "SOLAS",
|
| "environment": "SOLAS",
|
| "liberal arts": "SOLAS",
|
| "sociology": "SOLAS",
|
| "history": "SOLAS",
|
| "behavioral science":"SOLAS",
|
| "psychology intern": "SOLAS",
|
| "counseling": "SOLAS",
|
| "sociologist": "SOLAS",
|
| "research assistant":"SOLAS",
|
|
|
| "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",
|
| }
|
|
|
|
|
|
|
| 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,
|
| )
|
|
|
| return f"https://www.linkedin.com/jobs/search/?{urllib.parse.urlencode(params)}"
|
|
|
|
|
|
|
|
|
| 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 = []
|
|
|
|
|
| cards = soup.select("ul.jobs-search__results-list > li")
|
|
|
|
|
|
|
| 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_el = (
|
| card.select_one(".base-search-card__title")
|
| or card.select_one("h3")
|
| or card.select_one("[class*='title']")
|
| )
|
|
|
| company_el = (
|
| card.select_one(".base-search-card__subtitle a")
|
| or card.select_one(".base-search-card__subtitle")
|
| or card.select_one("h4 a")
|
| )
|
|
|
| loc_el = (
|
| card.select_one(".job-search-card__location")
|
| or card.select_one("[class*='location']")
|
| )
|
|
|
| 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']")
|
| )
|
|
|
| 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
|
|
|
|
|
| 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": [],
|
| })
|
| except Exception:
|
| continue
|
|
|
| return jobs
|
|
|
|
|
|
|
|
|
|
|
| _SERPER_API_KEY = os.environ.get("SERPER_API_KEY", "0820e51e9080c289b3849e30189e023774fdad87")
|
| _SERPER_URL = "https://google.serper.dev/search"
|
|
|
|
|
| _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.
|
| """
|
|
|
| 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]}
|
|
|
|
|
| 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}
|
|
|
|
|
| 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
|
|
|
|
|
| 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])
|
|
|
|
|
|
|
|
|
|
|
|
|
| _SEE_MORE_URL = (
|
| "https://www.linkedin.com/jobs-guest/jobs/api/seeMoreJobPostings/search"
|
| )
|
|
|
|
|
|
|
|
|
| _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
|
| 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:
|
|
|
| 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}")
|
|
|
|
|
| max_network_retries = 3
|
| for attempt in range(1, max_network_retries + 1):
|
| try:
|
| resp = await client.get(url)
|
| break
|
| except (httpx.TimeoutException, httpx.NetworkError, OSError) as e:
|
|
|
| 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
|
|
|
|
|
|
|
| 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
|
|
|
|
|
|
|
|
|
|
|
| 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)})"
|
| )
|
|
|
|
|
| start += len(page_jobs)
|
|
|
|
|
|
|
| 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
|
|
|
|
|
|
|
|
|
| _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)
|
|
|
|
|
|
|
| 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 [], {}
|
|
|
|
|
| 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)
|
|
|
|
|
| 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)
|
|
|
|
|
| try:
|
| await page.wait_for_selector(_JOB_CARD_CSS, timeout=15000)
|
| except Exception:
|
| logger.warning("Playwright: no job cards appeared within 15 s")
|
| return [], {}
|
|
|
|
|
| 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)
|
|
|
|
|
| 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()
|
|
|
|
|
|
|
|
|
| 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 = {}
|
|
|
|
|
| if paginate and params:
|
| max_http_attempts = 3
|
| for attempt in range(1, max_http_attempts + 1):
|
|
|
| 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
|
|
|
|
|
| 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"
|
| )
|
|
|
|
|
| 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}")
|
|
|
|
|
|
|
|
|
|
|
| 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),
|
| )
|
| 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
|
|
|
|
|
|
|
|
|
| @app.get("/schools")
|
| async def get_schools():
|
| """Return the full school registry used for job classification."""
|
| return SCHOOLS
|
|
|
|
|
| @app.get("/locations")
|
| async def get_locations():
|
| """Return the city β geoId mapping used for location checkboxes."""
|
| return INDIAN_METROS
|
|
|
|
|
|
|
|
|
|
|
|
|
| _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
|
|
|
| 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:
|
|
|
| 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}")
|
|
|
|
|
| max_network_retries = 3
|
| resp = None
|
| for attempt in range(1, max_network_retries + 1):
|
| try:
|
| resp = await client.get(url)
|
| break
|
| except (httpx.TimeoutException, httpx.NetworkError, OSError) as e:
|
|
|
| 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
|
|
|
|
|
| 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
|
|
|
| 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)
|
|
|
|
|
| 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"
|
| )
|
|
|
|
|
| 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()
|
|
|
|
|
| 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
|
|
|
|
|
| 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)
|
|
|
|
|
|
|
| if key not in [t.lower() for t in terms]:
|
|
|
| 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
|
|
|
|
|
| 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)
|
|
|
|
|
|
|
| @app.post("/scrape-internships")
|
| 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",
|
| ),
|
| ):
|
|
|
| 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
|
|
|
|
|
| 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:
|
|
|
| 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"
|
|
|
|
|
| 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:
|
|
|
|
|
|
|
|
|
| seen_links: set = set()
|
|
|
| async def _scrape_and_collect(kw: str, index: int):
|
| """Scrape one keyword with a stagger delay to avoid 429s."""
|
|
|
| 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
|
|
|
|
|
| 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
|
|
|
|
|
| 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)
|
|
|
|
|
| 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") or "").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)})"
|
| )
|
|
|
|
|
| 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:
|
|
|
|
|
| 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)
|
|
|
|
|
| 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") or "").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"
|
|
|
|
|
| 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"
|
| )
|
|
|
| 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")
|
|
|
|
|
| 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}")
|
|
|
|
|
| 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 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:
|
|
|
| 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():
|
|
|
| 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}")
|
|
|
|
|
| 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()
|
|
|
|
|
| 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",
|
| },
|
| )
|
|
|
| @app.get("/scrape-internships/status")
|
| 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
|
| }
|
|
|
| @app.get("/scrape-internships/stream")
|
| async def stream_scrape():
|
| """Reconnect to the active scrape session or replay the last session's history."""
|
| async def _reconnect_consumer():
|
|
|
| for item in global_scrape.history:
|
| yield item
|
|
|
|
|
| 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",
|
| },
|
| )
|
|
|
|
|
|
|
|
|
| @app.get("/student/jobs/recent")
|
| 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 = get_jobs_in_timeframe(1, 0)
|
| elif hours == 24:
|
|
|
| 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))
|
|
|
| @app.get("/student/jobs/all-timeframes")
|
| def get_all_timeframes(background_tasks: BackgroundTasks):
|
| """Get jobs for all timeframes for student dashboard and auto-cleanup old jobs"""
|
| try:
|
|
|
| background_tasks.add_task(cleanup_old_jobs, 30)
|
|
|
|
|
| return get_binned_jobs()
|
| except Exception as e:
|
| raise HTTPException(status_code=500, detail=str(e))
|
|
|
| @app.post("/admin/cleanup")
|
| 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))
|
|
|
|
|
| @app.post("/admin/trigger-scrape")
|
| 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,
|
| }
|
|
|
|
|
| @app.post("/admin/stop-scrape")
|
| 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."}
|
|
|
|
|
| @app.get("/admin/scrape-status")
|
| 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
|
|
|
|
|
| @app.get("/admin/scrape-history")
|
| 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)
|
|
|
| if item.get("started_at"):
|
| elapsed = (item.get("finished_at") or time.time()) - item["started_at"]
|
| item["elapsed_seconds"] = round(elapsed)
|
|
|
| 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", []):
|
|
|
| if "rating" in org and isinstance(org["rating"], (int, float)):
|
| return company, float(org["rating"])
|
|
|
| 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:
|
|
|
| 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:
|
|
|
| ratings[company] = round(max(1.0, min(5.0, rating)), 1)
|
|
|
|
|
| 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]
|
|
|
|
|
| @app.post("/rate-companies")
|
| 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)
|
|
|
| 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}")
|
|
|
|
|
|
|
|
|
| class SetRatingRequest(BaseModel):
|
| company: str
|
| rating: float
|
|
|
| @app.post("/set-company-rating")
|
| async def set_company_rating(req: SetRatingRequest):
|
| """Manually set a rating for a specific company (for testing/admin purposes)."""
|
| try:
|
|
|
| rating = max(1.0, min(5.0, req.rating))
|
|
|
|
|
| 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))
|
|
|
|
|
| @app.get("/get-company-ratings")
|
| 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))
|
|
|
|
|