perf: batch store_findings, dedup source_items in O(1), remove dead code (#206)

1. N+1 queries in store.store_findings()
   The old loop ran one SELECT per finding to check existence, then one
   INSERT or UPDATE. 100 findings cost 200 serial SQLite roundtrips.
   Now: one batch SELECT with WHERE source_url IN (...) builds a lookup
   dict, then executemany() handles all inserts and updates. Query count
   stays constant regardless of batch size. Benchmark on 500 findings:
   ~30ms to ~20ms; gap widens on slower storage.

2. O(n^2) source_items dedup in fusion.weighted_rrf()
   Merging an item into an existing candidate ran any(existing.source ==
   ... for existing in candidate.source_items), linearly scanning a list
   that grew with each merge. At 40 candidates with 20 source_items each,
   fusion went quadratic. Now tracks (source, item_id) tuples in a
   per-candidate set for O(1) lookup. The source_items list itself is
   unchanged since other code iterates it.

3. Dead code removal
   - providers.GeminiClient.ground_search() and .url_context_json(): zero
     callers. Deleted.
   - render._top_comment_excerpt(): zero callers. Deleted.
   - env.is_reddit_available(): one-line wrapper around get_reddit_source.
     Callers can check get_reddit_source(config) is not None directly.
This commit is contained in:
Ilia Alshanetsky
2026-04-25 17:17:17 -04:00
committed by GitHub
parent e6b89f2644
commit 2acbf8a869
5 changed files with 76 additions and 74 deletions
-8
View File
@@ -372,14 +372,6 @@ def config_exists() -> bool:
return False return False
def is_reddit_available(config: dict[str, Any]) -> bool:
"""Check if Reddit search is available.
v3 uses ScrapeCreators only.
"""
return bool(config.get('SCRAPECREATORS_API_KEY'))
def get_reddit_source(config: dict[str, Any]) -> str | None: def get_reddit_source(config: dict[str, Any]) -> str | None:
"""Determine which Reddit backend to use. """Determine which Reddit backend to use.
+6 -1
View File
@@ -116,6 +116,8 @@ def weighted_rrf(
"""Fuse ranked lists into a single candidate pool.""" """Fuse ranked lists into a single candidate pool."""
subqueries = {subquery.label: subquery for subquery in plan.subqueries} subqueries = {subquery.label: subquery for subquery in plan.subqueries}
candidates: dict[str, schema.Candidate] = {} candidates: dict[str, schema.Candidate] = {}
# Track (source, item_id) pairs already attached to each candidate for O(1) dedup.
seen_source_items: dict[str, set[tuple[str, str]]] = {}
for (label, source), items in streams.items(): for (label, source), items in streams.items():
subquery = subqueries[label] subquery = subqueries[label]
@@ -154,6 +156,7 @@ def weighted_rrf(
] ]
}, },
) )
seen_source_items[key] = {(item.source, item.item_id)}
continue continue
candidate = candidates[key] candidate = candidates[key]
@@ -179,7 +182,9 @@ def weighted_rrf(
candidate.subquery_labels.append(label) candidate.subquery_labels.append(label)
if item.source not in candidate.sources: if item.source not in candidate.sources:
candidate.sources.append(item.source) candidate.sources.append(item.source)
if not any(existing.source == item.source and existing.item_id == item.item_id for existing in candidate.source_items): source_item_key = (item.source, item.item_id)
if source_item_key not in seen_source_items[key]:
seen_source_items[key].add(source_item_key)
candidate.source_items.append(item) candidate.source_items.append(item)
candidate.metadata.setdefault("provenance", []).append( candidate.metadata.setdefault("provenance", []).append(
{ {
@@ -93,13 +93,6 @@ class GeminiClient(ReasoningClient):
) )
return extract_gemini_text(payload) return extract_gemini_text(payload)
def ground_search(self, model: str, prompt: str) -> dict[str, Any]:
return self._generate_content(model, prompt, tools=[{"google_search": {}}])
def url_context_json(self, model: str, prompt: str) -> dict[str, Any]:
return self.generate_json(model, prompt, tools=[{"url_context": {}}])
class OpenAIClient(ReasoningClient): class OpenAIClient(ReasoningClient):
name = "openai" name = "openai"
-10
View File
@@ -1505,16 +1505,6 @@ def _top_comments_list(item: schema.SourceItem | None, limit: int = 3, min_score
return [c for c in comments if (c.get("score") or 0) >= min_score][:limit] return [c for c in comments if (c.get("score") or 0) >= min_score][:limit]
def _top_comment_excerpt(item: schema.SourceItem | None) -> str | None:
if not item:
return None
comments = item.metadata.get("top_comments") or []
if not comments or not isinstance(comments[0], dict):
return None
top = comments[0]
return str(top.get("excerpt") or top.get("text") or "").strip() or None
def _comment_insight(item: schema.SourceItem | None) -> str | None: def _comment_insight(item: schema.SourceItem | None) -> str | None:
if not item: if not item:
return None return None
+59 -37
View File
@@ -346,46 +346,50 @@ def store_findings(
findings: List[Dict[str, Any]], findings: List[Dict[str, Any]],
) -> Dict[str, int]: ) -> Dict[str, int]:
"""Store findings with URL-based dedup. Returns counts of new/updated.""" """Store findings with URL-based dedup. Returns counts of new/updated."""
conn = _connect() # Collect findings that have a URL, preserving order.
new_count = 0 with_urls: List[tuple[str, Dict[str, Any]]] = []
updated_count = 0
try:
for f in findings: for f in findings:
url = f.get("source_url") or f.get("url") url = f.get("source_url") or f.get("url")
if not url: if url:
continue with_urls.append((url, f))
existing = conn.execute( if not with_urls:
"SELECT id, engagement_score, sighting_count FROM findings WHERE source_url = ?", conn = _connect()
(url,), try:
).fetchone()
if existing:
# Update engagement and re-sighting info
new_engagement = f.get("engagement_score", 0)
conn.execute( conn.execute(
"""UPDATE findings SET "UPDATE research_runs SET findings_new = 0, findings_updated = 0 WHERE id = ?",
last_seen = datetime('now'), (run_id,),
sighting_count = sighting_count + 1, )
engagement_score = ?, conn.commit()
run_id = ? finally:
WHERE id = ?""", conn.close()
( return {"new": 0, "updated": 0}
conn = _connect()
try:
# Single batch SELECT to find existing findings by URL.
urls = [url for url, _ in with_urls]
placeholders = ",".join("?" for _ in urls)
rows = conn.execute(
f"SELECT id, source_url, engagement_score FROM findings WHERE source_url IN ({placeholders})",
urls,
).fetchall()
existing_by_url = {row["source_url"]: row for row in rows}
update_rows: List[tuple] = []
insert_rows: List[tuple] = []
for url, f in with_urls:
existing = existing_by_url.get(url)
new_engagement = f.get("engagement_score", 0)
if existing:
update_rows.append((
max(new_engagement, existing["engagement_score"] or 0), max(new_engagement, existing["engagement_score"] or 0),
run_id, run_id,
existing["id"], existing["id"],
), ))
)
updated_count += 1
else: else:
# New finding insert_rows.append((
conn.execute(
"""INSERT INTO findings
(run_id, topic_id, source, source_url, source_title,
author, content, summary, engagement_score, relevance_score)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(
run_id, run_id,
topic_id, topic_id,
f.get("source", "unknown"), f.get("source", "unknown"),
@@ -394,13 +398,31 @@ def store_findings(
f.get("author", ""), f.get("author", ""),
f.get("content") or f.get("text", ""), f.get("content") or f.get("text", ""),
f.get("summary", ""), f.get("summary", ""),
f.get("engagement_score", 0), new_engagement,
f.get("relevance_score", 0), f.get("relevance_score", 0),
), ))
)
new_count += 1
# Update run stats if update_rows:
conn.executemany(
"""UPDATE findings SET
last_seen = datetime('now'),
sighting_count = sighting_count + 1,
engagement_score = ?,
run_id = ?
WHERE id = ?""",
update_rows,
)
if insert_rows:
conn.executemany(
"""INSERT INTO findings
(run_id, topic_id, source, source_url, source_title,
author, content, summary, engagement_score, relevance_score)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
insert_rows,
)
new_count = len(insert_rows)
updated_count = len(update_rows)
conn.execute( conn.execute(
"UPDATE research_runs SET findings_new = ?, findings_updated = ? WHERE id = ?", "UPDATE research_runs SET findings_new = ?, findings_updated = ? WHERE id = ?",
(new_count, updated_count, run_id), (new_count, updated_count, run_id),