import hashlib from datetime import datetime from app.dsp.enums import DSPId from app.dsp.types import ( FollowedArtistsResult, PlaylistsResult, ProfileResult, RecentlyPlayedResult, SavedAlbumsResult, SavedTracksResult, TopArtistsResult, TopTracksResult, ) from app.fandata.models import ( Fan, FanFollowedArtist, FanPlaylist, FanRecentlyPlayed, FanSavedAlbum, FanSavedTrack, FanTopArtist, FanTopTrack, ) from app.fandata.sinks.base import DataSink from app.fandata.types import FanRecord class PostgresSink(DataSink): def write_fans( self, *, fan: FanRecord, dsp_id: DSPId, profile: ProfileResult, collected_at: datetime, ) -> None: email = profile.email or "" fan_id = hashlib.sha256(email.encode()).hexdigest() Fan.query.upsert( [ { "dsp_user_id": fan.dsp_user_id, "dsp_id": dsp_id, "email": email, "fan_id": fan_id, "country": profile.country, "product": profile.product, "display_name": profile.display_name, "collected_at": collected_at, } ] ) def write_fan_top_artists( self, *, fan: FanRecord, dsp_id: DSPId, result: TopArtistsResult, collected_at: datetime, ) -> None: if not result.items: return FanTopArtist.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "artist_id": item.id, "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_top_tracks( self, *, fan: FanRecord, dsp_id: DSPId, result: TopTracksResult, collected_at: datetime, ) -> None: if not result.items: return FanTopTrack.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.id, "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_recently_played( self, *, fan: FanRecord, dsp_id: DSPId, result: RecentlyPlayedResult, collected_at: datetime, ) -> None: if not result.items: return FanRecentlyPlayed.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.track_id, "played_at": datetime.fromisoformat(item.played_at), "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_playlists( self, *, fan: FanRecord, dsp_id: DSPId, result: PlaylistsResult, collected_at: datetime, ) -> None: if not result.items: return FanPlaylist.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "playlist_id": item.id, "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_saved_albums( self, *, fan: FanRecord, dsp_id: DSPId, result: SavedAlbumsResult, collected_at: datetime, ) -> None: if not result.items: return FanSavedAlbum.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "album_id": item.album_id, "added_at": datetime.fromisoformat(item.added_at), "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_saved_tracks( self, *, fan: FanRecord, dsp_id: DSPId, result: SavedTracksResult, collected_at: datetime, ) -> None: if not result.items: return FanSavedTrack.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.track_id, "added_at": datetime.fromisoformat(item.added_at), "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def write_fan_followed_artists( self, *, fan: FanRecord, dsp_id: DSPId, result: FollowedArtistsResult, collected_at: datetime, ) -> None: if not result.items: return FanFollowedArtist.query.upsert( [ { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "artist_id": item.id, "data": item.data, "collected_at": collected_at, } for item in result.items ] ) def flush(self) -> None: pass