"""Monitoring run via LLM agent — re-scrape known posts for all active campaigns.""" 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.storage.campaign_config import ( list_active_campaign_ids, list_posts, ) from marketing_intelligence.storage.checkpoint import ( clear_checkpoint, load_checkpoint, save_checkpoint, ) from marketing_intelligence.storage.factory import get_backend log = structlog.get_logger("monitor_agent") ALLOWED_TOOLS: set[str] = { "create_browser_session", "close_browser_session", "scrape_posts_batch", "scrape_url", "scrape_artist_profile", "apply_behavior_tactics", "persist_post_metrics", "write_run_record", "write_artist_snapshot", } async def run_monitor_agent(run_id: str) -> None: campaign_ids = list_active_campaign_ids() if not campaign_ids: log.info("monitor_agent.empty") return watch_model_id = ( settings.anthropic_watch_model_id if settings.agent_backend == "anthropic" else settings.bedrock_watch_model_id ) log.info("monitor_agent.start", campaigns=len(campaign_ids), model=watch_model_id) for campaign_id in campaign_ids: campaign_run_id = f"{run_id}_{campaign_id}" ck_key = f"{campaign_id}_agent" posts = list_posts(campaign_id) if not posts: log.info("monitor_agent.no_posts", campaign_id=campaign_id) continue checkpoint = load_checkpoint(ck_key) completed_posts: set[str] = set(checkpoint.get("completed", [])) pending = [p for p in posts if p.post_id not in completed_posts] if not pending: log.info( "monitor_agent.skipped", campaign_id=campaign_id, reason="all posts done", ) continue log.info( "monitor_agent.campaign_start", campaign_id=campaign_id, total=len(posts), pending=len(pending), ) cckey = pending[0].campaign_config_key prev_metrics = await get_backend().read_latest_metrics(cckey) def _post_line(p: Any, _pm: dict[str, Any] = prev_metrics) -> str: line = ( f" - post_id={p.post_id} url={p.url} " f"sound_id={p.sound_id or ''} " f"campaign_config_key={p.campaign_config_key}" ) prev = _pm.get(p.post_id) if prev: line += ( f"\n prev: views={prev.get('views')} likes={prev.get('likes')} " f"comments={prev.get('comments')} shares={prev.get('shares')} " f"favorites={prev.get('favorites')}" ) return line post_lines = "\n".join(_post_line(p) for p in pending) parts = cckey.split("_", 1) campaign_key = parts[0] if parts else campaign_id artist_key = parts[1] if len(parts) > 1 else "" task = ( f"Monitoring run for campaign '{campaign_id}' " f"— re-scrape {len(pending)} known posts.\n\n" f"Posts to scrape:\n{post_lines}" ) params = { "campaign_id": campaign_id, "campaign_key": campaign_key, "artist_key": artist_key, "artist_name": "", "tiktok_handle": "", "run_id": campaign_run_id, "run_type": "watch", } try: result = await run_agent( task, params, run_id=campaign_run_id, model_id=watch_model_id, allowed_tools=ALLOWED_TOOLS, ) clear_checkpoint(ck_key) log.info( "monitor_agent.done", campaign_id=campaign_id, result=(result or "")[:200], ) except Exception as e: completed_posts.update(p.post_id for p in pending) save_checkpoint(ck_key, {"completed": list(completed_posts), "failed": []}) log.error( "monitor_agent.failed", campaign_id=campaign_id, error=str(e)[:200] ) log.info("monitor_agent.all_done", campaigns=len(campaign_ids))