"""Application Handlers. Requests are redirected to handlers, which are responsible for getting information from the URL and passing it down to the logic layer. The way each layer talks to each other is through Response objects which defines the type status of the data and the data itself. Please note: the Orchard uses the term handlers over views as convention for clarity See: oto.response for more details. """ from connector_neo4j import Neo4jSession from connector_neo4j import exceptions as neo4j_connector_exceptions from flask import g from flask import jsonify import neo4j from owsresponse import response from owsresponse.adaptors.flask import flaskify from sound_recordings import config from sound_recordings.api import app from sound_recordings.constants import error from sound_recordings.logic import acrids from sound_recordings.logic import fingerprint_rules from sound_recordings.logic import orchard_sound_recordings from sound_recordings.models.ows_track import BadGateway as OwsTrackBadGateway from sound_recordings.utils import request_access from sound_recordings.validation import schemas from webargs import ValidationError from webargs import fields from webargs import validate from webargs.flaskparser import use_args def _require_one_of(*fields): def checker(args): fields_set = any( field in args and bool(args[field]) for field in fields ) if not fields_set: raise ValidationError( f"At least one field required from {','.join(fields)}" ) return checker @app.route(config.HEALTH_CHECK, methods=['GET']) def health(): """Check the health of the application.""" return jsonify({'status': 'ok'}) @app.route('/acrids', methods=['POST']) @use_args( { 'acr_id': fields.Str(required=True), 'asset_id': fields.UUID(required=True) } ) def acrids_post(args): """Upsert OrchardAsset -> ACRID in neo4j. Args: see @use_args above Returns: dict: track_id(s) attached to asset """ (tuids, updates) = acrids.upsert(args['acr_id'], args['asset_id']) status_code = 201 if updates else 200 message = {'track_ids': tuids} return flaskify( response.Response( status=status_code, message=message ) ) @app.route('/sound_recordings', methods=['POST']) @use_args( { 'track_id': fields.Int(required=True) } ) def sound_recordings_post(args): """Upsert OrchardSoundRecording in neo4j. Args: see @use_args above Returns: dict: orchard_sound_recording_id from OrchardSoundRecording.id """ try: (sr_id, updates) = orchard_sound_recordings.upsert( args['track_id'] ) status_code = 201 if updates else 200 message = {'orchard_sound_recording_id': sr_id} except ( orchard_sound_recordings.AssetNotFound, orchard_sound_recordings.TrackISRCNotSet, orchard_sound_recordings.MultipleAcrOsrRelationships ) as e: status_code = 409 error_name = type(e).__name__ message = {'error': error_name} # Enrich the autolog response line with the 409 cause and track_id # (exposed as searchable @resources.* facets in Datadog) rather than # emitting a separate log entry. g.log.extra_message = f'{error_name} for track_id={args["track_id"]}' g.log.resources.update({ 'error': error_name, 'track_id': args['track_id'] }) return flaskify( response.Response( status=status_code, message=message ) ) @app.route('/sound_recordings/', methods=['PATCH']) @use_args({'osr_id': fields.UUID()}, location='view_args') @use_args( { 'primary_track_id': fields.Int(required=True) } ) def sound_recordings_patch(self, args, osr_id): """Update OrchardSoundRecording in neo4j. Args: see @use_args above Returns: dict: OrchardSoundRecording """ try: osr = orchard_sound_recordings.update( osr_id, args['primary_track_id'] ) status_code = 200 message = osr except ( orchard_sound_recordings.OrchardSoundRecordingNotFound ) as e: status_code = 404 message = {'error': type(e).__name__} except ( orchard_sound_recordings.PrimaryTrackNotRelated, orchard_sound_recordings.OrchardSoundRecordingWithoutAssets ) as ex: status_code = 409 message = {'error': type(ex).__name__} return flaskify( response.Response( status=status_code, message=message ) ) @app.route('/sound_recordings', methods=['GET']) @use_args( { 'ids': fields.DelimitedList(fields.UUID(), delimiter=',', load_default=[]), 'tuids': fields.DelimitedList(fields.Int(), delimiter=',', load_default=[]), 'upcs': fields.DelimitedList(fields.Int(), delimiter=',', load_default=[]), 'product_ids': fields.DelimitedList( fields.Int(), delimiter=',', load_default=[] ), 'project_ids': fields.DelimitedList( fields.Int(), delimiter=',', load_default=[] ), 'asset_ids': fields.DelimitedList(fields.UUID(), delimiter=',', load_default=[]), # noqa:E501 'track_isrcs': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'include_deleted': fields.Boolean(load_default=False), 'include_transfer_to_content': fields.Boolean(load_default=False), 'term': fields.Str(load_default=''), 'include_inactive': fields.Boolean(load_default=False) }, validate=_require_one_of( 'ids', 'tuids', 'upcs', 'product_ids', 'project_ids', 'asset_ids', 'track_isrcs', 'term' ), location='query' ) @Neo4jSession(use_v2=True, database=config.NEO4J_DATABASE_NAME) def sound_recordings_get(args): """Fetch list of orchard sound recording, with acrid and asset info. Args: see @use_args above Returns: list: orchard sound recordings """ include_inactive = args['include_inactive'] results = orchard_sound_recordings.fetch( orchard_sound_recording_ids=args['ids'], track_ids=args['tuids'], upcs=args['upcs'], product_ids=args['product_ids'], project_ids=args['project_ids'], asset_ids=args['asset_ids'], track_isrcs=args['track_isrcs'], # if we use include_inactive, we also need to include deleted assets include_deleted=include_inactive or args['include_deleted'], include_transfer_to_content=args['include_transfer_to_content'], term=args['term'], include_inactive=include_inactive ) return flaskify( response.Response( status=200, message=results ) ) @app.route('/sound_recordings/cr', methods=['GET']) @use_args( { 'product_id': fields.Int(load_default=None), 'track_isrcs': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'include_transfer_to_content': fields.Boolean(load_default=False) }, location='query' ) @Neo4jSession(use_v2=True, database=config.NEO4J_DATABASE_NAME) def sound_recordings_product_get(args): """Fetch list of orchard sound recordings for Content Review (CR). Args: see @use_args above Returns: list: orchard sound recordings """ results = orchard_sound_recordings.fetch_product( args['product_id'], args['track_isrcs'], args['include_transfer_to_content'] ) return flaskify( response.Response( status=200, message=results ) ) @app.route('/sound_recordings//versions/', methods=['GET']) @use_args( { 'osr_id': fields.UUID(required=True), 'version_id': fields.Str(required=True) }, location='view_args' ) def sound_recordings_version_get(self, osr_id, version_id): """Fetch orchard sound recording for a specific version. Args: osr_id (str): The Orchard Sound Recording ID version_id (str): The version ID for the Orchard Sound Recording Returns: dict: Orchard sound recording data for a specific version """ result = orchard_sound_recordings.fetch_version(osr_id, version_id) if not result: return flaskify( response.create_not_found_response(error.ERROR_MESSAGE_SR_NOT_FOUND) ) return flaskify( response.Response( status=200, message=result ) ) @app.route('/sound_recording//full_delivery_history', methods=['GET']) @use_args( { 'execution_type': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'service': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'event_type': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'ack_status': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 }, location='query' ) def sound_recordings_full_delivery_history_get(query_args, osr_id): """Fetch orchard sound recording full delivery history events. Args: osr_id (str): The Orchard Sound Recording ID execution_type (list): List of execution types to filter by (e.g. ['FULL_DELIVERY']) service (list): List of UGC services to filter by event_type (list): List of event types to filter by ack_status (list): List of ack statuses to filter by (e.g. ['success', 'awaited', 'error', 'missing']) Returns: dict: Orchard sound recording delivery history data """ result = orchard_sound_recordings.fetch_full_delivery_history( [osr_id], query_args['execution_type'], query_args['service'], query_args['event_type'], query_args['ack_status'] ) return flaskify( response.Response( status=200, message=result ) ) @app.route('/sound_recordings/full_delivery_history', methods=['POST']) @use_args( { 'execution_type': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'service': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'event_type': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'ack_status': fields.DelimitedList(fields.Str(), delimiter=',', load_default=[]), # noqa:E501 'osr_ids': fields.List(fields.Str(), load_default=[]), # noqa:E501 'limit': fields.Int(validate=validate.Range(min=1, max=1000), load_default=None), # noqa:E501 'offset': fields.Int(validate=validate.Range(min=0), load_default=None), 'start_date': fields.DateTime(load_default=None) } ) def sound_recordings_full_delivery_history(args): """Fetch for many orchard sound recordings full delivery history events. Args: osr_ids (list): The Orchard Sound Recording IDs execution_type (list): List of execution types to filter by (e.g. ['FULL_DELIVERY']) service (list): List of UGC services to filter by event_type (list): List of event types to filter by ack_status (list): List of ack statuses to filter by (e.g. ['success', 'awaited', 'error', 'missing']) limit (int): Maximum number of records to return (default: None) offset (int): Number of records to skip (default: None) start_date (datetime): Filter records from this date onwards (default: None) Returns: dict: Orchard sound recording delivery history data """ result = orchard_sound_recordings.fetch_full_delivery_history( args['osr_ids'], args['execution_type'], args['service'], args['event_type'], args['ack_status'], args['limit'], args['offset'], args['start_date'] ) return flaskify( response.Response( status=200, message=result ) ) @app.route('/sound_recordings/touch', methods=['PATCH']) @use_args( { 'resource_id': fields.Int(required=True), 'resource_type': fields.Str( required=True, validate=validate.OneOf(['vendor', 'subaccount'])), 'modified_before': fields.DateTime(required=True), 'limit': fields.Int(validate=validate.Range(min=1, max=10000), required=True) } ) def sound_recordings_touch(args): """Update OrchardSoundRecording(s) based on vendor/date in neo4j. Args: see @use_args above Returns: dict: OrchardSoundRecording """ result = orchard_sound_recordings.touch( args['resource_id'], args['resource_type'], args['modified_before'], args['limit'] ) return flaskify( response.Response( status=200, message=result ) ) @app.route('/vendors/rules', methods=['GET']) @use_args( { 'ids': fields.DelimitedList(fields.Int(), delimiter=',', load_default=[]) }, location='query' ) @Neo4jSession(use_v2=True, database=config.NEO4J_DATABASE_NAME) def vendor_fingerprint_rules_get(args): """Return set of fingerprint rules for vendors. Args: see @use_args above Returns: request.Response """ vendor_ids = list(set(args['ids'])) (has_access, profile) = request_access.check_access( vendor_ids, 'Vendor' ) if not has_access: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) rules = fingerprint_rules.match_bulk('Vendor', vendor_ids) return flaskify( response.Response( status=200, message=rules ) ) @app.route('/subaccounts/rules', methods=['GET']) @use_args( { 'ids': fields.DelimitedList(fields.Int(), delimiter=',', load_default=[]) }, location='query' ) @Neo4jSession(use_v2=True, database=config.NEO4J_DATABASE_NAME) def subaccount_fingerprint_rules_get(args): """Return set of fingerprint rules for subaccounts. Args: see @use_args above Returns: request.Response """ subaccount_ids = list(set(args['ids'])) (has_access, profile) = request_access.check_access( subaccount_ids, 'Subaccount' ) if not has_access: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) rules = fingerprint_rules.match_bulk('Subaccount', subaccount_ids) return flaskify( response.Response( status=200, message=rules ) ) @app.route('/tracks/rules', methods=['GET']) @use_args( { 'ids': fields.DelimitedList(fields.Int(), delimiter=',', load_default=[]) }, location='query' ) @Neo4jSession(use_v2=True, database=config.NEO4J_DATABASE_NAME) def track_fingerprint_rules_get(args): """Return set of fingerprint rules for tracks. Args: see @use_args above Returns: request.Response """ track_ids = list(set(args['ids'])) (has_access, profile) = request_access.check_access( track_ids, 'Track' ) if not has_access: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) rules = fingerprint_rules.match_bulk('Track', track_ids) return flaskify( response.Response( status=200, message=rules ) ) @app.route('//bulk/rules', methods=['POST']) @use_args(schemas.FingerprintBulkRuleURLSchema(), location='view_args') @use_args(schemas.FingerprintBulkRuleQuerySchema(), location='query') @use_args(schemas.FingerprintBulkRuleSchema()) @Neo4jSession(transaction=True, use_v2=True, database=config.NEO4J_DATABASE_NAME) def fingerprint_bulk_rules_post(url_parts, query_args, body, **kwargs): """Save set of fingerprint rules for the several objects. When `services` query parameter is provided, only rules for those services will be affected. Rules for other services remain untouched. If omitted, all services are affected (full replacement behavior). Args: see @use_args above Returns: request.Response """ obj_ids = body.keys() affected_services = query_args.get('services') if affected_services and '*' in affected_services: affected_services = None elif affected_services is not None: if affected_services == []: return flaskify( response.Response( status=422, message={ 'error': 'EmptyAffectedServices', 'message': "Query parameter 'services' cannot be empty if provided." # noqa:E501 } ) ) invalid_services = { rule['service'] for rules in body.values() for rule in rules if rule['service'] not in affected_services } if invalid_services: return flaskify( response.Response( status=422, message={ 'error': 'ServiceNotInAffectedServices', 'services': sorted(invalid_services) } ) ) (has_access, profile) = request_access.check_access( list(map(int, obj_ids)), url_parts['obj_type'] ) if not has_access or not profile: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) total_num_created = 0 total_num_deleted = 0 for obj_id in obj_ids: (num_created, num_deleted) = fingerprint_rules.upsert( url_parts['obj_type'], int(obj_id), new_rules=body[obj_id], profile=profile, affected_services=affected_services ) total_num_created += num_created total_num_deleted += num_deleted return flaskify( response.Response( status=200, message={ 'num_created': total_num_created, 'num_deleted': total_num_deleted } ) ) @app.route('///rules', methods=['POST']) @use_args(schemas.FingerprintRuleURLSchema(), location='view_args') @use_args(schemas.FingerprintRuleSchema(many=True)) @Neo4jSession(transaction=True, use_v2=True, database=config.NEO4J_DATABASE_NAME) def fingerprint_rules_post(url_parts, body, **kwargs): """Save set of fingerprint rules for this object. Args: see @use_args above Returns: request.Response """ (has_access, profile) = request_access.check_access( [url_parts['obj_id']], url_parts['obj_type'] ) if not has_access or not profile: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) (num_created, num_deleted) = fingerprint_rules.upsert( url_parts['obj_type'], url_parts['obj_id'], body, profile ) return flaskify( response.Response( status=200, message={ 'num_created': num_created, 'num_deleted': num_deleted } ) ) @app.route('///rules/history', methods=['GET']) @use_args(schemas.FingerprintRuleURLSchema(), location='view_args') @use_args( { 'service': fields.Str(load_default=''), 'policy': fields.Str(load_default='') }, location='query' ) @Neo4jSession(transaction=False, use_v2=True, database=config.NEO4J_DATABASE_NAME) def sound_recordings_rules_history(url_parts, query_args, **kwargs): """Get set of rules history for this object. Args: see @use_args above Returns: request.Response """ (has_access, profile) = request_access.check_access( [url_parts['obj_id']], url_parts['obj_type'] ) if not has_access: return flaskify( response.Response( status=403, message={'error': 'NotAuthorizedOperation'} ) ) history = fingerprint_rules.fetch_history( url_parts['obj_type'], url_parts['obj_id'], query_args ) message = {'history': history or []} return flaskify( response.Response( status=200, message=message ) ) @app.errorhandler(OwsTrackBadGateway) def handle_ows_track_bad_gateway(e): """Handle bad gateway errors from ows-track.""" g.log.exception(e) return flaskify( response.create_error_response( 'bad_gateway', 'ows-track', 502 ) ) # https://github.com/neo4j/neo4j-python-driver/blob/4.3/neo4j/exceptions.py#L141 @app.errorhandler(neo4j.exceptions.DatabaseError) def handler_database_error_neo4j(e): """Handle database neo4j errors.""" g.log.exception(e) return flaskify( response.create_error_response( 'database_error', 'neo4j database error', 502 ) ) # https://github.com/neo4j/neo4j-python-driver/blob/4.3/neo4j/exceptions.py#L146 @app.errorhandler(neo4j.exceptions.TransientError) def handle_transient_neo4j(e): """Handle transient neo4j errors.""" if e.is_retryable(): g.log.exception(e) return flaskify( response.create_error_response( str(e.code), str(e.message), 502 ) ) raise e # https://github.com/neo4j/neo4j-python-driver/blob/4.3/neo4j/exceptions.py#L247 @app.errorhandler(neo4j.exceptions.SessionExpired) def handle_session_expired_neo4j(e): """Handle session neo4j errors.""" g.log.exception(e) return flaskify( response.create_error_response( 'session_expired', 'neo4j session error', 502 ) ) # https://github.com/neo4j/neo4j-python-driver/blob/4.3/neo4j/exceptions.py#L274 @app.errorhandler(neo4j.exceptions.ServiceUnavailable) def handle_service_unavailable_neo4j(e): """Handle service neo4j errors.""" g.log.exception(e) return flaskify( response.create_error_response( 'service_unavailable', 'neo4j service error', 502 ) ) @app.errorhandler(neo4j_connector_exceptions.SessionNotCreated) def handle_session_not_created_neo4j(e): """Handle connector session neo4j errors.""" g.log.exception(e) return flaskify( response.create_error_response( 'session_not_created', 'neo4j session error', 502 ) ) @app.errorhandler(422) def handle_bad_request(e): """Handle exceptions that should result in HTTP 400.""" return flaskify( response.create_error_response( 'bad_request', e.data.get('messages', ['Invalid Input.']), 422 ) ) @app.errorhandler(404) def handle_not_found(e): """Handle exceptions that should result in HTTP 404.""" return flaskify( response.create_not_found_response(e.description) ) @app.errorhandler(500) def exception_handler(error): """Handle error when uncaught exception is raised. Default exception handler. Note: Exception will also be sent to Sentry if config.SENTRY is set. Returns: flask.Response: A 500 response with JSON 'code' & 'message' payload. """ message = ( 'The server encountered an internal error ' 'and was unable to complete your request.' ) g.log.exception(error) return flaskify(response.create_fatal_response(message))