diff --git a/homework 2/ClusterClient/Clients/ParallelClusterClient.cs b/homework 2/ClusterClient/Clients/ParallelClusterClient.cs index 5531800..f1613bd 100644 --- a/homework 2/ClusterClient/Clients/ParallelClusterClient.cs +++ b/homework 2/ClusterClient/Clients/ParallelClusterClient.cs @@ -1,7 +1,9 @@ using System; using System.Collections.Generic; +using System.ComponentModel.DataAnnotations; using System.Linq; using System.Text; +using System.Threading; using System.Threading.Tasks; using log4net; @@ -13,9 +15,36 @@ 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 tasks = ReplicaAddresses.Select(async replica => + { + var webRequest = CreateRequest(replica + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + return await ProcessRequestAsync(webRequest); + }).ToList(); + + var timeoutTask = Task.Delay(timeout); + + while (tasks.Count != 0) + { + var completed = await Task.WhenAny(tasks.Append(timeoutTask)); + + if (completed == timeoutTask) + throw new TimeoutException(); + + var result = (Task)completed; + tasks.Remove(result); + try + { + return await result; + } + catch (Exception e) + { + + } + } + throw new Exception(); } protected override ILog Log => LogManager.GetLogger(typeof(ParallelClusterClient)); diff --git a/homework 2/ClusterClient/Clients/ReplicaStatistics.cs b/homework 2/ClusterClient/Clients/ReplicaStatistics.cs new file mode 100644 index 0000000..3d51c7d --- /dev/null +++ b/homework 2/ClusterClient/Clients/ReplicaStatistics.cs @@ -0,0 +1,34 @@ +using System.Collections.Concurrent; +using System.Linq; + +namespace ClusterClient.Clients +{ + public static class ReplicaStatistics + { + private static readonly ConcurrentDictionary> _responseTimes = new(); + private const int MaxStoredTimes = 100; + + public static void RecordResponseTime(string replicaAddress, long elapsedMilliseconds) + { + var bag = _responseTimes.GetOrAdd(replicaAddress, _ => new ConcurrentBag()); + bag.Add(elapsedMilliseconds); + + var newBag = new ConcurrentBag(bag.TakeLast(MaxStoredTimes)); + _responseTimes.TryUpdate(replicaAddress, newBag, bag); + } + + private static double GetAverageResponseTime(string replicaAddress) + { + if (!_responseTimes.TryGetValue(replicaAddress, out var times) || times.IsEmpty) + { + return double.MaxValue; + } + return times.Average(); + } + + public static string[] OrderBySpeed(string[] replicaAddresses) + { + return replicaAddresses.OrderBy(GetAverageResponseTime).ToArray(); + } + } +} diff --git a/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs b/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs index 0293628..e1931e2 100644 --- a/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs +++ b/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs @@ -1,7 +1,5 @@ using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; +using System.Diagnostics; using System.Threading.Tasks; using log4net; @@ -13,11 +11,37 @@ public RoundRobinClusterClient(string[] replicaAddresses) : base(replicaAddresse { } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + var orderedReplicas = ReplicaStatistics.OrderBySpeed(ReplicaAddresses); + + var deadline = DateTime.UtcNow + timeout; + var addressesLeft = orderedReplicas.Length; + + foreach (var replica in orderedReplicas) + { + var remaining = deadline - DateTime.UtcNow; + var replicaTimeout = remaining / addressesLeft; + + var webRequest = CreateRequest(replica + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + + var timer = Stopwatch.StartNew(); + var resultTask = ProcessRequestAsync(webRequest); + + await Task.WhenAny(resultTask, Task.Delay(replicaTimeout)); + timer.Stop(); + + if (resultTask.IsCompletedSuccessfully) + { + ReplicaStatistics.RecordResponseTime(replica, timer.ElapsedMilliseconds); + return await resultTask; + } + addressesLeft--; + } + throw new TimeoutException(); } protected override ILog Log => LogManager.GetLogger(typeof(RoundRobinClusterClient)); } -} +} \ No newline at end of file diff --git a/homework 2/ClusterClient/Clients/SmartClusterClient.cs b/homework 2/ClusterClient/Clients/SmartClusterClient.cs index eb06d8b..c64427f 100644 --- a/homework 2/ClusterClient/Clients/SmartClusterClient.cs +++ b/homework 2/ClusterClient/Clients/SmartClusterClient.cs @@ -1,7 +1,7 @@ using System; using System.Collections.Generic; +using System.Diagnostics; using System.Linq; -using System.Text; using System.Threading.Tasks; using log4net; @@ -12,12 +12,74 @@ public class SmartClusterClient : ClusterClientBase public SmartClusterClient(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 orderedReplicas = ReplicaStatistics.OrderBySpeed(ReplicaAddresses); + + var deadline = DateTime.UtcNow + timeout; + var runningTasks = new List<(Task task, string address, Stopwatch timer)>(); + var addressesLeft = orderedReplicas.Length; + foreach (var replica in orderedReplicas) + { + var remaining = deadline - DateTime.UtcNow; + var replicaTimeout = remaining / addressesLeft; + + var webRequest = CreateRequest(replica + "?query=" + query); + Log.InfoFormat($"Processing {webRequest.RequestUri}"); + + var timer = Stopwatch.StartNew(); + var task = ProcessRequestAsync(webRequest); + runningTasks.Add((task, replica, timer)); + + var finished = await WaitOneAsync(runningTasks, replicaTimeout); + if(finished != null) + { + var (resultTask, resultReplica, resultTimer) = finished.Value; + resultTimer.Stop(); + ReplicaStatistics.RecordResponseTime(resultReplica, resultTimer.ElapsedMilliseconds); + return resultTask.Result; + } + addressesLeft--; + } + + while (runningTasks.Count > 0) + { + var remaining = deadline - DateTime.UtcNow; + if (remaining <= TimeSpan.Zero) + throw new TimeoutException(); + + var finished = await WaitOneAsync(runningTasks, remaining); + if (finished == null) continue; + var (resultTask, resultReplica, resultTimer) = finished.Value; + resultTimer.Stop(); + ReplicaStatistics.RecordResponseTime(resultReplica, resultTimer.ElapsedMilliseconds); + return resultTask.Result; + } + throw new TimeoutException(); + } + protected override ILog Log => LogManager.GetLogger(typeof(SmartClusterClient)); + + private async Task<(Task, string, Stopwatch)?> WaitOneAsync + (List<(Task task, string address, Stopwatch timer)> runningTasks, TimeSpan timeout) + { + if (timeout <= TimeSpan.Zero) + return null; + + var delayTask = Task.Delay(timeout); + var tasksList = runningTasks.Select(Task (x) => x.task).Append(delayTask).ToList(); + var completed = await Task.WhenAny(tasksList); + + if (completed == delayTask) + return null; + + var finishedTask = (Task)completed; + var result = runningTasks.First(x => x.task == finishedTask); + runningTasks.Remove(result); + + return finishedTask.IsCompletedSuccessfully ? (finishedTask, result.address, result.timer) : null; + } } -} +} \ No newline at end of file