401 lines
15 KiB
Python
401 lines
15 KiB
Python
import os
|
|
import hashlib
|
|
import uuid
|
|
import datetime
|
|
import sqlite3
|
|
from pathlib import Path
|
|
|
|
def get_sqlite_path() -> Path:
|
|
base_dir = Path(__file__).resolve().parent.parent
|
|
sqlite_db = base_dir / "web" / "prisma" / "dev.db"
|
|
return sqlite_db
|
|
|
|
def generate_job_hash(job_url: str) -> str:
|
|
return hashlib.sha256(job_url.encode("utf-8")).hexdigest()
|
|
|
|
def is_quality_active_job(job: dict) -> bool:
|
|
"""
|
|
Anti-Ghost & Quality Filter Gate:
|
|
1. Title must be substantial and not placeholder text.
|
|
2. Company cannot be generic or unknown.
|
|
3. Description must have at least 150 characters of actual role details (rejects ghost stubs).
|
|
4. URL must be a valid canonical HTTP link.
|
|
5. Rejects typical recruiter spam strings.
|
|
"""
|
|
title = (job.get("title") or "").strip()
|
|
company = (job.get("company") or "").strip()
|
|
desc = (job.get("description") or "").strip()
|
|
url = (job.get("job_url") or "").strip()
|
|
|
|
if len(title) < 3 or title.lower() in ["untitled", "various", "general", "n/a"]:
|
|
return False
|
|
|
|
if len(company) < 2 or company.lower() in ["unknown", "confidential", "confidential company", "n/a"]:
|
|
return False
|
|
|
|
if len(desc) < 140:
|
|
return False
|
|
|
|
if not url.startswith("http"):
|
|
return False
|
|
|
|
# Filter out obvious third-party spam aggregators / ghost posting traps
|
|
low_desc = desc.lower()
|
|
spam_indicators = [
|
|
"earn up to $500/day stuffing envelopes",
|
|
"wire transfer assistant",
|
|
"mystery shopper wanted"
|
|
]
|
|
if any(s in low_desc for s in spam_indicators):
|
|
return False
|
|
|
|
return True
|
|
|
|
def upsert_jobs(jobs_list: list):
|
|
if not jobs_list:
|
|
return 0
|
|
|
|
# Apply quality filter gate
|
|
filtered_jobs = [j for j in jobs_list if is_quality_active_job(j)]
|
|
rejected_count = len(jobs_list) - len(filtered_jobs)
|
|
if rejected_count > 0:
|
|
print(f"[Quality Gate] Filtered out {rejected_count} low-quality/stub/ghost postings.")
|
|
|
|
db_url = os.getenv("DATABASE_URL", "")
|
|
if "postgresql" in db_url:
|
|
return _upsert_postgres(filtered_jobs, db_url)
|
|
else:
|
|
return _upsert_sqlite(filtered_jobs)
|
|
|
|
def _upsert_sqlite(jobs_list: list):
|
|
db_path = get_sqlite_path()
|
|
if not db_path.parent.exists():
|
|
db_path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
conn = sqlite3.connect(str(db_path))
|
|
cursor = conn.cursor()
|
|
|
|
query = """
|
|
INSERT INTO Job (
|
|
id, jobUrlHash, title, company, location, isRemote,
|
|
department, experienceLevel, description, salaryMin, salaryMax, jobUrl, source, datePosted,
|
|
createdAt, updatedAt
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT (jobUrlHash) DO UPDATE SET
|
|
title = excluded.title,
|
|
company = excluded.company,
|
|
location = excluded.location,
|
|
isRemote = excluded.isRemote,
|
|
department = COALESCE(excluded.department, Job.department),
|
|
experienceLevel = COALESCE(excluded.experienceLevel, Job.experienceLevel),
|
|
description = excluded.description,
|
|
salaryMin = COALESCE(excluded.salaryMin, Job.salaryMin),
|
|
salaryMax = COALESCE(excluded.salaryMax, Job.salaryMax),
|
|
datePosted = COALESCE(excluded.datePosted, Job.datePosted),
|
|
updatedAt = excluded.updatedAt;
|
|
"""
|
|
|
|
def format_iso(dt_val):
|
|
if not dt_val:
|
|
dt_val = datetime.datetime.now(datetime.timezone.utc)
|
|
if isinstance(dt_val, str):
|
|
dt_val = dt_val.strip()
|
|
if len(dt_val) == 10 and "-" in dt_val:
|
|
return dt_val + "T00:00:00.000Z"
|
|
return dt_val
|
|
if isinstance(dt_val, (datetime.datetime, datetime.date)):
|
|
return dt_val.isoformat()
|
|
return datetime.datetime.now(datetime.timezone.utc).isoformat()
|
|
|
|
now_iso = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
|
seen_hashes = set()
|
|
records = []
|
|
|
|
# 1. First ensure all distinct companies exist in Company table
|
|
distinct_companies = {}
|
|
for j in jobs_list:
|
|
c_name = (j.get("company") or "").strip()
|
|
c_loc = (j.get("location") or "").strip()
|
|
if c_name and len(c_name) >= 2 and c_name.lower() not in ["unknown", "confidential", "confidential company", "n/a"]:
|
|
if c_name not in distinct_companies:
|
|
distinct_companies[c_name] = c_loc
|
|
|
|
for c_name, c_loc in distinct_companies.items():
|
|
comp_id = "comp_" + hashlib.sha256(c_name.lower().encode("utf-8")).hexdigest()[:20]
|
|
cursor.execute("""
|
|
INSERT OR IGNORE INTO Company (
|
|
id, name, location, verificationStatus, trustStatus, createdAt, updatedAt
|
|
) VALUES (?, ?, ?, 'UNCLAIMED', 'UNVERIFIED', ?, ?)
|
|
""", (comp_id, c_name, c_loc or "USA", now_iso, now_iso))
|
|
|
|
# Fetch mapping of company name to id
|
|
cursor.execute("SELECT id, name FROM Company")
|
|
comp_map = {row[1]: row[0] for row in cursor.fetchall()}
|
|
|
|
for j in jobs_list:
|
|
job_url = j.get("job_url", "")
|
|
if not job_url:
|
|
continue
|
|
job_hash = generate_job_hash(job_url)
|
|
if job_hash in seen_hashes:
|
|
continue
|
|
seen_hashes.add(job_hash)
|
|
job_id = "job_" + str(uuid.uuid4()).replace("-", "")[:20]
|
|
c_name = j.get("company", "Unknown Company")[:255]
|
|
c_id = comp_map.get(c_name)
|
|
|
|
records.append((
|
|
job_id,
|
|
job_hash,
|
|
j.get("title", "Untitled Position")[:255],
|
|
c_name,
|
|
j.get("location", "Not Specified")[:255],
|
|
1 if j.get("is_remote") else 0,
|
|
j.get("department") or "Other",
|
|
j.get("experience_level") or "Mid-Level",
|
|
j.get("description", "") or "No description provided.",
|
|
j.get("salary_min"),
|
|
j.get("salary_max"),
|
|
job_url,
|
|
j.get("source", "jobspy"),
|
|
format_iso(j.get("date_posted")),
|
|
c_id,
|
|
now_iso,
|
|
now_iso
|
|
))
|
|
|
|
query = """
|
|
INSERT INTO Job (
|
|
id, jobUrlHash, title, company, location, isRemote,
|
|
department, experienceLevel, description, salaryMin, salaryMax, jobUrl, source, datePosted,
|
|
companyId, createdAt, updatedAt
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT (jobUrlHash) DO UPDATE SET
|
|
title = excluded.title,
|
|
company = excluded.company,
|
|
location = excluded.location,
|
|
isRemote = excluded.isRemote,
|
|
department = COALESCE(excluded.department, Job.department),
|
|
experienceLevel = COALESCE(excluded.experienceLevel, Job.experienceLevel),
|
|
description = excluded.description,
|
|
salaryMin = COALESCE(excluded.salaryMin, Job.salaryMin),
|
|
salaryMax = COALESCE(excluded.salaryMax, Job.salaryMax),
|
|
datePosted = COALESCE(excluded.datePosted, Job.datePosted),
|
|
companyId = COALESCE(excluded.companyId, Job.companyId),
|
|
updatedAt = excluded.updatedAt;
|
|
"""
|
|
|
|
try:
|
|
cursor.executemany(query, records)
|
|
# Self-healing sanitize step: reset any existing jobs falsely tagged as remote if location or title contains negative indicators
|
|
cursor.execute("""
|
|
UPDATE Job
|
|
SET isRemote = 0
|
|
WHERE isRemote = 1
|
|
AND (
|
|
LOWER(location) LIKE '%hybrid%'
|
|
OR LOWER(location) LIKE '%onsite%'
|
|
OR LOWER(location) LIKE '%on-site%'
|
|
OR LOWER(location) LIKE '%not remote%'
|
|
OR LOWER(location) LIKE '%in-office%'
|
|
OR LOWER(location) LIKE '%in office%'
|
|
OR LOWER(title) LIKE '%hybrid%'
|
|
OR LOWER(title) LIKE '%onsite%'
|
|
OR LOWER(title) LIKE '%on-site%'
|
|
OR LOWER(title) LIKE '%not remote%'
|
|
)
|
|
AND LOWER(location) NOT LIKE '%100% remote%'
|
|
AND LOWER(location) NOT LIKE '%fully remote%'
|
|
AND LOWER(title) NOT LIKE '%100% remote%'
|
|
AND LOWER(title) NOT LIKE '%fully remote%';
|
|
""")
|
|
conn.commit()
|
|
count = len(records)
|
|
cursor.close()
|
|
conn.close()
|
|
return count
|
|
except Exception as e:
|
|
cursor.close()
|
|
conn.close()
|
|
print(f"[SQLite Error] Failed to upsert: {e}")
|
|
return 0
|
|
|
|
def _upsert_postgres(jobs_list: list, db_url: str):
|
|
import psycopg2
|
|
from psycopg2.extras import execute_values
|
|
import time
|
|
|
|
if "?schema=" in db_url:
|
|
db_url = db_url.split("?schema=")[0]
|
|
|
|
conn = None
|
|
for attempt in range(40):
|
|
try:
|
|
conn = psycopg2.connect(db_url)
|
|
cursor = conn.cursor()
|
|
cursor.execute("SELECT to_regclass('\"Job\"');")
|
|
result = cursor.fetchone()
|
|
table_exists = result[0] if result else None
|
|
if table_exists:
|
|
break
|
|
cursor.close()
|
|
conn.close()
|
|
conn = None
|
|
print(f"[Scraper] Waiting for database tables to initialize (attempt {attempt + 1}/40)...")
|
|
time.sleep(5)
|
|
except Exception as conn_err:
|
|
print(f"[Scraper] Waiting for PostgreSQL readiness ({conn_err}) (attempt {attempt + 1}/40)...")
|
|
time.sleep(5)
|
|
|
|
if not conn or conn.closed:
|
|
try:
|
|
conn = psycopg2.connect(db_url)
|
|
cursor = conn.cursor()
|
|
except Exception as err:
|
|
print(f"[Postgres Error] Could not connect to database: {err}")
|
|
return 0
|
|
else:
|
|
cursor = conn.cursor()
|
|
|
|
now = datetime.datetime.now(datetime.timezone.utc)
|
|
|
|
# 1. First ensure all distinct companies exist in "Company" table
|
|
distinct_companies = {}
|
|
for j in jobs_list:
|
|
c_name = (j.get("company") or "").strip()
|
|
c_loc = (j.get("location") or "").strip()
|
|
if c_name and len(c_name) >= 2 and c_name.lower() not in ["unknown", "confidential", "confidential company", "n/a"]:
|
|
if c_name not in distinct_companies:
|
|
distinct_companies[c_name] = c_loc
|
|
|
|
if distinct_companies:
|
|
company_records = []
|
|
for c_name, c_loc in distinct_companies.items():
|
|
comp_id = "comp_" + hashlib.sha256(c_name.lower().encode("utf-8")).hexdigest()[:20]
|
|
company_records.append((
|
|
comp_id,
|
|
c_name,
|
|
c_loc or "USA",
|
|
"UNCLAIMED",
|
|
"UNVERIFIED",
|
|
now,
|
|
now
|
|
))
|
|
comp_upsert_query = """
|
|
INSERT INTO "Company" ("id", "name", "location", "verificationStatus", "trustStatus", "createdAt", "updatedAt")
|
|
VALUES %s
|
|
ON CONFLICT ("name") DO UPDATE SET
|
|
"location" = COALESCE("Company"."location", EXCLUDED."location"),
|
|
"updatedAt" = EXCLUDED."updatedAt";
|
|
"""
|
|
try:
|
|
execute_values(cursor, comp_upsert_query, company_records)
|
|
conn.commit()
|
|
except Exception as comp_err:
|
|
conn.rollback()
|
|
print(f"[Postgres Warning] Company pre-population error: {comp_err}")
|
|
|
|
# Fetch mapping of company name to id
|
|
comp_map = {}
|
|
try:
|
|
cursor.execute('SELECT "id", "name" FROM "Company"')
|
|
for row in cursor.fetchall():
|
|
comp_map[row[1]] = row[0]
|
|
except Exception as map_err:
|
|
print(f"[Postgres Warning] Failed to fetch company map: {map_err}")
|
|
|
|
query = """
|
|
INSERT INTO "Job" (
|
|
"id", "jobUrlHash", "title", "company", "location", "isRemote",
|
|
"department", "experienceLevel", "description", "salaryMin", "salaryMax", "jobUrl", "source", "datePosted",
|
|
"companyId", "lifecycleStatus", "lastSeenAt", "createdAt", "updatedAt"
|
|
) VALUES %s
|
|
ON CONFLICT ("jobUrlHash") DO UPDATE SET
|
|
"title" = EXCLUDED."title",
|
|
"company" = EXCLUDED."company",
|
|
"location" = EXCLUDED."location",
|
|
"isRemote" = EXCLUDED."isRemote",
|
|
"department" = COALESCE(EXCLUDED."department", "Job"."department"),
|
|
"experienceLevel" = COALESCE(EXCLUDED."experienceLevel", "Job"."experienceLevel"),
|
|
"description" = EXCLUDED."description",
|
|
"salaryMin" = COALESCE(EXCLUDED."salaryMin", "Job"."salaryMin"),
|
|
"salaryMax" = COALESCE(EXCLUDED."salaryMax", "Job"."salaryMax"),
|
|
"datePosted" = COALESCE(EXCLUDED."datePosted", "Job"."datePosted"),
|
|
"companyId" = COALESCE(EXCLUDED."companyId", "Job"."companyId"),
|
|
"lifecycleStatus" = 'ACTIVE',
|
|
"lastSeenAt" = NOW(),
|
|
"updatedAt" = NOW();
|
|
"""
|
|
|
|
records = []
|
|
seen_hashes = set()
|
|
|
|
for j in jobs_list:
|
|
job_url = j.get("job_url", "")
|
|
if not job_url:
|
|
continue
|
|
job_hash = generate_job_hash(job_url)
|
|
if job_hash in seen_hashes:
|
|
continue
|
|
seen_hashes.add(job_hash)
|
|
job_id = "job_" + str(uuid.uuid4()).replace("-", "")[:20]
|
|
c_name = j.get("company", "Unknown Company")[:255]
|
|
c_id = comp_map.get(c_name)
|
|
|
|
records.append((
|
|
job_id,
|
|
job_hash,
|
|
j.get("title", "Untitled Position")[:255],
|
|
c_name,
|
|
j.get("location", "Not Specified")[:255],
|
|
bool(j.get("is_remote", False)),
|
|
j.get("department") or "Other",
|
|
j.get("experience_level") or "Mid-Level",
|
|
j.get("description", "") or "No description provided.",
|
|
j.get("salary_min"),
|
|
j.get("salary_max"),
|
|
job_url,
|
|
j.get("source", "jobspy"),
|
|
now,
|
|
c_id,
|
|
"ACTIVE",
|
|
now,
|
|
now,
|
|
now
|
|
))
|
|
|
|
try:
|
|
execute_values(cursor, query, records)
|
|
# Self-healing sanitize step: reset any existing jobs falsely tagged as remote if location or title contains negative indicators
|
|
cursor.execute("""
|
|
UPDATE "Job"
|
|
SET "isRemote" = FALSE
|
|
WHERE "isRemote" = TRUE
|
|
AND (
|
|
LOWER("location") LIKE '%hybrid%'
|
|
OR LOWER("location") LIKE '%onsite%'
|
|
OR LOWER("location") LIKE '%on-site%'
|
|
OR LOWER("location") LIKE '%not remote%'
|
|
OR LOWER("location") LIKE '%in-office%'
|
|
OR LOWER("location") LIKE '%in office%'
|
|
OR LOWER("title") LIKE '%hybrid%'
|
|
OR LOWER("title") LIKE '%onsite%'
|
|
OR LOWER("title") LIKE '%on-site%'
|
|
OR LOWER("title") LIKE '%not remote%'
|
|
)
|
|
AND LOWER("location") NOT LIKE '%100% remote%'
|
|
AND LOWER("location") NOT LIKE '%fully remote%'
|
|
AND LOWER("title") NOT LIKE '%100% remote%'
|
|
AND LOWER("title") NOT LIKE '%fully remote%';
|
|
""")
|
|
conn.commit()
|
|
count = len(records)
|
|
cursor.close()
|
|
conn.close()
|
|
return count
|
|
except Exception as e:
|
|
conn.rollback()
|
|
cursor.close()
|
|
conn.close()
|
|
print(f"[Postgres Error] Failed to upsert: {e}")
|
|
return 0
|