# pylint: disable=redefined-outer-name,unused-argument import json import os from time import sleep from unittest.mock import create_autospec import boto3 import pytest from boto3_type_annotations.kinesis import Client as KinesisClient from boto3_type_annotations.s3 import Client as S3Client from dapd_db_schema.schemas.etl import ( DimAlbum, DimAlbumArtist, DimArtist, DimDsp, DimMarket, DimPlaylist, DimPlaylistMeta, DimTrack, DimTrackAlbum, DimTrackArtist, FactArtistFollowers, ) from moto import mock_kinesis, mock_s3 from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker from structlog.stdlib import BoundLogger from dapd_transformation_service.entities.dimension_meta import DimT, DspEnum from dapd_transformation_service.services.config import ConfigService @pytest.fixture(scope='session') def config(): config_service = ConfigService() yield config_service.load_config() @pytest.fixture(scope='function') def db(): if os.getenv('IS_INTEGRATION_TESTING'): test_engine = create_engine(os.getenv('DB_URL')) test_session_class = sessionmaker(bind=test_engine) test_session = test_session_class(autocommit=True, autoflush=True) # Simple hack for waiting migrations to be completed while True: result = test_session.execute( # pylint: disable=no-member 'SELECT id FROM databasechangelog ORDER BY dateexecuted DESC LIMIT 1' ) if result.first(): break sleep(1) yield test_session else: yield @pytest.fixture(scope='function', autouse=True) def clean_db(db): yield if os.getenv('IS_INTEGRATION_TESTING'): for model in [ FactArtistFollowers, # FactPlaylistTrackDynamics, DimTrackArtist, DimTrackAlbum, DimAlbumArtist, DimTrack, DimAlbum, DimArtist, DimPlaylistMeta, DimPlaylist, DimMarket, DimDsp, ]: db.query(model).delete() @pytest.fixture def dim_dsp(db: Session): instance = DimDsp(dsp_name='apple_music') db.add(instance) db.flush() yield instance @pytest.fixture def dim_market(db: Session): instance = DimMarket( **{ 'market_id': 1, 'market_code': 'global', 'market_name': 'foo_name', 'market_full_name': None, 'territory_type_id': None, 'parent_market_id': None, 'created_at': '2020-10-30 09:51:33.191551', 'updated_at': '2020-10-30 09:51:33.191000' } ) db.add(instance) db.flush() yield instance @pytest.fixture def dim_playlist(db: Session, dim_dsp: DimDsp, dim_market: DimMarket): instance = DimPlaylist( **{ 'dsp_id': dim_dsp.dsp_id, 'market_id': dim_market.market_id, 'playlist_id': '0076R65zbUgOgxM5tzWtsR', 'uri': 'spotify:playlist:0076R65zbUgOgxM5tzWtsR', 'type': None, 'is_public': False, 'is_personalized': False, 'collect_historical_data': False, 'updated_at': '2020-10-15 05:42:24.000000', 'created_at': '2020-10-21 22:52:53.335658' } ) db.add(instance) db.flush() yield instance @pytest.fixture def dim_artist(db: Session, dim_dsp: DimDsp, dim_market: DimMarket): instance = DimArtist( **{ 'dsp_id': dim_dsp.dsp_id, 'market_id': dim_market.market_id, 'artist_id': 'foo_id', 'artist_name': 'Foo Name', 'genres': None, 'image_path': None, 'artist_uri': 'foo_id_uri', 'is_blacklisted': False, 'created_at': '2020-10-30 09:49:05.606156', 'updated_at': '2020-10-30 09:49:05.606156', } ) db.add(instance) db.flush() yield instance @pytest.fixture def dim_album(db: Session, dim_dsp: DimDsp, dim_market: DimMarket): instance = DimAlbum( **{ 'dsp_id': dim_dsp.dsp_id, 'album_id': 'foo_id', 'album_name': 'Foo Name', 'popularity': None, 'release_date': None, 'release_date_precision': None, 'market_id': dim_market.market_id, 'is_mastered_for_itunes': False, 'is_complete': True, 'track_count': None, 'copyright': None, 'lbl': None, 'image_path': None, 'upc': None, 'created_at': '2020-10-30 09:53:57.138828', 'updated_at': '2020-10-30 09:53:57.138000' } ) db.add(instance) db.flush() yield instance @pytest.fixture def dim_track(db: Session, dim_dsp: DimDsp, dim_market: DimMarket): instance = DimTrack( **{ 'dsp_id': dim_dsp.dsp_id, 'market_id': dim_market.market_id, 'track_id': '2dNxQWaUJp1FWa7vhcI597', 'track_name': '062 - Spuk am Himmel - Teil 02', 'track_duration': 104, 'isrc': 'DEC711904690', 'audio_features': None, 'composer': None, 'upc': None, 'updated_at': '2020-10-24 22:11:00.385252', 'created_at': '2020-10-24 22:11:00.385252', 'release_date': '2020-10-24' } ) db.add(instance) db.flush() yield instance @pytest.fixture def raw_data(dsp: DspEnum, dim_type: DimT): current_file_dir = os.path.dirname(os.path.realpath(__file__)) sample_dir = os.path.join(current_file_dir, 'samples') sample_file_name = f'{dsp.value}_{dim_type.__name__[3:].lower()}.json' sample_path = os.path.join(sample_dir, sample_file_name) with open(sample_path, 'rb') as sample_fp: yield sample_fp.read() @pytest.fixture(scope='function', autouse=True) def aws_credentials(): os.environ['AWS_ACCESS_KEY_ID'] = 'testing' os.environ['AWS_SECRET_ACCESS_KEY'] = 'testing' os.environ['AWS_SECURITY_TOKEN'] = 'testing' os.environ['AWS_SESSION_TOKEN'] = 'testing' os.environ['AWS_DEFAULT_REGION'] = 'us-east-1' boto3.setup_default_session( aws_access_key_id='testing', aws_secret_access_key='testing', aws_session_token='testing', region_name='testing', ) @pytest.fixture(scope='function') def kinesis(aws_credentials, config, dim_dsp): with mock_kinesis(): kinesis_client: KinesisClient = boto3.client('kinesis') for entity in ['playlists', 'albums', 'artists', 'tracks']: stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' data = {'entity': entity} kinesis_client.create_stream(StreamName=stream_name, ShardCount=10) kinesis_client.put_record( StreamName=stream_name, Data=json.dumps(data).encode(), PartitionKey='foo_key' ) yield kinesis_client @pytest.fixture def logger(): yield create_autospec(BoundLogger) @pytest.fixture def s3_client() -> S3Client: with mock_s3(): test_bucket_name = 'test_bucket' _s3_client: S3Client = boto3.client('s3') _s3_client.create_bucket(Bucket=test_bucket_name) yield _s3_client s3_resource = boto3.resource('s3') bucket = s3_resource.Bucket(test_bucket_name) bucket.objects.all().delete() bucket.delete()