返回 last30days-skill
reddit_keyless.py
根目录 / skills / last30days / scripts / lib / reddit_keyless.py
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
391 lines PYTHON