JB/scraper/db.py

209 lines
7.3 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 upsert_jobs(jobs_list: list):
if not jobs_list:
return 0
db_url = os.getenv("DATABASE_URL", "")
if "postgresql" in db_url:
return _upsert_postgres(jobs_list, db_url)
else:
return _upsert_sqlite(jobs_list)
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"
if dt_val.endswith("+00:00"):
return dt_val[:-6] + "Z"
if not dt_val.endswith("Z"):
return dt_val + "Z"
return dt_val
if isinstance(dt_val, (datetime.datetime, datetime.date)):
if isinstance(dt_val, datetime.date) and not isinstance(dt_val, datetime.datetime):
dt_val = datetime.datetime.combine(dt_val, datetime.time.min)
return dt_val.strftime("%Y-%m-%dT%H:%M:%S.000Z")
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%S.000Z")
now_iso = format_iso(datetime.datetime.now(datetime.timezone.utc))
inserted_count = 0
for j in jobs_list:
job_url = j.get("job_url", "")
if not job_url:
continue
job_hash = generate_job_hash(job_url)
job_id = "job_" + str(uuid.uuid4()).replace("-", "")[:20]
date_posted_raw = j.get("date_posted")
date_posted_str = format_iso(date_posted_raw)
try:
cursor.execute(query, (
job_id,
job_hash,
j.get("title", "Untitled Position")[:255],
j.get("company", "Unknown Company")[:255],
j.get("location", "Not Specified")[:255],
1 if j.get("is_remote", False) 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"),
date_posted_str,
now_iso,
now_iso
))
inserted_count += 1
except Exception as e:
print(f"[SQLite Warning] Failed to insert job {job_hash[:8]}: {e}")
conn.commit()
conn.close()
return inserted_count
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]
# Wait for postgres and Prisma migrations to be ready
conn = None
for attempt in range(15):
try:
conn = psycopg2.connect(db_url)
cursor = conn.cursor()
cursor.execute("SELECT to_regclass('\"Job\"');")
table_exists = cursor.fetchone()[0]
if table_exists:
break
cursor.close()
conn.close()
print(f"[Scraper] Waiting for database tables to initialize (attempt {attempt + 1}/15)...")
time.sleep(3)
except Exception as conn_err:
print(f"[Scraper] Waiting for PostgreSQL readiness ({conn_err}) (attempt {attempt + 1}/15)...")
time.sleep(3)
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()
query = """
INSERT INTO "Job" (
"id", "jobUrlHash", "title", "company", "location", "isRemote",
"department", "experienceLevel", "description", "salaryMin", "salaryMax", "jobUrl", "source", "datePosted",
"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"),
"updatedAt" = NOW();
"""
records = []
now = datetime.datetime.now(datetime.timezone.utc)
for j in jobs_list:
job_url = j.get("job_url", "")
if not job_url:
continue
job_hash = generate_job_hash(job_url)
job_id = "job_" + str(uuid.uuid4()).replace("-", "")[:20]
records.append((
job_id,
job_hash,
j.get("title", "Untitled Position")[:255],
j.get("company", "Unknown Company")[:255],
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,
now,
now
))
try:
execute_values(cursor, query, records)
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