import logging from collections.abc import Sequence from anydi import singleton from dmp.ad_accounts.enums import AdAccountPlatform from dmp.ad_accounts.repositories import AdAccountDbtRepository from dmp.adapters.fivetran import FivetranClient from dmp.adapters.fivetran.enums import AdsAccountsSyncMode, ConnectionSetupState from dmp.app_connections.enums import AppConnectionStatus from dmp.tiktok.models import ( TikTokAdAccount, TikTokAdReportingConnection, ) from dmp.tiktok.repositories import ( TikTokAdReportingConnectionRepository, TikTokUserConnectionRepository, ) from dmp.tiktok.services import TikTokAdAccountService logger = logging.getLogger(__name__) @singleton class TikTokAdReportingConnectionService: def __init__( self, fivetran_client: FivetranClient, ad_reporting_connection_repository: TikTokAdReportingConnectionRepository, user_connection_repository: TikTokUserConnectionRepository, ad_account_service: TikTokAdAccountService, ad_account_dbt_repository: AdAccountDbtRepository, ) -> None: self.fivetran_client = fivetran_client self.ad_reporting_connection_repository = ad_reporting_connection_repository self.user_connection_repository = user_connection_repository self.ad_account_service = ad_account_service self.ad_account_dbt_repository = ad_account_dbt_repository def get_connection_by_identity_id_and_user_id( self, identity_id: str, user_id: str ) -> TikTokAdReportingConnection | None: connection = self.ad_reporting_connection_repository.get_connection_by_identity_id_and_user_id( identity_id, user_id ) return connection def refresh_connection( self, connection: TikTokAdReportingConnection ) -> TikTokAdReportingConnection: fivetran_connection = self.fivetran_client.get_ads_connection( connection.fivetran_connector_id ) if fivetran_connection.config.sync_mode is None: logger.warning( "TikTok connection %s is missing sync_mode in config", fivetran_connection.id, ) status = AppConnectionStatus.from_connection_status(fivetran_connection.status) if status == AppConnectionStatus.CONNECTED: if connection.synced_at is None: self.fivetran_client.sync_connection_data( connection.fivetran_connector_id ) status = AppConnectionStatus.PENDING elif not self.is_connection_processed(connection): status = AppConnectionStatus.WAITING_FOR_PROCESSING if ( fivetran_connection.status.setup_state == ConnectionSetupState.CONNECTED or connection.ad_accounts ): if ( fivetran_connection.config.sync_mode == AdsAccountsSyncMode.SPECIFIC_ACCOUNTS and fivetran_connection.config.accounts ): connection_ad_accounts = fivetran_connection.config.accounts else: connection_ad_accounts = [ ad_account.id for ad_account in self.ad_account_dbt_repository.find_by_schemas_and_platform( schemas=[fivetran_connection.name], platform=AdAccountPlatform.TIKTOK, ) ] if ( not connection_ad_accounts and connection.status in AppConnectionStatus.pre_processed_statuses() ): user_connection = ( self.user_connection_repository.get_by_identity_id_and_user_id( connection.identity_id, connection.user_id ) ) if user_connection: connection_ad_accounts = [ ad_account.external_id for ad_account in user_connection.ad_accounts ] if connection_ad_accounts: ad_accounts = self.ad_account_service.get_ad_accounts_for_connection( connection, connection_ad_accounts ) connection.ad_accounts = ad_accounts connection.status = status connection.synced_at = fivetran_connection.succeeded_at self.ad_reporting_connection_repository.save(connection) return connection def find_to_notify_on_initial_sync_completion( self, ) -> list[TikTokAdReportingConnection]: not_notified_ad_reporting_connections = ( self.ad_reporting_connection_repository.find_not_notified_connections() ) processed_accounts = ( self.ad_account_dbt_repository.find_by_schemas_and_platform( schemas=[ connection.fivetran_schema for connection in not_notified_ad_reporting_connections ], platform=AdAccountPlatform.TIKTOK, ) ) processed_schemas = [account.source_schema for account in processed_accounts] ad_reporting_connections_to_notify = [ conn for conn in not_notified_ad_reporting_connections if conn.fivetran_schema in processed_schemas ] return ad_reporting_connections_to_notify def is_connection_processed(self, connection: TikTokAdReportingConnection) -> bool: return self.ad_account_dbt_repository.exists_by_schema_and_platform( schema=connection.fivetran_schema, platform=AdAccountPlatform.TIKTOK, ) def delete_ad_account( self, connection: TikTokAdReportingConnection, /, *, ad_account: TikTokAdAccount ) -> None: connection.remove_ad_account(ad_account) if not connection.ad_accounts: self.ad_reporting_connection_repository.delete(connection) def delete_ad_account_by_identity_id( self, identity_id: str, ad_account: TikTokAdAccount ) -> Sequence[TikTokAdReportingConnection] | None: ad_reporting_connections = ( self.ad_reporting_connection_repository.get_by_identity_id( identity_id=identity_id ) ) if ad_reporting_connections: for ad_reporting_connection in ad_reporting_connections: self.delete_ad_account(ad_reporting_connection, ad_account=ad_account) return ad_reporting_connections return None def delete_user_ad_accounts( self, identity_id: str, ad_account: TikTokAdAccount ) -> Sequence[TikTokAdReportingConnection] | None: ad_reporting_connections = ( self.ad_reporting_connection_repository.get_by_identity_id(identity_id) ) if ad_reporting_connections: for ad_reporting_connection in ad_reporting_connections: self.delete_ad_account(ad_reporting_connection, ad_account=ad_account) return ad_reporting_connections return None