diff --git a/src/ad_seller/clients/gam_soap_client.py b/src/ad_seller/clients/gam_soap_client.py index 51b940f..edc82e3 100644 --- a/src/ad_seller/clients/gam_soap_client.py +++ b/src/ad_seller/clients/gam_soap_client.py @@ -232,9 +232,7 @@ def create_order( try: order["externalOrderId"] = int(external_order_id) except (ValueError, TypeError): - order["notes"] = ( - f"{order.get('notes', '')} [deal_id:{external_order_id}]".strip() - ) + order["notes"] = f"{order.get('notes', '')} [deal_id:{external_order_id}]".strip() result = order_service.createOrders([order]) order_data = result[0] @@ -747,7 +745,7 @@ def _fmt_date(dt: Any) -> str: d = getattr(dt, "date", None) if d is None: return str(dt) - return f"{getattr(d, 'year', '')-0:04d}-{getattr(d, 'month', ''):02d}-{getattr(d, 'day', ''):02d}" + return f"{getattr(d, 'year', '') - 0:04d}-{getattr(d, 'month', ''):02d}-{getattr(d, 'day', ''):02d}" return { "id": str(getattr(o, "id", "")), @@ -778,9 +776,7 @@ def list_line_items_for_order(self, order_id: str) -> list[dict[str, Any]]: "id": str(getattr(li, "id", "")), "name": getattr(li, "name", ""), "status": str(getattr(li, "status", "")), - "impressions_goal": getattr( - getattr(li, "primaryGoal", None), "units", -1 - ), + "impressions_goal": getattr(getattr(li, "primaryGoal", None), "units", -1), "cost_type": str(getattr(li, "costType", "CPM")), } for li in (getattr(result, "results", None) or []) @@ -848,9 +844,7 @@ def run_delivery_report( downloader = self._client.GetDataDownloader(version=self.api_version) buf = io.BytesIO() - downloader.DownloadReportToFile( - job_id, "CSV_EXCEL", buf, use_gzip_compression=False - ) + downloader.DownloadReportToFile(job_id, "CSV_EXCEL", buf, use_gzip_compression=False) buf.seek(0) lines = buf.read().decode("utf-8").splitlines() diff --git a/src/ad_seller/clients/ssp_deals_api_mcp_client.py b/src/ad_seller/clients/ssp_deals_api_mcp_client.py new file mode 100644 index 0000000..adf5d50 --- /dev/null +++ b/src/ad_seller/clients/ssp_deals_api_mcp_client.py @@ -0,0 +1,218 @@ +# Author: Green Mountain Systems AI Inc. +# Donated to IAB Tech Lab + +"""SSP client for IAB deals-api-mcp (HTTP Streamable MCP transport). + +Connects to a running deals-api-mcp server via MCP Streamable HTTP and maps +the IAB Deal Sync API v1.0 tool schema to the generic SSPClient interface. + +deals-api-mcp tools used: + - deals_create: create a new deal (required: name, origin, seller, dealFloor, startDate) + - deals_status: get deal + all buyer seat statuses + history + - deals_list: list deals with optional status filter + - deals_update: update mutable deal fields (blocked after deals_send) + - deals_pause: pause an active deal (propagates to provider) + - deals_resume: resume a paused deal (propagates to provider) +""" + +import logging +from datetime import datetime, timezone +from typing import Any, Optional + +from .freewheel_mcp_client import FreeWheelMCPClient +from .ssp_base import ( + SSPClient, + SSPDeal, + SSPDealCreateRequest, + SSPDealStatus, + SSPTroubleshootResult, + SSPType, +) + +logger = logging.getLogger(__name__) + +# deals-api-mcp sellerStatus integer → SSPDealStatus +# SellerStatus enum: 0=Active, 1=Paused, 2=Pending, 4=Complete, 5=Archived +_SELLER_STATUS_MAP: dict[int, SSPDealStatus] = { + 0: SSPDealStatus.ACTIVE, + 1: SSPDealStatus.PAUSED, + 2: SSPDealStatus.CREATED, + 4: SSPDealStatus.EXPIRED, + 5: SSPDealStatus.ARCHIVED, +} + +_SELLER_STATUS_LABEL: dict[int, str] = { + 0: "Active", + 1: "Paused", + 2: "Pending", + 4: "Complete", + 5: "Archived", +} + + +class DealsAPIMCPClient(SSPClient): + """SSP connector for deals-api-mcp via MCP Streamable HTTP. + + Wraps FreeWheelMCPClient for transport and maps structured IAB tool + arguments to/from the generic SSPClient interface. + """ + + ssp_type: SSPType = SSPType.CUSTOM + ssp_name: str = "IAB Deals MCP" + + def __init__( + self, + *, + mcp_url: str, + api_key: Optional[str] = None, + seller_origin: str = "publisher.example.com", + ) -> None: + self.ssp_type = SSPType.CUSTOM + self.ssp_name = "IAB Deals MCP" + self._mcp_url = mcp_url + self._api_key = api_key + self._seller_origin = seller_origin + self._mcp_client = FreeWheelMCPClient() + + # ── Lifecycle ────────────────────────────────────────────────────────── + + async def connect(self) -> None: + auth_params = {"api_key": self._api_key} if self._api_key else None + await self._mcp_client.connect( + url=self._mcp_url, + auth_params=auth_params, + ) + logger.info("Connected to deals-api-mcp at %s", self._mcp_url) + + async def disconnect(self) -> None: + await self._mcp_client.disconnect() + + # ── Deal Operations ──────────────────────────────────────────────────── + + async def create_deal(self, request: SSPDealCreateRequest) -> SSPDeal: + """Map SSPDealCreateRequest → deals_create structured args.""" + now_iso = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") + + args: dict[str, Any] = { + "name": getattr(request, "name", None) or "Untitled Deal", + "origin": self._seller_origin, + "seller": getattr(request, "advertiser", None) or self.ssp_name, + "dealFloor": getattr(request, "cpm", None) or 1.0, + "startDate": getattr(request, "start_date", None) or now_iso, + } + + if getattr(request, "end_date", None): + args["endDate"] = request.end_date + if getattr(request, "impressions_goal", None): + args["units"] = request.impressions_goal + if getattr(request, "buyer_seat_ids", None): + args["wseat"] = request.buyer_seat_ids + if getattr(request, "description", None): + args["description"] = request.description + if getattr(request, "currency", None): + args["currency"] = request.currency + + raw = await self._mcp_client.call_tool("deals_create", args) + return self._parse_deal(raw) + + async def get_deal(self, deal_id: str) -> SSPDeal: + raw = await self._mcp_client.call_tool("deals_status", {"dealId": deal_id}) + return self._parse_deal(raw) + + async def list_deals( + self, + *, + status: Optional[SSPDealStatus] = None, + limit: int = 100, + ) -> list[SSPDeal]: + args: dict[str, Any] = {"pageSize": min(limit, 100)} + raw = await self._mcp_client.call_tool("deals_list", args) + if isinstance(raw, dict): + items = raw.get("deals", raw.get("items", [])) + return [self._parse_deal({"deal": d}) for d in items] + return [] + + async def clone_deal( + self, + source_deal_id: str, + overrides: Optional[dict[str, Any]] = None, + ) -> SSPDeal: + """Clone by fetching the source deal and creating a new one with overrides.""" + source_raw = await self._mcp_client.call_tool("deals_status", {"dealId": source_deal_id}) + + # Build create args from source deal's terms, apply overrides + source_deal = source_raw.get("deal", {}) if isinstance(source_raw, dict) else {} + terms = source_deal.get("terms", {}) if isinstance(source_deal.get("terms"), dict) else {} + now_iso = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") + + args: dict[str, Any] = { + "name": f"Copy of {source_deal.get('name', source_deal_id)}", + "origin": self._seller_origin, + "seller": source_deal.get("seller", self.ssp_name), + "dealFloor": terms.get("dealFloor", 1.0), + "startDate": terms.get("startDate", now_iso), + } + if terms.get("endDate"): + args["endDate"] = terms["endDate"] + if overrides: + args.update(overrides) + + raw = await self._mcp_client.call_tool("deals_create", args) + return self._parse_deal(raw) + + async def update_deal(self, deal_id: str, updates: dict[str, Any]) -> SSPDeal: + raw = await self._mcp_client.call_tool("deals_update", {"id": deal_id, **updates}) + return self._parse_deal(raw) + + async def troubleshoot_deal(self, deal_id: str) -> SSPTroubleshootResult: + raw = await self._mcp_client.call_tool("deals_status", {"dealId": deal_id}) + + issues: list[str] = [] + if isinstance(raw, dict): + seats = raw.get("buyerSeats", []) + for seat in seats: + if isinstance(seat, dict) and seat.get("buyerStatusLabel") == "Rejected": + issues.append(f"Buyer seat {seat.get('seatId', '?')} was rejected by provider") + + return SSPTroubleshootResult( + deal_id=deal_id, + status=self._seller_status_label(raw), + primary_issues=issues, + ssp_type=self.ssp_type, + raw=raw, + ) + + # ── Response Parsing ─────────────────────────────────────────────────── + + def _parse_deal(self, raw: Any) -> SSPDeal: + """Parse a deals-api-mcp tool response into SSPDeal.""" + if not isinstance(raw, dict): + return SSPDeal(deal_id="unknown", ssp_type=self.ssp_type, ssp_name=self.ssp_name) + + # deals_create wraps in {"success": true, "deal": {...}} + # deals_status wraps in {"deal": {...}, "buyerSeats": [...]} + deal = raw.get("deal", raw) + if not isinstance(deal, dict): + deal = raw + + terms = deal.get("terms", {}) if isinstance(deal.get("terms"), dict) else {} + seller_status_int = deal.get("sellerStatus") + + return SSPDeal( + deal_id=str(deal.get("externalDealId", deal.get("id", "unknown"))), + name=deal.get("name"), + status=_SELLER_STATUS_MAP.get(seller_status_int, SSPDealStatus.CREATED), + cpm=terms.get("dealFloor"), + currency=terms.get("currency", "USD"), + ssp_type=self.ssp_type, + ssp_name=self.ssp_name, + raw=raw, + ) + + def _seller_status_label(self, raw: Any) -> str: + if isinstance(raw, dict): + deal = raw.get("deal", raw) + if isinstance(deal, dict): + code = deal.get("sellerStatus") + return _SELLER_STATUS_LABEL.get(code, "unknown") + return "unknown" diff --git a/src/ad_seller/clients/ssp_factory.py b/src/ad_seller/clients/ssp_factory.py index fe4591a..c89587e 100644 --- a/src/ad_seller/clients/ssp_factory.py +++ b/src/ad_seller/clients/ssp_factory.py @@ -113,6 +113,19 @@ def _create_ssp_client(name: str, settings: Any) -> Any: api_key=settings.index_exchange_api_key, ) + elif name_lower == "deals_api_mcp": + from .ssp_deals_api_mcp_client import DealsAPIMCPClient + + if not settings.deals_api_mcp_url: + logger.warning("deals_api_mcp configured but DEALS_API_MCP_URL not set") + return None + + return DealsAPIMCPClient( + mcp_url=settings.deals_api_mcp_url, + api_key=settings.deals_api_mcp_key, + seller_origin=settings.deals_api_mcp_seller_origin, + ) + else: logger.warning( "Unknown SSP '%s'. To add support, create a client in ssp_factory.py " diff --git a/src/ad_seller/config/settings.py b/src/ad_seller/config/settings.py index e23e8c8..62e3783 100644 --- a/src/ad_seller/config/settings.py +++ b/src/ad_seller/config/settings.py @@ -127,7 +127,7 @@ class Settings(BaseSettings): # SSP Connectors (publishers can configure multiple SSPs) # Comma-separated list of SSP names to enable - ssp_connectors: str = "" # e.g. "pubmatic,magnite" + ssp_connectors: str = "" # e.g. "pubmatic,magnite,deals_api_mcp" # Routing rules: inventory_type:ssp_name pairs, comma-separated ssp_routing_rules: str = "" # e.g. "ctv:pubmatic,display:magnite" # PubMatic SSP @@ -139,6 +139,10 @@ class Settings(BaseSettings): # Index Exchange SSP (REST API) index_exchange_api_url: Optional[str] = None index_exchange_api_key: Optional[str] = None + # IAB Deals MCP (deals-api-mcp server — HTTP Streamable transport) + deals_api_mcp_url: Optional[str] = None # e.g. http://localhost:3100/mcp + deals_api_mcp_key: Optional[str] = None # IAB_DEALS_API_KEY on the MCP server + deals_api_mcp_seller_origin: str = "publisher.example.com" # origin field for deals_create # Pricing Configuration default_currency: str = "USD" diff --git a/src/ad_seller/interfaces/api/main.py b/src/ad_seller/interfaces/api/main.py index 8d7e109..ba83f41 100644 --- a/src/ad_seller/interfaces/api/main.py +++ b/src/ad_seller/interfaces/api/main.py @@ -3930,14 +3930,12 @@ async def get_deal_performance(deal_id: str): else 0.0 ) revenue = summary.get("revenue_usd", 0.0) - avg_cpm = ( - round(revenue / impressions_served * 1000, 2) - if impressions_served - else 0.0 - ) + avg_cpm = round(revenue / impressions_served * 1000, 2) if impressions_served else 0.0 pacing = ( - "not_started" if impressions_served == 0 - else "on_track" if fill_rate >= 40 + "not_started" + if impressions_served == 0 + else "on_track" + if fill_rate >= 40 else "behind" ) diff --git a/src/ad_seller/storage/quote_history.py b/src/ad_seller/storage/quote_history.py index 7f4a666..fcbbbd5 100644 --- a/src/ad_seller/storage/quote_history.py +++ b/src/ad_seller/storage/quote_history.py @@ -80,9 +80,7 @@ async def record_quote( } # Store by quote_id - await self._storage.set( - f"quote_history:{quote_id}", record - ) + await self._storage.set(f"quote_history:{quote_id}", record) # Maintain a buyer+product index for fast lookup index_key = f"quote_history_index:{buyer_id}:{product_id}" @@ -185,7 +183,7 @@ async def verify_pricing( reason=( f"Proposed CPM ${proposed_cpm:.2f} matches quote " f"{quote['quote_id']} (${quoted_cpm:.2f}, " - f"{relative_diff*100:.1f}% difference)." + f"{relative_diff * 100:.1f}% difference)." ), )