from collections.abc import Iterator import redis import sqlalchemy import sqlalchemy.pool from anydi import Container, Module, Provider, provider from confluent_kafka import Producer from fansifter_common import context from fansifter_common.adapters.aws.secretsmanager import SecretsManager from fansifter_common.adapters.db.base import Database from fansifter_common.constants import PROD_ENVIRONMENT, QA_ENVIRONMENT from fansifter_common.m2m_token import M2MTokenManager from fansifter_common.utils.functional import lazy_proxy from owsclient import OwsClient from app.config import Settings, settings from app.connectors.aws.s3 import S3Client from app.connectors.db.repositories import ( AccountToFanResponseEmailAddressMappingRepository, ) from app.connectors.ows_preference_center import OwsPreferenceCenterClient from app.connectors.sendgrid import SendgridClient from app.handlers import SendgridWebhooksReplyHandler, SendgridWebhooksUnsubHandler class AppModule(Module): @provider(scope="singleton") def secrets_manager(self, settings: Settings) -> SecretsManager: return SecretsManager(region_name=settings.aws_region_name) @provider(scope="singleton") def m2m_token_manager( self, settings: Settings, secrets_manager: SecretsManager ) -> M2MTokenManager: return M2MTokenManager( secrets_manager=secrets_manager, secret_name_key=settings.m2m_token_secret_key_name, secret_expire_name_key=settings.m2m_token_secret_expiry_key_name, ) @provider(scope="singleton") def ows_client( self, settings: Settings, m2m_token_manager: M2MTokenManager ) -> OwsClient: return OwsClient( environment=settings.environment, service_name=settings.service_name, m2m_token_manager=m2m_token_manager if settings.environment in [QA_ENVIRONMENT, PROD_ENVIRONMENT] else None, correlation_id_getter=context.get_correlation_id, request_context_getter=context.get_request_context, ) @provider(scope="singleton") def ows_preference_center( self, ows_client: OwsClient, settings: Settings ) -> OwsPreferenceCenterClient: return OwsPreferenceCenterClient( ows_client=ows_client, ows_client_timeout=settings.ows_client_timeout, ) @provider(scope="singleton") def kafka_producer(self, settings: Settings) -> Producer: producer = Producer( { "bootstrap.servers": settings.kafka_bootstrap_servers, "security.protocol": "ssl" if settings.kafka_use_ssl else "plaintext", } ) return producer @provider(scope="singleton") def db(self, settings: Settings) -> Iterator[Database]: with Database( url=settings.snowflake_url, engine_args={ "echo": settings.snowflake_echo, "poolclass": sqlalchemy.pool.QueuePool, "pool_size": settings.snowflake_pool_size, "max_overflow": settings.snowflake_pool_max_overflow, "pool_recycle": settings.snowflake_pool_recycle, "pool_pre_ping": settings.snowflake_pool_pre_ping, "pool_reset_on_return": settings.snowflake_pool_reset_on_return, "connect_args": settings.snowflake_connect_args, }, ) as db: yield db @provider(scope="singleton") def account_to_fan_response_email_address_mapping_repository( self, db: Database ) -> AccountToFanResponseEmailAddressMappingRepository: return AccountToFanResponseEmailAddressMappingRepository(db) @provider(scope="singleton") def s3_client(self, settings: Settings) -> S3Client: client = S3Client(region_name=settings.aws_region_name) return client @provider(scope="singleton") def sendgrid_client(self, settings: Settings) -> SendgridClient: return SendgridClient( api_key=settings.sendgrid_api_key, timeout=settings.sendgrid_timeout, ) @provider(scope="singleton") def redis_client(self, settings: Settings) -> redis.Redis: return redis.Redis( host=settings.redis_host, port=settings.redis_port, db=settings.redis_db, ssl=settings.redis_ssl, decode_responses=True, ) @provider(scope="singleton") def unsub_handler( self, settings: Settings, kafka_producer: Producer, s3_client: S3Client, ows_preference_center: OwsPreferenceCenterClient, ) -> SendgridWebhooksUnsubHandler: return SendgridWebhooksUnsubHandler( kafka_producer=kafka_producer, kafka_sendgrid_inbound_topic=settings.kafka_sendgrid_inbound_topic, s3_client=s3_client, failed_events_s3_bucket=settings.failed_events_s3_bucket, ows_preference_center=ows_preference_center, ) @provider(scope="singleton") def reply_handler( self, settings: Settings, kafka_producer: Producer, db: Database, repository: AccountToFanResponseEmailAddressMappingRepository, sendgrid_client: SendgridClient, redis_client: redis.Redis, s3_client: S3Client, ) -> SendgridWebhooksReplyHandler: return SendgridWebhooksReplyHandler( kafka_producer=kafka_producer, kafka_sendgrid_reply_topic=settings.kafka_sendgrid_reply_topic, db=db, repository=repository, sendgrid_client=sendgrid_client, redis_client=redis_client, cache_key_template=settings.reply_redis_cache_key_template, requests_count_limit=settings.reply_requests_count_limit, requests_count_reset_seconds=settings.reply_requests_count_reset_seconds, s3_client=s3_client, failed_events_s3_bucket=settings.failed_events_s3_bucket, default_forward_email=settings.default_forward_email, from_email_template=settings.from_email_template, accounts_to_forward_on_qa_env=settings.accounts_to_forward_on_qa_env, env=settings.environment, ) def make_container() -> Container: # Configure DI container return Container( providers=[ Provider(Settings, factory=lambda: settings, scope="singleton"), ], modules=[AppModule], ) container = lazy_proxy(make_container)