返回 last30days-skill
pipeline.py
根目录 / skills / last30days / scripts / lib / pipeline.py
1 """v3.0.0 orchestration pipeline."""
2
3 from __future__ import annotations
4
5 import contextlib
6 import copy
7 from collections.abc import Iterable, Mapping
8 import math
9 import queue
10 import re
11 import sqlite3
12 import sys
13 import threading
14 import time
15 from collections import Counter
16 from concurrent.futures import ThreadPoolExecutor, as_completed
17 from dataclasses import dataclass, field, replace
18 from datetime import date, datetime, timedelta, timezone
19 from pathlib import Path
20 from shutil import which
21 from typing import Any
22
23 from . import (
24 amazon,
25 arxiv,
26 bird_x,
27 bluesky,
28 brightdata,
29 corpus,
30 dates,
31 dedupe,
32 digg,
33 dripstack,
34 entity_extract,
35 env,
36 github,
37 grok_x,
38 grounding,
39 hackernews,
40 health,
41 hiring_signals,
42 http,
43 instagram,
44 jobs,
45 linkedin,
46 library,
47 library_index,
48 log,
49 meta_ads,
50 normalize,
51 permission_preflight,
52 perplexity,
53 pinterest,
54 planner,
55 polymarket,
56 providers,
57 query,
58 reddit,
59 reddit_listing,
60 reddit_public,
61 relevance,
62 rerank,
63 schema,
64 signals,
65 snippet,
66 stocktwits,
67 techmeme,
68 telegram,
69 threads,
70 tiktok,
71 topic_shape,
72 truthsocial,
73 trustpilot,
74 x_api,
75 x_envelope,
76 x_judge,
77 xai_x,
78 xiaohongshu_api,
79 xquik,
80 xurl_x,
81 youtube_yt,
82 )
83 from .cluster import cluster_candidates
84 from . import fusion
85 from . import render
86 from .fusion import collapse_duplicate_urls, weighted_rrf
87
88 DISCOVERY_SOURCES = ("reddit", "hackernews", "digg", "x")
89 _DISCOVERY_GENERIC_DOMAIN_TERMS = {
90 "ai", "artificial", "intelligence", "tech", "technology", "trending", "trend",
91 }
92
93 DEPTH_SETTINGS = {
94 "quick": {"per_stream_limit": 6, "pool_limit": 15, "rerank_limit": 12},
95 "default": {"per_stream_limit": 12, "pool_limit": 40, "rerank_limit": 40},
96 "deep": {"per_stream_limit": 20, "pool_limit": 60, "rerank_limit": 60},
97 }
98
99 SEARCH_ALIAS = {
100 "hn": "hackernews",
101 "bsky": "bluesky",
102 "truth": "truthsocial",
103 "web": "grounding",
104 "xhs": "xiaohongshu",
105 "meta": "meta_ads",
106 "meta-ads": "meta_ads",
107 "xquik": "x", # xquik is a backend of the single "x" source, not its own source
108 }
109
110 # trustpilot is capped at 1: every subquery would use the identical company
111 # identifier, so N streams are pure redundancy -- and each extra stream risks
112 # its own WAF-cookie Chrome harvest.
113 # amazon is capped at 1 for the same reason as trustpilot: the model supplies
114 # one product keyword for the run, so every subquery would issue the identical
115 # product search. Extra streams would be pure redundancy at one credit each.
116 # meta_ads is capped at 1 for the same reason as amazon: one advertiser page is
117 # resolved per run, so every subquery would issue the identical page fetch.
118 MAX_SOURCE_FETCHES: dict[str, int] = {
119 "x": 2, "jobs": 1, "linkedin": 1, "stocktwits": 1, "trustpilot": 1, "amazon": 1,
120 "telegram": 1, "meta_ads": 1,
121 }
122
123 # Sources whose thin result is their normal success state, so the "<3 items"
124 # retry would re-fetch them after every success -- bypassing
125 # MAX_SOURCE_FETCHES and, for the resolved-entity sources, re-resolving
126 # WITHOUT the caller's override (a lookalike-misattribution path).
127 # trustpilot returns at most ONE item by design.
128 # perplexity answers once per run.
129 # meta_ads resolves one advertiser page per run, so a brand that genuinely
130 # ran two creatives this month is complete; a retry would re-resolve the
131 # page and re-spend the discovery credit.
132 THIN_RETRY_EXEMPT: frozenset[str] = frozenset({"trustpilot", "perplexity", "meta_ads"})
133
134 # Stream-artifact keys promoted to named top-level report artifacts. A stream
135 # artifact only ever reaches the report as an anonymous entry in the grounding
136 # list, so anything the renderer needs by name has to be lifted out of it --
137 # most importantly on a zero-item run, which is exactly when naming the
138 # resolved advertiser and its counts matters most.
139 STREAM_ARTIFACT_LIFT_KEYS: tuple[str, ...] = ("meta_ads_page", "meta_ads_tally")
140
141
142 def _lift_stream_artifacts(bundle) -> None:
143 """Promote per-stream artifacts the renderer reads by name."""
144 for stream_artifact in bundle.artifacts.get("grounding", []):
145 if not isinstance(stream_artifact, dict):
146 continue
147 for key in STREAM_ARTIFACT_LIFT_KEYS:
148 value = stream_artifact.get(key)
149 if value:
150 bundle.artifacts[key] = value
151
152
153 _FAILURE_SPECIFICITY = {
154 health.AUTH_FAILED: 0,
155 health.PAYMENT_REQUIRED: 1,
156 health.RATE_LIMITED: 2,
157 health.SCHEMA_DRIFT: 3,
158 health.TIMEOUT: 4,
159 health.UNREACHABLE: 5,
160 health.ERROR: 6,
161 }
162
163
164 @dataclass
165 class PaidSourceBudget:
166 """Command-wide, thread-safe budget for paid source adapter calls."""
167
168 used: int = 0
169 owner: str | None = None
170 _lock: Any = field(default_factory=threading.Lock, repr=False)
171
172 def try_consume(self, limit: int, *, claimant: str | None = None) -> bool:
173 with self._lock:
174 if self.owner is not None and claimant != self.owner:
175 return False
176 if self.used >= limit:
177 return False
178 self.used += 1
179 return True
180
181
182 def _source_fetch_cap(source: str, config: dict[str, Any]) -> int | None:
183 """Return the effective per-run cap for one source.
184
185 Every Perplexity adapter call is paid, and ``both`` performs two paid POSTs.
186 A generic fetch-cap override must not multiply either normal or Deep
187 Research mode across planner subqueries.
188 """
189 override = config.get("_max_source_fetches")
190 if source == "perplexity":
191 return 1 if override is None else min(1, int(override))
192 cap = MAX_SOURCE_FETCHES.get(source)
193 if cap is not None and override is not None:
194 return int(override)
195 return cap
196
197
198 def _resolve_depth_settings(depth: str, config: dict[str, Any]) -> dict[str, int]:
199 """Depth profile with optional CLI cap overrides applied (issue #716).
200
201 Returns a copy so the module-level DEPTH_SETTINGS is never mutated. Overrides
202 are set directly (not max()) so callers can also lower a cap. `--max-results`
203 raises the final ranked pool (pool_limit/rerank_limit); `--max-per-source`
204 raises the per-stream truncation applied before pooling. The per-source fetch
205 cap (`--max-source-fetches`) is applied separately at the fetch site.
206 """
207 settings = dict(DEPTH_SETTINGS[depth])
208 # `is not None` (not truthiness) so an explicit 0 is honored as a real lower
209 # bound rather than ignored as "unset" — matches how main() stashes these.
210 max_per_source = config.get("_max_per_source")
211 if max_per_source is not None:
212 settings["per_stream_limit"] = int(max_per_source)
213 max_results = config.get("_max_results")
214 if max_results is not None:
215 settings["pool_limit"] = int(max_results)
216 settings["rerank_limit"] = int(max_results)
217 return settings
218
219 # Per-handle result caps for the X handle-search lanes. The FROM lane (the
220 # subject's own timeline) is the single best source for a person topic, so it
221 # gets the highest cap; the ABOUT (mention) and related-handle lanes stay
222 # modest so total volume and request budget don't balloon.
223 FROM_LANE_COUNT_PER = 8
224 MENTION_LANE_COUNT_PER = 5
225 RELATED_HANDLE_COUNT_PER = 3
226
227
228 def _has_perplexity_provider(config: dict[str, Any]) -> bool:
229 # Prefer direct Agent/Search APIs, but preserve the synchronous OpenRouter
230 # Sonar fallback for existing installs.
231 return bool(
232 config.get("PERPLEXITY_API_KEY") or config.get("OPENROUTER_API_KEY")
233 )
234
235 MOCK_AVAILABLE_SOURCES = [
236 "reddit",
237 "x",
238 "youtube",
239 "tiktok",
240 "instagram",
241 "hackernews",
242 "bluesky",
243 "truthsocial",
244 "polymarket",
245 "grounding",
246 "xiaohongshu",
247 "github",
248 "perplexity",
249 "threads",
250 "pinterest",
251 "digg",
252 "arxiv",
253 "techmeme",
254 "trustpilot",
255 "amazon",
256 "meta_ads",
257 "jobs",
258 "linkedin",
259 "corpus",
260 "dripstack",
261 "telegram",
262 ]
263
264
265 def normalize_requested_sources(sources: list[str] | None) -> list[str] | None:
266 if not sources:
267 return None
268 normalized = []
269 for source in sources:
270 key = SEARCH_ALIAS.get(source.lower(), source.lower())
271 if key not in normalized:
272 normalized.append(key)
273 return normalized
274
275
276 def available_sources(
277 config: dict[str, Any],
278 requested_sources: list[str] | None = None,
279 *,
280 x_pending: bool | None = None,
281 local_only: bool = False,
282 x_envelope: bool = False,
283 suppress_x_host_lane: bool = False,
284 ) -> list[str]:
285 """List the sources the next run can serve.
286
287 ``local_only=True`` is the safe/diagnose flavor (doctor's permission
288 block): availability is answered from local evidence only, so the X
289 check never spawns xurl's live ``whoami`` network call. Research-time
290 callers keep the default live semantics.
291
292 X is listed when an engine backend is available, or browser auth is
293 pending, or the hosting model declared the X connector lane
294 (``env.x_host_lane_declared``), or a validated ``--x-posts`` envelope is
295 present for this run (``x_envelope``), in every cookie mode.
296 ``suppress_x_host_lane`` turns only the lane branch off (discovery
297 enrichment passes); an envelope still counts.
298 """
299 available: list[str] = []
300 # reddit_public needs no API key - always available
301 available.append("reddit")
302 if corpus.resolve_directories(
303 config.get("_CORPUS_DIRS"), config.get("LAST30DAYS_CORPUS_DIRS")
304 ):
305 available.append("corpus")
306 if config.get("SCRAPECREATORS_API_KEY"):
307 available.extend(["tiktok", "instagram"])
308 if env.get_x_source(config, local_only=local_only):
309 available.append("x")
310 elif x_envelope or (
311 not suppress_x_host_lane and env.x_host_lane_declared(config)
312 ):
313 # Host-fetched X lane: the model passes connector results through
314 # --x-posts, so X is served without an engine backend.
315 available.append("x")
316 else:
317 # Safe inspection (--diagnose/--preflight) skips browser-cookie
318 # extraction, so get_x_source is None even though a real run would
319 # authenticate X via FROM_BROWSER. Report it as available so consumers
320 # of available_sources (SKILL.md ACTIVE_SOURCES_LIST) don't under-report.
321 # diagnose() precomputes the predicate and passes it via x_pending to
322 # avoid evaluating it twice in one diagnose() call.
323 if x_pending is None:
324 x_pending = env.x_pending_browser_auth(config)
325 if x_pending:
326 available.append("x")
327 if which("yt-dlp") or env.is_youtube_sc_available(config):
328 available.append("youtube")
329 available.extend(["hackernews", "polymarket"])
330 # StockTwits is gated to ticker/crypto topics only (flag set in run()).
331 if config.get("_financial_topic"):
332 available.append("stocktwits")
333 # GitHub is reachable via the unauthenticated REST tier too, so it is
334 # available even without a token/gh CLI (a token only raises rate limits).
335 available.append("github")
336 # DripStack is opt-in only (owner decision, #791): a commercial
337 # third-party API must never receive default-run traffic. Opt in per run
338 # (--search dripstack) or persistently (INCLUDE_SOURCES=dripstack in
339 # .env, the LinkedIn/Perplexity pattern); the search API is free and
340 # public (no key), so the opt-in itself is the gate.
341 include_sources = {
342 token.strip()
343 for token in (config.get("INCLUDE_SOURCES") or "").lower().split(",")
344 if token.strip()
345 }
346 if "dripstack" in include_sources or (
347 requested_sources and "dripstack" in requested_sources
348 ):
349 available.append("dripstack")
350 if which("digg-pp-cli"):
351 available.append("digg")
352 # arXiv is default-on when its Printing Press CLI is installed (zero auth).
353 # The adapter relevance-and-recency gates so it stays quiet off-topic.
354 if which("arxiv-pp-cli"):
355 available.append("arxiv")
356 # Techmeme is default-on when its CLI is installed (zero auth; sub-second
357 # local sync before each run's first search).
358 if which("techmeme-pp-cli"):
359 available.append("techmeme")
360 if env.is_bluesky_available(config):
361 available.append("bluesky")
362 if env.is_truthsocial_available(config):
363 available.append("truthsocial")
364 # Grounding (general web) is available when a paid backend is configured OR
365 # the keyless floor is permitted (i.e. the host has no native search). On a
366 # native-search host with no paid key, keyless_web_allowed is False and the
367 # engine leaves general web to the model's own search.
368 if (config.get("BRAVE_API_KEY") or config.get("EXA_API_KEY")
369 or config.get("SERPER_API_KEY") or config.get("PARALLEL_API_KEY")
370 or env.keyless_web_allowed(config)):
371 available.append("grounding")
372 if requested_sources and "jobs" in requested_sources:
373 available.append("jobs")
374 # Perplexity Agent API: opt-in additive source via INCLUDE_SOURCES=perplexity
375 if _has_perplexity_provider(config) and (
376 "perplexity" in include_sources or (requested_sources and "perplexity" in requested_sources)
377 ):
378 available.append("perplexity")
379 # LinkedIn: opt-in additive source via INCLUDE_SOURCES=linkedin (same
380 # consent pattern as Perplexity). Unlike tiktok/instagram, which are
381 # offered during SKILL.md Step 0 onboarding, LinkedIn is power-user-only
382 # and must not silently activate for existing SCRAPECREATORS_API_KEY
383 # holders.
384 if config.get("SCRAPECREATORS_API_KEY") and (
385 "linkedin" in include_sources or (requested_sources and "linkedin" in requested_sources)
386 ):
387 available.append("linkedin")
388 # Trustpilot: opt-in additive source via INCLUDE_SOURCES=trustpilot (same
389 # consent pattern as Perplexity/LinkedIn). Off by default -- unlike arXiv and
390 # Techmeme, which are zero-auth, it can spawn a one-time headless-Chrome WAF
391 # cookie harvest on a brand topic, so activating it is the user's choice.
392 if which("trustpilot-pp-cli") and (
393 "trustpilot" in include_sources or (requested_sources and "trustpilot" in requested_sources)
394 ):
395 available.append("trustpilot")
396 # Amazon: opt-in additive source, dual-gated. The Bright Data CLI must be
397 # on the agent subprocess PATH and carry a credential signal, AND the run
398 # must ask for it -- the model per-run via --search, or the user durably
399 # via INCLUDE_SOURCES=amazon. Never inferred from topic shape: the engine
400 # misroutes most shopping phrasings, and auto-firing would spend a CLI
401 # owner's credits on runs that have nothing to do with products.
402 if brightdata.is_available(config) and (
403 "amazon" in include_sources or (requested_sources and "amazon" in requested_sources)
404 ):
405 available.append("amazon")
406 # Meta Ads: opt-in additive source on the Amazon precedent. The
407 # ScrapeCreators key must be present AND the run must ask for it -- the
408 # model per-run via --search, or the user durably via
409 # INCLUDE_SOURCES=meta_ads. Never inferred from topic shape: keyword ad
410 # search on a non-brand topic returns a wrong-entity advertiser, and
411 # auto-firing would spend credits resolving it.
412 if config.get("SCRAPECREATORS_API_KEY") and (
413 "meta_ads" in include_sources
414 or (requested_sources and "meta_ads" in requested_sources)
415 ):
416 available.append("meta_ads")
417 if (
418 "xiaohongshu" in include_sources
419 or (requested_sources and "xiaohongshu" in requested_sources)
420 ) and env.is_xiaohongshu_available(config):
421 available.append("xiaohongshu")
422 # Threads: opt-in via INCLUDE_SOURCES (same pattern as perplexity/linkedin).
423 # Was auto-on with the key; gated so the onboarding "Everything" tier is a
424 # real choice vs the "Recommended" (TikTok/Instagram) tier.
425 if env.is_threads_available(config) and (
426 "threads" in include_sources or (requested_sources and "threads" in requested_sources)
427 ):
428 available.append("threads")
429 # Pinterest: opt-in via INCLUDE_SOURCES. Previously read requested_sources
430 # only, so a persisted INCLUDE_SOURCES=pinterest never activated it; now it
431 # honors both the per-run --sources list and the saved config.
432 if env.is_pinterest_available(config) and (
433 "pinterest" in include_sources or (requested_sources and "pinterest" in requested_sources)
434 ):
435 available.append("pinterest")
436 # Telegram: opt-in via INCLUDE_SOURCES AND requires a channel list. The
437 # channel list (TELEGRAM_SOURCES env or --telegram-sources CLI) is the gate:
438 # without named channels there is no discovery endpoint to call.
439 if config.get("SCRAPECREATORS_API_KEY") and (
440 "telegram" in include_sources or (requested_sources and "telegram" in requested_sources)
441 ):
442 if telegram.is_telegram_configured(config):
443 available.append("telegram")
444 # xquik is a backend of the single "x" source (see env.x_backend_chain),
445 # not a separate parallel source — registered via the "x" entry above.
446 exclude = {s.strip().lower() for s in (config.get("EXCLUDE_SOURCES") or "").split(",") if s.strip()}
447 if exclude:
448 available = [s for s in available if s not in exclude]
449 return available
450
451
452 def _mock_discovery_items(
453 source: str,
454 domain: str,
455 to_date: str,
456 ) -> list[dict[str, Any]]:
457 """Deterministic listing fixtures for the public --mock CLI contract."""
458 labels = [
459 "Agent memory protocols",
460 "Browser-using agents",
461 "Local agent runtimes",
462 "Multi-agent orchestration",
463 "Agent security sandboxes",
464 "Voice agent latency",
465 ]
466 end = datetime.fromisoformat(to_date).date()
467 items: list[dict[str, Any]] = []
468 for index, label in enumerate(labels, start=1):
469 published = (end - timedelta(days=index)).isoformat()
470 slug = re.sub(r"[^a-z0-9]+", "-", label.lower()).strip("-")
471 if source == "reddit":
472 items.append({
473 "id": f"discovery-r-{index}",
474 "title": label,
475 "url": f"https://reddit.com/r/example/comments/{slug}",
476 "subreddit": "example",
477 "date": published,
478 "engagement": {"score": 180 - index * 10, "num_comments": 30 + index},
479 "selftext": label,
480 "relevance": 0.9,
481 "why_relevant": "Mock discovery listing",
482 })
483 elif source == "hackernews":
484 items.append({
485 "id": f"discovery-hn-{index}",
486 "title": label,
487 "url": f"https://example.com/{slug}",
488 "hn_url": f"https://news.ycombinator.com/item?id={index}",
489 "author": f"example{index}",
490 "date": published,
491 "engagement": {"points": 120 - index * 8, "comments": 20 + index},
492 "relevance": 0.88,
493 "why_relevant": "Mock HN discovery listing",
494 })
495 elif source == "digg":
496 items.append({
497 "id": f"discovery-d-{index}",
498 "title": label,
499 "url": f"https://di.gg/ai/{slug}",
500 "tldr": label,
501 "date": published,
502 "engagement": {"postCount": 30 - index, "uniqueAuthors": 12 - index},
503 "relevance": 0.9,
504 "why_relevant": "Mock Digg discovery cluster",
505 })
506 elif source == "x":
507 items.append({
508 "id": f"discovery-x-{index}",
509 "text": label,
510 "url": f"https://x.com/example{index}/status/{index}",
511 "author_handle": f"example{index}",
512 "date": published,
513 "engagement": {"likes": 140 - index * 9, "reposts": 18 + index},
514 "relevance": 0.9,
515 "why_relevant": "Mock X discovery activity",
516 })
517 return items
518
519
520 def _matches_discovery_domain(domain: str, text: str) -> bool:
521 """Require a distinctive domain term, not a generic token such as ``AI``."""
522 def terms(value: str) -> set[str]:
523 # Keep BOTH the surface form and the naive stem: replacing the token
524 # broke non-plurals ("bias" -> "bia", "crisis" -> "crisi") so in-domain
525 # listings stopped intersecting. The union preserves plural matching
526 # without corrupting the anchor.
527 words: set[str] = set()
528 for word in relevance.tokenize(value):
529 words.add(word)
530 if len(word) > 4 and word.endswith("s") and not word.endswith("ss"):
531 words.add(word[:-1])
532 return words
533
534 domain_terms = terms(domain)
535 anchors = domain_terms - _DISCOVERY_GENERIC_DOMAIN_TERMS
536 return bool((anchors or domain_terms) & terms(text))
537
538
539 def _fetch_discovery_source(
540 source: str,
541 plan: schema.DiscoveryPlan,
542 *,
543 from_date: str,
544 to_date: str,
545 depth: str,
546 mock: bool,
547 config: dict[str, Any],
548 keyword_gate: bool = True,
549 ) -> tuple[list[dict[str, Any]], str | None]:
550 """Fetch one listing/river source for the nominate stage.
551
552 ``keyword_gate`` controls whether items are filtered to the domain by
553 ``_matches_discovery_domain``. Domain-scoped discovery (``--discover X``)
554 keeps the gate on; global trending (``--discover`` with no domain) turns it
555 off, because there is no keyword to gate against - the river feeds ARE the
556 "what is hot right now" signal, and the confidence floor downstream is what
557 keeps junk out, not a keyword match.
558 """
559 if mock:
560 return _mock_discovery_items(source, plan.domain, to_date), None
561 if source == "reddit":
562 result = reddit_listing.fetch_discovery_listings(
563 plan.subreddits, depth=depth, query=plan.domain,
564 )
565 items = result.get("items") or []
566 if keyword_gate:
567 items = [
568 item for item in items
569 if _matches_discovery_domain(
570 plan.domain,
571 f"{item.get('title') or ''} {item.get('selftext') or ''}",
572 )
573 ]
574 return items, "; ".join(result.get("errors") or []) or None
575 if source == "hackernews":
576 result = hackernews.fetch_discovery_listings(from_date, to_date, depth=depth)
577 items = result.get("items") or []
578 for item in items:
579 item["relevance"] = relevance.token_overlap_relevance(
580 plan.domain,
581 str(item.get("title") or ""),
582 )
583 # HN is a broad technology listing, so keep only domain-bearing stories
584 # when a domain is in play; global trending keeps the whole front page.
585 if keyword_gate:
586 items = [
587 item for item in items
588 if _matches_discovery_domain(plan.domain, str(item.get("title") or ""))
589 ]
590 errors = result.get("errors") or []
591 return items, "; ".join(errors) or None
592 if source == "digg":
593 result = digg.search_digg(plan.domain, from_date, to_date, depth=depth)
594 items = digg.parse_digg_response(result, query=plan.domain)
595 # Digg is an AI-focused broad listing, so keep only domain-bearing
596 # clusters when scoped; global trending keeps the whole feed.
597 if keyword_gate:
598 items = [
599 item for item in items
600 if _matches_discovery_domain(plan.domain, str(item.get("title") or ""))
601 ]
602 return items, result.get("error")
603 if source == "x":
604 # Discovery uses domain directly as query (no planner search_query)
605 query = plan.domain
606 last_error = ""
607 for backend in env.x_backend_chain(config):
608 items, error = _fetch_x_backend(
609 backend, query, from_date, to_date, depth, config,
610 )
611 if items:
612 # Earlier failed-over backends' errors are observability, not
613 # degradation - but the producing backend's own error means
614 # these items are partial and must surface as such.
615 if last_error:
616 print(f"[x] earlier backend failed: {last_error}", file=sys.stderr)
617 return items, error or None
618 if error:
619 last_error = f"{backend}: {error}"
620 return [], last_error or None
621 raise ValueError(f"Unsupported discovery source: {source}")
622
623
624 def _discovery_engagement(
625 items: list[schema.SourceItem],
626 ) -> dict[str, dict[str, float | int]]:
627 totals: dict[str, dict[str, float | int]] = {}
628 for item in items:
629 bucket = totals.setdefault(item.source, {})
630 for field, value in item.engagement.items():
631 if not isinstance(value, (int, float)) or isinstance(value, bool):
632 continue
633 # Rank/score/reach metadata is not additive engagement: summing
634 # Digg ranks across items fabricates a metric (agent-export uses
635 # the same counter-field rule).
636 if not schema._is_counter_field(field):
637 continue
638 bucket[field] = bucket.get(field, 0) + value
639 return {
640 source: dict(sorted(metrics.items()))
641 for source, metrics in sorted(totals.items())
642 }
643
644
645 def _discovery_momentum(items: list[schema.SourceItem], to_date: str) -> str:
646 as_of = datetime.fromisoformat(to_date).date()
647 ages: list[int] = []
648 for item in items:
649 try:
650 published = datetime.fromisoformat((item.published_at or "").replace("Z", "+00:00")).date()
651 except (TypeError, ValueError):
652 continue
653 ages.append(max(0, (as_of - published).days))
654 return "new-this-week" if ages and max(ages) < 7 else "building"
655
656
657 def nominate_candidates(
658 plan: schema.DiscoveryPlan,
659 *,
660 from_date: str,
661 to_date: str,
662 depth: str,
663 mock: bool,
664 config: dict[str, Any],
665 lookback_days: int,
666 keyword_gate: bool = True,
667 ) -> schema.RetrievalBundle:
668 """Stage 1 of discovery: fetch, normalize, and bundle candidate hot items
669 from the river/listing feeds.
670
671 This is the topic-nomination pass. For domain discovery ``keyword_gate`` is
672 on and the feeds are filtered to the domain; for global trending it is off
673 and the feeds' own hot ranking IS the signal. The returned bundle feeds the
674 clustering + enrichment stages downstream. Every source's failure is
675 recorded on the bundle (never raised) so a single dead feed cannot sink the
676 run - the confidence floor decides whether the surviving evidence is enough.
677 """
678 bundle = schema.RetrievalBundle()
679 with ThreadPoolExecutor(max_workers=max(1, len(plan.sources))) as executor:
680 futures = {
681 executor.submit(
682 _fetch_discovery_source,
683 source,
684 plan,
685 from_date=from_date,
686 to_date=to_date,
687 depth=depth,
688 mock=mock,
689 config=config,
690 keyword_gate=keyword_gate,
691 ): source
692 for source in plan.sources
693 }
694 for future in as_completed(futures):
695 source = futures[future]
696 bundle.mark_attempted(source)
697 try:
698 raw_items, partial_error = future.result()
699 normalized = normalize.normalize_source_items(
700 source,
701 raw_items,
702 from_date,
703 to_date,
704 freshness_mode="breaking",
705 )
706 # Global trending has no domain; annotate against a neutral
707 # phrase so snippet extraction still works without biasing
708 # relevance toward any keyword.
709 prepared = relevance.PreparedQuery(plan.domain or "trending now")
710 normalized = signals.annotate_stream(
711 normalized,
712 prepared,
713 "breaking",
714 reference_date=to_date,
715 max_days=lookback_days,
716 )
717 normalized = dedupe.dedupe_items(normalized)
718 for item in normalized:
719 item.snippet = snippet.extract_best_snippet(item, prepared)
720 bundle.add_items("discovery-listings", source, normalized)
721 if partial_error:
722 failure_state = (
723 bird_x.classify_run_failure(partial_error)
724 if source == "x" and partial_error.startswith("bird:")
725 else http.classify_failure(message=partial_error)
726 )
727 bundle.record_failure(
728 source,
729 failure_state,
730 partial_error,
731 )
732 except Exception as exc:
733 state, attempted = _classify_source_failure(exc)
734 bundle.record_failure(source, state, str(exc), attempted=attempted)
735 return bundle
736
737
738 @dataclass(frozen=True)
739 class Nomination:
740 """A named candidate topic produced by the nominate stage.
741
742 ``seed_score`` is the cheap pre-enrichment rank - seed velocity on the
743 nominate stage, blended with the HOST judge's content-worthiness on the
744 protocol resume leg (see ``rerank.judge_blended_score``). Enough to
745 decide WHICH candidates deserve a full pipeline pass, but not the final
746 ranking signal (that comes from enriched evidence downstream).
747 ``junk_shape`` flags help-me/beginner/musing shapes that should not
748 become content topics; ``worthiness`` is the host judge's 0-100 content
749 score, None on the heuristic path.
750 """
751
752 name: str
753 seed_score: float
754 items: list[schema.SourceItem] = field(default_factory=list)
755 summary: str = ""
756 junk_shape: bool = False
757 worthiness: float | None = None
758
759
760 def _cluster_entity_counts(
761 cluster: schema.Cluster,
762 candidate_map: dict[str, schema.Candidate],
763 ) -> Counter:
764 """Entity-token frequencies across a cluster's members (title + snippet)."""
765 counts: Counter = Counter()
766 for candidate_id in cluster.candidate_ids:
767 candidate = candidate_map.get(candidate_id)
768 if candidate:
769 counts.update(entity_extract.extract_text_entities(
770 f"{candidate.title} {candidate.snippet}"
771 ))
772 return counts
773
774
775 # Bound on how many distinguishing entity tokens a colliding cluster may try
776 # before it is treated as indistinguishable from the earlier story. Keeps a
777 # pathological cluster (dozens of unique tokens, every resulting name already
778 # taken) from scanning its whole vocabulary.
779 _DISAMBIGUATION_TOKEN_LIMIT = 5
780
781
782 def _disambiguated_topic_name(
783 name: str,
784 cluster: schema.Cluster,
785 earlier_cluster: schema.Cluster,
786 candidate_map: dict[str, schema.Candidate],
787 entity_counts_cache: dict[str, Counter],
788 taken_names: dict[str, schema.Cluster],
789 ) -> str | None:
790 """Disambiguate a colliding topic name by appending the later cluster's
791 strongest entity token that the earlier cluster does not share.
792
793 Distinguishing tokens are tried in descending strength order (bounded at
794 ``_DISAMBIGUATION_TOKEN_LIMIT``) and the first resulting name not already
795 present in ``taken_names`` (casefolded keys) wins: a first-choice suffix
796 colliding with an already-taken name must not drop a distinct story while
797 another distinguishing token remains.
798
799 ``entity_counts_cache`` (keyed by cluster id, owned by the caller) memoizes
800 per-cluster entity counts so repeated collisions against the same cluster
801 never recompute them.
802
803 Returns None when no distinguishing entity yields an unused name - the
804 clusters cannot be told apart by content, so the caller treats them as the
805 same story.
806 """
807 def cached_counts(target: schema.Cluster) -> Counter:
808 counts = entity_counts_cache.get(target.cluster_id)
809 if counts is None:
810 counts = _cluster_entity_counts(target, candidate_map)
811 entity_counts_cache[target.cluster_id] = counts
812 return counts
813
814 later_counts = cached_counts(cluster)
815 earlier_entities = set(cached_counts(earlier_cluster))
816 name_tokens = {token.casefold() for token in name.split()}
817 choices = [
818 (count, token) for token, count in later_counts.items()
819 if token not in earlier_entities and token.casefold() not in name_tokens
820 ]
821 # Strongest first = most frequent across the cluster; alphabetical
822 # tie-break keeps the result deterministic.
823 ranked = sorted(choices, key=lambda entry: (-entry[0], entry[1]))
824 for _, token in ranked[:_DISAMBIGUATION_TOKEN_LIMIT]:
825 display = token
826 for candidate_id in cluster.candidate_ids:
827 candidate = candidate_map.get(candidate_id)
828 if candidate is None:
829 continue
830 match = next(
831 (
832 word.strip("\"'`()[]{}.,:;!?")
833 for word in f"{candidate.title} {candidate.snippet}".split()
834 if word.strip("\"'`()[]{}.,:;!?").lower() == token
835 ),
836 None,
837 )
838 if match:
839 display = match
840 break
841 resolved = f"{name} {display}"
842 if resolved.casefold() not in taken_names:
843 return resolved
844 return None
845
846
847 def nominate_topic_pool(
848 bundle: schema.RetrievalBundle,
849 query_plan: schema.QueryPlan,
850 plan: schema.DiscoveryPlan,
851 *,
852 from_date: str,
853 to_date: str,
854 limit: int,
855 ) -> list[tuple[Nomination, str]]:
856 """Stage 1b of discovery: cluster nominated items into named candidate
857 topics, rank them, and pair each with its source cluster id.
858
859 This is the shared core behind ``nominate_topics`` (the one-shot path,
860 which drops the cluster ids) and the leg-1 nominate-only sweep (which
861 keys nominations-bundle rows on them, see ``run_discover_nominate``).
862
863 Naming and junk classification are the deterministic ``topic_shape``
864 heuristics and ranking is velocity-only - the engine runs no LLM here.
865 Reasoning-model judgment lives in the host-judged protocol: the host
866 renames, junk-filters, and worthiness-scores this pool from the leg-1
867 bundle, and ``run_discover_resume`` applies those verdicts. The one-shot
868 path ships the heuristic names as-is.
869
870 Casefold name collisions are disambiguated (the later cluster's strongest
871 non-shared entity token is appended, trying successive tokens when the
872 first-choice suffix is itself already taken) rather than blindly dropped:
873 short distilled names collide far more often than raw 96-char titles, and
874 a silent drop hides a distinct story. A colliding cluster is dropped only
875 when it shares a representative candidate with the earlier one (the same
876 story surfacing twice) or when no distinguishing entity token yields an
877 unused name.
878
879 Returns at most ``limit`` ``(nomination, cluster_id)`` pairs, never
880 padded - fewer clusters than ``limit`` means a shorter list, and the
881 confidence floor downstream decides whether what survived is worth
882 showing.
883 """
884 candidates = weighted_rrf(
885 bundle.items_by_source_and_query,
886 query_plan,
887 pool_limit=80,
888 range_from=from_date,
889 range_to=to_date,
890 )
891 for candidate in candidates:
892 velocity = rerank.discovery_velocity_score(candidate.source_items, as_of_date=to_date)
893 candidate.final_score = min(100.0, 12.0 * math.log1p(velocity)) if velocity else 0.0
894 candidates.sort(key=lambda candidate: (-candidate.final_score, candidate.title.lower()))
895 clusters = cluster_candidates(candidates, query_plan)
896 candidate_map = {candidate.candidate_id: candidate for candidate in candidates}
897
898 ranked_clusters: list[tuple[float, schema.Cluster, list[schema.SourceItem]]] = []
899 for cluster in clusters:
900 cluster_items: list[schema.SourceItem] = []
901 for candidate_id in cluster.candidate_ids:
902 candidate = candidate_map.get(candidate_id)
903 if candidate:
904 cluster_items.extend(candidate.source_items)
905 score = rerank.discovery_velocity_score(cluster_items, as_of_date=to_date)
906 if score <= 0:
907 continue
908 ranked_clusters.append((score, cluster, cluster_items))
909 ranked_clusters.sort(key=lambda entry: (-entry[0], entry[1].title.lower()))
910
911 # Heuristic naming from each cluster's leader text (title + snippet).
912 named: list[tuple[float, schema.Cluster, list[schema.SourceItem], str, bool]] = []
913 for score, cluster, cluster_items in ranked_clusters:
914 leader = candidate_map.get(cluster.representative_ids[0]) if cluster.representative_ids else None
915 title = (leader.title if leader else cluster.title) or ""
916 snip = (leader.snippet if leader else "") or ""
917 name = topic_shape.distill_topic_name(title, snip) or plan.domain or title
918 junk_shape = topic_shape.is_junk_shape(title, snip)
919 named.append((score, cluster, cluster_items, name, junk_shape))
920 named.sort(key=lambda entry: (-entry[0], entry[3].lower()))
921
922 pool: list[tuple[Nomination, str]] = []
923 taken_names: dict[str, schema.Cluster] = {}
924 entity_counts_cache: dict[str, Counter] = {}
925 for score, cluster, cluster_items, name, junk_shape in named:
926 name_key = name.casefold()
927 if name_key in taken_names:
928 earlier_cluster = taken_names[name_key]
929 if set(cluster.representative_ids) & set(earlier_cluster.representative_ids):
930 continue # same story surfacing twice
931 resolved = _disambiguated_topic_name(
932 name, cluster, earlier_cluster, candidate_map, entity_counts_cache,
933 taken_names,
934 )
935 if resolved is None:
936 continue # indistinguishable by content: treat as the same story
937 name = resolved
938 name_key = name.casefold()
939 taken_names[name_key] = cluster
940 leader = candidate_map.get(cluster.representative_ids[0]) if cluster.representative_ids else None
941 summary = (leader.snippet if leader else "") or (leader.title if leader else name)
942 pool.append((Nomination(
943 name=name,
944 seed_score=score,
945 items=cluster_items,
946 summary=summary,
947 junk_shape=junk_shape,
948 ), cluster.cluster_id))
949 if len(pool) >= limit:
950 break
951 return pool
952
953
954 def nominate_topics(
955 bundle: schema.RetrievalBundle,
956 query_plan: schema.QueryPlan,
957 plan: schema.DiscoveryPlan,
958 *,
959 from_date: str,
960 to_date: str,
961 limit: int,
962 ) -> list[Nomination]:
963 """``nominate_topic_pool`` without the cluster ids: the one-shot
964 discovery path's contract (see that function for the full semantics)."""
965 return [
966 nomination
967 for nomination, _cluster_id in nominate_topic_pool(
968 bundle, query_plan, plan, from_date=from_date, to_date=to_date, limit=limit,
969 )
970 ]
971
972
973 # Enrichment fan-out bounds. Sub-runs hit the same upstream APIs as a normal
974 # research pass, so parallelism stays low and the whole batch runs against a
975 # wall-clock budget - a slow topic is dropped, never fatal.
976 ENRICH_LIMIT = 6
977 ENRICH_DEPTH = "quick"
978 ENRICH_MAX_WORKERS = 3
979 ENRICH_BUDGET_SECONDS = 240.0
980
981
982 @dataclass
983 class EnrichedTopic:
984 """A nomination plus the full-pipeline evidence gathered for it.
985
986 ``report`` is None when enrichment for this topic failed or ran past the
987 batch budget - the topic survives as nomination-only and the confidence
988 floor downstream decides whether its seed evidence is enough to show.
989 """
990
991 nomination: Nomination
992 report: schema.Report | None = None
993 error: str | None = None
994
995
996 def enrich_nominations(
997 nominations: list[Nomination],
998 *,
999 config: dict[str, Any],
1000 requested_sources: list[str] | None = None,
1001 mock: bool = False,
1002 depth: str = ENRICH_DEPTH,
1003 lookback_days: int = 30,
1004 as_of_date: str | None = None,
1005 max_workers: int = ENRICH_MAX_WORKERS,
1006 budget_seconds: float = ENRICH_BUDGET_SECONDS,
1007 ) -> list[EnrichedTopic]:
1008 """Stage 2 of discovery: run the real research pipeline on each nomination.
1009
1010 Each nominated topic gets a full ``run()`` pass (``internal_subrun=True``,
1011 same lane as comparison-mode sub-runs), which buys the whole multi-source
1012 corpus - Reddit with comments, X, YouTube, Techmeme, arXiv, HN, Polymarket,
1013 web - plus clustering and ranking, with zero bespoke fetch code.
1014
1015 Failure containment: a topic whose sub-run raises is returned with
1016 ``report=None`` and the error recorded; topics still unfinished when the
1017 batch budget expires are likewise dropped to nomination-only. The batch
1018 never raises and preserves nomination order.
1019 """
1020 if not nominations:
1021 return []
1022
1023 def _run_one(nomination: Nomination) -> schema.Report:
1024 # Per-worker copy: run() mutates config in place
1025 # (config["_financial_topic"] = ...), so sharing one dict across
1026 # daemon threads lets topic A's flag overwrite topic B's mid-run,
1027 # with stragglers mutating past budget expiry. Same idiom as the
1028 # competitor runner (entity_config = dict(config)).
1029 return run(
1030 topic=nomination.name,
1031 config=dict(config),
1032 depth=depth,
1033 requested_sources=requested_sources,
1034 mock=mock,
1035 lookback_days=lookback_days,
1036 as_of_date=as_of_date,
1037 internal_subrun=True,
1038 # Enrichment passes never carry a connector envelope, so the
1039 # per-session lane signal must not plan X in and record a
1040 # spurious X error on every nominated topic.
1041 suppress_x_host_lane=True,
1042 )
1043
1044 # Daemon threads + a semaphore instead of ThreadPoolExecutor: executor
1045 # threads are non-daemon and joined at interpreter shutdown, so one hung
1046 # sub-run could keep the whole process alive long after its topic was
1047 # downgraded to nomination-only. Daemon workers make the wall-clock budget
1048 # real - stragglers cannot delay process exit. Abandonment is safe because
1049 # internal_subrun passes write nothing to disk (no save, no library sync,
1050 # no store), and every fetch layer inside run() carries its own timeout.
1051 youtube_yt.reset_search_cache()
1052 enriched: dict[str, EnrichedTopic] = {}
1053 results_queue: queue.Queue[tuple[Nomination, schema.Report | None, Exception | None]] = queue.Queue()
1054 slots = threading.Semaphore(max(1, max_workers))
1055
1056 def _worker(nomination: Nomination) -> None:
1057 with slots:
1058 try:
1059 results_queue.put((nomination, _run_one(nomination), None))
1060 except Exception as exc: # noqa: BLE001 - containment is the contract
1061 results_queue.put((nomination, None, exc))
1062
1063 for nomination in nominations:
1064 threading.Thread(
1065 target=_worker,
1066 args=(nomination,),
1067 name=f"discover-enrich-{nomination.name[:32]}",
1068 daemon=True,
1069 ).start()
1070
1071 deadline = time.monotonic() + max(1.0, budget_seconds)
1072 pending = len(nominations)
1073 while pending and (remaining := deadline - time.monotonic()) > 0:
1074 try:
1075 nomination, report, exc = results_queue.get(timeout=min(remaining, 0.5))
1076 except queue.Empty:
1077 continue
1078 pending -= 1
1079 if exc is None:
1080 enriched[nomination.name] = EnrichedTopic(
1081 nomination=nomination, report=report,
1082 )
1083 else:
1084 enriched[nomination.name] = EnrichedTopic(
1085 nomination=nomination,
1086 error=f"{type(exc).__name__}: {exc}",
1087 )
1088 print(
1089 f"[Discover] enrichment failed for {nomination.name!r}: "
1090 f"{type(exc).__name__}: {exc}",
1091 file=sys.stderr,
1092 )
1093 # Budget expired (or all done): unfinished topics fall through below as
1094 # nomination-only; their daemon workers are abandoned and cannot block exit.
1095
1096 results: list[EnrichedTopic] = []
1097 for nomination in nominations:
1098 entry = enriched.get(nomination.name)
1099 if entry is None:
1100 entry = EnrichedTopic(
1101 nomination=nomination,
1102 error="enrichment budget exhausted",
1103 )
1104 print(
1105 f"[Discover] enrichment budget exhausted before {nomination.name!r} "
1106 "finished; keeping nomination-only evidence",
1107 file=sys.stderr,
1108 )
1109 results.append(entry)
1110 return results
1111
1112
1113 def _enriched_evidence_items(entry: EnrichedTopic) -> list[schema.SourceItem]:
1114 """The items a topic is judged on: the enriched corpus when the pipeline
1115 pass succeeded, the nomination's seed items otherwise."""
1116 if entry.report is not None:
1117 flattened: list[schema.SourceItem] = []
1118 for source_items in entry.report.items_by_source.values():
1119 flattened.extend(source_items)
1120 if flattened:
1121 return flattened
1122 return entry.nomination.items
1123
1124
1125 def _best_community_comment(items: list[schema.SourceItem]) -> str | None:
1126 """The strongest verbatim community comment across a topic's evidence,
1127 formatted with attribution - the voice-of-the-people line on a trend card.
1128
1129 Vote strength is per-platform-normalized (signals.normalized_comment_vote)
1130 so one viral platform's counts don't drown out the rest.
1131 """
1132 best: tuple[float, str, str | None, float | int | None] | None = None
1133 for item in items:
1134 comments = item.metadata.get("top_comments") or []
1135 for comment in comments:
1136 if not isinstance(comment, dict):
1137 continue
1138 body = (comment.get("excerpt") or comment.get("text") or comment.get("body") or "").strip()
1139 if len(body) < 12:
1140 continue
1141 strength = signals.normalized_comment_vote(item.source, comment.get("score"))
1142 if best is None or strength > best[0]:
1143 best = (strength, body, comment.get("author"), comment.get("score"))
1144 if best is None:
1145 return None
1146 _, body, author, score = best
1147 # Comment bodies that themselves start/end with quote characters would
1148 # render as doubled quotes inside our wrapping quotes.
1149 body = body.strip('"“”‘’\'').strip()
1150 if len(body) > 200:
1151 body = body[:197].rsplit(" ", 1)[0] + "..."
1152 attribution = f" - {author}" if author else ""
1153 votes = (
1154 f" ({int(score):,} votes)"
1155 if isinstance(score, (int, float)) and not isinstance(score, bool) and score > 0
1156 else ""
1157 )
1158 return f'"{body}"{attribution}{votes}'
1159
1160
1161 @dataclass(frozen=True)
1162 class _DiscoverySweep:
1163 """The shared front half of both discovery entry points: the resolved
1164 plan and window, the swept listing bundle, and finalized per-source
1165 status. Everything downstream (judging, enrichment, floor, queue)
1166 belongs to the caller's leg."""
1167
1168 plan: schema.DiscoveryPlan
1169 query_plan: schema.QueryPlan
1170 from_date: str
1171 to_date: str
1172 bundle: schema.RetrievalBundle
1173 source_status: dict[str, schema.SourceOutcome]
1174
1175
1176 def _discovery_sweep(
1177 *,
1178 domain: str,
1179 config: dict[str, Any],
1180 depth: str,
1181 requested_sources: list[str] | None,
1182 mock: bool,
1183 subreddits: list[str] | None,
1184 lookback_days: int,
1185 as_of_date: str | None,
1186 ) -> _DiscoverySweep:
1187 """Resolve the momentum window, validate/bound the listing sources, build
1188 the discovery plan, sweep the river feeds, and finalize source status.
1189
1190 Shared verbatim by ``run_discover`` (one-shot) and
1191 ``run_discover_nominate`` (protocol leg 1) so the two paths can never
1192 drift on what a sweep means."""
1193 from_date, to_date = dates.get_date_range(lookback_days, as_of_date=as_of_date)
1194 requested = normalize_requested_sources(requested_sources)
1195 unsupported = sorted(set(requested or []) - set(DISCOVERY_SOURCES))
1196 if unsupported:
1197 raise ValueError(
1198 "Discovery supports listing sources only: reddit, hackernews, digg "
1199 f"(unsupported: {', '.join(unsupported)})"
1200 )
1201 available = list(DISCOVERY_SOURCES) if mock else [
1202 source for source in available_sources(config, requested, x_pending=False)
1203 if source in DISCOVERY_SOURCES
1204 ]
1205 if requested:
1206 available = [source for source in available if source in requested]
1207 plan = planner.build_discovery_plan(
1208 domain,
1209 available_sources=available,
1210 subreddits=subreddits,
1211 )
1212
1213 global_mode = not plan.domain
1214 domain_label = plan.domain or "everything"
1215 query_plan = schema.QueryPlan(
1216 intent="breaking_news",
1217 freshness_mode="breaking",
1218 cluster_mode="story",
1219 raw_topic=plan.domain,
1220 subqueries=[schema.SubQuery(
1221 label="discovery-listings",
1222 search_query=plan.domain,
1223 ranking_query=f"What is accelerating in {domain_label}?",
1224 sources=list(plan.sources),
1225 )],
1226 source_weights={source: 1.0 for source in plan.sources},
1227 notes=["discover-mode", "listing-sweep"],
1228 )
1229
1230 bundle = nominate_candidates(
1231 plan,
1232 from_date=from_date,
1233 to_date=to_date,
1234 depth=depth,
1235 mock=mock,
1236 config=config,
1237 lookback_days=lookback_days,
1238 # Global trending has no keyword to gate against - the river feeds' own
1239 # hot ranking is the signal and the confidence floor culls the junk.
1240 keyword_gate=not global_mode,
1241 )
1242
1243 source_status: dict[str, schema.SourceOutcome] = {}
1244 for source in DISCOVERY_SOURCES:
1245 if source in bundle.source_status:
1246 continue
1247 detail = (
1248 "Source is not configured for discovery."
1249 )
1250 source_status[source] = schema.SourceOutcome(
1251 source=source,
1252 state=schema.SKIPPED_UNCONFIGURED,
1253 attempted=False,
1254 detail=detail,
1255 fix_hint="doctor",
1256 )
1257 source_status.update(_finalize_source_status(bundle.source_status, bundle.items_by_source))
1258 return _DiscoverySweep(
1259 plan=plan,
1260 query_plan=query_plan,
1261 from_date=from_date,
1262 to_date=to_date,
1263 bundle=bundle,
1264 source_status=source_status,
1265 )
1266
1267
1268 def _degraded_discovery_sources(
1269 source_status: dict[str, schema.SourceOutcome],
1270 ) -> list[str]:
1271 """Sources whose outcome is neither clean nor an expected skip."""
1272 return [
1273 source for source, outcome_state in source_status.items()
1274 if outcome_state.state not in {health.OK, schema.NO_RESULTS, schema.SKIPPED_UNCONFIGURED}
1275 ]
1276
1277
1278 @dataclass(frozen=True)
1279 class DiscoverNominateResult:
1280 """Leg 1 output of the host-judged discovery protocol: the ranked judge
1281 pool as ``(nomination, cluster_id)`` pairs plus the sweep context the CLI
1282 needs to write the nominations bundle - or to render the nothing-solid
1283 brief when the pool is empty."""
1284
1285 plan: schema.DiscoveryPlan
1286 from_date: str
1287 to_date: str
1288 source_status: dict[str, schema.SourceOutcome]
1289 pool: list[tuple[Nomination, str]]
1290
1291
1292 def run_discover_nominate(
1293 *,
1294 domain: str,
1295 config: dict[str, Any],
1296 depth: str = "default",
1297 requested_sources: list[str] | None = None,
1298 mock: bool = False,
1299 subreddits: list[str] | None = None,
1300 lookback_days: int = 30,
1301 as_of_date: str | None = None,
1302 ) -> DiscoverNominateResult:
1303 """Protocol leg 1: sweep the listings and build the FULL judge pool.
1304
1305 Same sweep and clustering as ``run_discover``, but the pool is cut at
1306 ``rerank.JUDGE_POOL_LIMIT`` (not the enrichment limit). Like every
1307 discovery path it is deterministic-heuristic: no provider is ever
1308 resolved, so names and junk flags are the ``topic_shape`` baselines the
1309 host judges against. No enrichment, no confidence floor, no queue
1310 writes - those belong to legs 2 and 3.
1311 """
1312 sweep = _discovery_sweep(
1313 domain=domain,
1314 config=config,
1315 depth=depth,
1316 requested_sources=requested_sources,
1317 mock=mock,
1318 subreddits=subreddits,
1319 lookback_days=lookback_days,
1320 as_of_date=as_of_date,
1321 )
1322 pool = nominate_topic_pool(
1323 sweep.bundle, sweep.query_plan, sweep.plan,
1324 from_date=sweep.from_date,
1325 to_date=sweep.to_date,
1326 limit=rerank.JUDGE_POOL_LIMIT,
1327 )
1328 return DiscoverNominateResult(
1329 plan=sweep.plan,
1330 from_date=sweep.from_date,
1331 to_date=sweep.to_date,
1332 source_status=sweep.source_status,
1333 pool=pool,
1334 )
1335
1336
1337 def nominate_nothing_solid_report(result: DiscoverNominateResult) -> schema.DiscoveryReport:
1338 """The honest-empty leg-1 report: a zero-nomination sweep renders the
1339 same nothing-solid brief a one-shot run would (and writes no bundle)."""
1340 warnings = [
1341 "The listing sweep nominated no topics this window; reporting "
1342 "nothing solid instead of ranked noise."
1343 ]
1344 failed = _degraded_discovery_sources(result.source_status)
1345 if failed:
1346 warnings.append(f"Some discovery sources degraded: {', '.join(sorted(failed))}.")
1347 return schema.DiscoveryReport(
1348 domain=result.plan.domain,
1349 range_from=result.from_date,
1350 range_to=result.to_date,
1351 generated_at=datetime.now(timezone.utc).isoformat(),
1352 plan=result.plan,
1353 topics=[],
1354 source_status=result.source_status,
1355 warnings=warnings,
1356 outcome="nothing-solid",
1357 weak_signal=None,
1358 )
1359
1360
1361 def _floor_survivor_records(
1362 enriched_entries: list[EnrichedTopic],
1363 *,
1364 to_date: str,
1365 topic_limit: int,
1366 ) -> tuple[
1367 list[dict[str, Any]],
1368 tuple[float, str] | None,
1369 tuple[float, str] | None,
1370 ]:
1371 """Apply the discovery confidence floor to enriched entries in order,
1372 returning the survivor records plus the strongest non-junk and junk weak
1373 signals among the failures.
1374
1375 Shared verbatim by ``run_discover`` (one-shot) and ``run_discover_resume``
1376 (protocol leg 2) so floor semantics can never drift between the paths.
1377 """
1378 survivors: list[dict[str, Any]] = []
1379 weak_signal: tuple[float, str] | None = None
1380 junk_weak_signal: tuple[float, str] | None = None
1381 for entry in enriched_entries:
1382 nomination = entry.nomination
1383 evidence_items = _enriched_evidence_items(entry)
1384 sources = sorted({item.source for item in evidence_items})
1385 native_total = sum(
1386 rerank.discovery_engagement_total(item) for item in evidence_items
1387 )
1388 score = rerank.discovery_velocity_score(evidence_items, as_of_date=to_date)
1389 if not rerank.passes_discovery_floor(
1390 source_count=len(sources),
1391 engagement_total=native_total,
1392 item_count=len(evidence_items),
1393 junk_shape=nomination.junk_shape,
1394 # Junk corroboration counts distinct SEED listing sources, never
1395 # the enriched corpus - a successful enrichment pass is
1396 # multi-source for almost any topic, so it would never bind.
1397 seed_source_count=len({item.source for item in nomination.items}),
1398 ):
1399 # Sub-floor evidence never ranks; remember what came closest so a
1400 # nothing-solid brief can still name the strongest weak signal.
1401 # Junk-shaped failures are tracked separately: the brief prefers
1402 # the strongest NON-junk failure and names a junk one only when
1403 # every failure is junk-shaped (never empty when failures exist).
1404 if nomination.junk_shape:
1405 if junk_weak_signal is None or score > junk_weak_signal[0]:
1406 junk_weak_signal = (score, nomination.name)
1407 elif weak_signal is None or score > weak_signal[0]:
1408 weak_signal = (score, nomination.name)
1409 continue
1410 if len(survivors) >= topic_limit:
1411 break
1412 source_phrase = ", ".join(sources[:-1]) + (
1413 f" and {sources[-1]}" if len(sources) > 1 else (sources[0] if sources else "the listings")
1414 )
1415 noun = "evidence item" if entry.report is not None else "listing item"
1416 why = (
1417 f"{len(evidence_items)} {noun}{'s' if len(evidence_items) != 1 else ''} on "
1418 f"{source_phrase} generated {native_total:,.0f} native interactions. "
1419 f"{nomination.summary[:220]}"
1420 )
1421 top_comment = _best_community_comment(evidence_items) if entry.report is not None else None
1422 # Stage-2 angle input: the survivor's strongest evidence, enriched
1423 # corpus when the pipeline pass succeeded, seed items otherwise
1424 # (evidence_items already resolves that).
1425 top_titles = [
1426 item.title.strip()
1427 for item in sorted(
1428 evidence_items,
1429 key=rerank.discovery_engagement_total,
1430 reverse=True,
1431 )
1432 if item.title and item.title.strip()
1433 ][:3]
1434 survivors.append({
1435 "name": nomination.name,
1436 "why": why,
1437 "momentum": _discovery_momentum(evidence_items, to_date),
1438 "velocity_score": round(score, 2),
1439 "sources": sources,
1440 "engagement_by_source": _discovery_engagement(evidence_items),
1441 "evidence_urls": list(dict.fromkeys(item.url for item in evidence_items if item.url))[:5],
1442 "top_comment": top_comment,
1443 "titles": "; ".join(top_titles),
1444 "engagement_phrase": f"{native_total:,.0f} native interactions across {source_phrase}",
1445 })
1446 return survivors, weak_signal, junk_weak_signal
1447
1448
1449 def _fold_same_story_records(survivors: list[dict[str, Any]]) -> list[dict[str, Any]]:
1450 """Same-story fold + velocity ordering over floor-survivor records.
1451
1452 Floor survivors that share enriched evidence are the SAME story wearing
1453 two judged names (the real-run failure: two topics quoting the identical
1454 1,635-vote comment). Duplicates = identical non-None top comment OR >= 2
1455 shared evidence URLs; the lower-velocity twin is dropped, and a winning
1456 replacement re-scans the kept list to a fixpoint so chained overlap
1457 (A~C~B) still collapses to one survivor. Selection stays seed-ordered
1458 upstream; this only prunes, then sorts by displayed velocity (stable) so
1459 rank 1 is the highest velocity_score.
1460 """
1461 def _same_story(a: dict[str, Any], b: dict[str, Any]) -> bool:
1462 if a["top_comment"] is not None and a["top_comment"] == b["top_comment"]:
1463 return True
1464 return len(set(a["evidence_urls"]) & set(b["evidence_urls"])) >= 2
1465
1466 folded: list[dict[str, Any]] = []
1467 for record in survivors:
1468 # Fold to a fixpoint: when the incoming record REPLACES a kept one,
1469 # the replacement may share evidence with entries the dropped record
1470 # never matched (three-way chains: A kept, C shares the comment with
1471 # A and URLs with B). The winner re-scans the remaining kept entries
1472 # until nothing matches, so one story always yields one survivor.
1473 incoming: dict[str, Any] | None = record
1474 while incoming is not None:
1475 dup_index = next(
1476 (index for index, kept in enumerate(folded) if _same_story(incoming, kept)),
1477 None,
1478 )
1479 if dup_index is None:
1480 folded.append(incoming)
1481 break
1482 kept = folded[dup_index]
1483 if incoming["velocity_score"] > kept["velocity_score"]:
1484 folded.pop(dup_index)
1485 dropped_name, kept_name = kept["name"], incoming["name"]
1486 else:
1487 dropped_name, kept_name = incoming["name"], kept["name"]
1488 incoming = None # dropped; the kept entry stays in place
1489 log.source_log(
1490 "Discover",
1491 f"folded duplicate story {dropped_name!r} into {kept_name!r} (shared evidence)",
1492 tty_only=False,
1493 )
1494
1495 folded.sort(key=lambda record: record["velocity_score"], reverse=True)
1496 return folded
1497
1498
1499 def _records_to_discovery_topics(
1500 folded: list[dict[str, Any]],
1501 ) -> list[schema.DiscoveryTopic]:
1502 """Folded survivor records to ranked topics (ranks = 1-based positions)."""
1503 return [
1504 schema.DiscoveryTopic(
1505 rank=position,
1506 name=record["name"],
1507 why_spiking=record["why"],
1508 momentum=record["momentum"],
1509 velocity_score=record["velocity_score"],
1510 sources=record["sources"],
1511 engagement_by_source=record["engagement_by_source"],
1512 command=f'/last30days "{record["name"].replace(chr(34), chr(39))}"',
1513 evidence_urls=record["evidence_urls"],
1514 top_comment=record["top_comment"],
1515 corroboration_count=len(record["sources"]),
1516 )
1517 for position, record in enumerate(folded, start=1)
1518 ]
1519
1520
1521 def _discovery_report_warnings(
1522 topics: list[schema.DiscoveryTopic],
1523 outcome: str,
1524 source_status: dict[str, schema.SourceOutcome],
1525 ) -> list[str]:
1526 """Coverage warnings shared by the one-shot and resume discovery paths.
1527 The resume leg never re-sweeps: it passes the bundle's RESTORED leg-1
1528 sweep status, so a degraded feed from the sweep still reaches the leg-2
1529 report exactly as the one-shot reports it."""
1530 warnings: list[str] = []
1531 if outcome == "nothing-solid":
1532 warnings.append(
1533 "No topic cleared the discovery confidence floor this window; "
1534 "reporting nothing solid instead of ranked noise."
1535 )
1536 elif len(topics) < 5:
1537 warnings.append("Fewer than five topic clusters cleared the confidence floor this window.")
1538 if topics and all(len(topic.sources) == 1 for topic in topics):
1539 warnings.append("Discovery evidence is single-source; configure Digg for broader confirmation.")
1540 failed = _degraded_discovery_sources(source_status)
1541 if failed:
1542 warnings.append(f"Some discovery sources degraded: {', '.join(sorted(failed))}.")
1543 return warnings
1544
1545
1546 def run_discover(
1547 *,
1548 domain: str,
1549 config: dict[str, Any],
1550 depth: str = "default",
1551 requested_sources: list[str] | None = None,
1552 mock: bool = False,
1553 subreddits: list[str] | None = None,
1554 lookback_days: int = 30,
1555 as_of_date: str | None = None,
1556 limit: int = 10,
1557 enrich: bool = False,
1558 enrich_requested_sources: list[str] | None = None,
1559 ) -> schema.DiscoveryReport:
1560 """Sweep category listings and rank the topics gaining velocity.
1561
1562 ``requested_sources`` bounds the listing sweep (discovery-capable feeds
1563 only). ``enrich_requested_sources`` bounds the per-topic research passes:
1564 None means every available source - which is what lets Techmeme, arXiv,
1565 YouTube, Polymarket, and community comments reach discovery despite having
1566 no river feed of their own. Pass the user's original --search list here so
1567 an explicit source boundary holds through enrichment too.
1568 """
1569 sweep = _discovery_sweep(
1570 domain=domain,
1571 config=config,
1572 depth=depth,
1573 requested_sources=requested_sources,
1574 mock=mock,
1575 subreddits=subreddits,
1576 lookback_days=lookback_days,
1577 as_of_date=as_of_date,
1578 )
1579 plan = sweep.plan
1580 from_date, to_date = sweep.from_date, sweep.to_date
1581 source_status = sweep.source_status
1582
1583 # The engine never names or angles topics with an LLM: the one-shot path
1584 # is deterministic-heuristic by design, and reasoning-model judgment
1585 # lives in the host-judged SKILL.md protocol. Say so loudly once per live
1586 # run; --mock stays silent (a deliberate mock run is not a degraded run).
1587 if not mock:
1588 log.source_log(
1589 "Discover",
1590 "one-shot run: topic names use deterministic heuristics and no "
1591 "content angles are generated - a reasoning-model host running "
1592 "the SKILL.md discovery protocol gets host-judged names, junk "
1593 "filtering, and podcast/X angles",
1594 tty_only=False,
1595 )
1596
1597 topic_limit = max(5, min(10, limit))
1598 nominations = nominate_topics(
1599 sweep.bundle, sweep.query_plan, plan,
1600 from_date=from_date,
1601 to_date=to_date,
1602 limit=ENRICH_LIMIT if enrich else topic_limit,
1603 )
1604
1605 if enrich and nominations:
1606 enriched_entries = enrich_nominations(
1607 nominations,
1608 config=config,
1609 requested_sources=enrich_requested_sources,
1610 mock=mock,
1611 lookback_days=lookback_days,
1612 as_of_date=as_of_date,
1613 )
1614 else:
1615 enriched_entries = [
1616 EnrichedTopic(nomination=nomination) for nomination in nominations
1617 ]
1618
1619 survivors, weak_signal, junk_weak_signal = _floor_survivor_records(
1620 enriched_entries, to_date=to_date, topic_limit=topic_limit,
1621 )
1622 folded = _fold_same_story_records(survivors)
1623 # One-shot topics ship without angles (podcast_angle / x_article_angle
1624 # stay None and the renderer omits those lines): content angles are a
1625 # host-judged protocol deliverable, written on the finalize leg.
1626 topics = _records_to_discovery_topics(folded)
1627
1628 if weak_signal is None:
1629 weak_signal = junk_weak_signal
1630
1631 outcome = "ok" if topics else "nothing-solid"
1632
1633 return schema.DiscoveryReport(
1634 domain=plan.domain,
1635 range_from=from_date,
1636 range_to=to_date,
1637 generated_at=datetime.now(timezone.utc).isoformat(),
1638 plan=plan,
1639 topics=topics,
1640 source_status=source_status,
1641 warnings=_discovery_report_warnings(topics, outcome, source_status),
1642 outcome=outcome,
1643 weak_signal=weak_signal[1] if weak_signal and not topics else None,
1644 )
1645
1646
1647 # Protocol leg 2 (resume) deep-tier enrichment bounds. The module-level
1648 # ENRICH_* constants above stay the one-shot --discover contract (quick depth,
1649 # 240s budget, 3 workers); a deep-tier bundle upgrades its per-topic sub-runs
1650 # to the default research depth with a wider wall-clock budget and one more
1651 # worker, because leg 2 is the protocol's only research pass. Shallow-tier
1652 # bundles keep the one-shot quick constants. Both tiers flow through
1653 # enrich_nominations' PARAMETERS - the constants themselves are never edited,
1654 # so neither tier can leak into the other path.
1655 RESUME_DEEP_ENRICH_DEPTH = "default"
1656 RESUME_DEEP_ENRICH_MAX_WORKERS = 4
1657 RESUME_DEEP_ENRICH_BUDGET_SECONDS = 450.0
1658
1659
1660 def _resume_enrich_budget_seconds(config: dict[str, Any]) -> float:
1661 """Deep-tier batch budget: LAST30DAYS_ENRICH_BUDGET_SECONDS from the
1662 RESOLVED config dict only (env.get_config already layers the process env
1663 over the .env files) - never read from bare os.environ. Blank,
1664 non-numeric, or non-positive values fall back to the 450s default."""
1665 raw = config.get("LAST30DAYS_ENRICH_BUDGET_SECONDS")
1666 if raw is None or str(raw).strip() == "":
1667 return RESUME_DEEP_ENRICH_BUDGET_SECONDS
1668 try:
1669 value = float(raw)
1670 except (TypeError, ValueError):
1671 return RESUME_DEEP_ENRICH_BUDGET_SECONDS
1672 return value if value > 0 else RESUME_DEEP_ENRICH_BUDGET_SECONDS
1673
1674
1675 @dataclass(frozen=True)
1676 class DiscoverResumeResult:
1677 """Leg 2 output of the host-judged discovery protocol: the floored,
1678 folded, velocity-ranked report plus the per-topic angle inputs (keyed by
1679 surviving nomination id) that the host writes leg-3 angles from.
1680 ``report.source_status`` is the bundle's restored leg-1 sweep status -
1681 leg 2 never re-sweeps the listing feeds, so the sweep's degraded-coverage
1682 signal must survive the handoff instead of reading as clean."""
1683
1684 report: schema.DiscoveryReport
1685 angle_inputs: dict[str, dict[str, str]]
1686
1687
1688 def run_discover_resume(
1689 bundle: Any,
1690 judgments: dict[str, Any],
1691 *,
1692 config: dict[str, Any],
1693 mock: bool = False,
1694 ) -> DiscoverResumeResult:
1695 """Protocol leg 2: apply host judgments to the leg-1 bundle, enrich the
1696 slot winners, and floor/fold/rank on the same code path as the one-shot
1697 run.
1698
1699 ``bundle`` is a ``discovery_handoff.NominationsBundle`` and ``judgments``
1700 the mapping ``discovery_handoff.read_judgments`` returns (annotated
1701 loosely because discovery_handoff imports this module at load time).
1702
1703 Judgment application is per field: an absent host name falls back to the
1704 bundle's heuristic name, an absent junk flag to the heuristic junk flag,
1705 and absent worthiness to the neutral blend default (None -> 50 inside
1706 ``rerank.judge_blended_score`` - the same treatment the judge-absent path
1707 always used). Applied names are collision-resolved over the whole pool
1708 before anything keys on them, and the applied name IS the enrichment
1709 sub-run topic.
1710
1711 Slot selection: host-junk rows never contend for enrichment slots, and a
1712 heuristic-junk fallback row with fewer than ``rerank.FLOOR_MIN_SOURCES``
1713 distinct seed sources is skipped pre-enrichment (it structurally cannot
1714 pass the floor's seed-corroboration rule). Both stay eligible to be the
1715 junk-tracked weak signal of a nothing-solid brief, and the brief prefers
1716 a non-junk weak signal exactly like the one-shot path. At the floor,
1717 host-judged rows pass ``junk_shape=False`` (host-junk never earned a
1718 slot) while heuristic-fallback rows keep their heuristic flag with the
1719 existing seed-source corroboration.
1720
1721 Velocity, momentum, and the enrichment window all score against the
1722 bundle's momentum window (from_date/to_date), never the resume-time
1723 clock: the host may judge up to the handoff TTL after the sweep, and the
1724 numbers must describe the window the sweep captured.
1725 """
1726 # Runtime-only import: discovery_handoff imports pipeline at module load,
1727 # so the reverse import must happen at call time (no import-time cycle).
1728 from . import discovery_handoff
1729
1730 to_date = bundle.to_date
1731 verdicts = [
1732 discovery_handoff.judgment_for(judgments, entry.nomination_id)
1733 for entry in bundle.nominations
1734 ]
1735 applied_names = discovery_handoff.resolve_name_collisions([
1736 (
1737 entry.nomination,
1738 verdict.name or entry.heuristic_name or entry.nomination.name,
1739 )
1740 for entry, verdict in zip(bundle.nominations, verdicts)
1741 ])
1742
1743 ranked: list[tuple[float, str, Nomination]] = []
1744 junk_weak_signal: tuple[float, str] | None = None
1745 for entry, verdict, name in zip(bundle.nominations, verdicts, applied_names):
1746 items = entry.nomination.items
1747 velocity = rerank.discovery_velocity_score(items, as_of_date=to_date)
1748 seed_source_count = len({item.source for item in items})
1749 host_junk = verdict.junk is True
1750 fallback_junk = verdict.junk is None and entry.heuristic_junk
1751 if host_junk or (
1752 fallback_junk and seed_source_count < rerank.FLOOR_MIN_SOURCES
1753 ):
1754 if junk_weak_signal is None or velocity > junk_weak_signal[0]:
1755 junk_weak_signal = (velocity, name)
1756 continue
1757 worthiness = (
1758 float(verdict.worthiness) if verdict.worthiness is not None else None
1759 )
1760 blended = rerank.judge_blended_score(velocity, worthiness)
1761 ranked.append((
1762 blended,
1763 entry.nomination_id,
1764 replace(
1765 entry.nomination,
1766 name=name,
1767 seed_score=blended,
1768 junk_shape=(
1769 False if verdict.junk is not None else entry.heuristic_junk
1770 ),
1771 worthiness=worthiness,
1772 ),
1773 ))
1774
1775 ranked.sort(key=lambda row: (-row[0], row[2].name.lower()))
1776 selected = ranked[:ENRICH_LIMIT]
1777 nominations = [nomination for _blended, _nomination_id, nomination in selected]
1778
1779 if bundle.tier == "shallow":
1780 depth, max_workers, budget_seconds = (
1781 ENRICH_DEPTH, ENRICH_MAX_WORKERS, ENRICH_BUDGET_SECONDS,
1782 )
1783 else:
1784 depth = RESUME_DEEP_ENRICH_DEPTH
1785 max_workers = RESUME_DEEP_ENRICH_MAX_WORKERS
1786 budget_seconds = _resume_enrich_budget_seconds(config)
1787
1788 enriched_entries = enrich_nominations(
1789 nominations,
1790 config=config,
1791 requested_sources=bundle.enrichment_source_boundary,
1792 mock=mock,
1793 depth=depth,
1794 lookback_days=bundle.lookback_days,
1795 as_of_date=to_date,
1796 max_workers=max_workers,
1797 budget_seconds=budget_seconds,
1798 ) if nominations else []
1799
1800 # topic_limit mirrors the one-shot default cap (limit=10); the slot cut
1801 # above already bounds the pool at ENRICH_LIMIT.
1802 survivors, weak_signal, floor_junk_weak_signal = _floor_survivor_records(
1803 enriched_entries, to_date=to_date, topic_limit=10,
1804 )
1805 if floor_junk_weak_signal is not None and (
1806 junk_weak_signal is None
1807 or floor_junk_weak_signal[0] > junk_weak_signal[0]
1808 ):
1809 junk_weak_signal = floor_junk_weak_signal
1810 folded = _fold_same_story_records(survivors)
1811 topics = _records_to_discovery_topics(folded)
1812
1813 nomination_id_by_name = {
1814 nomination.name: nomination_id
1815 for _blended, nomination_id, nomination in selected
1816 }
1817 angle_inputs = {
1818 nomination_id_by_name[record["name"]]: {
1819 "name": record["name"],
1820 "titles": record["titles"],
1821 "top_comment": record["top_comment"] or "",
1822 "engagement": record["engagement_phrase"],
1823 }
1824 for record in folded
1825 }
1826
1827 if weak_signal is None:
1828 weak_signal = junk_weak_signal
1829 outcome = "ok" if topics else "nothing-solid"
1830 plan = schema.DiscoveryPlan(
1831 domain=bundle.domain,
1832 category=None,
1833 subreddits=[],
1834 sources=(
1835 list(bundle.requested_sources)
1836 if bundle.requested_sources
1837 else sorted({
1838 item.source
1839 for entry in bundle.nominations
1840 for item in entry.nomination.items
1841 })
1842 ),
1843 )
1844 # The bundle's restored leg-1 sweep status (empty for pre-field bundles):
1845 # degraded sweep coverage must reach this report's status map and its
1846 # degraded-sources warning exactly as the one-shot reports it.
1847 source_status = dict(getattr(bundle, "source_status", None) or {})
1848 report = schema.DiscoveryReport(
1849 domain=bundle.domain,
1850 range_from=bundle.from_date,
1851 range_to=to_date,
1852 generated_at=datetime.now(timezone.utc).isoformat(),
1853 plan=plan,
1854 topics=topics,
1855 source_status=source_status,
1856 warnings=_discovery_report_warnings(topics, outcome, source_status),
1857 outcome=outcome,
1858 weak_signal=weak_signal[1] if weak_signal and not topics else None,
1859 )
1860 return DiscoverResumeResult(report=report, angle_inputs=angle_inputs)
1861
1862
1863 def diagnose(
1864 config: dict[str, Any],
1865 requested_sources: list[str] | None = None,
1866 *,
1867 safe: bool = False,
1868 x_envelope: bool = False,
1869 ) -> dict[str, Any]:
1870 # ``x_envelope`` is True when a validated --x-posts envelope is present for
1871 # this invocation, so available_sources lists x even without a backend
1872 # and the optional-source omission note does not fire.
1873 requested_sources = normalize_requested_sources(requested_sources)
1874 google_key = _google_key(config)
1875 x_status = env.get_x_source_status(config, probe=not safe)
1876 # Compute once and reuse for both the diag flag and available_sources below.
1877 # safe=True (doctor/--diagnose/--preflight) must stay network-free.
1878 x_pending = env.x_pending_browser_auth(config, local_only=safe)
1879 native_web_backend = None
1880 if config.get("BRAVE_API_KEY"):
1881 native_web_backend = "brave"
1882 elif config.get("EXA_API_KEY"):
1883 native_web_backend = "exa"
1884 elif config.get("SERPER_API_KEY"):
1885 native_web_backend = "serper"
1886 elif config.get("PARALLEL_API_KEY"):
1887 native_web_backend = "parallel"
1888 providers_status = {
1889 "google": bool(google_key),
1890 "openai": bool(config.get("OPENAI_API_KEY")) and config.get("OPENAI_AUTH_STATUS") == env.AUTH_STATUS_OK,
1891 "xai": bool(config.get("XAI_API_KEY")),
1892 "openrouter": bool(config.get("OPENROUTER_API_KEY")),
1893 "perplexity": bool(config.get("PERPLEXITY_API_KEY")),
1894 }
1895 reasoning_provider_available = any(
1896 providers_status[name] for name in ("google", "openai", "xai", "openrouter")
1897 )
1898 external_commands = {
1899 "yt-dlp": bool(which("yt-dlp")),
1900 "digg-pp-cli": bool(which("digg-pp-cli")),
1901 "arxiv-pp-cli": bool(which("arxiv-pp-cli")),
1902 "techmeme-pp-cli": bool(which("techmeme-pp-cli")),
1903 "trustpilot-pp-cli": bool(which("trustpilot-pp-cli")),
1904 "brightdata": bool(which("brightdata")),
1905 "gh": bool(which("gh")),
1906 }
1907 # Network-free two-field probe (bird_installed/bird_authenticated
1908 # precedent): "installed" is PATH resolution, "authenticated" is a
1909 # presence-only credential signal that never reads the secret.
1910 brightdata_status = brightdata.gate_status(config)
1911 credential_destinations = {
1912 "global_env": str(env.CONFIG_FILE) if env.CONFIG_FILE else None,
1913 }
1914 browser_cookies = {
1915 "mode": config.get("_BROWSER_COOKIE_MODE", "off"),
1916 "browsers": list(config.get("_BROWSER_COOKIE_BROWSERS") or []),
1917 "reads_values": False if safe else config.get("_BROWSER_COOKIE_MODE") == "read",
1918 }
1919 ignored_project_keys = list(config.get("_IGNORED_PROJECT_CONFIG_KEYS") or [])
1920 ignored_endpoint_overrides = [
1921 key for key in ignored_project_keys if key in permission_preflight.ENDPOINT_OVERRIDE_KEYS
1922 ]
1923 local_writes: list[dict[str, str]] = []
1924 if config.get("LAST30DAYS_MEMORY_DIR"):
1925 local_writes.append({"kind": "report", "path": str(config.get("LAST30DAYS_MEMORY_DIR"))})
1926 diag = {
1927 "providers": providers_status,
1928 "local_mode": not reasoning_provider_available,
1929 "reasoning_provider": (config.get("LAST30DAYS_REASONING_PROVIDER") or "auto").lower(),
1930 # The host-fetched connector lane serves X when no engine backend
1931 # exists and the model declared the lane.
1932 "x_backend": x_status["source"] or (
1933 "connector" if env.x_host_lane_declared(config) else None
1934 ),
1935 "bird_installed": x_status["bird_installed"],
1936 "bird_authenticated": x_status["bird_authenticated"],
1937 "bird_username": x_status["bird_username"],
1938 "x_pending_browser_auth": x_pending,
1939 "xquik_available": x_status.get("xquik_available", False),
1940 "xquik_working": x_status.get("xquik_working"),
1941 "xquik_status": x_status.get("xquik_status", ""),
1942 "native_web_backend": native_web_backend,
1943 "native_search": env.is_native_search(config),
1944 "has_scrapecreators": bool(config.get("SCRAPECREATORS_API_KEY")),
1945 "has_github": bool(config.get("GITHUB_TOKEN") or which("gh")),
1946 "brightdata_installed": brightdata_status["brightdata_installed"],
1947 "brightdata_authenticated": brightdata_status["brightdata_authenticated"],
1948 # safe=True (doctor/--diagnose/--preflight) must stay network-free:
1949 # answer X availability from local evidence only. x_pending is
1950 # precomputed by diagnose() to avoid double evaluation.
1951 "available_sources": available_sources(
1952 config, requested_sources, x_pending=x_pending, local_only=safe,
1953 x_envelope=x_envelope,
1954 ),
1955 "safe": safe,
1956 "config_source": config.get("_CONFIG_SOURCE"),
1957 "ignored_project_config": config.get("_IGNORED_PROJECT_CONFIG"),
1958 "ignored_project_config_keys": ignored_project_keys,
1959 "ignored_endpoint_overrides": ignored_endpoint_overrides,
1960 "browser_cookies": browser_cookies,
1961 "external_commands": external_commands,
1962 "credential_destinations": credential_destinations,
1963 "local_writes": local_writes,
1964 }
1965 diag["permission_preflight"] = permission_preflight.build(config, diag)
1966 return diag
1967
1968
1969 def _inner_max_workers(stream_count: int, *, internal_subrun: bool) -> int:
1970 """Worker-pool size for the per-stream fanout inside a single pipeline run.
1971
1972 Top-level runs use up to 16 workers. Subruns of ``run_competitor_fanout``
1973 cap the inner pool to 4 so a six-way competitor fan-out stays below
1974 roughly 30 worker threads in aggregate instead of ~96.
1975 """
1976 if internal_subrun:
1977 return max(2, min(4, stream_count or 1))
1978 return max(4, min(16, stream_count or 1))
1979
1980
1981 def _load_library_context(
1982 *,
1983 topic: str,
1984 config: dict[str, Any],
1985 mock: bool,
1986 internal_subrun: bool,
1987 x_handle: str | None,
1988 github_user: str | None,
1989 github_repos: list[str] | None,
1990 save_dir: Path | str | None = None,
1991 ) -> tuple[list[schema.LibraryContext], str | None]:
1992 """Resolve compact prior-run context without making a research run depend on it."""
1993 setting = str(config.get("LAST30DAYS_LIBRARY_CONTEXT") or "off").strip().lower()
1994 if mock or internal_subrun or setting in {"0", "false", "no", "off"}:
1995 return [], None
1996 if save_dir == "":
1997 return [], None
1998
1999 memory_dir = (
2000 save_dir
2001 if save_dir is not None
2002 else config.get("LAST30DAYS_MEMORY_DIR") or library.DEFAULT_MEMORY_DIR
2003 )
2004 briefs_dir = config.get("_LAST30DAYS_LIBRARY_BRIEFS_DIR") or (
2005 Path(memory_dir).expanduser() / "briefings"
2006 if save_dir is not None
2007 else library.DEFAULT_BRIEFS_DIR
2008 )
2009 db_path = config.get("_LAST30DAYS_LIBRARY_DB")
2010 if not db_path:
2011 db_path = (
2012 Path(memory_dir).expanduser().resolve() / ".last30days-library.db"
2013 if save_dir is not None
2014 else library_index.DEFAULT_LIBRARY_DB
2015 )
2016 store_db = config.get("_LAST30DAYS_STORE_DB")
2017 if not store_db:
2018 # Scoped runs read only a store inside the save dir (usually absent);
2019 # the shared store would leak other scopes' sightings into this one.
2020 store_db = (
2021 Path(memory_dir).expanduser().resolve() / "research.db"
2022 if save_dir is not None
2023 else library_index.DEFAULT_STORE_DB
2024 )
2025 queries = [topic, x_handle or "", github_user or "", *(github_repos or [])]
2026 queries = list(dict.fromkeys(value.strip() for value in queries if value and value.strip()))
2027 try:
2028 library_index.sync_library(memory_dir, briefs_dir, db_path=db_path)
2029 matches: list[library_index.LibrarySearchMatch] = []
2030 for query_text in queries:
2031 matches.extend(
2032 library_index.search(
2033 query_text,
2034 limit=6,
2035 db_path=db_path,
2036 store_db_path=store_db,
2037 )
2038 )
2039 except (library_index.LibrarySearchUnavailable, OSError, sqlite3.DatabaseError) as exc:
2040 return [], f"Library context unavailable: {exc}"
2041
2042 contexts: list[schema.LibraryContext] = []
2043 seen_runs: set[tuple[str, date]] = set()
2044 for match in sorted(
2045 matches,
2046 key=lambda item: (-item.published_date.toordinal(), item.rank, item.topic.casefold()),
2047 ):
2048 if match.run_key in seen_runs:
2049 continue
2050 seen_runs.add(match.run_key)
2051 contexts.append(
2052 schema.LibraryContext(
2053 topic=match.topic,
2054 published_date=match.published_date.isoformat(),
2055 headline=match.headline,
2056 summary=match.snippet or match.headline,
2057 source_kind=match.source_kind,
2058 )
2059 )
2060 if len(contexts) == 3:
2061 break
2062 return contexts, None
2063
2064
2065 def run(
2066 *,
2067 topic: str,
2068 config: dict[str, Any],
2069 depth: str,
2070 requested_sources: list[str] | None = None,
2071 mock: bool = False,
2072 x_handle: str | None = None,
2073 x_related: list[str] | None = None,
2074 web_backend: str = "auto",
2075 external_plan: dict | None = None,
2076 subreddits: list[str] | None = None,
2077 tiktok_hashtags: list[str] | None = None,
2078 tiktok_creators: list[str] | None = None,
2079 ig_creators: list[str] | None = None,
2080 lookback_days: int = 30,
2081 as_of_date: str | None = None,
2082 github_user: str | None = None,
2083 github_repos: list[str] | None = None,
2084 trustpilot_domain: str | None = None,
2085 trustpilot_domain_is_hint: bool = False,
2086 hiring_signals_mode: bool = False,
2087 internal_subrun: bool = False,
2088 suppress_x_host_lane: bool = False,
2089 save_dir: Path | str | None = None,
2090 corpus_dirs: list[str] | None = None,
2091 corpus_all_time: bool = False,
2092 x_posts: x_envelope.Envelope | None = None,
2093 ) -> schema.Report:
2094 # ``suppress_x_host_lane`` is distinct from ``internal_subrun``: comparison
2095 # entities share the latter and must still honor the connector lane;
2096 # only discovery enrichment passes set the former.
2097 # ``x_posts`` is a validated ``--x-posts`` envelope: when present
2098 # it replaces the engine's X fetch for this run and is served once.
2099 # Standalone runs (not competitor/discover sub-runs) own the YouTube
2100 # search-cache lifecycle. Comparison fan-out clears once before submit so
2101 # parallel entity sub-runs can still share in-run hits.
2102 if not internal_subrun:
2103 youtube_yt.reset_search_cache()
2104 settings = _resolve_depth_settings(depth, config)
2105 requested_sources = normalize_requested_sources(requested_sources)
2106 # Wall-clock origin for budget-aware enrichment lanes. Amazon review
2107 # enrichment starts at search time (inside _retrieve_stream_impl) so it
2108 # overlaps other sources instead of waiting for them all to finish.
2109 run_started = time.monotonic()
2110 from_date, to_date = dates.get_date_range(lookback_days, as_of_date=as_of_date)
2111 resolved_corpus_dirs = corpus.resolve_directories(
2112 corpus_dirs or config.get("_CORPUS_DIRS"),
2113 config.get("LAST30DAYS_CORPUS_DIRS"),
2114 )
2115 excluded_sources = {
2116 source.strip().lower()
2117 for source in str(config.get("EXCLUDE_SOURCES") or "").split(",")
2118 if source.strip()
2119 }
2120 corpus_enabled = bool(resolved_corpus_dirs) and "corpus" not in excluded_sources
2121 corpus_requested = bool(requested_sources and "corpus" in requested_sources)
2122 if corpus_enabled and requested_sources and "corpus" not in requested_sources:
2123 requested_sources = [*requested_sources, "corpus"]
2124
2125 # Host-fetched X lane. EXCLUDE_SOURCES=x or a --search list without
2126 # x wins: the envelope is ignored with a receipt line and stays unconsumed.
2127 envelope = x_posts
2128 if envelope is not None and (
2129 "x" in excluded_sources
2130 or (requested_sources and "x" not in requested_sources)
2131 ):
2132 log.source_log(
2133 "x", "host-fetched X: envelope ignored (x is excluded from this run)",
2134 tty_only=False,
2135 )
2136 envelope = None
2137 # The lane signal without an envelope is a broken handoff, not a reason to
2138 # spend a backup backend: X records the fixed not-passed outcome.
2139 x_lane_missing = (
2140 envelope is None
2141 and not mock
2142 and not suppress_x_host_lane
2143 and env.x_host_lane_declared(config)
2144 )
2145 if envelope is not None or x_lane_missing:
2146 # Ride the config dict (the _polymarket_keywords idiom) so the stream
2147 # workers and the handle-lane section see it without widening their
2148 # signatures. Copy first: comparison entities shallow-copy the shared
2149 # config and must never inherit another entity's envelope.
2150 config = dict(config)
2151 config["_x_envelope"] = envelope
2152 config["_x_lane_missing"] = x_lane_missing
2153
2154 # Gate StockTwits to ticker/crypto topics. Single chokepoint: when False,
2155 # available_sources() never registers stocktwits, so the planner can't
2156 # assign it (eligible_sources = available ∩ capabilities).
2157 config["_financial_topic"] = stocktwits.is_financial_topic(topic)
2158
2159 if mock:
2160 runtime = providers.mock_runtime(config, depth)
2161 reasoning_provider = None
2162 available = list(requested_sources or MOCK_AVAILABLE_SOURCES)
2163 if corpus_enabled and "corpus" not in available:
2164 available.append("corpus")
2165 if not corpus_enabled and not corpus_requested:
2166 available = [source for source in available if source != "corpus"]
2167 if not requested_sources and not hiring_signals_mode and not _company_topic_likely(topic):
2168 available = [source for source in available if source != "jobs"]
2169 else:
2170 runtime, reasoning_provider = providers.resolve_runtime(config, depth)
2171 available = available_sources(
2172 config, requested_sources,
2173 suppress_x_host_lane=suppress_x_host_lane,
2174 x_envelope=envelope is not None,
2175 )
2176 if requested_sources:
2177 available = [source for source in available if source in requested_sources]
2178 # Keep an explicitly requested but unconfigured corpus in the plan long
2179 # enough to record its skipped-unconfigured source outcome. It is never
2180 # submitted to the network executor below.
2181 if corpus_requested and "corpus" not in excluded_sources and "corpus" not in available:
2182 available.append("corpus")
2183 if web_backend == "none":
2184 available = [s for s in available if s != "grounding"]
2185 elif web_backend in ("brave", "exa", "serper", "parallel", "parallel-mcp", "keyless") and "grounding" not in available:
2186 available.append("grounding")
2187 if (
2188 hiring_signals_mode
2189 or (not requested_sources and _company_topic_likely(topic))
2190 ) and "jobs" not in available:
2191 available.append("jobs")
2192 if hiring_signals_mode:
2193 config = dict(config)
2194 config["_hiring_signals_mode"] = True
2195 if not requested_sources:
2196 available = ["jobs"]
2197 if not available:
2198 raise RuntimeError("No sources are available for this run.")
2199
2200 planner_requested_sources = requested_sources
2201 if hiring_signals_mode and not planner_requested_sources:
2202 planner_requested_sources = ["jobs"]
2203
2204 if external_plan is not None:
2205 # External plan provided (e.g., from Claude Code via --plan flag).
2206 # Explicit input is a contract: validate it before permissive sanitization.
2207 planner.validate_external_plan(external_plan)
2208 plan = planner._sanitize_plan(
2209 external_plan, topic, available, planner_requested_sources, depth,
2210 honor_plan_sources=True,
2211 )
2212 plan_source = "external"
2213 else:
2214 plan = planner.plan_query(
2215 topic=topic,
2216 available_sources=available,
2217 requested_sources=planner_requested_sources,
2218 depth=depth,
2219 provider=None if mock else reasoning_provider,
2220 model=None if mock else runtime.planner_model,
2221 context=config.get("_auto_resolve_context", ""),
2222 internal_subrun=internal_subrun,
2223 )
2224 # Source labelling: the fallback path annotates notes with "fallback-plan"
2225 # or "deterministic-comparison-plan"; anything else came from the LLM.
2226 if any("fallback" in note or "deterministic" in note for note in (plan.notes or [])):
2227 plan_source = "deterministic"
2228 elif not mock and reasoning_provider and runtime.planner_model:
2229 plan_source = "llm"
2230 else:
2231 plan_source = "deterministic"
2232
2233 # Safety net: ensure grounding appears in all subqueries even if the planner
2234 # omits it. This is redundant when the planner includes grounding via
2235 # SOURCE_CAPABILITIES, but kept as a fallback.
2236 if (
2237 web_backend != "none"
2238 and "grounding" in available
2239 and "drill-mode" not in plan.notes
2240 ):
2241 for sq in plan.subqueries:
2242 if "grounding" not in sq.sources:
2243 sq.sources.append("grounding")
2244 if "drill-mode" not in plan.notes:
2245 # Drill plans re-fetch only the sources that contributed to the matched
2246 # cluster; the company-topic jobs injection must not widen that set.
2247 _ensure_jobs_in_plan(plan, available, explicit=hiring_signals_mode, topic=topic)
2248 if "corpus" in available and plan.subqueries:
2249 # Corpus is deterministic and user-registered, so it always gets one
2250 # bounded stream even when a quick/LLM plan omits it. Reuse the primary
2251 # subquery instead of multiplying local scans across every subquery.
2252 if "corpus" not in plan.subqueries[0].sources:
2253 plan.subqueries[0].sources.append("corpus")
2254 if "corpus" not in plan.source_weights:
2255 plan.source_weights["corpus"] = 1.0
2256 plan.source_weights = planner._normalize_weights(plan.source_weights)
2257
2258 # Add the paid-only Perplexity lane after all normal-source safety nets.
2259 # This preserves the planner's primary subquery, gives the bounded paid
2260 # call the whole user topic, and prevents grounding, jobs, or corpus from
2261 # being attached to the dedicated lane.
2262 _ensure_perplexity_in_plan(
2263 plan,
2264 topic,
2265 available,
2266 force=bool(config.get("_deep_research")),
2267 )
2268
2269 # Always-on planner trace. Emits one summary line plus one per subquery
2270 # so retrieval-breadth failures like the 2026-04-19 Hermes Agent Use Cases
2271 # disaster are visible without --debug. Stderr only; does not leak into
2272 # the user-facing stdout synthesis.
2273 print(
2274 f"[Planner] Plan: intent={plan.intent}, freshness={plan.freshness_mode}, "
2275 f"cluster_mode={plan.cluster_mode}, subqueries={len(plan.subqueries)}, "
2276 f"source={plan_source}",
2277 file=sys.stderr,
2278 )
2279 if plan.subqueries:
2280 for index, sq in enumerate(plan.subqueries, start=1):
2281 sources_str = ",".join(sq.sources) if sq.sources else "(none)"
2282 print(
2283 f"[Planner] sq{index} label={sq.label} "
2284 f'search="{sq.search_query}" sources=[{sources_str}]',
2285 file=sys.stderr,
2286 )
2287 else:
2288 print("[Planner] (no subqueries in plan)", file=sys.stderr)
2289
2290 bundle = schema.RetrievalBundle(artifacts={"grounding": []})
2291 if envelope is not None:
2292 # The footer's X provenance reads "via X connector" (render._render_stats).
2293 bundle.artifacts["x_provenance"] = "connector"
2294 # Handles the user named explicitly. Available before any retrieval, unlike
2295 # the entity-extracted set, so Phase 1 and quick-depth runs get first-party
2296 # protection too. Without this the exemption reached only the Phase 2
2297 # supplement path -- which quick runs skip entirely -- so a subject-authored
2298 # post retrieved in Phase 1 was still pruned before fusion, which is exactly
2299 # the evidence loss this change exists to prevent.
2300 explicit_first_party = {
2301 h.lstrip("@").strip().lower()
2302 for h in ([x_handle, github_user, *(x_related or [])])
2303 if h and h.strip()
2304 }
2305 # Creator accounts named via --ig-creators / --creators carry the same
2306 # explicit intent: the run is searching those accounts, and a creator's
2307 # caption rarely repeats the topic's literal tokens, so without an
2308 # exemption the relevance floor prunes them as third-party noise (issue
2309 # #1101: 36 creator reels fetched, 0 reported). The exemption is scoped
2310 # per platform, NOT merged into the global set: each flag names accounts
2311 # on one platform, and an unrelated same-name account elsewhere must not
2312 # bypass the floors.
2313 creator_first_party = _creator_first_party_by_source(tiktok_creators, ig_creators)
2314 # Real X handles: --x-handle, --x-related, or @mentions in the topic. These
2315 # determine whether the deferred X floor applies. Topic words like "peter"
2316 # are NOT real handles and should not trigger the floor — when no real
2317 # handle is identified, the floor is skipped entirely (policy: a noisier
2318 # report beats losing the subject's evidence).
2319 explicit_x_handles = {
2320 h.lstrip("@").strip().lower()
2321 for h in ([x_handle, *(x_related or [])])
2322 if h and h.strip()
2323 } | _topic_handle_mentions(topic)
2324 # Plus handle-shaped tokens from the topic. Phase 1 and quick-depth runs
2325 # never reach automatic handle resolution, so without this a quick search
2326 # naming a subject still discards everything that subject wrote.
2327 explicit_first_party |= _topic_first_party_candidates(topic)
2328
2329 for source in (requested_sources or []):
2330 if source not in available:
2331 bundle.record_failure(
2332 source,
2333 schema.SKIPPED_UNCONFIGURED,
2334 "Source was requested but is not configured for this run.",
2335 attempted=False,
2336 )
2337 if corpus_requested and not corpus_enabled:
2338 bundle.record_failure(
2339 "corpus",
2340 schema.SKIPPED_UNCONFIGURED,
2341 "Corpus was requested but no readable directory was configured.",
2342 attempted=False,
2343 )
2344 # Expose plan_source to the renderer so render_compact can emit the
2345 # DEGRADED RUN banner when a named-entity topic was invoked bare
2346 # (source=deterministic AND no pre-research flags). LAW 7 backstop.
2347 bundle.artifacts["plan_source"] = plan_source
2348 bundle.artifacts["corpus_in_export"] = bool(config.get("_CORPUS_IN_EXPORT"))
2349 # Hiring-signals is deliberately jobs-only with no multi-source --plan, so
2350 # the LAW 7 degraded-run and Step 0.55 pre-research banners do not apply -
2351 # they would contradict the documented jobs-scoped flow. Suppress them.
2352 bundle.artifacts["hiring_signals_mode"] = hiring_signals_mode
2353 # Record the resolved Amazon keyword whenever the lane is active, so the
2354 # footer can name it on an empty result. A search that matched nothing
2355 # still spent a credit, and the fix is almost always the keyword -- a
2356 # suppressed line means nobody ever learns it was wrong.
2357 if "amazon" in (available or []):
2358 bundle.artifacts["amazon_query"] = (
2359 str(config.get("_amazon_query") or "").strip() or topic
2360 )
2361
2362 # Project-mode or person-mode GitHub: run once before the main subquery loop
2363 _github_custom_done = False
2364 _github_enriched_repos: set[str] = set()
2365
2366 # Project mode takes priority over person mode
2367 if github_repos and "github" in available:
2368 bundle.mark_attempted("github")
2369 try:
2370 project_items = github.search_github_project(
2371 github_repos, from_date, to_date,
2372 depth=depth, token=config.get("GITHUB_TOKEN"),
2373 )
2374 if project_items:
2375 normalized = _normalize_score_dedupe(
2376 "github", project_items, from_date, to_date,
2377 freshness_mode=plan.freshness_mode,
2378 ranking_query=f"What are {', '.join(github_repos)} doing on GitHub?",
2379 )
2380 primary_label = plan.subqueries[0].label if plan.subqueries else "primary"
2381 bundle.add_items(primary_label, "github", normalized)
2382 _github_custom_done = True
2383 _github_enriched_repos = {r.lower() for r in github_repos}
2384 except Exception as exc:
2385 bundle.errors_by_source["github"] = f"Project-mode failed: {exc}"
2386 state, attempted = _classify_source_failure(exc)
2387 bundle.record_failure("github", state, str(exc), attempted=attempted)
2388
2389 _github_person_done = False
2390 if github_user and "github" in available and not _github_custom_done:
2391 bundle.mark_attempted("github")
2392 _github_person_done = True
2393 try:
2394 person_items = github.search_github_person(
2395 github_user, from_date, to_date,
2396 depth=depth, token=config.get("GITHUB_TOKEN"),
2397 )
2398 if person_items:
2399 normalized = _normalize_score_dedupe(
2400 "github", person_items, from_date, to_date,
2401 freshness_mode=plan.freshness_mode,
2402 ranking_query=f"What is @{github_user} doing on GitHub?",
2403 )
2404 # Use the first subquery's label so RRF can look up the weight
2405 primary_label = plan.subqueries[0].label if plan.subqueries else "primary"
2406 bundle.add_items(primary_label, "github", normalized)
2407 else:
2408 # A pinned --github-user that yields nothing must not be
2409 # silently backfilled by generic keyword search: the report
2410 # would then present unrelated repos as this person's work.
2411 bundle.record_failure(
2412 "github",
2413 "no-results",
2414 f"Person mode found no activity for @{github_user} in the window",
2415 )
2416 except Exception as exc:
2417 bundle.errors_by_source["github"] = f"Person-mode failed: {exc}"
2418 state, attempted = _classify_source_failure(exc)
2419 bundle.record_failure("github", state, str(exc), attempted=attempted)
2420
2421 # Trustpilot session warm-up happens inside search_trustpilot at the
2422 # first (capped, single) fetch -- lazily, so it never delays the other
2423 # sources' streams and never fires for runs whose plan fetches no
2424 # Trustpilot. The module-level lock in lib/trustpilot.py serializes
2425 # concurrent vs-mode sub-runs so they never race Chrome harvests.
2426
2427 # Thread-safe set prevents redundant fetches after a source returns 429
2428 rate_limited_sources: set[str] = set()
2429 rate_limit_lock = threading.Lock()
2430
2431 # Local corpus retrieval is intentionally outside the network executor and
2432 # retry budget. One bounded stream participates in the same signal scoring,
2433 # fusion, reranking, and per-source result cap as remote sources.
2434 if corpus_enabled and plan.subqueries:
2435 primary = plan.subqueries[0]
2436 bundle.mark_attempted("corpus")
2437 result = corpus.search(
2438 topic,
2439 resolved_corpus_dirs,
2440 from_date=from_date,
2441 to_date=to_date,
2442 all_time=corpus_all_time,
2443 limit=settings["per_stream_limit"],
2444 cache_dir=env.CONFIG_DIR,
2445 )
2446 prepared_query = relevance.PreparedQuery(primary.ranking_query)
2447 lookback_window_days = (
2448 datetime.strptime(to_date, "%Y-%m-%d").date()
2449 - datetime.strptime(from_date, "%Y-%m-%d").date()
2450 ).days
2451 corpus_items = signals.annotate_stream(
2452 result.items,
2453 prepared_query,
2454 plan.freshness_mode,
2455 reference_date=to_date,
2456 max_days=lookback_window_days,
2457 )
2458 corpus_items = signals.prune_low_relevance(corpus_items)
2459 corpus_items = dedupe.dedupe_items(corpus_items)
2460 for item in corpus_items:
2461 item.snippet = snippet.extract_best_snippet(item, prepared_query)
2462 bundle.add_items(primary.label, "corpus", corpus_items)
2463 if result.notes:
2464 outcome = bundle.source_status["corpus"]
2465 bundle.source_status["corpus"] = schema.SourceOutcome(
2466 source="corpus",
2467 state=outcome.state,
2468 items_returned=outcome.items_returned,
2469 attempted=True,
2470 detail="; ".join(result.notes),
2471 )
2472 bundle.artifacts["corpus"] = {
2473 "files_scanned": result.files_scanned,
2474 "cache_hits": result.cache_hits,
2475 "all_time": corpus_all_time,
2476 }
2477
2478 futures = {}
2479 # Per-source fetch budget prevents redundant API calls
2480 source_fetch_count: dict[str, int] = {}
2481 stream_count = sum(
2482 1
2483 for subquery in plan.subqueries
2484 for source in subquery.sources
2485 if source in available and source != "corpus"
2486 )
2487 max_workers = _inner_max_workers(stream_count, internal_subrun=internal_subrun)
2488 with ThreadPoolExecutor(max_workers=max_workers) as executor:
2489 for subquery in plan.subqueries:
2490 for source in subquery.sources:
2491 if source not in available:
2492 continue
2493 if source == "corpus":
2494 continue
2495 # Skip GitHub keyword search if person-mode already ran
2496 if source == "github" and (_github_person_done or _github_custom_done):
2497 continue
2498 # Enforce per-source fetch cap. A CLI override (issue #716) raises
2499 # the cap for capped sources so every X subquery in a multi-angle
2500 # --plan fetches, instead of only the first two.
2501 cap = _source_fetch_cap(source, config)
2502 if cap is not None:
2503 if cap <= 0:
2504 continue
2505 current = source_fetch_count.get(source, 0)
2506 if current >= cap:
2507 continue
2508 shared_paid_budget = config.get("_perplexity_paid_budget")
2509 if (
2510 source == "perplexity"
2511 and isinstance(shared_paid_budget, PaidSourceBudget)
2512 and not shared_paid_budget.try_consume(
2513 cap,
2514 claimant=topic,
2515 )
2516 ):
2517 bundle.artifacts.setdefault("paid_source_budget", {})[
2518 "perplexity"
2519 ] = {
2520 "state": "skipped-budget",
2521 "attempted": False,
2522 "owner": shared_paid_budget.owner,
2523 "claimant": topic,
2524 }
2525 continue
2526 source_fetch_count[source] = current + 1
2527 bundle.mark_attempted(source)
2528 futures[
2529 executor.submit(
2530 _retrieve_stream,
2531 topic=topic,
2532 subquery=subquery,
2533 source=source,
2534 config=config,
2535 depth=depth,
2536 date_range=(from_date, to_date),
2537 runtime=runtime,
2538 mock=mock,
2539 rate_limited_sources=rate_limited_sources,
2540 rate_limit_lock=rate_limit_lock,
2541 web_backend=web_backend,
2542 raw_topic=topic,
2543 subreddits=subreddits,
2544 tiktok_hashtags=tiktok_hashtags,
2545 tiktok_creators=tiktok_creators,
2546 ig_creators=ig_creators,
2547 trustpilot_domain=trustpilot_domain,
2548 trustpilot_domain_is_hint=trustpilot_domain_is_hint,
2549 run_started=run_started,
2550 )
2551 ] = (subquery, source)
2552
2553 for future in as_completed(futures):
2554 subquery, source = futures[future]
2555 try:
2556 raw_items, artifact = future.result()
2557 except Exception as exc:
2558 # Share 429 signal so pending futures skip this source
2559 if _is_rate_limit_error(exc):
2560 with rate_limit_lock:
2561 rate_limited_sources.add(source)
2562 bundle.errors_by_source[source] = str(exc)
2563 state, attempted = _classify_source_failure(exc)
2564 bundle.record_failure(source, state, str(exc), attempted=attempted)
2565 continue
2566 # Retry once for transient 5xx errors
2567 if _is_transient_error(exc):
2568 time.sleep(3)
2569 try:
2570 raw_items, artifact = _retrieve_stream(
2571 topic=topic, subquery=subquery, source=source,
2572 config=config, depth=depth, date_range=(from_date, to_date),
2573 runtime=runtime, mock=mock,
2574 rate_limited_sources=rate_limited_sources,
2575 rate_limit_lock=rate_limit_lock,
2576 web_backend=web_backend,
2577 raw_topic=topic,
2578 subreddits=subreddits,
2579 tiktok_hashtags=tiktok_hashtags,
2580 tiktok_creators=tiktok_creators,
2581 ig_creators=ig_creators,
2582 trustpilot_domain=trustpilot_domain,
2583 trustpilot_domain_is_hint=trustpilot_domain_is_hint,
2584 run_started=run_started,
2585 )
2586 except Exception as retry_exc:
2587 detail = f"{exc} (retried once, still failed: {retry_exc})"
2588 bundle.errors_by_source[source] = detail
2589 state, attempted = _classify_source_failure(retry_exc)
2590 bundle.record_failure(source, state, detail, attempted=attempted)
2591 continue
2592 else:
2593 bundle.errors_by_source[source] = str(exc)
2594 state, attempted = _classify_source_failure(exc)
2595 bundle.record_failure(source, state, str(exc), attempted=attempted)
2596 continue
2597 outcome_note = None
2598 if isinstance(artifact, dict) and artifact.get("_source_outcome"):
2599 artifact = dict(artifact)
2600 outcome_note = artifact.pop("_source_outcome")
2601 bundle.record_failure(
2602 source,
2603 outcome_note["state"],
2604 outcome_note["detail"],
2605 attempted=outcome_note.get("attempted", True),
2606 )
2607 if isinstance(artifact, dict) and artifact.get("_source_outcome_detail"):
2608 artifact = dict(artifact)
2609 lane_state = artifact.pop("_source_outcome_detail_state", None)
2610 bundle.record_detail(
2611 source, artifact.pop("_source_outcome_detail"), state=lane_state
2612 )
2613 if lane_state == health.RATE_LIMITED:
2614 # Do not re-fan-out against a host still inside its window.
2615 with rate_limit_lock:
2616 rate_limited_sources.add(source)
2617 normalized = _normalize_score_dedupe(
2618 source, raw_items, from_date, to_date,
2619 freshness_mode=plan.freshness_mode,
2620 ranking_query=subquery.ranking_query,
2621 first_party_handles=explicit_first_party,
2622 first_party_by_source=creator_first_party,
2623 # X defers its relevance floor until resolved_handles exists.
2624 # Everything else prunes here as before.
2625 defer_relevance_prune=(source == "x"),
2626 )
2627 # Jobs is exempt from per_stream_limit: a careers board is a complete
2628 # snapshot of open roles, and truncating it to the default 12 drops
2629 # strategic postings (the whole point of hiring-signals coverage).
2630 if source != "jobs":
2631 normalized = _apply_reddit_stream_keepers(
2632 source, normalized, settings["per_stream_limit"], topic
2633 )
2634 bundle.add_items(subquery.label, source, normalized)
2635 if artifact:
2636 bundle.artifacts.setdefault("grounding", []).append(artifact)
2637
2638 # Phase 2: supplemental entity-based searches
2639 supplemental_handles: list[str] = []
2640 _run_supplemental_searches(
2641 topic=topic,
2642 bundle=bundle,
2643 plan=plan,
2644 config=config,
2645 depth=depth,
2646 date_range=(from_date, to_date),
2647 runtime=runtime,
2648 mock=mock,
2649 rate_limited_sources=rate_limited_sources,
2650 rate_limit_lock=rate_limit_lock,
2651 x_handle=x_handle,
2652 x_related=x_related,
2653 resolved_handles_out=supplemental_handles,
2654 )
2655
2656 # Phase 2b: retry thin sources with simplified query
2657 # Note: _github_skip_sources tells the retry to not re-run GitHub keyword search
2658 # when project-mode or person-mode already provided authoritative data.
2659 _github_skip_retry = {"corpus"}
2660 if _github_person_done or _github_custom_done:
2661 _github_skip_retry.add("github")
2662 _retry_thin_sources(
2663 topic=topic,
2664 bundle=bundle,
2665 plan=plan,
2666 config=config,
2667 depth=depth,
2668 date_range=(from_date, to_date),
2669 runtime=runtime,
2670 mock=mock,
2671 rate_limited_sources=rate_limited_sources,
2672 rate_limit_lock=rate_limit_lock,
2673 settings=settings,
2674 web_backend=web_backend,
2675 skip_sources=_github_skip_retry,
2676 subreddits=subreddits,
2677 tiktok_hashtags=tiktok_hashtags,
2678 tiktok_creators=tiktok_creators,
2679 ig_creators=ig_creators,
2680 first_party_handles=explicit_first_party,
2681 first_party_by_source=creator_first_party,
2682 run_started=run_started,
2683 )
2684
2685 # Reclassify partial failures as DEGRADED instead of silently dropping them.
2686 # A source that 429'd on one subquery but succeeded on another is not a hard
2687 # failure, but it is not healthy either: it likely returned fewer results
2688 # than it should have. Move it out of errors_by_source (so it isn't reported
2689 # as "failed") and into degraded_by_source (so it survives into warnings),
2690 # rather than deleting the signal outright as the engine used to.
2691 degraded_by_source: dict[str, str] = {}
2692 for source in list(bundle.errors_by_source):
2693 if bundle.items_by_source.get(source):
2694 degraded_by_source[source] = bundle.errors_by_source[source]
2695 del bundle.errors_by_source[source]
2696
2697 hiring_summary = _apply_hiring_signal_gate(
2698 bundle,
2699 explicit=hiring_signals_mode,
2700 topic=topic,
2701 )
2702 if hiring_summary:
2703 bundle.artifacts["hiring_signals"] = hiring_summary
2704
2705 items_by_source = _finalize_items_by_source(
2706 bundle.items_by_source, topic=topic, config=config, depth=depth, mock=mock,
2707 elapsed=time.monotonic() - run_started,
2708 )
2709 source_status = _finalize_source_status(bundle.source_status, items_by_source)
2710 # Normalized set of handles this run resolved for the topic. A candidate
2711 # authored by one of these is first-party and is exempted from the
2712 # entity-miss demotion in rerank (a post never repeats its own author's
2713 # name, so the body-text grounding check would otherwise zero out the
2714 # subject's own highest-signal posts). Built before fusion so the
2715 # per-author cap can give the topic's subject a higher allowance than an
2716 # incidental third-party account.
2717 resolved_handles = explicit_first_party | {
2718 h.lstrip("@").strip().lower()
2719 for h in supplemental_handles
2720 if h and h.strip()
2721 } | {h for handles in creator_first_party.values() for h in handles}
2722 # resolved_handles feeds rerank/fusion, where the first-party marks only
2723 # resist demotion of items that already passed the inclusion floors and a
2724 # named account is treated as the subject across surfaces. The inclusion
2725 # gate (prune_low_relevance) instead gets the platform-scoped
2726 # creator_first_party map so a cross-platform name collision cannot
2727 # bypass the floors.
2728 # Real X handles from explicit flags, @mentions in topic, or Phase 2 discovery.
2729 # When no real handle is identified, skip the X floor entirely — a noisier
2730 # report beats losing the subject's evidence. Topic tokens like "peter" are
2731 # NOT real handles: they populate resolved_handles for downstream first-party
2732 # protection but should NOT trigger the floor.
2733 real_x_handles = explicit_x_handles | {
2734 h.lstrip("@").strip().lower()
2735 for h in supplemental_handles
2736 if h and h.strip()
2737 }
2738 # Deferred X relevance floor. Phase 1 skipped it so this could run with the
2739 # run's actual resolved handles rather than a guess made before anyone knew
2740 # who the subject was. Applied per subquery stream so fusion sees the same
2741 # shape it always has. Only applied when we have real X handles — topic
2742 # tokens alone cannot identify the subject.
2743 if real_x_handles:
2744 # IG/TikTok creator exemptions stay on their own platforms: a creator
2745 # handle must not exempt a same-name X account from this floor. Only
2746 # creator-ONLY handles are subtracted, where X provenance means
2747 # real_x_handles (explicit X flags, @mentions in the topic, Phase 2
2748 # discovery) - a plain topic token or --github-user match is NOT X
2749 # provenance and does not preserve the exemption.
2750 x_floor_handles = resolved_handles - _creator_only_handles(
2751 creator_first_party, real_x_handles
2752 )
2753 for key, stream in list(bundle.items_by_source_and_query.items()):
2754 if key[1] != "x" or not stream:
2755 continue
2756 pruned = signals.prune_low_relevance(
2757 stream,
2758 first_party_handles=x_floor_handles,
2759 first_party_by_source=creator_first_party,
2760 )
2761 _log_prune_drop("x", len(stream), len(pruned), scope="per-query stream")
2762 bundle.items_by_source_and_query[key] = pruned
2763 if bundle.items_by_source.get("x"):
2764 x_stream = bundle.items_by_source["x"]
2765 pruned = signals.prune_low_relevance(
2766 x_stream,
2767 first_party_handles=x_floor_handles,
2768 first_party_by_source=creator_first_party,
2769 )
2770 _log_prune_drop("x", len(x_stream), len(pruned), scope="merged stream")
2771 bundle.items_by_source["x"] = pruned
2772
2773 candidates = weighted_rrf(
2774 bundle.items_by_source_and_query,
2775 plan,
2776 pool_limit=settings["pool_limit"],
2777 range_from=from_date,
2778 range_to=to_date,
2779 first_party_handles=resolved_handles,
2780 )
2781 private_candidates = [
2782 candidate
2783 for candidate in candidates
2784 if candidate.source == "corpus"
2785 or any(item.source == "corpus" for item in candidate.source_items)
2786 ]
2787 private_candidate_ids = {id(candidate) for candidate in private_candidates}
2788 public_candidates = [
2789 candidate for candidate in candidates if id(candidate) not in private_candidate_ids
2790 ]
2791 ranked_public = rerank.rerank_candidates(
2792 topic=topic,
2793 plan=plan,
2794 candidates=public_candidates,
2795 provider=None if mock else reasoning_provider,
2796 model=None if mock else runtime.rerank_model,
2797 shortlist_size=settings["rerank_limit"],
2798 resolved_handles=resolved_handles,
2799 )
2800 # Corpus titles/snippets must never enter a hosted reasoning prompt. Score
2801 # every candidate carrying corpus evidence with the deterministic fallback,
2802 # even when the rest of the run uses a remote reranker.
2803 ranked_private = rerank.rerank_candidates(
2804 topic=topic,
2805 plan=plan,
2806 candidates=private_candidates,
2807 provider=None,
2808 model=None,
2809 shortlist_size=settings["rerank_limit"],
2810 resolved_handles=resolved_handles,
2811 )
2812 ranked_public = rerank.prune_fallback_entity_misses(ranked_public, topic=topic)
2813 if hiring_summary:
2814 # The diagnostic source dump retains pruned jobs; hiring aggregation
2815 # needs their rejection identities after they leave the ranked pool.
2816 retained_ids = {candidate.candidate_id for candidate in ranked_public}
2817 hiring_summary["rejected_job_keys"] = sorted({
2818 fusion.candidate_key(item)
2819 for candidate in public_candidates
2820 if candidate.candidate_id not in retained_ids
2821 for item in candidate.source_items
2822 if item.source == "jobs"
2823 })
2824 # Private corpus already cleared a body-aware retrieval floor; do not apply
2825 # the public title/snippet visibility gate (filenames often omit the head
2826 # token even when the document body matched).
2827 ranked_candidates = sorted(
2828 [*ranked_public, *ranked_private],
2829 key=lambda candidate: (
2830 1 if schema.candidate_out_of_window(candidate) else 0,
2831 -candidate.final_score,
2832 -(candidate.engagement or -1),
2833 min(candidate.native_ranks.values(), default=999),
2834 candidate.title,
2835 ),
2836 )
2837 rerank.score_fun(
2838 topic=topic,
2839 candidates=ranked_public,
2840 provider=None if mock else reasoning_provider,
2841 model=None if mock else runtime.rerank_model,
2842 )
2843 rerank.score_fun(
2844 topic=topic,
2845 candidates=ranked_private,
2846 provider=None,
2847 model=None,
2848 )
2849
2850 # Phase 3: post-rerank GitHub star enrichment. Record/replay-aware so the
2851 # eval harness stays fully offline: this path calls the GitHub API (and the
2852 # gh-credential fallback) outside the _retrieve_stream seam, so it gets its
2853 # own fixture exchange keyed by phase.
2854 if "github" in available and not mock:
2855 star_request = {
2856 "source": "github",
2857 "phase": "post_rerank_star_enrichment",
2858 "topic": topic,
2859 "depth": depth,
2860 }
2861 star_matched, star_replayed = http.fixture_source_replay(star_request)
2862 if star_matched:
2863 star_map = star_replayed if isinstance(star_replayed, dict) else {}
2864 github.apply_star_map(ranked_candidates, star_map)
2865 else:
2866 collected_star_map: dict[str, int] = {}
2867 github.enrich_candidates_with_stars(
2868 ranked_candidates,
2869 token=config.get("GITHUB_TOKEN"),
2870 already_enriched=_github_enriched_repos,
2871 collect_map=collected_star_map,
2872 )
2873 http.fixture_source_record(star_request, collected_star_map)
2874
2875 clusters = cluster_candidates(ranked_candidates, plan)
2876 warnings = _warnings(items_by_source, ranked_candidates, bundle.errors_by_source, degraded_by_source)
2877 # One-sided entity coverage is a reporting warning, not a source failure:
2878 # marking the source PARTIAL would trip LAST30DAYS_STRICT_EXIT on runs that
2879 # returned good X results.
2880 warnings.extend(bundle.artifacts.get("x_partial_coverage", []))
2881 # Backend receipts that are not failures (xapi's truncated window), and
2882 # the Meta Ads footer inputs. A stream artifact only ever reaches the
2883 # report as an anonymous entry in this list, so the advertiser and the
2884 # pre-truncation counts have to be lifted to named top-level artifacts or
2885 # the footer cannot render them -- least of all on a zero-item run, which
2886 # is exactly when naming the advertiser matters most.
2887 for stream_artifact in bundle.artifacts.get("grounding", []):
2888 if isinstance(stream_artifact, dict):
2889 warnings.extend(stream_artifact.get("x_receipts", []))
2890 _lift_stream_artifacts(bundle)
2891 library_context, library_warning = _load_library_context(
2892 topic=topic,
2893 config=config,
2894 mock=mock,
2895 internal_subrun=internal_subrun,
2896 x_handle=x_handle,
2897 github_user=github_user,
2898 github_repos=github_repos,
2899 save_dir=save_dir,
2900 )
2901 if library_warning:
2902 warnings.append(library_warning)
2903
2904 return schema.Report(
2905 topic=topic,
2906 range_from=from_date,
2907 range_to=to_date,
2908 generated_at=datetime.now(timezone.utc).isoformat(),
2909 provider_runtime=runtime,
2910 query_plan=plan,
2911 clusters=clusters,
2912 ranked_candidates=ranked_candidates,
2913 items_by_source=items_by_source,
2914 errors_by_source=bundle.errors_by_source,
2915 source_status=source_status,
2916 warnings=warnings,
2917 artifacts=bundle.artifacts,
2918 library_context=library_context,
2919 )
2920
2921
2922 def _candidate_is_duplicate(
2923 candidate: schema.Candidate,
2924 kept: list[schema.Candidate],
2925 ) -> bool:
2926 if any(existing.candidate_id == candidate.candidate_id for existing in kept):
2927 return True
2928 if candidate.url and any(existing.url == candidate.url for existing in kept):
2929 return True
2930 candidate_text = " ".join((candidate.title, candidate.snippet)).strip()
2931 return bool(candidate_text) and any(
2932 dedupe.hybrid_similarity(
2933 candidate_text,
2934 " ".join((existing.title, existing.snippet)).strip(),
2935 ) >= 0.7
2936 for existing in kept
2937 )
2938
2939
2940 def merge_drill_report(
2941 report: schema.Report,
2942 drill_report: schema.Report,
2943 matched_clusters: list[schema.Cluster],
2944 *,
2945 target: str,
2946 ) -> schema.Report:
2947 """Merge a narrow follow-up into its cached report while preserving other clusters."""
2948 merged = copy.deepcopy(report)
2949 selected_cluster_ids = {cluster.cluster_id for cluster in matched_clusters}
2950 selected_candidate_ids = {
2951 candidate_id
2952 for cluster in matched_clusters
2953 for candidate_id in cluster.candidate_ids
2954 }
2955 original_candidates = {
2956 candidate.candidate_id: candidate for candidate in merged.ranked_candidates
2957 }
2958 unrelated_candidates = [
2959 candidate for candidate in merged.ranked_candidates
2960 if candidate.candidate_id not in selected_candidate_ids
2961 ]
2962 original_summary = ""
2963 for cluster in matched_clusters:
2964 for candidate_id in cluster.representative_ids:
2965 candidate = original_candidates.get(candidate_id)
2966 if candidate:
2967 original_summary = candidate.snippet or candidate.explanation or candidate.title
2968 if original_summary:
2969 break
2970 if original_summary:
2971 break
2972
2973 unrelated_candidate_indexes = {
2974 candidate.candidate_id: index
2975 for index, candidate in enumerate(unrelated_candidates)
2976 }
2977 focused_candidates: list[schema.Candidate] = []
2978 for candidate in [
2979 *copy.deepcopy(drill_report.ranked_candidates),
2980 *[
2981 copy.deepcopy(candidate)
2982 for candidate in merged.ranked_candidates
2983 if candidate.candidate_id in selected_candidate_ids
2984 ],
2985 ]:
2986 unrelated_index = unrelated_candidate_indexes.get(candidate.candidate_id)
2987 if unrelated_index is not None:
2988 candidate.cluster_id = unrelated_candidates[unrelated_index].cluster_id
2989 unrelated_candidates[unrelated_index] = candidate
2990 continue
2991 if not _candidate_is_duplicate(candidate, focused_candidates):
2992 focused_candidates.append(candidate)
2993
2994 primary_cluster = matched_clusters[0]
2995 for candidate in focused_candidates:
2996 candidate.cluster_id = primary_cluster.cluster_id
2997 focused_ids = [candidate.candidate_id for candidate in focused_candidates]
2998 focused_sources = sorted({
2999 source
3000 for candidate in focused_candidates
3001 for source in schema.candidate_sources(candidate)
3002 })
3003 replacement_cluster = schema.Cluster(
3004 cluster_id=primary_cluster.cluster_id,
3005 title=primary_cluster.title,
3006 candidate_ids=focused_ids,
3007 representative_ids=focused_ids[:3],
3008 sources=focused_sources,
3009 score=max((candidate.final_score for candidate in focused_candidates), default=0.0),
3010 uncertainty="single-source" if len(focused_sources) == 1 else None,
3011 )
3012
3013 first_selected_index = min(
3014 index
3015 for index, cluster in enumerate(merged.clusters)
3016 if cluster.cluster_id in selected_cluster_ids
3017 )
3018 remaining_clusters = [
3019 cluster for cluster in merged.clusters
3020 if cluster.cluster_id not in selected_cluster_ids
3021 ]
3022 remaining_clusters.insert(first_selected_index, replacement_cluster)
3023 merged.clusters = remaining_clusters
3024
3025 merged.ranked_candidates = focused_candidates + unrelated_candidates
3026
3027 all_sources = set(merged.items_by_source) | set(drill_report.items_by_source)
3028 new_item_count = 0
3029 merged_items: dict[str, list[schema.SourceItem]] = {}
3030 for source in sorted(all_sources):
3031 old_items = merged.items_by_source.get(source, [])
3032 new_items = drill_report.items_by_source.get(source, [])
3033 # Collapse exact URL matches first, preferring the drill's copy (it
3034 # carries fresh transcripts/comments); fuzzy dedupe alone keeps both
3035 # when enrichment changed the text substantially.
3036 new_urls = {item.url for item in new_items if item.url}
3037 kept_old = [item for item in old_items if not (item.url and item.url in new_urls)]
3038 combined = dedupe.dedupe_items([*copy.deepcopy(new_items), *kept_old])
3039 old_unique = dedupe.dedupe_items(old_items)
3040 new_item_count += max(0, len(combined) - len(old_unique))
3041 merged_items[source] = combined
3042 merged.items_by_source = merged_items
3043
3044 merged.generated_at = drill_report.generated_at
3045 merged.query_plan = drill_report.query_plan
3046 # The drill's retrieval window is the report's window now (a --days/--as-of
3047 # override on the drill must not be mislabeled with the cached range).
3048 merged.range_from = drill_report.range_from
3049 merged.range_to = drill_report.range_to
3050 attempted_sources = {
3051 source
3052 for source, outcome in drill_report.source_status.items()
3053 if outcome.attempted or outcome.state == schema.SKIPPED_UNCONFIGURED
3054 }
3055 for source in attempted_sources:
3056 if source in drill_report.errors_by_source:
3057 merged.errors_by_source[source] = drill_report.errors_by_source[source]
3058 else:
3059 merged.errors_by_source.pop(source, None)
3060 merged.source_status[source] = drill_report.source_status[source]
3061 merged.source_status = _finalize_source_status(
3062 merged.source_status,
3063 merged.items_by_source,
3064 )
3065 degraded_by_source = {
3066 source: outcome.detail or "partial results"
3067 for source, outcome in merged.source_status.items()
3068 if outcome.state == schema.PARTIAL
3069 }
3070 merged.warnings = _warnings(
3071 merged.items_by_source,
3072 merged.ranked_candidates,
3073 merged.errors_by_source,
3074 degraded_by_source,
3075 )
3076 merged.artifacts.update(copy.deepcopy(drill_report.artifacts))
3077 history = list(merged.artifacts.get("drill_history") or [])
3078 history.append({
3079 "target": target,
3080 "clusters": [cluster.title for cluster in matched_clusters],
3081 "new_items": new_item_count,
3082 "generated_at": drill_report.generated_at,
3083 })
3084 merged.artifacts["drill_history"] = history
3085 merged.artifacts["drill_context"] = {
3086 "target": target,
3087 "cluster_titles": [cluster.title for cluster in matched_clusters],
3088 "original_summary": original_summary,
3089 "new_items": new_item_count,
3090 "sources": focused_sources,
3091 }
3092 merged.drill_of = primary_cluster.title
3093 return merged
3094
3095
3096 def _batch_subject_handles(raw_items: list[dict], *, top_n: int = 2) -> set[str]:
3097 """Most-mentioned handles in a batch of X items, as first-party candidates.
3098
3099 Mirrors entity_extract's ranking but runs before pruning rather than after,
3100 and keys on *mentions only* rather than mentions plus authors. That
3101 distinction is the safety property: a prolific commentator inflates the
3102 author count, but being mentioned by other accounts is what identifies the
3103 subject of a topic. Capped at the top few so a busy thread cannot exempt
3104 the whole batch.
3105 """
3106 counts: Counter = Counter()
3107 for item in raw_items or []:
3108 text = str((item or {}).get("text") or "")
3109 for mention in re.findall(r"@([A-Za-z0-9_]{1,15})", text):
3110 counts[mention.lower()] += 1
3111 if not counts:
3112 return set()
3113 return {handle for handle, _ in counts.most_common(top_n)}
3114
3115
3116 # Reddit engagement keepers: per stream, the top-N threads by upvotes plus
3117 # comments that clear the relevance floor and name the primary entity survive
3118 # per_stream_limit truncation even when their local rank score is low. The
3119 # stream order is 65% title relevance, so the month's most-discussed on-topic
3120 # thread (16K upvotes, 0.19 relevance) was otherwise cut behind one-upvote
3121 # posts with better title overlap.
3122 REDDIT_STREAM_KEEPERS = 3
3123
3124
3125 def _apply_reddit_stream_keepers(
3126 source: str,
3127 items: list[schema.SourceItem],
3128 limit: int,
3129 topic: str,
3130 ) -> list[schema.SourceItem]:
3131 """Truncate a stream to *limit*, holding slots for Reddit engagement keepers."""
3132 kept = list(items[:limit])
3133 if source != "reddit" or len(items) <= limit:
3134 return kept
3135 entity = rerank._primary_entity(topic or "") if topic else ""
3136 floor = fusion.relevance_floor_for_entity(entity)
3137 keepers = [
3138 item
3139 for item in sorted(items, key=fusion.raw_engagement, reverse=True)
3140 if fusion.reddit_thread_qualifies(item, entity, floor)
3141 ][:REDDIT_STREAM_KEEPERS]
3142 keeper_ids = {id(item) for item in keepers}
3143 for keeper in keepers:
3144 if any(item is keeper for item in kept):
3145 continue
3146 # Displace the lowest-ranked non-keeper so the slice stays at limit;
3147 # when the slice is already all keepers there is nothing to trade.
3148 displaced = False
3149 for index in range(len(kept) - 1, -1, -1):
3150 if id(kept[index]) not in keeper_ids:
3151 del kept[index]
3152 displaced = True
3153 break
3154 if displaced or len(kept) < limit:
3155 kept.append(keeper)
3156 return kept[:limit]
3157
3158
3159
3160 def _creator_first_party_by_source(
3161 tiktok_creators: Iterable[str] | None,
3162 ig_creators: Iterable[str] | None,
3163 ) -> dict[str, set[str]]:
3164 """Platform-scoped creator handles for the relevance prune.
3165
3166 --ig-creators names Instagram accounts and --creators names TikTok
3167 accounts; each exemption applies only on its own platform. Merging both
3168 into one global handle set would let an unrelated same-name account on
3169 another platform bypass the relevance and engagement floors.
3170 """
3171
3172 def _norm(handles: Iterable[str] | None) -> set[str]:
3173 return {
3174 h.lstrip("@").strip().lower()
3175 for h in (handles or [])
3176 if h and h.strip()
3177 }
3178
3179 return {"instagram": _norm(ig_creators), "tiktok": _norm(tiktok_creators)}
3180
3181
3182 def _creator_only_handles(
3183 creator_first_party: Mapping[str, set[str]],
3184 *x_provenance_sets: Iterable[str],
3185 ) -> set[str]:
3186 """Creator handles that carry no X provenance.
3187
3188 A handle named ONLY via --ig-creators / --creators must not exempt a
3189 same-name X account from the deferred X floor. But the same person is
3190 often named on both surfaces (--x-handle foo --ig-creators foo): the
3191 normalized sets collapse that to one string, so subtracting the whole
3192 creator set would strip the explicitly requested X exemption too. The
3193 subtraction therefore covers only handles absent from every
3194 X-provenance set (explicit flags, topic mentions, Phase 2 discovery).
3195 """
3196 creator_flat = {h for handles in creator_first_party.values() for h in handles}
3197 x_provenance = {h for handles in x_provenance_sets for h in handles}
3198 return creator_flat - x_provenance
3199
3200
3201 def _log_prune_drop(
3202 source: str, before: int, after: int, scope: str | None = None
3203 ) -> None:
3204 """Log when the relevance prune removes items from a stream.
3205
3206 The prune is silent by design inside ``signals`` (a pure function), but a
3207 silent drop is invisible to the user: issue #1101 fetched 36 creator reels
3208 and reported zero with no line explaining why. Log the count and the reason
3209 class here, next to the other per-stream retrieval logs. The all-weak
3210 ``filtered or items`` rescue keeps the originals, so before == after and
3211 nothing is logged - a rescue is not a drop.
3212 """
3213 dropped = before - after
3214 if dropped <= 0:
3215 return
3216 scope_note = f" ({scope})" if scope else ""
3217 log.source_log(
3218 render.SOURCE_LABELS.get(source, source.capitalize()),
3219 f"relevance prune dropped {dropped} of {before} items below the "
3220 f"relevance/engagement floor{scope_note}",
3221 tty_only=False,
3222 )
3223
3224
3225 def _normalize_score_dedupe(
3226 source: str,
3227 raw_items: list[dict],
3228 from_date: str,
3229 to_date: str,
3230 freshness_mode: str,
3231 ranking_query: str,
3232 first_party_handles: Iterable[str] | None = None,
3233 first_party_by_source: Mapping[str, Iterable[str]] | None = None,
3234 defer_relevance_prune: bool = False,
3235 ) -> list[schema.SourceItem]:
3236 """Normalize, annotate, prune, dedupe, and extract snippets for a batch of raw items.
3237
3238 ``defer_relevance_prune`` skips the relevance floor here so the caller can
3239 apply it once the run has resolved who the topic's subject is. Pruning X
3240 before handle resolution is the ordering bug behind the whole first-party
3241 evidence loss: the floor cannot exempt an author nobody has identified yet,
3242 and no amount of guessing at prune time substitutes for knowing.
3243
3244 ``first_party_handles`` names accounts this run is explicitly searching, so
3245 their own posts survive the relevance floor (see signals.prune_low_relevance).
3246 """
3247 normalized = normalize.normalize_source_items(
3248 source, raw_items, from_date, to_date,
3249 freshness_mode=freshness_mode,
3250 )
3251 prepared_query = relevance.PreparedQuery(ranking_query)
3252 lookback_window_days = (
3253 datetime.strptime(to_date, "%Y-%m-%d").date()
3254 - datetime.strptime(from_date, "%Y-%m-%d").date()
3255 ).days
3256 normalized = signals.annotate_stream(
3257 normalized,
3258 prepared_query,
3259 freshness_mode,
3260 reference_date=to_date,
3261 max_days=lookback_window_days,
3262 )
3263 if source != "jobs" and not defer_relevance_prune:
3264 floor_handles = set(first_party_handles or ())
3265 if source == "x":
3266 # Union, never a fallback. The caller's set is derived partly from
3267 # topic tokens, so it is non-empty for essentially every real topic
3268 # -- gating this on "no handles supplied" would make it dead code
3269 # and leave the name-only case exactly as broken as before.
3270 #
3271 # Reuses the engine's own resolution signal on the batch already in
3272 # hand: posts *about* a subject mention their handle, so the
3273 # most-mentioned account in a topic's own results is the subject.
3274 # Costs nothing extra -- no search, no network -- and closes the
3275 # case where the handle never appears in the topic at all
3276 # ("Peter Steinberger" -> @steipete).
3277 floor_handles |= _batch_subject_handles(raw_items)
3278 pre_prune_count = len(normalized)
3279 normalized = signals.prune_low_relevance(
3280 normalized,
3281 first_party_handles=floor_handles,
3282 first_party_by_source=first_party_by_source,
3283 )
3284 _log_prune_drop(source, pre_prune_count, len(normalized))
3285 normalized = dedupe.dedupe_items(normalized)
3286 for item in normalized:
3287 item.snippet = snippet.extract_best_snippet(item, prepared_query)
3288 return normalized
3289
3290
3291 def _finalize_items_by_source(
3292 items_by_source_raw: dict[str, list[schema.SourceItem]],
3293 topic: str = "",
3294 config: dict | None = None,
3295 depth: str = "default",
3296 mock: bool = False,
3297 elapsed: float = 0.0,
3298 ) -> dict[str, list[schema.SourceItem]]:
3299 finalized = {}
3300 for source, items in items_by_source_raw.items():
3301 items = sorted(items, key=lambda item: item.local_rank_score or 0.0, reverse=True)
3302 # Same thread from two subquery streams: fold the enriched copy into
3303 # the first before the text-similarity dedupe, which would otherwise
3304 # keep whichever copy ranked higher and drop its comments.
3305 items = collapse_duplicate_urls(items)
3306 items = dedupe.dedupe_items(items)
3307 enrichment_request = {
3308 "source": source,
3309 "phase": "post_ranking_enrichment",
3310 "topic": topic,
3311 "depth": depth,
3312 }
3313 if source == "youtube" and items and not mock:
3314 # Same budget-at-the-survivors principle as the digg branch
3315 # below: retrieval-time transcripts go to each search's
3316 # top-by-views candidates, while final selection ranks by
3317 # relevance. Backfill survivors that arrived without one so the
3318 # transcript budget lands on videos the brief actually shows
3319 # (#542).
3320 matched, replayed = http.fixture_source_replay(enrichment_request)
3321 if matched:
3322 items = _merge_replayed_enrichment(items, replayed)
3323 else:
3324 sc_token = (
3325 config.get("SCRAPECREATORS_API_KEY")
3326 if config and env.is_youtube_sc_available(config) else None
3327 )
3328 youtube_yt.backfill_transcripts(
3329 items, topic=topic, depth=depth, token=sc_token,
3330 )
3331 http.fixture_source_record(enrichment_request, schema.to_dict(items))
3332 # Post-merge topic-relevance filter for Polymarket: comparison queries
3333 # fan out into per-entity subqueries ("Hermes", "OpenClaw") whose topic
3334 # is too narrow for Gamma API to filter meaningfully. Re-validating the
3335 # merged list against the full original topic drops off-topic markets
3336 # (e.g., WTI crude oil, Elon tweet counts) before footer emission.
3337 if source == "polymarket" and topic:
3338 items = polymarket.filter_items_against_topic(topic, items)
3339 # --polymarket-keywords (via config): additional keyword filter
3340 # for ambiguous single-token topics (e.g., "Warriors" → nba,gsw).
3341 keywords = config.get("_polymarket_keywords") if isinstance(config, dict) else None
3342 if keywords:
3343 items = polymarket.filter_items_against_keywords(items, keywords)
3344 if source == "digg" and items:
3345 # Pull top-ranked X posts only for the survivors that will appear
3346 # in the brief. Spending the enrichment budget here (rather than
3347 # at retrieval time) keeps the inline 'via Digg' quotes
3348 # paired with the clusters dedupe actually kept.
3349 matched, replayed = http.fixture_source_replay(enrichment_request)
3350 if matched:
3351 items = _merge_replayed_enrichment(items, replayed)
3352 else:
3353 digg.enrich_source_items(items, top_k=3)
3354 http.fixture_source_record(enrichment_request, schema.to_dict(items))
3355 if source == "amazon" and items and not mock:
3356 # Attach-if-missing: review enrichment now runs at search time in
3357 # _retrieve_stream_impl, so items arriving here should already have
3358 # top_comments. enrich_source_items no-ops when top_comments is set.
3359 # This path handles fixture replay and any edge cases where retrieve
3360 # didn't enrich (e.g., run_started was not passed).
3361 matched, replayed = http.fixture_source_replay(enrichment_request)
3362 if matched:
3363 items = _merge_replayed_enrichment(items, replayed)
3364 else:
3365 amazon.enrich_source_items(
3366 items,
3367 depth=depth,
3368 config=config,
3369 keyword=str((config or {}).get("_amazon_query") or "").strip() or topic,
3370 elapsed=elapsed,
3371 )
3372 http.fixture_source_record(enrichment_request, schema.to_dict(items))
3373 finalized[source] = items
3374 return finalized
3375
3376
3377 def _merge_replayed_enrichment(
3378 items: list[schema.SourceItem],
3379 replayed: list[dict],
3380 ) -> list[schema.SourceItem]:
3381 """Apply recorded post-ranking enrichment onto freshly computed items.
3382
3383 Enrichment (transcripts, Digg posts) only mutates ``metadata``. Merging by
3384 item_id instead of replacing the list keeps normalization, scoring, and
3385 dedupe regressions visible to the eval - fixture state must not overwrite
3386 what the current pipeline computed.
3387 """
3388 replayed_by_id = {
3389 entry.get("item_id"): entry for entry in replayed if isinstance(entry, dict)
3390 }
3391 for item in items:
3392 record = replayed_by_id.get(item.item_id)
3393 if record and record.get("metadata"):
3394 item.metadata.update(record["metadata"])
3395 return items
3396
3397
3398 def _apply_hiring_signal_gate(
3399 bundle: schema.RetrievalBundle,
3400 *,
3401 explicit: bool,
3402 topic: str,
3403 ) -> dict[str, Any] | None:
3404 jobs_items = bundle.items_by_source.get("jobs") or []
3405 if not jobs_items:
3406 if explicit:
3407 return hiring_signals.analyze([], explicit=True, topic=topic)
3408 return None
3409
3410 summary = hiring_signals.analyze(jobs_items, explicit=explicit, topic=topic)
3411 if not explicit and not summary.get("include"):
3412 bundle.items_by_source.pop("jobs", None)
3413 for key in list(bundle.items_by_source_and_query):
3414 if key[1] == "jobs":
3415 del bundle.items_by_source_and_query[key]
3416 return summary
3417
3418
3419 def _ensure_jobs_in_plan(
3420 plan: schema.QueryPlan,
3421 available: list[str],
3422 *,
3423 explicit: bool,
3424 topic: str,
3425 ) -> None:
3426 if "jobs" not in available:
3427 return
3428 if not (explicit or _company_topic_likely(topic)):
3429 return
3430 if "jobs" not in plan.source_weights:
3431 plan.source_weights["jobs"] = 1.0
3432 for subquery in plan.subqueries:
3433 if "jobs" not in subquery.sources:
3434 subquery.sources.append("jobs")
3435
3436
3437 def _ensure_perplexity_in_plan(
3438 plan: schema.QueryPlan,
3439 topic: str,
3440 available: list[str],
3441 *,
3442 force: bool,
3443 ) -> None:
3444 """Route a bounded paid Perplexity action through the whole topic.
3445
3446 Deep Research forces its explicit lane. Normal modes are rerouted only when
3447 the sanitized plan already selected Perplexity.
3448 """
3449 if "perplexity" not in available:
3450 return
3451 planned = any(
3452 "perplexity" in subquery.sources for subquery in plan.subqueries
3453 )
3454 if not force and not planned:
3455 return
3456 retained: list[schema.SubQuery] = []
3457 for subquery in plan.subqueries:
3458 sources = [
3459 source for source in subquery.sources if source != "perplexity"
3460 ]
3461 if sources:
3462 retained.append(replace(subquery, sources=sources))
3463 retained.append(
3464 schema.SubQuery(
3465 label="deep-research" if force else "perplexity-whole-topic",
3466 search_query=topic,
3467 ranking_query=f"What current source-grounded evidence matters for {topic}?",
3468 sources=["perplexity"],
3469 weight=1.0,
3470 ),
3471 )
3472 plan.subqueries = planner._normalize_subquery_weights(retained)
3473 plan.source_weights.setdefault("perplexity", 1.0)
3474 plan.source_weights = planner._normalize_weights(plan.source_weights)
3475
3476
3477 def _company_topic_likely(topic: str) -> bool:
3478 text = topic.strip()
3479 if not text:
3480 return False
3481 lower = text.lower()
3482 if "?" in text or len(text.split()) > 4:
3483 return False
3484 generic = {
3485 "how", "what", "why", "best", "top", "tutorial", "guide", "prompts",
3486 "news", "latest", "ideas", "examples",
3487 }
3488 if any(word in generic for word in lower.split()):
3489 return False
3490 known_single_word_companies = {
3491 "apple", "uber", "google", "microsoft", "amazon", "meta", "netflix",
3492 "openai", "anthropic", "qualtrics", "stripe", "brex",
3493 }
3494 if " vs " in lower or " versus " in lower:
3495 parts = re.split(r"\s+(?:vs|versus)\s+", text, maxsplit=1, flags=re.IGNORECASE)
3496 if len(parts) != 2:
3497 return False
3498 return _comparison_side_company_like(parts[0], known_single_word_companies) or _comparison_side_company_like(
3499 parts[1], known_single_word_companies
3500 )
3501 return bool(text[:1].isupper() or lower in known_single_word_companies)
3502
3503
3504 def _comparison_side_company_like(side: str, known_companies: set[str]) -> bool:
3505 token = re.sub(r"[^\w.+#-]", "", side.strip().split()[0] if side.strip() else "")
3506 if not token:
3507 return False
3508 lower = token.lower()
3509 common_tech_terms = {
3510 "python", "ruby", "javascript", "typescript", "java", "go", "golang",
3511 "rust", "php", "swift", "kotlin", "scala", "clojure", "elixir",
3512 "react", "vue", "angular", "svelte", "node", "django", "rails",
3513 "postgres", "mysql", "redis", "kubernetes", "docker",
3514 }
3515 if lower in common_tech_terms:
3516 return False
3517 return bool(token[:1].isupper() or lower in known_companies)
3518
3519
3520 def _warnings(
3521 items_by_source: dict[str, list[schema.SourceItem]],
3522 candidates: list[schema.Candidate],
3523 errors_by_source: dict[str, str],
3524 degraded_by_source: dict[str, str] | None = None,
3525 ) -> list[str]:
3526 warnings: list[str] = []
3527 if not candidates:
3528 warnings.append("No candidates survived retrieval and ranking.")
3529 if len(candidates) < 5:
3530 warnings.append("Evidence is thin for this topic.")
3531 top_sources = {
3532 source
3533 for candidate in candidates[:5]
3534 for source in schema.candidate_sources(candidate)
3535 }
3536 if len(top_sources) <= 1 and len(candidates) >= 3:
3537 warnings.append("Top evidence is highly concentrated in one source.")
3538 if errors_by_source:
3539 warnings.append(f"Some sources failed: {', '.join(sorted(errors_by_source))}")
3540 if degraded_by_source:
3541 # Partial failures: the source returned some items but errored/timed out
3542 # on at least one subquery, so its coverage is likely incomplete. Kept
3543 # distinct from hard failures so the signal is not silently dropped.
3544 warnings.append(
3545 f"Some sources returned partial results (degraded): {', '.join(sorted(degraded_by_source))}"
3546 )
3547 if not items_by_source:
3548 warnings.append("No source returned usable items.")
3549 return warnings
3550
3551
3552 _STATUS_PREFIX = r"\b(?:https?(?:/\d(?:\.\d)?)?(?:\s+error)?|status(?:[\s_]*code)?|code)\s*[:=#]?\s*"
3553
3554
3555 def _mentions_status(msg: str, code_pattern: str, phrases: tuple[str, ...]) -> bool:
3556 """True when ``msg`` names an HTTP status matching ``code_pattern``.
3557
3558 A bare number is not enough (``"batch 429 failed"`` is not a rate
3559 limit): the code must follow an HTTP/status/code marker, or co-occur
3560 with one of ``phrases`` (``"rate limited (429)"``). Digit-aware
3561 boundaries keep ``"14293"`` from reading as a 429.
3562 """
3563 if not msg:
3564 return False
3565 code = r"(?:" + code_pattern + r")(?!\d)"
3566 if re.search(_STATUS_PREFIX + code, msg, re.IGNORECASE):
3567 return True
3568 lowered = msg.lower()
3569 return any(p in lowered for p in phrases) and re.search(r"(?<!\d)" + code, msg) is not None
3570
3571
3572 _RATE_LIMIT_PHRASES = ("rate limit", "rate-limit", "ratelimit", "too many requests")
3573 _SERVER_ERROR_PHRASES = (
3574 "server error",
3575 "bad gateway",
3576 "service unavailable",
3577 "gateway timeout",
3578 "gateway time-out",
3579 )
3580
3581
3582 def _is_rate_limit_error(exc: Exception) -> bool:
3583 """Detect 429 rate-limit errors by status code or message text."""
3584 if hasattr(exc, "status_code") and getattr(exc, "status_code", None) == 429:
3585 return True
3586 return _mentions_status(str(exc), "429", _RATE_LIMIT_PHRASES)
3587
3588
3589 class SourceRunError(RuntimeError):
3590 """Source-specific failure that survived a module's fallback logic."""
3591
3592 def __init__(self, message: str, state: schema.RunOutcomeState | None = None):
3593 super().__init__(message)
3594 self.outcome_state = state or http.classify_failure(message=message)
3595
3596
3597 def _classify_source_failure(exc: Exception) -> tuple[schema.RunOutcomeState, bool]:
3598 """Classify HTTP, subprocess, and module-specific failures consistently."""
3599 detail = str(exc)
3600 lowered = detail.lower()
3601 if any(marker in lowered for marker in ("not configured", "no api key", "not installed")):
3602 return schema.SKIPPED_UNCONFIGURED, False
3603 if any(
3604 marker in lowered
3605 for marker in (
3606 "cookie expired",
3607 "expired cookie",
3608 "login required",
3609 "not logged in",
3610 "grok session expired",
3611 "session expired or was revoked",
3612 "invalid_grant",
3613 "not signed in",
3614 )
3615 ):
3616 return schema.AUTH_FAILED, True
3617 state = getattr(exc, "outcome_state", None) or http.classify_failure(
3618 status_code=getattr(exc, "status_code", None),
3619 message=detail,
3620 )
3621 return state, True
3622
3623
3624 def _outcome_artifact(
3625 state: schema.RunOutcomeState,
3626 detail: str,
3627 *,
3628 attempted: bool = True,
3629 ) -> dict[str, Any]:
3630 return {
3631 "_source_outcome": {
3632 "state": state,
3633 "detail": detail,
3634 "attempted": attempted,
3635 }
3636 }
3637
3638
3639 def _result_outcome_artifact(source: str, result: Any) -> dict[str, Any]:
3640 """Convert a legacy ``{"error": ...}`` source result into typed status."""
3641 if not isinstance(result, dict) or not result.get("error"):
3642 return {}
3643 detail = str(result["error"])
3644 if source == "reddit":
3645 state = reddit.classify_run_failure(detail)
3646 attempted = True
3647 elif source == "youtube":
3648 state = youtube_yt.classify_run_failure(detail)
3649 attempted = state != schema.SKIPPED_UNCONFIGURED
3650 elif source == "x":
3651 state = bird_x.classify_run_failure(detail)
3652 attempted = True
3653 elif source == "truthsocial" and detail == "Truth Social token expired":
3654 state = schema.AUTH_FAILED
3655 attempted = True
3656 elif source == "bluesky" and "network-level block" in detail.lower():
3657 state = schema.UNREACHABLE
3658 attempted = True
3659 else:
3660 state, attempted = _classify_source_failure(SourceRunError(detail))
3661 return _outcome_artifact(state, detail, attempted=attempted)
3662
3663
3664 def _legacy_artifact_outcome(
3665 source: str,
3666 artifact: Any,
3667 ) -> dict[str, Any] | None:
3668 """Map known pre-outcome artifact contracts to a typed outcome note."""
3669 if not isinstance(artifact, dict):
3670 return None
3671 explicit = artifact.get("_source_outcome")
3672 if isinstance(explicit, dict):
3673 return explicit
3674 if source == "perplexity":
3675 candidates: list[tuple[str | None, dict[str, Any]]] = [(None, artifact)]
3676 if artifact.get("mode") == "both":
3677 for leg in ("search", "agent"):
3678 value = artifact.get(leg)
3679 if isinstance(value, dict):
3680 candidates.append((leg, value))
3681 outcomes: list[dict[str, Any]] = []
3682 for leg, candidate in candidates:
3683 if not candidate.get("error"):
3684 continue
3685 error = str(candidate["error"])
3686 detail = str(
3687 candidate.get("backgroundErrorMessage")
3688 or candidate.get("backgroundPollError")
3689 or candidate.get("agentErrorMessage")
3690 or candidate.get("asyncErrorMessage")
3691 or candidate.get("message")
3692 or error
3693 )
3694 if leg:
3695 detail = f"{leg} leg: {detail}"
3696 status_code = candidate.get("statusCode")
3697 if status_code is None:
3698 status_code = candidate.get("backgroundPollStatusCode")
3699 state = (
3700 health.TIMEOUT
3701 if error.lower() == "timeout"
3702 else http.classify_failure(
3703 status_code=status_code,
3704 message=f"{error}: {detail}",
3705 )
3706 )
3707 outcomes.append(_outcome_artifact(state, detail)["_source_outcome"])
3708 if outcomes:
3709 return min(
3710 outcomes,
3711 key=lambda outcome: _FAILURE_SPECIFICITY.get(outcome["state"], 9),
3712 )
3713 if (
3714 source == "grounding"
3715 and artifact.get("reason") == "keyless-search-unavailable"
3716 ):
3717 return _outcome_artifact(
3718 schema.UNREACHABLE,
3719 "Keyless web search unavailable",
3720 )["_source_outcome"]
3721 return None
3722
3723
3724 def _summarize_lane_failures(failures: list[http.HTTPError], source: str = "") -> str:
3725 """One line naming what a source lost to swallowed sub-request failures.
3726
3727 ``"3 sub-requests rate-limited (HTTP 429); 1 sub-request blocked (HTTP 403)"``.
3728 Used as ``SourceOutcome.detail`` on a source that still delivered items,
3729 so the loss is visible to ``doctor --postmortem`` without branding the
3730 source partial (issue #985 wording; PR #959 semantics).
3731 """
3732 counts: dict[tuple[str, int | None], int] = {}
3733 for failure in failures:
3734 state = getattr(failure, "outcome_state", None) or health.ERROR
3735 code = getattr(failure, "status_code", None)
3736 counts[(state, code)] = counts.get((state, code), 0) + 1
3737 labels = {
3738 health.RATE_LIMITED: "rate-limited",
3739 health.AUTH_FAILED: "blocked",
3740 health.PAYMENT_REQUIRED: health.credits_exhausted_label(source),
3741 health.TIMEOUT: "timed out",
3742 health.UNREACHABLE: "unreachable",
3743 health.SCHEMA_DRIFT: "returned an unexpected shape",
3744 }
3745 parts = []
3746 for (state, code), n in sorted(counts.items(), key=lambda kv: -kv[1]):
3747 noun = "sub-request" if n == 1 else "sub-requests"
3748 label = labels.get(state, "failed")
3749 suffix = f" (HTTP {code})" if code else ""
3750 parts.append(f"{n} {noun} {label}{suffix}")
3751 return "; ".join(parts)
3752
3753
3754 def _resolve_stream_outcome(
3755 source: str,
3756 artifact: Any,
3757 failures: list[http.HTTPError],
3758 ) -> dict[str, Any] | None:
3759 """Choose the most specific artifact or captured HTTP outcome."""
3760 artifact_outcome = _legacy_artifact_outcome(source, artifact)
3761 if not failures:
3762 return artifact_outcome
3763 # Pick the most specific failure rather than the last-appended one:
3764 # parallel workers append in nondeterministic order, and an auth failure
3765 # must not be masked by a later 429 (wrong doctor prescription).
3766 failure = min(
3767 failures,
3768 key=lambda f: _FAILURE_SPECIFICITY.get(f.outcome_state, 9),
3769 )
3770 captured_outcome = _outcome_artifact(
3771 failure.outcome_state,
3772 str(failure),
3773 )["_source_outcome"]
3774 if artifact_outcome is None:
3775 return captured_outcome
3776 if (
3777 artifact_outcome.get("state") == health.ERROR
3778 and failure.outcome_state != health.ERROR
3779 ):
3780 return captured_outcome
3781 return artifact_outcome
3782
3783
3784 def _finalize_source_status(
3785 outcomes: dict[str, schema.SourceOutcome],
3786 items_by_source: dict[str, list[schema.SourceItem]],
3787 ) -> dict[str, schema.SourceOutcome]:
3788 """Sync outcome counts to the final post-filter evidence set."""
3789 finalized: dict[str, schema.SourceOutcome] = {}
3790 for source, outcome in outcomes.items():
3791 count = len(items_by_source.get(source, []))
3792 state = outcome.state
3793 detail = outcome.detail
3794 fix_hint = outcome.fix_hint
3795 if state == schema.NO_RESULTS and count:
3796 state = health.OK
3797 detail = None
3798 fix_hint = None
3799 elif state == health.OK and not count:
3800 state = outcome.lane_failure_state or schema.NO_RESULTS
3801 elif state == schema.PARTIAL and not count:
3802 state = http.classify_failure(message=detail or "")
3803 finalized[source] = schema.SourceOutcome(
3804 source=source,
3805 state=state,
3806 items_returned=count,
3807 attempted=outcome.attempted,
3808 detail=detail,
3809 at=outcome.at,
3810 fix_hint=fix_hint,
3811 lane_failure_state=outcome.lane_failure_state,
3812 )
3813 return finalized
3814
3815
3816 def _is_transient_error(exc: Exception) -> bool:
3817 """Detect 5xx server errors that are worth retrying."""
3818 status = getattr(exc, "status_code", None)
3819 if isinstance(status, int) and 500 <= status < 600:
3820 return True
3821 return _mentions_status(str(exc), r"5\d\d", _SERVER_ERROR_PHRASES)
3822
3823
3824 def _topic_handle_mentions(topic: str) -> set[str]:
3825 """@mentions in the topic, which are real X handles.
3826
3827 These are used to determine whether the subject was identified: an
3828 @mention like "@steipete" is a real handle that can exempt its owner from
3829 the relevance floor. Regular words like "Peter" are not real handles.
3830 """
3831 return {
3832 mention.lower()
3833 for mention in re.findall(r"@([A-Za-z0-9_]{1,15})", topic or "")
3834 }
3835
3836
3837 def _topic_first_party_candidates(topic: str) -> set[str]:
3838 """Handle-shaped tokens in the topic itself, usable before any retrieval.
3839
3840 Phase 1 runs before automatic handle resolution, and a quick-depth run
3841 skips that resolution entirely, so neither has access to the extracted
3842 handle set. Without this a quick search for "Peter Steinberger steipete"
3843 still drops every post steipete wrote, which is the exact failure this
3844 branch exists to fix.
3845
3846 Deliberately permissive about what looks like a handle and strict about
3847 what it does: a candidate only ever matters if a retrieved post's *author*
3848 matches it, so an ordinary word like "lunch" costs nothing -- no author is
3849 named "lunch". The realistic false positive is an account named after a
3850 topic word, which the frequency-ranked path could surface anyway.
3851 """
3852 tokens = set()
3853 for mention in re.findall(r"@([A-Za-z0-9_]{1,15})", topic or ""):
3854 tokens.add(mention.lower())
3855 for word in re.findall(r"[A-Za-z0-9_]{3,15}", topic or ""):
3856 lowered = word.lower()
3857 if lowered not in relevance.STOPWORDS:
3858 tokens.add(lowered)
3859 return tokens
3860
3861
3862 def _name_lane_subject(topic: str) -> str:
3863 """Resolve the entity name to search for by name, not the whole topic.
3864
3865 Phrase-quoting a raw topic ("Peter Steinberger steipete") matches nothing
3866 on X: nobody writes the handle and the display name together. Prefer a
3867 title-cased proper noun the way the planner's keyword query does, and fall
3868 back to the first compound term, then to the topic.
3869 """
3870 import re as _re
3871 compounds = query.extract_compound_terms(topic) or []
3872 title_cased = [
3873 term for term in compounds
3874 if _re.match(r"^(?:[A-Z][a-z]+\s+){1,}[A-Z][a-z]+$", term)
3875 ]
3876 if title_cased:
3877 return title_cased[0]
3878 if compounds:
3879 return compounds[0]
3880 return topic.strip()
3881
3882
3883 def _run_supplemental_searches(
3884 *,
3885 topic: str,
3886 bundle: schema.RetrievalBundle,
3887 plan: schema.QueryPlan,
3888 config: dict[str, Any],
3889 depth: str,
3890 date_range: tuple[str, str],
3891 runtime: schema.ProviderRuntime,
3892 mock: bool,
3893 rate_limited_sources: set[str],
3894 rate_limit_lock: threading.Lock,
3895 x_handle: str | None = None,
3896 x_related: list[str] | None = None,
3897 resolved_handles_out: list[str] | None = None,
3898 ) -> None:
3899 """Phase 2: extract entities from Phase 1 results, run targeted supplemental searches."""
3900 # The sanitized plan already intersects requested, available, and excluded
3901 # sources. A handle or host envelope must not expand that source boundary.
3902 if not any("x" in subquery.sources for subquery in plan.subqueries):
3903 return
3904 from_date, to_date = date_range
3905
3906 # Host-fetched X lane: the envelope's lane calls replace the backend
3907 # lanes and are served at every depth (the host already paid for them),
3908 # before the quick/mock return and before the chain is recomputed.
3909 # Extracted-handle promotion is skipped on envelope runs; a declared lane
3910 # without an envelope runs no lane at all (the topic stream already
3911 # recorded the not-passed outcome).
3912 if config.get("_x_lane_missing"):
3913 return
3914 envelope = config.get("_x_envelope")
3915 if envelope is not None:
3916 _serve_envelope_lanes(
3917 envelope, bundle=bundle, plan=plan, x_handle=x_handle, x_related=x_related,
3918 from_date=from_date, to_date=to_date,
3919 )
3920 return
3921
3922 if depth == "quick" or mock:
3923 return
3924
3925 # Convert SourceItems to dicts for entity_extract. All X items (whatever
3926 # backend fetched them — bird, xai, xurl, xquik) land under the single "x"
3927 # slug, so this reads the whole X corpus.
3928 x_dicts = [
3929 {"author_handle": item.author or "", "text": item.body or ""}
3930 for item in bundle.items_by_source.get("x", [])
3931 ]
3932 reddit_dicts = [
3933 {
3934 "subreddit": item.container or "",
3935 "comment_insights": item.metadata.get("comment_insights", []),
3936 "top_comments": [
3937 {"excerpt": c.get("excerpt", c.get("text", ""))}
3938 for c in (item.metadata.get("top_comments") or [])
3939 if isinstance(c, dict)
3940 ],
3941 }
3942 for item in bundle.items_by_source.get("reddit", [])
3943 ]
3944
3945 if not x_dicts and not reddit_dicts and not x_handle and not x_related:
3946 return
3947
3948 entities = entity_extract.extract_entities(
3949 reddit_dicts, x_dicts,
3950 max_handles=3, max_subreddits=3,
3951 )
3952
3953 handles = entities.get("x_handles", [])
3954
3955 # Add explicit --x-handle if provided
3956 if x_handle:
3957 handle_clean = x_handle.lstrip("@").lower()
3958 if handle_clean not in [h.lower() for h in handles]:
3959 handles.insert(0, handle_clean)
3960
3961 # Collect related handles (searched separately with lower weight)
3962 related_handles = []
3963 if x_related:
3964 primary_lower = x_handle.lstrip("@").lower() if x_handle else ""
3965 for rh in x_related:
3966 rh_clean = rh.lstrip("@").lower().strip()
3967 if rh_clean and rh_clean != primary_lower and rh_clean not in [h.lower() for h in handles]:
3968 related_handles.append(rh_clean)
3969
3970 # Surface every handle this run resolved back to the caller. resolved_handles
3971 # is built later from --x-handle / --github-user / --x-related only, so
3972 # without this an auto-discovered subject handle never reaches it and every
3973 # downstream first-party protection (entity-miss exemption, FIRST_PARTY_FLOOR,
3974 # interaction floor) stays inert on any run that did not pass --x-handle.
3975 # Populated before the early return below so a run whose lanes cannot execute
3976 # still contributes its resolved handles.
3977 if resolved_handles_out is not None:
3978 # Only corroborated handles get first-party status. The extracted set is
3979 # frequency-ranked over retrieved post text, so a prolific commentator --
3980 # or an engagement-farming account that posts on every topic -- lands in
3981 # it without being the subject. First-party status is strong: it exempts
3982 # an author from the relevance floor entirely and raises their per-author
3983 # cap, so granting it on frequency alone would let a spam account buy
3984 # immunity from filtering. Require the handle to look like the topic's
3985 # subject, or to have been named explicitly by the user.
3986 explicit = {
3987 h.lstrip("@").strip().lower()
3988 for h in ([x_handle] + list(x_related or []))
3989 if h and h.strip()
3990 }
3991 topic_tokens = {t for t in re.findall(r"[a-z0-9]+", topic.lower()) if len(t) > 2}
3992 seen = {h.lower() for h in resolved_handles_out}
3993 for h in [*handles, *related_handles]:
3994 clean = h.lstrip("@").strip().lower()
3995 if not clean or clean in seen:
3996 continue
3997 corroborated = clean in explicit or any(
3998 token in clean or clean in token for token in topic_tokens
3999 )
4000 if corroborated:
4001 resolved_handles_out.append(clean)
4002 seen.add(clean)
4003
4004 if not handles and not related_handles:
4005 return
4006
4007 # Pick the X handle-search backend: the first handle-capable backend in the
4008 # chain (grok, bird, xapi, or xquik). These supplemental from:/mentions lanes are
4009 # complementary to the topic search, so when the topic primary can't run
4010 # them (xai/xurl have no handle-lane implementation) but a capable backend
4011 # is available, use it rather than skipping Phase 2. bird scrapes X GraphQL
4012 # with the user's browser cookies; xquik runs the same lanes over its REST
4013 # API. All items land under the single "x" slug.
4014 x_slug = "x"
4015
4016 def _record_handle_lane_failures(lane: str, failures: list[str]) -> bool:
4017 """Surface handle-lane failures the adapter would otherwise swallow.
4018
4019 bird and xquik report a per-handle failure by returning no items, so an
4020 empty lane is indistinguishable from a subject who simply did not post.
4021 Left unrecorded, the run reports the X source as a clean zero and the
4022 report states as fact that nothing was posted.
4023
4024 Returns True when the failure is auth-shaped, so the caller's existing
4025 AUTH_FAILED branch owns the message and the fix hint. Everything else
4026 (timeout, spawn failure, non-zero exit, bad JSON) is recorded here.
4027 ``record_failure`` keeps already-returned items as partial.
4028 """
4029 if not failures:
4030 return False
4031 detail = "; ".join(failures[:3])
4032 if len(failures) > 3:
4033 detail += f" (+{len(failures) - 3} more)"
4034 # Only xquik reports an auth-shaped handle-lane failure; bird's are all
4035 # transport (timeout / spawn / exit / JSON). Match the two phrases
4036 # xquik._execute_search actually emits rather than adding a fifth
4037 # copy of this repo's auth-marker vocabulary.
4038 lowered = detail.lower()
4039 if "auth failed" in lowered or "key unpaid" in lowered:
4040 return True
4041 bundle.record_failure(
4042 x_slug,
4043 health.UNREACHABLE,
4044 f"Phase 2 {lane}-lane: {detail}",
4045 attempted=True,
4046 )
4047 return False
4048
4049 chain = env.x_backend_chain(config)
4050 # Trust an explicit runtime backend as the head of the chain.
4051 pinned = runtime.x_search_backend
4052 if pinned:
4053 chain = [pinned] + [b for b in chain if b != pinned]
4054 primary = next((b for b in chain if b in ("grok", "bird", "xapi", "xquik")), None)
4055
4056 # Name lane (posts naming the subject in plain text, no @-mention) is
4057 # grok-only for now: it needs phrase-quoting and negation operators the
4058 # other handle-capable backends do not expose uniformly. It is NOT a
4059 # fallback for the mention lane -- most discussion of a person or company
4060 # never @-mentions them, so the two lanes reach disjoint sets.
4061 _name_lane = None
4062
4063 if primary == "grok":
4064 # One budget shared by all three lanes, started here rather than per
4065 # lane: the point is to bound the total, not each part.
4066 lane_deadline = time.monotonic() + grok_x.LANE_BUDGET_SECONDS
4067
4068 def _from_lane(hs: list, count: int, and_topic: bool = False) -> tuple[list, bool]:
4069 items, revoked = grok_x.search_handles(
4070 hs, topic, from_date, to_date, count_per=count,
4071 deadline=lane_deadline, and_topic=and_topic,
4072 )
4073 return items, revoked
4074
4075 def _about_lane(hs: list, count: int) -> tuple[list, bool]:
4076 items, revoked = grok_x.search_mentions(
4077 hs, from_date, to_date, topic=topic, count_per=count,
4078 deadline=lane_deadline,
4079 )
4080 return items, revoked
4081
4082 def _name_lane(hs: list, count: int) -> tuple[list, bool]:
4083 # Use the resolved entity name, not the raw topic. Phrase-quoting
4084 # the whole topic ("Peter Steinberger steipete") matches nothing on
4085 # X; the subject's name is what other people actually write.
4086 subject = _name_lane_subject(topic)
4087 if not subject.strip():
4088 return [], False
4089 items, revoked = grok_x.search_name(
4090 subject, from_date, to_date, exclude_handles=hs, count_per=count,
4091 deadline=lane_deadline,
4092 )
4093 return items, revoked
4094 elif primary == "bird":
4095 def _from_lane(hs: list, count: int, and_topic: bool = False) -> tuple[list, bool]:
4096 # bird_x.search_handles doesn't support and_topic yet
4097 failures: list[str] = []
4098 items = bird_x.search_handles(
4099 hs, topic, from_date, count_per=count, failure_out=failures,
4100 to_date=to_date,
4101 )
4102 return items, _record_handle_lane_failures("FROM", failures)
4103
4104 def _about_lane(hs: list, count: int) -> tuple[list, bool]:
4105 failures: list[str] = []
4106 items = bird_x.search_mentions(
4107 hs, from_date, count_per=count, failure_out=failures,
4108 to_date=to_date,
4109 )
4110 return items, _record_handle_lane_failures("ABOUT", failures)
4111 elif primary == "xapi":
4112 # Direct X API v2 with the app-only bearer: from:/@ lanes run over
4113 # search/all with the recent-search fallback. One budget shared by
4114 # every lane below (same shape as the grok lanes): a slow key bounds
4115 # the whole supplemental phase, not each call.
4116 xapi_token = config.get("X_BEARER_TOKEN") or ""
4117 xapi_deadline = time.monotonic() + x_api.LANE_BUDGET_SECONDS
4118
4119 def _xapi_lane_receipt(lane_warnings: list[str]) -> None:
4120 # A deadline stop is incomplete coverage, reported in
4121 # report.warnings (the x_partial_coverage artifact), never a
4122 # healthy-looking silence and never a source failure.
4123 sink = bundle.artifacts.setdefault("x_partial_coverage", [])
4124 for note in lane_warnings:
4125 line = f"X handle lanes: {note}"
4126 if line not in sink:
4127 sink.append(line)
4128
4129 def _from_lane(hs: list, count: int, and_topic: bool = False) -> tuple[list, bool]:
4130 # x_api.search_handles doesn't support and_topic; topic ranks only
4131 lane_warnings: list[str] = []
4132 items = x_api.search_handles(
4133 hs, topic, from_date, to_date, count_per=count, token=xapi_token,
4134 deadline=xapi_deadline, warnings=lane_warnings,
4135 )
4136 _xapi_lane_receipt(lane_warnings)
4137 return items, False
4138
4139 def _about_lane(hs: list, count: int) -> tuple[list, bool]:
4140 lane_warnings: list[str] = []
4141 items = x_api.search_mentions(
4142 hs, from_date, to_date, topic=topic, count_per=count, token=xapi_token,
4143 deadline=xapi_deadline, warnings=lane_warnings,
4144 )
4145 _xapi_lane_receipt(lane_warnings)
4146 return items, False
4147 elif primary == "xquik":
4148 xquik_token = env.get_xquik_token(config)
4149
4150 def _from_lane(hs: list, count: int, and_topic: bool = False) -> tuple[list, bool]:
4151 # xquik.search_handles doesn't support and_topic yet
4152 failures: list[str] = []
4153 items = xquik.search_handles(
4154 hs, topic, from_date, to_date, count_per=count, token=xquik_token,
4155 failure_out=failures,
4156 )
4157 return items, _record_handle_lane_failures("FROM", failures)
4158
4159 def _about_lane(hs: list, count: int) -> tuple[list, bool]:
4160 failures: list[str] = []
4161 items = xquik.search_mentions(
4162 hs, from_date, to_date, topic=topic, count_per=count,
4163 token=xquik_token, failure_out=failures,
4164 )
4165 return items, _record_handle_lane_failures("ABOUT", failures)
4166 else:
4167 return # primary X backend has no handle-lane support (xai/xurl) or none configured
4168
4169 # Skip if the X source is rate-limited.
4170 if x_slug in rate_limited_sources:
4171 return
4172
4173 # Collect existing URLs for deduplication
4174 existing_urls = {
4175 item.url
4176 for items in bundle.items_by_source.values()
4177 for item in items
4178 if item.url
4179 }
4180
4181 ranking_query = plan.subqueries[0].ranking_query if plan.subqueries else topic
4182 primary_label = plan.subqueries[0].label if plan.subqueries else "primary"
4183
4184 # Split FROM promotion: determine which handles get FROM lane and how.
4185 # - Primary explicit handle (--x-handle): always FROM, no AND topic, full weight
4186 # - x_related handles: searched separately with lower weight (0.3), kept in
4187 # related_handles variable for the supplemental-related section below
4188 # - Extracted handles: FROM only if ≥2 on-topic hits AND ratio ≥0.5,
4189 # and those pulls DO AND the topic (from:handle Rome)
4190 primary_explicit = [x_handle] if x_handle else []
4191
4192 explicit_promotable, extracted_promotable = x_judge.promotable_handles(
4193 x_dicts, # Phase 1 X items for judging
4194 topic,
4195 handles, # entity_extract handles
4196 explicit_handles=primary_explicit,
4197 ranking_query=ranking_query,
4198 )
4199
4200 # All promotable handles for ABOUT and NAME lanes (primary only, not related)
4201 all_promotable = list(set(explicit_promotable + extracted_promotable))
4202
4203 # Search primary handles (full weight): FROM lane (their own tweets) +
4204 # ABOUT lane (tweets mentioning them). Both engagement-weighted and deduped
4205 # by URL at normalize time.
4206 any_revoked = False # Track auth revocation across lanes
4207 if all_promotable:
4208 # Independent try/except per lane so a failure in one does not discard
4209 # the other's already-computed results.
4210 from_items: list = []
4211 about_items: list = []
4212 about_revoked = False
4213 name_revoked = False
4214
4215 # FROM lane: explicit handles without AND topic (person posts omit their own name)
4216 if explicit_promotable:
4217 try:
4218 explicit_items, explicit_revoked = _from_lane(explicit_promotable, FROM_LANE_COUNT_PER, and_topic=False)
4219 from_items.extend(explicit_items)
4220 if explicit_revoked:
4221 any_revoked = True
4222 bundle.record_failure(
4223 x_slug, schema.AUTH_FAILED,
4224 f"Phase 2 FROM-lane (explicit): {primary} authentication failed (session expired, revoked, or key unpaid)",
4225 attempted=True,
4226 )
4227 except Exception as exc:
4228 print(f"[Pipeline] Phase 2 FROM-lane (explicit) failed: {exc}", file=sys.stderr)
4229 state, attempted = _classify_source_failure(exc)
4230 bundle.record_failure(
4231 x_slug, state, f"Phase 2 FROM-lane (explicit): {exc}", attempted=attempted,
4232 )
4233
4234 # FROM lane: extracted handles WITH AND topic (from:handle Rome)
4235 if extracted_promotable:
4236 try:
4237 extracted_items, extracted_revoked = _from_lane(extracted_promotable, FROM_LANE_COUNT_PER, and_topic=True)
4238 from_items.extend(extracted_items)
4239 if extracted_revoked:
4240 any_revoked = True
4241 bundle.record_failure(
4242 x_slug, schema.AUTH_FAILED,
4243 f"Phase 2 FROM-lane (extracted): {primary} authentication failed (session expired, revoked, or key unpaid)",
4244 attempted=True,
4245 )
4246 except Exception as exc:
4247 print(f"[Pipeline] Phase 2 FROM-lane (extracted) failed: {exc}", file=sys.stderr)
4248 state, attempted = _classify_source_failure(exc)
4249 bundle.record_failure(
4250 x_slug, state, f"Phase 2 FROM-lane (extracted): {exc}", attempted=attempted,
4251 )
4252 if not bundle.items_by_source.get(x_slug):
4253 bundle.errors_by_source[x_slug] = f"Phase 2 FROM-lane: {exc}"
4254
4255 try:
4256 about_items, about_revoked = _about_lane(all_promotable, MENTION_LANE_COUNT_PER)
4257 if about_revoked:
4258 any_revoked = True
4259 bundle.record_failure(
4260 x_slug, schema.AUTH_FAILED,
4261 f"Phase 2 ABOUT-lane: {primary} authentication failed (session expired, revoked, or key unpaid)",
4262 attempted=True,
4263 )
4264 except Exception as exc:
4265 print(f"[Pipeline] Phase 2 ABOUT-lane search failed: {exc}", file=sys.stderr)
4266 state, attempted = _classify_source_failure(exc)
4267 bundle.record_failure(
4268 x_slug,
4269 state,
4270 f"Phase 2 ABOUT-lane: {exc}",
4271 attempted=attempted,
4272 )
4273 name_items: list = []
4274 if _name_lane is not None:
4275 try:
4276 name_items, name_revoked = _name_lane(all_promotable, MENTION_LANE_COUNT_PER)
4277 if name_revoked:
4278 any_revoked = True
4279 bundle.record_failure(
4280 x_slug, schema.AUTH_FAILED,
4281 f"Phase 2 NAME-lane: {primary} authentication failed (session expired, revoked, or key unpaid)",
4282 attempted=True,
4283 )
4284 except Exception as exc:
4285 print(f"[Pipeline] Phase 2 NAME-lane search failed: {exc}", file=sys.stderr)
4286 state, attempted = _classify_source_failure(exc)
4287 bundle.record_failure(
4288 x_slug, state, f"Phase 2 NAME-lane: {exc}", attempted=attempted,
4289 )
4290
4291 raw_items = from_items + about_items + name_items
4292
4293 # Partial coverage is a reportable outcome, not a normal result: a
4294 # report carrying only one side of an entity topic is incomplete, and
4295 # without this it looks indistinguishable from genuinely thin
4296 # discussion.
4297 if _name_lane is not None:
4298 empty = [
4299 label for label, items in
4300 (("by", from_items), ("mention", about_items), ("name", name_items))
4301 if not items
4302 ]
4303 if empty and len(empty) < 3:
4304 # A warning, not a source outcome. record_failure would set the
4305 # X source to PARTIAL, which is outside _STRICT_EXIT_OK_STATES
4306 # and would make wrappers using LAST30DAYS_STRICT_EXIT exit 3 on
4307 # runs that returned perfectly good X coverage. An empty lane is
4308 # common and legitimate: the name lane carries an engagement
4309 # floor and the mention lane is empty for most non-famous
4310 # handles.
4311 bundle.artifacts.setdefault("x_partial_coverage", []).append(
4312 f"X partial coverage: {', '.join(empty)} lane(s) returned "
4313 "nothing; the report may show only one side of this entity."
4314 )
4315
4316 if raw_items:
4317 # First-party handles: only primary explicit handle, not promoted commentators
4318 # (first-party exempts from relevance floor; granting to commentators
4319 # would let junk become un-prunable)
4320 first_party_for_normalize = list(set(
4321 h.lower().lstrip("@") for h in primary_explicit if h
4322 ))
4323 normalized = _normalize_score_dedupe(
4324 x_slug, raw_items, from_date, to_date,
4325 freshness_mode=plan.freshness_mode,
4326 ranking_query=ranking_query,
4327 first_party_handles=first_party_for_normalize,
4328 )
4329 # Deduplicate against Phase 1 URLs
4330 normalized = [item for item in normalized if item.url not in existing_urls]
4331 if normalized:
4332 bundle.add_items(primary_label, x_slug, normalized)
4333 # Update existing URLs for related-handle dedup
4334 for item in normalized:
4335 if item.url:
4336 existing_urls.add(item.url)
4337
4338 # Search related handles with lower weight (0.3)
4339 # Related handles are explicit (--x-related), so FROM without AND topic.
4340 if related_handles:
4341 try:
4342 raw_items, rel_revoked = _from_lane(related_handles, RELATED_HANDLE_COUNT_PER, and_topic=False)
4343 if rel_revoked:
4344 any_revoked = True
4345 bundle.record_failure(
4346 x_slug, schema.AUTH_FAILED,
4347 f"Phase 2 related handle search: {primary} authentication failed (session expired, revoked, or key unpaid)",
4348 attempted=True,
4349 )
4350 except Exception as exc:
4351 print(f"[Pipeline] Phase 2 related handle search failed: {exc}", file=sys.stderr)
4352 state, attempted = _classify_source_failure(exc)
4353 bundle.record_failure(
4354 x_slug,
4355 state,
4356 f"Phase 2 related handle search: {exc}",
4357 attempted=attempted,
4358 )
4359 raw_items = []
4360
4361 if raw_items:
4362 normalized = _normalize_score_dedupe(
4363 x_slug, raw_items, from_date, to_date,
4364 freshness_mode=plan.freshness_mode,
4365 ranking_query=ranking_query,
4366 first_party_handles=related_handles,
4367 )
4368 # Deduplicate against all existing URLs (Phase 1 + primary handles)
4369 normalized = [item for item in normalized if item.url not in existing_urls]
4370 if normalized:
4371 # Use a separate subquery label with lower weight so RRF
4372 # scores related-handle results below primary results.
4373 bundle.add_items("supplemental-related", x_slug, normalized)
4374 # Register the supplemental-related label in the plan for fusion
4375 if not any(sq.label == "supplemental-related" for sq in plan.subqueries):
4376 plan.subqueries.append(
4377 schema.SubQuery(
4378 label="supplemental-related",
4379 search_query=", ".join(related_handles),
4380 ranking_query=ranking_query,
4381 sources=[x_slug],
4382 weight=0.3,
4383 )
4384 )
4385
4386
4387 def _retry_thin_sources(
4388 *,
4389 topic: str,
4390 bundle: schema.RetrievalBundle,
4391 plan: schema.QueryPlan,
4392 config: dict[str, Any],
4393 depth: str,
4394 date_range: tuple[str, str],
4395 runtime: schema.ProviderRuntime,
4396 mock: bool,
4397 rate_limited_sources: set[str],
4398 rate_limit_lock: threading.Lock,
4399 settings: dict[str, Any],
4400 web_backend: str = "auto",
4401 skip_sources: set[str] | None = None,
4402 subreddits: list[str] | None = None,
4403 tiktok_hashtags: list[str] | None = None,
4404 tiktok_creators: list[str] | None = None,
4405 ig_creators: list[str] | None = None,
4406 first_party_handles: Iterable[str] | None = None,
4407 first_party_by_source: Mapping[str, Iterable[str]] | None = None,
4408 run_started: float | None = None,
4409 ) -> None:
4410 """Retry sources with thin results using simplified core subject query."""
4411 if depth == "quick":
4412 return
4413
4414 planned_sources: list[str] = []
4415 for subquery in plan.subqueries:
4416 for source in subquery.sources:
4417 if source not in planned_sources:
4418 planned_sources.append(source)
4419 _skip = (skip_sources or set()) | THIN_RETRY_EXEMPT
4420 thin_sources = [
4421 source
4422 for source in planned_sources
4423 if len(bundle.items_by_source.get(source, [])) < 3
4424 and source not in bundle.errors_by_source
4425 and source not in _skip
4426 ]
4427
4428 if not thin_sources:
4429 return
4430
4431 core = query.extract_core_subject(topic, max_words=3)
4432 if not core:
4433 return
4434 # Note: we intentionally do NOT skip when core == topic. For short topics
4435 # like "Kanye West", the 3-word core IS the topic — but the planner may
4436 # have sent a different (worse) query to the source. Retrying with the
4437 # raw core subject is still valuable.
4438
4439 from_date, to_date = date_range
4440
4441 # Create a retry subquery with the simplified core subject
4442 retry_subquery = schema.SubQuery(
4443 label="retry",
4444 search_query=core,
4445 ranking_query=f"What recent evidence from the last 30 days matters for {core}?",
4446 sources=thin_sources,
4447 weight=0.3,
4448 )
4449
4450 def _retry_one_source(
4451 source: str,
4452 ) -> tuple[str, list[schema.SourceItem], dict[str, Any] | None]:
4453 raw_items, artifact = _retrieve_stream(
4454 topic=topic,
4455 subquery=retry_subquery,
4456 source=source,
4457 config=config,
4458 depth=depth,
4459 date_range=date_range,
4460 runtime=runtime,
4461 mock=mock,
4462 rate_limited_sources=rate_limited_sources,
4463 rate_limit_lock=rate_limit_lock,
4464 web_backend=web_backend,
4465 raw_topic=topic,
4466 subreddits=subreddits,
4467 tiktok_hashtags=tiktok_hashtags,
4468 tiktok_creators=tiktok_creators,
4469 ig_creators=ig_creators,
4470 run_started=run_started,
4471 # Skip Amazon review enrichment here to avoid duplicate Bright Data
4472 # pulls for ASINs already enriched in Phase 1. Finalize will enrich
4473 # any genuinely new products that weren't in Phase 1.
4474 skip_amazon_enrichment=True,
4475 )
4476 outcome_note = artifact.get("_source_outcome") if isinstance(artifact, dict) else None
4477 detail_note = artifact.get("_source_outcome_detail") if isinstance(artifact, dict) else None
4478 detail_state = artifact.get("_source_outcome_detail_state") if isinstance(artifact, dict) else None
4479 normalized = _normalize_score_dedupe(
4480 source,
4481 raw_items,
4482 from_date,
4483 to_date,
4484 freshness_mode=plan.freshness_mode,
4485 ranking_query=retry_subquery.ranking_query,
4486 first_party_handles=first_party_handles,
4487 first_party_by_source=first_party_by_source,
4488 # Match Phase 1: X defers its relevance floor until the run has
4489 # resolved handles. Applying it here would discard a subject-
4490 # authored post that does not repeat the subject's name, and the
4491 # later resolved-handle floor cannot recover a post that never
4492 # entered the bundle.
4493 defer_relevance_prune=(source == "x"),
4494 )
4495 if source == "jobs":
4496 return source, normalized, outcome_note, (detail_note, detail_state)
4497 normalized = _apply_reddit_stream_keepers(
4498 source, normalized, settings["per_stream_limit"], topic
4499 )
4500 return source, normalized, outcome_note, (detail_note, detail_state)
4501
4502 retryable = [s for s in thin_sources if s not in rate_limited_sources]
4503
4504 from concurrent.futures import ThreadPoolExecutor, as_completed
4505 with ThreadPoolExecutor(max_workers=min(4, len(retryable) or 1)) as executor:
4506 futures = {executor.submit(_retry_one_source, s): s for s in retryable}
4507 for future in as_completed(futures):
4508 source = futures[future]
4509 try:
4510 source, normalized, outcome_note, (detail_note, detail_state) = future.result()
4511 if outcome_note:
4512 bundle.record_failure(
4513 source,
4514 outcome_note["state"],
4515 outcome_note["detail"],
4516 attempted=outcome_note.get("attempted", True),
4517 )
4518 if detail_note:
4519 bundle.record_detail(source, detail_note, state=detail_state)
4520 existing_urls = {item.url for item in bundle.items_by_source.get(source, []) if item.url}
4521 new_items = [item for item in normalized if item.url not in existing_urls]
4522
4523 if new_items:
4524 primary_label = plan.subqueries[0].label if plan.subqueries else "primary"
4525 bundle.add_items(primary_label, source, new_items)
4526 except Exception as exc:
4527 print(f"[Pipeline] Retry failed for {source}: {type(exc).__name__}: {exc}", file=sys.stderr)
4528 state, attempted = _classify_source_failure(exc)
4529 bundle.record_failure(
4530 source,
4531 state,
4532 f"Simplified-query retry failed: {exc}",
4533 attempted=attempted,
4534 )
4535
4536
4537 def _fetch_x_backend(backend, query, from_date, to_date, depth, config, warnings=None):
4538 """Fetch X items from a single backend. Returns (items, error_str).
4539
4540 ``warnings``, when given, collects backend receipts that are not
4541 failures (xapi's "window truncated to 7 days" after the recent-search
4542 fallback) so the X branch can surface them as run artifacts.
4543
4544 Backends are tried in priority order by the caller (env.x_backend_chain);
4545 a non-empty error_str signals a hard failure (auth/payment/etc.) so the
4546 caller can fail over to the next backend or surface the error honestly.
4547
4548 For grok, auth_revoked signals mid-run session revocation: the error
4549 string includes "grok session expired" so _classify_source_failure maps
4550 it to AUTH_FAILED with a proper fix hint, distinct from "never signed in".
4551
4552 The ``query`` parameter is the compiled search query - typically
4553 ``raw_topic or topic`` (like Reddit/YouTube), NOT the planner's
4554 ``search_query`` which may contain operator strings like "Rome Italy".
4555 """
4556 if backend == "bird":
4557 result = bird_x.search_x(query, from_date, to_date, depth=depth)
4558 items = bird_x.parse_bird_response(result, query=query)
4559 elif backend == "grok":
4560 result = grok_x.search_x(query, from_date, to_date, depth=depth)
4561 items = result.get("items", []) if isinstance(result, dict) else []
4562 if isinstance(result, dict) and result.get("auth_revoked"):
4563 err = result.get("error") or "grok session expired or was revoked"
4564 return items, f"grok: {err}"
4565 elif backend == "xai":
4566 model = config.get("LAST30DAYS_X_MODEL") or config.get("XAI_MODEL_PIN") or providers.XAI_DEFAULT
4567 result = xai_x.search_x(config["XAI_API_KEY"], model, query, from_date, to_date, depth=depth)
4568 items = xai_x.parse_x_response(result)
4569 elif backend == "xurl":
4570 result = xurl_x.search_x(query, depth=depth)
4571 items = xurl_x.parse_x_response(result, topic=query)
4572 elif backend == "xquik":
4573 result = xquik.search_xquik(query, from_date, to_date, depth=depth, token=env.get_xquik_token(config))
4574 items = xquik.parse_xquik_response(result)
4575 elif backend == "xapi":
4576 result = x_api.search_x(config.get("X_BEARER_TOKEN") or "", query, from_date, to_date, depth=depth)
4577 items = result.get("items", []) if isinstance(result, dict) else []
4578 warning = result.get("warning") if isinstance(result, dict) else None
4579 if warning:
4580 print(f"[X] xapi: {warning}", file=sys.stderr)
4581 if warnings is not None:
4582 warnings.append(f"X: xapi {warning}")
4583 else:
4584 return [], f"unknown X backend: {backend}"
4585 err = result.get("error") if isinstance(result, dict) else ""
4586 return items, (err or "")
4587
4588
4589 def _reddit_post_key(item: dict) -> str:
4590 """Stable per-thread dedupe key (base36 post id from the url/permalink)."""
4591 url = item.get("url") or item.get("permalink") or ""
4592 m = re.search(r"/comments/([A-Za-z0-9]+)", url)
4593 return m.group(1) if m else url
4594
4595
4596 def _merge_reddit_items(free: list[dict], sc: list[dict]) -> list[dict]:
4597 """Merge free + ScrapeCreators Reddit items, free first, deduped by post id.
4598
4599 Used when the thinness-floor trigger backfills a thin free run with SC, so a
4600 thread present in both is never double-listed.
4601 """
4602 merged = list(free)
4603 seen = {_reddit_post_key(it) for it in free}
4604 for it in sc:
4605 key = _reddit_post_key(it)
4606 if key and key not in seen:
4607 seen.add(key)
4608 merged.append(it)
4609 return merged
4610
4611
4612 def _retrieve_stream(*args, **kwargs) -> tuple[list[dict], dict]:
4613 """Run one stream and retain HTTP failures swallowed by source adapters."""
4614 # run_started is passed through but not used here; it goes to _retrieve_stream_impl
4615 source = str(kwargs.get("source") or "")
4616 fixture_request = {
4617 "source": source,
4618 "topic": kwargs.get("topic") or "",
4619 "search_query": getattr(kwargs.get("subquery"), "search_query", ""),
4620 "date_range": list(kwargs.get("date_range") or ()),
4621 "depth": kwargs.get("depth") or "",
4622 }
4623 module_backed = source in {
4624 "reddit",
4625 "x",
4626 "youtube",
4627 "stocktwits",
4628 "digg",
4629 "arxiv",
4630 "techmeme",
4631 "trustpilot",
4632 "github",
4633 }
4634 if module_backed:
4635 matched, replayed = http.fixture_source_replay(fixture_request)
4636 if matched:
4637 return replayed[0], replayed[1]
4638 try:
4639 with http.capture_failures() as failures, \
4640 http.fixture_module_capture(module_backed):
4641 items, artifact = _retrieve_stream_impl(*args, **kwargs)
4642 except Exception as exc:
4643 recorded_exc = exc
4644 if failures and not getattr(exc, "outcome_state", None):
4645 failure = min(
4646 failures,
4647 key=lambda f: _FAILURE_SPECIFICITY.get(f.outcome_state, 9),
4648 )
4649 recorded_exc = SourceRunError(str(exc), failure.outcome_state)
4650 if module_backed:
4651 http.fixture_source_record_error(fixture_request, recorded_exc)
4652 if recorded_exc is not exc:
4653 raise recorded_exc from exc
4654 raise
4655 outcome_note = _resolve_stream_outcome(
4656 str(kwargs.get("source") or ""),
4657 artifact,
4658 failures,
4659 )
4660 if outcome_note:
4661 # Lane-level HTTP failures (e.g. a blocked shreddit partial on a
4662 # datacenter IP) are captured by the sink even when the source
4663 # delivered items. Only attach them when the run produced nothing,
4664 # or when the impl attached its own explicit outcome artifact (e.g.
4665 # "primary failed; fallback returned N items"). A swallowed lane
4666 # failure must not brand a successful source auth-failed/partial.
4667 # An adapter-declared outcome (typed ``_source_outcome`` or a legacy
4668 # ``{"error": ...}`` / per-leg artifact) is explicit and always
4669 # brands the source, even with items; only failures the adapter
4670 # swallowed into the capture sink are demoted to detail.
4671 explicit = isinstance(artifact, dict) and (
4672 bool(artifact.get("_source_outcome"))
4673 or _legacy_artifact_outcome(str(kwargs.get("source") or ""), artifact) is not None
4674 )
4675 if explicit or not items:
4676 artifact = dict(artifact or {})
4677 artifact["_source_outcome"] = outcome_note
4678 elif failures:
4679 # The source delivered items. Keep it ``ok`` but carry what the
4680 # swallowed sub-requests lost, so doctor can still show it, and
4681 # the most specific failure state so a later empty filter result
4682 # or the thin-source retry can act on it.
4683 # Append to any detail the impl already set (e.g. a Reddit
4684 # backfill note) rather than overwriting it.
4685 artifact = dict(artifact or {})
4686 summary = _summarize_lane_failures(failures, str(kwargs.get("source") or ""))
4687 prior = artifact.get("_source_outcome_detail")
4688 artifact["_source_outcome_detail"] = f"{prior}; {summary}" if prior else summary
4689 states = [f.outcome_state for f in failures]
4690 if artifact.get("_source_outcome_detail_state"):
4691 states.append(artifact["_source_outcome_detail_state"])
4692 artifact["_source_outcome_detail_state"] = min(
4693 states, key=lambda state: _FAILURE_SPECIFICITY.get(state, 9)
4694 )
4695 if module_backed:
4696 http.fixture_source_record(fixture_request, [items, artifact])
4697 return items, artifact
4698
4699
4700 def _serve_envelope_topic(envelope: x_envelope.Envelope) -> tuple[list[dict], dict]:
4701 """Serve the envelope's topic-lane rows once.
4702
4703 The first X subquery takes the rows and the envelope-status outcome;
4704 every later call (a second planner subquery, judge-retry, thin-retry)
4705 gets no items and no error, and no backend is ever consulted.
4706 """
4707 items = envelope.take_topic()
4708 if items is None:
4709 return [], {}
4710 artifact: dict[str, Any] = {}
4711 if envelope.warnings:
4712 # A narrower host window is a receipt (report.warnings), not a failure.
4713 artifact["x_receipts"] = [f"X: {warning}" for warning in envelope.warnings]
4714 outcome = envelope.outcome()
4715 if outcome is not None:
4716 state, detail = outcome
4717 artifact.update(_outcome_artifact(state, detail))
4718 return items, artifact
4719
4720
4721 def _serve_envelope_lanes(
4722 envelope: x_envelope.Envelope,
4723 *,
4724 bundle: schema.RetrievalBundle,
4725 plan: schema.QueryPlan,
4726 x_handle: str | None,
4727 x_related: list[str] | None,
4728 from_date: str,
4729 to_date: str,
4730 ) -> None:
4731 """Serve the envelope's from/mention/related calls into the lane merge.
4732
4733 Mirrors the backend lanes: primary-handle rows (from + mention) join the
4734 primary subquery with first-party handling for the explicit handle and
4735 the per-handle lane counts; related rows join ``supplemental-related``
4736 at the 0.3 weight. Lane claims were already validated at read time.
4737 """
4738 calls = envelope.take_lanes()
4739 if not calls:
4740 return
4741 x_slug = "x"
4742 existing_urls = {
4743 item.url
4744 for items in bundle.items_by_source.values()
4745 for item in items
4746 if item.url
4747 }
4748 ranking_query = plan.subqueries[0].ranking_query if plan.subqueries else ""
4749 primary_label = plan.subqueries[0].label if plan.subqueries else "primary"
4750 primary_handles = sorted(
4751 {x_handle.lstrip("@").strip().lower()} if x_handle and x_handle.strip() else set()
4752 )
4753 related_handles = [
4754 h.lstrip("@").strip().lower()
4755 for h in (x_related or [])
4756 if h.strip() and h.lstrip("@").strip().lower() not in primary_handles
4757 ]
4758
4759 def _cap_per_author(posts: list[dict], cap: int) -> list[dict]:
4760 seen: Counter[str] = Counter()
4761 kept: list[dict] = []
4762 for post in posts:
4763 author = str(post.get("author_handle") or "").lower()
4764 if seen[author] >= cap:
4765 continue
4766 seen[author] += 1
4767 kept.append(post)
4768 return kept
4769
4770 primary_items: list[dict] = []
4771 related_items: list[dict] = []
4772 for call in calls:
4773 if call.lane == "from":
4774 primary_items.extend(_cap_per_author(call.posts, FROM_LANE_COUNT_PER))
4775 elif call.lane == "mention":
4776 primary_items.extend(
4777 call.posts[: MENTION_LANE_COUNT_PER * max(1, len(call.handles))]
4778 )
4779 elif call.lane == "related":
4780 related_items.extend(_cap_per_author(call.posts, RELATED_HANDLE_COUNT_PER))
4781
4782 if primary_items:
4783 normalized = _normalize_score_dedupe(
4784 x_slug, primary_items, from_date, to_date,
4785 freshness_mode=plan.freshness_mode,
4786 ranking_query=ranking_query,
4787 first_party_handles=primary_handles,
4788 )
4789 normalized = [item for item in normalized if item.url not in existing_urls]
4790 if normalized:
4791 bundle.add_items(primary_label, x_slug, normalized)
4792 existing_urls.update(item.url for item in normalized if item.url)
4793
4794 if related_items:
4795 normalized = _normalize_score_dedupe(
4796 x_slug, related_items, from_date, to_date,
4797 freshness_mode=plan.freshness_mode,
4798 ranking_query=ranking_query,
4799 first_party_handles=related_handles,
4800 )
4801 normalized = [item for item in normalized if item.url not in existing_urls]
4802 if normalized:
4803 bundle.add_items("supplemental-related", x_slug, normalized)
4804 if not any(sq.label == "supplemental-related" for sq in plan.subqueries):
4805 plan.subqueries.append(
4806 schema.SubQuery(
4807 label="supplemental-related",
4808 search_query=", ".join(related_handles),
4809 ranking_query=ranking_query,
4810 sources=[x_slug],
4811 weight=0.3,
4812 )
4813 )
4814
4815
4816 def _retrieve_stream_impl(
4817 *,
4818 topic: str,
4819 subquery: schema.SubQuery,
4820 source: str,
4821 config: dict[str, Any],
4822 depth: str,
4823 date_range: tuple[str, str],
4824 runtime: schema.ProviderRuntime,
4825 mock: bool,
4826 rate_limited_sources: set[str] | None = None,
4827 rate_limit_lock: threading.Lock | None = None,
4828 web_backend: str = "auto",
4829 raw_topic: str = "",
4830 subreddits: list[str] | None = None,
4831 tiktok_hashtags: list[str] | None = None,
4832 tiktok_creators: list[str] | None = None,
4833 ig_creators: list[str] | None = None,
4834 trustpilot_domain: str | None = None,
4835 trustpilot_domain_is_hint: bool = False,
4836 run_started: float | None = None,
4837 skip_amazon_enrichment: bool = False,
4838 ) -> tuple[list[dict], dict]:
4839 # Early exit if source was rate-limited by a sibling future
4840 if rate_limited_sources is not None and source in rate_limited_sources:
4841 return [], {}
4842 from_date, to_date = date_range
4843 if mock:
4844 return _mock_stream_results(source, subquery)
4845 if source == "grounding":
4846 return grounding.web_search(
4847 subquery.search_query, date_range, config, backend=web_backend)
4848 if source == "jobs":
4849 return jobs.search_jobs(
4850 raw_topic or topic or subquery.search_query,
4851 date_range,
4852 config,
4853 depth=depth,
4854 web_backend=web_backend,
4855 explicit=bool(config.get("_hiring_signals_mode")),
4856 )
4857 if source == "reddit":
4858 # Use raw_topic so expand_reddit_queries() generates diverse variants
4859 # from the original user topic, not the planner's narrowed search_query.
4860 reddit_query = raw_topic or subquery.search_query
4861 dedicated_subreddits = config.get("_dedicated_subreddits") or None
4862 has_sc_key = bool(config.get("SCRAPECREATORS_API_KEY"))
4863 sc_first = (
4864 has_sc_key
4865 and (config.get(env.REDDIT_BACKEND_PIN_VAR) or "").lower()
4866 == "scrapecreators"
4867 )
4868 if sc_first:
4869 # env.REDDIT_BACKEND_PIN_VAR=scrapecreators: SC primary, public fallback
4870 primary_failure: Exception | None = None
4871 try:
4872 result = reddit.search_and_enrich_memo(
4873 reddit_query, from_date, to_date, depth=depth,
4874 token=config.get("SCRAPECREATORS_API_KEY"),
4875 subreddits=subreddits,
4876 )
4877 items = reddit.parse_reddit_response(result)
4878 if items:
4879 return items, {}
4880 sys.stderr.write(
4881 "[Reddit] ScrapeCreators primary returned no items, "
4882 "using public fallback\n"
4883 )
4884 except Exception as exc:
4885 primary_failure = exc
4886 sys.stderr.write(
4887 f"[Reddit] ScrapeCreators primary failed "
4888 f"({type(exc).__name__}: {exc}), using public fallback\n"
4889 )
4890 public_failure: Exception | None = None
4891 try:
4892 public_results = reddit_public.search_reddit_public(
4893 reddit_query, from_date, to_date, depth=depth,
4894 subreddits=subreddits,
4895 )
4896 if public_results:
4897 if primary_failure is not None:
4898 state = reddit.classify_run_failure(str(primary_failure))
4899 return public_results, _outcome_artifact(
4900 state,
4901 f"Reddit primary failed; public fallback returned "
4902 f"{len(public_results)} items: {primary_failure}",
4903 )
4904 return public_results, {}
4905 sys.stderr.write(
4906 "[Reddit] Public fallback returned no items after "
4907 "ScrapeCreators primary miss\n"
4908 )
4909 except Exception as exc:
4910 public_failure = exc
4911 sys.stderr.write(
4912 f"[Reddit] Public fallback also failed "
4913 f"({type(exc).__name__}: {exc})\n"
4914 )
4915 failure = public_failure or primary_failure
4916 if failure is not None:
4917 state = reddit.classify_run_failure(str(failure))
4918 raise SourceRunError(
4919 f"Reddit primary and fallback produced no results after failure: {failure}",
4920 state,
4921 )
4922 return [], {}
4923
4924 # Default: public Reddit first (free). ScrapeCreators backfills when the
4925 # free path returns fewer than the thinness floor (env.reddit_sc_min_items:
4926 # unset -> env.REDDIT_SC_MIN_ITEMS_DEFAULT, explicit 0 -> empty-only,
4927 # malformed -> 0 so a typo never spends extra credits).
4928 min_items = env.reddit_sc_min_items(config)
4929 public_results: list[dict] = []
4930 public_failure: Exception | None = None
4931 try:
4932 public_results = reddit_public.search_reddit_public(
4933 reddit_query, from_date, to_date, depth=depth,
4934 subreddits=subreddits, dedicated_subreddits=dedicated_subreddits,
4935 ) or []
4936 except Exception as exc:
4937 public_failure = exc
4938 sys.stderr.write(
4939 f"[Reddit] Public search failed ({type(exc).__name__}: {exc})"
4940 )
4941 if not has_sc_key:
4942 sys.stderr.write("\n")
4943 state = reddit.classify_run_failure(str(exc))
4944 raise SourceRunError(f"Reddit public search failed: {exc}", state) from exc
4945 sys.stderr.write(", using ScrapeCreators backup\n")
4946 # Enough free results, or no key to backfill with -> done. max(min_items,
4947 # 1) keeps the default (min_items=0) as empty-only AND treats exactly
4948 # `min_items` results as acceptable (no backfill) for min_items > 0.
4949 if len(public_results) >= max(min_items, 1) or not has_sc_key:
4950 return public_results, {}
4951 if public_results:
4952 sys.stderr.write(
4953 f"[Reddit] Free path returned {len(public_results)} "
4954 f"(below the {min_items}-item floor); backfilling with ScrapeCreators\n"
4955 )
4956 # With free items in hand, scope the backfill's own HTTP failures
4957 # away from the stream's capture sink: a ScrapeCreators 429 is not a
4958 # reddit.com rate limit, and a RATE_LIMITED lane state would put
4959 # 'reddit' in rate_limited_sources and skip later Reddit streams and
4960 # the thin-source retry. Its failures become the detail note instead.
4961 # With no free items the backfill IS the source, so failures reach
4962 # the stream sink as before.
4963 backfill_scope = (
4964 http.capture_failures() if public_results else contextlib.nullcontext([])
4965 )
4966 try:
4967 with backfill_scope as backfill_failures:
4968 result = reddit.search_and_enrich_memo(
4969 reddit_query, from_date, to_date, depth=depth,
4970 token=config.get("SCRAPECREATORS_API_KEY"),
4971 subreddits=subreddits,
4972 )
4973 sc_items = reddit.parse_reddit_response(result)
4974 except Exception as exc:
4975 sys.stderr.write(
4976 f"[Reddit] ScrapeCreators backup also failed "
4977 f"({type(exc).__name__}: {exc})\n"
4978 )
4979 state = reddit.classify_run_failure(str(exc))
4980 if public_results:
4981 # The free path delivered: Reddit is working. The failed
4982 # backfill is a detail note, never the source's outcome.
4983 artifact = {
4984 "_source_outcome_detail": (
4985 f"ScrapeCreators backfill failed after "
4986 f"{len(public_results)} free items: {exc}"
4987 ),
4988 }
4989 if state != health.RATE_LIMITED:
4990 # A ScrapeCreators rate limit says nothing about
4991 # reddit.com (and the run memo already stops a retry).
4992 artifact["_source_outcome_detail_state"] = state
4993 return public_results, artifact
4994 return public_results, _outcome_artifact(
4995 state,
4996 f"Reddit backup failed after {len(public_results)} public items: {exc}",
4997 )
4998 merged = _merge_reddit_items(public_results, sc_items)
4999 if public_failure is not None:
5000 state = reddit.classify_run_failure(str(public_failure))
5001 return merged, _outcome_artifact(
5002 state,
5003 f"Reddit public search failed; backup returned {len(sc_items)} items: "
5004 f"{public_failure}",
5005 )
5006 trigger = (
5007 f"below the {min_items}-item floor" if min_items > 0 else "free path empty"
5008 )
5009 backfill_detail = (
5010 f"ScrapeCreators backfill ran ({len(public_results)} free items, "
5011 f"{trigger}); added {len(merged) - len(public_results)} items"
5012 )
5013 if backfill_failures:
5014 summary = _summarize_lane_failures(backfill_failures, "reddit")
5015 backfill_detail = f"{backfill_detail}; ScrapeCreators backfill lost {summary}"
5016 return merged, {"_source_outcome_detail": backfill_detail}
5017 if source == "x":
5018 if config.get("_x_lane_missing"):
5019 # The model declared the connector lane but passed no envelope.
5020 return [], _outcome_artifact(health.ERROR, x_envelope.DETAIL_NOT_PASSED)
5021 envelope = config.get("_x_envelope")
5022 if envelope is not None:
5023 # Host-fetched lane: the envelope replaces the backend chain and
5024 # is single-serve, so no backend runs and no judge-retry follows.
5025 return _serve_envelope_topic(envelope)
5026
5027 # Compile X query from raw_topic (like Reddit/YouTube), not planner's
5028 # search_query which may contain operator strings like "Rome Italy".
5029 x_query = raw_topic or topic or subquery.search_query
5030 ranking_query = subquery.ranking_query
5031
5032 # One X source, an ordered chain of interchangeable backends. Try the
5033 # primary; fall through to the next only if it returns nothing or errors.
5034 chain = env.x_backend_chain(config)
5035 # Trust an explicit runtime backend as the primary (already resolved as
5036 # available), keeping the rest of the chain as failover backups.
5037 pinned = runtime.x_search_backend
5038 if pinned:
5039 chain = [pinned] + [b for b in chain if b != pinned]
5040 if not chain:
5041 raise SourceRunError(
5042 "No X backend is available (not configured).",
5043 schema.SKIPPED_UNCONFIGURED,
5044 )
5045 last_error = ""
5046 chain_errors: list[str] = []
5047 items = []
5048 used_backend = None
5049 x_warnings: list[str] = []
5050 for i, backend in enumerate(chain):
5051 items, err = _fetch_x_backend(
5052 backend, x_query, from_date, to_date, depth, config, warnings=x_warnings,
5053 )
5054 if items:
5055 if i > 0:
5056 # xapi is metered: name the spend when it served as a backup.
5057 spend = " (spends X API credits)" if backend == "xapi" else ""
5058 print(
5059 f"[X] primary backend(s) returned nothing; used fallback '{backend}'{spend}",
5060 file=sys.stderr,
5061 )
5062 # Check for auth errors before proceeding to judge-retry
5063 if last_error:
5064 # Fallback succeeded after earlier backend failed. Classify
5065 # the original error: if it was AUTH_FAILED (grok revoked),
5066 # preserve that state so user gets re-login guidance.
5067 prior_state = http.classify_failure(message=last_error)
5068 if prior_state == schema.AUTH_FAILED:
5069 # Keep AUTH_FAILED visible so host shows re-login hint
5070 return items, _outcome_artifact(
5071 schema.AUTH_FAILED,
5072 f"X served via {backend} after {last_error}; re-login needed for primary backend",
5073 )
5074 # Prior error was non-auth. Check if *current* backend also
5075 # reported an error (e.g., grok returned items + revocation).
5076 if err:
5077 current_state = http.classify_failure(message=err)
5078 if current_state == schema.AUTH_FAILED:
5079 return items, _outcome_artifact(
5080 schema.AUTH_FAILED,
5081 f"X served {len(items)} items via {backend} but also errored: {err}; re-login needed",
5082 )
5083 # Non-auth prior error, no current auth error → fallback OK
5084 return items, _outcome_artifact(
5085 health.OK,
5086 f"X served via {backend} after {last_error}",
5087 )
5088 if err:
5089 # Mixed result: backend returned items BUT also hit an error
5090 # (e.g., grok got some posts then auth was revoked mid-fanout).
5091 # Surface the error so the user gets re-login guidance.
5092 state = http.classify_failure(message=err)
5093 return items, _outcome_artifact(
5094 state,
5095 f"X returned {len(items)} items but also errored: {err}",
5096 )
5097 # No auth issues and no prior errors - proceed to judge-retry
5098 used_backend = backend
5099 break
5100 if err:
5101 last_error = f"{backend}: {err}"
5102 chain_errors.append(last_error)
5103 print(f"[X] backend '{backend}' failed ({err}); trying next", file=sys.stderr)
5104
5105 if not items and last_error:
5106 # A credit-exhaustion failure earlier in the chain is the most
5107 # specific outcome (top up, not re-authenticate); a later
5108 # backend's generic failure must not mask it.
5109 for candidate in chain_errors:
5110 if http.classify_failure(message=candidate) == health.PAYMENT_REQUIRED:
5111 last_error = candidate
5112 break
5113 state = (
5114 bird_x.classify_run_failure(last_error)
5115 if last_error.startswith("bird:")
5116 else http.classify_failure(message=last_error)
5117 )
5118 raise SourceRunError(f"All X backends failed — {last_error}", state)
5119
5120 # Retrieve-judge-retry: judge corpus and retry if off-topic flood.
5121 # Skip retry on quick/mock (same as Phase 2).
5122 artifact = {}
5123 if x_warnings:
5124 # e.g. xapi's "window truncated to 7 days": a receipt that reaches
5125 # report.warnings (see the grounding artifacts walk in
5126 # _build_report), never a source failure.
5127 artifact["x_receipts"] = list(x_warnings)
5128 if items and depth != "quick" and not mock:
5129 items_for_judge = [
5130 {"author_handle": it.get("author_handle", ""), "text": it.get("text", "")}
5131 for it in items
5132 ]
5133 if x_judge.should_retry_x_search(items_for_judge, x_query, ranking_query=ranking_query, depth=depth):
5134 # Retry with cleaned query (1 retry, ≤2 extra grok calls)
5135 # Strip noise words but preserve all significant terms to avoid
5136 # losing disambiguating terms (e.g., "react server components")
5137 core_tokens = query.extract_core_subject(x_query)
5138 retry_query = core_tokens or x_query
5139 print(f"[X] corpus off-topic; retrying with '{retry_query}'", file=sys.stderr)
5140
5141 if used_backend:
5142 retry_items, retry_err = _fetch_x_backend(
5143 used_backend, retry_query, from_date, to_date, depth, config
5144 )
5145 if retry_items:
5146 # Judge retry corpus
5147 retry_for_judge = [
5148 {"author_handle": it.get("author_handle", ""), "text": it.get("text", "")}
5149 for it in retry_items
5150 ]
5151 retry_judgment = x_judge.judge_x_corpus(
5152 retry_for_judge, x_query, ranking_query=ranking_query
5153 )
5154 orig_judgment = x_judge.judge_x_corpus(
5155 items_for_judge, x_query, ranking_query=ranking_query
5156 )
5157 # Use retry if better on-topic ratio
5158 if retry_judgment["on_topic_ratio"] > orig_judgment["on_topic_ratio"]:
5159 print(
5160 f"[X] retry improved on-topic ratio: "
5161 f"{orig_judgment['on_topic_ratio']:.0%} -> "
5162 f"{retry_judgment['on_topic_ratio']:.0%}",
5163 file=sys.stderr,
5164 )
5165 items = retry_items
5166
5167 # Prune off-topic items before the pool. Eight on-topic → ok with 8.
5168 # Zero on-topic after retry → no-results, not ok with 40 junk.
5169 # Only prune items that have text to judge; items without text pass through.
5170 original_count = len(items)
5171 items_with_text = [(i, it) for i, it in enumerate(items) if it.get("text", "").strip()]
5172
5173 if items_with_text:
5174 items_for_prune = [
5175 {"author_handle": it.get("author_handle", ""), "text": it.get("text", "")}
5176 for _, it in items_with_text
5177 ]
5178 judgment = x_judge.judge_x_corpus(
5179 items_for_prune, x_query, ranking_query=ranking_query
5180 )
5181 # Build set of indices for on-topic items
5182 on_topic_indices = set()
5183 for (orig_idx, _), pruned_item in zip(items_with_text, items_for_prune):
5184 if pruned_item in judgment["on_topic_items"]:
5185 on_topic_indices.add(orig_idx)
5186
5187 # Keep items that are on-topic OR have no text (can't judge)
5188 items = [
5189 it for i, it in enumerate(items)
5190 if i in on_topic_indices or not it.get("text", "").strip()
5191 ]
5192
5193 # Record warning if significant pruning occurred (artifact, not failure)
5194 if len(items) < original_count:
5195 pruned = original_count - len(items)
5196 artifact.setdefault("_warnings", []).append(
5197 f"X: pruned {pruned} off-topic items; {len(items)} on-topic remain"
5198 )
5199
5200 if last_error and items:
5201 state = (
5202 bird_x.classify_run_failure(last_error)
5203 if last_error.startswith("bird:")
5204 else http.classify_failure(message=last_error)
5205 )
5206 return items, _outcome_artifact(
5207 state,
5208 f"X fallback '{used_backend}' returned {len(items)} items after {last_error}",
5209 )
5210 return items, artifact
5211 if source == "youtube":
5212 # Use raw_topic so expand_youtube_queries() generates diverse variants
5213 # from the original user topic, not the planner's narrowed search_query.
5214 yt_query = raw_topic or subquery.search_query
5215 result = None
5216 youtube_failure: str | None = None
5217 # ScrapeCreators key (when present) is the default-on backup tier: it
5218 # powers the per-video transcript fallback, the SC search fallback, and
5219 # comment enrichment. None when no key, which keeps everything keyless.
5220 sc_token = (
5221 config.get("SCRAPECREATORS_API_KEY", "")
5222 if env.is_youtube_sc_available(config) else None
5223 )
5224 # Try yt-dlp first; the SC transcript fallback covers per-video failures.
5225 if which("yt-dlp"):
5226 try:
5227 result = youtube_yt.search_and_transcribe(
5228 yt_query, from_date, to_date, depth=depth, token=sc_token,
5229 )
5230 if result.get("error"):
5231 youtube_failure = str(result["error"])
5232 except Exception as exc:
5233 youtube_failure = str(exc)
5234 result = None
5235 # Fall back to SC YouTube search if yt-dlp failed or isn't installed.
5236 if (result is None or not result.get("items")) and sc_token:
5237 try:
5238 result = youtube_yt.search_youtube_sc(
5239 yt_query, from_date, to_date, depth=depth, token=sc_token,
5240 )
5241 if result.get("error"):
5242 youtube_failure = str(result["error"])
5243 except Exception as exc:
5244 youtube_failure = str(exc)
5245 result = None
5246 if result is None:
5247 result = {"items": []}
5248 # Enrich top videos with comments (default-on when a key is present).
5249 items = youtube_yt.parse_youtube_response(result)
5250 if items and env.is_youtube_comments_available(config):
5251 youtube_yt.enrich_with_comments(
5252 items, token=config.get("SCRAPECREATORS_API_KEY", ""),
5253 )
5254 if youtube_failure:
5255 state = youtube_yt.classify_run_failure(youtube_failure)
5256 attempted = state != schema.SKIPPED_UNCONFIGURED
5257 return items, _outcome_artifact(state, youtube_failure, attempted=attempted)
5258 return items, {}
5259 if source == "tiktok":
5260 # Use raw_topic so expand_tiktok_queries() generates diverse variants
5261 # from the original user topic, not the planner's narrowed search_query.
5262 tiktok_query = raw_topic or subquery.search_query
5263 result = tiktok.search_and_enrich(
5264 tiktok_query,
5265 from_date,
5266 to_date,
5267 depth=depth,
5268 token=env.get_tiktok_token(config),
5269 hashtags=tiktok_hashtags,
5270 creators=tiktok_creators,
5271 )
5272 items = tiktok.parse_tiktok_response(result)
5273 if items and env.is_tiktok_comments_available(config):
5274 sc_token = config.get("SCRAPECREATORS_API_KEY", "")
5275 tiktok.enrich_with_comments(items, token=sc_token)
5276 return items, _result_outcome_artifact(source, result)
5277 if source == "instagram":
5278 # Use raw_topic so expand_instagram_queries() generates diverse variants
5279 # from the original user topic, not the planner's narrowed search_query.
5280 ig_query = raw_topic or subquery.search_query
5281 result = instagram.search_and_enrich(
5282 ig_query,
5283 from_date,
5284 to_date,
5285 depth=depth,
5286 token=env.get_instagram_token(config),
5287 ig_creators=ig_creators,
5288 )
5289 items = instagram.parse_instagram_response(result)
5290 if items and env.is_instagram_comments_available(config):
5291 instagram.enrich_with_comments(
5292 items, token=config.get("SCRAPECREATORS_API_KEY", ""),
5293 )
5294 return items, _result_outcome_artifact(source, result)
5295 if source == "linkedin":
5296 token = config.get("SCRAPECREATORS_API_KEY", "")
5297 result = linkedin.search_linkedin(
5298 subquery.search_query,
5299 from_date,
5300 to_date,
5301 depth=depth,
5302 token=token,
5303 )
5304 items = linkedin.parse_linkedin_response(
5305 result, from_date=from_date, to_date=to_date
5306 )
5307 # Articles never appear in post search — surface them (high signal)
5308 # via a bounded profile-enrichment lane on person topics.
5309 items += linkedin.enrich_articles(
5310 items, raw_topic or topic, token, from_date=from_date, to_date=to_date
5311 )
5312 return items, _result_outcome_artifact(source, result)
5313 if source == "hackernews":
5314 result = hackernews.search_hackernews(subquery.search_query, from_date, to_date, depth=depth)
5315 return (
5316 hackernews.parse_hackernews_response(result, query=subquery.search_query),
5317 _result_outcome_artifact(source, result),
5318 )
5319 if source == "stocktwits":
5320 # Pass raw_topic so symbol detection sees the full topic, not the
5321 # narrowed per-subquery search_query (same rationale as reddit).
5322 result = stocktwits.search_stocktwits(
5323 raw_topic or topic or subquery.search_query, from_date, to_date, depth=depth)
5324 return (
5325 stocktwits.parse_stocktwits_response(result, query=subquery.search_query),
5326 _result_outcome_artifact(source, result),
5327 )
5328 if source == "dripstack":
5329 result = dripstack.search_dripstack(
5330 subquery.search_query, from_date, to_date, depth=depth)
5331 relevance_topic = raw_topic or topic or subquery.search_query
5332 return (
5333 dripstack.parse_dripstack_response(result, query=relevance_topic),
5334 _result_outcome_artifact(source, result),
5335 )
5336 if source == "digg":
5337 result = digg.search_digg(subquery.search_query, from_date, to_date, depth=depth)
5338 items = digg.parse_digg_response(result, query=subquery.search_query)
5339 # Enrichment with attached X posts is deferred to
5340 # _finalize_items_by_source so it runs on the items that actually
5341 # survive dedupe rather than on top-K of the raw fanout.
5342 return items, _result_outcome_artifact(source, result)
5343 if source == "arxiv":
5344 result = arxiv.search_arxiv(subquery.search_query, from_date, to_date, depth=depth)
5345 # Relevance keys off the stable research topic, not the per-subquery
5346 # search_query, so off-topic narrowing does not let weak matches through.
5347 relevance_topic = raw_topic or topic or subquery.search_query
5348 return (
5349 arxiv.parse_arxiv_response(result, query=relevance_topic),
5350 _result_outcome_artifact(source, result),
5351 )
5352 if source == "techmeme":
5353 result = techmeme.search_techmeme(subquery.search_query, from_date, to_date, depth=depth)
5354 relevance_topic = raw_topic or topic or subquery.search_query
5355 return (
5356 techmeme.parse_techmeme_response(result, query=relevance_topic),
5357 _result_outcome_artifact(source, result),
5358 )
5359 if source == "trustpilot":
5360 # Brand-shape gate keys off the stable research topic, not the narrowed
5361 # per-subquery search_query, so the company is detected consistently.
5362 relevance_topic = raw_topic or topic or subquery.search_query
5363 result = trustpilot.search_trustpilot(
5364 relevance_topic, from_date, to_date, depth=depth, config=config,
5365 explicit_domain=trustpilot_domain,
5366 domain_is_hint=trustpilot_domain_is_hint,
5367 )
5368 return (
5369 trustpilot.parse_trustpilot_response(result, query=relevance_topic),
5370 _result_outcome_artifact(source, result),
5371 )
5372 if source == "amazon":
5373 # The search keyword is model-supplied and may differ from the topic
5374 # ("Matt Van Horn" searches "June Oven"), so it keys off the stable
5375 # research topic rather than the narrowed per-subquery search_query.
5376 keyword = (
5377 str((config or {}).get("_amazon_query") or "").strip()
5378 or raw_topic or topic or subquery.search_query
5379 )
5380 domain = str((config or {}).get("LAST30DAYS_AMAZON_DOMAIN") or amazon.DEFAULT_DOMAIN)
5381 result = amazon.search_products(keyword, domain=domain, config=config)
5382 products = amazon.parse_search_response(result, keyword, domain=domain)
5383 artifact = _result_outcome_artifact(source, result)
5384
5385 # Skip enrichment when called from thin retry (_retry_thin_sources) to
5386 # avoid duplicate Bright Data pulls for ASINs already enriched in Phase 1.
5387 # Finalize will enrich any NEW products (enrich_source_items no-ops when
5388 # top_comments is already set, so duplicates get skipped there too).
5389 if skip_amazon_enrichment:
5390 return products, artifact
5391
5392 # Start review enrichment now, while other sources are still running.
5393 # Elapsed is measured from run_started so multi-source runs that finish
5394 # search quickly (30-90s) still have 190-250s of budget (clamped to 180).
5395 # This replaces the old deferred-to-finalize path which left only crumbs
5396 # (e.g. 11s) after long retrieval phases.
5397 elapsed = time.monotonic() - run_started if run_started else 0.0
5398 enriched, review_status = amazon.enrich_with_reviews(
5399 products,
5400 depth=depth,
5401 config=config,
5402 elapsed=elapsed,
5403 keyword=keyword,
5404 )
5405
5406 # Record PARTIAL status if review lane was skipped or all pulls dropped
5407 if review_status:
5408 artifact = artifact or {}
5409 artifact = dict(artifact) if artifact else {}
5410 artifact["_source_outcome"] = {
5411 "state": schema.PARTIAL,
5412 "detail": review_status,
5413 "attempted": True,
5414 }
5415
5416 return enriched, artifact
5417 if source == "meta_ads":
5418 # The advertiser is resolved from the stable research topic, not the
5419 # narrowed per-subquery search_query: a subquery like "kettle reviews"
5420 # would resolve a different page than the brand the run is about.
5421 brand = raw_topic or topic or subquery.search_query
5422 result = meta_ads.search_meta_ads(
5423 brand,
5424 from_date,
5425 to_date,
5426 depth=depth,
5427 token=(config or {}).get("SCRAPECREATORS_API_KEY") or "",
5428 country=str(
5429 (config or {}).get("LAST30DAYS_META_ADS_COUNTRY")
5430 or meta_ads.DEFAULT_COUNTRY
5431 ),
5432 page_override=str((config or {}).get("_meta_ads_page") or "").strip(),
5433 )
5434 if result.get("partial"):
5435 # A partial lane carries `error` too, so the generic classifier
5436 # would run and have its verdict overwritten here regardless.
5437 artifact = {
5438 "_source_outcome": {
5439 "state": schema.PARTIAL,
5440 "detail": str(result.get("error") or "partial"),
5441 "attempted": True,
5442 }
5443 }
5444 else:
5445 artifact = dict(_result_outcome_artifact(source, result) or {})
5446 # The footer needs the resolved advertiser and the pre-truncation
5447 # counts even on a run that produced zero items, and stream artifacts
5448 # only reach the report through the grounding list, so they ride here
5449 # and are lifted to top-level artifacts after retrieval.
5450 artifact["meta_ads_page"] = result.get("page") or {}
5451 artifact["meta_ads_tally"] = result.get("tally") or {}
5452 return result.get("ads") or [], artifact
5453 if source == "bluesky":
5454 result = bluesky.search_bluesky(subquery.search_query, from_date, to_date, depth=depth, config=config)
5455 return bluesky.parse_bluesky_response(result), _result_outcome_artifact(source, result)
5456 if source == "threads":
5457 result = threads.search_threads(
5458 subquery.search_query, from_date, to_date,
5459 depth=depth,
5460 token=config.get("SCRAPECREATORS_API_KEY"),
5461 )
5462 return threads.parse_threads_response(result), _result_outcome_artifact(source, result)
5463 if source == "telegram":
5464 result = telegram.search_telegram(
5465 subquery.search_query, from_date, to_date,
5466 depth=depth,
5467 token=config.get("SCRAPECREATORS_API_KEY"),
5468 config=config,
5469 )
5470 return telegram.parse_telegram_response(result), _result_outcome_artifact(source, result)
5471 if source == "truthsocial":
5472 result = truthsocial.search_truthsocial(subquery.search_query, from_date, to_date, depth=depth, config=config)
5473 return truthsocial.parse_truthsocial_response(result), _result_outcome_artifact(source, result)
5474 if source == "polymarket":
5475 result = polymarket.search_polymarket(subquery.search_query, from_date, to_date, depth=depth)
5476 # Relevance filtering keys off the stable original research topic, not the
5477 # per-subquery search_query (which narrows differently on each fanout pass
5478 # and would let off-topic markets through on broad subqueries while dropping
5479 # everything on narrow ones).
5480 relevance_topic = raw_topic or topic or subquery.search_query
5481 return (
5482 polymarket.parse_polymarket_response(result, topic=relevance_topic),
5483 _result_outcome_artifact(source, result),
5484 )
5485 if source == "github":
5486 # Resolve once at the pipeline boundary so search and enrich
5487 # share the result; otherwise each call would re-run the env
5488 # lookup and gh-CLI subprocess fallback (up to 5s timeout each).
5489 token = github.resolve_token(config.get("GITHUB_TOKEN"))
5490 response = github.search_github(subquery.search_query, from_date, to_date, depth=depth, token=token)
5491 items = github.parse_github_response(response)
5492 # Note: an unauth rate-limit (response["error"]) is expected on the
5493 # tokenless anon tier and returns empty here rather than raising — github
5494 # is now always eligible, so raising would spam "github failed" on every
5495 # tokenless run. The condition is logged in github.search_github.
5496 items = github.enrich_with_comments(items, depth=depth, token=token)
5497 return items, _result_outcome_artifact(source, response)
5498 if source == "pinterest":
5499 result = pinterest.search_pinterest(
5500 subquery.search_query, from_date, to_date,
5501 depth=depth,
5502 token=env.get_pinterest_token(config),
5503 )
5504 return pinterest.parse_pinterest_response(result), _result_outcome_artifact(source, result)
5505 if source == "xiaohongshu":
5506 return xiaohongshu_api.search_feeds(
5507 subquery.search_query,
5508 from_date,
5509 to_date,
5510 env.get_xiaohongshu_api_base(config),
5511 depth=depth,
5512 ), {}
5513 if source == "perplexity":
5514 return perplexity.search(subquery.search_query, date_range, config, deep=config.get("_deep_research", False))
5515 raise RuntimeError(f"Unsupported source: {source}")
5516
5517
5518 def _google_key(config: dict[str, Any]) -> str | None:
5519 return config.get("GOOGLE_API_KEY") or config.get("GEMINI_API_KEY") or config.get("GOOGLE_GENAI_API_KEY")
5520
5521
5522
5523
5524 def _mock_stream_results(source: str, subquery: schema.SubQuery) -> tuple[list[dict], dict]:
5525 # Namespace URLs and the canned comment by topic: real runs never hand two
5526 # distinct stories byte-identical evidence, and discovery's same-story fold
5527 # (correctly) collapses topics that share it. Mock enrichment sub-runs feed
5528 # this fixture one topic per subquery, so the slug keeps them distinct.
5529 slug = re.sub(r"[^a-z0-9]+", "-", subquery.search_query.lower()).strip("-") or "topic"
5530 payloads = {
5531 "reddit": [
5532 {
5533 "id": "R1",
5534 "title": f"{subquery.search_query} discussion thread",
5535 "url": f"https://reddit.com/r/example/comments/{slug}-1",
5536 "subreddit": "example",
5537 "date": dates.get_date_range(5)[0],
5538 "engagement": {"score": 120, "num_comments": 48, "upvote_ratio": 0.91},
5539 "selftext": f"Community discussion about {subquery.search_query}.",
5540 "top_comments": [{"excerpt": f"Strong firsthand feedback from {subquery.search_query} users."}],
5541 "relevance": 0.82,
5542 "why_relevant": "Mock Reddit result",
5543 }
5544 ],
5545 "x": [
5546 {
5547 "id": "X1",
5548 "text": f"People on X are discussing {subquery.search_query} right now.",
5549 "url": f"https://x.com/example/status/{slug}-1",
5550 "author_handle": "example",
5551 "date": dates.get_date_range(2)[0],
5552 "engagement": {"likes": 200, "reposts": 35, "replies": 18, "quotes": 4},
5553 "relevance": 0.79,
5554 "why_relevant": "Mock X result",
5555 }
5556 ],
5557 "grounding": [
5558 {
5559 "id": "WB1",
5560 "title": f"{subquery.search_query} article",
5561 "url": f"https://example.com/article/{slug}",
5562 "source_domain": "example.com",
5563 "snippet": f"Recent web reporting about {subquery.search_query}.",
5564 "date": dates.get_date_range(7)[0],
5565 "relevance": 0.88,
5566 "why_relevant": "Brave web search",
5567 }
5568 ],
5569 "digg": [
5570 {
5571 "id": "mock1abc",
5572 "title": f"Digg cluster about {subquery.search_query}",
5573 "url": f"https://di.gg/ai/mock1abc-{slug}",
5574 "tldr": f"Curated cluster summarizing recent {subquery.search_query} discussion across the AI 1000.",
5575 "author": "",
5576 "date": dates.get_date_range(3)[0],
5577 "engagement": {"postCount": 8, "uniqueAuthors": 5, "rank": 2, "rank_score": 49.0},
5578 "first_post_age": "3d",
5579 "posts": [
5580 {
5581 "username": "exampledev",
5582 "display_name": "Example Dev",
5583 "category": "Engineer",
5584 "rank": 142,
5585 "body": f"Quote from the AI 1000 about {subquery.search_query}.",
5586 "post_type": "tweet",
5587 "x_url": "https://x.com/exampledev/status/1",
5588 "posted_at": dates.get_date_range(3)[0],
5589 },
5590 ],
5591 "relevance": 0.84,
5592 "why_relevant": "Mock Digg cluster",
5593 },
5594 {
5595 "id": "mock2def",
5596 "title": f"Second Digg cluster on {subquery.search_query}",
5597 "url": f"https://di.gg/ai/mock2def-{slug}",
5598 "tldr": f"Another angle on {subquery.search_query}.",
5599 "author": "",
5600 "date": dates.get_date_range(8)[0],
5601 "engagement": {"postCount": 3, "uniqueAuthors": 2, "rank": 18, "rank_score": 33.0},
5602 "first_post_age": "8d",
5603 "posts": [],
5604 "relevance": 0.71,
5605 "why_relevant": "Mock Digg cluster",
5606 },
5607 ],
5608 "arxiv": [
5609 {
5610 "id": f"http://arxiv.org/abs/2606.00001v1-{slug}",
5611 "title": f"A Survey of {subquery.search_query}",
5612 "url": f"https://arxiv.org/abs/2606.00001v1-{slug}",
5613 "summary": f"We present a comprehensive study of {subquery.search_query} and its recent advances.",
5614 "author": "Ada Lovelace et al.",
5615 "authors": ["Ada Lovelace", "Alan Turing"],
5616 "date": dates.get_date_range(20)[0],
5617 "engagement": {},
5618 "relevance": 0.86,
5619 "why_relevant": "Mock arXiv paper",
5620 },
5621 ],
5622 "techmeme": [
5623 {
5624 "id": f"https://www.techmeme.com/260627/p1-{slug}",
5625 "title": f"Major development in {subquery.search_query} reshapes the industry",
5626 "url": f"https://www.techmeme.com/260627/p1-{slug}",
5627 "source_name": "techcrunch.com",
5628 "date": dates.get_date_range(1)[0],
5629 "engagement": {},
5630 "relevance": 0.83,
5631 "why_relevant": "Mock Techmeme headline",
5632 },
5633 ],
5634 "dripstack": [
5635 {
5636 "id": "DS1",
5637 "title": f"Deep dive: {subquery.search_query} from a paid newsletter",
5638 "url": f"https://newsletter.example.com/deep-dive-{slug}",
5639 "author": "newsletter.example.com",
5640 "date": dates.get_date_range(3)[0],
5641 "engagement": {},
5642 "relevance": 0.85,
5643 "why_relevant": "Mock DripStack newsletter result",
5644 "snippet": f"Professional analyst coverage of {subquery.search_query}.",
5645 "metadata": {
5646 "publication_slug": "newsletter.example.com",
5647 "post_slug": "deep-dive",
5648 "relevance_score": 85,
5649 "match_confidence": "strong",
5650 },
5651 },
5652 ],
5653 "trustpilot": [
5654 {
5655 "id": "example.com",
5656 "title": f"{subquery.search_query}: TrustScore 3.4",
5657 "url": f"https://www.trustpilot.com/review/{slug}.example.com",
5658 "summary": f"Across recent reviews, customers were split on {subquery.search_query}: some praised support, others cited delays.",
5659 "name": subquery.search_query,
5660 "trustScore": 3.4,
5661 "reviewCount": 128,
5662 "date": dates.get_date_range(1)[0],
5663 "engagement": {"reviews": 128, "trustScore": 3.4},
5664 "relevance": 0.8,
5665 "why_relevant": "Mock Trustpilot sentiment",
5666 },
5667 ],
5668 # Three products spanning the drift states the footer renders: one
5669 # sagging below its all-time average (with enough in-window reviews
5670 # to clear the arrow threshold), one steady, and one too new to have
5671 # a baseline. Mock runs exercise the full R1c line without a CLI.
5672 "amazon": [
5673 {
5674 "asin": "B0MOCK00X1",
5675 "date": dates.get_date_range(1)[1],
5676 "name": f"{subquery.search_query} Pro Model | Flagship Edition",
5677 "short_name": "Pro Model",
5678 "brand": subquery.search_query.split()[0].title() if subquery.search_query else "Example",
5679 "url": "https://www.amazon.com/dp/B0MOCK00X1",
5680 "rating": 4.4,
5681 "num_ratings": 459,
5682 "price": 39.99,
5683 "currency": "USD",
5684 "badge": "Best Seller",
5685 "sponsored": False,
5686 "relevance": 0.85,
5687 "why_relevant": "Mock Amazon product",
5688 "product_rating": 4.4,
5689 "product_rating_count": 459,
5690 "star_distribution": {
5691 "one_star": 28, "two_star": 9, "three_star": 28,
5692 "four_star": 60, "five_star": 335,
5693 },
5694 "top_comments": [
5695 {
5696 "score": 3, "rating": 2, "verified": True,
5697 "date": dates.get_date_range(3)[1],
5698 "excerpt": "The tray shifts in transit and the lid jams shut.",
5699 "title": "Lid jams",
5700 },
5701 {
5702 "score": 1, "rating": 4, "verified": True,
5703 "date": dates.get_date_range(9)[1],
5704 "excerpt": "Solid build, but arrived with a dented panel.",
5705 "title": "Shipping dent",
5706 },
5707 {
5708 "score": 0, "rating": 5, "verified": True,
5709 "date": dates.get_date_range(14)[1],
5710 "excerpt": "Keeps everything cold through a full school day.",
5711 "title": "Works great",
5712 },
5713 {
5714 "score": 0, "rating": 4, "verified": True,
5715 "date": dates.get_date_range(19)[1],
5716 "excerpt": "Good size for the price.",
5717 "title": "Good value",
5718 },
5719 {
5720 "score": 0, "rating": 4, "verified": False,
5721 "date": dates.get_date_range(24)[1],
5722 "excerpt": "Does the job, nothing fancy.",
5723 "title": "Fine",
5724 },
5725 ],
5726 },
5727 {
5728 "asin": "B0MOCK00X2",
5729 "date": dates.get_date_range(1)[1],
5730 "name": f"{subquery.search_query} Classic | Everyday Model",
5731 "short_name": "Classic",
5732 "brand": subquery.search_query.split()[0].title() if subquery.search_query else "Example",
5733 "url": "https://www.amazon.com/dp/B0MOCK00X2",
5734 "rating": 4.7,
5735 "num_ratings": 8446,
5736 "price": 24.99,
5737 "currency": "USD",
5738 "sponsored": False,
5739 "relevance": 0.8,
5740 "why_relevant": "Mock Amazon product",
5741 "product_rating": 4.7,
5742 "product_rating_count": 8446,
5743 "star_distribution": {
5744 "one_star": 120, "two_star": 90, "three_star": 300,
5745 "four_star": 1010, "five_star": 6926,
5746 },
5747 "top_comments": [
5748 {
5749 "score": 12, "rating": 5, "verified": True,
5750 "date": dates.get_date_range(4)[1],
5751 "excerpt": "Third one we've bought. They last.",
5752 "title": "Repeat buyer",
5753 },
5754 ],
5755 },
5756 {
5757 "asin": "B0MOCK00X3",
5758 "date": dates.get_date_range(1)[1],
5759 "name": f"{subquery.search_query} Mini | New Release",
5760 "short_name": "Mini",
5761 "brand": subquery.search_query.split()[0].title() if subquery.search_query else "Example",
5762 "url": "https://www.amazon.com/dp/B0MOCK00X3",
5763 "rating": None,
5764 "num_ratings": 57,
5765 "price": 19.99,
5766 "currency": "USD",
5767 "sponsored": False,
5768 "relevance": 0.72,
5769 "why_relevant": "Mock Amazon product",
5770 },
5771 ],
5772 "jobs": [
5773 {
5774 "id": "J1",
5775 "title": "Founding Enterprise Solutions Engineer",
5776 "url": f"https://boards.greenhouse.io/example/jobs/{slug}-1",
5777 "description": (
5778 f"Work with enterprise customers on SSO, SOC 2, security, "
5779 f"and procurement workflows for {subquery.search_query}."
5780 ),
5781 "department": "Sales",
5782 "location": "San Francisco, CA",
5783 "date": dates.get_date_range(4)[0],
5784 "provider": "mock",
5785 "relevance": 0.8,
5786 "why_relevant": "Mock public job posting",
5787 },
5788 {
5789 "id": "J2",
5790 "title": "Security Platform Engineer",
5791 "url": f"https://boards.greenhouse.io/example/jobs/{slug}-2",
5792 "description": "Build enterprise security, audit, and admin workflows.",
5793 "department": "Engineering",
5794 "location": "Remote",
5795 "date": dates.get_date_range(6)[0],
5796 "provider": "mock",
5797 "relevance": 0.78,
5798 "why_relevant": "Mock public job posting",
5799 },
5800 ],
5801 }
5802 if source == "grounding":
5803 return payloads.get(source, []), {
5804 "label": subquery.label,
5805 "mock": True,
5806 "webSearchQueries": [subquery.search_query],
5807 "resultCount": 1,
5808 }
5809 return payloads.get(source, []), {}
5810
5810 lines PYTHON