using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Net; using System.Text; using System.Threading.Tasks; using Amazon.CloudSearchDomain; using Amazon.CloudSearchDomain.Model; using Amazon.Runtime; using ByteSizeLib; using Newtonsoft.Json; using Sony.Filtr.Contracts.Abstractions; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Contracts.Entities.Genre; using Sony.Filtr.Core.LastModified; using Sony.Filtr.ErrorLogging; using Sony.Filtr.Search.Json; using Sony.Filtr.Search.SearchItems; using Sony.Filtr.Utility.Extensions; using SearchResponse = Amazon.CloudSearchDomain.Model.SearchResponse; namespace Sony.Filtr.Search { public class SearchManager : ISearchManager { private readonly AWSCredentials _awsCredentials; private readonly AmazonCloudSearchDomainConfig _amazonCloudSearchDomainConfig; private readonly ErrorLoggingManager _errorLoggingManager; private readonly LastModifiedManager _lastModifiedManager; private readonly CloudSearchMapper _cloudSearchMapper; private readonly JsonSerializerSettings _jsonSerializerSettings; public SearchManager(AWSCredentials credentials, AmazonCloudSearchDomainConfig amazonCloudSearchDomainConfig, ErrorLoggingManager errorLoggingManager, LastModifiedManager lastModifiedManager) { _awsCredentials = credentials; _amazonCloudSearchDomainConfig = amazonCloudSearchDomainConfig; _errorLoggingManager = errorLoggingManager; _lastModifiedManager = lastModifiedManager; _cloudSearchMapper = new CloudSearchMapper(); _jsonSerializerSettings = new JsonSerializerSettings() { Formatting = Formatting.None, NullValueHandling = NullValueHandling.Ignore, ContractResolver = new LowercaseContractResolver(), }; _jsonSerializerSettings.Converters.Add(new BoolConverter()); } public async Task> SearchPlaylistsAsync(PlaylistSearchRequest playlistSearchRequest) { var client = new AmazonCloudSearchDomainClient(_awsCredentials, _amazonCloudSearchDomainConfig); const string returnFields = "typename,name,spotifyuri,applicationid,priority,imageurl,countrycode,description,followers,buzzcategoryid"; SearchResponse playlistSearchResult; if (string.IsNullOrWhiteSpace(playlistSearchRequest.SearchQuery)) { playlistSearchResult = await client.SearchAsync(new SearchRequest() { Query = BuildPlaylistSearchFilter(playlistSearchRequest), QueryParser = "structured", QueryOptions = "{ fields: ['name^5', 'description^2', 'tags^1', 'artists^0.5', 'tracks^0.1'] }", Return = returnFields, Size = 10000, }); } else { var queryString = playlistSearchRequest.SearchQuery; if (!queryString.Contains(" ")) queryString = queryString + "~2"; else queryString = string.Format("(and {0})", string.Join(" ", queryString.Split(' ').Select(a => "'" + a + "~1'"))); playlistSearchResult = await client.SearchAsync(new SearchRequest() { Query = queryString, QueryOptions = "{ fields: ['name^5', 'description^2', 'tags^1', 'artists^0.5', 'tracks^0.1'] }", FilterQuery = BuildPlaylistSearchFilter(playlistSearchRequest), Return = returnFields, Size = 10000, }); } return _cloudSearchMapper.Map(playlistSearchResult.Hits).ToList(); } private string BuildPlaylistSearchFilter(PlaylistSearchRequest playlistSearchRequest) { List conditions = new List(); if (playlistSearchRequest.ServiceType != null) conditions.Add(BuildOrIntCondition("servicetype", playlistSearchRequest.ServiceType.Select(s=> (int)s))); if (playlistSearchRequest.ArtistQuery != null) conditions.Add(string.Format("artists:'{0}'", playlistSearchRequest.ArtistQuery)); if (playlistSearchRequest.TagQuery != null) conditions.Add(string.Format("tags:'{0}'", playlistSearchRequest.TagQuery)); if (playlistSearchRequest.MarketIds != null) conditions.Add(BuildOrIntCondition("applicationid", playlistSearchRequest.MarketIds.Select(s => (int)s))); if (playlistSearchRequest.CountryCode != null) conditions.Add(string.Format("countrycode:'{0}'", playlistSearchRequest.CountryCode)); if (playlistSearchRequest.BuzzCategoryIds != null) conditions.Add(BuildOrIntCondition("buzzcategoryid", playlistSearchRequest.BuzzCategoryIds)); if (playlistSearchRequest.Owner != null) conditions.Add(string.Format("owner:'{0}'", playlistSearchRequest.Owner)); return string.Format("(and typename:'{0}' {1} )", typeof (PlaylistSearchIndexItem).Name, string.Join(" ", conditions)); } private static string BuildOrIntCondition(string field, IEnumerable ints) { var serviceTypeCondition = "(or "; foreach (var integer in ints) { serviceTypeCondition += field + ":" + integer; } serviceTypeCondition += ")"; return serviceTypeCondition; } public class Condition { public string Field { get; set; } public string ConditionValue { get; set; } } public async Task SearchStuffAsync(Application application, ServiceType serviceType, string searchString) { var client = new AmazonCloudSearchDomainClient(_awsCredentials, _amazonCloudSearchDomainConfig); var searchResult = new SearchResult(); var queryString = searchString; if (!queryString.Contains(" ")) queryString = queryString + "~2"; else queryString = string.Format("(and {0})", string.Join(" ", queryString.Split(' ').Select(a => "'" + a + "~1'"))); var playlistSearchResult = client.SearchAsync(new SearchRequest() { Query = queryString, QueryOptions = "{ fields: ['name^5', 'description^2', 'tags^1', 'artists^0.5', 'tracks^0.1'] }", FilterQuery = string.Format("(and applicationid:'{0}' typename:'{1}' iseditorialplaylist:'1' servicetype:'{2}')", application.ID, typeof(PlaylistSearchIndexItem).Name, (int)serviceType), Size = 20, }); var releaseSearchResult = client.SearchAsync(new SearchRequest() { Query = queryString, FilterQuery = string.Format("(and applicationid:'{0}' typename:'{1}' servicetype:'{2}')", application.ID, typeof(ReleaseSearchIndexItem).Name, (int)serviceType), Size = 20, }); var artistSearchResult = client.SearchAsync(new SearchRequest() { Query = queryString, FilterQuery = string.Format("(and typename:'{0}')", typeof(ArtistSearchIndexItem).Name), Sort = "artistboost desc", Size = 20, }); var tagSearchResult = client.SearchAsync(new SearchRequest() { Query = queryString, FilterQuery = string.Format("(and applicationid:'{0}' typename:'{1}')", application.ID, typeof(TagSearchIndexItem).Name), Size = 20, }); searchResult.Playlists = _cloudSearchMapper.Map((await playlistSearchResult).Hits).ToList(); searchResult.Releases = _cloudSearchMapper.Map((await releaseSearchResult).Hits).ToList(); searchResult.Artists = _cloudSearchMapper.Map((await artistSearchResult).Hits).ToList(); searchResult.Tags = _cloudSearchMapper.Map((await tagSearchResult).Hits).ToList(); return searchResult; } public async Task IndexPlaylistsAsync(IEnumerable playlists, IDictionary> playlistCategories) { var extendedEditorialPlaylists = playlists as IList ?? playlists.ToList(); var searchIndexDocuments = extendedEditorialPlaylists.Select(p => new SearchIndexDocument() { Type = "add", Id = "p" + p.EditorialPlaylist.ID, Fields = new PlaylistSearchIndexItem(p, (playlistCategories.ContainsKey(p.EditorialPlaylist.ID)) ?playlistCategories[p.EditorialPlaylist.ID] : null) }); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(extendedEditorialPlaylists.Select(p=> p.EditorialPlaylist.ApplicationID)); } private async Task SetLastModifiedDate(IEnumerable applicationIds) { foreach (var applicationId in applicationIds) { await _lastModifiedManager.SetLastModifyDateAsync(applicationId, EntityType.Search); } } public async Task DeletePlaylistsFromIndexAsync(IEnumerable playlists) { var extendedEditorialPlaylists = playlists as IList ?? playlists.ToList(); var searchIndexDocuments = extendedEditorialPlaylists.Select(p => new SearchIndexDocument() { Type = "delete", Id = "p" + p.ID, }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(extendedEditorialPlaylists.Select(p=> p.ApplicationID)); } public async Task IndexTagsAsync(IEnumerable tags) { var enumerable = tags as IList ?? tags.ToList(); var searchIndexDocuments = enumerable.Select(p => new SearchIndexDocument() { Type = "add", Id = "t" + p.ID, Fields = new TagSearchIndexItem(p) }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(enumerable.Select(p => p.ApplicationID)); } public async Task IndexReleasesAsync(IEnumerable releases) { var enumerable = releases as IList ?? releases.ToList(); var searchIndexDocuments = enumerable.Select(r => new SearchIndexDocument() { Type = "add", Id = "sr:" + r.TrackLink + ":" + r.ApplicationID, Fields = new ReleaseSearchIndexItem(r) }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(enumerable.Select(r => r.ApplicationID)); } public async Task DeleteReleasesFromIndexAsync(IEnumerable releases) { var enumerable = releases as IList ?? releases.ToList(); var searchIndexDocuments = enumerable.Select(r => new SearchIndexDocument() { Type = "delete", Id = "sr:" + r.TrackLink + ":" + r.ApplicationID, }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(enumerable.Select(r => r.ApplicationID)); } public async Task IndexDeezerReleasesAsync(IEnumerable releases) { var deezerReleases = releases as IList ?? releases.ToList(); var searchIndexDocuments = deezerReleases.Select(r => new SearchIndexDocument() { Type = "add", Id = "dr" + r.ID, Fields = new ReleaseSearchIndexItem(r) }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(deezerReleases.Select(r => r.ApplicationID)); } public async Task DeleteDeezerReleasesAsync(IEnumerable releases) { var deezerReleases = releases as IList ?? releases.ToList(); var searchIndexDocuments = deezerReleases.Select(r => new SearchIndexDocument() { Type = "delete", Id = "dr" + r.ID, }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); await SetLastModifiedDate(deezerReleases.Select(r => r.ApplicationID)); } private async Task IndexSearchIndexDocuments(IEnumerable searchIndexDocuments) { var domainClient = new AmazonCloudSearchDomainClient(_awsCredentials, _amazonCloudSearchDomainConfig); try { await BatchBySize(searchIndexDocuments).ForEachAsync(1, async documentBatch => { var uploadDocumentsStream = GetSerializedUploadDocumentsStream(documentBatch); var uploadResult = await domainClient.UploadDocumentsAsync(new UploadDocumentsRequest() { ContentType = ContentType.ApplicationJson, Documents = uploadDocumentsStream, }); if (uploadResult.HttpStatusCode != HttpStatusCode.OK) { _errorLoggingManager.LogError(new Exception("SearchIndexing error, recieved statuscode:" + uploadResult.HttpStatusCode)); } }); } catch (Exception ex) { _errorLoggingManager.LogError(new Exception("SearchIndexing exception", ex)); throw; } } private MemoryStream GetSerializedUploadDocumentsStream(object documentBatch) { var serializedObject = JsonConvert.SerializeObject(documentBatch, Formatting.None, _jsonSerializerSettings); var uploadDocumentsStream = new MemoryStream(Encoding.UTF8.GetBytes(serializedObject)); return uploadDocumentsStream; } private IEnumerable> BatchBySize(IEnumerable searchIndexDocuments) { //const int maxBatchSize = 5242880; // 5 MB in bytes int maxBatchSize = (int)ByteSize.FromMegaBytes(2).Bytes; //First try with all, only then batch ICollection searchIndexDocumentsCollection = searchIndexDocuments as ICollection; if (searchIndexDocumentsCollection != null) { var allStream = GetSerializedUploadDocumentsStream(searchIndexDocuments); if (allStream.Length < maxBatchSize) { yield return searchIndexDocumentsCollection.ToList(); yield break; } } var testBatch = new List(); var serializedObjectSizes = searchIndexDocuments.Select(s => new { IndexObject = s, SizeOfSerializedObject = GetSerializedUploadDocumentsStreamSize(s), }); //var serializationTime = sw2.ElapsedMilliseconds; const int serializedBatchSizeInitialValue = 2; //Assume 2 bytes for [] long serializedBatchSize = serializedBatchSizeInitialValue; testBatch.Clear(); foreach (var serializedObjectSize in serializedObjectSizes) { testBatch.Add(serializedObjectSize.IndexObject); serializedBatchSize += serializedObjectSize.SizeOfSerializedObject + 1; //Assume 1 byte for , if (serializedBatchSize > maxBatchSize) { testBatch.Remove(serializedObjectSize.IndexObject); yield return new List(testBatch); // Reset, start of new batch testBatch.Clear(); testBatch.Add(serializedObjectSize.IndexObject); serializedBatchSize = serializedBatchSizeInitialValue; } } if (testBatch.Any()) { yield return testBatch; } } private long GetSerializedUploadDocumentsStreamSize(SearchIndexDocument searchIndexDocument) { return GetSerializedUploadDocumentsStream(searchIndexDocument).Length; } public async Task DeleteAllReleasesFromIndexAsync() { var client = new AmazonCloudSearchDomainClient(_awsCredentials, _amazonCloudSearchDomainConfig); long releaseToRemove = 1; string cursor = "initial"; int loop = 0; while (releaseToRemove != 0) { Console.WriteLine("Begin search for loop " + loop++); var releaseSearchResult = await client.SearchAsync(new SearchRequest() { Query = "lolzcat|-lolzcat", FilterQuery = string.Format("(and typename:'{0}' servicetype:'{1}')", typeof(ReleaseSearchIndexItem).Name, (int)ServiceType.Spotify), Size = 10000, Return = "name", Cursor = cursor, }); releaseToRemove = releaseSearchResult.Hits.Hit.Count; Console.WriteLine("Search done. ReleaseToRemove:" + releaseToRemove); if (releaseToRemove > 0) { cursor = releaseSearchResult.Hits.Cursor; Console.WriteLine("Begin delete"); var searchIndexDocuments = releaseSearchResult.Hits.Hit.Select(r => new SearchIndexDocument() { Type = "delete", Id = r.Id, }).ToList(); await IndexSearchIndexDocuments(searchIndexDocuments); Console.WriteLine("Done delete"); } } } } }