using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Data.Common; using System.Globalization; using System.IO; using System.IO.Compression; using System.Linq; using System.Text; using System.Threading.Tasks; using Amazon; using Amazon.S3; using Amazon.S3.Transfer; using MoreLinq; using NLog; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Database; using Sony.Filtr.Playlists.Spotify; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.Spotify { public class ExportSpotifyPlaylistTracklistHistoryTask : IScheduledTask { private readonly SpotifyPlaylistHistoricTrackListManager _spotifyPlaylistHistoricTrackListManager; private readonly S3BucketReference _destinationBucket; private readonly Logger _logger; private const string _FilenameBase = "SpotifyPlaylistTracklistHistory-"; private const string _FilenamePersonalized = "PersonalizedPlaylists-"; public ExportSpotifyPlaylistTracklistHistoryTask(SpotifyPlaylistHistoricTrackListManager spotifyPlaylistHistoricTrackListManager, S3BucketReference destinationBucket) { _spotifyPlaylistHistoricTrackListManager = spotifyPlaylistHistoricTrackListManager; _destinationBucket = destinationBucket; _logger = LogManager.GetLogger("ExportSpotifyPlaylistTracklistHistoryTask"); } public async Task ExecuteAsync(Guid scheduledTaskLogId) { await Retry.Do(async () => await ProcessTask(scheduledTaskLogId), TimeSpan.FromSeconds(60), _logger); return null; } public async Task ProcessTask(Guid scheduledTaskLogId) { DateTime startDate = new DateTime(2016, 5, 24); var endDate = DateTime.Today.AddDays(-1); //Ensure we are done processing everything for today. var dateRange = startDate.GetDateRangeTo(endDate).Reverse().ToList(); var alreadyProcessed = await GetAlreadyProcessedDatesAsync(); dateRange = dateRange.ExceptBy(alreadyProcessed, d => d.Date.Date).ToList(); dateRange.Insert(0, DateTime.Now); foreach (var date in dateRange) { _logger.Debug($"Begin processing {date.ToShortDateString()}"); var isToday = DateTime.Today == date.Date; List playlistTracks = null; List playlistTracksToday = null; if (isToday) { playlistTracksToday = await GetPlaylistHistoryDataTodayAsync(); } else { playlistTracks = await GetPlaylistHistoryDataAsync(date); } if (playlistTracks == null && playlistTracksToday != null && playlistTracksToday.Count() == 0 || playlistTracksToday == null && playlistTracks != null && playlistTracks.Count() == 0) { continue; } var personalizedPlaylists = await GetPersonalizedPlaylists(); var filename = await WriteToDiskAsync(date, playlistTracks, playlistTracksToday, isToday); var semaphorePath = CreateSemaphore(filename); var personalizedPlaylistsPath = await CreatePersonalizedPlaylistsFile(personalizedPlaylists, filename, isToday); await UploadFileAsync(filename); await UploadFileAsync(semaphorePath); if (!string.IsNullOrEmpty(personalizedPlaylistsPath)) { await UploadFileAsync(personalizedPlaylistsPath); File.Delete(personalizedPlaylistsPath); } File.Delete(filename); File.Delete(semaphorePath); _logger.Debug($"Done processing {date.ToShortDateString()}"); await SetDateAsExportedAsync(date); } return null; } private async Task UploadFileAsync(string filepath) { _logger.Debug("Uploading to S3"); var s3Client = new AmazonS3Client(_destinationBucket.AccessKey, _destinationBucket.SecretKey, RegionEndpoint.USEast1); var transferUtility = new TransferUtility(s3Client); var filename = Path.GetFileName(filepath); await transferUtility.UploadAsync(new TransferUtilityUploadRequest() { BucketName = _destinationBucket.Bucket, FilePath = filepath, Key = _destinationBucket.BucketPath + "/delphi_feed/spotify_playlists/" + filename, CannedACL = S3CannedACL.BucketOwnerFullControl, }); _logger.Debug("Done uploading to S3"); } private async Task WriteToDiskAsync(DateTime date, List tracklistHistory, List tracklistHistoryToday, bool isToday = false) { var datestring = date.ToString("yyyy-MM-ddTHH.mm.ss"); var dateStringInFile = date.ToString("yyyy-MM-dd"); var path = GetOrCreateTempDirectory(); var filename = $"{_FilenameBase}{datestring}.csv"; var filePath = Path.Combine(path, filename); FileInfo file = new FileInfo(filePath); _logger.Debug("Writing to disk"); using (var fileStream = file.OpenWrite()) { fileStream.Write(Encoding.UTF8.GetPreamble(), 0, Encoding.UTF8.GetPreamble().Length); using (var fileWriter = new StreamWriter(fileStream)) { fileWriter.WriteLine(string.Join("\t", "PlaylistUri", "Date", "TrackId", "PlaylistIndex", "EarliestAdded", "ISRC", "TrackName", "ArtistsName")); if (!isToday) { foreach (var playlistTrackData in tracklistHistory) { if (playlistTrackData == null || playlistTrackData.PlaylistId == null || playlistTrackData.Tracks == null) { continue; } var playlistId = playlistTrackData.PlaylistId != null ? playlistTrackData.PlaylistId.Trim() : "null"; var playlistUri = $"spotify:playlist:{playlistId}"; foreach (var trackOrder in playlistTrackData.Tracks.ItemIndex()) { if (trackOrder == null || trackOrder.Item == null) { continue; } var row = string.Join("\t", playlistUri, dateStringInFile, trackOrder.Item.TrackId, trackOrder.Index, String.Empty, String.Empty, String.Empty, String.Empty); await fileWriter.WriteLineAsync(row); } } } else { foreach (var trackDataToday in tracklistHistoryToday) { if (trackDataToday == null || trackDataToday.PlaylistId == null || trackDataToday.Track == null) { continue; } var playlistId = trackDataToday.PlaylistId; var playlistUri = $"spotify:playlist:{playlistId}"; var row = string.Join("\t", playlistUri, dateStringInFile, trackDataToday.Track.TrackId, trackDataToday.Track.PlaylistIndex, trackDataToday.Track.EarliestDate, trackDataToday.Track.ISRC, trackDataToday.Track.TrackName, trackDataToday.Track.ArtistsName); await fileWriter.WriteLineAsync(row); } } } } var compressedFilename = CompressFileAsync(filePath); File.Delete(filePath); _logger.Debug("Done writing to disk"); return compressedFilename; } private string CreateSemaphore(string filePath) { filePath += "_SUCCESS"; FileInfo file = new FileInfo(filePath); using (StreamWriter sw = new StreamWriter(filePath, true)) { _logger.Debug("Semaphore created"); } return filePath; } private async Task CreatePersonalizedPlaylistsFile(List personalizedPlaylists, string filePath, bool isToday = true) { if (!isToday) { return string.Empty; } var newFilePath = filePath.Replace(_FilenameBase, _FilenamePersonalized); newFilePath = newFilePath.Remove(newFilePath.Length - 3); FileInfo file = new FileInfo(newFilePath); _logger.Debug("Writing PersonalizedPlaylists to disk"); using (var fileStream = file.OpenWrite()) { fileStream.Write(Encoding.UTF8.GetPreamble(), 0, Encoding.UTF8.GetPreamble().Length); using (var fileWriter = new StreamWriter(fileStream)) { fileWriter.WriteLine(string.Join("\t", "PlaylistId")); foreach (var playlistId in personalizedPlaylists) { if (string.IsNullOrEmpty(playlistId)) { continue; } var row = string.Join("\t", playlistId); await fileWriter.WriteLineAsync(row); } } } var compressedFilename = CompressFileAsync(newFilePath); File.Delete(newFilePath); _logger.Debug("Done writing PersonalizedPlaylists to disk"); return compressedFilename; } private string GetOrCreateTempDirectory() { var dirPath = "TrackListHistoryExport"; if (Directory.Exists(dirPath)) { return dirPath; } Directory.CreateDirectory(dirPath); return dirPath; } private string CompressFileAsync(string filepath) { var originalFilename = Path.GetFileName(filepath); var compressedFilename = $"{originalFilename}.gz"; var directory = GetOrCreateTempDirectory(); var compressedFilePath = Path.Combine(directory, compressedFilename); using (FileStream originalFileStream = File.OpenRead(filepath)) { using (FileStream destinationCompressedFile = File.OpenWrite(compressedFilePath)) { using (GZipStream compressionStream = new GZipStream(destinationCompressedFile, CompressionMode.Compress)) { originalFileStream.CopyTo(compressionStream); } } } _logger.Debug($"Done compressing {filepath} to {compressedFilePath}"); return compressedFilePath; } private async Task> GetPlaylistHistoryDataAsync(DateTime date) { var tracklistHistoryList = new List(); var playlistIds = await _spotifyPlaylistHistoricTrackListManager.GetPlaylistWithHistoricTracklistsAsync(date); playlistIds = playlistIds.Where(s => !string.IsNullOrWhiteSpace(s)).ToList(); var playlistIdCount = playlistIds.Count; SpotifyPlaylistHistory tracklistHistory; await playlistIds.ItemIndex().ForEachAsync(2, async playlist => { var playlistId = playlist.Item; Console.WriteLine($"{playlist.Index} / {playlistIdCount}"); tracklistHistory = await GetTracksInPlaylistAsync(playlistId, date); tracklistHistory.PlaylistId = playlistId; tracklistHistoryList.Add(tracklistHistory); }); return tracklistHistoryList; } private async Task> GetPlaylistHistoryDataTodayAsync() { var tracklistHistoryList = new List(); tracklistHistoryList = await GetTracksInPlaylistForTodayAsync(); return tracklistHistoryList; } private async Task GetTracksInPlaylistAsync(string playlistId, DateTime date) { SpotifyPlaylistHistory tracklistHistory = new SpotifyPlaylistHistory(); tracklistHistory.Tracks = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT REPLACE(REPLACE(REPLACE(trackId, '\r', ' '), '\n', ' '), '\t', ' ') AS TrackId " + "FROM tblSpotifyPlaylistTrackListHistory2 " + "WHERE PlaylistId=@playlistId AND Date = @date " + "ORDER BY PlaylistIndex"; cmd.Parameters.AddWithValue("@playlistId", playlistId); cmd.Parameters.AddWithValue("@date", date.Date); var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { var trackId = reader.GetString(0); tracklistHistory.Tracks.Add(new TrackWithEarliestDate { TrackId = trackId }); } } tracklistHistory.PlaylistId = playlistId; return tracklistHistory; } private async Task> GetTracksInPlaylistForTodayAsync() { List tracklistHistoryList = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = @"SELECT REPLACE(REPLACE(REPLACE(sptl.trackId, '\r', ' '), '\n', ' '), '\t', ' ') AS TrackId, sptl.EarliestAdded as EarliestAdded, sptl.PlaylistId as PlaylistId, sptl.PlaylistIndex as PlaylistIndex, st.ISRC as ISRC, st.Name as TrackName, REPLACE(REPLACE(REPLACE(GROUP_CONCAT(distinct sa.Name), '\r', ''), '\n', ''), '\t', '') as ArtistsName FROM tblSpotifyPlaylistTrackList2 AS sptl INNER JOIN (SELECT distinct playlistId from tblSpotifyPlaylistStatistics AS sps WHERE DATE(sps.TracklistLastUpdate) > (NOW() - INTERVAL 180 DAY) ) as sps_id on sptl.PlaylistId = sps_id.PlaylistId INNER JOIN tblSpotifyTrack2 AS st ON sptl.TrackId = st.TrackId INNER JOIN tblSpotifyTrackArtist AS ta ON st.TrackId = ta.TrackId INNER JOIN tblSpotifyArtist AS sa ON ta.ArtistId = sa.ArtistId GROUP BY sptl.PlaylistId, st.ISRC"; var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { SpotifyPlaylistHistoryToday tracklistHistory = new SpotifyPlaylistHistoryToday(); var trackId = reader.GetSafeString("TrackId"); var earliestAdded = reader.GetDateTimeOrDefault("EarliestAdded"); var playlistIndex = reader.GetIntOrDefault("PlaylistIndex"); var isrc = reader.GetSafeString("ISRC"); var trackName = reader.GetSafeString("TrackName"); var artistsName = reader.GetSafeString("ArtistsName"); tracklistHistory.Track = new TrackWithEarliestDate { TrackId = trackId, EarliestDate = earliestAdded, PlaylistIndex = playlistIndex, ArtistsName = artistsName, ISRC = isrc, TrackName = trackName }; tracklistHistory.PlaylistId = reader.GetSafeString("PlaylistId"); tracklistHistoryList.Add(tracklistHistory); } } return tracklistHistoryList; } private async Task> GetAlreadyProcessedDatesAsync() { List dates = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT date FROM tblPlaylistHistoryExport WHERE MusicServiceId = @musicServiceId"; cmd.Parameters.AddWithValue("@musicServiceId", (int)MusicService.Spotify); var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { var date = reader.GetDateTime(0); dates.Add(date); } } return dates; } private async Task> GetPersonalizedPlaylists() { List personalizedPlaylists = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT PlaylistId FROM tblSpotifyPersonalizedPlaylist"; var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { var personalizedPlaylist = reader.GetSafeString("PlaylistId"); personalizedPlaylists.Add(personalizedPlaylist); } } return personalizedPlaylists; } private async Task SetDateAsExportedAsync(DateTime date) { using (var connection = await DatabaseHandler.GetOpenConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "INSERT INTO tblPlaylistHistoryExport (Date, MusicServiceId) VALUES (@date, @MusicServiceId)"; cmd.Parameters.AddWithValue("@date", date); cmd.Parameters.AddWithValue("@musicServiceId", (int)MusicService.Spotify); await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteNonQueryAsync()); } } } }