"""Campaign management endpoints.""" from __future__ import annotations from typing import Any import structlog from fastapi import APIRouter, BackgroundTasks, HTTPException from marketing_intelligence.agent_workflows.discovery import run_discovery from marketing_intelligence.agent_workflows.discovery_relevant import ( run_discovery_relevant, ) from marketing_intelligence.api.schemas.campaigns import CampaignStartedResponse from marketing_intelligence.core.models import SentimentRequest from marketing_intelligence.storage.manifest import count_runs, get_manifest logger = structlog.get_logger("api") router = APIRouter(prefix="/campaigns", tags=["campaigns"]) async def _next_run_id(campaign_id: str) -> str: existing = await count_runs(campaign_id) return f"{campaign_id}_{existing + 1}" @router.get( "/{campaign_id}/{run_id}/status", responses={404: {"description": "No run found for the given campaign / run ID"}}, ) async def campaign_run_status(campaign_id: str, run_id: str) -> dict[str, Any]: manifest = await get_manifest(campaign_id, run_id) if manifest is None: raise HTTPException( status_code=404, detail=f"no run found for {campaign_id}/{run_id}" ) return manifest @router.post("") async def create_campaign( brief: SentimentRequest, background_tasks: BackgroundTasks, ) -> CampaignStartedResponse: campaign = brief.to_campaign() run_id = await _next_run_id(campaign.campaign_id or campaign.campaign_key) background_tasks.add_task(run_discovery, campaign, run_id) return CampaignStartedResponse( status="started", run_id=run_id, artist=campaign.artist_name or "", tracks=len(campaign.tags), ) @router.post("/relevant") async def create_relevant_campaign( request: SentimentRequest, background_tasks: BackgroundTasks, ) -> CampaignStartedResponse: run_id = await _next_run_id(request.artist_key) background_tasks.add_task(run_discovery_relevant, request, run_id) return CampaignStartedResponse( status="started", run_id=run_id, artist=request.artist_name, tracks=len(request.tracks), )