返回 last30days-skill
reddit_search.py
根目录 / skills / last30days / scripts / lib / reddit_search.py
1 """Keyless Reddit discovery via Reddit's own site search fragment.
2
3 ``/svc/shreddit/search/?q=...&type=posts`` is the server-rendered HTML the
4 reddit.com search page loads. It serves HTTP 200 with no API key, and
5 ``/svc/shreddit/r/{sub}/search/`` restricts the same search to one subreddit.
6 Each post result unit carries:
7
8 - a ``<search-telemetry-tracker>`` whose JSON tracking context has
9 ``action_info.type == "post"``, the post id, title, and subreddit,
10 - the post-title link (the canonical permalink),
11 - a ``<faceplate-timeago ts=...>`` post date,
12 - a counter row of ``<faceplate-number>`` votes and comments.
13
14 So every discovered post arrives dated and scored, with no extra listing fetch.
15 The next page is a lazy ``<faceplate-partial>`` whose ``src`` carries a
16 ``cursor``. Output dicts match ``reddit_listing.parse_cards`` so downstream
17 code is unaffected.
18 """
19
20 import html as _html
21 import json
22 import re
23 import sys
24 from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
25 from typing import Any, Dict, List, Optional, Tuple
26 from urllib.parse import parse_qs, quote, urlencode, urlsplit
27
28 from . import http, reddit_listing
29 from .relevance import token_overlap_relevance
30
31 BASE = "https://www.reddit.com"
32
33 # Post caps per run by depth.
34 DEPTH_LIMITS = {"quick": 10, "default": 25, "deep": 50}
35 # Global search pages per depth (about 7 posts per page).
36 GLOBAL_PAGE_CAPS = {"quick": 2, "default": 4, "deep": 8}
37 # Targeted subreddits get one page each.
38 SUB_PAGE_CAP = 1
39 MAX_WORKERS = 4
40 SEARCH_TIMEOUT = 15
41
42 # A recognizable page has post result units or Reddit's explicit no-results
43 # unit. Anything else (challenge page, changed markup) is a lane failure.
44 RESULTS_MARKER = 'view-events="search/view/post"'
45 NO_RESULTS_MARKER = 'view-events="search/view/no_results"'
46
47 _TRACKER = re.compile(
48 r'<search-telemetry-tracker\b[^>]*?\bdata-faceplate-tracking-context="([^"]*)"'
49 )
50 _TITLE_HREF = re.compile(
51 r'<a\b[^>]*?data-testid="post-title"[^>]*?\bhref="([^"]*)"'
52 r'|<a\b[^>]*?\bhref="([^"]*)"[^>]*?data-testid="post-title"'
53 )
54 _TIMEAGO = re.compile(r'<faceplate-timeago\b[^>]*?\bts="([^"]+)"')
55 _COUNTER = re.compile(
56 r'<faceplate-number\b[^>]*?\bnumber="(\d+)"[^>]*>\s*</faceplate-number>\s*(vote|comment)'
57 )
58 _NEXT_PAGE = re.compile(
59 r'<faceplate-partial\b[^>]*?\bsrc="([^"]*/svc/shreddit/[^"]*search/[^"]*)"'
60 )
61
62
63 def _log(msg: str) -> None:
64 sys.stderr.write(f"[RedditSearch] {msg}\n")
65 sys.stderr.flush()
66
67
68 def unrecognized_body(text: str) -> Optional[str]:
69 """Validator for the fetch hook: None for a search page, else a reason."""
70 if NO_RESULTS_MARKER in text:
71 return None
72 if RESULTS_MARKER in text:
73 # Marker kept but tracking context changed shape: a clean zero-result
74 # parse would hide the drift, so treat it as a lane failure.
75 if _post_units(text):
76 return None
77 return "Reddit search result units no longer parse"
78 return "not a Reddit search results page"
79
80
81 def search_url(
82 query: str,
83 time_filter: str = "month",
84 subreddit: Optional[str] = None,
85 cursor: Optional[str] = None,
86 ) -> str:
87 """Build a deterministic search URL so concurrent streams share the memo."""
88 if subreddit:
89 sub = subreddit.strip().removeprefix("r/").strip("/")
90 path = f"/svc/shreddit/r/{quote(sub, safe='')}/search/"
91 else:
92 path = "/svc/shreddit/search/"
93 params = [("q", query), ("type", "posts"), ("t", time_filter), ("sort", "relevance")]
94 if cursor:
95 params.append(("cursor", cursor))
96 return f"{BASE}{path}?{urlencode(params)}"
97
98
99 def _post_units(html_text: str) -> List[Tuple[int, Dict[str, Any]]]:
100 """(offset, tracking context) for each post-title tracker, in page order."""
101 units = []
102 for m in _TRACKER.finditer(html_text):
103 try:
104 ctx = json.loads(_html.unescape(m.group(1)))
105 except (ValueError, TypeError):
106 continue
107 if not isinstance(ctx, dict):
108 continue
109 action = ctx.get("action_info")
110 post = ctx.get("post")
111 # A field that is no longer an object is drift: skip the unit, and a
112 # page left with none is flagged by unrecognized_body.
113 if not isinstance(action, dict) or not isinstance(post, dict):
114 continue
115 if action.get("type") == "post" and post.get("id"):
116 units.append((m.start(), ctx))
117 return units
118
119
120 def _next_cursor(html_text: str) -> Optional[str]:
121 for m in _NEXT_PAGE.finditer(html_text):
122 query = urlsplit(_html.unescape(m.group(1))).query
123 cursor = (parse_qs(query).get("cursor") or [None])[0]
124 if cursor:
125 return cursor
126 return None
127
128
129 def parse_page(html_text: str, query: str = "") -> Tuple[List[Dict[str, Any]], Optional[str]]:
130 """Parse one search fragment into (normalized posts, next-page cursor)."""
131 html_text = html_text or ""
132 units = _post_units(html_text)
133 posts: List[Dict[str, Any]] = []
134 seen: set = set()
135 for i, (start, ctx) in enumerate(units):
136 post_id = str(ctx["post"]["id"]).removeprefix("t3_")
137 if post_id in seen:
138 continue
139 seen.add(post_id)
140 end = units[i + 1][0] if i + 1 < len(units) else len(html_text)
141 chunk = html_text[start:end]
142
143 sub_ctx = ctx.get("subreddit")
144 subreddit = str(sub_ctx.get("name") or "") if isinstance(sub_ctx, dict) else ""
145 title = str(ctx["post"].get("title") or "")
146 href_match = _TITLE_HREF.search(chunk)
147 href = _html.unescape((href_match.group(1) or href_match.group(2))) if href_match else ""
148 href = href.split("?", 1)[0]
149 if f"/comments/{post_id}/" in href and href.startswith("/r/"):
150 url = f"{BASE}{href}"
151 else:
152 url = f"{BASE}/r/{subreddit}/comments/{post_id}/"
153
154 ts_match = _TIMEAGO.search(chunk)
155 ts = ts_match.group(1) if ts_match else None
156 # First counter of each kind wins: this post's own counters precede any
157 # that belong to a skipped (malformed) neighbour sharing the chunk.
158 counts: Dict[str, int] = {}
159 for number, label in _COUNTER.findall(chunk):
160 counts.setdefault(label, int(number))
161 score, num_comments = counts.get("vote", 0), counts.get("comment", 0)
162
163 posts.append({
164 "id": "",
165 "title": title,
166 "url": url,
167 "score": score,
168 "num_comments": num_comments,
169 "subreddit": subreddit,
170 "created_utc": reddit_listing._to_epoch(ts),
171 "author": "",
172 "selftext": "",
173 "date": reddit_listing._to_date(ts),
174 "engagement": {
175 "score": score,
176 "num_comments": num_comments,
177 "upvote_ratio": None,
178 },
179 "relevance": round(token_overlap_relevance(query, title), 3) if query else 0.0,
180 "why_relevant": "Reddit search",
181 "metadata": {"post_id": post_id},
182 })
183 return posts, _next_cursor(html_text)
184
185
186 def _fetch_page(url: str) -> Tuple[Optional[str], Optional[str]]:
187 """(body, error). Shared limiter, run memo, one 429 retry, failure sink."""
188 try:
189 return http.reddit_keyless_get_text_retry_429(
190 url,
191 timeout=SEARCH_TIMEOUT,
192 accept="text/html",
193 validate=unrecognized_body,
194 )
195 except Exception as e: # defensive: one bad page must not sink the run
196 _log(f"search fetch failed for {url}: {e}")
197 return None, str(e)
198
199
200 def _search_stream(
201 query: str,
202 time_filter: str,
203 max_pages: int,
204 subreddit: Optional[str] = None,
205 ) -> List[Dict[str, Any]]:
206 """Page one search sequentially. Keeps earlier pages when a later one fails."""
207 label = f"r/{subreddit}" if subreddit else "global"
208 posts: List[Dict[str, Any]] = []
209 seen_ids: set = set()
210 used_cursors: set = set()
211 cursor: Optional[str] = None
212 for page in range(1, max_pages + 1):
213 text, error = _fetch_page(search_url(query, time_filter, subreddit, cursor))
214 if text is None:
215 _log(f"{label} page {page} failed: {error or 'no response'}")
216 break
217 page_posts, next_cursor = parse_page(text, query)
218 new = [p for p in page_posts if p["metadata"]["post_id"] not in seen_ids]
219 for p in new:
220 seen_ids.add(p["metadata"]["post_id"])
221 posts.extend(new)
222 if not new or not next_cursor or next_cursor in used_cursors:
223 break
224 used_cursors.add(next_cursor)
225 cursor = next_cursor
226 return posts
227
228
229 def _result_timeout(batch_size: int, pages: int = 1) -> float:
230 """Per-future wait: each sequential page's own timeout plus the bucket's queue."""
231 return pages * (SEARCH_TIMEOUT + 5) + http.reddit_keyless_wait_allowance(batch_size)
232
233
234 def _time_filter(from_date: Optional[str], to_date: Optional[str]) -> str:
235 # Same mapping as the ScrapeCreators path; falls back to "month".
236 from .reddit import _window_to_time_filter
237 return _window_to_time_filter(from_date or "", to_date or "")
238
239
240 def search(
241 query: str,
242 depth: str = "default",
243 subreddits: Optional[List[str]] = None,
244 from_date: Optional[str] = None,
245 to_date: Optional[str] = None,
246 ) -> List[Dict[str, Any]]:
247 """Discover Reddit posts for a query via the site search fragment.
248
249 Runs a paged global search plus one page of per-subreddit search for each
250 targeted subreddit. Returns normalized, dated, scored posts, deduped by URL
251 and capped by depth. Streams are interleaved round-robin before the cap so
252 neither the targeted subreddits nor the global search crowds out the other. Returns ``[]`` on total failure and never raises; failures
253 land in the pipeline's failure sink.
254 """
255 try:
256 limit = DEPTH_LIMITS.get(depth, DEPTH_LIMITS["default"])
257 global_pages = GLOBAL_PAGE_CAPS.get(depth, GLOBAL_PAGE_CAPS["default"])
258 time_filter = _time_filter(from_date, to_date)
259 subs = []
260 for raw in subreddits or []:
261 sub = (raw or "").strip().removeprefix("r/").strip("/")
262 if sub and sub.lower() not in {s.lower() for s in subs}:
263 subs.append(sub)
264
265 jobs = [(sub, SUB_PAGE_CAP) for sub in subs] + [(None, global_pages)]
266 batch = sum(pages for _sub, pages in jobs)
267 streams: List[List[Dict[str, Any]]] = []
268 with ThreadPoolExecutor(max_workers=min(MAX_WORKERS, len(jobs))) as executor:
269 # submit_with_context, not executor.submit: a plain submit starts the
270 # worker with an empty context, dropping the pipeline's
271 # capture_failures() sink so a 429/403 or a challenge page is
272 # silently discarded (issue #899).
273 futures = [
274 (http.submit_with_context(executor, _search_stream, query, time_filter, pages, sub), sub, pages)
275 for sub, pages in jobs
276 ]
277 for future, sub, pages in futures:
278 try:
279 streams.append(future.result(timeout=_result_timeout(batch, pages)))
280 except (Exception, FuturesTimeoutError) as e:
281 _log(f"{'r/' + sub if sub else 'global'} search future failed: {e}")
282
283 # Round-robin across streams, each in its own relevance order, so the
284 # depth cap keeps a share of every stream.
285 results = [
286 stream[i]
287 for i in range(max((len(s) for s in streams), default=0))
288 for stream in streams
289 if i < len(stream)
290 ]
291 seen: set = set()
292 unique: List[Dict[str, Any]] = []
293 for post in results:
294 if post["url"] not in seen:
295 seen.add(post["url"])
296 unique.append(post)
297 unique = unique[:limit]
298 for i, post in enumerate(unique):
299 post["id"] = f"R{i + 1}"
300 _log(f"{len(unique)} posts (t={time_filter}, {len(subs)} targeted subs)")
301 return unique
302 except Exception as e: # defensive: discovery must never raise into the pipeline
303 _log(f"search failed: {e}")
304 return []
305
305 lines PYTHON