-
Notifications
You must be signed in to change notification settings - Fork 39
/
Metrics.cs
90 lines (81 loc) · 2.93 KB
/
Metrics.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
//---------------------------------------------------------------------------------
// Copyright (c) Microsoft Corporation. All rights reserved.
//
// THIS CODE AND INFORMATION ARE PROVIDED "AS IS" WITHOUT WARRANTY OF ANY KIND,
// EITHER EXPRESSED OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE IMPLIED WARRANTIES
// OF MERCHANTABILITY AND/OR FITNESS FOR A PARTICULAR PURPOSE.
//---------------------------------------------------------------------------------
namespace ThroughputTest
{
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
sealed class Metrics : IObservable<SendMetrics>, IObservable<ReceiveMetrics>
{
object sendMetricsObserverLock = new object();
List<IObserver<SendMetrics>> sendMetricsObservers = new List<IObserver<SendMetrics>>();
object receiveMetricsObserverLock = new object();
List<IObserver<ReceiveMetrics>> receiveMetricsObservers = new List<IObserver<ReceiveMetrics>>();
public Metrics(Settings settings)
{
}
public void PushSendMetrics(SendMetrics sendMetrics)
{
Task.Run(() =>
{
IObserver<SendMetrics>[] observers;
lock (sendMetricsObserverLock)
{
observers = sendMetricsObservers.ToArray();
}
foreach (var observer in observers)
{
observer.OnNext(sendMetrics);
}
}).Fork();
}
public void PushReceiveMetrics(ReceiveMetrics receiveMetrics)
{
Task.Run(() =>
{
IObserver<ReceiveMetrics>[] observers;
lock (receiveMetricsObserverLock)
{
observers = receiveMetricsObservers.ToArray();
}
foreach (var observer in observers)
{
observer.OnNext(receiveMetrics);
}
}).Fork();
}
public IDisposable Subscribe(IObserver<ReceiveMetrics> observer)
{
lock (receiveMetricsObserverLock)
{
receiveMetricsObservers.Add(observer);
}
return System.Reactive.Disposables.Disposable.Create(() =>
{
lock (receiveMetricsObserverLock)
{
receiveMetricsObservers.Remove(observer);
}
});
}
public IDisposable Subscribe(IObserver<SendMetrics> observer)
{
lock (sendMetricsObserverLock)
{
sendMetricsObservers.Add(observer);
}
return System.Reactive.Disposables.Disposable.Create(() =>
{
lock (sendMetricsObserverLock)
{
sendMetricsObservers.Remove(observer);
}
});
}
}
}