"""Overlay-navigation comment streaming for sound pages. Exclusive to discovery_relevant's scrape_relevant_stage (sound path). Opens the first video from the sound page listing, then chains forward via the in-overlay "Next video" button (same strategy as tag_stream) for as long as it keeps working — only falls back to the slower click→Escape/go_back→re-click-from-listing path when the Next button is missing or breaks. """ from __future__ import annotations import asyncio import base64 import contextlib from datetime import UTC, datetime from typing import Any import structlog from anthropic import APIError as _AnthropicAPIError, Anthropic as _Anthropic from playwright.async_api import Error as PlaywrightError from marketing_intelligence.browser.pool import _pool from marketing_intelligence.core.config import settings 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 ScrapeTagCommentsResponse from marketing_intelligence.mcp.tools.scraping_tools.base import ( claim_video_id, click_and_await_comments, release_video_id, wait_for_comments_stable, ) from marketing_intelligence.scraper.extractor import TikTokExtractor from marketing_intelligence.scraper.url_builder import sound_url from marketing_intelligence.storage.factory import get_backend logger = structlog.get_logger("mcp.tools") async def _return_to_sound_page(page: Any, entry_url: str) -> bool: """Close overlay and confirm we are back on the sound listing page.""" # 1. Try Escape try: await page.keyboard.press("Escape") await page.wait_for_function( "() => !window.location.href.includes('/video/')", timeout=5000, ) except PlaywrightError: # 2. Escape didn't move URL — try go_back (faster, no network re-fetch if cached) try: await page.go_back(wait_until="domcontentloaded", timeout=15000) await page.wait_for_function( "() => !window.location.href.includes('/video/')", timeout=5000, ) except PlaywrightError: # 3. Last resort: full reload of entry URL try: await page.goto(entry_url, wait_until="domcontentloaded", timeout=30000) except (PlaywrightError, asyncio.TimeoutError) as e: logger.warning("sound_stream.return_failed", error=str(e)[:100]) return False try: await page.wait_for_selector('a[href*="/video/"]', timeout=15000) except PlaywrightError: return False return True async def _find_next_sound_href(page: Any, seen: set[str]) -> str | None: """Find first video href on the current (listing) page that we haven't processed yet.""" for _ in range(6): hrefs: list[str] = await page.evaluate("""() => Array.from(document.querySelectorAll('a[href*="/video/"]')) .map(a => a.href) .filter(h => /\\/video\\/\\d+/.test(h)) """) for href in hrefs: _, vid = TikTokExtractor.author_and_video_id(href) if vid and vid not in seen: return href # All visible videos already processed — scroll for more await page.evaluate("window.scrollTo(0, document.body.scrollHeight)") await asyncio.sleep(2) return None async def _click_next_in_sound_overlay(page: Any) -> bool: """Advance to the next video via the in-overlay Next button. Returns False if the button is missing or the click fails.""" next_btn = page.locator( 'button[aria-label="Go to next video"], button[aria-label="Next video"]' ) if await next_btn.count() == 0: return False try: await click_and_await_comments(page, next_btn.first.click) await page.wait_for_function( "() => window.location.href.includes('/video/')", timeout=10000, ) await wait_for_comments_stable(page) return True except (PlaywrightError, asyncio.TimeoutError) as e: logger.warning("sound_stream.next_btn_failed", error=str(e)[:100]) return False async def _process_sound_video( page: Any, vid: str, current_url: str, sound_id: str, original_sound_id: str | None, campaign_id: str, run_id: str, track_name: str | None, location: str | None, backend: Any, ) -> tuple[int, int, str | None]: """Scrape comments, caption, and vision for one sound video and write to storage. Calls release_video_id on the except path. Returns (stored_delta, errors_delta, video_id_appended).""" try: comment_texts = await page.evaluate("""() => { const panels = document.querySelectorAll( 'div[class*="CommentList"], div[class*="comment-list"], div[class*="DivCommentListContainer"]' ); for (const p of panels) { if (p.scrollHeight > p.clientHeight) p.scrollTop = p.scrollHeight; } const skip = new Set(["Reply", "See translation", "View more", "View more replies", "Add comment..."]); const dedup = new Set(); const results = []; const candidates = document.querySelectorAll( '[data-e2e="comment-level-1"] p, ' + '[class*="CommentItem"] p, ' + '[class*="comment-item"] p, ' + '[class*="DivCommentContentContainer"] p, ' + '[class*="CommentText"] span, ' + '[data-e2e="comment-text"]' ); for (const el of candidates) { const t = (el.innerText || el.textContent || "").trim(); if (t && t.length > 1 && !skip.has(t) && !dedup.has(t)) { dedup.add(t); results.push(t); } } return results; }""") logger.info("sound_stream.comments", video_id=vid, count=len(comment_texts)) caption = None hashtags: list[str] = [] try: cap_data = await page.evaluate("""() => { const captionEl = document.querySelector( '[data-e2e="video-desc"], span[data-e2e*="desc"], [class*="DivVideoDescription"] span' ); const caption = captionEl ? (captionEl.innerText || captionEl.textContent || "").trim() || null : null; const tagEls = document.querySelectorAll('a[href*="/tag/"], [data-e2e="video-tag"]'); const hashtags = Array.from(tagEls) .map(el => (el.innerText || el.textContent || "").trim()) .filter(Boolean); return { caption, hashtags }; }""") caption = cap_data.get("caption") hashtags = cap_data.get("hashtags", []) except (PlaywrightError, AttributeError, ValueError) as cap_err: logger.warning( "sound_stream.caption_failed", video_id=vid, error=str(cap_err)[:100], ) video_description = None try: with contextlib.suppress(Exception): await page.wait_for_function( "() => { const v = document.querySelector('video'); " "return v && v.readyState >= 1 && !!v.currentSrc; }", timeout=4000, ) shot = await page.screenshot(type="jpeg", quality=50) b64 = base64.b64encode(shot).decode() _anth = _Anthropic(api_key=settings.anthropic_api_key) resp = _anth.messages.create( model="claude-haiku-4-5-20251001", max_tokens=60, messages=[ { "role": "user", "content": [ { "type": "image", "source": { "type": "base64", "media_type": "image/jpeg", "data": b64, }, }, { "type": "text", "text": ( "What is happening in this TikTok video? " "Reply with ONE plain sentence, no markdown, no heading. " "If the video has not loaded" " (black/blank screen or loading spinner)," " reply with exactly: not_loaded" ), }, ], } ], ) text_block = next((b for b in resp.content if hasattr(b, "text")), None) raw = ( getattr(text_block, "text", "not_loaded").strip() if text_block else "not_loaded" ) video_description = None if raw.lower().startswith("not_loaded") else raw logger.info( "sound_stream.video_described", video_id=vid, description=(video_description or "not_loaded")[:80], ) except ( _AnthropicAPIError, PlaywrightError, asyncio.TimeoutError, OSError, ) as vis_err: logger.warning( "sound_stream.vision_failed", video_id=vid, error=str(vis_err)[:100], ) is_original = bool(original_sound_id and sound_id == original_sound_id) await backend.write_video_comments( campaign_id, run_id, vid or "", { "video_id": vid, "url": current_url, "track_name": track_name, "caption": caption, "hashtags": hashtags, "sound_id": sound_id, "is_original_sound": is_original, "location": location, "video_description": video_description, "comment_count": len(comment_texts), "comment_texts": comment_texts, "observed_at": datetime.now(UTC).isoformat(), }, ) return 1, 0, vid except ( PlaywrightError, asyncio.TimeoutError, OSError, ValueError, RuntimeError, ) as e: logger.error("sound_stream.video_error", video_id=vid, error=str(e)[:200]) if vid: await release_video_id(campaign_id, run_id, track_name or "", vid) return 0, 1, None async def _skip_sound_positions(page: Any, n: int, seen: set[str]) -> None: """Mark n videos at the top of the sound listing as seen without opening them.""" skipped = 0 for _ in range(n + 10): # small margin for scroll passes if skipped >= n: break hrefs: list[str] = await page.evaluate("""() => Array.from(document.querySelectorAll('a[href*="/video/"]')) .map(a => a.href) .filter(h => /\\/video\\/\\d+/.test(h)) """) for h in hrefs: _, vid = TikTokExtractor.author_and_video_id(h) if vid and vid not in seen: seen.add(vid) skipped += 1 if skipped >= n: break if skipped < n: await page.evaluate("window.scrollTo(0, document.body.scrollHeight)") await asyncio.sleep(1) logger.info("sound_stream.skip_ahead_done", requested=n, actual=skipped) async def _chain_next_sound_videos( page: Any, sound_id: str, seen: set[str], stored: int, errors: int, video_ids: list[str], sample_size: int, original_sound_id: str | None, campaign_id: str, run_id: str, track_name: str | None, location: str | None, backend: Any, ) -> tuple[int, int]: """Chain forward via the in-overlay Next button until it breaks or sample_size is reached. Mutates seen and video_ids in-place. Returns updated (stored, errors).""" consecutive_repeats = 0 next_video_id: str | None = None while stored < sample_size: advanced = await _click_next_in_sound_overlay(page) if not advanced: logger.info( "sound_stream.next_btn_unavailable", sound_id=sound_id, stored=stored ) break current_url = page.url _, next_video_id = TikTokExtractor.author_and_video_id(current_url) if not next_video_id or next_video_id in seen: consecutive_repeats += 1 if consecutive_repeats >= 3: logger.info( "sound_stream.next_btn_stuck", sound_id=sound_id, stored=stored ) break continue consecutive_repeats = 0 assert next_video_id is not None if not await claim_video_id( campaign_id, run_id, track_name or "", next_video_id ): logger.info("sound_stream.skip_claimed", video_id=next_video_id) seen.add(next_video_id) continue ds, de, vid_out = await _process_sound_video( page, next_video_id, current_url, sound_id, original_sound_id, campaign_id, run_id, track_name, location, backend, ) stored += ds errors += de if vid_out: video_ids.append(vid_out) seen.add(next_video_id) return stored, errors @mcp.tool() @tool_guard(ScrapeTagCommentsResponse, echo=("sound_id",)) async def scrape_sound_comments_stream( campaign_id: str, run_id: str, sound_id: str, original_sound_id: str | None = None, session_id: str | None = None, track_name: str | None = None, sample_size: int = 20, skip_video_ids: list[str] | None = None, location: str | None = None, skip_positions: int = 0, ) -> ScrapeTagCommentsResponse: """Scrape comments from a TikTok sound page. Opens the first video from the sound page listing, then chains forward via the in-overlay "Next video" button as long as it keeps working (no listing round-trip per video). Falls back to click→Escape/go_back→re-click-from-listing only when the Next button is missing or breaks. sound_id: TikTok sound ID (e.g. "6969882811419330561") original_sound_id: pass same as sound_id to mark all videos is_original_sound=True. skip_video_ids: video IDs already collected (tag scrape pass). location: label stored on every record (e.g. "us", "gb"). skip_positions: mark this many videos at the top of the listing as seen (without opening them) before starting — pass on the first call of a fresh browser session to offset this shard's starting point, reducing overlap with sibling shards. """ proxy = None if session_id: sess = get_session(session_id) if sess: proxy = sess.proxy page = await _pool.get_page( session_id or "__no_session__", proxy, headless=settings.browser_headless ) entry_url = sound_url(sound_id) seen: set[str] = set(skip_video_ids or []) # ── 1. Load sound page ────────────────────────────────────────────────── logger.info( "sound_stream.load", sound_id=sound_id, location=location, skipping=len(seen) ) try: await page.goto(entry_url, wait_until="domcontentloaded", timeout=30000) except (PlaywrightError, asyncio.TimeoutError) as e: logger.error("sound_stream.goto_failed", sound_id=sound_id, error=str(e)[:200]) return ScrapeTagCommentsResponse( success=False, tag=sound_id, error=f"sound_page_goto_failed: {e}" ) try: await page.wait_for_selector('a[href*="/video/"]', timeout=20000) except PlaywrightError: body_snippet = (await page.inner_text("body"))[:300] logger.error("sound_stream.page_blank", sound_id=sound_id, body=body_snippet) return ScrapeTagCommentsResponse( success=False, tag=sound_id, error="sound_page_blank: no video thumbnails found — likely blocked", ) if skip_positions > 0: await _skip_sound_positions(page, skip_positions, seen) backend = get_backend() video_ids: list[str] = [] stored = 0 errors = 0 consecutive_open_failures = 0 # ── 2. Open from listing, then chain via Next button; fall back to listing on break ── while stored < sample_size: href = await _find_next_sound_href(page, seen) if not href: logger.info("sound_stream.no_more_videos", stored=stored) break _, video_id = TikTokExtractor.author_and_video_id(href) if not video_id: seen.add("") continue if not await claim_video_id(campaign_id, run_id, track_name or "", video_id): logger.info("sound_stream.skip_claimed", video_id=video_id) seen.add(video_id) continue logger.info("sound_stream.video_start", video_id=video_id, stored=stored) # Click video to open overlay try: link = page.locator(f'a[href*="{video_id}"]').first await click_and_await_comments(page, link.click) await page.wait_for_function( "() => window.location.href.includes('/video/')", timeout=15000, ) await wait_for_comments_stable(page) except (PlaywrightError, asyncio.TimeoutError) as e: logger.warning( "sound_stream.click_failed", video_id=video_id, error=str(e)[:100] ) seen.add(video_id) await release_video_id(campaign_id, run_id, track_name or "", video_id) errors += 1 consecutive_open_failures += 1 if consecutive_open_failures >= 3: break if not await _return_to_sound_page(page, entry_url): break continue consecutive_open_failures = 0 ds, de, vid_out = await _process_sound_video( page, video_id, page.url, sound_id, original_sound_id, campaign_id, run_id, track_name, location, backend, ) stored += ds errors += de if vid_out: video_ids.append(vid_out) seen.add(video_id) stored, errors = await _chain_next_sound_videos( page, sound_id, seen, stored, errors, video_ids, sample_size, original_sound_id, campaign_id, run_id, track_name, location, backend, ) if stored < sample_size and not await _return_to_sound_page(page, entry_url): logger.warning("sound_stream.return_failed_stopping", stored=stored) break logger.info( "sound_stream.done", sound_id=sound_id, stored=stored, errors=errors, location=location, ) return ScrapeTagCommentsResponse( tag=sound_id, videos_processed=stored, errors=errors, video_ids=video_ids )