using MoreLinq; using MySql.Data.MySqlClient; using NLog; using PetaPoco.Business; using Sony.Filtr.ApolloAPI; using Sony.Filtr.ApolloAPI.Models; using Sony.Filtr.Core.Buzz; using Sony.Filtr.Database; using Sony.Filtr.ErrorLogging; using Sony.Filtr.Spotify.Artists; using Sony.Filtr.SpotifyWebAPI; using Sony.Filtr.Utility; using Sony.Filtr.Utility.Extensions; using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Globalization; using System.IO; using System.Linq; using System.Threading; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; using Artist = Sony.Filtr.SpotifyWebAPI.Model.Artist; namespace Sony.Filtr.Tasks.Tasks.Spotify { public class ImportArtistFollowersTask : IScheduledTask { private readonly ErrorLoggingManager _errorLoggingManager; private readonly BuzzArtistManager _buzzArtistManager; protected readonly Logger Logger; protected readonly SpotifyArtistManager SpotifyArtistManager; private readonly IApolloWebApi _vendorApi; private const string _FileDirectory = @"artistfollowers\"; private const string _TempFileName = @"artistfollowersimport.csv"; private const string _FilePath = _FileDirectory + _TempFileName; public ImportArtistFollowersTask(IApolloWebApi vendorApi, ErrorLoggingManager errorLoggingManager, BuzzArtistManager buzzArtistManager, SpotifyArtistManager spotifyArtistManager, string loggerId = nameof(ImportArtistFollowersTask)) { _vendorApi = vendorApi; _errorLoggingManager = errorLoggingManager; _buzzArtistManager = buzzArtistManager; SpotifyArtistManager = spotifyArtistManager; Logger = LogManager.GetLogger(loggerId); } private static int FetchFromSpotifyConcurrencyCount { get { return Maybe.GetAppSettingsIntOrDefault("ImportArtistFollowersTask_FetchFromSpotifyConcurrencyCount", 1); } } public async Task ExecuteAsync(Guid scheduledTaskLogId) { Directory.CreateDirectory(_FileDirectory); //CreateDirectory checks if directory exists internally var artists = new ConcurrentBag(); var notFoundArtistIds = new ConcurrentBag(); Logger.Info("Starting import of artist followers..."); //In order to keep rate limit and time spent importing in mind, full imports are only run when specified var artistIds = await GetArtistIdsAsync(); var date = DateTime.Today; var artistWithData = await GetArtistsWithFollowerDataAsync(date); var artistIdsToFetch = artistIds.Except(artistWithData).ToList(); int currentBatchIndex = 0; var batches = artistIdsToFetch.Batch(SpotifyWebApi.MaxArtistBatchSize).ToList(); Logger.Debug($"Importing follower data for { artistIdsToFetch.Count } artists..."); var getArtistsBlock = new TransformBlock, ArtistsDto>(async batch => { Logger.Debug($"Fetching batch {Interlocked.Increment(ref currentBatchIndex)} of {batches.Count}."); try { var receivedArtists = await _vendorApi.GetAllArtistsAsync(batch); return new ArtistsDto(batch, receivedArtists); } catch (Exception ex) { Logger.Error(ex, "Could not load data from Spotify"); ErrorLoggingManager.Instance.LogError(ex); return null; } }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = FetchFromSpotifyConcurrencyCount }); var createResultingCollectionsBlock = new ActionBlock(dto => { foreach (var artist in dto.ReceivedArtists) { artists.Add(artist); } foreach (var missingArtistId in dto.MissingArtistsIds) { notFoundArtistIds.Add(missingArtistId); } }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1 }); var linkOptions = new DataflowLinkOptions { PropagateCompletion = true }; getArtistsBlock.LinkTo(createResultingCollectionsBlock, linkOptions, dto => dto != null); getArtistsBlock.LinkTo(DataflowBlock.NullTarget(), linkOptions); foreach (var batch in batches) { getArtistsBlock.Post(batch); } getArtistsBlock.Complete(); await createResultingCollectionsBlock.Completion; if (artists.Any()) { try { Logger.Debug("Saving follower data to database..."); await UpdateArtistFollowersAsync(artists); await SaveFollowerDataToCsvAsync(artists, date); await BulkAddFollowerDataAsync(); UpdateBuzzArtistFollowers(artists); File.Delete(_FilePath); } catch (Exception ex) { _errorLoggingManager.LogError(new Exception("ImportArtistFollowers", ex)); Logger.Error(ex, "Exception :( {0}", ex); } if (notFoundArtistIds.Any()) { Logger.Info($"{ notFoundArtistIds.Count } artists not found:\n{string.Join("|", notFoundArtistIds)}"); } Logger.Debug("All done!"); } return null; } protected virtual Task> GetArtistIdsAsync() { return Task.FromResult(SpotifyArtistManager.GetArtists().Select(x => x.ArtistId).ToList()); } private async Task> GetArtistsWithFollowerDataAsync(DateTime date) { var artistFollowerData = await _buzzArtistManager.GetSpotifyArtistFollowerLogAsync(date); return artistFollowerData.Select(a => a.ArtistId).ToList(); } private static async Task UpdateArtistFollowersAsync(IEnumerable artists) { var artistFollowers = artists.Where(x => x?.followers.total != null).Select(x => new Filtr.Spotify.Artists.Data.SpotifyArtistFollowers() { ArtistId = x.id, Followers = (int)x.followers.total }); await PetaPocoRepository.Instance.ImportBulkFileLoaderAsync(artistFollowers, MySqlBulkLoaderConflictOption.Replace); } private void UpdateBuzzArtistFollowers(IEnumerable artistInfo) { var allBuzzArtists = _buzzArtistManager.GetBuzzArtists(); var joinedArtists = allBuzzArtists.Join(artistInfo, buzzArtist => buzzArtist.SpotifyArtistId, apiArtist => apiArtist.id, (buzzArtist, apiArtist) => new { BuzzArtist = buzzArtist, apiArtist = apiArtist }); foreach (var joinedArtist in joinedArtists) { var buzzArtist = joinedArtist.BuzzArtist; buzzArtist.Followers = joinedArtist.apiArtist.followers.total ?? 0; _buzzArtistManager.UpdateBuzzArtist(buzzArtist); } } private async Task SaveFollowerDataToCsvAsync(IEnumerable artistInformation, DateTime importDate) { var cultureInfo = CultureInfo.GetCultureInfo("sv-SE"); using (var saveFile = File.Open(_FilePath, FileMode.OpenOrCreate)) using (var fileWriter = new StreamWriter(saveFile)) { foreach (var artist in artistInformation.Where(a => a != null)) { var csvLine = string.Join("|", artist.id, importDate.ToString(cultureInfo.DateTimeFormat.ShortDatePattern), artist.followers.total); await fileWriter.WriteLineAsync(csvLine); } } } private async Task BulkAddFollowerDataAsync() { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var bulkLoader = new MySqlBulkLoader(conn) { Local = true, FileName = _FilePath, TableName = "tblSpotifyArtistFollowerLog", FieldTerminator = "|", LineTerminator = Environment.NewLine, ConflictOption = MySqlBulkLoaderConflictOption.Replace }; await bulkLoader.LoadAsync(); } } } internal class ArtistsDto { public readonly string[] InputArtistIds; public readonly SpotifyArtist[] ReceivedArtists; public ArtistsDto(IEnumerable input, IEnumerable receivedArtists) { this.InputArtistIds = input.ToArray(); this.ReceivedArtists = receivedArtists.ToArray(); } public string[] MissingArtistsIds { get { return this.InputArtistIds.Except(this.ReceivedArtists.Select(a => a.id)).ToArray(); } } } }