import logging from collections.abc import Iterable from typing import Any from anydi import singleton from fansifter_common.utils.dictutil import compare_seq_of_dicts from dmp.adapters.fivetran import FivetranClient from dmp.adapters.fivetran.custom_reports import ( FACEBOOK_CUSTOM_REPORTS, GOOGLE_CUSTOM_REPORTS, TIKTOK_CUSTOM_REPORTS, ) from dmp.adapters.fivetran.exceptions import FivetranClientTimeoutError from dmp.adapters.fivetran.models import ( AdsConnection, SchemaUpdate, StandardConfig, StandardConfigUpdate, TableUpdate, ) from dmp.app_connections.constants import APP_CONNECTION_ENABLED_TABLES from dmp.app_connections.enums import AppConnectionPlatform, AppConnectionStatus from dmp.app_connections.exceptions import AppConnectionNotFoundError from dmp.app_connections.types import AppConnectionRepository logger = logging.getLogger(__name__) @singleton class AppConnectionService: def __init__( self, fivetran_client: FivetranClient, ) -> None: self.fivetran_client = fivetran_client @staticmethod def _build_schema_config( existing_schema_config: StandardConfig, enabled_tables: Iterable[str] ) -> StandardConfigUpdate: new_schema_config = StandardConfigUpdate(schemas={}) for schema_name, schema in existing_schema_config.schemas.items(): new_schema = SchemaUpdate(tables={}) new_tables_config = {} for table_name, table in schema.tables.items(): if ( table.enabled_patch_settings and not table.enabled_patch_settings.allowed ): continue if table_name.upper() in enabled_tables: new_tables_config[table_name] = TableUpdate(enabled=True) elif table.enabled_patch_settings.allowed: new_tables_config[table_name] = TableUpdate(enabled=False) else: new_tables_config[table_name] = TableUpdate(enabled=table.enabled) new_schema.tables = new_tables_config new_schema_config.schemas[schema_name] = new_schema return new_schema_config @staticmethod def _is_config_schema_changed( old_config: StandardConfig, new_config: StandardConfigUpdate ) -> bool: for old_schema_name, old_schema in old_config.schemas.items(): new_schema = new_config.schemas.get(old_schema_name) if new_schema is None: return True for old_table_name, old_table in old_schema.tables.items(): if not old_table.enabled_patch_settings.allowed: continue new_table = new_schema.tables.get(old_table_name) if new_table is None or old_table.enabled != new_table.enabled: return True return False def sync_custom_reports( # noqa: C901 self, connection_id: str, platform: AppConnectionPlatform ) -> bool: def _sync_meta_reports(connection: AdsConnection) -> bool: if compare_seq_of_dicts( connection.config.custom_tables or [], FACEBOOK_CUSTOM_REPORTS, ): return False self.fivetran_client.patch_ads_connection( connection_id, payload={ "config": { "custom_tables": FACEBOOK_CUSTOM_REPORTS, }, }, ) return True def _sync_tiktok_reports(connection: AdsConnection) -> bool: if compare_seq_of_dicts( connection.config.custom_reports or [], TIKTOK_CUSTOM_REPORTS, ): return False self.fivetran_client.patch_ads_connection( connection_id, payload={ "config": { "custom_reports": TIKTOK_CUSTOM_REPORTS, }, }, ) return True def _sync_google_reports(connection: AdsConnection) -> bool: if compare_seq_of_dicts( connection.config.reports or [], GOOGLE_CUSTOM_REPORTS, ): return False self.fivetran_client.patch_ads_connection( connection_id, payload={ "config": { "reports": GOOGLE_CUSTOM_REPORTS, }, }, ) return True if platform == AppConnectionPlatform.SHOPIFY: return False connection = self.fivetran_client.get_ads_connection(connection_id) try: if platform == AppConnectionPlatform.TIKTOK: return _sync_tiktok_reports(connection) elif platform == AppConnectionPlatform.META: return _sync_meta_reports(connection) elif platform == AppConnectionPlatform.GOOGLE: return _sync_google_reports(connection) return False except FivetranClientTimeoutError: logger.warning( "Fivetran patch ads connection timeout.", extra={ "connection_id": connection_id, "platform": platform, }, ) return False def sync_fivetran_tables_state( self, platform: AppConnectionPlatform, connection_repository: AppConnectionRepository, fivetran_connection_id: str | None = None, ) -> Iterable[str]: logger.info( "Syncing Fivetran tables state for %s", platform, extra={ "platform": platform, "fivetran_connection_id": fivetran_connection_id, }, ) updated_connections = [] platform_enabled_tables = APP_CONNECTION_ENABLED_TABLES.get(platform) if not platform_enabled_tables: logger.info( "No enabled tables configuration found for platform %s; skipping sync", platform, extra={"platform": platform}, ) return [] if fivetran_connection_id: logger.info( "Syncing for Fivetran connection id %s", fivetran_connection_id, extra={ "platform": platform, "fivetran_connection_id": fivetran_connection_id, }, ) connection = connection_repository.get_by_fivetran_connection_id( fivetran_connection_id=fivetran_connection_id ) if not connection: logger.error( "Connection not found for connection id %s", fivetran_connection_id, extra={ "platform": platform, "fivetran_connection_id": fivetran_connection_id, }, ) raise AppConnectionNotFoundError connections: Iterable[Any] = (connection,) else: logger.info("Syncing all connections for platform %s", platform) connections = connection_repository.all() for connection in connections: if connection.status != AppConnectionStatus.CONNECTED: logger.debug( "Skipping connection %s due to status %s", connection.fivetran_connector_id, connection.status, extra={ "platform": platform, "fivetran_connection_id": connection.fivetran_connector_id, }, ) continue if self.sync_custom_reports(connection.fivetran_connector_id, platform): logger.info( "Custom reports synced for connection %s", connection.fivetran_connector_id, extra={ "platform": platform, "fivetran_connection_id": connection.fivetran_connector_id, }, ) updated_connections.append(connection.fivetran_connector_id) connection_id = connection.fivetran_connector_id connection_schema_config = ( self.fivetran_client.get_connection_schema_config(connection_id) ) if connection_schema_config is not None: schema_config = self._build_schema_config( connection_schema_config, enabled_tables=platform_enabled_tables, ) if self._is_config_schema_changed( connection_schema_config, schema_config ): self.fivetran_client.patch_connection_schema_config( connection_id, schema_config ) updated_connections.append(connection_id) logger.info( "Updated tables state for connection %s", connection_id, extra={ "platform": platform, "fivetran_connection_id": connection_id, }, ) else: logger.debug( "No tables state changes required for connection %s", connection_id, extra={ "platform": platform, "fivetran_connection_id": connection_id, }, ) else: logger.warning( "Connection schema config missing for connection %s", connection_id, extra={ "platform": platform, "fivetran_connection_id": connection_id, }, ) return set(updated_connections)