using System; using System.Collections.Generic; using System.Globalization; using System.IO; using System.Linq; using System.Threading.Tasks; using MySql.Data.MySqlClient; using NLog; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Core.Factory; using Sony.Filtr.Database; namespace Sony.Filtr.Tasks.Tasks.SonyReleases { public class SonyIsrcReleaseImport { private readonly StorageFactory _storageFactory; private readonly SonyReleaseParser _sonyReleaseParser; private Logger _logger; private const string _FieldTerminator = "|"; public SonyIsrcReleaseImport(StorageFactory storageFactory, SonyReleaseParser sonyReleaseParser) { _storageFactory = storageFactory; _sonyReleaseParser = sonyReleaseParser; _logger = LogManager.GetLogger("ImportNewSonyReleasesTask"); } public async Task ExecuteAsync() { var settings = new ReleaseImportSettings() { Name = "Sony", IsrcBucket = "filtr-new-releases", UpcBucket = "filtr-new-releases-upc", IncomingFolder = "incoming/", ProcessedFolder = "processed/", SpotifyAnalyticsAccount = SpotifyAnalyticsAccount.Sony, }; _logger.Debug($"Begin ImportISRCReleaseFilesAsync for Sony"); await ImportIsrcReleaseFilesAsync(settings); _logger.Debug($"Done ImportISRCReleaseFilesAsync for Sony"); } private async Task ImportIsrcReleaseFilesAsync(ReleaseImportSettings releaseImportSetting) { var s3Paths = await GetIsrcFilesToProcessAsync(releaseImportSetting); if (s3Paths.Any()) { var releases = await GetIsrcReleasesToImportAsync(releaseImportSetting, s3Paths); if (releases.Any()) { var filepath = await SaveIsrcReleasesToFileAsync(releaseImportSetting, releases); await BulkAddIsrcReleaseDataAsync(filepath); File.Delete(filepath); await MarkIsrcFilesAsProcessedAsync(releaseImportSetting, s3Paths); } } } private async Task> GetIsrcFilesToProcessAsync(ReleaseImportSettings releaseImportSetting) { var directory = releaseImportSetting.IncomingFolder; var files = await _storageFactory.GetObjectsAsync(releaseImportSetting.IsrcBucket, directory); return files; } private async Task> GetIsrcReleasesToImportAsync(ReleaseImportSettings releaseImportSetting, List files) { List releases = new List(); foreach (var file in files) { var fileContent = await _storageFactory.GetObjectAsync(file, releaseImportSetting.IsrcBucket); releases.AddRange(_sonyReleaseParser.ParseIsrcReleaseFile(fileContent)); } return releases; } private async Task SaveIsrcReleasesToFileAsync(ReleaseImportSettings releaseImportSetting, List releases) { const string fileDirectory = "/ImportSonyReleases/"; if (!Directory.Exists(fileDirectory)) { Directory.CreateDirectory(fileDirectory); } var filePath = fileDirectory + releaseImportSetting.Name + "-" + DateTime.Now.ToString("yyyy-MM-dd-HH-mm", CultureInfo.InvariantCulture) + ".csv"; using (var saveFile = File.Open(filePath, FileMode.OpenOrCreate)) { using (var fileWriter = new StreamWriter(saveFile)) { foreach (var release in releases.SelectMany(r => r.Regions, (release, s) => new { ISRC = release.Isrc, Region = s })) { var csvLine = string.Join(_FieldTerminator, release.ISRC, release.Region); await fileWriter.WriteLineAsync(csvLine); } } } return filePath; } private async Task MarkIsrcFilesAsProcessedAsync(ReleaseImportSettings releaseImportSetting, List filepaths) { foreach (var filePath in filepaths) { var destinationPath = Flurl.Url.Combine(releaseImportSetting.ProcessedFolder, filePath.Split('/').Last()); await _storageFactory.CopyObjectAsync(releaseImportSetting.IsrcBucket, filePath, destinationPath); await _storageFactory.DeleteObjectAsync(releaseImportSetting.IsrcBucket, filePath); } } private async Task BulkAddIsrcReleaseDataAsync(string filePath) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var bulkLoader = new MySqlBulkLoader(conn); bulkLoader.Local = true; bulkLoader.FileName = filePath; bulkLoader.TableName = "tblSonyRelease"; bulkLoader.FieldTerminator = _FieldTerminator; bulkLoader.LineTerminator = Environment.NewLine; bulkLoader.ConflictOption = MySqlBulkLoaderConflictOption.Replace; await bulkLoader.LoadAsync(); } } } }