from __future__ import annotations import asyncio import base64 import json as jsonlib import logging import pathlib from collections.abc import AsyncIterator from types import TracebackType from typing import Any, Literal, TypeVar from urllib.parse import urljoin import httpx import pydantic from httpx._types import QueryParamTypes # noqa from campaigns.connectors.facebook import constants from campaigns.connectors.facebook.arguments import ( AdCreativeArgs, CreateAdCreativeArgs, CreateAdsetArgs, CreateEngagementCustomAudienceArgs, CreateImageAdCreativeArgs, CreatePostAdCreativeArgs, CreateVideoAdCreativeArgs, GetAdPreviewArgs, GetAudienceSummaryArgs, UpdateAdArgs, UpdateAdsetArgs, UpdateCampaignArgs, ) from campaigns.connectors.facebook.constants import ( AGENCY_TASKS, ENGAGEMENT_AUDIENCE_PREFILL_VALUE, ) from campaigns.connectors.facebook.enums import ( AdCampaignObjective, CampaignBidStrategy, ObjectStatus, ) from campaigns.connectors.facebook.exceptions import FacebookClientError from campaigns.connectors.facebook.models import ( AccessToken, Ad, AdCampaign, AdCreative, AdImage, AdPreviewData, AdPreviewResponse, AdSet, AdVideo, AudienceSummary, AudienceSummaryResponse, CreatedAdImage, CreatedAdImageResponse, CreatedAdVideo, DebugToken, DeleteResponse, EngagementAudience, ErrorResponse, InstagramBusinessAccount, InstagramBusinessAccountResponse, InstagramBusinessAccountWithId, InstagramMedia, InstagramUser, InstagramUserForPageResponse, Page, PagePost, Response, TargetingGeolocation, TargetingInterest, User, ) from campaigns.core.types import PlainToken from campaigns.utils.image import get_image_content T = TypeVar("T") logger = logging.getLogger(__name__) class FacebookClient: base_url = "https://graph.facebook.com" graph_api_version = "v15.0" sensitive_param_keys = { "client_secret", "fb_exchange_token", "input_token", "access_token", } DEFAULT_REQUEST_TIMEOUT = 10 CREATE_ADVIDEO_TIMEOUT = 60 def __init__( self, client_id: str, client_secret: str, app_system_user_token: PlainToken, app_system_user_id: str, admin_system_user_token: PlainToken, ad_account_id: str, business_manager_id: str, default_timeout: int | None = None, ) -> None: self.client_id = client_id self.client_secret = client_secret self.app_system_user_token = app_system_user_token self.app_system_user_id = app_system_user_id self.admin_system_user_token = admin_system_user_token self.ad_account_id = f"act_{ad_account_id}" self.business_manager_id = business_manager_id self.default_timeout = default_timeout or self.DEFAULT_REQUEST_TIMEOUT self._client: httpx.AsyncClient | None = None @property def client(self) -> httpx.AsyncClient: if self._client is None: raise AttributeError("Facebook client is not started") return self._client async def start(self) -> None: if self._client is None: self._client = httpx.AsyncClient(timeout=self.default_timeout) async def close(self) -> None: await self.client.aclose() async def __aenter__(self) -> FacebookClient: await self.start() return self async def __aexit__( self, exc_type: type[BaseException] | None, exc_value: BaseException | None, traceback: TracebackType | None, ) -> None: await self.close() def build_url(self, path: str) -> str: if not path.startswith("/"): path = f"/{path}" return urljoin(self.base_url, f"{self.graph_api_version}{path}") def sanitize_url(self, url: httpx.URL) -> httpx.URL: for key, _ in url.params.items(): if key in self.sensitive_param_keys: url = url.copy_remove_param(key) return url async def make_request( self, method: str, path: str, *, params: QueryParamTypes | None = None, json: Any = None, timeout: int | None = None, ) -> httpx.Response: try: response = await self.client.request( method, self.build_url(path), params=params, json=json, timeout=timeout, ) except httpx.RequestError as exc: raise FacebookClientError( f"Failed to request `facebook-graphapi`: {exc.request.url.path}", request=exc.request, ) from exc # Sanitize request url for exception stacktrace context response.request.url = self.sanitize_url(response.url) try: response.raise_for_status() except httpx.HTTPStatusError as exc: context: dict[str, Any] = { "status_code": exc.response.status_code, } error_message = ( "Invalid `facebook-graphapi` " f"response status: {exc.response.status_code}" ) try: error_response_obj = pydantic.parse_obj_as( ErrorResponse, response.json(), ) context["fb_error"] = error_response_obj.error.dict() fb_error_message = ( error_response_obj.error.error_user_msg or error_response_obj.error.message ) error_message = f"`facebook-graphapi` client error: {fb_error_message}" except (pydantic.ValidationError, jsonlib.JSONDecodeError): pass raise FacebookClientError( error_message, context=context, request=exc.request, response=exc.response, ) from exc return response async def make_request_typed( self, method: str, path: str, type_: type[T], *, params: QueryParamTypes | None = None, json: Any = None, timeout: int | None = None, ) -> T: response = await self.make_request( method, path, params=params, json=json, timeout=timeout ) try: return pydantic.parse_obj_as(type_, response.json()) except pydantic.ValidationError as exc: raise FacebookClientError( "Failed to deserialize `facebook-graphapi` " f"response into {type_}: {exc.errors()}" ) from exc async def get_oauth_access_token(self, exchange_token: PlainToken) -> AccessToken: return await self.make_request_typed( "POST", "/oauth/access_token", type_=AccessToken, json={ "grant_type": "fb_exchange_token", "client_id": self.client_id, "client_secret": self.client_secret, "fb_exchange_token": exchange_token, }, ) async def debug_token(self, access_token: PlainToken) -> DebugToken: response_obj = await self.make_request_typed( "GET", "/debug_token", type_=Response[DebugToken], params={ "input_token": access_token, "access_token": f"{self.client_id}|{self.client_secret}", }, ) return response_obj.data async def debug_tokens( self, access_tokens: list[PlainToken] ) -> dict[PlainToken, DebugToken]: if not access_tokens: return {} tasks = [] for access_token in access_tokens: tasks.append(self.debug_token(access_token)) return dict(zip(access_tokens, await asyncio.gather(*tasks), strict=True)) async def get_user(self, user_id: str, user_access_token: PlainToken) -> User: return await self.make_request_typed( "GET", f"/{user_id}", type_=User, params={ "access_token": user_access_token, "fields": ",".join(User.__fields__), }, ) async def get_user_pages( self, user_id: str, user_access_token: PlainToken, after: str | None = None ) -> AsyncIterator[Page]: params = { "access_token": user_access_token, "fields": Page.RESPONSE_FIELDS, } if after: params["after"] = after response_obj = await self.make_request_typed( "GET", f"/{user_id}/accounts", type_=Response[list[Page]], params=params, ) for account in response_obj.data: yield account if response_obj.paging and response_obj.paging.next: async for user_account in self.get_user_pages( user_id=user_id, user_access_token=user_access_token, after=response_obj.paging.cursors.after, ): yield user_account async def delete_user_permissions( self, *, user_id: str, user_access_token: PlainToken ) -> None: await self.make_request( "DELETE", f"/{user_id}/permissions", params={"access_token": user_access_token}, ) async def get_targeting_geolocations( self, query: str, limit: int | None = None ) -> list[TargetingGeolocation]: params: dict[str, Any] = { "access_token": self.app_system_user_token, "q": query, "type": constants.GEOLOCATION, "location_types": jsonlib.dumps(constants.TARGETING_LOCATION_TYPES), } if limit is not None: params["limit"] = limit response_obj = await self.make_request_typed( "GET", "/search", Response[list[TargetingGeolocation]], params=params ) return response_obj.data async def get_targeting_interests( self, query: str, limit: int | None = None ) -> list[TargetingInterest]: params: dict[str, Any] = { "access_token": self.app_system_user_token, "q": query, "type": constants.INTEREST, } if limit is not None: params["limit"] = limit response_obj = await self.make_request_typed( "GET", "/search", Response[list[TargetingInterest]], params=params ) return response_obj.data async def create_campaign( self, *, name: str, objective: AdCampaignObjective, lifetime_budget: int ) -> AdCampaign: data = { "status": ObjectStatus.PAUSED, "access_token": self.app_system_user_token, "special_ad_categories": "NONE", "fields": AdCampaign.RESPONSE_FIELDS, "bid_strategy": CampaignBidStrategy.LOWEST_COST_WITHOUT_CAP, "name": name, "objective": objective, "lifetime_budget": lifetime_budget, } return await self.make_request_typed( "POST", f"/{self.ad_account_id}/campaigns", AdCampaign, json=data ) async def update_campaign( self, campaign_id: str, *, args: UpdateCampaignArgs ) -> AdCampaign: data = { "fields": AdCampaign.RESPONSE_FIELDS, "access_token": self.app_system_user_token, **args.dict(exclude_unset=True, exclude_none=True), } return await self.make_request_typed( "POST", f"/{campaign_id}", AdCampaign, json=data ) async def delete_campaign(self, campaign_id: str) -> None: params = { "access_token": self.app_system_user_token, "fields": AdCampaign.RESPONSE_FIELDS, } await self.make_request("DELETE", f"/{campaign_id}", params=params) async def create_ad_set(self, campaign_id: str, *, args: CreateAdsetArgs) -> AdSet: data = { "access_token": self.app_system_user_token, "fields": AdSet.RESPONSE_FIELDS, "status": ObjectStatus.PAUSED, "campaign_id": campaign_id, "name": args.name, "optimization_goal": args.optimization_goal, "billing_event": args.billing_event, "start_time": args.start_time.isoformat(), "end_time": args.end_time.isoformat(), } if args.targeting: data["targeting"] = args.targeting.json( exclude_unset=True, exclude_none=True ) return await self.make_request_typed( "POST", f"/{self.ad_account_id}/adsets", AdSet, json=data ) async def update_adset(self, adset_id: str, *, args: UpdateAdsetArgs) -> AdSet: data = { "access_token": self.app_system_user_token, "fields": AdSet.RESPONSE_FIELDS, **args.dict( exclude_unset=True, exclude_none=True, exclude={"start_time", "end_time"}, ), } if args.start_time: data["start_time"] = args.start_time.isoformat() if args.end_time: data["end_time"] = args.end_time.isoformat() return await self.make_request_typed("POST", f"/{adset_id}", AdSet, json=data) async def delete_adset(self, adset_id: str) -> None: params = {"access_token": self.app_system_user_token} await self.make_request("DELETE", f"/{adset_id}", params=params) async def create_ad_creative(self, args: CreateAdCreativeArgs) -> AdCreative: data = { "access_token": self.app_system_user_token, "fields": AdCreative.RESPONSE_FIELDS, **args.dict(exclude_unset=True), } return await self.make_request_typed( "POST", f"/{self.ad_account_id}/adcreatives", AdCreative, json=data ) async def create_video_ad_creative( self, args: CreateVideoAdCreativeArgs ) -> AdCreative: return await self.create_ad_creative( args=AdCreativeArgs( object_story_spec=args.object_story_spec, ) ) async def create_image_ad_creative( self, args: CreateImageAdCreativeArgs ) -> AdCreative: return await self.create_ad_creative( args=AdCreativeArgs( object_story_spec=args.object_story_spec, ) ) async def create_post_adcreative( self, args: CreatePostAdCreativeArgs ) -> AdCreative: # When creating ad creatives, if the object_story_id being used # is already in use by an existing creative, then the API will return # the value of the existing creative_id instead of creating a new one. return await self.create_ad_creative( args=AdCreativeArgs( object_story_id=f"{args.page_id}_{args.post_id}", ) ) async def delete_adcreative(self, adcreative_id: str) -> None: params = {"access_token": self.app_system_user_token} await self.make_request("DELETE", f"/{adcreative_id}", params=params) async def create_ad(self, name: str, ad_set_id: str, ad_creative_id: str) -> Ad: data = { "access_token": self.app_system_user_token, "fields": Ad.RESPONSE_FIELDS, "status": ObjectStatus.PAUSED, "name": name, "adset_id": ad_set_id, "creative": { "creative_id": ad_creative_id, }, } return await self.make_request_typed( "POST", f"/{self.ad_account_id}/ads", Ad, json=data ) async def delete_ad(self, ad_id: str) -> None: params = {"access_token": self.app_system_user_token} await self.make_request("DELETE", f"/{ad_id}", params=params) async def update_ad(self, ad_id: str, *, args: UpdateAdArgs) -> Ad: data = { "access_token": self.app_system_user_token, "fields": Ad.RESPONSE_FIELDS, **args.dict(exclude_unset=True, exclude_none=True), } return await self.make_request_typed("POST", f"/{ad_id}", Ad, json=data) async def get_ad(self, ad_id: str) -> Ad: params = { "access_token": self.app_system_user_token, "fields": Ad.RESPONSE_FIELDS, } return await self.make_request_typed("GET", f"/{ad_id}", Ad, params=params) async def assign_agency_to_business_manager( self, *, page_id: str, page_access_token: PlainToken ) -> None: data = { "business": self.business_manager_id, "permitted_tasks": AGENCY_TASKS, "access_token": page_access_token, } await self.make_request("POST", f"/{page_id}/agencies", json=data) async def assign_agency_to_app_system_user(self, page_id: str) -> None: data = { "access_token": self.admin_system_user_token, "user": self.app_system_user_id, "tasks": AGENCY_TASKS, } await self.make_request("POST", f"/{page_id}/assigned_users", json=data) async def get_published_posts_for_page( self, *, page_id: str, page_access_token: PlainToken, limit: int | None = None ) -> list[PagePost]: params: dict[str, Any] = { "access_token": page_access_token, "fields": ",".join(PagePost.__fields__), } if limit is not None: params["limit"] = limit response_obj = await self.make_request_typed( "GET", f"/{page_id}/published_posts", Response[list[PagePost]], params=params, ) return response_obj.data async def get_post( self, *, post_id: str, page_access_token: PlainToken ) -> PagePost: params: dict[str, Any] = { "access_token": page_access_token, "fields": ",".join(PagePost.__fields__), } return await self.make_request_typed( "GET", f"/{post_id}/", PagePost, params=params ) async def create_ad_image( self, target: pathlib.Path | bytes | str ) -> CreatedAdImage: image_bytes = await get_image_content(target) data = { "access_token": self.app_system_user_token, "bytes": base64.b64encode(image_bytes).decode(), } response_obj = await self.make_request_typed( "POST", f"/{self.ad_account_id}/adimages", type_=CreatedAdImageResponse, json=data, ) return response_obj.images.bytes async def get_ad_image(self, image_hash: str) -> AdImage | None: params = { "access_token": self.app_system_user_token, "hashes": f'["{image_hash}"]', "fields": AdImage.RESPONSE_FIELDS, } response_obj = await self.make_request_typed( "GET", f"/{self.ad_account_id}/adimages", type_=Response[list[AdImage]], params=params, ) if response_obj.data: return response_obj.data[0] return None async def delete_ad_image(self, image_hash: str) -> bool: params = { "access_token": self.app_system_user_token, "hash": image_hash, } response_obj = await self.make_request_typed( "DELETE", f"/{self.ad_account_id}/adimages", type_=DeleteResponse, params=params, ) return response_obj.success async def create_ad_video(self, url: str) -> CreatedAdVideo: data = { "access_token": self.app_system_user_token, "file_url": url, } return await self.make_request_typed( "POST", f"/{self.ad_account_id}/advideos", type_=CreatedAdVideo, json=data, timeout=self.CREATE_ADVIDEO_TIMEOUT, ) async def get_ad_video(self, video_id: str) -> AdVideo: params = { "access_token": self.app_system_user_token, "fields": AdVideo.RESPONSE_FIELDS, } return await self.make_request_typed( "GET", f"/{video_id}", type_=AdVideo, params=params ) async def delete_ad_video(self, video_id: str) -> bool: params = { "access_token": self.app_system_user_token, "video_id": video_id, } response_obj = await self.make_request_typed( "DELETE", f"/{self.ad_account_id}/advideos", type_=DeleteResponse, params=params, ) return response_obj.success async def get_instagram_business_account_for_page( self, page_id: str, page_access_token: PlainToken ) -> InstagramBusinessAccountWithId | None: params = { "access_token": page_access_token, "fields": "instagram_business_account", } response_obj = await self.make_request_typed( "GET", f"/{page_id}", params=params, type_=InstagramBusinessAccountResponse ) return response_obj.instagram_business_account async def get_instagram_user_for_page( self, page_id: str, page_access_token: PlainToken ) -> InstagramUser | None: params = { "access_token": page_access_token, "fields": "id", } response_obj = await self.make_request_typed( "GET", f"/{page_id}/instagram_accounts", params=params, type_=InstagramUserForPageResponse, ) if data := response_obj.data: return data[0] return None async def get_instagram_business_account( self, instagram_business_account_id: str, page_access_token: PlainToken ) -> InstagramBusinessAccount: params = { "access_token": page_access_token, "fields": ",".join(InstagramBusinessAccount.__fields__), } return await self.make_request_typed( "GET", f"/{instagram_business_account_id}", params=params, type_=InstagramBusinessAccount, ) async def get_ad_preview(self, args: GetAdPreviewArgs) -> AdPreviewData: params = { "access_token": self.app_system_user_token, "ad_format": args.ad_format, "creative": args.creative.json(exclude_unset=True, exclude_none=True), } response_obj = await self.make_request_typed( "GET", f"/{self.ad_account_id}/generatepreviews", type_=AdPreviewResponse, params=params, ) return response_obj.data[0] # Fb will return iframe object even if parameters were incorrect async def get_audience_summary( self, args: GetAudienceSummaryArgs ) -> AudienceSummary: params = { "access_token": self.app_system_user_token, "targeting_spec": args.json(exclude_unset=True, exclude_none=True), } response_obj = await self.make_request_typed( "GET", f"/{self.ad_account_id}/reachestimate", type_=AudienceSummaryResponse, params=params, ) return response_obj.data async def get_instagram_media( self, instagram_account_id: str, page_access_token: PlainToken, after: str | None = None, ) -> AsyncIterator[InstagramMedia]: async for media in self._get_instagram_media( "media", instagram_account_id=instagram_account_id, page_access_token=page_access_token, after=after, ): yield media async def get_instagram_stories( self, instagram_account_id: str, page_access_token: PlainToken, after: str | None = None, ) -> AsyncIterator[InstagramMedia]: async for media in self._get_instagram_media( "stories", instagram_account_id=instagram_account_id, page_access_token=page_access_token, after=after, ): yield media async def get_instagram_media_item( self, post_id: str, page_access_token: PlainToken ) -> InstagramMedia: params = { "access_token": page_access_token, "fields": ",".join(InstagramMedia.__fields__), } return await self.make_request_typed( "GET", f"/{post_id}", type_=InstagramMedia, params=params ) async def _get_instagram_media( self, path: Literal["media", "stories"], *, instagram_account_id: str, page_access_token: PlainToken, after: str | None = None, ) -> AsyncIterator[InstagramMedia]: params = { "access_token": page_access_token, "fields": ",".join(InstagramMedia.__fields__), } if after: params["after"] = after response_obj = await self.make_request_typed( "GET", f"/{instagram_account_id}/{path}", params=params, type_=Response[list[InstagramMedia]], ) for post in response_obj.data: yield post if response_obj.paging and response_obj.paging.next: async for post in self._get_instagram_media( path, instagram_account_id=instagram_account_id, page_access_token=page_access_token, after=response_obj.paging.cursors.after, ): yield post async def create_engagement_audience( self, args: CreateEngagementCustomAudienceArgs ) -> EngagementAudience: data = { "access_token": self.app_system_user_token, "fields": EngagementAudience.RESPONSE_FIELDS, "name": args.name, "prefil": ENGAGEMENT_AUDIENCE_PREFILL_VALUE, "rule": args.rule.json(), } return await self.make_request_typed( "POST", f"/{self.ad_account_id}/customaudiences", EngagementAudience, json=data, ) async def get_engagement_audience( self, audience_external_id: str ) -> EngagementAudience: params = { "access_token": self.app_system_user_token, "fields": EngagementAudience.RESPONSE_FIELDS, } return await self.make_request_typed( "GET", f"/{audience_external_id}", params=params, type_=EngagementAudience, )