diff --git a/adapters/_store.py b/adapters/_store.py new file mode 100644 index 0000000..d2c4ba4 --- /dev/null +++ b/adapters/_store.py @@ -0,0 +1,127 @@ +"""Shared DB storage helper for Athena adapters. + +Fixes the first_seen re-stamp bug (2026-07-13, Tony): + - Previously every adapter did INSERT OR REPLACE with first_seen = harvest + time, so re-harvesting an existing URL OVERWROTE the true publish date + with today's date. The recency guard then believed stale stories were new. + - true_first_seen() derives the REAL publish date from raw_metadata: + hackernews -> raw_meta['time'] (unix epoch) + reddit -> raw_meta['published'] (ISO) + rss -> raw_meta['published'] (RFC822 / ISO) + arxiv -> raw_meta['published'] (ISO) + huggingface-> raw_meta['createdAt'] (ISO) + Falls back to harvest time only if no source date exists. + - upsert_entries() is idempotent: first sight stores the true first_seen; + re-encounter PRESERVES the original first_seen and only bumps last_updated. +""" +import json +import sqlite3 +from datetime import datetime, timezone + + +def _epoch_to_iso(ts): + try: + return datetime.fromtimestamp(float(ts), tz=timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + except (TypeError, ValueError): + return None + + +def _norm_iso(s): + """Best-effort normalize an arbitrary date string to our ISO 'Z' format.""" + if not s: + return None + s = str(s).strip() + # already ISO-ish + try: + return datetime.fromisoformat(s.replace("Z", "+00:00")).astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + except ValueError: + pass + # RFC822 e.g. 'Sat, 11 Jul 2026 14:13:00 +000' -> normalize offset to +0000 + import re as _re + s_norm = _re.sub(r"([+-]\d{2})(\d{2})$", r"\1:\2", s) # +0000 -> +00:00 + if _re.search(r"[+-]\d{3}$", s_norm): # +000 -> +0000 + s_norm = s_norm[:-3] + "0" + s_norm[-3:] + for candidate in (s_norm, s): + for fmt in ("%a, %d %b %Y %H:%M:%S %z", "%a, %d %b %Y %H:%M:%S %Z"): + try: + return datetime.strptime(candidate, fmt).astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + except ValueError: + continue + # Email parser fallback (most robust for RFC822) + try: + from email.utils import parsedate_to_datetime + d = parsedate_to_datetime(s) + if d is not None: + return d.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + except Exception: + pass + return None + + +def true_first_seen(raw_meta, source, now_str): + """Derive the real publish date (ISO 'Z') for an entry, or harvest time.""" + if isinstance(raw_meta, str): + try: + raw_meta = json.loads(raw_meta) + except (json.JSONDecodeError, TypeError): + raw_meta = {} + if not isinstance(raw_meta, dict): + raw_meta = {} + src = (source or "").lower() + try: + if src == "hackernews": + return _epoch_to_iso(raw_meta.get("time")) or now_str + if src == "reddit": + return _norm_iso(raw_meta.get("published")) or now_str + if src == "rss": + return _norm_iso(raw_meta.get("published")) or now_str + if src == "arxiv": + return _norm_iso(raw_meta.get("published")) or now_str + if src == "huggingface": + return _norm_iso(raw_meta.get("createdAt")) or now_str + if src == "github": + return _norm_iso(raw_meta.get("created_at") or raw_meta.get("pushed_at") or raw_meta.get("published_at")) or now_str + except Exception: + return now_str + return now_str + + +UPSERT_SQL = """ +INSERT INTO entries + (source, source_id, url, title, extracted_text, summary, + category_tags, signal_score, raw_metadata, first_seen, last_updated) +VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) +ON CONFLICT(url) DO UPDATE SET + source = excluded.source, + source_id = excluded.source_id, + title = excluded.title, + extracted_text = excluded.extracted_text, + summary = excluded.summary, + category_tags = excluded.category_tags, + signal_score = excluded.signal_score, + raw_metadata = excluded.raw_metadata, + last_updated = excluded.last_updated, + first_seen = COALESCE((SELECT first_seen FROM entries WHERE url = excluded.url), excluded.first_seen) +""" + + +def upsert_entries(conn, entries): + """Idempotent store. Preserves original first_seen on re-harvest. + + `entries` is the list of dicts as built by each adapter; each dict must + already have first_seen set to the TRUE publish date (via true_first_seen) + and last_updated to the harvest time. + Returns count of rows written. + """ + cur = conn.cursor() + written = 0 + for e in entries: + cur.execute(UPSERT_SQL, ( + e["source"], e["source_id"], e["url"], e["title"], + e.get("extracted_text"), e.get("summary"), + e.get("category_tags"), e.get("signal_score"), + e.get("raw_metadata"), e["first_seen"], e["last_updated"], + )) + written += 1 + conn.commit() + return written diff --git a/adapters/arxiv.py b/adapters/arxiv.py index 2002d7a..adcb9fd 100644 --- a/adapters/arxiv.py +++ b/adapters/arxiv.py @@ -33,6 +33,7 @@ from datetime import datetime, timedelta, timezone from html import unescape from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries # arXiv API ARXIV_API = "http://export.arxiv.org/api/query" @@ -443,6 +444,7 @@ class ArxivAdapter(SourceAdapter): raw_meta["applied_domain"] = paper["_applied_domain"] now_str = now.strftime("%Y-%m-%dT%H:%M:%SZ") + first_seen = true_first_seen(raw_meta, "arxiv", now_str) entries.append({ "source": "arxiv", "source_id": source_id, @@ -453,7 +455,7 @@ class ArxivAdapter(SourceAdapter): "category_tags": json.dumps(tags), "signal_score": score, "raw_metadata": json.dumps(raw_meta), - "first_seen": now_str, + "first_seen": first_seen, "last_updated": now_str, }) @@ -491,24 +493,7 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") - + stored = upsert_entries(conn, entries) conn.commit() conn.close() print(f" Stored {stored} entries") diff --git a/adapters/github.py b/adapters/github.py index dda90a5..a38dd36 100644 --- a/adapters/github.py +++ b/adapters/github.py @@ -18,6 +18,7 @@ import urllib.parse from datetime import datetime, timedelta, timezone from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries class GitHubAdapter(SourceAdapter): @@ -256,17 +257,18 @@ class GitHubAdapter(SourceAdapter): } now_str = now.strftime("%Y-%m-%dT%H:%M:%SZ") + first_seen = true_first_seen(raw_meta, "github", now_str) entries.append({ "source": "github", "source_id": source_id, "url": url, "title": title, "extracted_text": readme_text, - "summary": None, # LLM later + "summary": None, "category_tags": json.dumps(tags), "signal_score": round(score, 2), "raw_metadata": json.dumps(raw_meta), - "first_seen": now_str, + "first_seen": first_seen, "last_updated": now_str, }) @@ -304,24 +306,7 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") - + stored = upsert_entries(conn, entries) conn.commit() conn.close() print(f" Stored {stored} entries") diff --git a/adapters/hackernews.py b/adapters/hackernews.py index 6680c6b..1677306 100644 --- a/adapters/hackernews.py +++ b/adapters/hackernews.py @@ -21,6 +21,7 @@ import urllib.error from datetime import datetime, timezone from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries class HackerNewsAdapter(SourceAdapter): @@ -243,6 +244,8 @@ class HackerNewsAdapter(SourceAdapter): } now_str = now.strftime("%Y-%m-%dT%H:%M:%SZ") + # first_seen = TRUE publish date (HN 'time'), not harvest time + first_seen = true_first_seen(raw_meta, "hackernews", now_str) entries.append({ "source": "hackernews", "source_id": source_id, @@ -253,7 +256,7 @@ class HackerNewsAdapter(SourceAdapter): "category_tags": json.dumps(tags), "signal_score": score, "raw_metadata": json.dumps(raw_meta), - "first_seen": now_str, + "first_seen": first_seen, "last_updated": now_str, }) @@ -288,27 +291,8 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") - - conn.commit() - conn.close() - print(f" Stored {stored} entries") + stored = upsert_entries(conn, entries) + print(f"\n Stored {stored} entries") # Print top 5 print(f"\n Top entries:") diff --git a/adapters/huggingface.py b/adapters/huggingface.py index 76c4746..dfe313f 100644 --- a/adapters/huggingface.py +++ b/adapters/huggingface.py @@ -35,6 +35,7 @@ import urllib.error from datetime import datetime, timedelta, timezone from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries class HuggingFaceAdapter(SourceAdapter): @@ -328,6 +329,7 @@ class HuggingFaceAdapter(SourceAdapter): } now_str = now.strftime("%Y-%m-%dT%H:%M:%SZ") + first_seen = true_first_seen(raw_meta, "huggingface", now_str) entries.append({ "source": "huggingface", "source_id": source_id, @@ -338,7 +340,7 @@ class HuggingFaceAdapter(SourceAdapter): "category_tags": json.dumps(tags), "signal_score": score, "raw_metadata": json.dumps(raw_meta), - "first_seen": now_str, + "first_seen": first_seen, "last_updated": now_str, }) @@ -374,24 +376,7 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") - + stored = upsert_entries(conn, entries) conn.commit() conn.close() print(f" Stored {stored} entries") diff --git a/adapters/reddit.py b/adapters/reddit.py index 2f3815d..25eb23d 100644 --- a/adapters/reddit.py +++ b/adapters/reddit.py @@ -25,6 +25,7 @@ from datetime import datetime, timezone from html import unescape from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries class RedditAdapter(SourceAdapter): @@ -469,6 +470,7 @@ class RedditAdapter(SourceAdapter): } now_str = now.strftime("%Y-%m-%dT%H:%M:%SZ") + first_seen = true_first_seen(raw_meta, "reddit", now_str) entries.append({ "source": "reddit", "source_id": source_id, @@ -479,7 +481,7 @@ class RedditAdapter(SourceAdapter): "category_tags": json.dumps(tags), "signal_score": score, "raw_metadata": json.dumps(raw_meta), - "first_seen": now_str, + "first_seen": first_seen, "last_updated": now_str, }) @@ -514,23 +516,8 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") + stored = upsert_entries(conn, entries) + print(f"\n Stored {stored} entries") conn.commit() conn.close() diff --git a/adapters/rss_feeds.py b/adapters/rss_feeds.py index 42b516e..e5614b6 100644 --- a/adapters/rss_feeds.py +++ b/adapters/rss_feeds.py @@ -25,6 +25,7 @@ from datetime import datetime, timedelta, timezone from email.utils import parsedate_to_datetime from adapters import SourceAdapter +from adapters._store import true_first_seen, upsert_entries # Curated feed list — AI-focused, reliable, diverse publishers. @@ -192,6 +193,9 @@ class RSSFeedsAdapter(SourceAdapter): score = self._score(entry_dict, label, age_hours) entry_dict["score"] = score + now_iso = now.strftime("%Y-%m-%dT%H:%M:%SZ") + first_seen = true_first_seen( + {"published": published}, "rss", now_iso) all_entries.append({ "source": "rss", "source_id": f"{source_key}:{entry.get('id', entry.get('link', ''))[-30:]}", @@ -210,8 +214,8 @@ class RSSFeedsAdapter(SourceAdapter): "tags": tags, "score_type": "estimated", }), - "first_seen": now.strftime("%Y-%m-%dT%H:%M:%SZ"), - "last_updated": now.strftime("%Y-%m-%dT%H:%M:%SZ"), + "first_seen": first_seen, + "last_updated": now_iso, }) time.sleep(0.5) # polite spacing @@ -260,23 +264,7 @@ if __name__ == "__main__": conn.commit() cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"], entry["summary"], - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" DB error: {e}") + stored = upsert_entries(conn, entries) conn.commit() conn.close() print(f"\n Stored {stored} entries") diff --git a/pipeline.py b/pipeline.py index 16b723f..9606730 100644 --- a/pipeline.py +++ b/pipeline.py @@ -25,6 +25,7 @@ from datetime import datetime, timezone sys.path.insert(0, os.path.dirname(__file__)) from adapters import SourceAdapter +from adapters._store import upsert_entries # Adapter registry — add new adapters here (one line each) ADAPTERS = { @@ -51,27 +52,14 @@ def init_db(db_path: str, schema_path: str) -> sqlite3.Connection: def store_entries(conn: sqlite3.Connection, entries: list[dict]) -> int: - """Store entries using INSERT OR REPLACE (dedup by source+source_id).""" - cur = conn.cursor() - stored = 0 - for entry in entries: - try: - cur.execute(""" - INSERT OR REPLACE INTO entries - (source, source_id, url, title, extracted_text, summary, - category_tags, signal_score, raw_metadata, first_seen, last_updated) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - """, ( - entry["source"], entry["source_id"], entry["url"], entry["title"], - entry["extracted_text"] or "", entry["summary"], # None → NULL in DB - entry["category_tags"], entry["signal_score"], - entry["raw_metadata"], entry["first_seen"], entry["last_updated"], - )) - stored += 1 - except Exception as e: - print(f" ⚠ DB error on {entry.get('source', '?')}/{entry.get('source_id', '?')}: {e}") - conn.commit() - return stored + """Store entries idempotently (dedup by url). + + FIX (2026-07-13, Tony): was INSERT OR REPLACE which OVERWROTE first_seen + with the harvest time on every re-harvest, turning stale stories into + "today". Now uses an UPSERT that preserves the ORIGINAL first_seen and only + bumps last_updated. Entries must already carry first_seen = true publish date. + """ + return upsert_entries(conn, entries) def verify_entries(conn: sqlite3.Connection, source: str, sample_size: int = 3):