| 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 |