using System; using System.Net.Http; using System.Threading; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; namespace Sony.Filtr.Utility.HttpHandlers { public class ThreadRateLimitHttpHandler : DelegatingHandler { private readonly int Capacity; private ActionBlock> ProcessBlock; public ThreadRateLimitHttpHandler(int capacity) { this.Capacity = capacity; this.ProcessBlock = new ActionBlock>(async input => { try { //Console.WriteLine($"LIMIT + {input.Input.RequestUri}"); var response = await base.SendAsync(input.Input, input.Token); //Console.WriteLine($"LIMIT - {response.StatusCode} {input.Input.RequestUri}"); input.TaskCompletionsSource.SetResult(response); } catch (Exception ex) { //Console.ForegroundColor = ConsoleColor.Green; //Console.WriteLine($"LIMIT !!!!!!!!!!!!! {input.Input.RequestUri} {ex.Message}"); //Console.ResetColor(); input.TaskCompletionsSource.SetException(ex); } }, new ExecutionDataflowBlockOptions { EnsureOrdered = false, MaxDegreeOfParallelism = capacity }); } protected override Task SendAsync(HttpRequestMessage request, CancellationToken cancellationToken) { var taskInput = TaskInput.Create(request, cancellationToken); //Console.WriteLine($"LIMIT +++ {request.RequestUri}"); this.ProcessBlock.SendAsync(taskInput); //Console.WriteLine($"LIMIT --- {request.RequestUri} {this.ProcessBlock.Completion.Status} {this.ProcessBlock.Completion.Exception}"); return taskInput.TaskCompletionsSource.Task; } private class TaskInput { public readonly TInput Input; public readonly TaskCompletionSource TaskCompletionsSource; public CancellationToken Token; private TaskInput(TaskCompletionSource tcs, TInput input, CancellationToken token) { this.Input = input; this.TaskCompletionsSource = tcs; this.Token = token; } public static TaskInput Create(TInput input, CancellationToken token) { return new TaskInput(new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously), input, token); } } } }