@cryptotaxi247 / netdata-1 / commits / 0cb8e0cee

reduce memcpy and memory usage on mqtt5 (#13450)

* reduce memcpy on mqtt5

Timotej S committed Aug 15, 2022 at 11:48 UTC 0cb8e0cee1dba7cfe54a62e2b6ad619707eedd84
1 file changed +11 -34
aclk/aclk_tx_msgs.c
+11 -34
@@ -43,28 +43,12 @@ uint16_t aclk_send_bin_message_subtopic_pid(mqtt_wss_client client, char *msg, s
43 return packet_id;
44 }
45
46 -/* UNUSED now but can be used soon MVP1?
47 -static void aclk_send_message_topic(mqtt_wss_client client, json_object *msg, const char *topic)
46 +// json_object_put returns int unfortunately :D
47 +// we need void(*fnc)(void *);
48 +static void json_object_put_wrapper(void *jsonobj)
49 {
49 - if (unlikely(!topic || topic[0] != '/')) {
50 - error ("Full topic required!");
51 - return;
52 - }
53 -
54 - const char *str = json_object_to_json_string_ext(msg, JSON_C_TO_STRING_PLAIN);
55 -
56 - mqtt_wss_publish(client, topic, str, strlen(str), MQTT_WSS_PUB_QOS1);
57 -#ifdef NETDATA_INTERNAL_CHECKS
58 - aclk_stats_msg_published();
59 -#endif
60 -#ifdef ACLK_LOG_CONVERSATION_DIR
61 -#define FN_MAX_LEN 1024
62 - char filename[FN_MAX_LEN];
63 - snprintf(filename, FN_MAX_LEN, ACLK_LOG_CONVERSATION_DIR "/%010d-tx.json", ACLK_GET_CONV_LOG_NEXT());
64 - json_object_to_file_ext(filename, msg, JSON_C_TO_STRING_PRETTY);
65 -#endif
50 + json_object_put(jsonobj);
51 }
67 -*/
52
53 #define TOPIC_MAX_LEN 512
54 #define V2_BIN_PAYLOAD_SEPARATOR "\x0D\x0A\x0D\x0A"
@@ -77,6 +61,7 @@ static int aclk_send_message_with_bin_payload(mqtt_wss_client client, json_objec
61
62 if (unlikely(!topic || topic[0] != '/')) {
63 error ("Full topic required!");
64 + json_object_put(msg);
65 return HTTP_RESP_INTERNAL_SERVER_ERROR;
66 }
67
@@ -87,32 +72,26 @@ static int aclk_send_message_with_bin_payload(mqtt_wss_client client, json_objec
72 full_msg = mallocz(len + strlen(V2_BIN_PAYLOAD_SEPARATOR) + payload_len);
73
74 memcpy(full_msg, str, len);
75 + json_object_put(msg);
76 + msg = NULL;
77 memcpy(&full_msg[len], V2_BIN_PAYLOAD_SEPARATOR, strlen(V2_BIN_PAYLOAD_SEPARATOR));
78 len += strlen(V2_BIN_PAYLOAD_SEPARATOR);
79 memcpy(&full_msg[len], payload, payload_len);
80 len += payload_len;
81 }
82
96 -/* TODO
97 -#ifdef ACLK_LOG_CONVERSATION_DIR
98 -#define FN_MAX_LEN 1024
99 - char filename[FN_MAX_LEN];
100 - snprintf(filename, FN_MAX_LEN, ACLK_LOG_CONVERSATION_DIR "/%010d-tx.json", ACLK_GET_CONV_LOG_NEXT());
101 - json_object_to_file_ext(filename, msg, JSON_C_TO_STRING_PRETTY);
102 -#endif */
103 -
83 if (use_mqtt_5)
105 - mqtt_wss_publish5(client, (char*)topic, NULL, (char*)(payload_len ? full_msg : str), NULL, len, MQTT_WSS_PUB_QOS1, &packet_id);
84 + mqtt_wss_publish5(client, (char*)topic, NULL, (char*)(payload_len ? full_msg : str), (payload_len ? &freez : &json_object_put_wrapper), len, MQTT_WSS_PUB_QOS1, &packet_id);
85 else {
86 rc = mqtt_wss_publish_pid_block(client, topic, payload_len ? full_msg : str, len, MQTT_WSS_PUB_QOS1, &packet_id, 5000);
87 + freez(full_msg);
88 + json_object_put(msg);
89 if (rc == MQTT_WSS_ERR_BLOCK_TIMEOUT) {
90 error("Timeout sending binpacked message");
110 - freez(full_msg);
91 return HTTP_RESP_BACKEND_FETCH_FAILED;
92 }
93 if (rc == MQTT_WSS_ERR_TX_BUF_TOO_SMALL) {
94 error("Message is bigger than allowed maximum");
115 - freez(full_msg);
95 return HTTP_RESP_FORBIDDEN;
96 }
97 }
@@ -120,7 +99,7 @@ static int aclk_send_message_with_bin_payload(mqtt_wss_client client, json_objec
99 #ifdef NETDATA_INTERNAL_CHECKS
100 aclk_stats_msg_published(packet_id);
101 #endif
123 - freez(full_msg);
102 +
103 return 0;
104 }
105
@@ -203,7 +182,6 @@ void aclk_http_msg_v2_err(mqtt_wss_client client, const char *topic, const char
182 if (aclk_send_message_with_bin_payload(client, msg, topic, payload, payload_len)) {
183 error("Failed to send cancelation message for http reply");
184 }
206 - json_object_put(msg);
185 }
186
187 void aclk_http_msg_v2(mqtt_wss_client client, const char *topic, const char *msg_id, usec_t t_exec, usec_t created, int http_code, const char *payload, size_t payload_len)
@@ -222,7 +200,6 @@ void aclk_http_msg_v2(mqtt_wss_client client, const char *topic, const char *msg
200 json_object_object_add(msg, "http-code", tmp);
201
202 int rc = aclk_send_message_with_bin_payload(client, msg, topic, payload, payload_len);
225 - json_object_put(msg);
203
204 switch (rc) {
205 case HTTP_RESP_FORBIDDEN: