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.Spotify.Artists; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.Spotify { public class ExportSpotifyArtistFollowersTask : IScheduledTask { private readonly SpotifyPlaylistHistoricTrackListManager _spotifyPlaylistHistoricTrackListManager; private readonly S3BucketReference _destinationBucket; private readonly Logger _logger; private const string _FilenameBase = "ExportSpotifyArtistFollowers-"; public ExportSpotifyArtistFollowersTask(SpotifyPlaylistHistoricTrackListManager spotifyPlaylistHistoricTrackListManager, S3BucketReference destinationBucket) { _spotifyPlaylistHistoricTrackListManager = spotifyPlaylistHistoricTrackListManager; _destinationBucket = destinationBucket; _logger = LogManager.GetLogger("ExportSpotifyArtistFollowersTask"); } public async Task ExecuteAsync(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; var spotifyArtistFolowersExport = await GetArtistFollowers(date.Date); if (spotifyArtistFolowersExport.Count() == 0) { continue; } var filename = await WriteToDiskAsync(DateTime.Now.Date, spotifyArtistFolowersExport); await UploadFileAsync(filename); File.Delete(filename); _logger.Debug($"Done processing"); 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 = "SpotifyArtistFollowers/" + filename, CannedACL = S3CannedACL.BucketOwnerFullControl, }); _logger.Debug("Done uploading to S3"); } private async Task WriteToDiskAsync(DateTime date, List spotifyArtistFollowersList) { 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", "artist_id", "followers_count", "as_of")); foreach (var spotifyArtistFollowers in spotifyArtistFollowersList) { var artist_id = spotifyArtistFollowers.artist_id; var followers_count = spotifyArtistFollowers.followers_count; var as_of = spotifyArtistFollowers.as_of; var row = string.Join("\t", artist_id, followers_count, as_of); await fileWriter.WriteLineAsync(row); } } } var compressedFilename = CompressFileAsync(filePath); File.Delete(filePath); _logger.Debug("Done writing to disk"); return compressedFilename; } private string GetOrCreateTempDirectory() { var dirPath = "ArtistFollowersExport"; 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> GetArtistFollowers(DateTime date) { List spotifyArtistsData = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT sa.ArtistId as artist_id, saf.Followers as followers_count, saf.Date as as_of FROM tblSpotifyArtist as sa INNER JOIN tblSpotifyArtistFollowerLog as saf on sa.ArtistId = saf.ArtistId where saf.Date = @date"; cmd.Parameters.AddWithValue("@date", date); var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { spotifyArtistsData.Add(new SpotifyArtistFollowers(reader)); } } return spotifyArtistsData; } private async Task> GetAlreadyProcessedDatesAsync() { List dates = new List(); using (var connection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT date FROM tblSpotifyArtistFollowersExportHistory"; var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { var date = reader.GetDateTime(0); dates.Add(date); } } return dates; } private async Task SetDateAsExportedAsync(DateTime date) { using (var connection = await DatabaseHandler.GetOpenConnectionAsync()) { var cmd = connection.CreateCommand(); cmd.CommandText = "INSERT INTO tblSpotifyArtistFollowersExportHistory (Date) VALUES (@date)"; cmd.Parameters.AddWithValue("@date", date); await cmd.ExecuteNonQueryAsync(); } } } }