"""Base scraping tools — one-shot page scrapers shared across every agent workflow.""" from __future__ import annotations import asyncio import functools import re import time from typing import Any, cast import structlog from camoufox.async_api import AsyncCamoufox from playwright.async_api import Error as PlaywrightError from marketing_intelligence.browser.base import ( BrowserAdapter, PageSnapshot, create_adapter, create_adapter_from_session, ) from marketing_intelligence.core.config import settings from marketing_intelligence.evasion.detection import PageSignals from marketing_intelligence.evasion.proxy import ProxyManager from marketing_intelligence.evasion.session import get_session from marketing_intelligence.mcp.app import mcp from marketing_intelligence.mcp.tools.guard import tool_guard from marketing_intelligence.mcp.tools.responses import ( ArtistProfileResponse, BatchScrapeResponse, BatchVideoResult, SoundPageResponse, SoundUrlResponse, TagPageResponse, UrlScrapeResponse, VideoScrapeResponse, ) from marketing_intelligence.scraper.extractor import TikTokExtractor from marketing_intelligence.scraper.url_builder import sound_url, tag_url logger = structlog.get_logger("mcp.tools") _STOP_MARKERS = frozenset( {"Log in", "You may like", "Remote-access", "Drag the slider", "Got it"} ) _COMMENT_SKIP = frozenset( {"Reply", "See translation", "View more", "View more replies"} ) _claimed_video_ids: dict[tuple[str, str, str], set[str]] = {} _claim_lock = asyncio.Lock() async def claim_video_id( campaign_id: str, run_id: str, track_name: str, video_id: str ) -> bool: """Atomically claim a video_id for (campaign, run, track) so parallel sibling sessions scraping the same track skip it instead of double-processing. False if already claimed.""" key = (campaign_id, run_id, track_name) async with _claim_lock: claimed = _claimed_video_ids.setdefault(key, set()) if video_id in claimed: return False claimed.add(video_id) return True async def release_video_id( campaign_id: str, run_id: str, track_name: str, video_id: str ) -> None: """Release a claim after a failed processing attempt so it isn't lost forever.""" key = (campaign_id, run_id, track_name) async with _claim_lock: _claimed_video_ids.get(key, set()).discard(video_id) async def wait_for_comments_stable( page: Any, max_iterations: int = 20, min_iterations: int = 3 ) -> None: """Poll the comment panel until its content stops growing, instead of a fixed sleep. min_iterations guards against a real race: on a slow-loading page the very first check can see 0 comments (nothing rendered yet), and 0 == 0 would trivially look "stable" — exiting after 1.5s despite the page not having loaded at all. Found live 2026-07-07 (a video whose vision-description call also came back "not_loaded" the same turn). Never treat a run of zeros as stable before min_iterations has elapsed. """ prev_count = 0 for i in range(max_iterations): await page.evaluate("""() => { const sel = [ 'div[class*="CommentList"]', 'div[class*="comment-list"]', 'div[class*="DivCommentListContainer"]', ].join(', '); const panels = document.querySelectorAll(sel); for (const p of panels) { if (p.scrollHeight > p.clientHeight) { p.scrollTop = p.scrollHeight; return; } } window.scrollTo(0, document.body.scrollHeight); }""") await asyncio.sleep(1.5) body_check = await page.inner_text("body") current_count = body_check.count("Reply") if current_count == prev_count and ( current_count > 0 or i + 1 >= min_iterations ): break prev_count = current_count async def click_and_await_comments( page: Any, action: Any, timeout: int = 8000 ) -> None: # NOSONAR """Perform `action` while listening for TikTok's internal comment-list API response. `/api/comment/list/` is a precise signal that the first batch of comments has arrived from the server. Falls back silently if no such response arrives within `timeout` — action has already run either way. Caller should still run `wait_for_comments_stable()` afterward to catch scroll-triggered pagination (this only confirms the FIRST batch loaded). """ try: async with page.expect_response( lambda r: "/api/comment/list/" in r.url, timeout=timeout ) as resp_info: await action() await resp_info.value except PlaywrightError, asyncio.TimeoutError: pass def _extract_comments_from_body(body: str) -> list[str]: """Parse comment texts out of a TikTok page body string.""" section: list[str] = [] in_comments = False for line in body.split("\n"): t = line.strip() if t.endswith("comments"): in_comments = True continue if in_comments and any(t.startswith(m) for m in _STOP_MARKERS): break if in_comments and t: section.append(t) return _parse_comments(section) def _is_comment_date(s: str) -> bool: return bool( re.match(r"^\d+[dwmhsy]+\s*ago$", s) or re.match(r"^\d{4}-\d{1,2}-\d{1,2}$", s) or re.match(r"^\d{1,2}-\d{1,2}$", s) ) def _find_comment_candidate(section: list[str], i: int) -> str | None: for back in (1, 2): if i - back < 0: break candidate = section[i - back] if candidate.startswith("·") or candidate in _COMMENT_SKIP: continue if " " in candidate: return candidate break return None def _parse_comments(section: list[str]) -> list[str]: """Date-anchored extraction: for each date line, take the closest preceding line that has a space (comment), skipping badges (·) and known noise tokens.""" results: list[str] = [] seen: set[str] = set() for i, line in enumerate(section): if not _is_comment_date(line): continue candidate = _find_comment_candidate(section, i) if candidate and candidate not in seen: results.append(candidate) seen.add(candidate) return results class ScrapingService: """Manages browser adapters and page scraping.""" def get_adapter( self, proxy: dict | None = None, session_id: str | None = None ) -> BrowserAdapter: """Return a BrowserAdapter for the given session or proxy settings.""" if session_id: session = get_session(session_id) if session: return create_adapter_from_session(session, session.proxy or proxy) if proxy is None and (settings.proxy_list or settings.proxy_lambda_name): try: proxy = ProxyManager().acquire().to_playwright_dict() except (OSError, RuntimeError, ValueError) as e: logger.warning("proxy.acquire_failed", error=str(e)[:200]) return create_adapter( engine=settings.browser_engine, profile_name=settings.browser_profile, proxy=proxy, human_behavior=settings.browser_human_behavior, ) def signals( self, snap: PageSnapshot, url: str, *, items_loaded: int = 0, count_text: str | None = None, ) -> PageSignals: """Build page-health signals from a fetched snapshot.""" return PageSignals.from_page( html=snap.html, requested_url=url, final_url=snap.final_url or url, response_time_ms=snap.response_time_ms, items_loaded=items_loaded, count_text=count_text, ) async def tag_page( self, tag: str, proxy: dict | None = None, session_id: str | None = None ) -> TagPageResponse: """Scrape a hashtag page and extract video references.""" url = tag_url(tag) snap = await self.get_adapter(proxy, session_id).fetch_with_scroll(url) videos, video_count_text = TikTokExtractor.tag_videos(snap.html) return TagPageResponse( tag=tag.lstrip("#"), tag_url=url, video_count_text=video_count_text, videos=videos, signals=self.signals( snap, url, items_loaded=len(videos), count_text=video_count_text ), ) async def sound_page( self, sound_id: str, proxy: dict | None = None, session_id: str | None = None ) -> SoundPageResponse: """Scrape a sound page and extract video URLs.""" url = sound_url(sound_id) snap = await self.get_adapter(proxy, session_id).fetch_with_scroll(url) title = snap.title.replace(" | TikTok", "").strip() if snap.title else None video_urls = TikTokExtractor.sound_video_urls(snap.html) video_count_text = TikTokExtractor.video_count(snap.html) return SoundPageResponse( sound_id=sound_id, sound_page_url=url, title=title, video_count_text=video_count_text, video_urls=video_urls, signals=self.signals( snap, url, items_loaded=len(video_urls), count_text=video_count_text ), ) async def video( self, url: str, proxy: dict | None = None, session_id: str | None = None ) -> VideoScrapeResponse: """Scrape a single TikTok video page via Camoufox.""" if proxy is None and session_id: sess = get_session(session_id) if sess: proxy = sess.proxy author, video_id = TikTokExtractor.author_and_video_id(url) t0 = time.monotonic() proxy_server = (proxy or {}).get("server", "none") logger.info( "video.browser_start", url=url, proxy=proxy_server, headless=settings.browser_headless, ) async with AsyncCamoufox( headless=settings.browser_headless, proxy=proxy ) as browser: page = await browser.new_page() logger.info("video.page_navigating", url=url) await page.goto(url, wait_until="domcontentloaded", timeout=30000) await asyncio.sleep(5) html = await page.content() metrics = TikTokExtractor.video_metrics(html) sound_links = await page.eval_on_selector_all( 'a[href*="/music/"]', "els => els.map(e => ({href: e.href, text: e.textContent}))", ) first_link = sound_links[0] if sound_links else None sound_url_val = first_link["href"] if first_link else None sound_title = ( (first_link.get("text") or "").strip() or None if first_link else None ) sid = ( TikTokExtractor.sound_id_from_url(sound_url_val) if sound_url_val else None ) caption_el = await page.query_selector('[data-e2e="video-desc"]') if not caption_el: caption_el = await page.query_selector('span[data-e2e*="desc"]') caption = ( await caption_el.inner_text() if caption_el else metrics.get("caption_json") ) metrics.pop("caption_json", None) logger.info( "video.page_loaded", url=url, views=metrics.get("views"), likes=metrics.get("likes"), ) await page.evaluate("""() => { const spans = document.querySelectorAll('span'); for (const el of spans) { if (el.innerText?.trim() === 'Comments') { const rect = el.getBoundingClientRect(); if (rect.width > 0 && rect.x > 400) { el.click(); return; } } } }""") await asyncio.sleep(4) await wait_for_comments_stable(page) body = await page.inner_text("body") final_html = await page.content() response_time_ms = int((time.monotonic() - t0) * 1000) comment_texts = _extract_comments_from_body(body) logger.info("video.comments_extracted", url=url, comments=len(comment_texts)) signals = PageSignals.from_page( html=final_html, requested_url=url, final_url=url, response_time_ms=response_time_ms, metrics=metrics, ) return VideoScrapeResponse( url=url, video_id=video_id, author=author, caption=caption, sound_url=sound_url_val, sound_id=sid, sound_title=sound_title, comment_texts=comment_texts[:15], signals=signals, **metrics, ) async def video_batch( self, urls: list[str], session_id: str | None = None ) -> BatchScrapeResponse: """Scrape a list of video URLs, collecting each result individually.""" results: list[BatchVideoResult] = [] for url in urls: try: res = await self.video(url, session_id=session_id) results.append( BatchVideoResult( url=url, success=res.success, video_id=res.video_id, views=res.views, likes=res.likes, comments=res.comments, shares=res.shares, favorites=res.favorites, sound_id=res.sound_id, error=res.error, ) ) except (PlaywrightError, asyncio.TimeoutError, OSError, RuntimeError) as e: results.append( BatchVideoResult(url=url, success=False, error=str(e)[:200]) ) failed = sum(1 for r in results if not r.success) return BatchScrapeResponse( results=results, count=len(results), failed_count=failed ) async def raw_page( self, url: str, proxy: dict | None = None, session_id: str | None = None ) -> UrlScrapeResponse: """Fetch any TikTok page and return its text content.""" async def _get_text( page: Any, _html: str, final_url: str, response_time_ms: int ) -> UrlScrapeResponse: await page.evaluate("window.scrollTo(0, document.body.scrollHeight / 2)") await asyncio.sleep(1) text = await page.inner_text("body") signals = PageSignals.from_page( html=_html, requested_url=url, final_url=final_url, response_time_ms=response_time_ms, ) return UrlScrapeResponse(url=url, content=text, signals=signals) return cast( UrlScrapeResponse, await self.get_adapter(proxy, session_id).fetch_interactive(url, _get_text), ) async def artist_profile( self, handle: str, proxy: dict | None = None, session_id: str | None = None ) -> ArtistProfileResponse: """Scrape a TikTok artist profile page.""" url = f"https://www.tiktok.com/@{handle.lstrip('@')}" snap = await self.get_adapter(proxy, session_id).fetch_with_scroll(url) profile = TikTokExtractor.artist_profile(snap.html) if not profile: return ArtistProfileResponse( success=False, handle=handle, error="Could not extract profile data", signals=self.signals(snap, url), ) return ArtistProfileResponse( handle=handle, nickname=profile.get("nickname"), followers=profile.get("followers"), following=profile.get("following"), likes=profile.get("likes"), video_count=profile.get("video_count"), signals=self.signals(snap, url), ) @staticmethod def build_sound_url(sound_id: str) -> SoundUrlResponse: """Build a TikTok sound page URL from a sound_id.""" return SoundUrlResponse(sound_id=sound_id, url=sound_url(sound_id)) @functools.lru_cache(maxsize=None) def _get_svc() -> ScrapingService: """Return the singleton ScrapingService instance.""" return ScrapingService() # ── MCP tool wrappers ──────────────────────────────────────────────────────── @mcp.tool() @tool_guard(TagPageResponse, echo=("tag",)) async def scrape_tag_page( tag: str, proxy: dict | None = None, session_id: str | None = None ) -> TagPageResponse: """Scrape a TikTok hashtag page. Returns post list {video_id, video_url, author}.""" return await _get_svc().tag_page(tag, proxy, session_id) @mcp.tool() @tool_guard(SoundPageResponse, echo=("sound_id",)) async def scrape_sound_page( sound_id: str, proxy: dict | None = None, session_id: str | None = None ) -> SoundPageResponse: """Scrape a TikTok sound page by sound_id. Returns video URLs.""" return await _get_svc().sound_page(sound_id, proxy, session_id) @mcp.tool() @tool_guard(VideoScrapeResponse, echo=("url",)) async def scrape_video( url: str, proxy: dict | None = None, session_id: str | None = None ) -> VideoScrapeResponse: """Scrape a TikTok video page. Returns views, likes, comments, shares, favorites, caption, hashtags, sound_id. """ return await _get_svc().video(url, proxy, session_id) @mcp.tool() @tool_guard(BatchScrapeResponse, echo=()) async def scrape_posts_batch( urls: list[str], session_id: str | None = None ) -> BatchScrapeResponse: """Scrape multiple TikTok video URLs in one call. Use for watch runs.""" return await _get_svc().video_batch(urls, session_id) @mcp.tool() @tool_guard(UrlScrapeResponse, echo=("url",)) async def scrape_url( url: str, proxy: dict | None = None, session_id: str | None = None ) -> UrlScrapeResponse: """Scrape any public TikTok page and return raw text.""" return await _get_svc().raw_page(url, proxy, session_id) @mcp.tool() @tool_guard(ArtistProfileResponse, echo=("handle",)) async def scrape_artist_profile( handle: str, proxy: dict | None = None, session_id: str | None = None ) -> ArtistProfileResponse: """Scrape a TikTok artist profile: followers, likes, video_count.""" return await _get_svc().artist_profile(handle, proxy, session_id) @mcp.tool() @tool_guard(SoundUrlResponse, echo=("sound_id",)) async def build_sound_url(sound_id: str) -> SoundUrlResponse: """Build a TikTok sound page URL from a sound ID.""" return _get_svc().build_sound_url(sound_id)