using System; using System.Collections.Concurrent; using System.Collections.Generic; 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.Model; using Amazon.S3.Transfer; using MoreLinq; using NLog; using Sony.Filtr.AppleMusic.Data; using Sony.Filtr.AppleMusic.Playlists; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Database; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.AppleMusic { public class ExportAppleMusicPlaylistTracklistHistoryTask : IScheduledTask { private readonly AppleMusicPlaylistManager _appleMusicPlaylistManager; private readonly S3BucketReference _destinationBucket; private readonly Logger _logger; private readonly AmazonS3Client _s3Client; private const string _FilenameBase = "AppleMusicPlaylistTracklistHistory-"; public ExportAppleMusicPlaylistTracklistHistoryTask(AppleMusicPlaylistManager appleMusicPlaylistManager, S3BucketReference destinationBucket) { _appleMusicPlaylistManager = appleMusicPlaylistManager; _destinationBucket = destinationBucket; _s3Client = new AmazonS3Client(_destinationBucket.AccessKey, _destinationBucket.SecretKey, RegionEndpoint.USEast1); _logger = LogManager.GetLogger("ExportAppleMusicPlaylistTracklistHistoryTask"); } 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(2017, 06, 2); var dateRange = startDate.GetDateRangeTo(DateTime.Today.AddDays(-1)).Reverse().ToList(); var alreadyProcessedDates = await GetAlreadyProcessedDatesAsync(); dateRange = dateRange.ExceptBy(alreadyProcessedDates, t => t.Date.Date).ToList(); dateRange.Insert(0, DateTime.Now); foreach (var date in dateRange) { _logger.Debug($"Begin processing {date.ToShortDateString()}"); List playlistTracks; if (date.Date != DateTime.Now.Date) { playlistTracks = await GetApplePlaylistTracklistAsync(date); } else { playlistTracks = await GetApplePlaylistTracklistForTodayAsync(date); } if (playlistTracks.Count() == 0) { continue; } var filename = await WriteToDiskAsync(date, playlistTracks); var semaphorePath = CreateSemaphore(filename); await UploadFileAsync(filename); await UploadFileAsync(semaphorePath); File.Delete(filename); File.Delete(semaphorePath); _logger.Debug($"Done processing {date.ToShortDateString()}"); await SetDateAsExportedAsync(date); } return null; } 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.AppleMusic); 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 UploadFileAsync(string filepath) { _logger.Debug("Uploading to S3"); 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/apple_playlists/" + filename, CannedACL = S3CannedACL.BucketOwnerFullControl, }); _logger.Debug("Done uploading to S3"); } private async Task WriteToDiskAsync(DateTime date, List playlistTracks) { var datestring = date.ToString("yyyy-MM-ddTHH.mm.ss"); 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", "PlaylistId", "Date", "PlaylistPosition", "TrackId", "TrackName", "TrackArtistsName", "TrackISRC", "PlaylistName")); foreach (var playlistTrackData in playlistTracks) { var row = string.Join("\t", playlistTrackData.PlaylistId, playlistTrackData.Date, playlistTrackData.PlaylistPosition, playlistTrackData.TrackId, playlistTrackData.TrackName, playlistTrackData.TrackArtistName, playlistTrackData.TrackISRC, playlistTrackData.PlaylistName); 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 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> GetApplePlaylistTracklistAsync(DateTime date) { List history = new List(); using (var connection = DatabaseHandler.GetOpenReadOnlyConnection()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT p.Id AS PlaylistId, pt.Date AS Date, pt.Position AS PlaylistPosition, s.Id AS TrackId, REPLACE(REPLACE(REPLACE(s.Name, '\r', ' '), '\n', ' '), '\t', ' ') AS TrackName, " + "REPLACE(REPLACE(REPLACE(s.ArtistName, '\r', ' '), '\n', ' '), '\t', ' ') AS TrackArtistName, s.ISRC AS TrackISRC, REPLACE(REPLACE(REPLACE(p.Name, '\r', ' '), '\n', ' '), '\t', ' ') AS PlaylistName FROM dbSony_dbo.tblAppleMusicPlaylist p " + "INNER JOIN dbSony_dbo.tblAppleMusicPlaylistTracklistHistory pt ON p.Id = pt.PlaylistId " + "INNER JOIN dbSony_dbo.tblAppleMusicSong s ON pt.SongId = s.Id " + "WHERE pt.Storefront = 'us' and s.Storefront = 'us' AND pt.Date = @date"; cmd.Parameters.AddWithValue("@date", date); cmd.CommandTimeout = 520; var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { var historyEntity = new AppleMusicPlaylistTracklistHistory(reader); history.Add(historyEntity); } } return history; } private async Task> GetApplePlaylistTracklistForTodayAsync(DateTime date) { List history = new List(); using (var connection = DatabaseHandler.GetOpenReadOnlyConnection()) { var cmd = connection.CreateCommand(); cmd.CommandText = @"SELECT pt.PlaylistId AS PlaylistId, NOW() AS Date, pt.Position AS PlaylistPosition, s.Id AS TrackId, REPLACE(REPLACE(REPLACE(s.Name, '\r', ' '), '\n', ' '), '\t', ' ') AS TrackName, REPLACE(REPLACE(REPLACE(GROUP_CONCAT(distinct ama.Name), '\r', ' '), '\n', ' '), '\t', ' ') AS TrackArtistName, s.ISRC AS TrackISRC, REPLACE(REPLACE(REPLACE(p.Name, '\r', ' '), '\n', ' '), '\t', ' ') AS PlaylistName FROM dbSony_dbo.tblAppleMusicPlaylist p INNER JOIN dbSony_dbo.tblAppleMusicPlaylistTracklist pt ON p.Id = pt.PlaylistId INNER JOIN dbSony_dbo.tblAppleMusicSong s ON pt.SongId = s.Id INNER JOIN tblAppleMusicSongArtist as amsa on s.Id = amsa.SongId INNER JOIN tblAppleMusicArtist as ama on ama.Id = amsa.ArtistId WHERE pt.Storefront = 'us' and s.Storefront = 'us' and amsa.Storefront = 'us' and ama.Storefront='us' group by PlaylistId, trackISRC"; cmd.CommandTimeout = 5500; var reader = await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteReaderAsync()); while (await reader.ReadAsync()) { var historyEntity = new AppleMusicPlaylistTracklistHistory(reader); history.Add(historyEntity); } } return history; } 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.AppleMusic); await FaultHandlingPolicy.MySqlRetryPolicy.Execute(async () => await cmd.ExecuteNonQueryAsync()); } } } }