using MoreLinq; using MySql.Data.MySqlClient; using NLog; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Core.Factory; using Sony.Filtr.Database; using Sony.Filtr.Utility.Extensions; using System; using System.Collections.Generic; using System.Globalization; using System.IO; using System.Linq; using System.Threading.Tasks; namespace Sony.Filtr.Tasks.Tasks.SonyReleases { public class SonyUpcReleasesImport { private readonly StorageFactory _storageFactory; private readonly SonyReleaseParser _sonyReleaseParser; private readonly Logger _logger; private const string _FieldTerminator = "|"; public SonyUpcReleasesImport(StorageFactory storageFactory, SonyReleaseParser sonyReleaseParser) { _storageFactory = storageFactory; _sonyReleaseParser = sonyReleaseParser; _logger = LogManager.GetLogger("ImportNewSonyReleasesTask"); } public async Task ExecuteAsync() { var releaseImportSetting = new ReleaseImportSettings() { Name = "Sony", IsrcBucket = "filtr-new-releases", UpcBucket = "filtr-new-releases-upc", IncomingFolder = "incoming/", ProcessedFolder = "processed/", SpotifyAnalyticsAccount = SpotifyAnalyticsAccount.Sony, }; await ImportAlbumUpcsUpdatesAsync(releaseImportSetting); } private async Task ImportAlbumUpcsUpdatesAsync(ReleaseImportSettings releaseImportSetting) { var files = await GetUpcFilesToProcessAsync(releaseImportSetting); var parsedFiles = files.Select(f => new { Filepath = f, Date = ParseDateFromUpcFilename(f), }).ToList(); var filesWithUnknownFilename = parsedFiles.Where(f => !f.Date.HasValue); foreach (var unknownFile in filesWithUnknownFilename) { _logger.Warn($"Could not parse date from file: {unknownFile.Filepath}"); } foreach (var file in parsedFiles.Where(f => f.Date.HasValue && !string.IsNullOrWhiteSpace(f.Filepath)).OrderBy(f => f.Date)) { var filepath = file.Filepath; var fileContent = _storageFactory.GetObject(filepath, releaseImportSetting.UpcBucket); var fileAlbums = await _sonyReleaseParser.ParseSonyAlbumUpdateFileAsync(filepath, fileContent); await SaveUpcAlbumUpdatesAsync(releaseImportSetting, fileAlbums); await MarkUpcFileAsProcessedAsync(releaseImportSetting, filepath); } } private async Task MarkUpcFileAsProcessedAsync(ReleaseImportSettings releaseImportSetting, string filePath) { var destinationPath = Flurl.Url.Combine(releaseImportSetting.ProcessedFolder, filePath.Split('/').Last()); await _storageFactory.CopyObjectAsync(releaseImportSetting.UpcBucket, filePath, destinationPath); await _storageFactory.DeleteObjectAsync(releaseImportSetting.UpcBucket, filePath); } private DateTime? ParseDateFromUpcFilename(string filePath) { //Filename looks like: Filtr_Admin_New_Releases_2017-03-23-16-31-39.xlsx var filename = Path.GetFileNameWithoutExtension(filePath); if (filename != null) { filename = filename.Replace("Filtr_Admin_New_Releases_", string.Empty); filename = filename.Replace("-", " "); DateTime date; if (DateTime.TryParseExact(filename, "yyyy MM dd HH mm ss", CultureInfo.InvariantCulture, DateTimeStyles.None, out date)) { return date; } } return null; } //private async Task> GetFullSonyAlbumUpcsToImportAsync() //{ // ConcurrentBag albums = new ConcurrentBag(); // var deezerPart1 = await _sonyReleaseParser.ParseSonyAlbumsAsync("sonyupc/DeezerFullProductCatalog_Part1_20170314.xlsx", MusicService.Deezer); // deezerPart1.ForEach(albums.Add); // var deezerPart2 = await _sonyReleaseParser.ParseSonyAlbumsAsync("sonyupc/DeezerFullProductCatalog_Part2_20170314.xlsx", MusicService.Deezer); // deezerPart2.ForEach(albums.Add); // var iTunes = await _sonyReleaseParser.ParseSonyAlbumsAsync("sonyupc/iTunesFullProductCatalog_20170314.xlsx", MusicService.AppleMusic); // iTunes.ForEach(albums.Add); // var sonyPart1 = await _sonyReleaseParser.ParseSonyAlbumsAsync("sonyupc/SpotifyFullProductCatalog_Part1_20170314.xlsx", MusicService.Spotify); // sonyPart1.ForEach(albums.Add); // var sonyPart2 = await _sonyReleaseParser.ParseSonyAlbumsAsync("sonyupc/SpotifyFullProductCatalog_Part2_20170314.xlsx", MusicService.Spotify); // sonyPart2.ForEach(albums.Add); // return albums.ToList(); //} //private async Task SaveSonyUpcAlbumsAsync(List sonyAlbums) //{ // var uniqueLabels = sonyAlbums.Select(a => a.Label).Distinct().ToList(); // _logger.Debug($"Found {uniqueLabels.Count} Labels to map"); // var labelMapping = await CreateLabelMappingAsync(uniqueLabels); // _logger.Debug($"Got label mapping for {labelMapping.Count} labels"); // await SaveUpcRegionDataAsync(sonyAlbums); // await SaveUpcLabelAsync(sonyAlbums, labelMapping); //} private async Task SaveUpcAlbumUpdatesAsync(ReleaseImportSettings releaseImportSetting, List sonyAlbums) { var uniqueLabels = sonyAlbums.Select(a => a.Label).Where(l=> !string.IsNullOrWhiteSpace(l)).Distinct().ToList(); _logger.Debug($"Found {uniqueLabels.Count} Labels to map"); var labelMapping = await CreateLabelMappingAsync(uniqueLabels); _logger.Debug($"Got label mapping for {labelMapping.Count} labels"); await SaveUpcLabelAsync(sonyAlbums, labelMapping); foreach (var musicServiceGroup in sonyAlbums.GroupBy(p => p.MusicService)) { var musicService = musicServiceGroup.Key; var ucpGroups = musicServiceGroup.GroupBy(a => a.UPC); await ucpGroups.ForEachAsync(10, async upcGroup => { var upc = upcGroup.Key; _logger.Debug($"Updating for UPC {upc}"); await FaultHandlingPolicy.MySqlRetryPolicyAsync.ExecuteAsync(async () => { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { using (var tran = await conn.BeginTransactionAsync()) { const string deleteSql = "DELETE FROM tblSonyUPCRegion WHERE Upc=@upc AND MusicServiceId=@musicServiceId AND ReleaseTypeId=@releaseTypeId"; using (var deleteCmd = new MySqlCommand(deleteSql, conn, tran)) { deleteCmd.Parameters.AddWithValue("upc", upc); deleteCmd.Parameters.AddWithValue("musicServiceId", (int)musicService); deleteCmd.Parameters.AddWithValue("releaseTypeId", (int)releaseImportSetting.SpotifyAnalyticsAccount); await deleteCmd.ExecuteNonQueryAsync(); } var revoked = upcGroup.Any(u => u.Revoked); if (!revoked) { var values = upcGroup.SelectMany(u => u.Regions).DistinctBy(u => u.Region).Select(r => "(" + string.Join(",", "'" + MySqlHelper.EscapeString(upc) + "'", "'" + MySqlHelper.EscapeString(r.Region) + "'", (int)musicService, (int)r.ReleaseType) + ")"); var insertSql = $"INSERT INTO tblSonyUPCRegion (Upc, Region, MusicServiceId, ReleaseTypeId) VALUES {string.Join(",", values)}"; using (var insertCmd = new MySqlCommand(insertSql, conn, tran)) { await insertCmd.ExecuteNonQueryAsync(); } } tran.Commit(); } } }); }); } } private async Task SaveUpcRegionDataAsync(List sonyAlbums) { if (!sonyAlbums.Any()) return; _logger.Debug($"Writing UPC Region data to file"); var filename = await SaveAlbumToFileAsync(sonyAlbums); if (filename != null) { _logger.Debug($"Importing UPC Region data to database"); await BulkAddAlbumDataAsync(filename); File.Delete(filename); } } private async Task SaveAlbumToFileAsync(List albums) { const string fileDirectory = "ImportSonyReleases/"; if (!Directory.Exists(fileDirectory)) { Directory.CreateDirectory(fileDirectory); } var filename = $"albums-" + DateTime.Now.ToString("yyyy-MM-dd-HH-mm", CultureInfo.InvariantCulture) + ".csv"; var filePath = Path.Combine(fileDirectory, filename); var albumsRegions = albums.SelectMany(r => r.Regions, (release, s) => new { upc = release.UPC, RegionData = s }).Select(release => string.Join(_FieldTerminator, release.upc, release.RegionData.Region, (int)release.RegionData.MusicService, (int)release.RegionData.ReleaseType)); File.WriteAllLines(filePath, albumsRegions); return filePath; } private async Task BulkAddAlbumDataAsync(string filePath) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var bulkLoader = new MySqlBulkLoader(conn); bulkLoader.Local = true; bulkLoader.FileName = filePath; bulkLoader.TableName = "tblSonyUPCRegion"; bulkLoader.FieldTerminator = _FieldTerminator; bulkLoader.LineTerminator = Environment.NewLine; bulkLoader.ConflictOption = MySqlBulkLoaderConflictOption.Ignore; await bulkLoader.LoadAsync(); } } private async Task BulkAddAlbumUPCLabelDataAsync(string filePath) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var bulkLoader = new MySqlBulkLoader(conn); bulkLoader.Local = true; bulkLoader.FileName = filePath; bulkLoader.TableName = "tblSonyUPC"; bulkLoader.FieldTerminator = _FieldTerminator; bulkLoader.LineTerminator = Environment.NewLine; bulkLoader.ConflictOption = MySqlBulkLoaderConflictOption.Replace; await bulkLoader.LoadAsync(); } } private async Task SaveUpcLabelAsync(List sonyAlbums, Dictionary labelMapping) { if (!sonyAlbums.Any()) return; _logger.Debug($"Writing UPC label data to file"); var filepath = await SaveUpcLabelFileAsync(sonyAlbums, labelMapping); if (filepath != null) { _logger.Debug($"Importing UPC label data to database"); await BulkAddAlbumUPCLabelDataAsync(filepath); File.Delete(filepath); } } private async Task SaveUpcLabelFileAsync(List sonyAlbums, Dictionary labelMapping) { const string fileDirectory = "ImportSonyReleases/"; if (!Directory.Exists(fileDirectory)) { Directory.CreateDirectory(fileDirectory); } var filename = $"upc-albums-" + DateTime.Now.ToString("yyyy-MM-dd-HH-mm", CultureInfo.InvariantCulture) + ".csv"; var filePath = Path.Combine(fileDirectory, filename); var albumsRegions = sonyAlbums .GroupBy(a => a.UPC) .Select(a => new { UPC = a.Key, LabelId = GetLabelId(labelMapping, a) }) .Where(a => a.LabelId.HasValue) .Select(album => string.Join(_FieldTerminator, album.UPC, album.LabelId)); File.WriteAllLines(filePath, albumsRegions); return filePath; } private int? GetLabelId(Dictionary labelMapping, IGrouping sonyUpcAlbums) { var labelName = sonyUpcAlbums.Select(p=> p.Label).FirstOrDefault(l => !string.IsNullOrWhiteSpace(l)); if (string.IsNullOrWhiteSpace(labelName)) { return null; } return labelMapping.GetValueOrDefault(labelName); } private async Task> CreateLabelMappingAsync(IEnumerable uniqueLabels) { Dictionary mapping = new Dictionary(); using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var insertSql = $"INSERT IGNORE INTO tblSonyUPCLabel (LabelName) VALUES {string.Join(",", uniqueLabels.Select(l => "('" + MySqlHelper.EscapeString(l) + "')"))} "; var insertCmd = new MySqlCommand(insertSql, conn); await insertCmd.ExecuteNonQueryAsync(); var getSql = $"SELECT Id, LabelName FROM tblSonyUPCLabel"; var getCmd = new MySqlCommand(getSql, conn); var reader = await getCmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { var id = reader.GetInt32("Id"); var labelName = reader.GetString("LabelName"); mapping.Add(labelName, id); } reader.Close(); } return mapping; } private async Task> GetUpcFilesToProcessAsync(ReleaseImportSettings releaseImportSetting) { var directory = releaseImportSetting.IncomingFolder; var files = await _storageFactory.GetObjectsAsync(releaseImportSetting.UpcBucket, directory); return files; } } }