| 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 | } |