-
Notifications
You must be signed in to change notification settings - Fork 402
Expand file tree
/
Copy pathCounterPipeline.cs
More file actions
111 lines (97 loc) · 3.81 KB
/
Copy pathCounterPipeline.cs
File metadata and controls
111 lines (97 loc) · 3.81 KB
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
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using Microsoft.Diagnostics.NETCore.Client;
using Microsoft.Diagnostics.Tracing;
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
namespace Microsoft.Diagnostics.Monitoring.EventPipe
{
internal class CounterPipeline : EventSourcePipeline<CounterPipelineSettings>
{
private readonly IEnumerable<ICountersLogger> _loggers;
private readonly CounterFilter _filter;
private string _sessionId;
public CounterPipeline(DiagnosticsClient client,
CounterPipelineSettings settings,
IEnumerable<ICountersLogger> loggers) : base(client, settings)
{
_loggers = loggers ?? throw new ArgumentNullException(nameof(loggers));
if (settings.CounterGroups.Length > 0)
{
_filter = new CounterFilter(Settings.CounterIntervalSeconds);
foreach (var counterGroup in settings.CounterGroups)
{
_filter.AddFilter(counterGroup.ProviderName, counterGroup.CounterNames);
}
}
else
{
_filter = CounterFilter.AllCounters(Settings.CounterIntervalSeconds);
}
}
protected override MonitoringSourceConfiguration CreateConfiguration()
{
var config = new MetricSourceConfiguration(Settings.CounterIntervalSeconds, _filter.GetProviders(), Settings.MaxHistograms, Settings.MaxTimeSeries);
_sessionId = config.SessionId;
return config;
}
protected override async Task OnEventSourceAvailable(EventPipeEventSource eventSource, Func<Task> stopSessionAsync, CancellationToken token)
{
await ExecuteCounterLoggerActionAsync((metricLogger) => metricLogger.PipelineStarted(token));
eventSource.Dynamic.All += traceEvent =>
{
try
{
if (traceEvent.TryGetCounterPayload(_filter, _sessionId, out List<ICounterPayload> counterPayload))
{
ExecuteCounterLoggerAction((metricLogger) => {
foreach (var payload in counterPayload)
{
metricLogger.Log(payload);
}
});
}
}
catch (Exception)
{
}
};
using var sourceCompletedTaskSource = new EventTaskSource<Action>(
taskComplete => taskComplete,
handler => eventSource.Completed += handler,
handler => eventSource.Completed -= handler,
token);
await sourceCompletedTaskSource.Task;
await ExecuteCounterLoggerActionAsync((metricLogger) => metricLogger.PipelineStopped(token));
}
private async Task ExecuteCounterLoggerActionAsync(Func<ICountersLogger, Task> action)
{
foreach (ICountersLogger logger in _loggers)
{
try
{
await action(logger);
}
catch (ObjectDisposedException)
{
}
}
}
private void ExecuteCounterLoggerAction(Action<ICountersLogger> action)
{
foreach (ICountersLogger logger in _loggers)
{
try
{
action(logger);
}
catch (ObjectDisposedException)
{
}
}
}
}
}