Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions .github/workflows/scrape.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name: scrape

on:
schedule:
# Offset from the top of the hour so it doesn't double up with the
# Railway :00 scheduler (and so a GH runner failure is distinguishable
# from a Railway one in the logs).
- cron: '17 * * * *'
workflow_dispatch:

jobs:
scrape:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: '3.12'
- name: Install backend deps
run: pip install -r backend/requirements.txt
- name: Scrape
working-directory: backend
run: python scraper.py
env:
DATABASE_URL: ${{ secrets.DATABASE_URL }}
22 changes: 21 additions & 1 deletion backend/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
import re
import base64
import requests
from datetime import datetime, timezone
from pathlib import Path
from urllib.parse import quote_plus, quote

Expand Down Expand Up @@ -153,20 +154,32 @@ def get_agent_identity(authorization: str = Header(None)):

logger = logging.getLogger(__name__)

# Last scrape outcome (in-process scheduler and manual /update), surfaced in
# /health so a silent scraper death shows up without digging through logs.
_last_scrape: dict = {}

def _record_scrape(result, error=None):
_last_scrape["ran_at"] = datetime.now(timezone.utc).isoformat()
_last_scrape["result"] = result if result is not None else "skipped (locked / OOM)"
_last_scrape["error"] = error

def run_scrape():
return scraper.update_database()

def scheduled_scrape():
try:
result = run_scrape()
if result is None:
_record_scrape(None)
logger.info("Scheduler: scrape skipped (already running, or out of memory)")
else:
_record_scrape(result)
# No invalidate_cache() here: in-process runs rewrite the snapshot
# file themselves (see scraper._update_database), so the next
# /recent picks the fresh data up by identity.
logger.info(f"Scheduler: {result}")
except Exception as e:
_record_scrape(None, error=str(e))
logger.error(f"Scheduler: scrape failed — {e}")

scheduler = BackgroundScheduler()
Expand Down Expand Up @@ -249,6 +262,7 @@ def health():
"next_scrape": str(next_run) if next_run else "unknown",
"rss_mb": _proc_status_mb("VmRSS:"),
"peak_rss_mb": _proc_status_mb("VmHWM:"),
"last_scrape": _last_scrape,
}

#Lists the data sources the Listing feed pulls from (fetched from the
Expand Down Expand Up @@ -302,9 +316,15 @@ def sources(request: Request):
@app.post("/update")
@limiter.limit("5/minute")
def update_base(request: Request, verified=Depends(verify_key)):
result = run_scrape()
try:
result = run_scrape()
except Exception as e:
_record_scrape(None, error=str(e))
raise HTTPException(status_code=500, detail=str(e))
if result is None:
_record_scrape(None)
raise HTTPException(status_code=409, detail="Scrape already running.")
_record_scrape(result)
read_db.invalidate_cache()
return {"result": read_db.recent_internships()}

Expand Down
173 changes: 159 additions & 14 deletions backend/scraper.py
Original file line number Diff line number Diff line change
Expand Up @@ -490,7 +490,7 @@ def update_database():
_UPSERT_SQL = """
INSERT INTO internships (company, role, location, date, link, type, season, ats, description, fingerprint, last_seen_at)
VALUES %s
ON CONFLICT (company, role, location, link)
ON CONFLICT (fingerprint)
DO UPDATE SET
date = EXCLUDED.date,
type = EXCLUDED.type,
Expand All @@ -502,6 +502,19 @@ def update_database():
"""


# Conflict target is `fingerprint`, not (company, role, location, link).
# fingerprint is norm(company)|norm(role)|norm(location), so any four-column
# conflict necessarily implies a fingerprint conflict -- arbitrating on the
# fingerprint covers both, which matters because Postgres accepts exactly one
# inference specification per ON CONFLICT and prod has two unique constraints
# on this table. Targeting only the four columns let a same-job/different-link
# listing reach internships_fingerprint_key and abort the scrape with
# UniqueViolation; a bare "ON CONFLICT" is rejected outright, since DO UPDATE
# requires an inference specification (only DO NOTHING may omit it).
# The SET list omits company/role/location/link, so a conflicting row keeps
# its identity -- its existing link survives.


