"""Logic of API responses.""" from typing import Dict import config from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from lambdacommon.common_config import logger from pydantic import ValidationError from src.models.spotify import SpotifyArtist from src.models.spotify import SpotifyAudioFeature from src.models.spotify import SpotifyGetAudioFeaturesRequest from src.models.spotify import SpotifySearchAlbumArtistMessage from src.models.spotify import SpotifySearchAlbumByUPCRequest from src.models.spotify import SpotifySearchAlbumMessage from src.models.spotify import SpotifySearchAlbumsResponse from src.models.spotify import SpotifySearchTrackArtistMessage from src.models.spotify import SpotifySearchTrackByISRCRequest from src.models.spotify import SpotifySearchTrackMessage from src.models.spotify import SpotifySearchTracksResponse string_serializer = StringSerializer() producer = EventProducer( config.KAFKA_BROKERS, string_serializer, string_serializer, 'SSL') def submit_spotify_api_product_search_response( api_request: SpotifySearchAlbumByUPCRequest, api_response: Dict): """Submit Spotify API response to Kafka. Args: api_request: Original API request. api_response: Spotify API response. """ try: spotify_search_album_response = ( SpotifySearchAlbumsResponse(**api_response) ) with producer: for album in spotify_search_album_response.albums.items: album_message = SpotifySearchAlbumMessage( api_request=api_request, album=album) producer.produce( 'event.spotifyapi.album', f'product-{api_request.payload.upc}', album_message.model_dump_json(), auto_flush=False) except ValidationError: logger.warning( f'Unexpected API response for the request {api_request}\n' f'Response: {api_response}' ) def submit_spotify_api_product_artists_search_response( api_request: SpotifySearchAlbumByUPCRequest, api_response: Dict): """Submit Spotify API response to Kafka. Args: api_request: Original API request. api_response: Spotify API response. """ try: with producer: for album_item in api_response.get('albums', {}).get('items', []): for artist_data in album_item.get('artists', []): artist = SpotifyArtist.model_validate(artist_data) artist_message = SpotifySearchAlbumArtistMessage( api_request=api_request, artist=artist) producer.produce( 'event.spotifyapi.album.artist', f'product-{api_request.payload.product_id}', artist_message.model_dump_json(), auto_flush=False) except ValidationError: logger.warning( f'Unexpected API response for the request {api_request}\n' f'Response: {api_response}' ) def submit_spotify_api_track_artists_search_response( api_request: SpotifySearchTrackByISRCRequest, api_response: Dict): """Submit Spotify API response to Kafka. Args: api_request: Original API request. api_response: Spotify API response. """ try: spotify_search_track_response = ( SpotifySearchTracksResponse(**api_response) ) try: track_ids = [t.id for t in spotify_search_track_response.tracks.items] except Exception: track_ids = [] logger.info( 'Spotify search returned %d track(s) for isrc=%s: %s', len(track_ids), api_request.payload.isrc, track_ids) with producer: for track in spotify_search_track_response.tracks.items: track_message = SpotifySearchTrackMessage( api_request=api_request, track=track) producer.produce( 'event.spotifyapi.track', f'isrc-{api_request.payload.isrc}', track_message.model_dump_json(), auto_flush=False) # Produce artist messages for artist in track.artists: artist_message = SpotifySearchTrackArtistMessage( api_request=api_request, artist=artist) producer.produce( 'event.spotifyapi.track.artist', f'isrc-{api_request.payload.isrc}', artist_message.model_dump_json(), auto_flush=False) except ValidationError: logger.warning( f'Unexpected API response for the request {api_request}\n' f'Response: {api_response}' ) def submit_spotify_api_get_audio_features_response( api_request: SpotifyGetAudioFeaturesRequest, api_response: Dict): """Submit Spotify API response to Kafka. Args: api_request: Original API request. api_response: Spotify API response. """ items = zip( api_request.payload.spotify_track_ids, api_response['audio_features']) with producer: for spotify_track_id, audio_features in items: try: if audio_features: features = SpotifyAudioFeature(**audio_features) producer.produce( 'event.spotifyapi.audioFeatures', features.id, features.model_dump_json(), auto_flush=False) else: producer.produce( 'event.spotifyapi.audioFeatures.missing', spotify_track_id, None, auto_flush=False) except (ValidationError, TypeError) as e: logger.debug(f'Unexpected /audio-features response. ' f'Error: {e} ' f'Record: {audio_features} ' f'Request: {api_request} ' f'Response: {api_response}')