using MediaBrowser.Common.Configuration; using MediaBrowser.Common.Extensions; using MediaBrowser.Common.IO; using MediaBrowser.Controller.Sync; using MediaBrowser.Model.Logging; using MediaBrowser.Model.Serialization; using MediaBrowser.Model.Sync; using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace MediaBrowser.Server.Implementations.Sync { public class TargetDataProvider : ISyncDataProvider { private readonly SyncTarget _target; private readonly IServerSyncProvider _provider; private readonly SemaphoreSlim _dataLock = new SemaphoreSlim(1, 1); private List _items; private readonly ILogger _logger; private readonly IJsonSerializer _json; private readonly IFileSystem _fileSystem; private readonly IApplicationPaths _appPaths; private readonly string _serverId; private readonly SemaphoreSlim _cacheFileLock = new SemaphoreSlim(1, 1); public TargetDataProvider(IServerSyncProvider provider, SyncTarget target, string serverId, ILogger logger, IJsonSerializer json, IFileSystem fileSystem, IApplicationPaths appPaths) { _logger = logger; _json = json; _provider = provider; _target = target; _fileSystem = fileSystem; _appPaths = appPaths; _serverId = serverId; } private string GetCachePath() { return Path.Combine(_appPaths.DataPath, "sync", _target.Id.GetMD5().ToString("N") + ".json"); } private string GetRemotePath() { var parts = new List { _serverId, "data.json" }; return _provider.GetFullPath(parts, _target); } private async Task CacheData(Stream stream) { var cachePath = GetCachePath(); await _cacheFileLock.WaitAsync().ConfigureAwait(false); try { Directory.CreateDirectory(Path.GetDirectoryName(cachePath)); using (var fileStream = _fileSystem.GetFileStream(cachePath, FileMode.Create, FileAccess.Write, FileShare.Read, true)) { await stream.CopyToAsync(fileStream).ConfigureAwait(false); } } catch (Exception ex) { _logger.ErrorException("Error saving sync data to {0}", ex, cachePath); } finally { _cacheFileLock.Release(); } } private async Task EnsureData(CancellationToken cancellationToken) { if (_items == null) { try { using (var stream = await _provider.GetFile(GetRemotePath(), _target, new Progress(), cancellationToken)) { _items = _json.DeserializeFromStream>(stream); } } catch (FileNotFoundException) { _items = new List(); } catch (DirectoryNotFoundException) { _items = new List(); } using (var memoryStream = new MemoryStream()) { _json.SerializeToStream(_items, memoryStream); // Now cache it memoryStream.Position = 0; await CacheData(memoryStream).ConfigureAwait(false); } } } private async Task SaveData(CancellationToken cancellationToken) { using (var stream = new MemoryStream()) { _json.SerializeToStream(_items, stream); // Save to sync provider stream.Position = 0; await _provider.SendFile(stream, GetRemotePath(), _target, new Progress(), cancellationToken).ConfigureAwait(false); // Now cache it stream.Position = 0; await CacheData(stream).ConfigureAwait(false); } } private async Task GetData(Func, T> dataFactory) { await _dataLock.WaitAsync().ConfigureAwait(false); try { await EnsureData(CancellationToken.None).ConfigureAwait(false); return dataFactory(_items); } finally { _dataLock.Release(); } } private async Task UpdateData(Func, List> action) { await _dataLock.WaitAsync().ConfigureAwait(false); try { await EnsureData(CancellationToken.None).ConfigureAwait(false); _items = action(_items); await SaveData(CancellationToken.None).ConfigureAwait(false); } finally { _dataLock.Release(); } } public Task> GetServerItemIds(SyncTarget target, string serverId) { return GetData(items => items.Where(i => string.Equals(i.ServerId, serverId, StringComparison.OrdinalIgnoreCase)).Select(i => i.ItemId).ToList()); } public Task AddOrUpdate(SyncTarget target, LocalItem item) { return UpdateData(items => { var list = items.Where(i => !string.Equals(i.Id, item.Id, StringComparison.OrdinalIgnoreCase)) .ToList(); list.Add(item); return list; }); } public Task Delete(SyncTarget target, string id) { return UpdateData(items => items.Where(i => !string.Equals(i.Id, id, StringComparison.OrdinalIgnoreCase)).ToList()); } public Task Get(SyncTarget target, string id) { return GetData(items => items.FirstOrDefault(i => string.Equals(i.Id, id, StringComparison.OrdinalIgnoreCase))); } private async Task> GetCachedData() { if (_items == null) { await _cacheFileLock.WaitAsync().ConfigureAwait(false); try { if (_items == null) { try { _items = _json.DeserializeFromFile>(GetCachePath()); } catch (FileNotFoundException) { _items = new List(); } catch (DirectoryNotFoundException) { _items = new List(); } } } finally { _cacheFileLock.Release(); } } return _items.ToList(); } public async Task> GetCachedServerItemIds(SyncTarget target, string serverId) { var items = await GetCachedData().ConfigureAwait(false); return items.Where(i => string.Equals(i.ServerId, serverId, StringComparison.OrdinalIgnoreCase)) .Select(i => i.ItemId) .ToList(); } public async Task GetCachedItem(SyncTarget target, string id) { var items = await GetCachedData().ConfigureAwait(false); return items.FirstOrDefault(i => string.Equals(i.Id, id, StringComparison.OrdinalIgnoreCase)); } } }