diff --git a/homework 2/ClusterClient/Clients/ParallelClusterClient.cs b/homework 2/ClusterClient/Clients/ParallelClusterClient.cs index 5531800..3dd32bc 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.Linq; +using System.Net; using System.Text; +using System.Threading; using System.Threading.Tasks; using log4net; @@ -13,9 +15,32 @@ 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 = new List>(); + var timeoutTask = Task.Delay(timeout); + + for (var i = 0; i < ReplicaAddresses.Length; ++i) + { + var request = CreateRequest(ReplicaAddresses[i] + "?query=" + query); + tasks.Add(ProcessRequestAsync(request)); + } + + while (tasks.Count > 0) + { + var completedTask = await Task.WhenAny(tasks.Concat(new[] {timeoutTask})); + if (completedTask == timeoutTask) + throw new TimeoutException(); + try + { + return await (Task)completedTask; + } + catch + { + tasks.Remove((Task)completedTask); + } + } + 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..ddc93b6 100644 --- a/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs +++ b/homework 2/ClusterClient/Clients/RoundRobinClusterClient.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Diagnostics; using System.Linq; using System.Text; using System.Threading.Tasks; @@ -9,13 +10,50 @@ namespace ClusterClient.Clients { public class RoundRobinClusterClient : ClusterClientBase { + private readonly ReplicasStatistics _replicasStatistics; public RoundRobinClusterClient(string[] replicaAddresses) : base(replicaAddresses) { + _replicasStatistics = new ReplicasStatistics(replicaAddresses); } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + var timer = Stopwatch.StartNew(); + var sortedReplicas = ReplicaAddresses + .OrderBy(x => _replicasStatistics.GetStats(x)) + .ToArray(); + var sentTasks = new Dictionary, (string, DateTime)>(); + for (int i = 0; i < sortedReplicas.Length; i++) + { + var remainingTime = timeout - timer.Elapsed; + if (remainingTime <= TimeSpan.Zero) + throw new TimeoutException(); + var timeoutForReplica = TimeSpan.FromTicks(remainingTime.Ticks / (sortedReplicas.Length - i)); + if (timeoutForReplica <= TimeSpan.Zero) + throw new TimeoutException(); + var timeoutTask = Task.Delay(timeoutForReplica); + var request = CreateRequest(sortedReplicas[i] + "?query=" + query); + var task = ProcessRequestAsync(request); + sentTasks.Add(task, (address: sortedReplicas[i], startTime: DateTime.Now)); + var requestResult = await Task.WhenAny(task, timeoutTask); + if (requestResult != timeoutTask) + { + var replicaData = sentTasks[task]; + var time = DateTime.Now - replicaData.Item2; + var result = (Task)requestResult; + if (result.Status == TaskStatus.RanToCompletion) + { + _replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds); + sentTasks.Remove(task); + return result.Result; + } + _replicasStatistics.UpdateStats(replicaData.Item1, timeout.TotalMilliseconds); + sentTasks.Remove(task); + } + else + _replicasStatistics.UpdateStats(sortedReplicas[i], timeoutForReplica.TotalMilliseconds); + } + throw new TimeoutException(); } 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..4c2cfb0 100644 --- a/homework 2/ClusterClient/Clients/SmartClusterClient.cs +++ b/homework 2/ClusterClient/Clients/SmartClusterClient.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Diagnostics; using System.Linq; using System.Text; using System.Threading.Tasks; @@ -9,15 +10,56 @@ namespace ClusterClient.Clients { public class SmartClusterClient : ClusterClientBase { + private readonly ReplicasStatistics _replicasStatistics; + public SmartClusterClient(string[] replicaAddresses) : base(replicaAddresses) { + _replicasStatistics = new ReplicasStatistics(replicaAddresses); } - public override Task ProcessRequestAsync(string query, TimeSpan timeout) + public override async Task ProcessRequestAsync(string query, TimeSpan timeout) { - throw new NotImplementedException(); + var timer = Stopwatch.StartNew(); + var sortedReplicas = ReplicaAddresses + .OrderBy(x => _replicasStatistics.GetStats(x)) + .ToArray(); + var requests = new Dictionary, (string, DateTime)>(); + for (int i = 0; i < sortedReplicas.Length; i++) + { + var remainingTime = timeout - timer.Elapsed; + if (remainingTime <= TimeSpan.Zero) + throw new TimeoutException(); + var timeoutForReplica = TimeSpan.FromTicks(remainingTime.Ticks / (sortedReplicas.Length - i)); + if (timeoutForReplica <= TimeSpan.Zero) + throw new TimeoutException(); + var timeoutTask = Task.Delay(timeoutForReplica); + var request = CreateRequest(sortedReplicas[i] + "?query=" + query); + requests.Add(ProcessRequestAsync(request), (sortedReplicas[i], DateTime.Now)); + var requestResult = await Task.WhenAny(requests.Keys.Concat(new[] { timeoutTask })); + if (requestResult != timeoutTask) + { + var result = (Task)requestResult; + var replicaData = requests[result]; + var time = DateTime.Now - replicaData.Item2; + if (result.Status == TaskStatus.RanToCompletion) + { + _replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds); + requests.Remove(result); + foreach (var badRequest in requests) + _replicasStatistics.UpdateStats(badRequest.Value.Item1, timeoutForReplica.TotalMilliseconds); + requests.Clear(); + return result.Result; + } + _replicasStatistics.UpdateStats(replicaData.Item1, time.TotalMilliseconds); + requests.Remove(result); + } + else + _replicasStatistics.UpdateStats(sortedReplicas[i], timeoutForReplica.TotalMilliseconds); + } + + throw new TimeoutException(); } protected override ILog Log => LogManager.GetLogger(typeof(SmartClusterClient)); } -} +} \ No newline at end of file diff --git a/homework 2/ClusterClient/ReplicasStatistics.cs b/homework 2/ClusterClient/ReplicasStatistics.cs new file mode 100644 index 0000000..b14b73a --- /dev/null +++ b/homework 2/ClusterClient/ReplicasStatistics.cs @@ -0,0 +1,31 @@ +using System.Collections.Concurrent; +using System.Linq; + +namespace ClusterClient; + +public class ReplicasStatistics +{ + private const double ConfidenceFactor = 0.8; + private readonly ConcurrentDictionary _replicaStatistics; + + public ReplicasStatistics(string[] replicaAddresses) + { + _replicaStatistics = new ConcurrentDictionary(); + foreach (var replicaAddress in replicaAddresses) + _replicaStatistics.TryAdd(replicaAddress, 0); + } + + public void UpdateStats(string replicaAddress, double workTime) + { + _replicaStatistics.AddOrUpdate( + replicaAddress, + workTime, + (_, previousValue) => previousValue * ConfidenceFactor + workTime * (1 - ConfidenceFactor) + ); + } + + public double GetStats(string replicaAddress) + { + return _replicaStatistics.TryGetValue(replicaAddress, out double value) ? value : 0; + } +} \ No newline at end of file