FIX: first_seen re-stamp bug — persist true source publish date, idempotent upsert, backfill DB
This commit is contained in:
@@ -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
|
||||
+4
-19
@@ -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")
|
||||
|
||||
+5
-20
@@ -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")
|
||||
|
||||
+6
-22
@@ -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:")
|
||||
|
||||
+4
-19
@@ -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")
|
||||
|
||||
+5
-18
@@ -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()
|
||||
|
||||
+7
-19
@@ -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")
|
||||
|
||||
+9
-21
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user