master
c 426 lines 16.4 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #include "database/rrd.h"
4
5 #ifdef ENABLE_DBENGINE
6
7 #define CHARTS 64
8 #define DIMS 16 // CHARTS * DIMS dimensions
9 #define REGIONS 11
10 #define POINTS_PER_REGION 16384
11 static const int REGION_UPDATE_EVERY[REGIONS] = {1, 15, 3, 20, 2, 6, 30, 12, 5, 4, 10};
12
13 #define START_TIMESTAMP MAX(2 * API_RELATIVE_TIME_MAX, 200000000)
14
15 static time_t region_start_time(time_t previous_region_end_time, time_t update_every) {
16 // leave a small gap between regions
17 // but keep them close together, so that cross-region queries will be fast
18
19 time_t rc = previous_region_end_time + update_every;
20 rc += update_every - (rc % update_every);
21 rc += update_every;
22 return rc;
23 }
24
25 static inline collected_number point_value_get(size_t region, size_t chart, size_t dim, size_t point) {
26 // calculate the value to be stored for each point in the database
27
28 collected_number r = (collected_number)region;
29 collected_number c = (collected_number)chart;
30 collected_number d = (collected_number)dim;
31 collected_number p = (collected_number)point;
32
33 return (r * CHARTS * DIMS * POINTS_PER_REGION +
34 c * DIMS * POINTS_PER_REGION +
35 d * POINTS_PER_REGION +
36 p) % 10000000;
37 }
38
39 static inline void storage_point_check(size_t region, size_t chart, size_t dim, size_t point, time_t now, time_t update_every, STORAGE_POINT sp, size_t *value_errors, size_t *time_errors, size_t *update_every_errors) {
40 // check the supplied STORAGE_POINT retrieved from the database
41 // against the computed timestamp, update_every and expected value
42
43 if(storage_point_is_gap(sp)) sp.min = sp.max = sp.sum = NAN;
44
45 collected_number expected = point_value_get(region, chart, dim, point);
46
47 if(roundndd(expected) != roundndd(sp.sum)) {
48 if(*value_errors < DIMS * 2) {
49 fprintf(stderr, " >>> DBENGINE: VALUE DOES NOT MATCH: "
50 "region %zu, chart %zu, dimension %zu, point %zu, time %ld: "
51 "expected %lld, found %f\n",
52 region, chart, dim, point, now, expected, sp.sum);
53 }
54
55 (*value_errors)++;
56 }
57
58 if(sp.start_time_s > now || sp.end_time_s < now) {
59 if(*time_errors < DIMS * 2) {
60 fprintf(stderr, " >>> DBENGINE: TIMESTAMP DOES NOT MATCH: "
61 "region %zu, chart %zu, dimension %zu, point %zu, timestamp %ld: "
62 "expected %ld, found %ld - %ld\n",
63 region, chart, dim, point, now, now, sp.start_time_s, sp.end_time_s);
64 }
65
66 (*time_errors)++;
67 }
68
69 if(update_every != sp.end_time_s - sp.start_time_s) {
70 if(*update_every_errors < DIMS * 2) {
71 fprintf(stderr, " >>> DBENGINE: UPDATE EVERY DOES NOT MATCH: "
72 "region %zu, chart %zu, dimension %zu, point %zu, timestamp %ld: "
73 "expected %ld, found %ld\n",
74 region, chart, dim, point, now, update_every, sp.end_time_s - sp.start_time_s);
75 }
76
77 (*update_every_errors)++;
78 }
79 }
80
81 static inline void rrddim_set_by_pointer_fake_time(RRDDIM *rd, collected_number value, time_t now) {
82 rd->collector.last_collected_time.tv_sec = now;
83 rd->collector.last_collected_time.tv_usec = 0;
84 rrddim_set_collected_int(rd, value);
85 rrddim_set_updated(rd);
86
87 rd->collector.counter++;
88
89 collected_number v = (value >= 0) ? value : -value;
90 if(unlikely(v > rd->collector.collected.i.collected_value_max)) rrddim_set_collected_max_int(rd, v);
91 }
92
93 static RRDHOST *dbengine_rrdhost_find_or_create(char *name) {
94 /* We don't want to drop metrics when generating load,
95 * we prefer to block data generation itself */
96
97 SYSTEM_TZ tz = system_tz_get();
98 RRDHOST *host = rrdhost_find_or_create(
99 name,
100 name,
101 name,
102 os_type,
103 tz.timezone,
104 tz.abbrev_timezone,
105 tz.utc_offset,
106 program_name,
107 NETDATA_VERSION,
108 nd_profile.update_every,
109 default_rrd_history_entries,
110 RRD_DB_MODE_DBENGINE,
111 health_plugin_enabled(),
112 stream_send.enabled,
113 stream_send.parents.destination,
114 stream_send.api_key,
115 stream_send.send_charts_matching,
116 stream_receive.replication.enabled,
117 stream_receive.replication.period,
118 stream_receive.replication.step,
119 NULL,
120 0
121 );
122 system_tz_free(&tz);
123 return host;
124 }
125
126 static void test_dbengine_create_charts(RRDHOST *host, RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
127 int update_every) {
128 fprintf(stderr, "DBENGINE Creating Test Charts...\n");
129
130 int i, j;
131 char name[101];
132
133 for (i = 0 ; i < CHARTS ; ++i) {
134 snprintfz(name, sizeof(name) - 1, "dbengine-chart-%d", i);
135
136 // create the chart
137 st[i] = rrdset_create(host, "netdata", name, name, "netdata", NULL, "Unit Testing", "a value", "unittest",
138 NULL, 1, update_every, RRDSET_TYPE_LINE);
139 rrdset_flag_set(st[i], RRDSET_FLAG_DEBUG);
140 rrdset_flag_set(st[i], RRDSET_FLAG_STORE_FIRST);
141 for (j = 0 ; j < DIMS ; ++j) {
142 snprintfz(name, sizeof(name) - 1, "dim-%d", j);
143
144 rd[i][j] = rrddim_add(st[i], name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
145 }
146 }
147
148 // Initialize DB with the very first entries
149 for (i = 0 ; i < CHARTS ; ++i) {
150 for (j = 0 ; j < DIMS ; ++j) {
151 rd[i][j]->collector.last_collected_time.tv_sec =
152 st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = START_TIMESTAMP - 1;
153 rd[i][j]->collector.last_collected_time.tv_usec =
154 st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
155 }
156 }
157 for (i = 0 ; i < CHARTS ; ++i) {
158 st[i]->usec_since_last_update = USEC_PER_SEC;
159
160 for (j = 0; j < DIMS; ++j) {
161 rrddim_set_by_pointer_fake_time(rd[i][j], 69, START_TIMESTAMP); // set first value to 69
162 }
163
164 struct timeval now;
165 now_realtime_timeval(&now);
166 rrdset_timed_done(st[i], now, false);
167 }
168 // Flush pages for subsequent real values
169 for (i = 0 ; i < CHARTS ; ++i) {
170 for (j = 0; j < DIMS; ++j) {
171 rrdeng_store_metric_flush_current_page((rd[i][j])->tiers[0].sch);
172 }
173 }
174 }
175
176 static time_t test_dbengine_create_metrics(
177 RRDSET *st[CHARTS],
178 RRDDIM *rd[CHARTS][DIMS],
179 size_t current_region,
180 time_t time_start) {
181
182 time_t update_every = REGION_UPDATE_EVERY[current_region];
183 fprintf(stderr, "DBENGINE Single Region Write to "
184 "region %zu, from %ld to %ld, with update every %ld...\n",
185 current_region, time_start, time_start + POINTS_PER_REGION * update_every, update_every);
186
187 // for the database to save the metrics at the right time, we need to set
188 // the last data collection time to be just before the first data collection.
189 time_t time_now = time_start;
190 for (size_t c = 0 ; c < CHARTS ; ++c) {
191 for (size_t d = 0 ; d < DIMS ; ++d) {
192 storage_engine_store_change_collection_frequency(rd[c][d]->tiers[0].sch, (int)update_every);
193
194 // setting these timestamps, to the data collection time, prevents interpolation
195 // during data collection, so that our value will be written as-is to the
196 // database.
197
198 rd[c][d]->collector.last_collected_time.tv_sec =
199 st[c]->last_collected_time.tv_sec = st[c]->last_updated.tv_sec = time_now;
200
201 rd[c][d]->collector.last_collected_time.tv_usec =
202 st[c]->last_collected_time.tv_usec = st[c]->last_updated.tv_usec = 0;
203 }
204 }
205
206 // set the samples to the database
207 for (size_t p = 0; p < POINTS_PER_REGION ; ++p) {
208 for (size_t c = 0 ; c < CHARTS ; ++c) {
209 st[c]->usec_since_last_update = USEC_PER_SEC * update_every;
210
211 for (size_t d = 0; d < DIMS; ++d)
212 rrddim_set_by_pointer_fake_time(rd[c][d], point_value_get(current_region, c, d, p), time_now);
213
214 rrdset_timed_done(st[c], (struct timeval){ .tv_sec = time_now, .tv_usec = 0 }, false);
215 }
216
217 time_now += update_every;
218 }
219
220 return time_now;
221 }
222
223 // Checks the metric data for the given region, returns number of errors
224 static size_t test_dbengine_check_metrics(
225 RRDSET *st[CHARTS] __maybe_unused,
226 RRDDIM *rd[CHARTS][DIMS],
227 size_t current_region,
228 time_t time_start,
229 time_t time_end) {
230
231 time_t update_every = REGION_UPDATE_EVERY[current_region];
232 fprintf(stderr, "DBENGINE Single Region Read from "
233 "region %zu, from %ld to %ld, with update every %ld...\n",
234 current_region, time_start, time_end, update_every);
235
236 // initialize all queries
237 struct storage_engine_query_handle handles[CHARTS * DIMS] = { 0 };
238 for (size_t c = 0 ; c < CHARTS ; ++c) {
239 for (size_t d = 0; d < DIMS; ++d) {
240 storage_engine_query_init(rd[c][d]->tiers[0].seb,
241 rd[c][d]->tiers[0].smh,
242 &handles[c * DIMS + d],
243 time_start,
244 time_end,
245 STORAGE_PRIORITY_NORMAL);
246 }
247 }
248
249 // check the stored samples
250 size_t value_errors = 0, time_errors = 0, update_every_errors = 0;
251 time_t time_now = time_start;
252 for(size_t p = 0; p < POINTS_PER_REGION ;p++) {
253 for (size_t c = 0 ; c < CHARTS ; ++c) {
254 for (size_t d = 0; d < DIMS; ++d) {
255 STORAGE_POINT sp = storage_engine_query_next_metric(&handles[c * DIMS + d]);
256 storage_point_check(current_region, c, d, p, time_now, update_every, sp,
257 &value_errors, &time_errors, &update_every_errors);
258 }
259 }
260
261 time_now += update_every;
262 }
263
264 // finalize the queries
265 for (size_t c = 0 ; c < CHARTS ; ++c) {
266 for (size_t d = 0; d < DIMS; ++d) {
267 storage_engine_query_finalize(&handles[c * DIMS + d]);
268 }
269 }
270
271 if(value_errors)
272 fprintf(stderr, "%zu value errors encountered (out of %d checks)\n", value_errors, POINTS_PER_REGION * CHARTS * DIMS);
273
274 if(time_errors)
275 fprintf(stderr, "%zu time errors encountered (out of %d checks)\n", time_errors, POINTS_PER_REGION * CHARTS * DIMS);
276
277 if(update_every_errors)
278 fprintf(stderr, "%zu update every errors encountered (out of %d checks)\n", update_every_errors, POINTS_PER_REGION * CHARTS * DIMS);
279
280 return value_errors + time_errors + update_every_errors;
281 }
282
283 static size_t dbengine_test_rrdr_single_region(
284 RRDSET *st[CHARTS],
285 RRDDIM *rd[CHARTS][DIMS],
286 size_t current_region,
287 time_t time_start,
288 time_t time_end) {
289
290 time_t update_every = REGION_UPDATE_EVERY[current_region];
291 fprintf(stderr, "RRDR Single Region Test on "
292 "region %zu, start time %lld, end time %lld, update every %ld, on %d dimensions...\n",
293 current_region, (long long)time_start, (long long)time_end, update_every, CHARTS * DIMS);
294
295 size_t errors = 0, value_errors = 0, time_errors = 0, update_every_errors = 0;
296 long points = (time_end - time_start) / update_every;
297 for(size_t c = 0; c < CHARTS ;c++) {
298 ONEWAYALLOC *owa = onewayalloc_create(0);
299 RRDR *r = rrd2rrdr_legacy(owa, st[c], points, time_start, time_end,
300 RRDR_GROUPING_AVERAGE, 0, RRDR_OPTION_NATURAL_POINTS,
301 NULL, NULL, 0, 0,
302 QUERY_SOURCE_UNITTEST, STORAGE_PRIORITY_NORMAL);
303 if (!r) {
304 fprintf(stderr, " >>> DBENGINE: %s: empty RRDR on region %zu\n", rrdset_name(st[c]), current_region);
305 onewayalloc_destroy(owa);
306 errors++;
307 continue;
308 }
309
310 if(r->internal.qt->request.st != st[c])
311 fatal("queried wrong chart");
312
313 if(rrdr_rows(r) != POINTS_PER_REGION)
314 fatal("query returned wrong number of points (expected %d, got %zu)", POINTS_PER_REGION, rrdr_rows(r));
315
316 time_t time_now = time_start;
317 for (size_t p = 0; p < rrdr_rows(r); p++) {
318 size_t d = 0;
319 RRDDIM *dim;
320 rrddim_foreach_read(dim, r->internal.qt->request.st) {
321 if(unlikely(d >= r->d))
322 fatal("got more dimensions (%zu) than expected (%zu)", d, r->d);
323
324 if(rd[c][d] != dim)
325 fatal("queried wrong dimension");
326
327 RRDR_VALUE_FLAGS *co = &r->o[ p * r->d ];
328 NETDATA_DOUBLE *cn = &r->v[ p * r->d ];
329
330 STORAGE_POINT sp = STORAGE_POINT_UNSET;
331 sp.min = sp.max = sp.sum = (co[d] & RRDR_VALUE_EMPTY) ? NAN :cn[d];
332 sp.count = 1;
333 sp.end_time_s = r->t[p];
334 sp.start_time_s = sp.end_time_s - r->view.update_every;
335
336 storage_point_check(current_region, c, d, p, time_now, update_every, sp, &value_errors, &time_errors, &update_every_errors);
337 d++;
338 }
339 rrddim_foreach_done(dim);
340 time_now += update_every;
341 }
342
343 rrdr_free(owa, r);
344 onewayalloc_destroy(owa);
345 }
346
347 if(value_errors)
348 fprintf(stderr, "%zu value errors encountered (out of %d checks)\n", value_errors, POINTS_PER_REGION * CHARTS * DIMS);
349
350 if(time_errors)
351 fprintf(stderr, "%zu time errors encountered (out of %d checks)\n", time_errors, POINTS_PER_REGION * CHARTS * DIMS);
352
353 if(update_every_errors)
354 fprintf(stderr, "%zu update every errors encountered (out of %d checks)\n", update_every_errors, POINTS_PER_REGION * CHARTS * DIMS);
355
356 return errors + value_errors + time_errors + update_every_errors;
357 }
358
359 int test_dbengine(void) {
360 // provide enough threads to dbengine
361 setenv("UV_THREADPOOL_SIZE", "48", 1);
362
363 size_t errors = 0, value_errors = 0, time_errors = 0;
364
365 nd_log_limits_unlimited();
366 fprintf(stderr, "\nRunning DB-engine test\n");
367
368 default_rrd_memory_mode = RRD_DB_MODE_DBENGINE;
369 fprintf(stderr, "Initializing localhost with hostname 'unittest-dbengine'");
370 RRDHOST *host = dbengine_rrdhost_find_or_create("unittest-dbengine");
371 if(!host)
372 fatal("Failed to initialize host");
373
374 RRDSET *st[CHARTS] = { 0 };
375 RRDDIM *rd[CHARTS][DIMS] = { 0 };
376 time_t time_start[REGIONS] = { 0 }, time_end[REGIONS] = { 0 };
377
378 // create the charts and dimensions we need
379 test_dbengine_create_charts(host, st, rd, REGION_UPDATE_EVERY[0]);
380
381 time_t now = START_TIMESTAMP;
382 time_t update_every_old = REGION_UPDATE_EVERY[0];
383 for(size_t current_region = 0; current_region < REGIONS ;current_region++) {
384 time_t update_every = REGION_UPDATE_EVERY[current_region];
385
386 if(update_every != update_every_old) {
387 for (size_t c = 0 ; c < CHARTS ; ++c)
388 rrdset_set_update_every_s(st[c], update_every);
389 }
390
391 time_start[current_region] = region_start_time(now, update_every);
392 now = time_end[current_region] = test_dbengine_create_metrics(st,rd, current_region, time_start[current_region]);
393
394 errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region], time_end[current_region]);
395 }
396
397 // check everything again
398 for(size_t current_region = 0; current_region < REGIONS ;current_region++)
399 errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region], time_end[current_region]);
400
401 // check again in reverse order
402 for(size_t current_region = 0; current_region < REGIONS ;current_region++) {
403 size_t region = REGIONS - 1 - current_region;
404 errors += test_dbengine_check_metrics(st, rd, region, time_start[region], time_end[region]);
405 }
406
407 // check all the regions using RRDR
408 // this also checks the query planner and the query engine of Netdata
409 for (size_t current_region = 0 ; current_region < REGIONS ; current_region++) {
410 errors += dbengine_test_rrdr_single_region(st, rd, current_region, time_start[current_region], time_end[current_region]);
411 }
412
413 // prevent closing the database before the test is finished
414 sleep(5);
415
416 rrd_wrlock();
417 rrdeng_quiesce((struct rrdengine_instance *)host->db[0].si);
418 rrdeng_flush_all((struct rrdengine_instance *)host->db[0].si);
419 rrdeng_exit((struct rrdengine_instance *)host->db[0].si);
420 rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_BEST_EFFORT, NULL, NULL);
421 rrd_wrunlock();
422
423 return (int)(errors + value_errors + time_errors);
424 }
425
426 #endif