"""Discovery (relevant) — three-stage shape, Stage 1 is adaptive: the agent picks tag vs sound_id itself and rotates proxy/browser/profile until it hits target. Stage 1: adaptive relevant-video scrape (tag or sound_id, agent's choice) → stored to disk Stage 2: per-video sentiment analysis → updates stored records Stage 3: campaign-level synthesis → writes fact_campaign_sentiment """ import asyncio from typing import Any import structlog from marketing_intelligence.agent.orchestrator import run as run_agent from marketing_intelligence.core.config import settings from marketing_intelligence.agent_workflows.stages import ( BROWSER_TOOLS, campaign_sentiment_stage, per_video_sentiment_stage, scrape_relevant_stage, ) from marketing_intelligence.core.models import Campaign, SentimentRequest from marketing_intelligence.mcp.tools.scraping_tools.base import claim_video_id from marketing_intelligence.storage.factory import get_backend from marketing_intelligence.storage.manifest import ( finalize, init_manifest, load_request, mark_interrupted_shards, now_iso, save_request, update_shard, update_stage, ) log = structlog.get_logger("discovery_relevant") STAGE1_MAX_CONCURRENT_UNITS = 5 STAGE1_SHARD_THRESHOLD = 100 STAGE1_SHARD_SIZE = 100 STAGE1_SHARD_TIMEOUT_SECONDS = ( 1800 # hard ceiling per shard — attempt budgets are prompt-based ) # (soft), this is the real code-enforced backstop def _shard_sizes(sample_size: int) -> list[int]: """Split a track's sample_size into ~STAGE1_SHARD_SIZE-video chunks for parallel scraping. Below the threshold, sharding overhead (extra browser sessions) isn't worth it — a track with a modest target (e.g. 80, already fine at 1 concurrent slot among several tracks) stays a single unit. Only large single/few-track targets (e.g. 500 on one track) get split. """ if sample_size <= STAGE1_SHARD_THRESHOLD: return [sample_size] shard_count = -(-sample_size // STAGE1_SHARD_SIZE) # ceil division base, remainder = divmod(sample_size, shard_count) return [base + (1 if i < remainder else 0) for i in range(shard_count)] def _build_units( tracks_and_targets: list[tuple[str, str, int]], label_suffix: str = "", ) -> list[tuple[str, str, int, str, int]]: """Shard each (tag, track_name, target) into (tag, track_name, shard_size, label, skip_positions) units. Shared by a fresh run (target = full sample_size) and resume (target = remaining).""" units: list[tuple[str, str, int, str, int]] = [] for tag, track_name, target in tracks_and_targets: if target <= 0: continue sizes = _shard_sizes(target) for i, size in enumerate(sizes, start=1): label = ( f"{track_name}{label_suffix}" if len(sizes) == 1 else f"{track_name}{label_suffix} (shard {i}/{len(sizes)})" ) skip_positions = (i - 1) * STAGE1_SHARD_SIZE units.append((tag, track_name, size, label, skip_positions)) return units async def _run_stage1_units( units: list[tuple[str, str, int, str, int]], campaign: Campaign, campaign_id: str, run_id: str, base_params: dict[str, Any], request: SentimentRequest, ) -> None: log.info( "discovery_relevant.stage1.start", campaign_id=campaign_id, total_units=len(units), max_concurrent=STAGE1_MAX_CONCURRENT_UNITS, ) semaphore = asyncio.Semaphore(STAGE1_MAX_CONCURRENT_UNITS) results = await asyncio.gather( *[ _run_stage1_track( tag, track_name, size, label, skip_positions, campaign, campaign_id, run_id, base_params, request, semaphore, ) for tag, track_name, size, label, skip_positions in units ], return_exceptions=True, ) failed = [ (label, exc) for (_, _, _, label, _), exc in zip(units, results, strict=True) if isinstance(exc, BaseException) ] if failed: log.error( "discovery_relevant.stage1.units_failed", campaign_id=campaign_id, failed_units=[u for u, _ in failed], errors=[str(e)[:200] for _, e in failed], ) log.info( "discovery_relevant.stage1.done", campaign_id=campaign_id, succeeded=len(units) - len(failed), failed=len(failed), ) async def _finish_stage2_stage3( campaign: Campaign, campaign_id: str, run_id: str, request: SentimentRequest, base_params: dict[str, Any], ) -> str | None: """Stage 2 (per-video sentiment) → Stage 3 (campaign synthesis) → PDF → finalize.""" backend = get_backend() collected_videos = len(await backend.read_all_video_sentiments(campaign_id, run_id)) await update_stage( campaign_id, run_id, "stage1", status="done", collected_videos=collected_videos, finished_at=now_iso(), ) # Haiku is the permanent choice here (and for Stage 1) — confirmed 2026-07-07 after # comparing joy_crookes_5/6/7. Stage 3 (below) stays on Sonnet. watch_model_id = ( settings.anthropic_watch_model_id if settings.agent_backend == "anthropic" else settings.bedrock_watch_model_id ) stage2 = per_video_sentiment_stage(campaign_id=campaign_id, run_id=run_id) log.info("discovery_relevant.stage2.start", campaign_id=campaign_id) await update_stage( campaign_id, run_id, "stage2", status="running", started_at=now_iso() ) await run_agent( stage2.flow, base_params, run_id=run_id, model_id=watch_model_id, allowed_tools=stage2.tools, ) await update_stage( campaign_id, run_id, "stage2", status="done", finished_at=now_iso() ) log.info("discovery_relevant.stage2.done", campaign_id=campaign_id) stage3 = campaign_sentiment_stage( campaign_key=campaign.campaign_key, artist_key=campaign.artist_key, campaign_id=campaign_id, run_id=run_id, prompt=request.prompt, sample_size=campaign.sample_size, ) log.info("discovery_relevant.stage3.start", campaign_id=campaign_id) await update_stage( campaign_id, run_id, "stage3", status="running", started_at=now_iso() ) result = await run_agent( stage3.flow, base_params, run_id=run_id, allowed_tools=stage3.tools ) await update_stage( campaign_id, run_id, "stage3", status="done", finished_at=now_iso() ) log.info("discovery_relevant.stage3.done", campaign_id=campaign_id) try: from marketing_intelligence.reporting.pdf_generator import save_pdf_temp campaign_config_key = f"{campaign.campaign_key}_{campaign.artist_key}" sentiment_data = await backend.read_sentiment(campaign_config_key, run_id) if sentiment_data: local_pdf_path = save_pdf_temp(sentiment_data) pdf_ref = await backend.write_pdf_report( campaign_config_key, run_id, local_pdf_path ) log.info("discovery_relevant.pdf_generated", run_id=run_id, path=pdf_ref) else: log.warning( "discovery_relevant.pdf_skipped", reason="sentiment record not found", run_id=run_id, ) except Exception as e: log.warning("discovery_relevant.pdf_failed", error=str(e)[:200]) await finalize(campaign_id, run_id, status="completed") return result async def _run_stage1_track( tag: str, track_name: str, shard_sample_size: int, shard_label: str, skip_positions: int, campaign: Campaign, campaign_id: str, run_id: str, base_params: dict[str, Any], request: SentimentRequest, semaphore: asyncio.Semaphore, ) -> None: async with semaphore: stage1 = scrape_relevant_stage( tracks=[(tag, track_name)], campaign_id=campaign_id, run_id=run_id, sample_size=shard_sample_size, prompt=request.prompt, tiktok_handle=campaign.tiktok_handle, skip_positions=skip_positions, ) track_params = {**base_params, "track_name": track_name} log.info( "discovery_relevant.stage1.unit.start", campaign_id=campaign_id, track=track_name, unit=shard_label, ) await update_shard( campaign_id, run_id, shard_label, status="running", started_at=now_iso() ) watch_model_id = ( settings.anthropic_watch_model_id if settings.agent_backend == "anthropic" else settings.bedrock_watch_model_id ) try: await asyncio.wait_for( run_agent( stage1.flow, track_params, run_id=run_id, model_id=watch_model_id, allowed_tools=stage1.tools | BROWSER_TOOLS, ), timeout=STAGE1_SHARD_TIMEOUT_SECONDS, ) await update_shard( campaign_id, run_id, shard_label, status="done", finished_at=now_iso() ) except TimeoutError: # Valid outcome — whatever this shard collected stands. Orphaned Camoufox processes # are a known, accepted rough edge here, not a correctness bug. log.warning( "discovery_relevant.stage1.unit.timeout", campaign_id=campaign_id, track=track_name, unit=shard_label, timeout_seconds=STAGE1_SHARD_TIMEOUT_SECONDS, ) await update_shard( campaign_id, run_id, shard_label, status="timeout", finished_at=now_iso(), ) log.info( "discovery_relevant.stage1.unit.done", campaign_id=campaign_id, track=track_name, unit=shard_label, ) async def run_discovery_relevant( request: SentimentRequest, run_id: str, stealth_hint: str = "conservative", ) -> str | None: campaign = request.to_campaign() campaign_id = campaign.campaign_id or campaign.campaign_key tracks = list(zip(campaign.tags, campaign.track_names, strict=False)) base_params = { "campaign_id": campaign_id, "campaign_key": campaign.campaign_key, "artist_key": campaign.artist_key, "run_type": "discovery_relevant", "run_id": run_id, "stealth_hint": stealth_hint, } log.info( "discovery_relevant.start", campaign_id=campaign_id, artist=campaign.artist_name, tracks=[t for _, t in tracks], run_id=run_id, ) try: units = _build_units( [(tag, track_name, campaign.sample_size) for tag, track_name in tracks] ) await init_manifest( campaign_id, run_id, campaign.sample_size, [label for _, _, _, label, _ in units], ) await save_request(campaign_id, run_id, request.model_dump()) await _run_stage1_units( units, campaign, campaign_id, run_id, base_params, request ) result = await _finish_stage2_stage3( campaign, campaign_id, run_id, request, base_params ) log.info( "discovery_relevant.done", campaign_id=campaign_id, run_id=run_id, result=(result or "")[:200], ) return result except Exception as e: await finalize(campaign_id, run_id, status="failed") log.error( "discovery_relevant.failed", campaign_id=campaign_id, run_id=run_id, error=str(e)[:300], ) raise async def resume_discovery_relevant(campaign_id: str, run_id: str) -> str | None: """Resume a run whose process died mid-flight (manifest.json still says "running").""" request_dict = await load_request(campaign_id, run_id) if request_dict is None: log.error( "discovery_relevant.resume.no_request", campaign_id=campaign_id, run_id=run_id, ) await finalize(campaign_id, run_id, status="failed") return None request = SentimentRequest(**request_dict) campaign = request.to_campaign() tracks = list(zip(campaign.tags, campaign.track_names, strict=False)) base_params = { "campaign_id": campaign_id, "campaign_key": campaign.campaign_key, "artist_key": campaign.artist_key, "run_type": "discovery_relevant", "run_id": run_id, "stealth_hint": "conservative", } log.info("discovery_relevant.resume.start", campaign_id=campaign_id, run_id=run_id) try: existing_records = await get_backend().read_all_video_sentiments( campaign_id, run_id ) remaining_targets: list[tuple[str, str, int]] = [] for tag, track_name in tracks: await mark_interrupted_shards(campaign_id, run_id, track_name) existing_ids = { r["video_id"] for r in existing_records if r.get("track_name") == track_name and r.get("video_id") } remaining = campaign.sample_size - len(existing_ids) log.info( "discovery_relevant.resume.track", campaign_id=campaign_id, track=track_name, already_collected=len(existing_ids), remaining=max(remaining, 0), ) for video_id in existing_ids: await claim_video_id(campaign_id, run_id, track_name, video_id) remaining_targets.append((tag, track_name, remaining)) units = _build_units(remaining_targets, label_suffix=" (resume)") if units: await _run_stage1_units( units, campaign, campaign_id, run_id, base_params, request ) else: log.info( "discovery_relevant.resume.stage1_already_done", campaign_id=campaign_id, run_id=run_id, ) result = await _finish_stage2_stage3( campaign, campaign_id, run_id, request, base_params ) log.info( "discovery_relevant.resume.done", campaign_id=campaign_id, run_id=run_id, result=(result or "")[:200], ) return result except Exception as e: await finalize(campaign_id, run_id, status="failed") log.error( "discovery_relevant.resume.failed", campaign_id=campaign_id, run_id=run_id, error=str(e)[:300], ) raise def run_discovery_relevant_sync( request: SentimentRequest, run_id: str, stealth_hint: str = "conservative", ) -> str | None: """Synchronous wrapper for CLI usage.""" return asyncio.run(run_discovery_relevant(request, run_id, stealth_hint))