import logging 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.google.models import ( GoogleAdReportingConnection, ) from dmp.google.repositories import ( GoogleAdReportingConnectionRepository, GoogleUserConnectionRepository, ) from dmp.google.services import GoogleAdAccountService logger = logging.getLogger(__name__) @singleton class GoogleAdReportingConnectionService: def __init__( self, fivetran_client: FivetranClient, ad_reporting_connection_repository: GoogleAdReportingConnectionRepository, user_connection_repository: GoogleUserConnectionRepository, ad_account_service: GoogleAdAccountService, 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 ) -> GoogleAdReportingConnection | None: connection = self.ad_reporting_connection_repository.get_connection_by_identity_id_and_user_id( identity_id, user_id ) return connection def is_connection_processed(self, connection: GoogleAdReportingConnection) -> bool: return self.ad_account_dbt_repository.exists_by_schema_and_platform( schema=connection.fivetran_schema, platform=AdAccountPlatform.GOOGLE, ) def refresh_connection( # noqa: C901 self, connection: GoogleAdReportingConnection ) -> GoogleAdReportingConnection: fivetran_connection = self.fivetran_client.get_ads_connection( connection.fivetran_connector_id ) if fivetran_connection.config.sync_mode is None: logger.warning( "Google 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 ): user_connection = ( self.user_connection_repository.get_by_identity_id_and_user_id( connection.identity_id, connection.user_id ) ) 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.GOOGLE, ) ] if ( not connection_ad_accounts and connection.status in AppConnectionStatus.pre_processed_statuses() ): if user_connection: connection_ad_accounts = [ connection_ad_account.ad_account.external_id for connection_ad_account in user_connection.connection_ad_accounts ] if connection_ad_accounts: if ( fivetran_connection.config.sync_mode in ( AdsAccountsSyncMode.ALL_ACCOUNTS, AdsAccountsSyncMode.MANAGER_ACCOUNTS, ) and user_connection ): ad_accounts = self.ad_account_service.get_ad_accounts_for_ad_reporting_connection( connection, connection_ad_accounts, { connection_ad_account.ad_account.external_id: connection_ad_account.ad_account for connection_ad_account in user_connection.connection_ad_accounts }, ) else: ad_accounts = self.ad_account_service.get_ad_accounts_for_ad_reporting_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[GoogleAdReportingConnection]: 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.GOOGLE, ) ) 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