"""Get new release data and write it to S3.""" import json import os from operator import itemgetter import boto3 import yaml from bravado_core.spec import Spec from bravado_core.validate import validate_object from owsrequest import request as owsrequest from snowflake_connector.snowflake_conn import get_session, set_default_sessionmaker from sqlalchemy import text import src.config as config from src.models import ows_assets def _get_swagger_spec_file(): """Gets the swagger spec file. Return: File: Swagger spec. """ dir_path = os.path.dirname(os.path.abspath(__file__)) spec_path = os.path.join(dir_path, config.LOCAL_SWAGGER_SPEC_PATH) with open(spec_path, "r") as spec: return spec.read() def _validate_against_swagger(response): """Validate the new release data json against the swagger spec. Args: response (dict): New release data. """ spec_dict = yaml.safe_load(_get_swagger_spec_file()) swagger_spec = Spec.from_dict(spec_dict) object_spec = spec_dict["definitions"]["ReleasesResponse"] validate_object(swagger_spec, object_spec, response) def _extract_project_artist(row): """Extract project artist fields from a us_physical_new_release_view row. Args: row (dict): A us_physical_new_release_view row. Return: dict: Project artist fields extracted from the us_physical_new_release_view row. """ project_artist_sites = json.loads(row["project_artist_urls"]) return { "artist_id": row["project_artist_id"], "name": row["project_artist_name"], "bio": row["project_artist_description"], "website_url": row["project_artist_website"], "facebook_url": project_artist_sites.get("facebook"), "twitter_url": project_artist_sites.get("twitter"), "youtube_channel_url": project_artist_sites.get("youtube"), "instagram_url": project_artist_sites.get("instagram"), "soundcloud_url": project_artist_sites.get("soundcloud"), "bandsintown_url": project_artist_sites.get("bandsintown"), "for_fans_of": json.loads(row["project_for_fans_of"]) or None, "home_market": row["project_artist_home_market"], } def _extract_project(row): """Extract project fields from a us_physical_new_release_view row. Args: row (dict): A us_physical_new_release_view row. Return: dict: Project fields extracted from the us_physical_new_release_view row. """ return { "project_id": row["project_id"], "name": row["project_name"], "type": "Stak", "book_dates": [], "book_order": 0, "description": row["project_description"], "highlights": row["project_highlights"], "video_urls": json.loads(row["project_artist_videos"]) or None, "label_contact": "sales@theorchard.com", "label_nm": row["product_imprint"], "parent_label_nm": row["vendor_name"] if row["subaccount_id"] else None, "theme": "Dark", "background_color": "131313", "font_color": "FFFFFF", "artist": _extract_project_artist(row), "products": [], } def _update_project_aggregate_fields(row, project): """Update a project's aggregate fields with a us_physical_new_release_view row. Args: row (dict): A us_physical_new_release_view row. project (dict): project whose aggregate fields should be updated. """ # Ensure book dates are unique. book_dates_set = set(project["book_dates"]) book_dates_set.add(str(row["product_sales_date"])) project["book_dates"] = list(book_dates_set) first_wk_ss = row["product_first_week_sales_estimate_sum"] project["book_order"] += first_wk_ss or 0 def _get_config_group_map(): """Get distribution_format_id to config_group mapping for 45Press. Return: dict: Mapping from distribution_format_id to 45Press config_group. """ product_configuration_response = owsrequest.process( config.APPLICATION_NAME, config.ENVIRONMENT, "GET", "ows-product-configuration", "/supply-chain/4/distribution-formats", headers={"User-Agent": config.APPLICATION_NAME}, ) if product_configuration_response.status_code != 200: raise Exception( "ows-product-configuration returned status: {}".format( product_configuration_response.status_code ) ) config_groups = product_configuration_response.json() config_group_map = {} for config_group in config_groups: distribution_format_id = config_group["distribution_format_id"] config_group_map[distribution_format_id] = config_group["config_group"] return config_group_map def _extract_product(row, config_group_map): """Extract product fields from a us_physical_new_release_view row. Args: row (dict): A us_physical_new_release_view row. config_group_map (dict): Mapping from distribution_format_id to 45Press config_group. Return: dict: Product fields extracted from the us_physical_new_release_view row. """ return { "product_id": row["product_id"], "product_code": row["product_code"], "version": (row["product_version"] or "").strip() or None, "sales_dt": str(row["product_sales_date"]), "order_due_dt": str(row["product_order_due_date"]), "release_dt": str(row["product_release_date"]), "price_code": row["product_price_code"], "upc": row["product_upc"], "box_lot": row["product_box_lot"], "returnable": row["product_is_returnable"], "config": row["product_display_configuration"], "config_group": config_group_map[str(row["product_distribution_format_id"])], "highlights": row["product_highlights"], "marketing_highlights": row["product_marketing_highlights"], "artist": row["product_primary_artist_name"], "title": row["product_name"], "genre": row["product_genre"], "subgenre": row["product_subgenre"], "country_of_origin": row["product_country_of_origin"], "tracks": sorted( json.loads(row["product_tracks"]), key=itemgetter("disc", "side", "track_number"), ) or None, } def _get_us_physical_new_release_view_rows(): """Query Snowflake for new release data. Return: array: All records in us_physical_new_release_view. """ set_default_sessionmaker( connect_args=config.SNOWFLAKE_CONNECT_ARGS, sf_config=config.SNOWFLAKE_OPTIONS ) with get_session() as session: rows = session.execute(text(config.US_PHYSICAL_NEW_RELEASE_VIEW_QUERY)) return rows.mappings().fetchall() def _save_to_s3(new_release_data): """Write new release data to S3. Args: new_release_data (dict): New release data to save to S3. """ s3 = boto3.client("s3") s3.put_object( Body=json.dumps(new_release_data), Bucket=config.S3_BUCKET, Key=config.NEW_RELEASE_DATA_S3_KEY, ) s3.put_object( Body=_get_swagger_spec_file(), Bucket=config.S3_BUCKET, Key=config.SWAGGER_SPEC_S3_KEY, ) def _get_new_release_data(lambda_event): """Get new release data. Args: lambda_event (dict): AWS Lambda Event Return: dict: New release data formatted according to spec/swagger.yml. """ config_group_map = _get_config_group_map() projects = {} product_ids = [] products = [] for row in _get_us_physical_new_release_view_rows(): project_id = row["project_id"] if not projects.get(project_id): projects[project_id] = _extract_project(row) product = _extract_product(row, config_group_map) products.append(product) projects[project_id]["products"].append(product) product_ids.append(product["product_id"]) _update_project_aggregate_fields(row, projects[project_id]) cover_360_urls = ows_assets.fetch_covers(product_ids, "large_cover") cover_2000_urls = ows_assets.fetch_covers(product_ids, "xlarge_cover") for product in products: product_id = str(product["product_id"]) product["cover_art_img_url"] = cover_360_urls[product_id] product["cover_art_img_url_2000"] = cover_2000_urls[product_id] return {"projects": list(projects.values()), "last_modified": lambda_event["time"]} def handler(event, context): """Get new release data and write it to S3. Args: event (dict): AWS Lambda Event context (dict): AWS Lambda Context """ new_release_data = _get_new_release_data(event) _validate_against_swagger(new_release_data) _save_to_s3(new_release_data) return {"status": "OK"}