diff --git a/.github/workflows/scrape.yml b/.github/workflows/scrape.yml new file mode 100644 index 0000000..4d50b5c --- /dev/null +++ b/.github/workflows/scrape.yml @@ -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 }} \ No newline at end of file diff --git a/backend/api.py b/backend/api.py index 493baf0..6be214d 100644 --- a/backend/api.py +++ b/backend/api.py @@ -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 @@ -153,6 +154,15 @@ 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() @@ -160,13 +170,16 @@ 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() @@ -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 @@ -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()} diff --git a/backend/scraper.py b/backend/scraper.py index 5b1725a..ca4f0a1 100644 --- a/backend/scraper.py +++ b/backend/scraper.py @@ -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, @@ -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.""" @@ -546,7 +559,7 @@ 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 @@ -554,22 +567,78 @@ def _backfill_fingerprints(cursor): 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(): @@ -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. @@ -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) @@ -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 diff --git a/backend/test_fingerprint_dedup.py b/backend/test_fingerprint_dedup.py new file mode 100644 index 0000000..80d9ac4 --- /dev/null +++ b/backend/test_fingerprint_dedup.py @@ -0,0 +1,162 @@ +"""The scraper's `seen` set and the prod unique constraint must agree on what +"the same job" means. They didn't: `seen` keyed on the raw lowercase +(company, role, location) triple while `internships_fingerprint_key` keys on +job_fingerprint(), which also strips punctuation. Two spellings of one company +therefore got distinct `seen` keys but one fingerprint, putting both rows in a +single execute_values page and aborting the scrape with a UniqueViolation. + +Run: python3 test_fingerprint_dedup.py +""" + +from read_db import job_fingerprint + +# Two listings for one job. Distinct raw keys, one fingerprint. +PAIRS = [ + ("Acme, Inc.", "Acme Inc."), + ("Acme Inc", "Acme Inc."), # whitespace collapse + ("Widgets, LLC", "Widgets LLC"), + ("Beta (USA) Co.", "Beta USA Co."), +] + +# Genuinely different jobs that must survive: different role, different +# location. Note these are NOT collapse pairs -- norm deletes punctuation +# without inserting a space and strips non-ASCII, so "Acme,Inc."/"Café Labs"/ +# "Foo-Bar" all fingerprint differently from their lookalikes. +DISTINCT = [ + ("Acme, Inc.", "Engineer", "NYC"), + ("Acme, Inc.", "Designer", "NYC"), + ("Acme, Inc.", "Engineer", "Remote"), + ("Foo-Bar", "Engineer", "NYC"), +] + + +def dedup_key(job): + """Mirrors the scraper's `seen` key. If this drifts from job_fingerprint, + the test below fails.""" + return job_fingerprint(job["company"], job["role"], job["location"]) + + +def main(): + # The upsert is a bare string literal, so py_compile says nothing about it. + # A duplicated clause here only surfaces as a psycopg2 SyntaxError against + # prod, which is how the "ON CONFLICT DO UPDATE / DO UPDATE SET" typo got + # pushed. Check the shape instead. + import re as _re + from scraper import _UPSERT_SQL + + flat = " ".join(_UPSERT_SQL.split()) + # Postgres rejects a bare "ON CONFLICT DO UPDATE" -- DO UPDATE requires an + # inference specification. Only DO NOTHING may omit it. That typo shipped + # once already and only failed against prod. + assert flat.count("DO UPDATE") == 1, flat + assert _re.search(r"\bON CONFLICT \(fingerprint\) DO UPDATE SET\b", flat), flat + # The fingerprint arbiter must cover the four-column one: identical + # (company, role, location) implies an identical fingerprint, so every + # four-column conflict is also a fingerprint conflict. Assert that + # implication rather than trusting the comment. + four_col = [("Acme, Inc.", "Engineer", "NYC", "http://x/1"), + ("Acme, Inc.", "Engineer", "NYC", "http://x/1")] + assert len({job_fingerprint(c, r, l) for c, r, l, _ in four_col}) == 1 + # And it must not over-cover: a different link with the same content is a + # conflict the arbiter is meant to catch. + diff_link = ("Acme, Inc.", "Engineer", "NYC", "http://x/2") + assert job_fingerprint(*diff_link[:3]) in {job_fingerprint(*q[:3]) for q in four_col} + + # The SET list must omit the columns the upsert arbitrates, so a conflict + # keeps the pre-existing row's identity (notably its link). + set_list = flat.split("DO UPDATE SET", 1)[1] + for col in ("company", "role", "location", "link"): + assert f"{col} = EXCLUDED" not in set_list, f"SET list overwrites {col}" + + for a, b in PAIRS: + ja = {"company": a, "role": "Engineer", "location": "NYC"} + jb = {"company": b, "role": "Engineer", "location": "NYC"} + raw_a = (a.lower(), "engineer", "nyc") + raw_b = (b.lower(), "engineer", "nyc") + assert job_fingerprint(**ja) == job_fingerprint(**jb), (a, b) + # The old key failed to collapse these; the new one must. + if raw_a != raw_b: + assert dedup_key(ja) == dedup_key(jb), f"dedup_key missed {a!r}/{b!r}" + + seen = set() + for company, role, loc in DISTINCT: + k = dedup_key({"company": company, "role": role, "location": loc}) + assert k not in seen, f"over-merged distinct jobs: {company}/{role}/{loc}" + seen.add(k) + + # The backfill must delete content duplicates, not leave them stale. A + # stale or NULL fingerprint is invisible to ON CONFLICT (fingerprint), so + # an insert matching such a row on (company, role, location, link) finds no + # arbiter and dies on internships_unique_job. Prod had 212 of these, and 52 + # more that were non-NULL but held a foreign normaliser's value. + def reconcile(rows): + """Returns (updates, doomed). Mirrors _backfill_fingerprints(). + + `rows` is (id, company, role, location, stored_fingerprint) in id order. + """ + updates, doomed, 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)) + return updates, doomed + + null_rows = [ + (1, "Acme, Inc.", "Engineer", "NYC", None), + (2, "Acme Inc.", "Engineer", "NYC", None), # same fingerprint as row 1 + (3, "Widgets, LLC", "Engineer", "NYC", None), + (4, "Widgets LLC", "Engineer", "NYC", None), # same fingerprint as row 3 + ] + updates, doomed = reconcile(null_rows) + assert [r[0] for r in updates] == [1, 3], updates + assert doomed == [2, 4], doomed + assert len({fp for _, fp in updates}) == len(updates), "reconcile emitted a duplicate" + # No surviving row may be left stale, or it becomes the next abort. + assert not (set(doomed) & {r[0] for r in updates}) + + # A row whose stored value is already correct is left alone — no write. + correct = job_fingerprint("Acme Inc.", "Engineer", "NYC") + updates, doomed = reconcile([(1, "Acme Inc.", "Engineer", "NYC", correct)]) + assert updates == [] and doomed == [], (updates, doomed) + + # A drifted row -- non-NULL, but not the normaliser's value for its own + # columns -- must be corrected. This is the 52-row case: the old code seeded + # `taken` from the stored column, so a drifted row could not be corrected if + # its stale value were anything but its own fingerprint, and if the stale + # value *was* its own fingerprint the row survived while remaining invisible + # to the arbiter. The DETAIL that proved it: + # stored 'tiktok|software engineer intern intelligent creation camera|san jose ca' + # normaliser 'tiktok|software engineer intern intelligent creationcamera|san jose ca' + drifted = [ + (1, "TikTok", "Software Engineer Intern, Intelligent Creation/Camera", "San Jose, CA", + "tiktok|software engineer intern intelligent creation camera|san jose ca"), + (2, "American Express", "Digital Product Analyst Intern", "NYC", + "american express|digital product analyst intern|ny"), + ] + updates, doomed = reconcile(drifted) + assert [r[0] for r in updates] == [1, 2], updates + assert doomed == [], doomed + # The corrections are the normaliser's values, punctuation deleted rather + # than spaced, and no location aliasing. + assert dict(updates) == { + 1: "tiktok|software engineer intern intelligent creationcamera|san jose ca", + 2: "american express|digital product analyst intern|nyc", + }, updates + + # And a drifted row whose corrected fingerprint is already held by an + # earlier row is deleted, not left holding a value the arbiter can't see. + updates, doomed = reconcile([(1, "Acme Inc.", "Engineer", "NYC", correct), + (2, "Acme, Inc.", "Engineer", "NYC", "acme inc old")]) + assert updates == [] and doomed == [2], (updates, doomed) + + print("OK: seen-key and fingerprint agree; distinct jobs kept; dupes deleted") + + +if __name__ == "__main__": + main()