using System; using System.Collections.Generic; using System.Globalization; using System.IO; using System.Linq; using System.Threading.Tasks; using Newtonsoft.Json; using NLog; using NodaTime.Serialization.JsonNet; using NodaTime.TimeZones; using PetaPoco.Business; using Sony.Filtr.Contracts.Abstractions; using Sony.Filtr.SpotifyAnalytics; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.SpotifyAnalytics.Models; namespace Sony.Filtr.Tasks.Tasks.Spotify.Analytics { public abstract class AggregatedImportTaskBase : IScheduledTask { private readonly SpotifyStreamingAggregatedReportApi _aggregatedReportApi; private readonly Logger _logger; protected const string TempFolder = "spotifyanalytics"; protected AggregatedImportTaskBase(SpotifyStreamingAggregatedReportApi aggregatedReportApi) { _aggregatedReportApi = aggregatedReportApi; _logger = LogManager.GetLogger("SpotifyStreamingAggregatedReportApi"); } public abstract Task ExecuteAsync(Guid scheduledTaskLogId); protected static void MarkAsProcessed(DateTime date, SpotifyS3File trackFile) { PetaPocoRepository.Instance.Insert(new SpotifyAnalyticsAggregatedDate() { Date = date, CountryCode = trackFile.Country, DistributerId = trackFile.DistributorId, FileType = trackFile.FileType, Version = trackFile.Version, }); } protected IEnumerable ReadNJsonData(string localPath, int trackFileDistributorId) where T : IAggregatedStreamFormat { var jsonSerializer = new JsonSerializer().ConfigureForNodaTime(new DateTimeZoneCache(new BclDateTimeZoneSource())); using (var fileStream = File.OpenRead(localPath)) { using (var jsonReader = new JsonTextReader(new StreamReader(fileStream)) { SupportMultipleContent = true }) { while (jsonReader.Read()) { var obj = jsonSerializer.Deserialize(jsonReader); obj.DistributerId = trackFileDistributorId; yield return obj; } } } } protected async Task> GetFilesToImportFromS3(DateTime date, IEnumerable supportedFileTypes) { _logger.Info($"Getting file list from S3 for date {date.ToShortDateString()}"); var files = await _aggregatedReportApi.GetFilesAsync(date); var filesWeCanImport = files.Where(p => p.Version == 2 && supportedFileTypes.Contains(p.FileType)); var alreadyImportedFiles = PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE date=@0", date); var notImportedFiles = filesWeCanImport.Where(p => !alreadyImportedFiles.Any(a => a.CountryCode == p.Country && a.DistributerId == p.DistributorId && a.FileType == p.FileType && a.Version == p.Version)).ToList(); return notImportedFiles; } protected async Task DownloadFileFromS3Async(SpotifyS3File playlistFile) { _logger.Debug($"Downloading {playlistFile.FilePath}"); var dateFilenamePart = playlistFile.Date.ToString("yyyy-dd-M--HH-mm-ss", CultureInfo.GetCultureInfo("sv-SE")); var localPath = $"{TempFolder}/track-{dateFilenamePart}-{playlistFile.DistributorId}-{playlistFile.Country}-{playlistFile.Version}-{playlistFile.FileType}.ndjson"; await _aggregatedReportApi.DownloadToFileAsync(localPath, playlistFile.FilePath); return localPath; } } }