from datetime import datetime, timezone from typing import List from connectors.schemas import ( CreateConnectorResponse, ConnectorResponse, WebhookEvent, FivetranStatuses, FivetranEventTypes ) from external_api.base.clients.fivetran_client import FivetranClient from external_api.base.clients.client_factory import ApiClientFactory from connectors.repositories.connectors_repository import ConnectorsRepository from config import FIVETRAN_GROUP_ID from labels.schemas import TeamMemberUserResponse from models.connector import ConnectorStatus, ConnectorPlatforms from services.labels_repository import LabelsRepository from services.users_repository import UsersRepository from utils import session_utils from config import DECIBEL_API_DOMAIN REDIRECT_URI = 'https://decibel.stream/{}' class ConnectorsService: connectors_repository: ConnectorsRepository users_repository: UsersRepository fivetran_client: FivetranClient labels_repository: LabelsRepository def __init__( self, connectors_repository: ConnectorsRepository = ConnectorsRepository(), fivetran_client: FivetranClient = ApiClientFactory.fivetran_client(), users_repository: UsersRepository = UsersRepository(), labels_repository: LabelsRepository = LabelsRepository() ): self.connectors_repository = connectors_repository self.fivetran_client = fivetran_client self.users_repository = users_repository self.labels_repository = labels_repository def __is_connector_decibel(self, connector_data): name_parts = connector_data.schema.split("_") try: if len(name_parts) < 3: return False if connector_data.service not in ConnectorPlatforms.all_values(): return False if not self.users_repository.is_user_exists(int(name_parts[2])): return False if not self.users_repository.is_label_exists(int(name_parts[1])): return False if not name_parts[0] == "decibel": return False return True except ValueError: return False def update_ads_accounts(self, connector_data): pass async def webhook_data_process(self, params: WebhookEvent): connector = self.connectors_repository.get_connector_by_external_id(params.connector_id) fivetran_conn = await self.fivetran_client.get_connector_by_id(params.connector_id) if connector is None and self.__is_connector_decibel(fivetran_conn): fivetran_conn_name_parts = fivetran_conn.schema.split("_") connector = self.connectors_repository.create_connector( external_id=fivetran_conn.schema, platform=fivetran_conn.service, label_id=int(fivetran_conn_name_parts[1]), user_id=int(fivetran_conn_name_parts[2]), name=fivetran_conn.schema ) if connector is not None and params.event == FivetranEventTypes.SYNC_END.value: if params.data.status == FivetranStatuses.SUCCESSFUL.value: sync_time = datetime.strptime(params.created, "%Y-%m-%dT%H:%M:%S.%fZ").replace(tzinfo=timezone.utc) if connector is not None and connector.last_synced_at < sync_time: connector.last_synced_at = sync_time connector.status = ConnectorStatus.ACTIVE.value elif params.data.status in (FivetranStatuses.FAILED.value, FivetranStatuses.FAILURE_WITH_TASK.value): connector.status = ConnectorStatus.RECONNECT.value connector.updated_at = datetime.now(timezone.utc) self.update_ads_accounts(fivetran_conn) session_utils.session_commit() async def get_fivetran_connectors(self): return await self.fivetran_client.get_connectors(FIVETRAN_GROUP_ID) async def get_connect_card(self, user_id: int, label_id: int, platform: str): conn = self.connectors_repository.get_incomplete_connector_for_user_in_label(user_id, label_id, platform) user = self.users_repository.get_user_by_id(user_id) if conn is not None: connect_card_uri = await self.fivetran_client.generate_connect_card( conn.external_id, redirect_uri=REDIRECT_URI.format(label_id) ) return CreateConnectorResponse( id=conn.id, name=conn.name, external_id=conn.external_id, status=conn.status, platform=conn.platform, owner=TeamMemberUserResponse(user_id, user.name, user.email, None), connect_card_uri=connect_card_uri ) conn_name = self.build_connector_name(user_id, label_id, platform) fivetran_conn = await self.fivetran_client.create_connector( FIVETRAN_GROUP_ID, platform, name=conn_name, redirect_uri=REDIRECT_URI.format(label_id) ) conn = self.connectors_repository.create_connector( external_id=fivetran_conn.id, platform=platform, label_id=label_id, name=conn_name, user_id=user_id ) session_utils.session_commit() return CreateConnectorResponse( id=conn.id, name=conn.name, external_id=conn.external_id, platform=platform, status=conn.status, owner=TeamMemberUserResponse(user_id, user.name, user.email, None), connect_card_uri=fivetran_conn.connect_card.uri ) def get_label_connectors(self, label_id: int): connectors = self.connectors_repository.get_label_connectors_without_incomplete(label_id=label_id) return [ ConnectorResponse( id=connector.id, name=connector.name, platform=connector.platform, external_id=connector.external_id, status=connector.status, owner=TeamMemberUserResponse(connector.owner.id, connector.owner.name, connector.owner.email, None) ) for connector in connectors] async def register_webhook(self): await self.fivetran_client.create_webhook( FIVETRAN_GROUP_ID, events=["sync_end", "sync_start"], url=DECIBEL_API_DOMAIN + "/connectors-webhook" ) @staticmethod def build_connector_name(user_id: int, label_id: int, platform: str): name_parts = [ "decibel", label_id, user_id, platform, str(datetime.now().timestamp()).replace(".", "") ] return "_".join(str(p) for p in name_parts)