diff --git a/homework 2/ClusterClient/Clients/ParallelClusterClient.cs b/homework 2/ClusterClient/Clients/ParallelClusterClient.cs index 5531800..0d9fc41 100644 --- a/homework 2/ClusterClient/Clients/ParallelClusterClient.cs +++ b/homework 2/ClusterClient/Clients/ParallelClusterClient.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Diagnostics; using System.Linq; using System.Text; using System.Threading.Tasks; @@ -13,9 +14,38 @@ public ParallelClusterClient(string[] replicaAddresses) : base(replicaAddresses) { } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + var timeoutTask = Task.Delay(timeout).ContinueWith(_ => throw new TimeoutException()); + + var uris = ReplicaAddresses.Select(async uri => + { + var webRequest = CreateRequest(uri + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + var resultTask = await ProcessRequestAsync(webRequest); + return resultTask; + + }).ToList(); + uris.Add(timeoutTask); + + while (uris.Any()) + { + try + { + var resultTask = await Task.WhenAny(uris); + if (resultTask == timeoutTask) + await timeoutTask; + return await resultTask; + } + catch (Exception e) + { + if (uris.Count == 2) + throw; + else + uris.Remove(uris.First()); + } + } + throw new TimeoutException(); } protected override ILog Log => LogManager.GetLogger(typeof(ParallelClusterClient)); diff --git a/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs b/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs index 0293628..b90ba48 100644 --- a/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs +++ b/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Net.Http; using System.Text; using System.Threading.Tasks; using log4net; @@ -9,13 +10,77 @@ namespace ClusterClient.Clients { public class RoundRobinClusterClient : ClusterClientBase { + private readonly Dictionary> stats = new (); + private readonly object lockObj = new (); + public RoundRobinClusterClient(string[] replicaAddresses) : base(replicaAddresses) { + foreach (var address in replicaAddresses) + stats[address] = new Queue(); } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + var orderedReplicas = ReplicaAddresses + .OrderBy(GetAverage) + .ToArray(); + + var partTimeout = timeout / orderedReplicas.Length; + var globalTimeout = Task.Delay(timeout); + + foreach (var replicaAddress in orderedReplicas) + { + var startTime = DateTime.UtcNow; + var timeoutTask = Task.Delay(partTimeout); + var webRequest = CreateRequest(replicaAddress + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + var replicaTask = ProcessRequestAsync(webRequest); + var resultTask = await Task.WhenAny(replicaTask, timeoutTask); + + if (resultTask == timeoutTask) + continue; + + if(resultTask.IsFaulted) + continue; + + var result = await replicaTask; + UpdateStatistics(replicaAddress, DateTime.UtcNow - startTime); + return result; + } + + var lastReplica = orderedReplicas.Last(); + var webRequestLast = CreateRequest(lastReplica + "?query=" + query); + Log.InfoFormat($"Processing {webRequestLast.RequestUri}"); + + var startLast = DateTime.UtcNow; + var replicaLastTask = ProcessRequestAsync(webRequestLast); + + var resultLastTask = await Task.WhenAny(replicaLastTask, globalTimeout); + + if (resultLastTask == globalTimeout) + throw new TimeoutException(); + + var resultFinal = await replicaLastTask; + UpdateStatistics(lastReplica, DateTime.UtcNow - startLast); + return resultFinal; + } + + private double GetAverage(string replica) + { + lock (lockObj) + { + var queue = stats[replica]; + return queue.Count == 0 ? double.MaxValue : queue.Average(); + } + } + + private void UpdateStatistics(string replica, TimeSpan time) + { + lock (lockObj) + { + var queue = stats[replica]; + queue.Enqueue(time.TotalMilliseconds); + } } protected override ILog Log => LogManager.GetLogger(typeof(RoundRobinClusterClient)); diff --git a/homework 2/ClusterClient/Clients/SmartClusterClient.cs b/homework 2/ClusterClient/Clients/SmartClusterClient.cs index eb06d8b..5427cfe 100644 --- a/homework 2/ClusterClient/Clients/SmartClusterClient.cs +++ b/homework 2/ClusterClient/Clients/SmartClusterClient.cs @@ -1,7 +1,9 @@ using System; using System.Collections.Generic; +using System.Globalization; using System.Linq; using System.Text; +using System.Threading; using System.Threading.Tasks; using log4net; @@ -9,15 +11,108 @@ namespace ClusterClient.Clients { public class SmartClusterClient : ClusterClientBase { + private readonly Dictionary> stats = new Dictionary>(); + private readonly object lockObj = new object(); + public SmartClusterClient(string[] replicaAddresses) : base(replicaAddresses) { + foreach (var address in replicaAddresses) + stats[address] = new Queue(); } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + { + var orderedReplicas = ReplicaAddresses + .OrderBy(GetAverage) + .ToArray(); + var partOfTimeout = timeout / orderedReplicas.Length; + var tasksAtWork = new List<(Task task, string replica, DateTime startTime)>(); + var timeoutGlobal = Task.Delay(timeout) + .ContinueWith(_ => throw new TimeoutException()); + var tasks = new List> { timeoutGlobal }; + + foreach (var replicaAddress in orderedReplicas) + { + var webRequest = CreateRequest(replicaAddress + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + var startTime = DateTime.UtcNow; + var replicaTask = ProcessRequestAsync(webRequest); + + tasksAtWork.Add((replicaTask, replicaAddress, startTime)); + tasks.Add(replicaTask); + + var timeoutDelay = Task.Delay(partOfTimeout).ContinueWith(_ => throw new TimeoutException()); + + tasks.Add(timeoutDelay); + + var taskInCycle = await Task.WhenAny(tasks); + + tasks.Remove(timeoutDelay); + + if (taskInCycle == timeoutGlobal) + throw new TimeoutException(); + + if (taskInCycle == timeoutDelay) + continue; + + try + { + var result = await taskInCycle; + var info = tasksAtWork.First(t => t.task == taskInCycle); + UpdateStatistics(info.replica, DateTime.UtcNow - info.startTime); + + return result; + } catch + { + tasks.Remove(taskInCycle); + tasksAtWork.RemoveAll(t => t.task == taskInCycle); + } + } + while (tasksAtWork.Count > 1) + { + var completedTask = await Task.WhenAny(tasks); + + if (completedTask == timeoutGlobal) + throw new TimeoutException(); + + try + { + var result = await completedTask; + + var info = tasksAtWork.First(t => t.task == completedTask); + UpdateStatistics(info.replica, DateTime.UtcNow - info.startTime); + + return result; + } + catch + { + tasks.Remove(completedTask); + tasks.RemoveAll(t => t == completedTask); + } + } + } + throw new TimeoutException(); } protected override ILog Log => LogManager.GetLogger(typeof(SmartClusterClient)); + + private double GetAverage(string replica) + { + lock (lockObj) + { + var queue = stats[replica]; + return queue.Count == 0 ? double.MaxValue : queue.Average(); + } + } + + private void UpdateStatistics(string replica, TimeSpan time) + { + lock (lockObj) + { + var queue = stats[replica]; + queue.Enqueue(time.TotalMilliseconds); + } + } } }