@cryptotaxi247 / netdata-1 / commits / 9adf4dd78

Configurable storage engine for Netdata agents: step 2 (#12808)

Adrien Béraud committed May 11, 2022 at 09:17 UTC 9adf4dd782b8587ca3daffeb46c6c3ea98fb76d8
10 files changed +277 -117
CMakeLists.txt
+5 -1
@@ -639,6 +639,10 @@ set(RRD_PLUGIN_FILES
639 database/rrdsetvar.h
640 database/rrdvar.c
641 database/rrdvar.h
642 + database/storage_engine.c
643 + database/storage_engine.h
644 + database/ram/rrddim_mem.c
645 + database/ram/rrddim_mem.h
646 database/sqlite/sqlite_functions.c
647 database/sqlite/sqlite_functions.h
648 database/sqlite/sqlite_aclk.c
@@ -1027,7 +1031,7 @@ IF(ENABLE_EXPORTING_PROMETHEUS_REMOTE_WRITE)
1031 message(STATUS "prometheus remote write exporting: enabled")
1032
1033 find_package(Protobuf REQUIRED)
1030 -
1034 +
1035 function(PROTOBUF_REMOTE_WRITE_GENERATE_CPP SRCS HDRS)
1036 if(NOT ARGN)
1037 message(SEND_ERROR "Error: PROTOBUF_REMOTE_WRITE_GENERATE_CPP() called without any proto files")
Makefile.am
+4
@@ -448,6 +448,10 @@ RRD_PLUGIN_FILES = \
448 database/rrdsetvar.h \
449 database/rrdvar.c \
450 database/rrdvar.h \
451 + database/storage_engine.c \
452 + database/storage_engine.h \
453 + database/ram/rrddim_mem.c \
454 + database/ram/rrddim_mem.h \
455 database/sqlite/sqlite_functions.c \
456 database/sqlite/sqlite_functions.h \
457 database/sqlite/sqlite_aclk.c \
database/ram/rrddim_mem.c new
+67
@@ -0,0 +1,67 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "rrddim_mem.h"
4 +
5 +// ----------------------------------------------------------------------------
6 +// RRDDIM legacy data collection functions
7 +
8 +void rrddim_collect_init(RRDDIM *rd) {
9 + rd->values[rd->rrdset->current_entry] = SN_EMPTY_SLOT;
10 + rd->state->handle = calloc(1, sizeof(struct mem_collect_handle));
11 +}
12 +void rrddim_collect_store_metric(RRDDIM *rd, usec_t point_in_time, storage_number number) {
13 + (void)point_in_time;
14 + rd->values[rd->rrdset->current_entry] = number;
15 +}
16 +int rrddim_collect_finalize(RRDDIM *rd) {
17 + free((struct mem_collect_handle*)rd->state->handle);
18 + return 0;
19 +}
20 +
21 +// ----------------------------------------------------------------------------
22 +// RRDDIM legacy database query functions
23 +
24 +void rrddim_query_init(RRDDIM *rd, struct rrddim_query_handle *handle, time_t start_time, time_t end_time) {
25 + handle->rd = rd;
26 + handle->start_time = start_time;
27 + handle->end_time = end_time;
28 + struct mem_query_handle* h = calloc(1, sizeof(struct mem_query_handle));
29 + h->slot = rrdset_time2slot(rd->rrdset, start_time);
30 + h->last_slot = rrdset_time2slot(rd->rrdset, end_time);
31 + h->finished = 0;
32 + handle->handle = (STORAGE_QUERY_HANDLE *)h;
33 +}
34 +
35 +storage_number rrddim_query_next_metric(struct rrddim_query_handle *handle, time_t *current_time) {
36 + RRDDIM *rd = handle->rd;
37 + struct mem_query_handle* h = (struct mem_query_handle*)handle->handle;
38 + long entries = rd->rrdset->entries;
39 + long slot = h->slot;
40 +
41 + (void)current_time;
42 + if (unlikely(h->slot == h->last_slot))
43 + h->finished = 1;
44 + storage_number n = rd->values[slot++];
45 +
46 + if(unlikely(slot >= entries)) slot = 0;
47 + h->slot = slot;
48 +
49 + return n;
50 +}
51 +
52 +int rrddim_query_is_finished(struct rrddim_query_handle *handle) {
53 + struct mem_query_handle* h = (struct mem_query_handle*)handle->handle;
54 + return h->finished;
55 +}
56 +
57 +void rrddim_query_finalize(struct rrddim_query_handle *handle) {
58 + freez(handle->handle);
59 +}
60 +
61 +time_t rrddim_query_latest_time(RRDDIM *rd) {
62 + return rrdset_last_entry_t_nolock(rd->rrdset);
63 +}
64 +
65 +time_t rrddim_query_oldest_time(RRDDIM *rd) {
66 + return rrdset_first_entry_t_nolock(rd->rrdset);
67 +}
database/ram/rrddim_mem.h new
+29
@@ -0,0 +1,29 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_RRDDIMMEM_H
4 +#define NETDATA_RRDDIMMEM_H
5 +
6 +#include "database/rrd.h"
7 +
8 +struct mem_collect_handle {
9 + long slot;
10 + long entries;
11 +};
12 +struct mem_query_handle {
13 + long slot;
14 + long last_slot;
15 + uint8_t finished;
16 +};
17 +
18 +extern void rrddim_collect_init(RRDDIM *rd);
19 +extern void rrddim_collect_store_metric(RRDDIM *rd, usec_t point_in_time, storage_number number);
20 +extern int rrddim_collect_finalize(RRDDIM *rd);
21 +
22 +extern void rrddim_query_init(RRDDIM *rd, struct rrddim_query_handle *handle, time_t start_time, time_t end_time);
23 +extern storage_number rrddim_query_next_metric(struct rrddim_query_handle *handle, time_t *current_time);
24 +extern int rrddim_query_is_finished(struct rrddim_query_handle *handle);
25 +extern void rrddim_query_finalize(struct rrddim_query_handle *handle);
26 +extern time_t rrddim_query_latest_time(RRDDIM *rd);
27 +extern time_t rrddim_query_oldest_time(RRDDIM *rd);
28 +
29 +#endif
database/rrd.c
+10 -14
@@ -2,6 +2,7 @@
2 #define NETDATA_RRD_INTERNALS 1
3
4 #include "rrd.h"
5 +#include "storage_engine.h"
6
7 // ----------------------------------------------------------------------------
8 // globals
@@ -47,24 +48,19 @@ inline const char *rrd_memory_mode_name(RRD_MEMORY_MODE id) {
48 return RRD_MEMORY_MODE_DBENGINE_NAME;
49 }
50
51 + STORAGE_ENGINE* eng = storage_engine_get(id);
52 + if (eng) {
53 + return eng->name;
54 + }
55 +
56 return RRD_MEMORY_MODE_SAVE_NAME;
57 }
58
59 RRD_MEMORY_MODE rrd_memory_mode_id(const char *name) {
54 - if(unlikely(!strcmp(name, RRD_MEMORY_MODE_RAM_NAME)))
55 - return RRD_MEMORY_MODE_RAM;
56 -
57 - else if(unlikely(!strcmp(name, RRD_MEMORY_MODE_MAP_NAME)))
58 - return RRD_MEMORY_MODE_MAP;
59 -
60 - else if(unlikely(!strcmp(name, RRD_MEMORY_MODE_NONE_NAME)))
61 - return RRD_MEMORY_MODE_NONE;
62 -
63 - else if(unlikely(!strcmp(name, RRD_MEMORY_MODE_ALLOC_NAME)))
64 - return RRD_MEMORY_MODE_ALLOC;
65 -
66 - else if(unlikely(!strcmp(name, RRD_MEMORY_MODE_DBENGINE_NAME)))
67 - return RRD_MEMORY_MODE_DBENGINE;
60 + STORAGE_ENGINE* eng = storage_engine_find(name);
61 + if (eng) {
62 + return eng->id;
63 + }
64
65 return RRD_MEMORY_MODE_SAVE;
66 }
database/rrd.h
-13
@@ -405,19 +405,6 @@ struct rrdset_volatile {
405 bool is_ar_chart;
406 };
407
408 -// RRDDIM legacy data collection structures
409 -
410 -struct mem_collect_handle {
411 - long slot;
412 - long entries;
413 -};
414 -
415 -struct mem_query_handle {
416 - long slot;
417 - long last_slot;
418 - uint8_t finished;
419 -};
420 -
408 // ----------------------------------------------------------------------------
409 // these loop macros make sure the linked list is accessed with the right lock
410
database/rrddim.c
+11 -89
@@ -2,6 +2,10 @@
2
3 #define NETDATA_RRD_INTERNALS
4 #include "rrd.h"
5 +#ifdef ENABLE_DBENGINE
6 +#include "database/engine/rrdengineapi.h"
7 +#endif
8 +#include "storage_engine.h"
9
10 // ----------------------------------------------------------------------------
11 // RRDDIM index
@@ -96,74 +100,6 @@ inline int rrddim_set_divisor(RRDSET *st, RRDDIM *rd, collected_number divisor)
100 return 1;
101 }
102
99 -// ----------------------------------------------------------------------------
100 -// RRDDIM legacy data collection functions
101 -
102 -static void rrddim_collect_init(RRDDIM *rd) {
103 - rd->state->handle = callocz(1, sizeof(struct mem_collect_handle));
104 - rd->values[rd->rrdset->current_entry] = SN_EMPTY_SLOT;
105 -}
106 -static void rrddim_collect_store_metric(RRDDIM *rd, usec_t point_in_time, storage_number number) {
107 - (void)point_in_time;
108 -
109 - rd->values[rd->rrdset->current_entry] = number;
110 -}
111 -static int rrddim_collect_finalize(RRDDIM *rd) {
112 - freez(rd->state->handle);
113 - return 0;
114 -}
115 -
116 -
117 -// ----------------------------------------------------------------------------
118 -// RRDDIM legacy database query functions
119 -
120 -static void rrddim_query_init(RRDDIM *rd, struct rrddim_query_handle *handle, time_t start_time, time_t end_time) {
121 - handle->rd = rd;
122 - handle->start_time = start_time;
123 - handle->end_time = end_time;
124 - struct mem_query_handle* mem_handle = callocz(1, sizeof(struct mem_query_handle));
125 - mem_handle->slot = rrdset_time2slot(rd->rrdset, start_time);
126 - mem_handle->last_slot = rrdset_time2slot(rd->rrdset, end_time);
127 - mem_handle->finished = 0;
128 - handle->handle = (STORAGE_QUERY_HANDLE *)mem_handle;
129 -}
130 -
131 -static storage_number rrddim_query_next_metric(struct rrddim_query_handle *handle, time_t *current_time) {
132 - RRDDIM *rd = handle->rd;
133 - struct mem_query_handle* mem_handle = (struct mem_query_handle*)handle->handle;
134 - long entries = rd->rrdset->entries;
135 - long slot = mem_handle->slot;
136 -
137 - (void)current_time;
138 - if (unlikely(mem_handle->slot == mem_handle->last_slot))
139 - mem_handle->finished = 1;
140 - storage_number n = rd->values[slot++];
141 -
142 - if(unlikely(slot >= entries)) slot = 0;
143 - mem_handle->slot = slot;
144 -
145 - return n;
146 -}
147 -
148 -static int rrddim_query_is_finished(struct rrddim_query_handle *handle) {
149 - struct mem_query_handle* mem_handle = (struct mem_query_handle*)handle->handle;
150 - return mem_handle->finished;
151 -}
152 -
153 -static void rrddim_query_finalize(struct rrddim_query_handle *handle) {
154 - freez(handle->handle);
155 - return;
156 -}
157 -
158 -static time_t rrddim_query_latest_time(RRDDIM *rd) {
159 - return rrdset_last_entry_t_nolock(rd->rrdset);
160 -}
161 -
162 -static time_t rrddim_query_oldest_time(RRDDIM *rd) {
163 - return rrdset_first_entry_t_nolock(rd->rrdset);
164 -}
165 -
166 -
103 // ----------------------------------------------------------------------------
104 // RRDDIM create a dimension
105
@@ -373,30 +309,16 @@ RRDDIM *rrddim_add_custom(RRDSET *st, const char *id, const char *name, collecte
309 rd->state->aclk_live_status = -1;
310 #endif
311 (void) find_dimension_uuid(st, rd, &(rd->state->metric_uuid));
376 - if(memory_mode == RRD_MEMORY_MODE_DBENGINE) {
312 +
313 + STORAGE_ENGINE* eng = storage_engine_get(memory_mode);
314 + rd->state->collect_ops = eng->api.collect_ops;
315 + rd->state->query_ops = eng->api.query_ops;
316 +
317 #ifdef ENABLE_DBENGINE
318 + if(memory_mode == RRD_MEMORY_MODE_DBENGINE) {
319 rrdeng_metric_init(rd);
379 - rd->state->collect_ops.init = rrdeng_store_metric_init;
380 - rd->state->collect_ops.store_metric = rrdeng_store_metric_next;
381 - rd->state->collect_ops.finalize = rrdeng_store_metric_finalize;
382 - rd->state->query_ops.init = rrdeng_load_metric_init;
383 - rd->state->query_ops.next_metric = rrdeng_load_metric_next;
384 - rd->state->query_ops.is_finished = rrdeng_load_metric_is_finished;
385 - rd->state->query_ops.finalize = rrdeng_load_metric_finalize;
386 - rd->state->query_ops.latest_time = rrdeng_metric_latest_time;
387 - rd->state->query_ops.oldest_time = rrdeng_metric_oldest_time;
388 -#endif
389 - } else {
390 - rd->state->collect_ops.init = rrddim_collect_init;
391 - rd->state->collect_ops.store_metric = rrddim_collect_store_metric;
392 - rd->state->collect_ops.finalize = rrddim_collect_finalize;
393 - rd->state->query_ops.init = rrddim_query_init;
394 - rd->state->query_ops.next_metric = rrddim_query_next_metric;
395 - rd->state->query_ops.is_finished = rrddim_query_is_finished;
396 - rd->state->query_ops.finalize = rrddim_query_finalize;
397 - rd->state->query_ops.latest_time = rrddim_query_latest_time;
398 - rd->state->query_ops.oldest_time = rrddim_query_oldest_time;
320 }
321 +#endif
322 store_active_dimension(&rd->state->metric_uuid);
323 rd->state->collect_ops.init(rd);
324 // append this dimension
database/storage_engine.c new
+120
@@ -0,0 +1,120 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "storage_engine.h"
4 +#include "ram/rrddim_mem.h"
5 +#ifdef ENABLE_DBENGINE
6 +#include "engine/rrdengineapi.h"
7 +#endif
8 +
9 +#define im_collect_ops { \
10 + .init = rrddim_collect_init,\
11 + .store_metric = rrddim_collect_store_metric,\
12 + .finalize = rrddim_collect_finalize\
13 +}
14 +
15 +#define im_query_ops { \
16 + .init = rrddim_query_init, \
17 + .next_metric = rrddim_query_next_metric, \
18 + .is_finished = rrddim_query_is_finished, \
19 + .finalize = rrddim_query_finalize, \
20 + .latest_time = rrddim_query_latest_time, \
21 + .oldest_time = rrddim_query_oldest_time \
22 +}
23 +
24 +static STORAGE_ENGINE engines[] = {
25 + {
26 + .id = RRD_MEMORY_MODE_NONE,
27 + .name = RRD_MEMORY_MODE_NONE_NAME,
28 + .api = {
29 + .collect_ops = im_collect_ops,
30 + .query_ops = im_query_ops
31 + }
32 + },
33 + {
34 + .id = RRD_MEMORY_MODE_RAM,
35 + .name = RRD_MEMORY_MODE_RAM_NAME,
36 + .api = {
37 + .collect_ops = im_collect_ops,
38 + .query_ops = im_query_ops
39 + }
40 + },
41 + {
42 + .id = RRD_MEMORY_MODE_MAP,
43 + .name = RRD_MEMORY_MODE_MAP_NAME,
44 + .api = {
45 + .collect_ops = im_collect_ops,
46 + .query_ops = im_query_ops
47 + }
48 + },
49 + {
50 + .id = RRD_MEMORY_MODE_SAVE,
51 + .name = RRD_MEMORY_MODE_SAVE_NAME,
52 + .api = {
53 + .collect_ops = im_collect_ops,
54 + .query_ops = im_query_ops
55 + }
56 + },
57 + {
58 + .id = RRD_MEMORY_MODE_ALLOC,
59 + .name = RRD_MEMORY_MODE_ALLOC_NAME,
60 + .api = {
61 + .collect_ops = im_collect_ops,
62 + .query_ops = im_query_ops
63 + }
64 + },
65 +#ifdef ENABLE_DBENGINE
66 + {
67 + .id = RRD_MEMORY_MODE_DBENGINE,
68 + .name = RRD_MEMORY_MODE_DBENGINE_NAME,
69 + .api = {
70 + .collect_ops = {
71 + .init = rrdeng_store_metric_init,
72 + .store_metric = rrdeng_store_metric_next,
73 + .finalize = rrdeng_store_metric_finalize
74 + },
75 + .query_ops = {
76 + .init = rrdeng_load_metric_init,
77 + .next_metric = rrdeng_load_metric_next,
78 + .is_finished = rrdeng_load_metric_is_finished,
79 + .finalize = rrdeng_load_metric_finalize,
80 + .latest_time = rrdeng_metric_latest_time,
81 + .oldest_time = rrdeng_metric_oldest_time
82 + }
83 + }
84 + },
85 +#endif
86 + { .id = RRD_MEMORY_MODE_NONE, .name = NULL }
87 +};
88 +
89 +STORAGE_ENGINE* storage_engine_find(const char* name)
90 +{
91 + for (STORAGE_ENGINE* it = engines; it->name; it++) {
92 + if (strcmp(it->name, name) == 0)
93 + return it;
94 + }
95 + return NULL;
96 +}
97 +
98 +STORAGE_ENGINE* storage_engine_get(RRD_MEMORY_MODE mmode)
99 +{
100 + for (STORAGE_ENGINE* it = engines; it->name; it++) {
101 + if (it->id == mmode)
102 + return it;
103 + }
104 + return NULL;
105 +}
106 +
107 +STORAGE_ENGINE* storage_engine_foreach_init()
108 +{
109 + // Assuming at least one engine exists
110 + return &engines[0];
111 +}
112 +
113 +STORAGE_ENGINE* storage_engine_foreach_next(STORAGE_ENGINE* it)
114 +{
115 + if (!it || !it->name)
116 + return NULL;
117 +
118 + it++;
119 + return it->name ? it : NULL;
120 +}
database/storage_engine.h new
+30
@@ -0,0 +1,30 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_STORAGEENGINEAPI_H
4 +#define NETDATA_STORAGEENGINEAPI_H
5 +
6 +#include "rrd.h"
7 +
8 +typedef struct storage_engine STORAGE_ENGINE;
9 +
10 +// ------------------------------------------------------------------------
11 +// function pointers for all APIs provided by a storge engine
12 +typedef struct storage_engine_api {
13 + struct rrddim_collect_ops collect_ops;
14 + struct rrddim_query_ops query_ops;
15 +} STORAGE_ENGINE_API;
16 +
17 +struct storage_engine {
18 + RRD_MEMORY_MODE id;
19 + const char* name;
20 + STORAGE_ENGINE_API api;
21 +};
22 +
23 +extern STORAGE_ENGINE* storage_engine_get(RRD_MEMORY_MODE mmode);
24 +extern STORAGE_ENGINE* storage_engine_find(const char* name);
25 +
26 +// Iterator over existing engines
27 +extern STORAGE_ENGINE* storage_engine_foreach_init();
28 +extern STORAGE_ENGINE* storage_engine_foreach_next(STORAGE_ENGINE* it);
29 +
30 +#endif
web/api/queries/query.c
+1
@@ -3,6 +3,7 @@
3 #include "query.h"
4 #include "web/api/formatters/rrd2json.h"
5 #include "rrdr.h"
6 +#include "database/ram/rrddim_mem.h"
7
8 #include "average/average.h"
9 #include "incremental_sum/incremental_sum.h"