internly / main.py
Sambhavvvvv's picture
Upload 2 files
5fb9df8 verified
Raw
History Blame Contribute Delete
96.2 kB
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}")
@app.on_event("startup")
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")
@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)
# ── 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 ─────────────────────────────────────────────
@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 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)
@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",
),
):
# ── 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",
},
)
@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():
# 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 ───────────────────────────────
@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 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))
@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:
# 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))
@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)
# 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]
@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)
# 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
@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:
# 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))
@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))