import asyncio from functools import partial import datadog import sentry_sdk import structlog from sqlalchemy import event from sqlalchemy.ext.asyncio import AsyncEngine, create_async_engine from dapd_api_scraper import __version__ from dapd_api_scraper.config import Config from dapd_api_scraper.entities import ( AppleMusicAlbumEntity, AppleMusicArtistEntity, AppleMusicPlaylistEntity, AppleMusicTrackEntity, SpotifyAlbumEntity, SpotifyArtistEntity, SpotifyPlaylistEntity, SpotifyTrackEntity, ) from dapd_api_scraper.scraper import process_entity from dapd_api_scraper.utils import ShutdownHandler async def entity_starter( config: Config, logger: structlog.BoundLogger, workflowdb_engine: AsyncEngine, etldb_engine: AsyncEngine, shutdown: ShutdownHandler, entity, ) -> None: """This is the entypoint for entity Here we create separate sentry hub and start processing """ with sentry_sdk.Hub(sentry_sdk.Hub.current): sentry_sdk.set_tag("data_source", entity.data_source) sentry_sdk.set_tag("entity", entity.entity_type) await process_entity(logger, shutdown, workflowdb_engine, etldb_engine, entity(config)) async def main_loop() -> None: """ Here we initialize common objects and start all needed entities in separate async coroutines""" config = Config() logger = structlog.get_logger( processors=[ structlog.processors.TimeStamper(fmt="iso", key="timestamp"), structlog.processors.add_log_level, structlog.processors.JSONRenderer(), ], ) shutdown = ShutdownHandler(logger) sentry_sdk.init( dsn=config.sentry_dsn, environment=config.env, release=__version__, debug=config.debug, max_breadcrumbs=10, ) datadog.initialize( statsd_host=config.datadog_host, statsd_port=config.datadog_port, ) workflowdb_engine = create_async_engine( config.workflowdb_uri, pool_size=64, max_overflow=128, connect_args={"server_settings": {"application_name": f"{config.env}-dapd-api-scraper"}}, isolation_level="AUTOCOMMIT", echo=config.debug, ) etldb_engine = create_async_engine( config.etldb_uri, pool_size=64, max_overflow=128, connect_args={"server_settings": {"application_name": f"{config.env}-dapd-api-scraper"}}, isolation_level="AUTOCOMMIT", echo=config.debug, ) @event.listens_for(workflowdb_engine.sync_engine, "do_connect") def workflowdb_engine_password_rotation_handler(dialect, conn_rec, cargs, cparams): credentials = config.get_secret(config.workflowdb_secret_name) cparams["user"] = credentials["username"] cparams["password"] = credentials["password"] cparams["host"] = credentials["active_endpoint"] cparams["port"] = credentials["port"] cparams["database"] = credentials["database"] @event.listens_for(etldb_engine.sync_engine, "do_connect") def etldb_engine_password_rotation_handler(dialect, conn_rec, cargs, cparams): credentials = config.get_secret(config.etldb_secret_name) cparams["user"] = credentials["username"] cparams["password"] = credentials["password"] cparams["host"] = credentials["active_endpoint"] cparams["port"] = credentials["port"] cparams["database"] = credentials["database"] entities_list = { "AppleMusicAlbum": AppleMusicAlbumEntity, "AppleMusicArtist": AppleMusicArtistEntity, "AppleMusicPlaylist": AppleMusicPlaylistEntity, "AppleMusicTrack": AppleMusicTrackEntity, "SpotifyAlbum": SpotifyAlbumEntity, "SpotifyArtist": SpotifyArtistEntity, "SpotifyPlaylist": SpotifyPlaylistEntity, "SpotifyTrack": SpotifyTrackEntity, } if config.entities_list: entities = [ entity for entity_name, entity in entities_list.items() if entity_name in config.entities_list ] else: entities = list(entities_list.values()) # healthcheck file init with open("healthcheck.json", "w", encoding="utf-8") as healthcheck: healthcheck.write("{}") await asyncio.gather( *map( partial( entity_starter, config, logger, workflowdb_engine, etldb_engine, shutdown, ), entities, ) ) def main(): asyncio.run(main_loop()) if __name__ == "__main__": main()