chart_stream.cc 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342
  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. free_label_list(instance->label_head);
  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. struct label *label;
  69. chart->set_id(update->id);
  70. chart->set_claim_id(update->claim_id);
  71. chart->set_node_id(update->node_id);
  72. chart->set_name(update->name);
  73. map = chart->mutable_chart_labels();
  74. label = update->label_head;
  75. while (label) {
  76. map->insert({label->key, label->value});
  77. label = label->next;
  78. }
  79. switch (update->memory_mode) {
  80. case RRD_MEMORY_MODE_NONE:
  81. chart->set_memory_mode(chart::v1::NONE);
  82. break;
  83. case RRD_MEMORY_MODE_RAM:
  84. chart->set_memory_mode(chart::v1::RAM);
  85. break;
  86. case RRD_MEMORY_MODE_MAP:
  87. chart->set_memory_mode(chart::v1::MAP);
  88. break;
  89. case RRD_MEMORY_MODE_SAVE:
  90. chart->set_memory_mode(chart::v1::SAVE);
  91. break;
  92. case RRD_MEMORY_MODE_ALLOC:
  93. chart->set_memory_mode(chart::v1::ALLOC);
  94. break;
  95. case RRD_MEMORY_MODE_DBENGINE:
  96. chart->set_memory_mode(chart::v1::DB_ENGINE);
  97. break;
  98. default:
  99. return 1;
  100. break;
  101. }
  102. chart->set_update_every_interval(update->update_every);
  103. chart->set_config_hash(update->config_hash);
  104. pos = chart->mutable_position();
  105. pos->set_sequence_id(update->position.sequence_id);
  106. pos->set_previous_sequence_id(update->position.previous_sequence_id);
  107. set_google_timestamp_from_timeval(update->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
  108. return 0;
  109. }
  110. static int set_chart_dim_updated(chart::v1::ChartDimensionUpdated *dim, const struct chart_dimension_updated *c_dim)
  111. {
  112. aclk_lib::v1::ACLKMessagePosition *pos;
  113. dim->set_id(c_dim->id);
  114. dim->set_chart_id(c_dim->chart_id);
  115. dim->set_node_id(c_dim->node_id);
  116. dim->set_claim_id(c_dim->claim_id);
  117. dim->set_name(c_dim->name);
  118. set_google_timestamp_from_timeval(c_dim->created_at, dim->mutable_created_at());
  119. set_google_timestamp_from_timeval(c_dim->last_timestamp, dim->mutable_last_timestamp());
  120. pos = dim->mutable_position();
  121. pos->set_sequence_id(c_dim->position.sequence_id);
  122. pos->set_previous_sequence_id(c_dim->position.previous_sequence_id);
  123. set_google_timestamp_from_timeval(c_dim->position.seq_id_creation_time, pos->mutable_seq_id_created_at());
  124. return 0;
  125. }
  126. 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)
  127. {
  128. chart::v1::ChartsAndDimensionsUpdated msg;
  129. chart::v1::ChartInstanceUpdated db_chart;
  130. chart::v1::ChartDimensionUpdated db_dim;
  131. aclk_lib::v1::ACLKMessagePosition *pos;
  132. msg.set_batch_id(batch_id);
  133. for (int i = 0; payloads[i]; i++) {
  134. if (is_dim[i]) {
  135. if (!db_dim.ParseFromArray(payloads[i], payload_sizes[i])) {
  136. error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
  137. return NULL;
  138. }
  139. pos = db_dim.mutable_position();
  140. pos->set_sequence_id(new_positions[i].sequence_id);
  141. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  142. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  143. chart::v1::ChartDimensionUpdated *dim = msg.add_dimensions();
  144. *dim = db_dim;
  145. } else {
  146. if (!db_chart.ParseFromArray(payloads[i], payload_sizes[i])) {
  147. error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
  148. return NULL;
  149. }
  150. pos = db_chart.mutable_position();
  151. pos->set_sequence_id(new_positions[i].sequence_id);
  152. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  153. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  154. chart::v1::ChartInstanceUpdated *chart = msg.add_charts();
  155. *chart = db_chart;
  156. }
  157. }
  158. *len = PROTO_COMPAT_MSG_SIZE(msg);
  159. char *bin = (char*)mallocz(*len);
  160. msg.SerializeToArray(bin, *len);
  161. return bin;
  162. }
  163. char *generate_charts_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
  164. {
  165. chart::v1::ChartsAndDimensionsUpdated msg;
  166. msg.set_batch_id(chart_batch_id);
  167. for (int i = 0; payloads[i]; i++) {
  168. chart::v1::ChartInstanceUpdated db_msg;
  169. chart::v1::ChartInstanceUpdated *chart;
  170. aclk_lib::v1::ACLKMessagePosition *pos;
  171. if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
  172. error("[ACLK] Could not parse chart::v1::ChartInstanceUpdated");
  173. return NULL;
  174. }
  175. pos = db_msg.mutable_position();
  176. pos->set_sequence_id(new_positions[i].sequence_id);
  177. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  178. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  179. chart = msg.add_charts();
  180. *chart = db_msg;
  181. }
  182. *len = PROTO_COMPAT_MSG_SIZE(msg);
  183. char *bin = (char*)mallocz(*len);
  184. msg.SerializeToArray(bin, *len);
  185. return bin;
  186. }
  187. char *generate_chart_dimensions_updated(size_t *len, char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
  188. {
  189. chart::v1::ChartsAndDimensionsUpdated msg;
  190. msg.set_batch_id(chart_batch_id);
  191. for (int i = 0; payloads[i]; i++) {
  192. chart::v1::ChartDimensionUpdated db_msg;
  193. chart::v1::ChartDimensionUpdated *dim;
  194. aclk_lib::v1::ACLKMessagePosition *pos;
  195. if (!db_msg.ParseFromArray(payloads[i], payload_sizes[i])) {
  196. error("[ACLK] Could not parse chart::v1::chart_dimension_updated");
  197. return NULL;
  198. }
  199. pos = db_msg.mutable_position();
  200. pos->set_sequence_id(new_positions[i].sequence_id);
  201. pos->set_previous_sequence_id(new_positions[i].previous_sequence_id);
  202. set_google_timestamp_from_timeval(new_positions[i].seq_id_creation_time, pos->mutable_seq_id_created_at());
  203. dim = msg.add_dimensions();
  204. *dim = db_msg;
  205. }
  206. *len = PROTO_COMPAT_MSG_SIZE(msg);
  207. char *bin = (char*)mallocz(*len);
  208. msg.SerializeToArray(bin, *len);
  209. return bin;
  210. }
  211. char *generate_chart_instance_updated(size_t *len, const struct chart_instance_updated *update)
  212. {
  213. chart::v1::ChartInstanceUpdated *chart = new chart::v1::ChartInstanceUpdated();
  214. if (set_chart_instance_updated(chart, update))
  215. return NULL;
  216. *len = PROTO_COMPAT_MSG_SIZE_PTR(chart);
  217. char *bin = (char*)mallocz(*len);
  218. chart->SerializeToArray(bin, *len);
  219. delete chart;
  220. return bin;
  221. }
  222. char *generate_chart_dimension_updated(size_t *len, const struct chart_dimension_updated *dim)
  223. {
  224. chart::v1::ChartDimensionUpdated *proto_dim = new chart::v1::ChartDimensionUpdated();
  225. if (set_chart_dim_updated(proto_dim, dim))
  226. return NULL;
  227. *len = PROTO_COMPAT_MSG_SIZE_PTR(proto_dim);
  228. char *bin = (char*)mallocz(*len);
  229. proto_dim->SerializeToArray(bin, *len);
  230. delete proto_dim;
  231. return bin;
  232. }
  233. using namespace google::protobuf;
  234. char *generate_retention_updated(size_t *len, struct retention_updated *data)
  235. {
  236. chart::v1::RetentionUpdated msg;
  237. msg.set_claim_id(data->claim_id);
  238. msg.set_node_id(data->node_id);
  239. switch (data->memory_mode) {
  240. case RRD_MEMORY_MODE_NONE:
  241. msg.set_memory_mode(chart::v1::NONE);
  242. break;
  243. case RRD_MEMORY_MODE_RAM:
  244. msg.set_memory_mode(chart::v1::RAM);
  245. break;
  246. case RRD_MEMORY_MODE_MAP:
  247. msg.set_memory_mode(chart::v1::MAP);
  248. break;
  249. case RRD_MEMORY_MODE_SAVE:
  250. msg.set_memory_mode(chart::v1::SAVE);
  251. break;
  252. case RRD_MEMORY_MODE_ALLOC:
  253. msg.set_memory_mode(chart::v1::ALLOC);
  254. break;
  255. case RRD_MEMORY_MODE_DBENGINE:
  256. msg.set_memory_mode(chart::v1::DB_ENGINE);
  257. break;
  258. default:
  259. return NULL;
  260. }
  261. for (int i = 0; i < data->interval_duration_count; i++) {
  262. Map<uint32, uint32> *map = msg.mutable_interval_durations();
  263. map->insert({data->interval_durations[i].update_every, data->interval_durations[i].retention});
  264. }
  265. set_google_timestamp_from_timeval(data->rotation_timestamp, msg.mutable_rotation_timestamp());
  266. *len = PROTO_COMPAT_MSG_SIZE(msg);
  267. char *bin = (char*)mallocz(*len);
  268. msg.SerializeToArray(bin, *len);
  269. return bin;
  270. }