-
Notifications
You must be signed in to change notification settings - Fork 848
/
Receiver.cs
96 lines (83 loc) · 3.59 KB
/
Receiver.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.
using System;
using System.Diagnostics;
using System.IO;
using System.Net;
using System.Net.Http;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Yarp.Kubernetes.Controller.Configuration;
using Yarp.Kubernetes.Controller.Hosting;
using Yarp.Kubernetes.Controller.Rate;
using Yarp.ReverseProxy.Forwarder;
namespace Yarp.Kubernetes.Protocol;
public class Receiver : BackgroundHostedService
{
private readonly ReceiverOptions _options;
private readonly Limiter _limiter;
private readonly IUpdateConfig _proxyConfigProvider;
public Receiver(
IOptions<ReceiverOptions> options,
IHostApplicationLifetime hostApplicationLifetime,
ILogger<Receiver> logger,
IUpdateConfig proxyConfigProvider) : base(hostApplicationLifetime, logger)
{
if (options is null)
{
throw new ArgumentNullException(nameof(options));
}
_options = options.Value;
_options.Client ??= new HttpMessageInvoker(new SocketsHttpHandler
{
ConnectTimeout = TimeSpan.FromSeconds(15),
});
// two requests per second after third failure
_limiter = new Limiter(new Limit(2), 3);
_proxyConfigProvider = proxyConfigProvider;
}
public override async Task RunAsync(CancellationToken cancellationToken)
{
while (!cancellationToken.IsCancellationRequested)
{
await _limiter.WaitAsync(cancellationToken).ConfigureAwait(false);
#pragma warning disable CA1303 // Do not pass literals as localized parameters
Logger.LogInformation("Connecting with {ControllerUrl}", _options.ControllerUrl.ToString());
try
{
var requestMessage = new HttpRequestMessage(HttpMethod.Get, _options.ControllerUrl);
var responseMessage = await _options.Client.SendAsync(requestMessage, cancellationToken).ConfigureAwait(false);
using var stream = await responseMessage.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false);
using var reader = new StreamReader(stream, Encoding.UTF8, leaveOpen: true);
using var cancellation = cancellationToken.Register(stream.Close);
while (true)
{
var json = await reader.ReadLineAsync().ConfigureAwait(false);
if (string.IsNullOrEmpty(json))
{
break;
}
var message = System.Text.Json.JsonSerializer.Deserialize<Message>(json);
Logger.LogInformation("Received {MessageType} for {MessageKey}", message.MessageType, message.Key);
Logger.LogInformation(json);
Logger.LogInformation(message.MessageType.ToString());
if (message.MessageType == MessageType.Update)
{
await _proxyConfigProvider.UpdateAsync(message.Routes, message.Cluster, cancellation.Token).ConfigureAwait(false);
}
}
}
#pragma warning disable CA1031 // Do not catch general exception types
catch (Exception ex)
#pragma warning restore CA1031 // Do not catch general exception types
{
Logger.LogInformation(ex, "Stream ended");
}
#pragma warning restore CA1303 // Do not pass literals as localized parameters
}
}
}