New Cloud chart related parsers and generators (#11393)
* adds message generators parsers and handlers for upcoming Chart stream implementation
Timotej S committed
Aug 6, 2021 at 13:50 UTC
14ce65525275fa8760fa4f20e10b5a4e5c34ccde
19 files changed
+861
-40
CMakeLists.txt
+1
@@ -805,6 +805,7 @@ set(ACLK_NG_FILES
805
aclk/schema-wrappers/node_creation.cc
806
aclk/schema-wrappers/node_creation.h
807
aclk/schema-wrappers/schema_wrappers.h
808
+ aclk/schema-wrappers/schema_wrapper_utils.cc
809
aclk/schema-wrappers/schema_wrapper_utils.h
810
)
811
Makefile.am
+38
@@ -567,6 +567,8 @@ ACLK_NG_FILES = \
567
aclk/aclk_rx_msgs.h \
568
aclk/https_client.c \
569
aclk/https_client.h \
570
+ aclk/aclk_charts_api.c \
571
+ aclk/aclk_charts_api.h \
572
mqtt_websockets/src/mqtt_wss_client.c \
573
mqtt_websockets/src/include/mqtt_wss_client.h \
574
mqtt_websockets/src/mqtt_wss_log.c \
@@ -584,7 +586,12 @@ ACLK_NG_FILES = \
586
aclk/schema-wrappers/node_connection.h \
587
aclk/schema-wrappers/node_creation.cc \
588
aclk/schema-wrappers/node_creation.h \
589
+ aclk/schema-wrappers/chart_stream.cc \
590
+ aclk/schema-wrappers/chart_stream.h \
591
+ aclk/schema-wrappers/chart_config.cc \
592
+ aclk/schema-wrappers/chart_config.h \
593
aclk/schema-wrappers/schema_wrappers.h \
594
+ aclk/schema-wrappers/schema_wrapper_utils.cc \
595
aclk/schema-wrappers/schema_wrapper_utils.h \
596
$(NULL)
597
@@ -594,10 +601,21 @@ ACLK_NG_PROTO_BUILT_FILES = aclk/aclk-schemas/proto/agent/v1/connection.pb.cc \
601
aclk/aclk-schemas/proto/nodeinstance/connection/v1/connection.pb.h \
602
aclk/aclk-schemas/proto/nodeinstance/create/v1/creation.pb.cc \
603
aclk/aclk-schemas/proto/nodeinstance/create/v1/creation.pb.h \
604
+ aclk/aclk-schemas/proto/chart/v1/stream.pb.cc \
605
+ aclk/aclk-schemas/proto/chart/v1/stream.pb.h \
606
+ aclk/aclk-schemas/proto/chart/v1/instance.pb.cc \
607
+ aclk/aclk-schemas/proto/chart/v1/instance.pb.h \
608
+ aclk/aclk-schemas/proto/chart/v1/dimension.pb.cc \
609
+ aclk/aclk-schemas/proto/chart/v1/dimension.pb.h \
610
+ aclk/aclk-schemas/proto/chart/v1/config.pb.cc \
611
+ aclk/aclk-schemas/proto/chart/v1/config.pb.h \
612
+ aclk/aclk-schemas/proto/aclk/v1/lib.pb.cc \
613
+ aclk/aclk-schemas/proto/aclk/v1/lib.pb.h \
614
$(NULL)
615
616
BUILT_SOURCES += $(ACLK_NG_PROTO_BUILT_FILES)
617
nodist_netdata_SOURCES += $(ACLK_NG_PROTO_BUILT_FILES)
618
+CLEANFILES += $(ACLK_NG_PROTO_BUILT_FILES)
619
620
aclk/aclk-schemas/proto/agent/v1/connection.pb.cc \
621
aclk/aclk-schemas/proto/agent/v1/connection.pb.h: aclk/aclk-schemas/proto/agent/v1/connection.proto
@@ -611,6 +629,26 @@ aclk/aclk-schemas/proto/nodeinstance/create/v1/creation.pb.cc \
629
aclk/aclk-schemas/proto/nodeinstance/create/v1/creation.pb.h: aclk/aclk-schemas/proto/nodeinstance/create/v1/creation.proto
630
$(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
631
632
+aclk/aclk-schemas/proto/chart/v1/stream.pb.cc \
633
+aclk/aclk-schemas/proto/chart/v1/stream.pb.h: aclk/aclk-schemas/proto/chart/v1/stream.proto
634
+ $(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
635
+
636
+aclk/aclk-schemas/proto/chart/v1/instance.pb.cc \
637
+aclk/aclk-schemas/proto/chart/v1/instance.pb.h: aclk/aclk-schemas/proto/chart/v1/instance.proto
638
+ $(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
639
+
640
+aclk/aclk-schemas/proto/chart/v1/dimension.pb.cc \
641
+aclk/aclk-schemas/proto/chart/v1/dimension.pb.h: aclk/aclk-schemas/proto/chart/v1/dimension.proto
642
+ $(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
643
+
644
+aclk/aclk-schemas/proto/chart/v1/config.pb.cc \
645
+aclk/aclk-schemas/proto/chart/v1/config.pb.h: aclk/aclk-schemas/proto/chart/v1/config.proto
646
+ $(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
647
+
648
+aclk/aclk-schemas/proto/aclk/v1/lib.pb.cc \
649
+aclk/aclk-schemas/proto/aclk/v1/lib.pb.h: aclk/aclk-schemas/proto/aclk/v1/lib.proto
650
+ $(PROTOC) -I=aclk/aclk-schemas --cpp_out=$(builddir)/aclk/aclk-schemas $^
651
+
652
endif #ACLK_NG
653
654
if ENABLE_ACLK
aclk/aclk-schemas
+1
-1
@@ -1 +1 @@
1
-Subproject commit b5fef3f3a84e6a5013b36b906f4677012c734416
1
+Subproject commit a0adf5b1e026ee8339d56cfa27af95bb26b53177
aclk/aclk_charts_api.c
new
+61
@@ -0,0 +1,61 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+#include "aclk_charts_api.h"
3
+
4
+#include "aclk_query_queue.h"
5
+
6
+#define CHART_DIM_UPDATE_NAME "ChartsAndDimensionsUpdated"
7
+
8
+void aclk_chart_inst_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
9
+{
10
+ aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
11
+ query->data.bin_payload.payload = generate_charts_updated(&query->data.bin_payload.size, payloads, payload_sizes, new_positions);
12
+ query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
13
+ if (query->data.bin_payload.payload)
14
+ aclk_queue_query(query);
15
+}
16
+
17
+void aclk_chart_dim_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
18
+{
19
+ aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
20
+ query->data.bin_payload.payload = generate_chart_dimensions_updated(&query->data.bin_payload.size, payloads, payload_sizes, new_positions);
21
+ query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
22
+ if (query->data.bin_payload.payload)
23
+ aclk_queue_query(query);
24
+}
25
+
26
+void aclk_chart_inst_and_dim_update(char **payloads, size_t *payload_sizes, int *is_dim, struct aclk_message_position *new_positions, uint64_t batch_id)
27
+{
28
+ aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
29
+ query->data.bin_payload.payload = generate_charts_and_dimensions_updated(&query->data.bin_payload.size, payloads, payload_sizes, is_dim, new_positions, batch_id);
30
+ query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
31
+ if (query->data.bin_payload.payload)
32
+ aclk_queue_query(query);
33
+}
34
+
35
+void aclk_chart_config_updated(struct chart_config_updated *config_list, int list_size)
36
+{
37
+ aclk_query_t query = aclk_query_new(CHART_CONFIG_UPDATED);
38
+ query->data.bin_payload.payload = generate_chart_configs_updated(&query->data.bin_payload.size, config_list, list_size);
39
+ query->data.bin_payload.msg_name = "ChartConfigsUpdated";
40
+ if (query->data.bin_payload.payload)
41
+ aclk_queue_query(query);
42
+}
43
+
44
+void aclk_chart_reset(chart_reset_t reset)
45
+{
46
+ aclk_query_t query = aclk_query_new(CHART_RESET);
47
+ query->data.bin_payload.payload = generate_reset_chart_messages(&query->data.bin_payload.size, reset);
48
+ query->data.bin_payload.msg_name = "ResetChartMessages";
49
+ if (query->data.bin_payload.payload)
50
+ aclk_queue_query(query);
51
+}
52
+
53
+void aclk_retention_updated(struct retention_updated *data)
54
+{
55
+ aclk_query_t query = aclk_query_new(RETENTION_UPDATED);
56
+ query->data.bin_payload.topic = ACLK_TOPICID_RETENTION_UPDATED;
57
+ query->data.bin_payload.payload = generate_retention_updated(&query->data.bin_payload.size, data);
58
+ query->data.bin_payload.msg_name = "RetentionUpdated";
59
+ if (query->data.bin_payload.payload)
60
+ aclk_queue_query(query);
61
+}
aclk/aclk_charts_api.h
new
+18
@@ -0,0 +1,18 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+#ifndef ACLK_CHARTS_H
3
+#define ACLK_CHARTS_H
4
+
5
+#include "../daemon/common.h"
6
+#include "schema-wrappers/schema_wrappers.h"
7
+
8
+void aclk_chart_inst_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions);
9
+void aclk_chart_dim_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions);
10
+void aclk_chart_inst_and_dim_update(char **payloads, size_t *payload_sizes, int *is_dim, struct aclk_message_position *new_positions, uint64_t batch_id);
11
+
12
+void aclk_chart_config_updated(struct chart_config_updated *config_list, int list_size);
13
+
14
+void aclk_chart_reset(chart_reset_t reset);
15
+
16
+void aclk_retention_updated(struct retention_updated *data);
17
+
18
+#endif /* ACLK_CHARTS_H */
aclk/aclk_query.c
+20
-9
@@ -273,16 +273,27 @@ static int node_state_update(struct aclk_query_thread *query_thr, aclk_query_t q
273
return 0;
274
}
275
276
+static int send_bin_msg(struct aclk_query_thread *query_thr, aclk_query_t query)
277
+{
278
+ // this will be simplified when legacy support is removed
279
+ aclk_send_bin_message_subtopic_pid(query_thr->client, query->data.bin_payload.payload, query->data.bin_payload.size, query->data.bin_payload.topic, query->data.bin_payload.msg_name);
280
+ return 0;
281
+}
282
+
283
aclk_query_handler aclk_query_handlers[] = {
277
- { .type = HTTP_API_V2, .name = "http api request v2", .fnc = http_api_v2 },
278
- { .type = ALARM_STATE_UPDATE, .name = "alarm state update", .fnc = alarm_state_update_query },
279
- { .type = METADATA_INFO, .name = "info metadata", .fnc = info_metadata },
280
- { .type = METADATA_ALARMS, .name = "alarms metadata", .fnc = alarms_metadata },
281
- { .type = CHART_NEW, .name = "chart new", .fnc = chart_query },
282
- { .type = CHART_DEL, .name = "chart delete", .fnc = info_metadata },
283
- { .type = REGISTER_NODE, .name = "register node", .fnc = register_node },
284
- { .type = NODE_STATE_UPDATE, .name = "node state update", .fnc = node_state_update },
285
- { .type = UNKNOWN, .name = NULL, .fnc = NULL }
284
+ { .type = HTTP_API_V2, .name = "http api request v2", .fnc = http_api_v2 },
285
+ { .type = ALARM_STATE_UPDATE, .name = "alarm state update", .fnc = alarm_state_update_query },
286
+ { .type = METADATA_INFO, .name = "info metadata", .fnc = info_metadata },
287
+ { .type = METADATA_ALARMS, .name = "alarms metadata", .fnc = alarms_metadata },
288
+ { .type = CHART_NEW, .name = "chart new", .fnc = chart_query },
289
+ { .type = CHART_DEL, .name = "chart delete", .fnc = info_metadata },
290
+ { .type = REGISTER_NODE, .name = "register node", .fnc = register_node },
291
+ { .type = NODE_STATE_UPDATE, .name = "node state update", .fnc = node_state_update },
292
+ { .type = CHART_DIMS_UPDATE, .name = "chart and dim update bin", .fnc = send_bin_msg },
293
+ { .type = CHART_CONFIG_UPDATED, .name = "chart config updated", .fnc = send_bin_msg },
294
+ { .type = CHART_RESET, .name = "reset chart messages", .fnc = send_bin_msg },
295
+ { .type = RETENTION_UPDATED, .name = "update retention info", .fnc = send_bin_msg },
296
+ { .type = UNKNOWN, .name = NULL, .fnc = NULL }
297
};
298
299
aclk/aclk_query_queue.c
+25
-10
@@ -139,27 +139,42 @@ aclk_query_t aclk_query_new(aclk_query_type_t type)
139
140
void aclk_query_free(aclk_query_t query)
141
{
142
- if (query->type == HTTP_API_V2) {
142
+ switch (query->type) {
143
+ case HTTP_API_V2:
144
freez(query->data.http_api_v2.payload);
145
if (query->data.http_api_v2.query != query->dedup_id)
146
freez(query->data.http_api_v2.query);
146
- }
147
+ break;
148
148
- if (query->type == CHART_NEW)
149
+ case CHART_NEW:
150
freez(query->data.chart_add_del.chart_name);
150
-
151
- if (query->type == ALARM_STATE_UPDATE && query->data.alarm_update)
152
- json_object_put(query->data.alarm_update);
153
-
154
- if (query->type == NODE_STATE_UPDATE) {
151
+ break;
152
+
153
+ case ALARM_STATE_UPDATE:
154
+ if (query->data.alarm_update)
155
+ json_object_put(query->data.alarm_update);
156
+ break;
157
+
158
+ case NODE_STATE_UPDATE:
159
freez((void*)query->data.node_update.claim_id);
160
freez((void*)query->data.node_update.node_id);
157
- }
161
+ break;
162
159
- if (query->type == REGISTER_NODE) {
163
+ case REGISTER_NODE:
164
freez((void*)query->data.node_creation.claim_id);
165
freez((void*)query->data.node_creation.hostname);
166
freez((void*)query->data.node_creation.machine_guid);
167
+ break;
168
+
169
+ case CHART_DIMS_UPDATE:
170
+ case CHART_CONFIG_UPDATED:
171
+ case CHART_RESET:
172
+ case RETENTION_UPDATED:
173
+ freez(query->data.bin_payload.payload);
174
+ break;
175
+
176
+ default:
177
+ break;
178
}
179
180
freez(query->dedup_id);
aclk/aclk_query_queue.h
+15
-1
@@ -7,6 +7,8 @@
7
#include "daemon/common.h"
8
#include "schema-wrappers/schema_wrappers.h"
9
10
+#include "aclk_util.h"
11
+
12
typedef enum {
13
UNKNOWN,
14
METADATA_INFO,
@@ -16,7 +18,11 @@ typedef enum {
18
CHART_DEL,
19
ALARM_STATE_UPDATE,
20
REGISTER_NODE,
19
- NODE_STATE_UPDATE
21
+ NODE_STATE_UPDATE,
22
+ CHART_DIMS_UPDATE,
23
+ CHART_CONFIG_UPDATED,
24
+ CHART_RESET,
25
+ RETENTION_UPDATED
26
} aclk_query_type_t;
27
28
struct aclk_query_metadata {
@@ -34,6 +40,13 @@ struct aclk_query_http_api_v2 {
40
char *query;
41
};
42
43
+struct aclk_bin_payload {
44
+ char *payload;
45
+ size_t size;
46
+ enum aclk_topics topic;
47
+ const char *msg_name;
48
+};
49
+
50
typedef struct aclk_query *aclk_query_t;
51
struct aclk_query {
52
aclk_query_type_t type;
@@ -61,6 +74,7 @@ struct aclk_query {
74
struct aclk_query_chart_add_del chart_add_del;
75
node_instance_creation_t node_creation;
76
node_instance_connection_t node_update;
77
+ struct aclk_bin_payload bin_payload;
78
json_object *alarm_update;
79
} data;
80
};
aclk/aclk_tx_msgs.c
+1
-1
@@ -36,7 +36,7 @@ static void aclk_send_message_subtopic(mqtt_wss_client client, json_object *msg,
36
#endif
37
}
38
39
-static uint16_t aclk_send_bin_message_subtopic_pid(mqtt_wss_client client, char *msg, size_t msg_len, enum aclk_topics subtopic, const char *msgname)
39
+uint16_t aclk_send_bin_message_subtopic_pid(mqtt_wss_client client, char *msg, size_t msg_len, enum aclk_topics subtopic, const char *msgname)
40
{
41
#ifndef ACLK_LOG_CONVERSATION_DIR
42
UNUSED(msgname);
aclk/aclk_tx_msgs.h
+3
@@ -7,6 +7,9 @@
7
#include "daemon/common.h"
8
#include "mqtt_wss_client.h"
9
#include "schema-wrappers/schema_wrappers.h"
10
+#include "aclk_util.h"
11
+
12
+uint16_t aclk_send_bin_message_subtopic_pid(mqtt_wss_client client, char *msg, size_t msg_len, enum aclk_topics subtopic, const char *msgname);
13
14
void aclk_send_info_metadata(mqtt_wss_client client, int metadata_submitted, RRDHOST *host);
15
void aclk_send_alarm_metadata(mqtt_wss_client client, int metadata_submitted);
aclk/aclk_util.c
+19
-9
@@ -13,6 +13,8 @@
13
int aclk_use_new_cloud_arch = 0;
14
usec_t aclk_session_newarch = 0;
15
16
+int chart_batch_id;
17
+
18
aclk_encoding_type_t aclk_encoding_type_t_from_str(const char *str) {
19
if (!strcmp(str, "json")) {
20
return ACLK_ENC_JSON;
@@ -110,15 +112,19 @@ struct topic_name {
112
// in answer to /password endpoint
113
const char *name;
114
} topic_names[] = {
113
- { .id = ACLK_TOPICID_CHART, .name = "chart" },
114
- { .id = ACLK_TOPICID_ALARMS, .name = "alarms" },
115
- { .id = ACLK_TOPICID_METADATA, .name = "meta" },
116
- { .id = ACLK_TOPICID_COMMAND, .name = "inbox-cmd" },
117
- { .id = ACLK_TOPICID_AGENT_CONN, .name = "agent-connection" },
118
- { .id = ACLK_TOPICID_CMD_NG_V1, .name = "inbox-cmd-v1" },
119
- { .id = ACLK_TOPICID_CREATE_NODE, .name = "create-node-instance" },
120
- { .id = ACLK_TOPICID_NODE_CONN, .name = "node-instance-connection" },
121
- { .id = ACLK_TOPICID_UNKNOWN, .name = NULL }
115
+ { .id = ACLK_TOPICID_CHART, .name = "chart" },
116
+ { .id = ACLK_TOPICID_ALARMS, .name = "alarms" },
117
+ { .id = ACLK_TOPICID_METADATA, .name = "meta" },
118
+ { .id = ACLK_TOPICID_COMMAND, .name = "inbox-cmd" },
119
+ { .id = ACLK_TOPICID_AGENT_CONN, .name = "agent-connection" },
120
+ { .id = ACLK_TOPICID_CMD_NG_V1, .name = "inbox-cmd-v1" },
121
+ { .id = ACLK_TOPICID_CREATE_NODE, .name = "create-node-instance" },
122
+ { .id = ACLK_TOPICID_NODE_CONN, .name = "node-instance-connection" },
123
+ { .id = ACLK_TOPICID_CHART_DIMS, .name = "chart-and-dims-updated" },
124
+ { .id = ACLK_TOPICID_CHART_CONFIGS_UPDATED, .name = "chart-configs-updated" },
125
+ { .id = ACLK_TOPICID_CHART_RESET, .name = "reset-charts" },
126
+ { .id = ACLK_TOPICID_RETENTION_UPDATED, .name = "chart-retention-updated" },
127
+ { .id = ACLK_TOPICID_UNKNOWN, .name = NULL }
128
};
129
130
enum aclk_topics compulsory_topics_legacy[] = {
@@ -139,6 +145,10 @@ enum aclk_topics compulsory_topics_new_cloud_arch[] = {
145
ACLK_TOPICID_CMD_NG_V1,
146
ACLK_TOPICID_CREATE_NODE,
147
ACLK_TOPICID_NODE_CONN,
148
+ ACLK_TOPICID_CHART_DIMS,
149
+ ACLK_TOPICID_CHART_CONFIGS_UPDATED,
150
+ ACLK_TOPICID_CHART_RESET,
151
+ ACLK_TOPICID_RETENTION_UPDATED,
152
ACLK_TOPICID_UNKNOWN
153
};
154
aclk/aclk_util.h
+15
-9
@@ -11,6 +11,8 @@
11
extern int aclk_use_new_cloud_arch;
12
extern usec_t aclk_session_newarch;
13
14
+extern int chart_batch_id;
15
+
16
typedef enum {
17
ACLK_ENC_UNKNOWN = 0,
18
ACLK_ENC_JSON,
@@ -54,15 +56,19 @@ void aclk_transport_desc_t_destroy(aclk_transport_desc_t *trp_desc);
56
void aclk_env_t_destroy(aclk_env_t *env);
57
58
enum aclk_topics {
57
- ACLK_TOPICID_UNKNOWN = 0,
58
- ACLK_TOPICID_CHART = 1,
59
- ACLK_TOPICID_ALARMS = 2,
60
- ACLK_TOPICID_METADATA = 3,
61
- ACLK_TOPICID_COMMAND = 4,
62
- ACLK_TOPICID_AGENT_CONN = 5,
63
- ACLK_TOPICID_CMD_NG_V1 = 6,
64
- ACLK_TOPICID_CREATE_NODE = 7,
65
- ACLK_TOPICID_NODE_CONN = 8
59
+ ACLK_TOPICID_UNKNOWN = 0,
60
+ ACLK_TOPICID_CHART = 1,
61
+ ACLK_TOPICID_ALARMS = 2,
62
+ ACLK_TOPICID_METADATA = 3,
63
+ ACLK_TOPICID_COMMAND = 4,
64
+ ACLK_TOPICID_AGENT_CONN = 5,
65
+ ACLK_TOPICID_CMD_NG_V1 = 6,
66
+ ACLK_TOPICID_CREATE_NODE = 7,
67
+ ACLK_TOPICID_NODE_CONN = 8,
68
+ ACLK_TOPICID_CHART_DIMS = 9,
69
+ ACLK_TOPICID_CHART_CONFIGS_UPDATED = 10,
70
+ ACLK_TOPICID_CHART_RESET = 11,
71
+ ACLK_TOPICID_RETENTION_UPDATED = 12
72
};
73
74
const char *aclk_get_topic(enum aclk_topics topic);
aclk/schema-wrappers/chart_config.cc
new
+105
@@ -0,0 +1,105 @@
1
+#include "chart_config.h"
2
+
3
+#include "proto/chart/v1/config.pb.h"
4
+
5
+#include "libnetdata/libnetdata.h"
6
+
7
+#include "schema_wrapper_utils.h"
8
+
9
+void destroy_update_chart_config(struct update_chart_config *cfg)
10
+{
11
+ freez(cfg->claim_id);
12
+ freez(cfg->node_id);
13
+ freez(cfg->hashes);
14
+}
15
+
16
+void destroy_chart_config_updated(struct chart_config_updated *cfg)
17
+{
18
+ freez(cfg->type);
19
+ freez(cfg->family);
20
+ freez(cfg->context);
21
+ freez(cfg->title);
22
+ freez(cfg->plugin);
23
+ freez(cfg->module);
24
+ freez(cfg->units);
25
+ freez(cfg->config_hash);
26
+}
27
+
28
+struct update_chart_config parse_update_chart_config(const char *data, size_t len)
29
+{
30
+ chart::v1::UpdateChartConfigs cfgs;
31
+ update_chart_config res;
32
+ memset(&res, 0, sizeof(res));
33
+
34
+ if (!cfgs.ParseFromArray(data, len))
35
+ return res;
36
+
37
+ res.claim_id = strdupz(cfgs.claim_id().c_str());
38
+ res.node_id = strdupz(cfgs.node_id().c_str());
39
+
40
+ // to not do bazillion tiny allocations for individual strings
41
+ // we calculate how much memory we will need for all of them
42
+ // and allocate at once
43
+ int hash_count = cfgs.config_hashes_size();
44
+ size_t total_strlen = 0;
45
+ for (int i = 0; i < hash_count; i++)
46
+ total_strlen += cfgs.config_hashes(i).length();
47
+ total_strlen += hash_count; //null bytes
48
+
49
+ res.hashes = (char**)callocz( 1,
50
+ (hash_count+1) * sizeof(char*) + //char * array incl. terminating NULL at the end
51
+ total_strlen //strings themselves incl. 1 null byte each
52
+ );
53
+
54
+ char* dest = ((char*)res.hashes) + (hash_count + 1 /* NULL ptr */) * sizeof(char *);
55
+ // now copy them strings
56
+ // null bytes handled by callocz
57
+ for (int i = 0; i < hash_count; i++) {
58
+ strcpy(dest, cfgs.config_hashes(i).c_str());
59
+ res.hashes[i] = dest;
60
+ dest += strlen(dest) + 1 /* end string null */;
61
+ }
62
+
63
+ return res;
64
+}
65
+
66
+char *generate_chart_configs_updated(size_t *len, const struct chart_config_updated *config_list, int list_size)
67
+{
68
+ chart::v1::ChartConfigsUpdated configs;
69
+ for (int i = 0; i < list_size; i++) {
70
+ chart::v1::ChartConfigUpdated *config = configs.add_configs();
71
+ config->set_type(config_list[i].type);
72
+ if (config_list[i].family)
73
+ config->set_family(config_list[i].family);
74
+ config->set_context(config_list[i].context);
75
+ config->set_title(config_list[i].title);
76
+ config->set_priority(config_list[i].priority);
77
+ config->set_plugin(config_list[i].plugin);
78
+
79
+ if (config_list[i].module)
80
+ config->set_module(config_list[i].module);
81
+
82
+ switch (config_list[i].chart_type) {
83
+ case RRDSET_TYPE_LINE:
84
+ config->set_chart_type(chart::v1::LINE);
85
+ break;
86
+ case RRDSET_TYPE_AREA:
87
+ config->set_chart_type(chart::v1::AREA);
88
+ break;
89
+ case RRDSET_TYPE_STACKED:
90
+ config->set_chart_type(chart::v1::STACKED);
91
+ break;
92
+ default:
93
+ return NULL;
94
+ }
95
+
96
+ config->set_units(config_list[i].units);
97
+ config->set_config_hash(config_list[i].config_hash);
98
+ }
99
+
100
+ *len = PROTO_COMPAT_MSG_SIZE(configs);
101
+ char *bin = (char*)mallocz(*len);
102
+ configs.SerializeToArray(bin, *len);
103
+
104
+ return bin;
105
+}
aclk/schema-wrappers/chart_config.h
new
+50
@@ -0,0 +1,50 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+
3
+#ifndef ACLK_SCHEMA_WRAPPER_CHART_CONFIG_H
4
+#define ACLK_SCHEMA_WRAPPER_CHART_CONFIG_H
5
+
6
+#include <stdlib.h>
7
+
8
+#include "database/rrd.h"
9
+
10
+#ifdef __cplusplus
11
+extern "C" {
12
+#endif
13
+
14
+struct update_chart_config {
15
+ char *claim_id;
16
+ char *node_id;
17
+ char **hashes;
18
+};
19
+
20
+enum chart_config_chart_type {
21
+ LINE,
22
+ AREA,
23
+ STACKED
24
+};
25
+
26
+struct chart_config_updated {
27
+ char *type;
28
+ char *family;
29
+ char *context;
30
+ char *title;
31
+ uint64_t priority;
32
+ char *plugin;
33
+ char *module;
34
+ RRDSET_TYPE chart_type;
35
+ char *units;
36
+ char *config_hash;
37
+};
38
+
39
+void destroy_update_chart_config(struct update_chart_config *cfg);
40
+void destroy_chart_config_updated(struct chart_config_updated *cfg);
41
+
42
+struct update_chart_config parse_update_chart_config(const char *data, size_t len);
43
+
44
+char *generate_chart_configs_updated(size_t *len, const struct chart_config_updated *config_list, int list_size);
45
+
46
+#ifdef __cplusplus
47
+}
48
+#endif
49
+
50
+#endif /* ACLK_SCHEMA_WRAPPER_CHART_CONFIG_H */
aclk/schema-wrappers/chart_stream.cc
new
+344
@@ -0,0 +1,344 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+
3
+#include "aclk/aclk_util.h"
4
+
5
+#include "proto/chart/v1/stream.pb.h"
6
+#include "chart_stream.h"
7
+
8
+#include "schema_wrapper_utils.h"
9
+
10
+#include <sys/time.h>
11
+#include <stdlib.h>
12
+
13
+stream_charts_and_dims_t parse_stream_charts_and_dims(const char *data, size_t len)
14
+{
15
+ chart::v1::StreamChartsAndDimensions msg;
16
+ stream_charts_and_dims_t res;
17
+ memset(&res, 0, sizeof(res));
18
+
19
+ if (!msg.ParseFromArray(data, len))
20
+ return res;
21
+
22
+ res.node_id = strdup(msg.node_id().c_str());
23
+ res.claim_id = strdup(msg.claim_id().c_str());
24
+ res.seq_id = msg.sequence_id();
25
+ res.batch_id = msg.batch_id();
26
+ set_timeval_from_google_timestamp(msg.seq_id_created_at(), &res.seq_id_created_at);
27
+
28
+ return res;
29
+}
30
+
31
+chart_and_dim_ack_t parse_chart_and_dimensions_ack(const char *data, size_t len)
32
+{
33
+ chart::v1::ChartsAndDimensionsAck msg;
34
+ chart_and_dim_ack_t res = { .claim_id = NULL, .node_id = NULL, .last_seq_id = 0 };
35
+
36
+ if (!msg.ParseFromArray(data, len))
37
+ return res;
38
+
39
+ res.node_id = strdup(msg.node_id().c_str());
40
+ res.claim_id = strdup(msg.claim_id().c_str());
41
+ res.last_seq_id = msg.last_sequence_id();
42
+
43
+ return res;
44
+}
45
+
46
+char *generate_reset_chart_messages(size_t *len, chart_reset_t reset)
47
+{
48
+ chart::v1::ResetChartMessages msg;
49
+
50
+ msg.set_claim_id(reset.claim_id);
51
+ msg.set_node_id(reset.node_id);
52
+ switch (reset.reason) {
53
+ case DB_EMPTY:
54
+ msg.set_reason(chart::v1::ResetReason::DB_EMPTY);
55
+ break;
56
+ case SEQ_ID_NOT_EXISTS:
57
+ msg.set_reason(chart::v1::ResetReason::SEQ_ID_NOT_EXISTS);
58
+ break;
59
+ case TIMESTAMP_MISMATCH:
60
+ msg.set_reason(chart::v1::ResetReason::TIMESTAMP_MISMATCH);
61
+ break;
62
+ default:
63
+ return NULL;
64
+ }
65
+
66
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
67
+ char *bin = (char*)malloc(*len);
68
+ if (bin)
69
+ msg.SerializeToArray(bin, *len);
70
+
71
+ return bin;
72
+}
73
+
74
+void chart_instance_updated_destroy(struct chart_instance_updated *instance)
75
+{
76
+ freez((char*)instance->id);
77
+ freez((char*)instance->claim_id);
78
+ freez((char*)instance->node_id);
79
+ freez((char*)instance->name);
80
+
81
+ free_label_list(instance->label_head);
82
+
83
+ freez((char*)instance->config_hash);
84
+}
85
+
86
+static int set_chart_instance_updated(chart::v1::ChartInstanceUpdated *chart, const struct chart_instance_updated *update)
87
+{
88
+ google::protobuf::Map<std::string, std::string> *map;
89
+ aclk_lib::v1::ACLKMessagePosition *pos;
90
+ struct label *label;
91
+
92
+ chart->set_id(update->id);
93
+ chart->set_claim_id(update->claim_id);
94
+ chart->set_node_id(update->node_id);
95
+ chart->set_name(update->name);
96
+
97
+ map = chart->mutable_chart_labels();
98
+ label = update->label_head;
99
+ while (label) {
100
+ map->insert({label->key, label->value});
101
+ label = label->next;
102
+ }
103
+
104
+ switch (update->memory_mode) {
105
+ case RRD_MEMORY_MODE_NONE:
106
+ chart->set_memory_mode(chart::v1::NONE);
107
+ break;
108
+ case RRD_MEMORY_MODE_RAM:
109
+ chart->set_memory_mode(chart::v1::RAM);
110
+ break;
111
+ case RRD_MEMORY_MODE_MAP:
112
+ chart->set_memory_mode(chart::v1::MAP);
113
+ break;
114
+ case RRD_MEMORY_MODE_SAVE:
115
+ chart->set_memory_mode(chart::v1::SAVE);
116
+ break;
117
+ case RRD_MEMORY_MODE_ALLOC:
118
+ chart->set_memory_mode(chart::v1::ALLOC);
119
+ break;
120
+ case RRD_MEMORY_MODE_DBENGINE:
121
+ chart->set_memory_mode(chart::v1::DB_ENGINE);
122
+ break;
123
+ default:
124
+ return 1;
125
+ break;
126
+ }
127
+
128
+ chart->set_update_every_interval(update->update_every);
129
+ chart->set_config_hash(update->config_hash);
130
+
131
+ pos = chart->mutable_position();
132
+ pos->set_sequence_id(update->position.sequence_id);
133
+ pos->set_previous_sequence_id(update->position.previous_sequence_id);
134
+ set_google_timestamp_from_timeval(update->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
135
+
136
+ return 0;
137
+}
138
+
139
+static int set_chart_dim_updated(chart::v1::ChartDimensionUpdated *dim, const struct chart_dimension_updated *c_dim)
140
+{
141
+ aclk_lib::v1::ACLKMessagePosition *pos;
142
+
143
+ dim->set_id(c_dim->id);
144
+ dim->set_chart_id(c_dim->chart_id);
145
+ dim->set_node_id(c_dim->node_id);
146
+ dim->set_claim_id(c_dim->claim_id);
147
+ dim->set_name(c_dim->name);
148
+
149
+ set_google_timestamp_from_timeval(c_dim->created_at, dim->mutable_created_at());
150
+ set_google_timestamp_from_timeval(c_dim->last_timestamp, dim->mutable_last_timestamp());
151
+
152
+ pos = dim->mutable_position();
153
+ pos->set_sequence_id(c_dim->position.sequence_id);
154
+ pos->set_previous_sequence_id(c_dim->position.previous_sequence_id);
155
+ set_google_timestamp_from_timeval(c_dim->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
156
+
157
+ return 0;
158
+}
159
+
160
+char *generate_charts_and_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, int *is_dim, struct aclk_message_position *new_positions, uint64_t batch_id)
161
+{
162
+ chart::v1::ChartsAndDimensionsUpdated msg;
163
+ chart::v1::ChartInstanceUpdated db_chart;
164
+ chart::v1::ChartDimensionUpdated db_dim;
165
+ aclk_lib::v1::ACLKMessagePosition *pos;
166
+
167
+ msg.set_batch_id(batch_id);
168
+
169
+ for (int i = 0; payloads[i]; i++) {
170
+ if (is_dim[i]) {
171
+ if (!db_dim.ParseFromArray(payloads[i], payload_sizes[i])) {
172
+ error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
173
+ return NULL;
174
+ }
175
+
176
+ pos = db_dim.mutable_position();
177
+ pos->set_sequence_id(new_positions[i].sequence_id);
178
+ pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
179
+ set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
180
+
181
+ chart::v1::ChartDimensionUpdated *dim = msg.add_dimensions();
182
+ *dim = db_dim;
183
+ } else {
184
+ if (!db_chart.ParseFromArray(payloads[i], payload_sizes[i])) {
185
+ error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
186
+ return NULL;
187
+ }
188
+
189
+ pos = db_chart.mutable_position();
190
+ pos->set_sequence_id(new_positions[i].sequence_id);
191
+ pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
192
+ set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
193
+
194
+ chart::v1::ChartInstanceUpdated *chart = msg.add_charts();
195
+ *chart = db_chart;
196
+ }
197
+ }
198
+
199
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
200
+ char *bin = (char*)mallocz(*len);
201
+ msg.SerializeToArray(bin, *len);
202
+
203
+ return bin;
204
+}
205
+
206
+char *generate_charts_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
207
+{
208
+ chart::v1::ChartsAndDimensionsUpdated msg;
209
+
210
+ msg.set_batch_id(chart_batch_id);
211
+
212
+ for (int i = 0; payloads[i]; i++) {
213
+ chart::v1::ChartInstanceUpdated db_msg;
214
+ chart::v1::ChartInstanceUpdated *chart;
215
+ aclk_lib::v1::ACLKMessagePosition *pos;
216
+
217
+ if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
218
+ error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
219
+ return NULL;
220
+ }
221
+
222
+ pos = db_msg.mutable_position();
223
+ pos->set_sequence_id(new_positions[i].sequence_id);
224
+ pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
225
+ set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
226
+
227
+ chart = msg.add_charts();
228
+ *chart = db_msg;
229
+ }
230
+
231
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
232
+ char *bin = (char*)mallocz(*len);
233
+ msg.SerializeToArray(bin, *len);
234
+
235
+ return bin;
236
+}
237
+
238
+char *generate_chart_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
239
+{
240
+ chart::v1::ChartsAndDimensionsUpdated msg;
241
+
242
+ msg.set_batch_id(chart_batch_id);
243
+
244
+ for (int i = 0; payloads[i]; i++) {
245
+ chart::v1::ChartDimensionUpdated db_msg;
246
+ chart::v1::ChartDimensionUpdated *dim;
247
+ aclk_lib::v1::ACLKMessagePosition *pos;
248
+
249
+ if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
250
+ error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
251
+ return NULL;
252
+ }
253
+
254
+ pos = db_msg.mutable_position();
255
+ pos->set_sequence_id(new_positions[i].sequence_id);
256
+ pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
257
+ set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
258
+
259
+ dim = msg.add_dimensions();
260
+ *dim = db_msg;
261
+ }
262
+
263
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
264
+ char *bin = (char*)mallocz(*len);
265
+ msg.SerializeToArray(bin, *len);
266
+
267
+ return bin;
268
+}
269
+
270
+char *generate_chart_instance_updated(size_t *len, const struct chart_instance_updated *update)
271
+{
272
+ chart::v1::ChartInstanceUpdated *chart = new chart::v1::ChartInstanceUpdated();
273
+
274
+ if (set_chart_instance_updated(chart, update))
275
+ return NULL;
276
+
277
+ *len = PROTO_COMPAT_MSG_SIZE_PTR(chart);
278
+ char *bin = (char*)mallocz(*len);
279
+ chart->SerializeToArray(bin, *len);
280
+
281
+ delete chart;
282
+ return bin;
283
+}
284
+
285
+char *generate_chart_dimension_updated(size_t *len, const struct chart_dimension_updated *dim)
286
+{
287
+ chart::v1::ChartDimensionUpdated *proto_dim = new chart::v1::ChartDimensionUpdated();
288
+
289
+ if (set_chart_dim_updated(proto_dim, dim))
290
+ return NULL;
291
+
292
+ *len = PROTO_COMPAT_MSG_SIZE_PTR(proto_dim);
293
+ char *bin = (char*)mallocz(*len);
294
+ proto_dim->SerializeToArray(bin, *len);
295
+
296
+ delete proto_dim;
297
+ return bin;
298
+}
299
+
300
+using namespace google::protobuf;
301
+
302
+char *generate_retention_updated(size_t *len, struct retention_updated *data)
303
+{
304
+ chart::v1::RetentionUpdated msg;
305
+
306
+ msg.set_claim_id(data->claim_id);
307
+ msg.set_node_id(data->node_id);
308
+
309
+ switch (data->memory_mode) {
310
+ case RRD_MEMORY_MODE_NONE:
311
+ msg.set_memory_mode(chart::v1::NONE);
312
+ break;
313
+ case RRD_MEMORY_MODE_RAM:
314
+ msg.set_memory_mode(chart::v1::RAM);
315
+ break;
316
+ case RRD_MEMORY_MODE_MAP:
317
+ msg.set_memory_mode(chart::v1::MAP);
318
+ break;
319
+ case RRD_MEMORY_MODE_SAVE:
320
+ msg.set_memory_mode(chart::v1::SAVE);
321
+ break;
322
+ case RRD_MEMORY_MODE_ALLOC:
323
+ msg.set_memory_mode(chart::v1::ALLOC);
324
+ break;
325
+ case RRD_MEMORY_MODE_DBENGINE:
326
+ msg.set_memory_mode(chart::v1::DB_ENGINE);
327
+ break;
328
+ default:
329
+ return NULL;
330
+ }
331
+
332
+ for (int i = 0; i < data->interval_duration_count; i++) {
333
+ Map<uint32, uint32> *map = msg.mutable_interval_durations();
334
+ map->insert({data->interval_durations[i].update_every, data->interval_durations[i].retention});
335
+ }
336
+
337
+ set_google_timestamp_from_timeval(data->rotation_timestamp, msg.mutable_rotation_timestamp());
338
+
339
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
340
+ char *bin = (char*)mallocz(*len);
341
+ msg.SerializeToArray(bin, *len);
342
+
343
+ return bin;
344
+}
aclk/schema-wrappers/chart_stream.h
new
+121
@@ -0,0 +1,121 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+
3
+#ifndef ACLK_SCHEMA_WRAPPER_CHART_STREAM_H
4
+#define ACLK_SCHEMA_WRAPPER_CHART_STREAM_H
5
+
6
+#ifdef __cplusplus
7
+extern "C" {
8
+#endif
9
+
10
+#include "database/rrd.h"
11
+
12
+typedef struct {
13
+ char* claim_id;
14
+ char* node_id;
15
+
16
+ uint64_t seq_id;
17
+ uint64_t batch_id;
18
+
19
+ struct timeval seq_id_created_at;
20
+} stream_charts_and_dims_t;
21
+
22
+stream_charts_and_dims_t parse_stream_charts_and_dims(const char *data, size_t len);
23
+
24
+typedef struct {
25
+ char* claim_id;
26
+ char* node_id;
27
+
28
+ uint64_t last_seq_id;
29
+} chart_and_dim_ack_t;
30
+
31
+chart_and_dim_ack_t parse_chart_and_dimensions_ack(const char *data, size_t len);
32
+
33
+enum chart_reset_reason {
34
+ DB_EMPTY,
35
+ SEQ_ID_NOT_EXISTS,
36
+ TIMESTAMP_MISMATCH
37
+};
38
+
39
+typedef struct {
40
+ char *claim_id;
41
+ char *node_id;
42
+
43
+ enum chart_reset_reason reason;
44
+} chart_reset_t;
45
+
46
+char *generate_reset_chart_messages(size_t *len, const chart_reset_t reset);
47
+
48
+struct aclk_message_position {
49
+ uint64_t sequence_id;
50
+ struct timeval seq_id_creation_time;
51
+ uint64_t previous_sequence_id;
52
+};
53
+
54
+struct chart_instance_updated {
55
+ const char *id;
56
+ const char *claim_id;
57
+ const char *node_id;
58
+ const char *name;
59
+
60
+ struct label *label_head;
61
+
62
+ RRD_MEMORY_MODE memory_mode;
63
+
64
+ uint32_t update_every;
65
+ const char * config_hash;
66
+
67
+ struct aclk_message_position position;
68
+};
69
+
70
+void chart_instance_updated_destroy(struct chart_instance_updated *instance);
71
+
72
+struct chart_dimension_updated {
73
+ const char *id;
74
+ const char *chart_id;
75
+ const char *node_id;
76
+ const char *claim_id;
77
+ const char *name;
78
+ struct timeval created_at;
79
+ struct timeval last_timestamp;
80
+ struct aclk_message_position position;
81
+};
82
+
83
+typedef struct {
84
+ struct chart_instance_updated *charts;
85
+ uint16_t chart_count;
86
+
87
+ struct chart_dimension_updated *dims;
88
+ uint16_t dim_count;
89
+
90
+ uint64_t batch_id;
91
+} charts_and_dims_updated_t;
92
+
93
+struct interval_duration {
94
+ uint32_t update_every;
95
+ uint32_t retention;
96
+};
97
+
98
+struct retention_updated {
99
+ char *claim_id;
100
+ char *node_id;
101
+
102
+ RRD_MEMORY_MODE memory_mode;
103
+
104
+ struct interval_duration *interval_durations;
105
+ int interval_duration_count;
106
+
107
+ struct timeval rotation_timestamp;
108
+};
109
+
110
+char *generate_charts_and_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, int *is_dim, struct aclk_message_position *new_positions, uint64_t batch_id);
111
+char *generate_charts_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions);
112
+char *generate_chart_instance_updated(size_t *len, const struct chart_instance_updated *update);
113
+char *generate_chart_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions);
114
+char *generate_chart_dimension_updated(size_t *len, const struct chart_dimension_updated *dim);
115
+char *generate_retention_updated(size_t *len, struct retention_updated *data);
116
+
117
+#ifdef __cplusplus
118
+}
119
+#endif
120
+
121
+#endif /* ACLK_SCHEMA_WRAPPER_CHART_STREAM_H */
aclk/schema-wrappers/schema_wrapper_utils.cc
new
+15
@@ -0,0 +1,15 @@
1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+
3
+#include "schema_wrapper_utils.h"
4
+
5
+void set_google_timestamp_from_timeval(struct timeval tv, google::protobuf::Timestamp *ts)
6
+{
7
+ ts->set_nanos(tv.tv_usec*1000);
8
+ ts->set_seconds(tv.tv_sec);
9
+}
10
+
11
+void set_timeval_from_google_timestamp(const google::protobuf::Timestamp &ts, struct timeval *tv)
12
+{
13
+ tv->tv_sec = ts.seconds();
14
+ tv->tv_usec = ts.nanos()/1000;
15
+}
aclk/schema-wrappers/schema_wrapper_utils.h
+7
@@ -3,10 +3,17 @@
3
#ifndef SCHEMA_WRAPPER_UTILS_H
4
#define SCHEMA_WRAPPER_UTILS_H
5
6
+#include <google/protobuf/timestamp.pb.h>
7
+
8
#if GOOGLE_PROTOBUF_VERSION < 3001000
9
#define PROTO_COMPAT_MSG_SIZE(msg) (size_t)msg.ByteSize();
10
+#define PROTO_COMPAT_MSG_SIZE_PTR(msg) (size_t)msg->ByteSize();
11
#else
12
#define PROTO_COMPAT_MSG_SIZE(msg) msg.ByteSizeLong();
13
+#define PROTO_COMPAT_MSG_SIZE_PTR(msg) msg->ByteSizeLong();
14
#endif
15
16
+void set_google_timestamp_from_timeval(struct timeval tv, google::protobuf::Timestamp *ts);
17
+void set_timeval_from_google_timestamp(const google::protobuf::Timestamp &ts, struct timeval *tv);
18
+
19
#endif /* SCHEMA_WRAPPER_UTILS_H */
aclk/schema-wrappers/schema_wrappers.h
+2
@@ -8,5 +8,7 @@
8
#include "connection.h"
9
#include "node_connection.h"
10
#include "node_creation.h"
11
+#include "chart_config.h"
12
+#include "chart_stream.h"
13
14
#endif /* SCHEMA_WRAPPERS_H */