using MySqlConnector; using System; using System.Collections.Generic; using System.Globalization; using System.Linq; using System.Threading.Tasks; using PetaPoco.Business; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Database; using Sony.Filtr.DistributedCaching; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.SpotifyAnalytics.Models; using Sony.Filtr.Utility.Extensions; using Dapper; using System.Text; using Sony.Filtr.Utility; using System.Collections.Concurrent; using TaskExtensions = Sony.Filtr.Utility.Extensions.TaskExtensions; namespace Sony.Filtr.SpotifyAnalytics { public class SpotifyAnalyticsManager { private readonly DistributedCacheHandler _distributedCacheHandler; private readonly SpotifyAnalyticsFactory _factory; public SpotifyAnalyticsManager(DistributedCacheHandler distributedCacheHandler, SpotifyAnalyticsFactory spotifyAnalyticsFactory) { _distributedCacheHandler = distributedCacheHandler; _factory = spotifyAnalyticsFactory; } public async Task GetLatestDayWithDataAsyncOld() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT Max(Date) FROM tblSpotifyAnalyticsAccountDate WHERE LogType = @logtype", conn); cmd.Parameters.AddWithValue("@logtype", AnalyticsLogType.Finished); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public async Task> GetFilesProcessedCountForDates(IEnumerable dates, IEnumerable supportedFileTypes, int version) { Dictionary result = new Dictionary(); using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand($@" SELECT Date, count(*) as Count FROM tblSpotifyAnalyticsDate WHERE Version = @version AND Date IN ({Maybe.ToCommaSeparated(dates.Select(d => d.ToString("yyyy-MM-dd")))}) AND FileType IN({String.Join(", ", supportedFileTypes.Select(i => (int)i))}) GROUP BY Date;", conn); cmd.Parameters.AddWithValue("@version", version); using (var reader = await cmd.ExecuteReaderAsync()) { while (await reader.ReadAsync()) { result.Add(reader.GetDateTime(0), reader.GetInt32(1)); } } } return result; } public async Task GetLatestDayWithDataAsyncNew() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT MAX(Date) FROM tblSpotifyAnalyticsDate WHERE FileType = @fileType AND Version = 2", conn); cmd.Parameters.AddWithValue("@fileType", SpotifyS3FileType.Playlists); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public async Task GetLatestTimestampWithDataAsync() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT MAX(TimeStamp) FROM tblSpotifyAnalyticsDate WHERE FileType = @fileType AND Version = 2", conn); cmd.Parameters.AddWithValue("@fileType", SpotifyS3FileType.Playlists); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public List GetStreamSummary(SpotifyAnalyticsAccount account, DateTime startDate, DateTime endDate) { return _factory.GetStreamSummary((int)account, startDate, endDate); } public List GetStreamSummary(DateTime startDate, DateTime endDate) { return _factory.GetStreamSummary(startDate, endDate); } public async Task SaveStreamSummaryAsync(int accountId, StreamDataSummaryForAnalysis summary) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { const string sqlString = "INSERT INTO tblSpotifyAnalyticsAccountStreamInfo VALUES (@AccountId, @Market, @Date, @TotalStreams, @DeviceDesktopStreams, @DeviceTabletStreams, @DeviceMobileStreams, " + "@OSAndroidStreams, @OSiOSStreams, @OSWindowsStreams, @OSMacoSStreams, @OSOtherStreams, @SourcePlaylist, @SourceSearch, @SourceArtist, @SourceAlbum, " + "@SourceCollection, @SourceOther, @StreamsFromSonyPlaylists, @StreamsFromLocalSonyPlaylists, @UniqueUsers, " + "@UniqueUsersSourcePlaylist, @UniqueUsersSourceSearch, @UniqueUsersSourceArtist, @UniqueUsersSourceAlbum, @UniqueUsersSourceCollection, @UniqueUsersSourceOther, " + "@UniqueUsersDeviceDesktop, @UniqueUsersDeviceMobile, @UniqueUsersDeviceTablet) " + "ON DUPLICATE KEY UPDATE TotalStreams = @TotalStreams, DeviceDesktopStreams = @DeviceDesktopStreams, DeviceTabletStreams = @DeviceTabletStreams, " + "DeviceMobileStreams = @DeviceMobileStreams, OSAndroidStreams = @OSAndroidStreams, OSiOSStreams = @OSiOSStreams, OSWindowsStreams = @OSWindowsStreams, " + "OSMacoSStreams = @OSMacoSStreams, OSOtherStreams = @OSOtherStreams, SourcePlaylist = @SourcePlaylist, SourceSearch = @SourceSearch, SourceArtist = @SourceArtist, " + "SourceAlbum = @SourceAlbum, SourceCollection = @SourceCollection, SourceOther = @SourceOther, StreamsFromSonyPlaylists = @StreamsFromSonyPlaylists, " + "StreamsFromLocalSonyPlaylists = @StreamsFromLocalSonyPlaylists, UniqueUsers = @UniqueUsers, " + "UniqueUsersSourcePlaylist = @UniqueUsersSourcePlaylist, UniqueUsersSourceSearch = @UniqueUsersSourceSearch, UniqueUsersSourceArtist = @UniqueUsersSourceArtist, " + "UniqueUsersSourceAlbum = @UniqueUsersSourceAlbum, UniqueUsersSourceCollection = @UniqueUsersSourceCollection, UniqueUsersSourceOther = @UniqueUsersSourceOther, " + "UniqueUsersDeviceDesktop = @UniqueUsersDeviceDesktop, UniqueUsersDeviceMobile = @UniqueUsersDeviceMobile, UniqueUsersDeviceTablet = @UniqueUsersDeviceTablet"; using (var command = new MySqlCommand(sqlString, conn)) { command.Parameters.AddWithValue("@AccountId", accountId); command.Parameters.AddWithValue("@Market", summary.Market); command.Parameters.AddWithValue("@Date", summary.Date); command.Parameters.AddWithValue("@TotalStreams", summary.TotalStreams); command.Parameters.AddWithValue("@DeviceDesktopStreams", summary.DeviceDesktopStreams); command.Parameters.AddWithValue("@DeviceTabletStreams", summary.DeviceTabletStreams); command.Parameters.AddWithValue("@DeviceMobileStreams", summary.DeviceMobileStreams); command.Parameters.AddWithValue("@OSAndroidStreams", summary.OSAndroidStreams); command.Parameters.AddWithValue("@OSiOSStreams", summary.OSiOSStreams); command.Parameters.AddWithValue("@OSWindowsStreams", summary.OSWindowsStreams); command.Parameters.AddWithValue("@OSMacoSStreams", summary.OSMacOSStreams); command.Parameters.AddWithValue("@OSOtherStreams", summary.OSOtherStreams); command.Parameters.AddWithValue("@SourcePlaylist", summary.SourcePlaylist); command.Parameters.AddWithValue("@SourceSearch", summary.SourceSearch); command.Parameters.AddWithValue("@SourceArtist", summary.SourceArtist); command.Parameters.AddWithValue("@SourceAlbum", summary.SourceAlbum); command.Parameters.AddWithValue("@SourceCollection", summary.SourceCollection); command.Parameters.AddWithValue("@SourceOther", summary.SourceOther); command.Parameters.AddWithValue("@StreamsFromSonyPlaylists", summary.StreamsFromSonyPlaylists); command.Parameters.AddWithValue("@StreamsFromLocalSonyPlaylists", summary.StreamsFromLocalSonyPlaylists); command.Parameters.AddWithValue("@UniqueUsers", summary.UniqueUsers); command.Parameters.AddWithValue("@UniqueUsersSourcePlaylist", summary.UniqueUsersSourcePlaylist); command.Parameters.AddWithValue("@UniqueUsersSourceSearch", summary.UniqueUsersSourceSearch); command.Parameters.AddWithValue("@UniqueUsersSourceArtist", summary.UniqueUsersSourceArtist); command.Parameters.AddWithValue("@UniqueUsersSourceAlbum", summary.UniqueUsersSourceAlbum); command.Parameters.AddWithValue("@UniqueUsersSourceCollection", summary.UniqueUsersSourceCollection); command.Parameters.AddWithValue("@UniqueUsersSourceOther", summary.UniqueUsersSourceOther); command.Parameters.AddWithValue("@UniqueUsersDeviceDesktop", summary.UniqueUsersDeviceDesktop); command.Parameters.AddWithValue("@UniqueUsersDeviceMobile", summary.UniqueUsersDeviceMobile); command.Parameters.AddWithValue("@UniqueUsersDeviceTablet", summary.UniqueUsersDeviceTablet); await command.ExecuteNonQueryAsync(); } } } public async Task SetSpotifyAnalyticsLogAsync(int accountId, string market, AnalyticsLogType logType, DateTime analyticsDate, int version) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { const string sqlString = "INSERT INTO tblSpotifyAnalyticsAccountDate (AccountId, Market, Date, LogType, Version) VALUES(@accountId, @market, @analyticsDate, @logType, @version) ON DUPLICATE KEY UPDATE Timestamp = CURRENT_TIMESTAMP"; var command = new MySqlCommand(sqlString, conn); command.Parameters.AddWithValue("@accountId", accountId); command.Parameters.AddWithValue("@market", market); command.Parameters.AddWithValue("@analyticsDate", analyticsDate); command.Parameters.AddWithValue("@logType", logType); command.Parameters.AddWithValue("@version", version); await command.ExecuteNonQueryAsync(); } } public async Task> GetSpotifyAnalyticsLogAsync(int accountId, DateTime analyticsDate, AnalyticsLogType logType, int version) { var dates = new Dictionary(); using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { const string sqlString = "SELECT Market, Date FROM tblSpotifyAnalyticsAccountDate WHERE AccountId=@accountId AND Date=@date AND LogType=@logType AND Version=@version ORDER BY Timestamp DESC"; var command = new MySqlCommand(sqlString, conn); command.Parameters.AddWithValue("@accountId", accountId); command.Parameters.AddWithValue("@date", analyticsDate); command.Parameters.AddWithValue("@logType", logType); command.Parameters.AddWithValue("@version", version); var reader = await command.ExecuteReaderAsync(); while (await reader.ReadAsync()) { dates.Add(reader.GetString(0), reader.GetDateTime(1)); } reader.Close(); } return dates; } public async Task> GetSpotifyAnalyticsLogAsync(AnalyticsLogType logType, int version) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE logType=@0 AND Version = @1", logType, version); } public async Task> GetPlaylistStreamDataAsync(string playlistId, string legacyPlaylistUri, DateTime startDate, DateTime endDate) { return await _factory.GetPlaylistStreamDataAsync(playlistId, legacyPlaylistUri, null, startDate, endDate); } public async Task> GetPlaylistStreamDataAsync(string playlistId, string legacyPlaylistUri, SpotifyAnalyticsAccount account, DateTime startDate, DateTime endDate) { return await _factory.GetPlaylistStreamDataAsync(playlistId, legacyPlaylistUri, (int)account, startDate, endDate); } public async Task> GetPlaylistsStreamSummaryAsync(List playlistIds) { var cachedStreamSummaries = await GetCachedPlaylistStreamSummariesAsync(playlistIds); var nonCachedStreamSummaries = playlistIds.Except(cachedStreamSummaries.Select(p => p.PlaylistId)).ToList(); if (nonCachedStreamSummaries.Any()) { var dbResults = await _factory.GetPlaylistsStreamSummaryAsync(nonCachedStreamSummaries); await CachePlaylistStreamSummariesAsync(dbResults); return cachedStreamSummaries.Union(dbResults).ToList(); } return cachedStreamSummaries; } private async Task CachePlaylistStreamSummariesAsync(List dbResults) { foreach (var playlistStreamSummary in dbResults) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlistStreamSummary.PlaylistId); _distributedCacheHandler.Put(cacheKey, playlistStreamSummary); } } private async Task> GetCachedPlaylistStreamSummariesAsync(List playlistIds) { List playlistStreamsSummaries = new List(); foreach (var playlist in playlistIds) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlist); var playlistResult = _distributedCacheHandler.Get(cacheKey) as List; if (playlistResult != null) playlistStreamsSummaries.AddRange(playlistResult); } return playlistStreamsSummaries; } public async Task> GetMarketPlaylistStreamsAsync(DateTime date, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync(date, buzzCategoryId); } public async Task> GetMarketPlaylistStreamsAsync(string market, DateTime startDate, DateTime endDate, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync(market, startDate, endDate, buzzCategoryId); } public async Task> GetMarketPlaylistStreamsAsync(SpotifyAnalyticsAccount spotifyAnalyticsAccount, string market, DateTime startDate, DateTime endDate, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync((int)spotifyAnalyticsAccount, market, startDate, endDate, buzzCategoryId); } public async Task> CalculatePlaylistStreamSummaryAsync(IEnumerable playlistUri, DateTime latestDate) { return await _factory.CalculatePlaylistStreamSummaryAsync(playlistUri, latestDate); } public async Task CalculatePlaylistCategoryStreams(int spotifyAnalyticsAccountId, DateTime streamsDate, Application app, List appRegions, StaticBuzzCategory buzzCategory) { var categoryStreams = new PlaylistCategoryStreams() { AccountId = spotifyAnalyticsAccountId, Market = app.SpotifyRegionCode, BuzzCategoryId = (int)buzzCategory, Date = streamsDate, TotalListStreams = 0, LocalListStreams = 0, TotalListListeners = 0, LocalListListeners = 0, }; foreach (var region in appRegions) { using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { const string sql = @" SELECT p.BuzzCategoryId AS StreamCategory, s.date AS StreamDate, SUM(s.TotalStreams) AS TotalListStreams, SUM(IF(p.CountryCode = @market, s.TotalStreams, 0)) AS LocalListStreams, SUM(s.UniqueUsers) AS TotalListListeners, SUM(IF(p.CountryCode = @market, s.UniqueUsers, 0)) AS LocalListListeners FROM tblSpotifyPlaylist AS p INNER JOIN tblSpotifyAnalyticsAccountPlaylistStreamInfo AS s ON p.playlistId = s.playlistUri AND s.Date = @date AND s.accountId = @accountId LEFT JOIN tblPlaylistIgnoredPlaylist AS ip ON ip.playlistId=p.playlistUri AND ip.musicServiceId=@spotifyMusicServiceId WHERE p.removed = 0 AND ip.playlistID IS NULL AND p.BuzzCategoryId = @buzzCategoryId AND (@market = '_gl' OR s.Market = @region) AND s.PlaylistUri <> '' GROUP BY s.date"; MySqlCommand cmd = new MySqlCommand(sql, conn); cmd.CommandTimeout = 3600; cmd.Parameters.AddWithValue("@market", app.SpotifyRegionCode); cmd.Parameters.AddWithValue("@accountId", spotifyAnalyticsAccountId); cmd.Parameters.AddWithValue("@region", region); cmd.Parameters.AddWithValue("@date", streamsDate); cmd.Parameters.AddWithValue("@buzzCategoryId", buzzCategory); cmd.Parameters.AddWithValue("@spotifyMusicServiceId", MusicService.Spotify); var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { categoryStreams.TotalListStreams += reader.GetIntOrFallback("TotalListStreams", 0); categoryStreams.LocalListStreams += reader.GetIntOrFallback("LocalListStreams", 0); categoryStreams.TotalListListeners += reader.GetIntOrFallback("TotalListListeners", 0); categoryStreams.LocalListListeners += reader.GetIntOrFallback("LocalListListeners", 0); } } } return categoryStreams; } private string GetPlaylistStreamSummaryCacheKey(string playlistId) { return "SpotifyAnalytics_PlaylistStreamSummary_2.0_" + playlistId; } public async Task ClearPlaylistStreamSummaryCacheAsync(IEnumerable playlistIds) { foreach (var playlistId in playlistIds) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlistId); _distributedCacheHandler.Remove(cacheKey); } } public async Task> GetTrackListenersAsync(List isrcs, DateTime startDate, DateTime endDate, string countryCode) { var distributerId = (int)SpotifyAnalyticsAccount.Sony; //We only get data this distributer so we specify it so no behaviour changes. return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE isrc IN (@isrcs) AND DistributerId = @distributerId AND date >= @startDate AND date <= @endDate AND (@countryCode IS NULL OR countryCode = @countryCode)", new { isrcs = isrcs, distributerId = distributerId, startDate = startDate, endDate = endDate, countryCode = countryCode }); } public List GetPlaylistProcessedDates(DateTime startDate, DateTime endDate) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE Date > @startDate AND Date < @endDate", new { startDate, endDate }); } public DateTime? GetLatestProccessedDateWithFullData() { return PetaPocoRepository.ReadOnlyInstance.FirstOrDefaultWithSql("WHERE AggregatedTimestamp IS NOT NULL AND CategoryStreamsSummary IS NOT NULL ORDER BY Date DESC")?.Date; } public List GetImportedFiles(DateTime date) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE date=@0", date); } public async Task SetMarketsAsCalculatedCategoryStreamsAsync(DateTime date) { string sql = @" INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, CategoryStreamsSummary) SELECT CountryCode, DistributerId, @date, Now() FROM tblSpotifyAnalyticsDate WHERE Date = @date ON DUPLICATE KEY UPDATE CategoryStreamsSummary = VALUES(CategoryStreamsSummary);"; using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { MySqlCommand cmd = conn.CreateCommand(); cmd.CommandText = sql; cmd.Parameters.AddWithValue("@date", date); await cmd.ExecuteNonQueryAsync(); } } public async Task SetMarketsAsCalculatedCategoryStreamsAsync(List importedFilesForToday, DateTime date) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var timestamp = DateTime.UtcNow; var sqlDateFormatCulture = CultureInfo.GetCultureInfo("sv-SE"); var sql = "INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, CategoryStreamsSummary) " + "VALUES " + string.Join(",", importedFilesForToday.Select(p => $"('{MySqlHelper.EscapeString(p.CountryCode)}', {p.DistributerId}, '{date.ToString(sqlDateFormatCulture)}', '{timestamp.ToString(sqlDateFormatCulture)}')")) + "ON DUPLICATE KEY UPDATE CategoryStreamsSummary = VALUES(CategoryStreamsSummary);"; MySqlCommand cmd = new MySqlCommand(sql, conn); await cmd.ExecuteNonQueryAsync(); } } public async Task SetMarketsAsAggregatedAsync(List importedFiles, DateTime date) { if (!importedFiles.Any()) { return; } using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var timestamp = DateTime.UtcNow; var sqlDateFormatCulture = CultureInfo.GetCultureInfo("sv-SE"); var sql = "INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, AggregatedTimestamp) " + "VALUES " + string.Join(",", importedFiles.Select(p => $"('{MySqlHelper.EscapeString(p.CountryCode)}', {p.DistributerId}, '{date.ToString(sqlDateFormatCulture)}', '{timestamp.ToString(sqlDateFormatCulture)}')")) + "ON DUPLICATE KEY UPDATE aggregatedTimestamp = VALUES(AggregatedTimestamp);"; MySqlCommand cmd = new MySqlCommand(sql, conn); await cmd.ExecuteNonQueryAsync(); } } } public enum AnalyticsLogType { Started, Finished } }