返回 last30days-skill
reddit.py
根目录 / skills / last30days / scripts / lib / reddit.py
1 """Reddit search via ScrapeCreators API for the v3 pipeline.
2
3 Uses ScrapeCreators REST API to search Reddit globally, discover relevant
4 subreddits, run targeted subreddit searches, and fetch comment trees.
5
6 Requires SCRAPECREATORS_API_KEY in config (same key as TikTok + Instagram).
7 API docs: https://scrapecreators.com/docs
8 """
9
10 import copy
11 import math
12 import re
13 import sys
14 import threading
15 import time
16 from collections import Counter
17 from concurrent.futures import ThreadPoolExecutor, as_completed, wait as futures_wait
18 from datetime import date, datetime, timezone
19 from typing import Any, Dict, List, Optional, Set
20
21 def _first_of(*values, default=None):
22 """Return first value that is not None."""
23 for v in values:
24 if v is not None:
25 return v
26 return default
27
28 from . import dates, health, http, log
29
30 SCRAPECREATORS_BASE = "https://api.scrapecreators.com/v1/reddit"
31
32 # Reddit's highest-upvote content (relationship drama, AITA, viral news) often
33 # has near-zero topic overlap. Engagement-only ranking floats it above on-topic
34 # posts, especially on bare global searches where no subreddits were resolved
35 # upstream. A relevance floor + relevance-first ranking (see _relevance_rank_key
36 # below and RELEVANCE_FLOOR / MIN_ON_TOPIC in relevance.py) keeps an off-topic
37 # viral post from ever outranking an on-topic one.
38
39 # Depth configurations: how many API calls per phase
40 DEPTH_CONFIG = {
41 "quick": {
42 "global_searches": 1,
43 "subreddit_searches": 2,
44 "comment_enrichments": 3,
45 },
46 "default": {
47 "global_searches": 2,
48 "subreddit_searches": 3,
49 "comment_enrichments": 5,
50 },
51 "deep": {
52 "global_searches": 3,
53 "subreddit_searches": 5,
54 "comment_enrichments": 8,
55 },
56 }
57
58 from .query import extract_core_subject as _query_extract, infer_query_intent
59 from .relevance import token_overlap_relevance, RELEVANCE_FLOOR, MIN_ON_TOPIC
60
61
62 # Reddit-specific noise words (preserves original smaller set)
63 NOISE_WORDS = frozenset({
64 'best', 'top', 'good', 'great', 'awesome', 'killer',
65 'latest', 'new', 'news', 'update', 'updates',
66 'trending', 'hottest', 'popular',
67 'practices', 'features', 'tips',
68 'recommendations', 'advice',
69 'prompt', 'prompts', 'prompting',
70 'methods', 'strategies', 'approaches',
71 'how', 'to', 'the', 'a', 'an', 'for', 'with',
72 'of', 'in', 'on', 'is', 'are', 'what', 'which',
73 'guide', 'tutorial', 'using',
74 })
75
76
77 def _log(msg: str):
78 log.source_log("Reddit", msg, tty_only=False)
79
80
81 def classify_run_failure(detail: str) -> str:
82 """Map Reddit auth and anti-bot responses that do not carry HTTP status."""
83 text = detail.lower()
84 if any(marker in text for marker in ("interstitial", "blocked by reddit", "too many requests")):
85 return health.RATE_LIMITED
86 if any(marker in text for marker in ("login required", "invalid token", "expired token")):
87 return health.AUTH_FAILED
88 return http.classify_failure(message=detail)
89
90
91 def _extract_core_subject(topic: str) -> str:
92 """Extract core subject from verbose query.
93
94 Strips meta/research words to keep only the core product/concept name.
95 """
96 return _query_extract(topic, noise=NOISE_WORDS)
97
98
99 def expand_reddit_queries(topic: str, depth: str) -> List[str]:
100 """Generate multiple Reddit search queries from a topic.
101
102 Uses local logic (no LLM call needed):
103 1. Extract core subject (strip noise words)
104 2. Include original topic if different from core
105 3. For default/deep: add casual/review variant
106 4. For deep: add problem/issues variant
107
108 Returns 1-4 query strings depending on depth.
109 """
110 core = _extract_core_subject(topic)
111 queries = [core]
112
113 # Broader variant: include more context from original topic
114 original_clean = topic.strip().rstrip('?!.')
115 if core.lower() != original_clean.lower() and len(original_clean.split()) <= 8:
116 queries.append(original_clean)
117
118 qtype = infer_query_intent(topic)
119
120 # Product queries: always include review-oriented variant to bias toward
121 # review communities instead of keyword-matching unrelated subreddits.
122 if qtype == "product":
123 queries.append(f"{core} review OR recommendation OR best")
124
125 # Comparison queries: include head-to-head discussion variant.
126 if qtype == "comparison":
127 queries.append(f"{core} worth it OR vs OR compared")
128
129 # Opinion/review variants for default/deep depth.
130 if depth in ("default", "deep") and qtype in ("product", "opinion"):
131 queries.append(f"{core} worth it OR thoughts OR review")
132
133 # Problem/bug variants are useful for tool workflows, not generic news.
134 if depth == "deep" and qtype in ("product", "opinion", "how_to"):
135 queries.append(f"{core} issues OR problems OR bug OR broken")
136
137 return queries
138
139
140 # Known utility/meta subreddits that match queries but aren't discussion subs.
141 # These get a 0.3x penalty (not banned) in subreddit discovery scoring.
142 UTILITY_SUBS = frozenset({
143 'namethatsong', 'findthatsong', 'tipofmytongue',
144 'whatisthissong', 'helpmefind', 'whatisthisthing',
145 'whatsthissong', 'findareddit', 'subredditdrama',
146 })
147
148
149 def discover_subreddits(
150 results: List[Dict[str, Any]],
151 topic: str = "",
152 max_subs: int = 5,
153 ) -> List[str]:
154 """Extract top subreddits from global search results with relevance weighting.
155
156 Uses frequency + topic-word matching + utility-sub penalties + engagement
157 bonus to find discussion subs rather than utility/meta subs.
158
159 Args:
160 results: List of post dicts from global search
161 topic: Original search topic (for relevance matching)
162 max_subs: Maximum subreddits to return
163
164 Returns:
165 Top subreddit names sorted by weighted score
166 """
167 core = _extract_core_subject(topic) if topic else ""
168 core_words = set(core.lower().split()) if core else set()
169
170 scores = Counter()
171 for post in results:
172 sub = _extract_subreddit_name(post.get("subreddit", ""))
173 if not sub:
174 continue
175
176 # Base: frequency count
177 base = 1.0
178
179 # Bonus: subreddit name contains a core topic word
180 sub_lower = sub.lower()
181 if core_words and any(w in sub_lower for w in core_words if len(w) > 2):
182 base += 2.0
183
184 # Penalty: known utility/meta subreddits
185 if sub_lower in UTILITY_SUBS:
186 base *= 0.3
187
188 # Bonus: post engagement (high-engagement posts = better sub)
189 ups = _first_of(post.get("ups"), post.get("score"), post.get("votes"), default=0)
190 if ups and ups > 100:
191 base += 0.5
192
193 scores[sub] += base
194
195 return [sub for sub, _ in scores.most_common(max_subs)]
196
197
198 def _parse_date(value) -> Optional[str]:
199 """Convert Unix timestamp or ISO-8601 string to YYYY-MM-DD.
200
201 Global search returns ``created_at`` as an ISO string
202 (e.g. "2018-05-03T01:09:17.620000+0000"); subreddit search returns
203 ``created_utc`` as a Unix timestamp. dates.parse_date() handles both,
204 plus edge cases like Z suffix and +0000 (no colon) offset.
205
206 Falsy inputs (None, "", 0) return None, matching the original behavior
207 where a Unix timestamp of 0 meant "no date" rather than epoch 0.
208 """
209 if not value:
210 return None
211 dt = dates.parse_date(str(value))
212 return dt.strftime("%Y-%m-%d") if dt else None
213
214
215 def _extract_subreddit_name(value: Any) -> str:
216 """Extract subreddit name from string or API object dict."""
217 if isinstance(value, dict):
218 return str(value.get("name") or value.get("display_name") or "").strip()
219 return str(value).strip()
220
221
222 def _extract_score(post: Dict[str, Any]) -> int:
223 """Extract post score from either API schema.
224
225 Global search uses ``votes``; subreddit search uses ``ups``/``score``.
226 """
227 return _first_of(post.get("ups"), post.get("score"), post.get("votes"), default=0)
228
229
230 def _extract_date(post: Dict[str, Any]) -> Optional[str]:
231 """Extract date from either API schema.
232
233 Global search uses ``created_at`` (ISO); subreddit search uses ``created_utc`` (Unix).
234 """
235 return _parse_date(
236 post.get("created_utc") or post.get("created_at") or post.get("created_at_iso")
237 )
238
239
240 def _normalize_reddit_id(raw_id: str) -> str:
241 """Strip Reddit fullname prefix (t3_) for consistent dedup."""
242 s = str(raw_id or "")
243 return s[3:] if s.startswith("t3_") else s
244
245
246 def _total_engagement(item: Dict[str, Any]) -> int:
247 """Combined engagement score: upvotes + comment count.
248
249 Used for selecting which threads to enrich with comments.
250 Threads with lots of comments are high-value even if upvote score is low.
251 """
252 eng = item.get("engagement", {})
253 score = eng.get("score", 0) or 0
254 num_comments = eng.get("num_comments", 0) or 0
255 return score + num_comments
256
257
258 def _relevance_rank_key(item: Dict[str, Any]) -> float:
259 """Rank by relevance first, with a bounded engagement bonus as tiebreaker.
260
261 The log-scaled bonus is capped at 0.25 so it orders similarly-relevant posts
262 by discussion volume but is too small to lift an off-topic post (relevance
263 ~0) above an on-topic one (relevance >= RELEVANCE_FLOOR).
264 """
265 rel = item.get("relevance") or 0.0
266 eng_bonus = min(0.25, math.log10(max(0, _total_engagement(item)) + 1) / 20.0)
267 return rel + eng_bonus
268
269
270 def _normalize_post(post: Dict[str, Any], idx: int, source_label: str = "global", query: str = "") -> Dict[str, Any]:
271 """Normalize a ScrapeCreators Reddit post to our internal format.
272
273 Handles both the global-search schema (``votes``, ``created_at``,
274 ``subreddit`` as dict) and the subreddit-search schema (``ups``/``score``,
275 ``created_utc``, ``subreddit`` as string).
276 """
277 permalink = post.get("permalink", "")
278 url = f"https://www.reddit.com{permalink}" if permalink else post.get("url", "")
279
280 # Ensure URL looks like a Reddit thread
281 if url and "reddit.com" not in url:
282 url = ""
283
284 title = str(post.get("title", "")).strip()
285 selftext = str(post.get("selftext", ""))
286
287 # Score the title first, then let the body provide limited support.
288 # This keeps long selftexts from overpowering the visible topic signal.
289 relevance = _compute_post_relevance(query, title, selftext) if query else 0.7
290
291 return {
292 "id": f"R{idx}",
293 "reddit_id": _normalize_reddit_id(post.get("id", "")),
294 "title": title,
295 "url": url,
296 "subreddit": _extract_subreddit_name(post.get("subreddit", "")),
297 "date": _extract_date(post),
298 "engagement": {
299 "score": _extract_score(post),
300 "num_comments": post.get("num_comments", 0),
301 "upvote_ratio": post.get("upvote_ratio"),
302 },
303 "relevance": relevance,
304 "why_relevant": f"Reddit {source_label} search",
305 "selftext": str(post.get("selftext", ""))[:500],
306 }
307
308
309 def _compute_post_relevance(query: str, title: str, selftext: str) -> float:
310 """Compute Reddit relevance with title-first weighting.
311
312 Title should carry most of the weight because it is the visible summary the
313 user sees. Selftext can lift a marginal match, but it should not rescue a
314 weak or ambiguous title into the top ranks.
315 """
316 title_score = token_overlap_relevance(query, title)
317 if not selftext.strip():
318 return title_score
319
320 body_score = token_overlap_relevance(query, selftext)
321 support_score = max(title_score, body_score)
322 return round(0.75 * title_score + 0.25 * support_score, 2)
323
324
325 def _global_search(
326 query: str,
327 token: str,
328 sort: str = "relevance",
329 timeframe: str = "month",
330 ) -> List[Dict[str, Any]]:
331 """Search across all of Reddit via ScrapeCreators global search.
332
333 Args:
334 query: Search query
335 token: ScrapeCreators API key
336 sort: Sort order (relevance, hot, top, new)
337 timeframe: Time filter (hour, day, week, month, year, all)
338
339 Returns:
340 List of post dicts
341 """
342 try:
343 data = http.get(
344 f"{SCRAPECREATORS_BASE}/search",
345 headers=http.scrapecreators_headers(token),
346 params={"query": query, "sort": sort, "timeframe": timeframe},
347 timeout=30,
348 retries=2,
349 )
350 return data.get("posts", data.get("data", []))
351 except http.HTTPError as e:
352 if e.status_code in (401, 402, 403):
353 raise
354 _log(f"Global search error: {e}")
355 return []
356 except Exception as e:
357 _log(f"Global search error: {e}")
358 return []
359
360
361 def _subreddit_search(
362 subreddit: str,
363 query: str,
364 token: str,
365 sort: str = "relevance",
366 timeframe: str = "month",
367 ) -> List[Dict[str, Any]]:
368 """Search within a specific subreddit via ScrapeCreators.
369
370 Args:
371 subreddit: Subreddit name (without r/)
372 query: Search query
373 token: ScrapeCreators API key
374 sort: Sort order
375 timeframe: Time filter
376
377 Returns:
378 List of post dicts
379 """
380 try:
381 data = http.get(
382 f"{SCRAPECREATORS_BASE}/subreddit/search",
383 headers=http.scrapecreators_headers(token),
384 params={
385 "subreddit": subreddit,
386 "query": query,
387 "sort": sort,
388 "timeframe": timeframe,
389 },
390 timeout=30,
391 retries=2,
392 )
393 return data.get("posts", data.get("data", []))
394 except http.HTTPError as e:
395 if e.status_code in (401, 402, 403):
396 raise
397 _log(f"Subreddit search error for r/{subreddit}: {e}")
398 return []
399 except Exception as e:
400 _log(f"Subreddit search error for r/{subreddit}: {e}")
401 return []
402
403
404 def fetch_post_comments(
405 url: str,
406 token: str,
407 ) -> List[Dict[str, Any]]:
408 """Fetch comments for a Reddit post via ScrapeCreators.
409
410 Args:
411 url: Reddit post URL or permalink
412 token: ScrapeCreators API key
413
414 Returns:
415 List of comment dicts with score, author, body, etc.
416 """
417 try:
418 data = http.get(
419 f"{SCRAPECREATORS_BASE}/post/comments",
420 headers=http.scrapecreators_headers(token),
421 params={"url": url},
422 timeout=30,
423 retries=2,
424 )
425 return data.get("comments", data.get("data", []))
426 except http.HTTPError as e:
427 if e.status_code in (401, 402, 403):
428 raise
429 _log(f"Comment fetch error: {e}")
430 return []
431 except Exception as e:
432 _log(f"Comment fetch error: {e}")
433 return []
434
435
436 def _dedupe_posts(posts: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
437 """Deduplicate posts by reddit_id, keeping first occurrence."""
438 seen_ids = set()
439 seen_urls = set()
440 unique = []
441 for post in posts:
442 rid = post.get("reddit_id", "")
443 url = post.get("url", "")
444 if rid and rid in seen_ids:
445 continue
446 if url and url in seen_urls:
447 continue
448 if rid:
449 seen_ids.add(rid)
450 if url:
451 seen_urls.add(url)
452 unique.append(post)
453 return unique
454
455
456 def _days_to_reddit_bucket(days: float) -> str:
457 """Map a day count onto the smallest Reddit rolling bucket that covers it.
458
459 Adds one day of slack so calendar windows that cross a day boundary still
460 fit inside Reddit's rolling ``t=`` buckets (a yesterday→today request needs
461 ``week``, not ``day``).
462 """
463 covered = days + 1
464 if covered <= 1:
465 return "day"
466 if covered <= 7:
467 return "week"
468 if covered <= 31:
469 return "month"
470 if covered <= 366:
471 return "year"
472 return "all"
473
474
475 def _window_to_time_filter(from_date: str, to_date: str) -> str:
476 """Map a requested YYYY-MM-DD window onto Reddit's coarse `t` param.
477
478 Reddit's ``t=day|week|month`` buckets are rolling windows ending *now*, not
479 calendar spans and not anchored to ``to_date``. Coverage therefore needs:
480
481 1. Span — a yesterday→today request needs more than rolling ``t=day``.
482 2. Historical reach — a one-day request ending two weeks ago still needs a
483 bucket that reaches ``from_date``; span-alone would pick ``week`` and
484 the API would omit the entire requested range.
485
486 Take the wider of the two, independently of retrieval depth.
487 Phase 5 still trims to ``from_date``/``to_date``. Falls back to ``month``
488 if the dates don't parse.
489 """
490 try:
491 start = date.fromisoformat(from_date)
492 end = date.fromisoformat(to_date)
493 except (ValueError, TypeError):
494 return "month"
495 span_days = max(0, (end - start).days)
496 # Age of from_date relative to "today" — Reddit always anchors to now.
497 age_days = max(0, (datetime.now(timezone.utc).date() - start).days)
498 return _days_to_reddit_bucket(max(span_days, age_days))
499
500
501 def search_reddit(
502 topic: str,
503 from_date: str,
504 to_date: str,
505 depth: str = "default",
506 token: str = None,
507 subreddits: List[str] | None = None,
508 ) -> Dict[str, Any]:
509 """Full Reddit search: multi-query global discovery + subreddit drill-down.
510
511 This is the main v3 Reddit entry point.
512
513 Args:
514 topic: Search topic
515 from_date: Start date (YYYY-MM-DD)
516 to_date: End date (YYYY-MM-DD)
517 depth: 'quick', 'default', or 'deep'
518 token: ScrapeCreators API key
519 subreddits: Optional list of subreddit names to search first (pre-resolved)
520
521 Returns:
522 Dict with 'items' list and optional 'error'.
523 """
524 if not token:
525 return {"items": [], "error": "No SCRAPECREATORS_API_KEY configured"}
526
527 config = DEPTH_CONFIG.get(depth, DEPTH_CONFIG["default"])
528 timeframe = _window_to_time_filter(from_date, to_date)
529 intent = infer_query_intent(topic)
530
531 # === Phase 1: Query Expansion ===
532 queries = expand_reddit_queries(topic, depth)
533 _log(f"Expanded '{topic}' into {len(queries)} queries: {queries}")
534
535 core = _extract_core_subject(topic)
536
537 # === Phase 1.5: Pre-resolved subreddit search (high-signal) ===
538 all_raw_posts = []
539 all_items: List[Dict[str, Any]] = []
540 if subreddits:
541 _log(f"Searching pre-resolved subreddits: {subreddits}")
542 with ThreadPoolExecutor(max_workers=min(5, len(subreddits))) as executor:
543 futures = {}
544 for sub in subreddits:
545 futures[http.submit_with_context(
546 executor, _subreddit_search, sub, core, token, "relevance", timeframe,
547 )] = sub
548 for future in as_completed(futures):
549 sub = futures[future]
550 sub_posts = future.result()
551 _log(f" -> {len(sub_posts)} results from pre-resolved r/{sub}")
552 for j, post in enumerate(sub_posts):
553 item = _normalize_post(post, len(all_items) + j + 1, f"r/{sub}", query=core)
554 all_items.append(item)
555
556 # === Phase 2: Global Discovery ===
557 max_global = config["global_searches"]
558
559 with ThreadPoolExecutor(max_workers=max_global or 1) as executor:
560 futures = {}
561 for i, query in enumerate(queries[:max_global]):
562 # Product/comparison queries: sort=top surfaces high-engagement posts
563 # from relevant communities instead of keyword-matched noise.
564 sort = "top" if intent in ("product", "comparison") else ("relevance" if i == 0 else "top")
565 _log(f"Global search {i+1}/{max_global}: '{query}' (sort={sort})")
566 futures[http.submit_with_context(
567 executor, _global_search, query, token, sort, timeframe,
568 )] = query
569 for future in as_completed(futures):
570 query = futures[future]
571 posts = future.result()
572 _log(f" -> {len(posts)} results for '{query}'")
573 all_raw_posts.extend(posts)
574
575 # Normalize all posts (with query for relevance scoring)
576 for i, post in enumerate(all_raw_posts):
577 item = _normalize_post(post, i + 1, "global", query=core)
578 all_items.append(item)
579
580 # === Phase 3: Subreddit Discovery + Targeted Search ===
581 subreddit_budget = 0 if intent == "how_to" else config["subreddit_searches"]
582 discovered_subs = discover_subreddits(all_raw_posts, topic=topic, max_subs=subreddit_budget)
583 _log(f"Discovered subreddits: {discovered_subs}")
584
585 subreddit_limit = subreddit_budget
586 if subreddit_limit > 0:
587 with ThreadPoolExecutor(max_workers=subreddit_limit) as executor:
588 futures = {}
589 for sub in discovered_subs[:subreddit_limit]:
590 _log(f"Subreddit search: r/{sub} for '{core}'")
591 futures[http.submit_with_context(
592 executor, _subreddit_search, sub, core, token, "relevance", timeframe,
593 )] = sub
594 for future in as_completed(futures):
595 sub = futures[future]
596 sub_posts = future.result()
597 _log(f" -> {len(sub_posts)} results from r/{sub}")
598 for j, post in enumerate(sub_posts):
599 item = _normalize_post(post, len(all_items) + j + 1, f"r/{sub}", query=core)
600 all_items.append(item)
601
602 # === Phase 4: Deduplicate ===
603 all_items = _dedupe_posts(all_items)
604 _log(f"After dedup: {len(all_items)} unique posts")
605
606 # === Phase 5: Date filter ===
607 in_range = []
608 out_of_range = 0
609 for item in all_items:
610 if item["date"] and from_date <= item["date"] <= to_date:
611 in_range.append(item)
612 elif item["date"] is None:
613 in_range.append(item) # Keep unknown dates
614 else:
615 out_of_range += 1
616
617 if in_range:
618 all_items = in_range
619 if out_of_range:
620 _log(f"Filtered {out_of_range} posts outside date range")
621 else:
622 _log(f"No posts within date range, keeping all {len(all_items)}")
623
624 # === Phase 6: Relevance floor + relevance-weighted ranking ===
625 # Drop the off-topic tail when enough on-topic posts remain (guard mirrors
626 # the date filter: keep all if too few clear the floor). When too few clear
627 # the soft floor, still strip zero-overlap posts (relevance exactly 0 = no
628 # title/body token match at all, never on-topic) whenever anything relevant
629 # remains, so viral high-upvote junk can't fill the section. Then rank by
630 # relevance with a bounded engagement bonus (see RELEVANCE_FLOOR note above).
631 before = len(all_items)
632 on_topic = [it for it in all_items if (it.get("relevance") or 0) >= RELEVANCE_FLOOR]
633 if len(on_topic) >= MIN_ON_TOPIC:
634 all_items = on_topic
635 else:
636 nonzero = [it for it in all_items if (it.get("relevance") or 0) > 0]
637 if nonzero:
638 all_items = nonzero
639 if len(all_items) < before:
640 _log(f"Relevance floor dropped {before - len(all_items)} off-topic posts")
641 all_items.sort(key=_relevance_rank_key, reverse=True)
642
643 # Re-index IDs
644 for i, item in enumerate(all_items):
645 item["id"] = f"R{i+1}"
646
647 _log(f"Final: {len(all_items)} Reddit posts")
648 return {"items": all_items}
649
650
651 def enrich_with_comments(
652 items: List[Dict[str, Any]],
653 token: str,
654 depth: str = "default",
655 budget_seconds: int = 60,
656 ) -> List[Dict[str, Any]]:
657 """Enrich top items with comment data from ScrapeCreators.
658
659 Args:
660 items: Reddit items from search_reddit()
661 token: ScrapeCreators API key
662 depth: Depth for comment limit
663 budget_seconds: Maximum total time for enrichment. If exceeded,
664 returns items with whatever enrichment completed. Never discards items.
665
666 Returns:
667 Items with top_comments and comment_insights added.
668 """
669 config = DEPTH_CONFIG.get(depth, DEPTH_CONFIG["default"])
670 max_comments = config["comment_enrichments"]
671
672 if not items or not token or max_comments <= 0:
673 return items
674
675 # Select the top threads by total engagement (upvotes + comment count),
676 # not by list position. This ensures high-comment threads like [FRESH ALBUM]
677 # always get enriched even if their upvote score is low.
678 ranked = sorted(items, key=_total_engagement, reverse=True)
679 top_items = ranked[:max_comments]
680 _log(f"Enriching comments for {len(top_items)} posts (by total engagement)")
681
682 start = time.monotonic()
683
684 with ThreadPoolExecutor(max_workers=min(4, len(top_items))) as executor:
685 futures = {
686 http.submit_with_context(
687 executor, fetch_post_comments, item.get("url", ""), token,
688 ): item
689 for item in top_items
690 if item.get("url")
691 }
692
693 # Wait with budget instead of unbounded as_completed
694 remaining = max(0, budget_seconds - (time.monotonic() - start))
695 done, not_done = futures_wait(futures, timeout=remaining)
696
697 enriched_count = 0
698 for future in done:
699 item = futures[future]
700 try:
701 raw_comments = future.result(timeout=0)
702 except Exception:
703 continue
704 if not raw_comments:
705 continue
706
707 top_comments = []
708 insights = []
709
710 for ci, c in enumerate(raw_comments[:10]):
711 body = c.get("body", "")
712 if not body or body in ("[deleted]", "[removed]"):
713 continue
714
715 score = c.get("ups") or c.get("score", 0)
716 author = c.get("author", "[deleted]")
717 permalink = c.get("permalink", "")
718 comment_url = f"https://reddit.com{permalink}" if permalink else ""
719
720 max_excerpt = 400 if ci == 0 else 300
721 top_comments.append({
722 "score": score,
723 "date": _parse_date(c.get("created_utc")),
724 "author": author,
725 "excerpt": body[:max_excerpt],
726 "url": comment_url,
727 })
728
729 if len(body) >= 30 and author not in ("[deleted]", "[removed]", "AutoModerator"):
730 insight = body[:150]
731 if len(body) > 150:
732 for i, char in enumerate(insight):
733 if char in '.!?' and i > 50:
734 insight = insight[:i+1]
735 break
736 else:
737 insight = insight.rstrip() + "..."
738 insights.append(insight)
739
740 top_comments.sort(key=lambda c: c.get("score", 0), reverse=True)
741 item["top_comments"] = top_comments[:10]
742 item["comment_insights"] = insights[:10]
743 enriched_count += 1
744
745 if not_done:
746 _log(f"Enrichment budget hit ({budget_seconds}s): {enriched_count}/{len(futures)} posts enriched, {len(not_done)} skipped")
747 for future in not_done:
748 future.cancel()
749 else:
750 elapsed = time.monotonic() - start
751 _log(f"Enriched {enriched_count}/{len(futures)} posts in {elapsed:.1f}s")
752
753 return items
754
755
756 def search_and_enrich(
757 topic: str,
758 from_date: str,
759 to_date: str,
760 depth: str = "default",
761 token: str = None,
762 subreddits: List[str] | None = None,
763 ) -> Dict[str, Any]:
764 """Full Reddit pipeline: search + comment enrichment.
765
766 This is the convenience function that does everything.
767
768 Args:
769 topic: Search topic
770 from_date: Start date (YYYY-MM-DD)
771 to_date: End date (YYYY-MM-DD)
772 depth: 'quick', 'default', or 'deep'
773 token: ScrapeCreators API key
774 subreddits: Optional list of subreddit names to search first (pre-resolved)
775
776 Returns:
777 Dict with 'items' list. Items include top_comments and comment_insights.
778 """
779 result = search_reddit(topic, from_date, to_date, depth, token, subreddits=subreddits)
780 items = result.get("items", [])
781
782 if items and token:
783 items = enrich_with_comments(items, token, depth)
784 result["items"] = items
785
786 return result
787
788
789 # Run-scoped memo for the paid ScrapeCreators Reddit call (R9). All subquery
790 # streams share the raw topic, and the thin-source retry repeats it, so without
791 # this one run could pay for the same query five times. Concurrent callers for
792 # one key wait on the first call; results AND failures are kept until the
793 # per-command reset, so a failed backfill is never retried within the run.
794 _SC_MEMO: Dict[tuple, tuple] = {}
795 _SC_INFLIGHT: Dict[tuple, threading.Event] = {}
796 _SC_MEMO_LOCK = threading.Lock()
797
798
799 def reset_scrapecreators_memo() -> None:
800 """Forget memoized ScrapeCreators Reddit calls. Called once per command, and by tests."""
801 with _SC_MEMO_LOCK:
802 _SC_MEMO.clear()
803 _SC_INFLIGHT.clear()
804
805
806 def _sc_memo_outcome(entry: tuple, replay_failures: bool = True) -> Dict[str, Any]:
807 ok, value, failures = entry
808 if replay_failures:
809 # search_and_enrich swallows ScrapeCreators HTTP errors into the
810 # caller's failure sink; a later caller must see them too, or an
811 # empty cached result reads as a clean no-results.
812 for failure in failures:
813 http._record_failure(failure)
814 if not ok:
815 raise value
816 # Each caller gets its own copy so one stream cannot mutate another's items.
817 return copy.deepcopy(value)
818
819
820 def search_and_enrich_memo(
821 topic: str,
822 from_date: str,
823 to_date: str,
824 depth: str = "default",
825 token: str = None,
826 subreddits: List[str] | None = None,
827 ) -> Dict[str, Any]:
828 """:func:`search_and_enrich`, at most once per key per command.
829
830 Keyed by query, date window, depth, and sorted subreddits.
831 """
832 key = (topic, from_date, to_date, depth, tuple(sorted(subreddits or ())))
833 while True:
834 with _SC_MEMO_LOCK:
835 entry = _SC_MEMO.get(key)
836 if entry is None:
837 gate = _SC_INFLIGHT.get(key)
838 owner = gate is None
839 if owner:
840 gate = threading.Event()
841 _SC_INFLIGHT[key] = gate
842 if entry is not None:
843 return _sc_memo_outcome(entry)
844 if owner:
845 break
846 gate.wait()
847 # Loop: read the owner's cached outcome, or re-elect if it was
848 # interrupted before caching one (e.g. KeyboardInterrupt).
849 try:
850 with http.tee_failures() as recorded:
851 try:
852 result = search_and_enrich(
853 topic, from_date, to_date, depth=depth, token=token,
854 subreddits=subreddits,
855 )
856 outcome = (True, result)
857 except Exception as exc:
858 outcome = (False, exc)
859 entry = (*outcome, list(recorded))
860 with _SC_MEMO_LOCK:
861 _SC_MEMO[key] = entry
862 finally:
863 with _SC_MEMO_LOCK:
864 _SC_INFLIGHT.pop(key, None)
865 gate.set()
866 # The owner's failures already reached its own sink through tee_failures.
867 return _sc_memo_outcome(entry, replay_failures=False)
868
869
870 def parse_reddit_response(response: Dict[str, Any]) -> List[Dict[str, Any]]:
871 """Parse ScrapeCreators response to item list.
872
873 Parse raw Reddit search output into the generic item shape.
874 """
875 return response.get("items", [])
876
876 lines PYTHON