from asyncio import Semaphore from collections.abc import AsyncIterator from aiodynamo.client import Client from aiodynamo.credentials import Credentials from aiodynamo.http.httpx import HTTPX from anydi import Container, Module, Provider, provider from anydi.testing import TestContainer from fansifter_common.utils.functional import lazy_proxy from httpx import AsyncClient from jwtauth import JWTAuth from jwtauth.utils import get_default_audience, get_default_issuer, get_default_jwks_url from url_shortener.adapters.aws.s3 import S3Client from url_shortener.adapters.kafka.kafka_client import KafkaClient from url_shortener.config import Settings, settings from url_shortener.repositories import PathsRepository from url_shortener.services import ( DeletePathsService, ShortenUrlService, ShortUrlRedirectService, UpdatePathService, ) class AppModule(Module): @provider(scope="singleton") def jwt_auth(self, settings: Settings) -> JWTAuth: return JWTAuth( jwks_url=get_default_jwks_url(settings.environment), audience=get_default_audience(settings.environment), issuer=get_default_issuer(settings.environment), ) @provider(scope="singleton") async def dynamodb_client(self, settings: Settings) -> AsyncIterator[Client]: async with AsyncClient() as web_client: yield Client( http=HTTPX(web_client), credentials=Credentials.auto(), region=settings.aws_region_name, # This function will be applied only to Number type (N) parameters # Default function is just a float which will convert "1" to 1.0 numeric_type=lambda value: float(value) if "." in value else int(value), ) @provider(scope="singleton") def dynamodb_write_semaphore(self, settings: Settings) -> Semaphore: return Semaphore(settings.dynamodb_write_threads_count) @provider(scope="singleton") def repository( self, dynamodb_client: Client, settings: Settings, dynamodb_write_semaphore: Semaphore, ) -> PathsRepository: return PathsRepository( dynamodb_client=dynamodb_client, dynamodb_tablename=settings.dynamodb_tablename, semaphore=dynamodb_write_semaphore, dynamodb_max_write_retry_count=settings.dynamodb_max_write_retry_count, ) @provider(scope="singleton") def s3_client(self, settings: Settings) -> S3Client: client: S3Client = S3Client(region_name=settings.aws_region_name) return client @provider(scope="singleton") async def kafka_client(self, settings: Settings) -> AsyncIterator[KafkaClient]: kafka_client = KafkaClient( bootstrap_servers=settings.kafka_bootstrap_servers, use_ssl=settings.kafka_use_ssl, topic=settings.kafka_analytics_topic, ) await kafka_client.start() yield kafka_client await kafka_client.stop() @provider(scope="singleton") def short_url_redirect_service( self, repository: PathsRepository, kafka_client: KafkaClient, s3_client: S3Client, ) -> ShortUrlRedirectService: return ShortUrlRedirectService( repository=repository, kafka_client=kafka_client, s3_client=s3_client, failed_events_s3_bucket=settings.failed_events_s3_bucket, ) @provider(scope="singleton") def shorten_url_service( self, repository: PathsRepository, settings: Settings ) -> ShortenUrlService: return ShortenUrlService( repository=repository, allowed_domains=settings.allowed_domains, short_path_allowed_chars=settings.short_path_allowed_chars, ) @provider(scope="singleton") def delete_paths_service(self, repository: PathsRepository) -> DeletePathsService: return DeletePathsService(repository=repository) @provider(scope="singleton") def update_path_service( self, repository: PathsRepository, settings: Settings, ) -> UpdatePathService: return UpdatePathService( repository=repository, allowed_domains=settings.allowed_domains, ) def setup_container() -> Container: """Configure the application.""" container = Container( providers=[ Provider(call=lambda: settings, scope="singleton", interface=Settings), ], modules=[AppModule], default_scope="singleton", ) if settings.environment == "test": return TestContainer.from_container(container) return container # Lazy container proxy container = lazy_proxy(setup_container)