def _iter_source(fetch, entry):
"""One source's rows, lazily. The feed streams; the README parsers return a
list, and yielding from one costs nothing."""
Expand Down Expand Up @@ -546,30 +559,86 @@ def _flush(cursor, batch, run_time):


def _backfill_fingerprints(cursor):
"""Fill fingerprint on rows written before the column existed.
"""Make every row's fingerprint the normaliser's own value for its columns.

In Python rather than SQL, deliberately. The value has to come out of the
exact same normaliser the upsert uses, and a SQL translation of those
regexes would quietly diverge from it — per-field .strip() versus a
single btrim over the joined string, for one. A mismatched fingerprint is
not a crash; it is a 404 on every tracked job until that row is re-scraped,
which is the worst possible failure for the feature that motivated this.
Idempotent: a no-op once every row is populated.

Every row, not just the NULL ones. Non-NULL is not the same as correct, and
being wrong is a hard failure rather than a cosmetic one: the upsert
arbitrates on (fingerprint) and nothing else, so a row whose stored
fingerprint is not what job_fingerprint computes for its own columns is
invisible to that arbiter while still holding
(company, role, location, link). An insert matching it there aborts the
run on internships_unique_job — which is exactly how 52 rows written by a
normaliser that does not exist anywhere in this repo ("intelligent
creation, camera" kept its space, "NYC" collapsed to "ny") killed every
scrape after the backfill was thought finished. Recomputing is a no-op for
a correct row, so the cost of covering them is one comparison.

Greedy, because prod holds rows whose content normalises to a fingerprint
another row already holds ("Acme, Inc." alongside "Acme Inc."), and the
unique index over `fingerprint` forbids two rows sharing one. Filling every
row blindly aborted the scrape on psycopg2.errors.UniqueViolation.

Those collisions are deleted rather than left stale. Leaving them NULL was
tried first and does not work: a NULL is invisible to ON CONFLICT
(fingerprint), so a later insert matching one of those rows on
(company, role, location, link) found no arbiter and died on
internships_unique_job instead. One error traded for another.

Deleting is the same reconciliation _ensure_schema already performs when
internships_unique_job's definition changes -- keep the lowest id, drop the
later copy. The surviving row is the one the frontend tracker already
resolved to, since a drifted or absent fingerprint is a 404 on every
tracked job. Nothing references internships by foreign key; job references
live in agent_proposals.payload as JSONB.
"""
# ORDER BY id so "keep the oldest" is deterministic rather than dependent
# on whatever order the heap happens to return.
cursor.execute(
"SELECT id, company, role, location FROM internships WHERE fingerprint IS NULL"
"SELECT id, company, role, location, fingerprint FROM internships ORDER BY id"
)
rows = cursor.fetchall()
if not rows:
return 0
execute_values(
cursor,
"UPDATE internships AS i SET fingerprint = v.fp "
"FROM (VALUES %s) AS v(id, fp) WHERE i.id = v.id",
((r[0], job_fingerprint(r[1], r[2], r[3])) for r in rows),
page_size=1000,
)
return len(rows)
updates = []
doomed = []
# Assigned in id order, seeded with nothing: every row passes through, so
# "already taken" means an earlier row, not a pre-existing value. Seeding
# from the stored column instead is what let a drifted row block its own
# correction.
taken = set()
for rid, company, role, loc, stored in rows:
fp = job_fingerprint(company, role, loc)
if fp == stored:
taken.add(fp)
continue
if fp in taken:
doomed.append(rid)
continue
taken.add(fp)
updates.append((rid, fp))
if doomed:
cursor.execute("DELETE FROM internships WHERE id = ANY(%s)", (doomed,))
print(
f" deleted {len(doomed)} duplicate row(s) sharing a fingerprint "
f"with an existing row (same job listed twice)",
flush=True,
)
if updates:
execute_values(
cursor,
"UPDATE internships AS i SET fingerprint = v.fp "
"FROM (VALUES %s) AS v(id, fp) WHERE i.id = v.id",
iter(updates),
page_size=1000,
)
return len(updates)


def _update_database():
Expand Down Expand Up @@ -645,7 +714,57 @@ def _update_database():

backfilled = _backfill_fingerprints(cursor)
if backfilled:
print(f" backfilled fingerprint on {backfilled} pre-existing rows", flush=True)
print(
f" reconciled fingerprint on {backfilled} row(s) whose stored "
f"value was not the normaliser's own value for their columns",
flush=True,
)

# What the upsert is actually arbitrating against. Worth one line per
# run: the fingerprint arbiter is the load-bearing piece of the write
# path, and "is it unique, and is it partial" is not something you want
# to infer from a UniqueViolation three layers up.
cursor.execute("""
SELECT indexname, indexdef FROM pg_indexes
WHERE tablename = 'internships' AND indexdef ILIKE '%UNIQUE%'
ORDER BY indexname
""")
for name, d in cursor.fetchall():
print(f" unique: {d}", flush=True)
cursor.execute(
"SELECT count(*) FROM internships WHERE fingerprint IS NULL"
)
print(f" rows still lacking a fingerprint: {cursor.fetchone()[0]}", flush=True)

# The upsert arbitrates on (fingerprint), so a unique index over that
# column has to exist -- ON CONFLICT infers from a unique index or an
# exclusion constraint, and from nothing else. Prod has had one since
# someone added it by hand, so a fresh or restored database failed
# every scrape with "no unique or exclusion constraint matching the
# ON CONFLICT specification". Created after the backfill, so the rows
# it covers are already populated.
#
# CREATE UNIQUE INDEX IF NOT EXISTS rather than ADD CONSTRAINT: the
# existing object is an index, so a constraint would collide on the
# name (DuplicateTable) even though the guard found nothing to add.
# IF NOT EXISTS is satisfied by either form, since a constraint is
# backed by an index of the same name.
try:
cursor.execute("""
CREATE UNIQUE INDEX IF NOT EXISTS internships_fingerprint_key
ON internships(fingerprint)
""")
except psycopg2.errors.UniqueViolation:
# Leftover content duplicates. The scrape still dedupes in Python
# via `seen`, so carry on without the index rather than failing
# every run until the rows are cleaned up by hand.
cursor.connection.rollback()
print(
" skipped the fingerprint unique index: duplicate fingerprints "
"still present (dedup continues without it)",

flush=True,
)

# Age purge is safe to run first: it keys off the stored date, not off
# anything a source told us this cycle.
Expand All @@ -669,7 +788,14 @@ def _update_database():
writing = False
try:
for job in _iter_source(fetch, (url, job_type, season)):
key = (job["company"].lower(), job["role"].lower(), job["location"].lower())
# Keyed on the fingerprint, not the raw lowercase triple,
# so this dedup and internships_fingerprint_key agree.
# They disagreed: "Acme, Inc." and "Acme Inc." are
# distinct here but collapse to one fingerprint, which
# would put both in one batch and make Postgres raise
# "ON CONFLICT DO UPDATE command cannot affect row a
# second time".
key = job_fingerprint(job["company"], job["role"], job["location"])
if key in seen:
continue
seen.add(key)
Expand All @@ -695,6 +821,25 @@ def _update_database():
# A failed write is not a failed source. Recording it and
# carrying on would report a successful run over a
# half-written table, and the purge would run.
#
# Print the server's own DETAIL verbatim. An earlier
# attempt parsed it with a regex anchored on "link=(",
# which never appears -- Postgres writes
# "Key (company, role, location, link)=(...) already
# exists", so the parse silently matched nothing and
# the diagnostic printed nothing on the one run it
# existed for. The server states the conflict better
# than a parser would.
detail = getattr(getattr(e, "diag", None), "detail_text", "") or ""
if detail:
print(f" ! write conflict -- {detail.strip()}", flush=True)
print(
" ! if the reconcile count above was 0, a stored "
"fingerprint that is not the normaliser's own value "
"for its row is what is invisible to "
"ON CONFLICT (fingerprint)",
flush=True,
)
raise
# One 429 or one empty source must not abort the run — that
# skips the upsert for every other source too. Recorded and
Expand Down
Loading
Loading