@cryptotaxi247 / netdata-1 / commits / d477f0a44

inform cloud about inability to deliver data payload (#12041)

Timotej S committed Mar 1, 2022 at 10:02 UTC d477f0a4469b4036866c899dc9c551661df72ecb
1 file changed +33 -14
aclk/aclk_tx_msgs.c
+33 -14
@@ -116,28 +116,30 @@ static void aclk_send_message_topic(mqtt_wss_client client, json_object *msg, co
116
117 #define TOPIC_MAX_LEN 512
118 #define V2_BIN_PAYLOAD_SEPARATOR "\x0D\x0A\x0D\x0A"
119 -static void aclk_send_message_with_bin_payload(mqtt_wss_client client, json_object *msg, const char *topic, const void *payload, size_t payload_len)
119 +static int aclk_send_message_with_bin_payload(mqtt_wss_client client, json_object *msg, const char *topic, const void *payload, size_t payload_len)
120 {
121 uint16_t packet_id;
122 const char *str;
123 - char *full_msg;
124 - int len;
123 + char *full_msg = NULL;
124 + int len, rc;
125
126 if (unlikely(!topic || topic[0] != '/')) {
127 error ("Full topic required!");
128 - return;
128 + return 500;
129 }
130
131 str = json_object_to_json_string_ext(msg, JSON_C_TO_STRING_PLAIN);
132 len = strlen(str);
133
134 - full_msg = mallocz(len + strlen(V2_BIN_PAYLOAD_SEPARATOR) + payload_len);
134 + if (payload_len) {
135 + full_msg = mallocz(len + strlen(V2_BIN_PAYLOAD_SEPARATOR) + payload_len);
136
136 - memcpy(full_msg, str, len);
137 - memcpy(&full_msg[len], V2_BIN_PAYLOAD_SEPARATOR, strlen(V2_BIN_PAYLOAD_SEPARATOR));
138 - len += strlen(V2_BIN_PAYLOAD_SEPARATOR);
139 - memcpy(&full_msg[len], payload, payload_len);
140 - len += payload_len;
137 + memcpy(full_msg, str, len);
138 + memcpy(&full_msg[len], V2_BIN_PAYLOAD_SEPARATOR, strlen(V2_BIN_PAYLOAD_SEPARATOR));
139 + len += strlen(V2_BIN_PAYLOAD_SEPARATOR);
140 + memcpy(&full_msg[len], payload, payload_len);
141 + len += payload_len;
142 + }
143
144 /* TODO
145 #ifdef ACLK_LOG_CONVERSATION_DIR
@@ -147,15 +149,22 @@ static void aclk_send_message_with_bin_payload(mqtt_wss_client client, json_obje
149 json_object_to_file_ext(filename, msg, JSON_C_TO_STRING_PRETTY);
150 #endif */
151
150 - int rc = mqtt_wss_publish_pid_block(client, topic, full_msg, len, MQTT_WSS_PUB_QOS1, &packet_id, 5000);
151 - if (rc == MQTT_WSS_ERR_BLOCK_TIMEOUT)
152 + rc = mqtt_wss_publish_pid_block(client, topic, payload_len ? full_msg : str, len, MQTT_WSS_PUB_QOS1, &packet_id, 5000);
153 + if (rc == MQTT_WSS_ERR_BLOCK_TIMEOUT) {
154 error("Timeout sending binpacked message");
153 - if (rc == MQTT_WSS_ERR_TX_BUF_TOO_SMALL)
155 + freez(full_msg);
156 + return 503;
157 + }
158 + if (rc == MQTT_WSS_ERR_TX_BUF_TOO_SMALL) {
159 error("Message is bigger than allowed maximum");
160 + freez(full_msg);
161 + return 403;
162 + }
163 #ifdef NETDATA_INTERNAL_CHECKS
164 aclk_stats_msg_published(packet_id);
165 #endif
166 freez(full_msg);
167 + return 0;
168 }
169
170 /*
@@ -331,8 +340,18 @@ void aclk_http_msg_v2(mqtt_wss_client client, const char *topic, const char *msg
340 tmp = json_object_new_int(http_code);
341 json_object_object_add(msg, "http-code", tmp);
342
334 - aclk_send_message_with_bin_payload(client, msg, topic, payload, payload_len);
343 + int rc = aclk_send_message_with_bin_payload(client, msg, topic, payload, payload_len);
344 json_object_put(msg);
345 +
346 + if (rc) {
347 + msg = create_hdr("http", msg_id, 0, 0, 2);
348 + tmp = json_object_new_int(rc);
349 + json_object_object_add(msg, "http-code", tmp);
350 + if (aclk_send_message_with_bin_payload(client, msg, topic, payload, payload_len)) {
351 + error("Failed to send cancelation message for http reply");
352 + }
353 + json_object_put(msg);
354 + }
355 }
356
357 void aclk_chart_msg(mqtt_wss_client client, RRDHOST *host, const char *chart)