using MoreLinq; using Sony.Filtr.Buzz; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Core.Buzz; using Sony.Filtr.ErrorLogging; using Sony.Filtr.Playlists; using Sony.Filtr.Playlists.Spotify; using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Sony.Filtr.Playlists.Spotify.Model; using System.Threading.Tasks.Dataflow; using System.Collections.ObjectModel; using Sony.Filtr.ApolloAPI; using Sony.Filtr.ApolloAPI.Models; using Sony.Filtr.Functional; using Sony.Filtr.Tasks.Helpers; using System.Threading; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.Spotify { public class ImportBuzzPlaylists : IScheduledTask { private readonly BuzzManager _buzzManager; private readonly ScheduledTaskManager _scheduledTaskManager; private readonly IgnoredPlaylistsManager _ignoredPlaylistsManager; private readonly SpotifyPlaylistManager _spotifyPlaylistManager; private readonly BuzzAccountManager _buzzAccountManager; private readonly AsyncLogger _logger; private readonly IApolloWebApi _vendorApi; public ImportBuzzPlaylists(BuzzManager buzzManager, IApolloWebApi vendorApi, BuzzAccountManager buzzAccountManager, ScheduledTaskManager scheduledTaskManager, IgnoredPlaylistsManager ignoredPlaylistsManager, SpotifyPlaylistManager spotifyPlaylistManager) { _buzzManager = buzzManager; _vendorApi = vendorApi; _scheduledTaskManager = scheduledTaskManager; _ignoredPlaylistsManager = ignoredPlaylistsManager; _spotifyPlaylistManager = spotifyPlaylistManager; _buzzAccountManager = buzzAccountManager; _logger = AsyncLogger.GetLogger("ImportBuzzPlaylists"); } public async Task ExecuteAsync(Guid scheduledTaskLogId) { return await ImportBuzzPlaylistsAsync(scheduledTaskLogId); } private async Task ImportBuzzPlaylistsAsync(Guid scheduledTaskLogId) { var taskLog = new ScheduledTaskLog(); var dataLog = new ScheduledTaskDataLog("ImportSpotifyBuzzPlaylist", DateTimeOffset.UtcNow, scheduledTaskLogId); var importUsers = await GetUsersToImportFromAsync(); var lastTaskInfo = await _scheduledTaskManager.GetLatestFinishedLogAsync("ImportBuzzPlaylists"); var lastImportDate = lastTaskInfo?.Finished ?? DateTimeOffset.Now.Date.AddDays(-1); var ignoredPlaylist = await _ignoredPlaylistsManager.GetIgnoredPlaylistsAsync(MusicService.Spotify); var ignoredPlaylistIds = ignoredPlaylist.Select(p => p.PlaylistId).Distinct().ToHashSet(); int totalProcessedPlaylistsCount = 0; Func> getUserPlaylists = async userName => { return await _vendorApi.GetAllPublicPlaylistsAsync(userName, "items(id,name,owner.id),total"); }; var loadPlaylistsBlock = new TransformBlock>( getUserPlaylists .Retry(2) .Duration((u, items, duration) => _logger.TraceAsync(() => $"Received {items.Length} playlists for user '{u}' in {duration.TotalSeconds}secs")) .TryCatch() .OnFailure((u, result) => _logger.ErrorAsync(result.Exception, () => $"Could not request playlists for user '{u}'")) .TapAsync((user, userPlaylists) => _logger.DebugAsync(() => $"Found {userPlaylists.Length} playlists for user '{user}'")) .BindAsync((user, userPlaylists) => userPlaylists.Where(p => p.owner.id.Equals(user, StringComparison.InvariantCultureIgnoreCase)).ToList()) .TapAsync((user, ownPlaylists) => _logger.DebugAsync(() => $"Found {ownPlaylists.Count} own playlists for user '{user}'")) .BindAsync(ownPlaylists => ownPlaylists.Select(p => BuildPlaylistReference(p)).ToList()) .TapAsync(refs => Interlocked.Add(ref totalProcessedPlaylistsCount, refs.Count)) .BindAsync((user, refs) => refs.Where(p => !ignoredPlaylistIds.Contains(p.PlaylistId)).ToList()) .BindAsync((user, refs) => new UserPlaylists(user, refs)) , new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 2 }); Func addOrUpdateUserPlaylists = async dto => { if (dto.Playlists.Any()) { _logger.DebugAsync(() => $"Begin saving {dto.Playlists.Count()} playlists for user: {dto.UserName}").FireAndForget(); await _spotifyPlaylistManager.AddOrUpdateSpotifyPlaylistsForTrackingAsync(dto.Playlists); } else { _logger.DebugAsync(() => $"Could not find any playlists for user {dto.UserName}").FireAndForget(); } }; var savePlaylistsBlock = new ActionBlock>( addOrUpdateUserPlaylists .ToUnit() .Map, Unit, UserPlaylists, Unit>( inputTransform: result => result.Value, outputTransform: (playlists, unit, result) => Unit.Default) .TryCatch() .OnFailure((dto, result) => _logger.ErrorAsync(() => $"Could not save Spotify playlists for user '{dto.Value.UserName}'")) , new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 10 }); loadPlaylistsBlock.LinkTo(savePlaylistsBlock, new DataflowLinkOptions { PropagateCompletion = true }, dto => dto.IsOk && dto.Value.Playlists != null); loadPlaylistsBlock.LinkTo(DataflowBlock.NullTarget>(), new DataflowLinkOptions { PropagateCompletion = true }); _logger.InfoAsync(() => $"Found {importUsers.Count} uses to process.").FireAndForget(); try { foreach (var userName in importUsers) { await loadPlaylistsBlock.SendAsync(userName); } loadPlaylistsBlock.Complete(); await savePlaylistsBlock.Completion; } catch (Exception ex) { _logger.ErrorAsync(ex, () => $"Exception: {ex}").FireAndForget(); ErrorLoggingManager.Instance.LogError(ex); } _logger.InfoAsync(() => $"Ended ImportBuzzPlaylists. Total processed playlists: {totalProcessedPlaylistsCount}").FireAndForget(); dataLog.Finished = DateTimeOffset.UtcNow; dataLog.Rows = importUsers.Count(); taskLog.DataLogs.Add(dataLog); return taskLog; } private async Task> GetUsersToImportFromAsync() { _logger.InfoAsync(() => "Begin ImportBuzzPlaylists").FireAndForget(); var oldBuzzUsers = _buzzManager.GetBuzzImportUsers(); _logger.DebugAsync(() => "Found " + oldBuzzUsers.Count() + " old import users.").FireAndForget(); List importFrom = new List(); importFrom.AddRange(oldBuzzUsers.Select(p => p.SpotifyUserName)); var buzzUsers = await _buzzAccountManager.GetBuzzUsersAsync(musicServiceId: (int)MusicService.Spotify); importFrom.AddRange(buzzUsers.Where(p => p.ImportPlaylists).Select(p => p.Username)); return importFrom.Distinct().ToList(); } private SpotifyPlaylistTrackingReference BuildPlaylistReference(SpotifyPlaylistItem playlistItem) { return new SpotifyPlaylistTrackingReference() { PlaylistId = playlistItem.id, Name = playlistItem.name, User = playlistItem.owner?.id, SaveTracklist = true, SnapshotId = playlistItem.snapshot_id }; } private class UserPlaylists { public readonly string UserName; public readonly IReadOnlyCollection Playlists; public UserPlaylists(string userName, IList playlists) { this.UserName = userName; this.Playlists = playlists == null ? null : new ReadOnlyCollection(playlists); } } } }