| 1 | """Keyless Reddit pipeline: free discovery + comment enrichment. |
| 2 | |
| 3 | ``search.json`` is permanently 403/429 keyless, so it is not used. Discovery |
| 4 | runs on the surfaces that still serve data without a key, then enrichment runs |
| 5 | on whatever was discovered: |
| 6 | |
| 7 | Dedicated lane entity-home subreddits (e.g. r/Kanye) pulled in full via the |
| 8 | shreddit listing partials (top+hot+new, real scores), kept |
| 9 | whole (floor-exempt) because the sub IS the topic. |
| 10 | Search lane Reddit's site search fragment (reddit_search): global keyword |
| 11 | search plus per-sub search for targeted subreddits. Results |
| 12 | arrive dated and scored. Targeted runs also merge the |
| 13 | subreddits' listing partials. Relevance-floored. |
| 14 | Enrichment shreddit comment + count enrichment (reddit_shreddit) for the |
| 15 | top-ranked posts (author + score + text + permalink). |
| 16 | |
| 17 | Returns ``[]`` (never raises) so ``pipeline.py`` can fall through to the |
| 18 | ScrapeCreators backup when every keyless lane comes up empty. |
| 19 | """ |
| 20 | |
| 21 | import concurrent.futures |
| 22 | import math |
| 23 | import sys |
| 24 | from concurrent.futures import ThreadPoolExecutor |
| 25 | from typing import Any, Dict, List, Optional |
| 26 | |
| 27 | from . import http |
| 28 | from . import reddit_search, reddit_shreddit, reddit_listing, reddit_arctic |
| 29 | # High-upvote posts from targeted-sub listings can bury on-topic search hits |
| 30 | # under an engagement-first final sort. A relevance floor + relevance-first |
| 31 | # final ranking keeps the section on-topic. Thresholds are shared with the |
| 32 | # keyed path (reddit.py) via relevance.py. |
| 33 | from .relevance import RELEVANCE_FLOOR, MIN_ON_TOPIC |
| 34 | |
| 35 | ENRICH_LIMITS = reddit_shreddit.ENRICH_LIMITS |
| 36 | ENRICH_BUDGET = 45 # seconds total across all enrichment threads |
| 37 | MAX_ENRICH_WORKERS = 4 |
| 38 | # Dedicated subreddits (the entity's home, e.g. r/Kanye for "Kanye West") are |
| 39 | # wholly on-topic, so pull top+hot+new — the top-of-month listing alone misses |
| 40 | # fresh threads — and keep every item (floor-exempt). |
| 41 | DEDICATED_SORTS = ["top", "hot", "new"] |
| 42 | |
| 43 | |
| 44 | def _relevance_rank_key(post: Dict[str, Any]) -> float: |
| 45 | """Rank by relevance first, with a bounded engagement bonus as tiebreaker. |
| 46 | |
| 47 | Mirrors reddit.py: the log-scaled bonus (capped at 0.25) orders |
| 48 | similarly-relevant posts by discussion volume but is too small to lift an |
| 49 | off-topic post (relevance ~0) above an on-topic one. |
| 50 | """ |
| 51 | eng = post.get("engagement", {}) |
| 52 | total = (eng.get("score", 0) or 0) + (eng.get("num_comments", 0) or 0) |
| 53 | return (post.get("relevance") or 0.0) + min(0.25, math.log10(max(0, total) + 1) / 20.0) |
| 54 | |
| 55 | |
| 56 | def _log(msg: str) -> None: |
| 57 | sys.stderr.write(f"[RedditKeyless] {msg}\n") |
| 58 | sys.stderr.flush() |
| 59 | |
| 60 | |
| 61 | def _apply_scores(post: Dict[str, Any], scored: Dict[str, int]) -> None: |
| 62 | post["score"] = scored["score"] |
| 63 | post["num_comments"] = scored["num_comments"] |
| 64 | post.setdefault("engagement", {})["score"] = scored["score"] |
| 65 | post["engagement"]["num_comments"] = scored["num_comments"] |
| 66 | |
| 67 | |
| 68 | def _scored_listings( |
| 69 | subreddits: List[str], |
| 70 | depth: str = "default", |
| 71 | query: str = "", |
| 72 | sorts: Optional[List[str]] = None, |
| 73 | ) -> List[Dict[str, Any]]: |
| 74 | """Scored subreddit listings: shreddit partials, arctic-shift supplement. |
| 75 | |
| 76 | The shreddit ``community-more-posts`` partials 403 from datacenter IPs |
| 77 | (and any host Reddit decides to block). Shreddit is tried first; arctic- |
| 78 | shift supplements with any posts shreddit missed. Individual sort lanes |
| 79 | can fail silently (shreddit's ``fetch_listings`` flattens results without |
| 80 | exposing per-sort status), so arctic is called for all requested subreddits |
| 81 | and merged via deduplication. This ensures fresh posts sought through |
| 82 | ``hot`` or ``new`` are recovered even when only ``top`` succeeded. Never |
| 83 | raises. |
| 84 | """ |
| 85 | posts = reddit_listing.fetch_listings(subreddits, depth=depth, query=query, sorts=sorts) |
| 86 | |
| 87 | # Supplement with arctic for all requested subreddits. Shreddit's per-sort |
| 88 | # success/failure is opaque, so arctic provides coverage for any failed |
| 89 | # sort lanes (e.g., hot/new failing while top succeeded). Deduplication |
| 90 | # ensures no redundant posts when shreddit fully succeeded. |
| 91 | if subreddits: |
| 92 | try: |
| 93 | arctic_posts = reddit_arctic.fetch_listings( |
| 94 | subreddits, depth=depth, query=query, sorts=sorts |
| 95 | ) |
| 96 | except Exception as exc: # the fallback must never break the pipeline |
| 97 | _log(f"arctic-shift listing supplement failed: {exc}") |
| 98 | arctic_posts = [] |
| 99 | if arctic_posts: |
| 100 | # Merge and dedupe by URL — shreddit posts take priority. |
| 101 | seen = {p["url"] for p in posts} |
| 102 | added = 0 |
| 103 | for p in arctic_posts: |
| 104 | if p["url"] not in seen: |
| 105 | seen.add(p["url"]) |
| 106 | posts.append(p) |
| 107 | added += 1 |
| 108 | if added: |
| 109 | _log(f"arctic-shift supplement: {added} new posts from {len(arctic_posts)} arctic results") |
| 110 | return posts |
| 111 | |
| 112 | |
| 113 | def _discover( |
| 114 | topic: str, |
| 115 | depth: str, |
| 116 | subreddits: Optional[List[str]], |
| 117 | dedicated_subreddits: Optional[List[str]] = None, |
| 118 | from_date: Optional[str] = None, |
| 119 | to_date: Optional[str] = None, |
| 120 | ) -> List[Dict[str, Any]]: |
| 121 | # Dedicated lane: the entity's home subs are wholly on-topic. Pull |
| 122 | # top+hot+new (real scores from the listing) and mark them floor-exempt so |
| 123 | # an on-topic post whose title lacks the entity name is never dropped. |
| 124 | dedicated_posts: List[Dict[str, Any]] = [] |
| 125 | if dedicated_subreddits: |
| 126 | dedicated_posts = _scored_listings( |
| 127 | dedicated_subreddits, depth=depth, query=topic, sorts=DEDICATED_SORTS |
| 128 | ) |
| 129 | for p in dedicated_posts: |
| 130 | p["dedicated"] = True |
| 131 | _log(f"Dedicated lane: {len(dedicated_posts)} posts from {dedicated_subreddits}") |
| 132 | |
| 133 | # search.json is permanently 403/429 keyless (no Tier 0). Discovery is |
| 134 | # Reddit's site search: global, plus per-sub for targeted subreddits. The |
| 135 | # dedicated subs keep their full-listing lane above and are not searched. |
| 136 | search_posts = reddit_search.search( |
| 137 | topic, depth=depth, subreddits=subreddits, from_date=from_date, to_date=to_date |
| 138 | ) |
| 139 | |
| 140 | # Targeted run: the caller chose these subreddits, so their listing cards |
| 141 | # are on-topic; include them as scored discovery. A bare run fetches no |
| 142 | # listings, so high-upvote off-topic posts cannot flood the results. |
| 143 | listing_posts = ( |
| 144 | _scored_listings(subreddits, depth=depth, query=topic) if subreddits else [] |
| 145 | ) |
| 146 | _log( |
| 147 | f"Tier 1 (site search) {len(search_posts)} posts; " |
| 148 | f"{'listing discovery ' + str(len(listing_posts)) if subreddits else 'no listings'}" |
| 149 | ) |
| 150 | |
| 151 | # Score lookup by post id, from the scored listing cards. |
| 152 | score_map: Dict[str, Dict[str, int]] = {} |
| 153 | for p in listing_posts: |
| 154 | pid = p.get("metadata", {}).get("post_id", "") |
| 155 | if pid: |
| 156 | score_map[pid] = {"score": p["score"], "num_comments": p["num_comments"]} |
| 157 | |
| 158 | # Merge: dedicated-sub posts first (floor-exempt), then scored listing posts |
| 159 | # (targeted only), then search results, taking a listing's live score where |
| 160 | # the post also appears there. First writer wins the dedupe, so a thread in |
| 161 | # both the dedicated lane and a listing keeps its floor-exempt status. |
| 162 | merged: List[Dict[str, Any]] = [] |
| 163 | seen: set = set() |
| 164 | for p in dedicated_posts + listing_posts: |
| 165 | if p["url"] not in seen: |
| 166 | seen.add(p["url"]) |
| 167 | merged.append(p) |
| 168 | for p in search_posts: |
| 169 | if p["url"] in seen: |
| 170 | continue |
| 171 | pid = reddit_listing._post_id(p["url"]) |
| 172 | if pid in score_map: |
| 173 | _apply_scores(p, score_map[pid]) |
| 174 | seen.add(p["url"]) |
| 175 | merged.append(p) |
| 176 | |
| 177 | # Backfill posts still at score 0 from the free arctic-shift archive. Posts |
| 178 | # already scored by search or a listing keep that live score; arctic only |
| 179 | # fills the gap, and is best-effort (never raises). |
| 180 | need = [pid for p in merged |
| 181 | if not (p.get("engagement", {}).get("score")) |
| 182 | for pid in [reddit_listing._post_id(p["url"])] if pid] |
| 183 | if need: |
| 184 | scores = reddit_arctic.fetch_scores(need) |
| 185 | filled = 0 |
| 186 | for p in merged: |
| 187 | if p.get("engagement", {}).get("score"): |
| 188 | continue |
| 189 | pid = reddit_listing._post_id(p["url"]) |
| 190 | if pid in scores: |
| 191 | _apply_scores(p, scores[pid]) |
| 192 | filled += 1 |
| 193 | if filled: |
| 194 | _log(f"arctic-shift backfilled {filled} post scores") |
| 195 | return merged |
| 196 | |
| 197 | |
| 198 | def _enrich_one(post: Dict[str, Any]) -> Dict[str, Any]: |
| 199 | """Attach shreddit comments + real comment count. Never raises.""" |
| 200 | try: |
| 201 | data = reddit_shreddit.fetch_comments(post.get("url", "")) |
| 202 | if data.get("top_comments"): |
| 203 | post["top_comments"] = data["top_comments"] |
| 204 | if data.get("comment_insights"): |
| 205 | post["comment_insights"] = data["comment_insights"] |
| 206 | num = data.get("num_comments") |
| 207 | if num is not None: |
| 208 | post["num_comments"] = num |
| 209 | post.setdefault("engagement", {})["num_comments"] = num |
| 210 | except Exception: |
| 211 | pass # keep the post with whatever discovery gave us |
| 212 | return post |
| 213 | |
| 214 | |
| 215 | def _enrich(posts: List[Dict[str, Any]], depth: str) -> List[Dict[str, Any]]: |
| 216 | """Enrich the top N posts with comments under a total time budget.""" |
| 217 | limit = ENRICH_LIMITS.get(depth, ENRICH_LIMITS["default"]) |
| 218 | to_enrich = posts[:limit] |
| 219 | rest = posts[limit:] |
| 220 | if not to_enrich: |
| 221 | return posts |
| 222 | |
| 223 | result_map: Dict[int, Dict[str, Any]] = {} |
| 224 | try: |
| 225 | with ThreadPoolExecutor(max_workers=min(limit, MAX_ENRICH_WORKERS)) as executor: |
| 226 | futures = { |
| 227 | http.submit_with_context(executor, _enrich_one, post): i |
| 228 | for i, post in enumerate(to_enrich) |
| 229 | } |
| 230 | # The budget covers the fetches; the allowance covers the shared |
| 231 | # bucket's queue (other lanes, other entities in compare mode). |
| 232 | done, not_done = concurrent.futures.wait( |
| 233 | futures, |
| 234 | timeout=ENRICH_BUDGET + http.reddit_keyless_wait_allowance(len(to_enrich)), |
| 235 | ) |
| 236 | for future in done: |
| 237 | idx = futures[future] |
| 238 | try: |
| 239 | result_map[idx] = future.result(timeout=0) |
| 240 | except Exception: |
| 241 | result_map[idx] = to_enrich[idx] |
| 242 | for future in not_done: |
| 243 | idx = futures[future] |
| 244 | result_map[idx] = to_enrich[idx] |
| 245 | future.cancel() |
| 246 | enriched = [result_map[i] for i in range(len(to_enrich))] |
| 247 | except Exception: |
| 248 | enriched = to_enrich |
| 249 | |
| 250 | return enriched + rest |
| 251 | |
| 252 | |
| 253 | def _by_comments(posts: List[Dict[str, Any]]) -> List[Dict[str, Any]]: |
| 254 | """Stable-sort posts by comment count descending for enrichment slots. |
| 255 | |
| 256 | Ties (equal counts, including unknown counts treated as 0) preserve the |
| 257 | incoming order, which search_and_enrich's provisional score-first sort |
| 258 | establishes. Mirrors _relevance_rank_key's `or 0` guard so a present-but- |
| 259 | None count is treated as 0 rather than raising. |
| 260 | """ |
| 261 | def _comment_count(post: Dict[str, Any]) -> int: |
| 262 | eng = post.get("engagement") or {} |
| 263 | return eng.get("num_comments") or post.get("num_comments") or 0 |
| 264 | |
| 265 | return sorted(posts, key=_comment_count, reverse=True) |
| 266 | |
| 267 | |
| 268 | def _slot_priority(topic: str, posts: List[Dict[str, Any]]) -> List[Dict[str, Any]]: |
| 269 | """Order posts for enrichment slots: entity-matching posts first. |
| 270 | |
| 271 | Comment slots (ENRICH_LIMITS) are scarce; spending them on high-upvote |
| 272 | posts that rerank later demotes as entity misses starves the on-topic |
| 273 | posts the user actually sees (2026-06-06 "OpenClaw vs Hermes" run: |
| 274 | 2,000+ upvote Gemma/GPU threads took every slot, then were demoted to |
| 275 | zero). Mirror rerank's demotion signal via the shared `_entity_grounded` |
| 276 | check (head token of the topic's stripped primary entity present in the |
| 277 | post text) so slots go to posts likely to survive final ranking — keying |
| 278 | on the same head token keeps the two paths from diverging. Falls back to |
| 279 | token-overlap relevance when the topic yields no usable primary entity. |
| 280 | Within each tier posts are ordered by comment count descending (stable: |
| 281 | equal or unknown counts preserve the incoming score-first order), so the |
| 282 | scarce slots go to the threads with the most discussion rather than to |
| 283 | near-empty threads that merely ranked higher by score. Never raises; on |
| 284 | any failure the incoming order is returned unchanged. |
| 285 | """ |
| 286 | try: |
| 287 | from . import relevance, rerank |
| 288 | |
| 289 | def _post_text(post: Dict[str, Any]) -> str: |
| 290 | return f"{post.get('title') or ''} {post.get('selftext') or ''}" |
| 291 | |
| 292 | entity = rerank._primary_entity(topic).lower() |
| 293 | if entity: |
| 294 | def _matches(post: Dict[str, Any]) -> bool: |
| 295 | return rerank._entity_grounded(_post_text(post), entity) |
| 296 | else: |
| 297 | prepared = relevance.PreparedQuery(topic) |
| 298 | |
| 299 | def _matches(post: Dict[str, Any]) -> bool: |
| 300 | return relevance.token_overlap_relevance(prepared, _post_text(post)) > 0.24 |
| 301 | |
| 302 | matches: List[Dict[str, Any]] = [] |
| 303 | misses: List[Dict[str, Any]] = [] |
| 304 | for post in posts: |
| 305 | (matches if _matches(post) else misses).append(post) |
| 306 | return _by_comments(matches) + _by_comments(misses) |
| 307 | except Exception: |
| 308 | return posts |
| 309 | |
| 310 | |
| 311 | def search_and_enrich( |
| 312 | topic: str, |
| 313 | from_date: str, |
| 314 | to_date: str, |
| 315 | depth: str = "default", |
| 316 | subreddits: Optional[List[str]] = None, |
| 317 | dedicated_subreddits: Optional[List[str]] = None, |
| 318 | ) -> List[Dict[str, Any]]: |
| 319 | """Full keyless Reddit pipeline: discover then enrich. |
| 320 | |
| 321 | Args: |
| 322 | topic: Search topic |
| 323 | from_date: Start date (YYYY-MM-DD) |
| 324 | to_date: End date (YYYY-MM-DD) |
| 325 | depth: 'quick', 'default', or 'deep' |
| 326 | subreddits: Optional pre-resolved broad/category subreddit names (no r/) |
| 327 | dedicated_subreddits: Optional entity-home subreddit names (no r/) pulled |
| 328 | in full (top+hot+new) and exempt from the relevance floor. |
| 329 | |
| 330 | Returns: |
| 331 | List of normalized item dicts matching the reddit_public output shape, |
| 332 | with top_comments/comment_insights attached on enriched posts. |
| 333 | Empty list when all keyless tiers fail (so SC backup can engage). |
| 334 | """ |
| 335 | posts = _discover( |
| 336 | topic, depth, subreddits, dedicated_subreddits, from_date=from_date, to_date=to_date |
| 337 | ) |
| 338 | if not posts: |
| 339 | return [] |
| 340 | |
| 341 | # Date filter: keep posts in range or with unknown dates (mirrors reddit_public). |
| 342 | posts = [ |
| 343 | p for p in posts |
| 344 | if p.get("date") is None or (from_date <= p["date"] <= to_date) |
| 345 | ] |
| 346 | |
| 347 | # Relevance floor: strip zero-overlap posts (relevance exactly 0 = no |
| 348 | # title/body token match at all) when anything relevant remains, so |
| 349 | # high-upvote listing posts can't bury on-topic search |
| 350 | # hits. Keep all only when nothing scored above zero. |
| 351 | before = len(posts) |
| 352 | # Dedicated-sub posts are floor-exempt: their whole subreddit is the topic, |
| 353 | # so an on-topic post whose title lacks the entity name must not be dropped. |
| 354 | on_topic = [p for p in posts if p.get("dedicated") or (p.get("relevance") or 0) >= RELEVANCE_FLOOR] |
| 355 | if len(on_topic) >= MIN_ON_TOPIC: |
| 356 | posts = on_topic |
| 357 | else: |
| 358 | nonzero = [p for p in posts if p.get("dedicated") or (p.get("relevance") or 0) > 0] |
| 359 | if nonzero: |
| 360 | posts = nonzero |
| 361 | if len(posts) < before: |
| 362 | _log(f"Relevance floor dropped {before - len(posts)} off-topic posts") |
| 363 | |
| 364 | # Provisional score-first order so enrichment-slot selection has a stable |
| 365 | # within-tier tiebreak order to preserve: within each entity tier, slots go |
| 366 | # to the most-commented threads first, and equal counts keep score order. |
| 367 | posts.sort( |
| 368 | key=lambda p: ( |
| 369 | p.get("engagement", {}).get("score", 0) or 0, |
| 370 | p.get("relevance", 0) or 0, |
| 371 | p.get("date") or "", |
| 372 | ), |
| 373 | reverse=True, |
| 374 | ) |
| 375 | |
| 376 | # Enrichment slot selection is comment-aware within entity tiers: |
| 377 | # entity-matching posts claim the scarce comment slots first, and within |
| 378 | # each tier the most-commented threads get slots first (score order is the |
| 379 | # stable tiebreak for equal counts). |
| 380 | posts = _enrich(_slot_priority(topic, posts), depth) |
| 381 | |
| 382 | # Final display order ranks relevance-first with a bounded engagement bonus, |
| 383 | # so an off-topic high-upvote post can't outrank an on-topic one in what the |
| 384 | # user sees. Enrichment above may have backfilled real comment counts. |
| 385 | posts.sort(key=_relevance_rank_key, reverse=True) |
| 386 | |
| 387 | for i, post in enumerate(posts): |
| 388 | post["id"] = f"R{i + 1}" |
| 389 | |
| 390 | return posts |
| 391 |