chart_stream.cc 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337
  1. // SPDX-License-Identifier: GPL-3.0-or-later
  2. #include "aclk/aclk_util.h"
  3. #include "proto/chart/v1/stream.pb.h"
  4. #include "chart_stream.h"
  5. #include "schema_wrapper_utils.h"
  6. #include <sys/time.h>
  7. #include <stdlib.h>
  8. stream_charts_and_dims_t parse_stream_charts_and_dims(const char *data, size_t len)
  9. {
  10. chart::v1::StreamChartsAndDimensions msg;
  11. stream_charts_and_dims_t res;
  12. memset(&res, 0, sizeof(res));
  13. if (!msg.ParseFromArray(data, len))
  14. return res;
  15. res.node_id = strdup(msg.node_id().c_str());
  16. res.claim_id = strdup(msg.claim_id().c_str());
  17. res.seq_id = msg.sequence_id();
  18. res.batch_id = msg.batch_id();
  19. set_timeval_from_google_timestamp(msg.seq_id_created_at(), &res.seq_id_created_at);
  20. return res;
  21. }
  22. chart_and_dim_ack_t parse_chart_and_dimensions_ack(const char *data, size_t len)
  23. {
  24. chart::v1::ChartsAndDimensionsAck msg;
  25. chart_and_dim_ack_t res = { .claim_id = NULL, .node_id = NULL, .last_seq_id = 0 };
  26. if (!msg.ParseFromArray(data, len))
  27. return res;
  28. res.node_id = strdup(msg.node_id().c_str());
  29. res.claim_id = strdup(msg.claim_id().c_str());
  30. res.last_seq_id = msg.last_sequence_id();
  31. return res;
  32. }
  33. char *generate_reset_chart_messages(size_t *len, chart_reset_t reset)
  34. {
  35. chart::v1::ResetChartMessages msg;
  36. msg.set_claim_id(reset.claim_id);
  37. msg.set_node_id(reset.node_id);
  38. switch (reset.reason) {
  39. case DB_EMPTY:
  40. msg.set_reason(chart::v1::ResetReason::DB_EMPTY);
  41. break;
  42. case SEQ_ID_NOT_EXISTS:
  43. msg.set_reason(chart::v1::ResetReason::SEQ_ID_NOT_EXISTS);
  44. break;
  45. case TIMESTAMP_MISMATCH:
  46. msg.set_reason(chart::v1::ResetReason::TIMESTAMP_MISMATCH);
  47. break;
  48. default:
  49. return NULL;
  50. }
  51. *len = PROTO_COMPAT_MSG_SIZE(msg);
  52. char *bin = (char*)malloc(*len);
  53. if (bin)
  54. msg.SerializeToArray(bin, *len);
  55. return bin;
  56. }
  57. void chart_instance_updated_destroy(struct chart_instance_updated *instance)
  58. {
  59. freez((char*)instance->id);
  60. freez((char*)instance->claim_id);
  61. rrdlabels_destroy(instance->chart_labels);
  62. freez((char*)instance->config_hash);
  63. }
  64. static int set_chart_instance_updated(chart::v1::ChartInstanceUpdated *chart, const struct chart_instance_updated *update)
  65. {
  66. google::protobuf::Map<std::string, std::string> *map;
  67. aclk_lib::v1::ACLKMessagePosition *pos;
  68. chart->set_id(update->id);
  69. chart->set_claim_id(update->claim_id);
  70. chart->set_node_id(update->node_id);
  71. chart->set_name(update->name);
  72. map = chart->mutable_chart_labels();
  73. rrdlabels_walkthrough_read(update->chart_labels, label_add_to_map_callback, map);
  74. switch (update->memory_mode) {
  75. case RRD_MEMORY_MODE_NONE:
  76. chart->set_memory_mode(chart::v1::NONE);
  77. break;
  78. case RRD_MEMORY_MODE_RAM:
  79. chart->set_memory_mode(chart::v1::RAM);
  80. break;
  81. case RRD_MEMORY_MODE_MAP:
  82. chart->set_memory_mode(chart::v1::MAP);
  83. break;
  84. case RRD_MEMORY_MODE_SAVE:
  85. chart->set_memory_mode(chart::v1::SAVE);
  86. break;
  87. case RRD_MEMORY_MODE_ALLOC:
  88. chart->set_memory_mode(chart::v1::ALLOC);
  89. break;
  90. case RRD_MEMORY_MODE_DBENGINE:
  91. chart->set_memory_mode(chart::v1::DB_ENGINE);
  92. break;
  93. default:
  94. return 1;
  95. break;
  96. }
  97. chart->set_update_every_interval(update->update_every);
  98. chart->set_config_hash(update->config_hash);
  99. pos = chart->mutable_position();
  100. pos->set_sequence_id(update->position.sequence_id);
  101. pos->set_previous_sequence_id(update->position.previous_sequence_id);
  102. set_google_timestamp_from_timeval(update->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
  103. return 0;
  104. }
  105. static int set_chart_dim_updated(chart::v1::ChartDimensionUpdated *dim, const struct chart_dimension_updated *c_dim)
  106. {
  107. aclk_lib::v1::ACLKMessagePosition *pos;
  108. dim->set_id(c_dim->id);
  109. dim->set_chart_id(c_dim->chart_id);
  110. dim->set_node_id(c_dim->node_id);
  111. dim->set_claim_id(c_dim->claim_id);
  112. dim->set_name(c_dim->name);
  113. set_google_timestamp_from_timeval(c_dim->created_at, dim->mutable_created_at());
  114. set_google_timestamp_from_timeval(c_dim->last_timestamp, dim->mutable_last_timestamp());
  115. pos = dim->mutable_position();
  116. pos->set_sequence_id(c_dim->position.sequence_id);
  117. pos->set_previous_sequence_id(c_dim->position.previous_sequence_id);
  118. set_google_timestamp_from_timeval(c_dim->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
  119. return 0;
  120. }
  121. 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)
  122. {
  123. chart::v1::ChartsAndDimensionsUpdated msg;
  124. chart::v1::ChartInstanceUpdated db_chart;
  125. chart::v1::ChartDimensionUpdated db_dim;
  126. aclk_lib::v1::ACLKMessagePosition *pos;
  127. msg.set_batch_id(batch_id);
  128. for (int i = 0; payloads[i]; i++) {
  129. if (is_dim[i]) {
  130. if (!db_dim.ParseFromArray(payloads[i], payload_sizes[i])) {
  131. error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
  132. return NULL;
  133. }
  134. pos = db_dim.mutable_position();
  135. pos->set_sequence_id(new_positions[i].sequence_id);
  136. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  137. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  138. chart::v1::ChartDimensionUpdated *dim = msg.add_dimensions();
  139. *dim = db_dim;
  140. } else {
  141. if (!db_chart.ParseFromArray(payloads[i], payload_sizes[i])) {
  142. error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
  143. return NULL;
  144. }
  145. pos = db_chart.mutable_position();
  146. pos->set_sequence_id(new_positions[i].sequence_id);
  147. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  148. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  149. chart::v1::ChartInstanceUpdated *chart = msg.add_charts();
  150. *chart = db_chart;
  151. }
  152. }
  153. *len = PROTO_COMPAT_MSG_SIZE(msg);
  154. char *bin = (char*)mallocz(*len);
  155. msg.SerializeToArray(bin, *len);
  156. return bin;
  157. }
  158. char *generate_charts_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
  159. {
  160. chart::v1::ChartsAndDimensionsUpdated msg;
  161. msg.set_batch_id(chart_batch_id);
  162. for (int i = 0; payloads[i]; i++) {
  163. chart::v1::ChartInstanceUpdated db_msg;
  164. chart::v1::ChartInstanceUpdated *chart;
  165. aclk_lib::v1::ACLKMessagePosition *pos;
  166. if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
  167. error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
  168. return NULL;
  169. }
  170. pos = db_msg.mutable_position();
  171. pos->set_sequence_id(new_positions[i].sequence_id);
  172. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  173. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  174. chart = msg.add_charts();
  175. *chart = db_msg;
  176. }
  177. *len = PROTO_COMPAT_MSG_SIZE(msg);
  178. char *bin = (char*)mallocz(*len);
  179. msg.SerializeToArray(bin, *len);
  180. return bin;
  181. }
  182. char *generate_chart_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
  183. {
  184. chart::v1::ChartsAndDimensionsUpdated msg;
  185. msg.set_batch_id(chart_batch_id);
  186. for (int i = 0; payloads[i]; i++) {
  187. chart::v1::ChartDimensionUpdated db_msg;
  188. chart::v1::ChartDimensionUpdated *dim;
  189. aclk_lib::v1::ACLKMessagePosition *pos;
  190. if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
  191. error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
  192. return NULL;
  193. }
  194. pos = db_msg.mutable_position();
  195. pos->set_sequence_id(new_positions[i].sequence_id);
  196. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  197. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  198. dim = msg.add_dimensions();
  199. *dim = db_msg;
  200. }
  201. *len = PROTO_COMPAT_MSG_SIZE(msg);
  202. char *bin = (char*)mallocz(*len);
  203. msg.SerializeToArray(bin, *len);
  204. return bin;
  205. }
  206. char *generate_chart_instance_updated(size_t *len, const struct chart_instance_updated *update)
  207. {
  208. chart::v1::ChartInstanceUpdated *chart = new chart::v1::ChartInstanceUpdated();
  209. if (set_chart_instance_updated(chart, update))
  210. return NULL;
  211. *len = PROTO_COMPAT_MSG_SIZE_PTR(chart);
  212. char *bin = (char*)mallocz(*len);
  213. chart->SerializeToArray(bin, *len);
  214. delete chart;
  215. return bin;
  216. }
  217. char *generate_chart_dimension_updated(size_t *len, const struct chart_dimension_updated *dim)
  218. {
  219. chart::v1::ChartDimensionUpdated *proto_dim = new chart::v1::ChartDimensionUpdated();
  220. if (set_chart_dim_updated(proto_dim, dim))
  221. return NULL;
  222. *len = PROTO_COMPAT_MSG_SIZE_PTR(proto_dim);
  223. char *bin = (char*)mallocz(*len);
  224. proto_dim->SerializeToArray(bin, *len);
  225. delete proto_dim;
  226. return bin;
  227. }
  228. using namespace google::protobuf;
  229. char *generate_retention_updated(size_t *len, struct retention_updated *data)
  230. {
  231. chart::v1::RetentionUpdated msg;
  232. msg.set_claim_id(data->claim_id);
  233. msg.set_node_id(data->node_id);
  234. switch (data->memory_mode) {
  235. case RRD_MEMORY_MODE_NONE:
  236. msg.set_memory_mode(chart::v1::NONE);
  237. break;
  238. case RRD_MEMORY_MODE_RAM:
  239. msg.set_memory_mode(chart::v1::RAM);
  240. break;
  241. case RRD_MEMORY_MODE_MAP:
  242. msg.set_memory_mode(chart::v1::MAP);
  243. break;
  244. case RRD_MEMORY_MODE_SAVE:
  245. msg.set_memory_mode(chart::v1::SAVE);
  246. break;
  247. case RRD_MEMORY_MODE_ALLOC:
  248. msg.set_memory_mode(chart::v1::ALLOC);
  249. break;
  250. case RRD_MEMORY_MODE_DBENGINE:
  251. msg.set_memory_mode(chart::v1::DB_ENGINE);
  252. break;
  253. default:
  254. return NULL;
  255. }
  256. for (int i = 0; i < data->interval_duration_count; i++) {
  257. Map<uint32, uint32> *map = msg.mutable_interval_durations();
  258. map->insert({data->interval_durations[i].update_every, data->interval_durations[i].retention});
  259. }
  260. set_google_timestamp_from_timeval(data->rotation_timestamp, msg.mutable_rotation_timestamp());
  261. *len = PROTO_COMPAT_MSG_SIZE(msg);
  262. char *bin = (char*)mallocz(*len);
  263. msg.SerializeToArray(bin, *len);
  264. return bin;
  265. }