"""Discovery run — assembles stages into an agent run.""" from typing import Any import structlog from marketing_intelligence.agent.orchestrator import run as run_agent from marketing_intelligence.agent_workflows.stages import ( BROWSER_TOOLS, RUN_RECORD_TOOLS, Stage, posts_stage, questions_stage, reasoning_stage, sentiment_stage, ) from marketing_intelligence.core.models import Campaign from marketing_intelligence.storage.campaign_config import deactivate_posts log = structlog.get_logger("discovery") def _build_params( campaign: Campaign, run_id: str, stealth_hint: str = "conservative" ) -> dict[str, Any]: campaign_id = campaign.campaign_id or campaign.campaign_key sample_size = campaign.sample_size or 10 tracks = ( list(zip(campaign.tags, campaign.track_names, strict=False)) if campaign.track_names else [(t, t) for t in campaign.tags] ) return { "campaign_id": campaign_id, "campaign_key": campaign.campaign_key, "artist_key": campaign.artist_key, "artist_name": campaign.artist_name or "", "tiktok_handle": campaign.tiktok_handle or "", "tracks": tracks, "questions": campaign.questions, "sample_size": sample_size, "run_type": "discovery", "run_id": run_id, "stealth_hint": stealth_hint, } def default_stages(params: dict[str, Any]) -> list[Stage]: """Standard discovery: posts (with inline comment sentiment) → reasoning → sentiment → questions.""" stages = [ posts_stage( tracks=params["tracks"], campaign_key=params["campaign_key"], artist_key=params["artist_key"], campaign_id=params["campaign_id"], run_id=params["run_id"], sample_size=params["sample_size"], ), reasoning_stage(params["campaign_key"], params["artist_key"], params["run_id"]), sentiment_stage( params["campaign_key"], params["artist_key"], params["run_id"], prompt=params["questions"][0] if params["questions"] else "", ), ] if params["questions"]: stages.append( questions_stage( params["questions"], params["campaign_key"], params["run_id"] ) ) return stages def collect_tools(stages: list[Stage]) -> set[str]: tools = set(BROWSER_TOOLS | RUN_RECORD_TOOLS) for stage in stages: tools |= stage.tools return tools async def run_discovery( campaign: Campaign, run_id: str, stealth_hint: str = "conservative", stages: list[Stage] | None = None, ) -> str | None: """Run discovery for a campaign. Stages default to full discovery pipeline.""" params = _build_params(campaign, run_id, stealth_hint) campaign_id = params["campaign_id"] if stages is None: stages = default_stages(params) allowed_tools = collect_tools(stages) track_tags = [tag for tag, _ in params["tracks"]] deactivated = deactivate_posts(campaign_id) log.info( "discovery.start", campaign_id=campaign_id, tracks=track_tags, run_id=run_id, stages=[s.name for s in stages], deactivated=deactivated, ) track_names = [name for _, name in params["tracks"]] task = ( f"Discovery run for '{campaign.campaign_nm}' " f"— {len(track_tags)} track(s): {', '.join(track_names)}. " f"Sample size: {params['sample_size']} posts per track." ) try: result = await run_agent( task, params, run_id=run_id, allowed_tools=allowed_tools ) log.info( "discovery.done", campaign_id=campaign_id, run_id=run_id, result=(result or "")[:200], ) return result except Exception as e: log.error( "discovery.failed", campaign_id=campaign_id, run_id=run_id, error=str(e)[:300], ) raise def run_discovery_sync( campaign: Campaign, run_id: str, stealth_hint: str = "conservative", stages: list[Stage] | None = None, ) -> str | None: """Synchronous wrapper for CLI usage.""" import asyncio return asyncio.run(run_discovery(campaign, run_id, stealth_hint, stages))