import logging from typing import Any from anydi import singleton from dmp.adapters.db.repositories import ModelT, Repository from dmp.adapters.fivetran import FivetranClient from dmp.adapters.fivetran.exceptions import FivetranClientTimeoutError from dmp.adapters.ows_notifications import ( AdReportingSyncCompletedNotificationArgs, OwsNotificationsClient, ) from dmp.app_connections.enums import AdsConnectionPlatform from dmp.app_connections.types import AppConnection, AppConnectionAdReportingService logger = logging.getLogger(__name__) @singleton class AdsConnectionService: def __init__( self, fivetran_client: FivetranClient, ows_notifications_client: OwsNotificationsClient, ) -> None: self.fivetran_client = fivetran_client self.ows_notifications_client = ows_notifications_client def remove_fivetran_ad_account( self, fivetran_connection_id: str, external_id: str ) -> None: connection = self.fivetran_client.get_ads_connection(fivetran_connection_id) existing_accounts = connection.config.accounts if existing_accounts and external_id in existing_accounts: existing_accounts.remove(external_id) try: self.fivetran_client.patch_ads_connection( fivetran_connection_id, payload={ "config": { "accounts": existing_accounts, }, }, ) except FivetranClientTimeoutError: logger.warning( "Fivetran patch facebook ads connection timeout.", extra={"connection_id": fivetran_connection_id}, ) def delete_fivetran_ad_account( self, connection: AppConnection, ad_account_external_id: str ) -> None: if not connection.ad_accounts: self.fivetran_client.delete_connection(connection.fivetran_connector_id) else: self.remove_fivetran_ad_account( connection.fivetran_connector_id, ad_account_external_id ) def notify_user_initial_sync_and_refresh_connection( self, platform: AdsConnectionPlatform, ad_reporting_connection_repository: Repository[ModelT], ad_reporting_connection_service: AppConnectionAdReportingService, ) -> list[Any]: ad_reporting_connection_to_update = [] ad_reporting_connections_to_notify = ( ad_reporting_connection_service.find_to_notify_on_initial_sync_completion() ) for connection in ad_reporting_connections_to_notify: is_notified = ( self.ows_notifications_client.notify_ad_reporting_sync_completed( AdReportingSyncCompletedNotificationArgs( identity_id=connection.identity_id, ad_accounts=[ account.name or account.id for account in connection.ad_accounts ], ad_reporting_platform=platform, ) ) ) if is_notified: connection.initial_email_sent = True connection = ad_reporting_connection_service.refresh_connection(connection) ad_reporting_connection_to_update.append(connection) ad_reporting_connection_repository.add_all(ad_reporting_connection_to_update) return ad_reporting_connection_to_update