using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.IO; using System.Linq; using System.Threading.Tasks; using MoreLinq; using NLog; using Sony.Filtr.Contracts.Abstractions; using Sony.Filtr.ErrorLogging; using Sony.Filtr.PlaylistSynchronization; using Sony.Filtr.Utility; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks { public class PlaylistSynchronizationTask : IScheduledTask { private readonly PlaylistSynchronizationManager _playlistSynchronizationManager; private readonly IServiceAccountManager _serviceAccountManager; private readonly Logger _logger; public PlaylistSynchronizationTask(PlaylistSynchronizationManager playlistSynchronizationManager, IServiceAccountManager serviceAccountManager) { _playlistSynchronizationManager = playlistSynchronizationManager; _serviceAccountManager = serviceAccountManager; _logger = LogManager.GetLogger("PlaylistSynchronization"); } public async Task ExecuteAsync(Guid scheduledTaskLogId) { _logger.Info("Begin playlist sync"); var playlistSyncs = this.GetSynchronizationsToProcess(); _logger.Info("Found {0} playlist sync definitions", playlistSyncs.Count); var processedPlaylistSyncs = await SyncPlaylistsAsync(playlistSyncs); playlistSyncs = this.GetSynchronizationsToProcess(); var playlistSyncsToProcess = playlistSyncs.ExceptBy(processedPlaylistSyncs, synchronization => synchronization.Id).ToList(); _logger.Info("Found {0} playlist sync definitions that had been added during the run or failed", playlistSyncsToProcess.Count); await SyncPlaylistsAsync(playlistSyncsToProcess); _logger.Info("Done with playlist sync"); return null; } private List GetSynchronizationsToProcess() { var allSyncs = _playlistSynchronizationManager.GetPlaylistSynchronizations(onlyActive: true).ToList(); return allSyncs.Where(new SynchronizationByAccountIdFilter().GetPredicate()).ToList(); } private async Task> SyncPlaylistsAsync(List playlistSyncs) { var processedPlaylistSyncs = new ConcurrentBag(); int totalToSync = playlistSyncs.Count; //foreach(var playlist in playlistSyncs) await playlistSyncs.ItemIndex().ForEachAsync( Maybe.GetAppSettingsIntOrDefault("PlaylistSynchronizationTask_NumberOfParallelThreadsToProcessSynchronizations", 1), async playlist => { try { //Fetch the playlist sync object so we get any changes that has occurred during the run var playlistSyncToUse = _playlistSynchronizationManager.GetPlaylistSynchronization(playlist.Item.Id); if (!playlistSyncToUse.Active) { _logger.Debug("Playlist sync has been disabled, skipping."); return; } _logger.Debug($"Synchronizing playlist ID: {playlistSyncToUse.Id}. From: {playlistSyncToUse.FromPlaylistId} to {playlistSyncToUse.ToPlaylistId}. {playlist.Index} out of {totalToSync}"); var serviceAccount = _serviceAccountManager.GetServiceAccount(playlistSyncToUse.ToServiceAccountId); var result = await _playlistSynchronizationManager.ExecutePlaylistSyncAsync(playlistSyncToUse, serviceAccount, forceAllSources: false); if (string.IsNullOrWhiteSpace(result.ErrorText)) { processedPlaylistSyncs.Add(playlistSyncToUse); } } catch (Exception ex) { ErrorLoggingManager.Instance.LogError(ex); _logger.Error($"Error when synchronizing playlistSync with id {playlist.Item.Id}, exception: {ex}"); } }); return processedPlaylistSyncs.ToList(); } //private async Task ValidateYoutubeServiceAccountsAsync() //{ // var youtubeServiceAccounts = _serviceAccountManager.GetServiceAccounts().Where(s => s.ServiceType == ServiceType.YouTube); // foreach (var youtubeServiceAccount in youtubeServiceAccounts) // { // try // { // var authClient = _youtubeApi.GetAuthenticatedYoutubeApi(youtubeServiceAccount); // var channelResponse = await authClient.GetMyChannelsAsync(); // var userInfo = await (new GoogleOAuth2Api("")).GetUserInfo(new GoogleOAuth2Api.AccessTokenResponse() // { // access_token = youtubeServiceAccount.AccessToken // }); // Console.WriteLine("ServiceAccount: " + youtubeServiceAccount.Id); // Console.WriteLine("Name: " + userInfo.name); // foreach (var channel in channelResponse.items) // { // Console.WriteLine("Channel: " + channel.id); // } // youtubeServiceAccount.DisplayName = userInfo.name; // _serviceAccountManager.SaveServiceAccount(youtubeServiceAccount); // } // catch (Exception ex) // { // Console.WriteLine("Error for {0}. Exception: {1}", youtubeServiceAccount.Id, ex); // } // } //} } public abstract class SynchronizationFilter { public abstract Func GetPredicate(); public abstract bool ShouldFilter(); } public abstract class SynchronizationFromFileFilter : SynchronizationFilter { public readonly string FilterFilePath; public SynchronizationFromFileFilter(string filePath) { this.FilterFilePath = filePath; } protected abstract T Converter(string line); public override bool ShouldFilter() { return File.Exists(FilterFilePath); } protected T[] GetFilterData() { return File.ReadAllLines(FilterFilePath) .Select(line => this.Converter(line)) .ToArray(); } } public class SynchronizationByAccountIdFilter : SynchronizationFromFileFilter { public SynchronizationByAccountIdFilter() : base(@"PlaylistSynchronizationByAccountId.txt") { } public override Func GetPredicate() { if (this.ShouldFilter()) { var accountIdsToUse = this.GetFilterData(); return sync => accountIdsToUse.Contains(sync.ToServiceAccountId); } return s => true; } protected override int Converter(string line) { return Int32.Parse(line); } } }