master
cc 227 lines 7.19 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #ifdef __cplusplus
4 extern "C" {
5 #endif
6
7 #include "database/rrd.h"
8
9 #ifdef __cplusplus
10 }
11 #endif
12
13 #include <random>
14 #include <thread>
15 #include <vector>
16
17 #define CONFIG_SECTION_PROFILE "plugin:profile"
18
19 class Generator {
20 public:
21 Generator(size_t N) : Offset(0) {
22 std::random_device RandDev;
23 std::mt19937 Gen(RandDev());
24 std::uniform_int_distribution<int> D(-16, 16);
25
26 V.reserve(N);
27 for (size_t Idx = 0; Idx != N; Idx++)
28 V.push_back(D(Gen));
29 }
30
31 double getRandValue() {
32 return V[Offset++ % V.size()];
33 }
34
35 private:
36 size_t Offset;
37 std::vector<double> V;
38 };
39
40 class Profiler {
41 public:
42 Profiler(size_t ID_arg, size_t NumCharts_arg, size_t NumDimsPerChart_arg, time_t SecondsToBackfill_arg, int UpdateEvery_arg) :
43 ID(ID_arg),
44 NumCharts(NumCharts_arg),
45 NumDimsPerChart(NumDimsPerChart_arg),
46 SecondsToBackfill(SecondsToBackfill_arg),
47 UpdateEvery(UpdateEvery_arg),
48 Gen(1024 * 1024)
49 {}
50
51 void create() {
52 char ChartId[1024];
53 char DimId[1024];
54
55 Charts.reserve(NumCharts);
56 for (size_t I = 0; I != NumCharts; I++) {
57 size_t CID = ID + Charts.size() + 1;
58
59 snprintfz(ChartId, 1024 - 1, "chart_%zu", CID);
60
61 RRDSET *RS = rrdset_create_localhost(
62 "profile", // type
63 ChartId, // id
64 nullptr, // name,
65 "profile_family", // family
66 "profile_context", // context
67 "profile_title", // title
68 "profile_units", // units
69 "profile_plugin", // plugin
70 "profile_module", // module
71 12345678 + CID, // priority
72 UpdateEvery, // update_every
73 RRDSET_TYPE_LINE // chart_type
74 );
75 if (I != 0)
76 rrdset_flag_set(RS, RRDSET_FLAG_HIDDEN);
77 Charts.push_back(RS);
78
79 Dimensions.reserve(NumDimsPerChart);
80 for (size_t J = 0; J != NumDimsPerChart; J++) {
81 snprintfz(DimId, 1024 - 1, "dim_%zu", J);
82
83 RRDDIM *RD = rrddim_add(
84 RS, // st
85 DimId, // id
86 nullptr, // name
87 1, // multiplier
88 1, // divisor
89 RRD_ALGORITHM_ABSOLUTE // algorithm
90 );
91
92 Dimensions.push_back(RD);
93 }
94 }
95 }
96
97 void update(const struct timeval &Now) {
98 for (RRDSET *RS: Charts) {
99 for (RRDDIM *RD : Dimensions) {
100 rrddim_timed_set_by_pointer(RS, RD, Now, Gen.getRandValue());
101 }
102
103 rrdset_timed_done(RS, Now, RS->counter_done != 0);
104 }
105 }
106
107 void run() {
108 #define WORKER_JOB_CREATE_CHARTS 0
109 #define WORKER_JOB_UPDATE_CHARTS 1
110 #define WORKER_JOB_METRIC_DURATION_TO_BACKFILL 2
111 #define WORKER_JOB_METRIC_POINTS_BACKFILLED 3
112
113 worker_register("PROFILER");
114 worker_register_job_name(WORKER_JOB_CREATE_CHARTS, "create charts");
115 worker_register_job_name(WORKER_JOB_UPDATE_CHARTS, "update charts");
116 worker_register_job_custom_metric(WORKER_JOB_METRIC_DURATION_TO_BACKFILL, "duration to backfill", "seconds", WORKER_METRIC_ABSOLUTE);
117 worker_register_job_custom_metric(WORKER_JOB_METRIC_POINTS_BACKFILLED, "points backfilled", "points", WORKER_METRIC_ABSOLUTE);
118
119 heartbeat_t HB;
120 heartbeat_init(&HB, UpdateEvery * USEC_PER_SEC);
121
122 worker_is_busy(WORKER_JOB_CREATE_CHARTS);
123 create();
124
125 struct timeval CollectionTV;
126 now_realtime_timeval(&CollectionTV);
127
128 if (SecondsToBackfill) {
129 CollectionTV.tv_sec -= SecondsToBackfill;
130 CollectionTV.tv_sec -= (CollectionTV.tv_sec % UpdateEvery);
131
132 CollectionTV.tv_usec = 0;
133 }
134
135 size_t BackfilledPoints = 0;
136 struct timeval NowTV, PrevTV;
137 now_realtime_timeval(&NowTV);
138 PrevTV = NowTV;
139
140 while (service_running(SERVICE_COLLECTORS)) {
141 worker_is_busy(WORKER_JOB_UPDATE_CHARTS);
142
143 update(CollectionTV);
144 CollectionTV.tv_sec += UpdateEvery;
145
146 now_realtime_timeval(&NowTV);
147
148 ++BackfilledPoints;
149 if (NowTV.tv_sec > PrevTV.tv_sec) {
150 PrevTV = NowTV;
151 worker_set_metric(WORKER_JOB_METRIC_POINTS_BACKFILLED, BackfilledPoints * NumCharts * NumDimsPerChart);
152 BackfilledPoints = 0;
153 }
154
155 size_t RemainingSeconds = (CollectionTV.tv_sec >= NowTV.tv_sec) ? 0 : (NowTV.tv_sec - CollectionTV.tv_sec);
156 worker_set_metric(WORKER_JOB_METRIC_DURATION_TO_BACKFILL, RemainingSeconds);
157
158 if (CollectionTV.tv_sec >= NowTV.tv_sec) {
159 worker_is_idle();
160 heartbeat_next(&HB);
161 }
162 }
163 }
164
165 private:
166 size_t ID;
167 size_t NumCharts;
168 size_t NumDimsPerChart;
169 size_t SecondsToBackfill;
170 int UpdateEvery;
171
172 Generator Gen;
173 std::vector<RRDSET *> Charts;
174 std::vector<RRDDIM *> Dimensions;
175 };
176
177 static void subprofile_main(void* Arg) {
178 Profiler *P = reinterpret_cast<Profiler *>(Arg);
179 P->run();
180 }
181
182 static void profile_main_cleanup(void *pptr) {
183 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)CLEANUP_FUNCTION_GET_PTR(pptr);
184 if(!static_thread) return;
185
186 static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
187
188 static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
189 }
190
191 extern "C" void *profile_main(void *ptr) {
192 CLEANUP_FUNCTION_REGISTER(profile_main_cleanup) cleanup_ptr = ptr;
193
194 int UpdateEvery = (int) inicfg_get_duration_seconds(&netdata_config, CONFIG_SECTION_PROFILE, "update every", 1);
195 if (UpdateEvery < localhost->rrd_update_every) {
196 UpdateEvery = localhost->rrd_update_every;
197 inicfg_set_duration_seconds(&netdata_config, CONFIG_SECTION_PROFILE, "update every", UpdateEvery);
198 }
199
200 // pick low-default values, in case this plugin is ever enabled accidentaly.
201 size_t NumThreads = inicfg_get_number(&netdata_config, CONFIG_SECTION_PROFILE, "number of threads", 2);
202 size_t NumCharts = inicfg_get_number(&netdata_config, CONFIG_SECTION_PROFILE, "number of charts", 2);
203 size_t NumDimsPerChart = inicfg_get_number(&netdata_config, CONFIG_SECTION_PROFILE, "number of dimensions per chart", 2);
204 size_t SecondsToBackfill = inicfg_get_number(&netdata_config, CONFIG_SECTION_PROFILE, "seconds to backfill", 10 * 60);
205
206 std::vector<Profiler> Profilers;
207
208 for (size_t Idx = 0; Idx != NumThreads; Idx++) {
209 Profiler P(1e8 + Idx * 1e6, NumCharts, NumDimsPerChart, SecondsToBackfill, UpdateEvery);
210 Profilers.push_back(P);
211 }
212
213 std::vector<ND_THREAD *> Threads(NumThreads);
214
215 for (size_t Idx = 0; Idx != NumThreads; Idx++) {
216 char Tag[NETDATA_THREAD_TAG_MAX + 1];
217
218 snprintfz(Tag, NETDATA_THREAD_TAG_MAX, "PROFILER[%zu]", Idx);
219 Threads[Idx] = nd_thread_create(Tag, NETDATA_THREAD_OPTION_DEFAULT,
220 subprofile_main, static_cast<void *>(&Profilers[Idx]));
221 }
222
223 for (size_t Idx = 0; Idx != NumThreads; Idx++)
224 nd_thread_join(Threads[Idx]);
225
226 return NULL;
227 }