import sqlalchemy as sa from sqlalchemy.orm import selectinload from dmp.adapters.db import Repository from dmp.app_connections.enums import AppConnectionStatus from dmp.meta.models import MetaAdReportingConnection class MetaAdReportingConnectionRepository(Repository[MetaAdReportingConnection]): default_options = [ selectinload(MetaAdReportingConnection.ad_accounts), ] def exists_by_identity_id(self, identity_id: str) -> bool: query = sa.select( sa.select(1) .exists() .where( MetaAdReportingConnection.identity_id == identity_id, ) ) result = self.db.session.execute(query) return bool(result.scalar_one()) def get_by_identity_id(self, identity_id: str) -> MetaAdReportingConnection | None: query = ( sa.select(MetaAdReportingConnection) .where( MetaAdReportingConnection.identity_id == identity_id, ) .options(*self.default_options) ) result = self.db.session.execute(query) return result.scalar_one_or_none() def get_by_fivetran_connection_id( self, fivetran_connection_id: str ) -> MetaAdReportingConnection | None: query = ( sa.select(MetaAdReportingConnection) .where( MetaAdReportingConnection.fivetran_connector_id == fivetran_connection_id ) .options(*self.default_options) ) result = self.db.session.execute(query) return result.scalar_one_or_none() def find_not_notified_connections( self, ) -> list[MetaAdReportingConnection]: query = ( sa.select(MetaAdReportingConnection) .where( sa.and_( MetaAdReportingConnection.status.in_( ( AppConnectionStatus.WAITING_FOR_PROCESSING, AppConnectionStatus.CONNECTED, ) ), MetaAdReportingConnection.initial_email_sent == sa.false(), ) ) .options(*self.default_options) ) result = self.db.session.execute(query) return list(result.scalars().all()) def get_fivetran_connections_ids(self) -> list[str]: query = sa.select(MetaAdReportingConnection.fivetran_connector_id).where( MetaAdReportingConnection.status.in_( AppConnectionStatus.schema_update_statuses() ) ) return list(self.db.session.execute(query).scalars().all())